Skip to main content

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/sec is marked legacy.

Flags​

FlagDefaultEffect
stream_confluent_wire_enabledoffProducers emit Confluent-framed payloads (consumers dual-read regardless).
streaming_diagnostics_routes_enabledoffRe-exposes the /streaming/* diagnostic routes.

See also​