Streaming SSH Session Output
SSH streaming from the netcli/scrapli/asyncssh transports requires hegemony-step-plugins 0.3.0 or newer, which the host's version floors and the demo stack's pinned wheels both require.
This is the design record for streaming SSH step output into the Run Detail Logs tab, alongside the container.run output that tab has always shown. It covers SSH-based steps (netcli.*, shell.execute, cisco.iosxe.upgrade.* - anything using an SSH transport), the phased plan that got there, and what each phase already delivers. Start with the status note above for what is live today.
Verdict: feasible, with no new infrastructure. The event pipeline that feeds the Logs tab is already handler-agnostic; "container output" is nothing more than kind=progress run-events carrying attrs={"milestone": "output", "stream": "stdout"|"stderr"}. The work is on the producer side (SSH transports return only complete command output today) plus removing a handful of container.run-only gates in the UI.
Plan / scope at a glance:
- API:
(run_id, seq)index onrun_events; SSE terminal-drain fix. No route or schema changes. - Worker: none in Phase A (host side); Phase B adds contextvar plumbing and an output sink (
apps/worker/output_sink.py, wired inflow_activities.py,step_handlers/services.py,shell_transport.py). - Plugins (sibling repo):
HandlerContext.emit_output— an additive SDK convenience wrapper over the existingemit_progresschannel (chunking, capping, UUID guard; no ABI change) — plus emissions from the SSH-based handlers. - UI: de-gate Logs tab / step panels from
handler === 'container.run'. - Tests/Docs: SSE drain regression tests, e2e non-container streaming test, this document.
Verification commands and pass criteria are listed in Verification.
Table of Contents
- How container streaming works today
- The SSH side before this work
- Why the same pipeline can be reused
- Constraints discovered
- Recommended phased design
- Batching parameters
- Repo split and release sequencing
- Risks and open questions
- Verification
1. How container streaming works today
End-to-end pipeline (all host-repo paths unless noted; plugin paths are relative to the sibling hegemony-step-plugins checkout):
docker start -a (attached pipes)
└─ RunContainerHandler._stream_and_capture() plugins/steps_container/…/run.py:1752-1870
├─ per-line readline on stdout/stderr
├─ batch: flush at 1 KB or 150 ms per stream
└─ await ctx.emit_progress(text,
attrs={"milestone": "output", "stream": "stdout"|"stderr"})
│
▼
emit_progress_callback apps/worker/flow_activities.py:1027-1062
├─ activity.heartbeat(message) (Temporal liveness)
└─ POST /internal/runs/{run_id}/events kind="progress" (best-effort; errors swallowed)
│
▼
create_run_event_internal apps/api/routers/internal.py:752-830
└─ SELECT runs.id … FOR UPDATE → seq = max(seq)+1 → INSERT run_events
│
▼
PostgreSQL run_events (no pubsub, no Redis for log data)
│
▼
GET /bff/runs/{run_id}/events/stream apps/api/routers/bff_runs.py:54-159
└─ DB poll every 1 s (clamped 0.5–30), LIMIT 200/tick, single-use ticket auth
│
▼
useRunEventStream apps/ui/src/lib/useRunEventStream.ts
├─ keep only attrs.milestone === 'output' && attrs.stream ∈ {stdout, stderr}
├─ split message on '\n' → TerminalLine{seq, subIndex, ts, stepId, stepRunId, stream, text}
└─ MAX_LINES = 5000 (FIFO)
│
▼
LogsTab → UnifiedTerminal apps/ui/src/pages/RunDetailComponents/LogsTab.tsx
apps/ui/src/components/UnifiedTerminal.tsxThe only container-specific parts were UI gates and copy — a containerStepRuns filter, an SSE-enable condition, per-step terminal gates, and "No Container Output" empty-state copy — all removed by the Phase A de-gating.
EventsTab.tsx and RunDetail.tsx exclude milestone === 'output' events generically, so the Events tab needed no change.
2. The SSH side before this work
The SDK contract had exactly one live channel
At the time of this analysis, ctx.emit_progress(message, device_id=None, attrs=None) (packages/step_sdk/…/contract.py:152-180 in the plugins repo) was the only way any handler code could push data mid-execution. HandlerServices had no log/event/stream API, and HandlerResult.evidence was terminal, whole-blob only.
Transports could not emit progress at all: they were constructed host-side (apps/worker/step_handlers/services.py:159-161) with only (DeviceConnectionSpec, cancellation_registry=, step_run_id=) — no ctx, no services. The Transport protocol (packages/step_sdk/…/services.py:38-77) exposed execute_command, execute_commands, execute_command_timing, scp_put, http_transfer — every one awaited the complete result.
What each transport could stream
| Transport | Model | Partial output available before this work |
|---|---|---|
netmiko (default) | blocking netmiko in run_in_executor thread | send_command is prompt-based — zero partial output. The execute_command_timing read loop (_execute_command_timing_sync, transport.py:942, loop ~:1069) already reads read_channel() every 2 s in a worker thread. Netmiko's unused session_log param is the only tap for plain send_command. |
scrapli | native async (asyncssh under scrapli) | send_command buffered; timing mode has a chunk reader _read_channel_chunk (transport.py:532) in its loop (:550+). |
asyncssh | native async | execute_commands (transport.py:243) uses fully-buffered connection.run() (:276); timing mode reads 4 KB PTY chunks (_read_pty_chunk, :376). Could switch the exec path to create_process + read loop. |
host shell_transport.py:46-65 (shell.execute) | native async | connection.run() fully buffered; persistent session per step — easy swap to create_process incremental reads (host-only change). Since v2.1.x the connection may ride a jump-host tunnel (:39-43, closed in close()); the swap must preserve that. |
Handler behavior before this work
netcli.execute/collect_evidence: oneexecute_commandscall for the whole command list (deliberate — config-mode detection needs the batch); output becamecli_outputartifacts after the fact, with noemit_progress.netcli.poll_until: per-attempt loop — a natural progress point that emitted nothing.shell.execute: per-commandtransport.run()loop — same.cisco.iosxe.upgrade.*: coarse milestoneemit_progressonly. Worst case:installranexecute_command_timing(read_timeout=1800, …)— up to 30 minutes with zero output reaching the UI (iosxe_install.py:301-307).
The Logs tab reads only SSE events, never artifacts, so at the time of this analysis SSH steps showed nothing there under any circumstances.
3. Why the same pipeline can be reused
The wire contract needs no change. Any code that can reach ctx.emit_progress and sends attrs={"milestone": "output", "stream": "stdout"} lights up the existing SSE stream, hook, and terminal. The UI line-splitter already handles multi-line batched messages. RunEvent already carries device_id, and attrs is free-form JSONB — so per-device labeling for multi-device SSH steps costs one extra attr (device_name), not a schema change.
Two structural facts make the "true streaming" phase cheaper than it looks:
- Transports are duck-typed (
packages/core/transports/registry.py:33-41; the registry never imports the SDK). The host contract test (tests/core/test_step_sdk_contract.py:61-67) pins exact signatures forHandlerServicesmethods only. An optional constructor kwarg on transport classes is invisible to both. - Transports are constructed host-side, so the host can inject an output sink at construction time without touching the
HandlerServicessignature — avoiding a lockstep host+plugins release.
There is now a live precedent for exactly this evolution path: step-plugins 0.2.0 added bastion support by extending DeviceConnectionSpec with an optional jump_host: JumpHostSpec field — an additive SDK change shipped across all three transport wheels with SDK_ABI_VERSION still 1 and no host lockstep (host merely populates the new field in services.py:connect). The output sink follows the same playbook (as a constructor kwarg rather than a spec field, since a live callable does not belong in a frozen connection spec).
4. Constraints discovered
Pre-existing host defects, both affecting container logs already — both since fixed (Phase 0):
- No
(run_id, seq)index onrun_events— every ingest computedmax(seq)under a per-runSELECT … FOR UPDATEand every SSE poll filteredseq > last_seq ORDER BY seq, both unindexed and degrading linearly with per-run event count. Fixed by migration037_run_events_run_seq_index(ix_run_events_run_seq). - SSE replay truncation for finished runs — the generator emitted one ≤
PAGE_LIMITpage, then the terminal-status check firedcompleteand returned, so a finished run with more events replayed only its first page. Fixed inbff_runs.py: full pages loop immediately without the terminal check;completefires only after a short page proves the backlog is drained. Regression tests:tests/api/test_bff_runs_stream.py.
Other constraints:
device_idmust be a valid UUID or the whole event POST is rejected with 400 (internal.py:793-796). Handlers usedevice.get("id", device.get("mgmt_host"))in artifact paths — that fallback must never be passed toemit_progress. Display names belong inattrs.device_name.- The worker services object is a module-level singleton shared by all concurrent steps (
flow_activities.py:62, bound at:1227-1228). Per-step sink state must ride acontextvarset aroundhandler.execute, not the services instance. - Worker egress gating (since v2.1.x): every device/bastion dial passes
enforce_worker_dialbefore connecting (step_handlers/services.py,shell_transport.py). Orthogonal to streaming — the sink emits over the existing worker→API internal channel, which is not policied traffic. - netmiko chunks arrive on an executor thread — they must be marshalled to the event loop via
asyncio.run_coroutine_threadsafe(loop reference captured at construction on the loop thread). - SSE delivery fetches at most 200 rows per query (
poll_intervalaccepts 0.5–30 s). Full pages drain in quick succession (a short fast-drain pause between pages); the configured poll delay applies only after a short page shows the backlog is caught up. Producers should still batch — every event is one row-locked insert on the way in. - Event retention exists (
apps/api/services/maintenance/jobs/run_event_pruning.py— age-based, terminal runs only), so added volume is bounded over time.
5. Recommended phased design
Phase 0 — Host hardening (also benefits container logs)
- Alembic migration
037_run_events_run_seq_indexaddsix_run_events_run_seq (run_id, seq). - Terminal-run SSE drain fixed in
bff_runs.py(completeonly after a short page; full pages refetch immediately).
Phase A — Handler-level emission (zero ABI change)
Host (UI de-gating): the Logs tab keys on output events, not handler id:
RunDetail.tsx:container.runfilter dropped; the stream is enabled whenever a run is loaded (runs without output events carry few events, so the finished-run replay stays cheap and closes oncomplete); terminal lanes are built from the streamed lines themselves, resolving step names fromstep_runsbystepRunIdwith a node-id fallback.StepDetailPanel.tsx/RunGraphTab.tsx: the per-step terminal renders on the panel's Execution tab for any step with lines (or while the step is RUNNING), whatever its handler.LogsTab.tsx: neutral copy ("No output was produced during this run").- Regression e2e:
e2e/unified-terminal.spec.tsincludes a non-container-handler streaming test.
Plugins (completed-output emission via existing ctx.emit_progress) — each event carries attrs.command and attrs.device_name (shipped in hegemony-step-plugins 0.3.0):
shell.execute: emit stdout/stderr after each command (real incremental value — one event per command).netcli.poll_until: emit per attempt with an attempt counter.netcli.execute/collect_evidence: emit per result after the batch call returns (honest limitation: batch-end granularity with netmiko). Config-mode sessions are excluded from streaming — see the secrets bullet under Risks.cisco.iosxe.*: emit accumulated output tails after timing calls return.- Shape: same
attrscontract plusdevice_name/command; UUID-onlydevice_id; split at ~8 KB on line boundaries; cap ~64 KB per command with a truncation marker (full output stays incli_outputartifacts).
Either repo can ship first; both orderings degrade gracefully.
Phase B — True mid-command streaming (transport sink injection)
As built:
- SDK: an additive
TransportOutputSinkprotocol (push,push_from_thread,flush,flush_from_thread;SDK_ABI_VERSIONstays 1). Transport wheels take an optionaloutput_sinkconstructor kwarg, feature-detected via a class attribute (supports_output_sink = True). Rejected alternatives: per-callon_outputkwargs (old wheels raiseTypeErroron unknown kwargs since the protocol is never runtime-checked, forcing feature detection into every handler) and a callable on the frozenDeviceConnectionSpec. - Host: a
contextvarset in_execute_step_scopedcarries the boundemit_progresscallback;WorkerHandlerServices.connectandopen_shell_transportbuild a per-device batching sink (apps/worker/output_sink.py) and pass it only to transports that advertise support. The sink strips ANSI, normalizes CRLF, redacts the step's resolved credential values (≥ 4 chars, longest-first, with a withheld tail so a secret split across two flush batches still lands whole in one redaction window), batches (1 KB or 500 ms, 16 KB hard split), serializes per-stream emission so events never reorder, and never raises into the transport. - netmiko:
_execute_command_timing_syncstreams the$ commandboundary marker, the initialsend_command_timingoutput, prompt-answer output, and each wait-loopread_channel()chunk via the sink's thread-safe entry point (push_from_threadmarshals withcall_soon_threadsafe). Plainsend_commandstays non-streaming (prompt-based, no partial output;session_logpiping is a documented future option). Config-mode (send_config_set) transcripts are never streamed at this transport-sink level — the permitted path for config sessions is the handler-level completed-transcript emission (netcli.executeemits the finished session viaemit_output, which flows through the redacting emit callback; see the redaction section below and Risks). Mid-command live streaming of a config session remains disabled. - scrapli / asyncssh: the timing-mode read loops push the marker and each chunk directly (natively async). The asyncssh exec path (
execute_commands) deliberately stays buffered: it batches command lists whose text may carry credentials (config pushes), and streaming needs per-command$ …markers that would echo them. Its completed outputs reach the Logs tab handler-level, through the same redacting callback. - Host
shell_transport.py:run()switches tocreate_processincremental reads when a sink is bound —shell.executegains true streaming with zero handler edits. The$ commandmarker exposes exactly what Phase A'sattrs.commandlabel (and thecli_outputevidence artifacts) already expose for every shell command; an operator-authored command carrying an inline literal secret was already visible on those surfaces, and the sink additionally redacts all platform-resolved credentials. - Dedup: a streaming transport sets duck-typed flags (
streams_timing_outputfor CLI timing mode,streams_exec_outputfor the shell transport); handlers skip their Phase A completed-output emissions when the flag is set, so output is never double-posted. The transport, not the handler, emits the boundary marker.
This unlocks the headline UX win: live output during the 30-minute install add file … window.
Flow-wide secret redaction
The sink's per-connection redaction generalized into a step-scoped secret registry (apps/worker/secret_registry.py):
- Every secret the platform resolves while a step executes registers its value — device/enable/bastion credentials at the
connect()/open_shell()sites, and everysecret()/file()template resolution via an observer hook the worker installs into the core template resolver (env()is infrastructure config, not observed). - The step's
emit_progresscallback redacts registered values from every message and attr before the Temporal heartbeat and the run-events POST — covering the sink's streamed chunks, Phase A handler emissions, and container logs alike. - Evidence artifacts are redacted at persistence (
content_text, recursivecontent_jsonincluding keys, display name) — credentials a device echoes no longer reach artifacts in clear either. - The step result (summary, error, metrics, output) is redacted before
execute_stepreturns it, on the handler-exception path too, because it lands in Temporal history, the step run and later steps' templates. When that changes anything, the worker logs the field paths (never values) and emits one progress line naming them. Usernames resolved fromsecret()refs are resolved undersecret_observer_suppressed()and never registered. - The scope is deliberate: used secrets only. A literal an operator types into a command is indistinguishable from ordinary text; we redact what the platform resolved, and cannot guarantee more. The
REDACT_MIN_CHARSfloor and its warning apply registry-wide.
With both surfaces redacted, the config-mode streaming exclusion is lifted: netcli.execute emits config-session transcripts to the Logs tab (plugins side), and all five cisco.iosxe.upgrade.* handlers emit their command outputs/transcripts (with streams_timing_output dedup guards where the transport may stream live).
Per-device parallelism (plugins)
Multi-device SSH handlers accept max_parallel_devices (default 1 = strictly sequential, the classic loop; clamped to 32). The SDK's for_each_target_device bounds concurrency with a semaphore and returns results in device order. Concurrent devices' live output interleaves in the Logs tab, where lines already carry attrs.device_name labels, and the terminal offers two device-filter levels (host UI): a per-step dropdown on each step badge (shown once a step has output from two or more devices) narrowing that step to specific devices, and a run-wide Devices dropdown filtering a device across all steps. Both derive their device lists from the streamed lines, so they always match what the terminal can show.
6. Batching parameters
| Producer | Recommendation |
|---|---|
| Container (existing) | 1 KB or 150 ms per stream — unchanged |
| Phase A handlers | one event per command result / attempt; split ~8 KB on line boundaries; cap ~64 KB per command + truncation marker |
| Phase B sink | flush per timing-poll chunk (2–2.5 s cadence is naturally one event); 1 KB-or-500 ms for fast producers; hard-split 16 KB; coalesce above ~4 events/s per device |
Rationale: each flush is one HTTP POST + one row-locked insert, and live viewers drain new rows on the poll_interval cadence (a backlog larger than one 200-row page delivers multiple pages in quick succession, with a short fast-drain pause between pages) — sub-second flushes buy nothing. Projected worst case (~20 devices in timing mode ≈ 10 events/s per run) is comfortable once the (run_id, seq) index exists.
7. Repo split and release sequencing
Host repo: Phase 0 (migration + SSE fix); Phase A UI de-gating (RunDetail.tsx, LogsTab.tsx, StepDetailPanel.tsx, RunGraphTab.tsx, useRunEventStream.ts, UnifiedTerminal.tsx); Phase B contextvar plumbing (flow_activities.py), sink construction (step_handlers/services.py), streaming shell_transport.py, plugin pin bumps.
Plugins repo: Phase A handler emissions (steps_shell, steps_netcli, optional steps_cisco_iosxe tails); Phase B additive SDK types (minor version, SDK_ABI_VERSION stays 1) + the three transport wheels.
No lockstep required. Transport wheels ship first and are inert on old hosts (the host never passes the kwarg). The host change ships second with pin bumps (.github/plugin-pins.json, pyproject.toml version floors, hegemony-demo-data/deploy/compose/demo-plugin-wheels.txt). Old host + new wheels: works, no streaming. New host + old wheels: feature-detection returns false, works.
8. Risks and open questions
- Secrets in output. Config-mode sessions echo lines like
username … password …, and that content reaches both the progress-event path andcli_outputartifacts — streaming adds prominence, not a new exposure class. Historically neither surface scrubbed anything; both are now covered by the control below. Chosen control (current): flow-wide redaction of used secrets. The step-scoped secret registry redacts every platform-resolved secret (device/enable/bastion credentials,secret()/file()template values) from progress events, Temporal heartbeats, AND evidence artifacts; the sink additionally holds back a tail across batch boundaries so split credentials still redact. Under that control the config-mode streaming exclusion is lifted and transcripts reach the Logs tab. Residual exposure: operator-typed inline literals (not resolvable, accepted), and values below the documented floor — values shorter thanREDACT_MIN_CHARS(4) are not redacted, because substring-replacing 1-3 char values corrupts ordinary output (everya, every port22) beyond the secrecy such a credential affords; a warning fires whenever the floor drops a value. The asyncssh exec path stays buffered at transport level (itsexecute_commandsbatch emission now flows through the redacting callback like everything else). - Multi-device steps. One step run, N devices;
UnifiedTerminallanes key bystepRunId(distinct lanes per loop iteration), not per device. Implemented in v1:attrs.device_nameis carried intoTerminalLine.deviceNameand lines render as[step @ device]; per-device filtering shipped as the per-step and run-wide device dropdowns described in the per-device parallelism section. - Noise. Command echo, pagers (
--More--), and VT100 escapes appear in raw chunks — strip ANSI in the sink; accept echo noise in v1.
9. Verification
Host repo — the repository gates, all passing means exit code 0 with no warnings and full-suite success:
task py:all(format, lint, typecheck) anduv run ty check --error-on-warning— clean.task test— full Python suite green, includingtests/api/test_bff_runs_stream.py(finished-run replay drains past one page, exact page multiples,after_seqresume),tests/worker/test_output_sink.py(sink batching/sanitization/redaction, thread marshalling, feature-detected injection, streaming shell transport), andtests/worker/test_secret_registry.py(registry lifecycle/isolation, resolver observer feed, redaction helpers incl. recursive artifact JSON, sink merge).cd apps/api && uv run alembic heads— exactly one head (037_run_events_run_seq_index).task ui:verifyandtask ui:audit— formatting, CSS/Tailwind lint, API-header checks, lint, typecheck, generated-type sync, build, and the Playwright smoke suite, all clean.npm --prefix apps/ui run test:e2e— full Playwright suite green, including "Non-container step output" (the same SSE events render under anetcli.executehandler, step lanes labelled, per-device lines shown as[step @ device]). CI uploads the Playwright HTML report on failure.
Plugins repo (all must pass):
uv run pytest— green, includinghegemony-step-plugins/tests/test_output_emission.py(attrs shape, chunking, UUID guard, truncation, config-transcript emission, per-handler wiring),hegemony-step-plugins/tests/test_output_streaming.py(sink kwarg feature detection, timing-loop marker/chunk pushes per transport, broken-sink resilience, exec-path exclusion, handler dedup flags), andhegemony-step-plugins/tests/test_device_concurrency.py(sequential default, bounded overlap, order preservation, handler-level parallelism).uv run ruff format --check . && uv run ruff check . && uv run ty check— clean.
End-to-end (manual, after both merge): run a flow with a netcli.execute or shell.execute step — its output appears in the Logs tab live and on the finished-run view; a config-mode netcli.execute step streams its transcript with every platform-resolved credential shown as [redacted], in the Logs tab and in the Artifacts alike. After the plugins release and host pin bump, a cisco.iosxe.upgrade.install step shows install output live during the timing window instead of one transcript at the end, and a step with max_parallel_devices > 1 interleaves its devices' labelled output.
Structured-record check: streamed output persists as run_events rows — each carries run_id, step_run_id/step_id, the emitting device_id (UUID-guarded; display names ride attrs.device_name), and attrs={"milestone": "output", "stream": …}. Those rows, not application logger lines, are the durable structured record of streamed output; acceptance checks should query run_events for the run.