Skip to content

betting_combat.consumers.kalshi.stream

Kalshi WebSocket: live order books and trades for the markets of one card.

One connection carries two subscriptions, orderbook_delta and trade. Every message on a subscription carries that subscription’s sequence number (seq): snapshots, deltas, trades and ok replies alike (recorded 2026-09-26, tests/fixtures/kalshi/ws_stream.jsonl). A missing number means a lost message, so the connection is rebuilt and every book re-snapshotted.

A book is served only while it is synced, and it is “as of” the last moment the connection proved alive (any message, or a pong). The caller treats an old as_of as stale and falls back to a REST snapshot, so a dead or lagging stream can cost speed but never correctness.

Classes

KalshiStream

Keeps a StreamState mirroring Kalshi until cancelled, rebuilding on any break.

prints

prints(ticker: str, since: dt.datetime, until: dt.datetime) -> list[Print] | None

run

run() -> None

status

status() -> dict[str, Any]

top

top(ticker: str) -> BookTop | None

volume

volume(ticker: str, since: dt.datetime, until: dt.datetime) -> float | None

LocalBook

One market’s resting bids on each side: price (pips) → contracts.

apply

apply(side: str, price: str, delta: str) -> None

from_snapshot

from_snapshot(body: dict[str, Any]) -> LocalBook

top

top(ticker: str, as_of: dt.datetime) -> BookTop

StreamResync

Bases: Exception

The mirror lost its place (sequence gap, delta before snapshot, negative size).

StreamState

The pure half of the stream: applies messages in order and checks sequence.

handle

handle(msg: dict[str, Any], received_at: dt.datetime) -> None

prints

prints(ticker: str, since: dt.datetime, until: dt.datetime) -> list[Print] | None

The trades in (since, until], oldest first; None if the stream has not seen all of them (it connected later, or since is past its memory).

reset

reset() -> None

A new connection: nothing is known until it is re-snapshotted.

top

top(ticker: str) -> BookTop | None

The mirrored top of ticker’s book, as of the last proof of life.

volume

volume(ticker: str, since: dt.datetime, until: dt.datetime) -> float | None

Contracts traded in (since, until]; None if the stream has not seen all of it.