Skip to content

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 on run_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 in flow_activities.py, step_handlers/services.py, shell_transport.py).
  • Plugins (sibling repo): HandlerContext.emit_output — an additive SDK convenience wrapper over the existing emit_progress channel (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 ​

  1. How container streaming works today
  2. The SSH side before this work
  3. Why the same pipeline can be reused
  4. Constraints discovered
  5. Recommended phased design
  6. Batching parameters
  7. Repo split and release sequencing
  8. Risks and open questions
  9. 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):

text
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.tsx

The 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 ​

TransportModelPartial output available before this work
netmiko (default)blocking netmiko in run_in_executor threadsend_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.
scraplinative async (asyncssh under scrapli)send_command buffered; timing mode has a chunk reader _read_channel_chunk (transport.py:532) in its loop (:550+).
asyncsshnative asyncexecute_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 asyncconnection.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: one execute_commands call for the whole command list (deliberate — config-mode detection needs the batch); output became cli_output artifacts after the fact, with no emit_progress.
  • netcli.poll_until: per-attempt loop — a natural progress point that emitted nothing.
  • shell.execute: per-command transport.run() loop — same.
  • cisco.iosxe.upgrade.*: coarse milestone emit_progress only. Worst case: install ran execute_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:

  1. 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 for HandlerServices methods only. An optional constructor kwarg on transport classes is invisible to both.
  2. Transports are constructed host-side, so the host can inject an output sink at construction time without touching the HandlerServices signature — 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):

  1. No (run_id, seq) index on run_events — every ingest computed max(seq) under a per-run SELECT … FOR UPDATE and every SSE poll filtered seq > last_seq ORDER BY seq, both unindexed and degrading linearly with per-run event count. Fixed by migration 037_run_events_run_seq_index (ix_run_events_run_seq).
  2. SSE replay truncation for finished runs — the generator emitted one ≤PAGE_LIMIT page, then the terminal-status check fired complete and returned, so a finished run with more events replayed only its first page. Fixed in bff_runs.py: full pages loop immediately without the terminal check; complete fires only after a short page proves the backlog is drained. Regression tests: tests/api/test_bff_runs_stream.py.

Other constraints:

  • device_id must be a valid UUID or the whole event POST is rejected with 400 (internal.py:793-796). Handlers use device.get("id", device.get("mgmt_host")) in artifact paths — that fallback must never be passed to emit_progress. Display names belong in attrs.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 a contextvar set around handler.execute, not the services instance.
  • Worker egress gating (since v2.1.x): every device/bastion dial passes enforce_worker_dial before 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_interval accepts 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.

Phase 0 — Host hardening (also benefits container logs) ​

  • Alembic migration 037_run_events_run_seq_index adds ix_run_events_run_seq (run_id, seq).
  • Terminal-run SSE drain fixed in bff_runs.py (complete only 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.run filter 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 on complete); terminal lanes are built from the streamed lines themselves, resolving step names from step_runs by stepRunId with 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.ts includes 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 attrs contract plus device_name / command; UUID-only device_id; split at ~8 KB on line boundaries; cap ~64 KB per command with a truncation marker (full output stays in cli_output artifacts).

Either repo can ship first; both orderings degrade gracefully.

Phase B — True mid-command streaming (transport sink injection) ​

As built:

  • SDK: an additive TransportOutputSink protocol (push, push_from_thread, flush, flush_from_thread; SDK_ABI_VERSION stays 1). Transport wheels take an optional output_sink constructor kwarg, feature-detected via a class attribute (supports_output_sink = True). Rejected alternatives: per-call on_output kwargs (old wheels raise TypeError on unknown kwargs since the protocol is never runtime-checked, forcing feature detection into every handler) and a callable on the frozen DeviceConnectionSpec.
  • Host: a contextvar set in _execute_step_scoped carries the bound emit_progress callback; WorkerHandlerServices.connect and open_shell_transport build 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_sync streams the $ command boundary marker, the initial send_command_timing output, prompt-answer output, and each wait-loop read_channel() chunk via the sink's thread-safe entry point (push_from_thread marshals with call_soon_threadsafe). Plain send_command stays non-streaming (prompt-based, no partial output; session_log piping 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.execute emits the finished session via emit_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 to create_process incremental reads when a sink is bound — shell.execute gains true streaming with zero handler edits. The $ command marker exposes exactly what Phase A's attrs.command label (and the cli_output evidence 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_output for CLI timing mode, streams_exec_output for 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 every secret() / 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_progress callback 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, recursive content_json including 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_step returns 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 from secret() refs are resolved under secret_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_CHARS floor 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 ​

ProducerRecommendation
Container (existing)1 KB or 150 ms per stream — unchanged
Phase A handlersone event per command result / attempt; split ~8 KB on line boundaries; cap ~64 KB per command + truncation marker
Phase B sinkflush 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 and cli_output artifacts — 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 than REDACT_MIN_CHARS (4) are not redacted, because substring-replacing 1-3 char values corrupts ordinary output (every a, every port 22) 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 (its execute_commands batch emission now flows through the redacting callback like everything else).
  • Multi-device steps. One step run, N devices; UnifiedTerminal lanes key by stepRunId (distinct lanes per loop iteration), not per device. Implemented in v1: attrs.device_name is carried into TerminalLine.deviceName and 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) and uv run ty check --error-on-warning — clean.
  • task test — full Python suite green, including tests/api/test_bff_runs_stream.py (finished-run replay drains past one page, exact page multiples, after_seq resume), tests/worker/test_output_sink.py (sink batching/sanitization/redaction, thread marshalling, feature-detected injection, streaming shell transport), and tests/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:verify and task 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 a netcli.execute handler, 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, including hegemony-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), and hegemony-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.

Released as open source under the AGPL-3.0-or-later license. Development is sponsored by Rexonix s.r.o.. Contact — [email protected].