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() -> NoneEnd every stream, stop the connection and write what is left (bounded wait).
start
start(tickers: Sequence[str], card_date: dt.date, session_id: str) -> NoneRecord 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) -> StreamRecordTapeState
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) -> NoneA heartbeat pong: everything the exchange sent before it has been received.
committed
committed(batch: Batch) -> Noneconnect
connect(now: dt.datetime) -> Nonediscard
discard(batch: Batch, now: dt.datetime) -> NoneThe 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) -> NoneThe connection is gone: every live stream ends at the last proof of life.
drain
drain() -> Batch | NoneEverything 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) -> Nonemaintain
maintain(now: dt.datetime) -> NonePeriodic: resubscribe after a drop once the queue has room; retire streams.
want
want(tickers: Iterable[str], now: dt.datetime) -> NoneFunctions
f_ended
f_ended(reason: str) -> strparse_delta
parse_delta(body: dict[str, Any], stream_id: int, seq: int, at: dt.datetime) -> DeltaRecordparse_snapshot
parse_snapshot(body: dict[str, Any], stream_id: int, seq: int, at: dt.datetime) -> SnapshotRecord