Streaming governance
In plain English: live market data flows through AlphaSwarm on Kafka, a message pipeline. Two classic failure modes with pipelines like this are (1) a producer changes the message format and every consumer downstream silently misreads data, and (2) code addresses topics by hard-coded name strings that drift from what actually exists. The ADR-0020/0021 work makes both impossible-by-construction: every message is framed with a schema id that consumers verify (refusing to guess when they don't recognize it), and topics are resolved through the governed data catalog instead of string literals.
Confluent wire framing with dual-read
Messages are framed in the Confluent wire format — a magic byte plus a
schema id ahead of the payload — implemented in
alphaswarm/streaming/schemas/wire.py.
The decoder is dual-read during the transition:
- Framed payloads resolve their schema by id — this works on mixed-schema topics and refuses to guess on unknown ids (fail closed, dead-letter rather than misparse).
- Legacy unframed bytes decode by topic, exactly as before.
So producers can cut over topic-by-topic without a coordinated big-bang consumer upgrade.
The bootstrap manifest
Schema ids come from a committed bootstrap manifest, registered into the schema registry by a bootstrap registrar at deploy time — the produce hot path never contacts a live registry, and framing fails closed if the manifest is missing. This keeps the latency-sensitive path free of network dependencies while preserving central governance.
Schema changes are gated in CI
(.github/workflows/streaming-schema-gates.yml).
Catalog-first stream addressing
KafkaDataFeed.from_feed_urn(...) resolves a feed URN through the
governed catalog binding to a concrete topic, instead of consumers
holding topic-name literals. A two-way parity gate checks the
catalog bindings against the wheel-bundled TOPIC_BY_SCHEMA canon, so
catalog drift is caught in CI rather than at 3 a.m.
Diagnostics posture
The 36 /streaming/* diagnostic HTTP routes are now default-off
(streaming_diagnostics_routes_enabled) and stripped from the public
OpenAPI spec. The sanctioned always-on diagnostics path is the MCP
tool surface, which carries governance (auth, audit, tenancy) that the
raw routes did not.
Topic governance also counts failed dead-letter produces —
alphaswarm_stream_deadletter_failed_total{origin_topic} — so a
dead-letter path that is itself failing becomes visible.
What was retired (ADR-0021)
- The curated Airbyte catalog (superseded by the connector control plane).
- Databento and Robinhood live-streaming scaffolds (historical Databento ingestion is unaffected).
data/sources/secis marked legacy.
Flags
| Flag | Default | Effect |
|---|---|---|
stream_confluent_wire_enabled | off | Producers emit Confluent-framed payloads (consumers dual-read regardless). |
streaming_diagnostics_routes_enabled | off | Re-exposes the /streaming/* diagnostic routes. |
See also
- Market-data vendors — where the data on these topics comes from.
- Connector control plane — governed connector onboarding.
- Dataops: Airbyte + Dagster runbook — operating the batch side of the data plane.
- Six-clock latency —
md.event_lagmeasures this pipeline's hand-off latency.