Skip to content

betting_combat.trading.book_recorder

Records the card’s Kalshi order books for research (Kalshi keeps no book history).

Off unless portfolio.record_orderbook is set. When it is, the card run hands it every market of the card’s six UFC series (winner, distance, method of victory, method of finish, rounds, victory round) and it records each book from the run’s start until the run ends (every fight settled): the orderbook_delta channel’s snapshots and deltas, into store/orderbook.py’s tables. It never touches trading:

  • its own WebSocket connection (the trading mirror, consumers/kalshi/stream.py, is not shared: a gap or a reconnect here never resets the trading books) and its own one-connection database pool (wiring.py), so it cannot take a connection the trading store is waiting for;
  • the event loop only parses a message and appends it to an in-memory queue; the batch write (every FLUSH_S) runs in a worker thread, one transaction per batch;
  • the queue is bounded (QUEUE_ROWS). When it is full (a slow or unreachable database) the message is dropped, never waited for: that stream ends (dropped, counted) and a fresh subscription is opened once the queue has room again;
  • every failure is counted in findings (merged into the card run’s findings at the end) and logged; nothing it raises reaches the card run.

Sequence numbers are per subscription (sid). A missing number (a gap) or a message that does not parse ends the stream (gap / malformed) and resubscribes at once for fresh snapshots; more than GAP_LIMIT such ends in GAP_WINDOW rebuild the connection instead. A broken connection ends every stream (reconnect) and the next connection resubscribes (backoff 1 s doubling to 15 s). A change of the card’s markets opens a new subscription and retires the old one once the new one has snapshotted its markets (or after RETIRE_AFTER), so the recording has no hole.

A stream’s complete_through is the last moment the connection proved alive (a message or a heartbeat pong) with every message of the stream so far recorded; it is written in the same transaction as the rows it covers, so it never claims rows that are not in the database.

Classes

Batch

BookRecorder

Runs a TapeState against Kalshi: one connection, a heartbeat, a sender, and the batch writer. start and close are the card run’s only calls.

close

close() -> None

End every stream, stop the connection and write what is left (bounded wait).

start

start(tickers: Sequence[str], card_date: dt.date, session_id: str) -> None

Record these markets (an empty list stops recording; the same list is a no-op).

status

status() -> dict[str, Any]

MalformedMessage

Bases: ValueError

RecorderResync

Bases: Exception

The connection must be rebuilt (an exchange error, too many gaps, garbage).

Stream

One subscription to orderbook_delta for a set of markets.

record

record(card_date: dt.date, session_id: str) -> StreamRecord

TapeState

The pure half of the recorder: applies messages in order, checks sequence, keeps the bounded queue and the streams’ state, and lists the commands to send. Everything here runs on the event loop and is O(message).

alive

alive(at: dt.datetime) -> None

A heartbeat pong: everything the exchange sent before it has been received.

committed

committed(batch: Batch) -> None

connect

connect(now: dt.datetime) -> None

discard

discard(batch: Batch, now: dt.datetime) -> None

The batch could not be written: its rows are given up and every stream that had rows in it ends where its committed record ends (write_failed).

disconnect

disconnect(now: dt.datetime, detail: str | None = None) -> None

The connection is gone: every live stream ends at the last proof of life.

drain

drain() -> Batch | None

Everything queued (with anything a failed write left) and every stream row that changed. A live stream is complete through the last proof of life: every message it received until then is in this batch or an earlier one.

handle

handle(msg: Any, received_at: dt.datetime) -> None

maintain

maintain(now: dt.datetime) -> None

Periodic: resubscribe after a drop once the queue has room; retire streams.

want

want(tickers: Iterable[str], now: dt.datetime) -> None

Functions

f_ended

f_ended(reason: str) -> str

parse_delta

parse_delta(body: dict[str, Any], stream_id: int, seq: int, at: dt.datetime) -> DeltaRecord

parse_snapshot

parse_snapshot(body: dict[str, Any], stream_id: int, seq: int, at: dt.datetime) -> SnapshotRecord