RFC: Bitnomial drop-copy ingestion on a Postgres-backed pipeline

Date: 2026-08-12

Status: Implemented, landing as a stacked series of PRs. This document is the map for reviewing them. Tracked as A-4547.

The problem this exists to remove

On 2026-08-02 the aiex-demo box lost two real fills, and nothing noticed (A-4465).

Bitnomial reset the drop-copy socket. Our consumer crashed, correctly, because a guard added a week earlier made it fail loudly instead of reconnecting into the void. On restart it asked ClickHouse "where did I get to?", got an answer that was behind what had actually been received, and rewound the FIX session's inbound sequence number to match. That rewind desynchronised the session: for the next two hours the engine rejected everything as out of sequence, and live execution reports arriving during that window were never written down at all. When it finally reconnected and asked the venue to resend, the venue replied with a gap fill -- "messages 140662 through 140694 are not coming" -- and 32 messages vanished.

Two fills, a buy and a sell of XAU-PERP, confirmed on the wire, produced zero rows in trades. Positions, realized PnL, fees and /fills were all wrong for that window.

Throughout, every service published healthy: true, because health was measured by whether bytes were moving on the FIX socket rather than by whether anything was being ingested. The order-entry gate stayed open. The consumer's progress marker even advanced -- while committing nothing.

Three separate patches were attempted and abandoned (#3416, #3418, #3486). Each fixed a symptom. The shared cause is structural, and it is worth stating plainly:

The system could not tell the difference between "I have processed everything" and "I have lost track of where I was" -- and the place it stored its position was not the place that held the truth.

The shape of the fix

One idea does most of the work: write down every message the venue sends, in order, before interpreting any of it.

Today the FIX session, the interpretation of messages, and the money all happen in one pass. If any part of that pass is disturbed, there is no record of what arrived, so recovery has to guess. The pipeline splits it in two:

flowchart LR
    V["Bitnomial<br/>FIX drop copy"] -->|"execution reports"| R["Recorder<br/><i>aiex-dropcopy-recorder</i>"]
    R -->|"append, in order"| L[("btnl.dropcopy_log<br/>position 1, 2, 3, ...")]
    L --> C1["Trade consumer<br/><i>aiex-trade-consumer</i>"]
    L -.planned.-> C2["Risk consumer"]
    L -.planned.-> C3["Orders consumer"]
    L -.planned.-> C4["ClickHouse projector"]
    C1 --> M[("trades, positions,<br/>fees, balances")]
    C4 -.-> CH[("ClickHouse<br/><i>reporting only</i>")]

The recorder owns the venue session and does exactly one job: put each message into btnl.dropcopy_log with a position number, unread and uninterpreted, along with its raw bytes. It never books anything.

Each consumer reads the log at its own pace and remembers how far it got. The trade consumer books fills. Others are planned and can be added without touching the recorder or each other.

ClickHouse is still there, but demoted. It is a reporting copy derived from the log, and it is no longer allowed to answer the question that started the incident. Nothing asks ClickHouse where to resume.

Concepts, defined

These terms appear throughout the code and the PRs. None of them are standard; they mean something specific here.

Log position. A number assigned by the recorder, counting from 1, in the order messages arrived. It is ours, not the venue's. The venue's own sequence numbers restart and are re-used; positions never are.

Watermark. How far one consumer has got, stored as a row in btnl.consumer_watermarks. Each consumer has its own, so a slow one cannot hold up a fast one, and a crashed one resumes exactly where it stopped. The watermark is written in the same database transaction as the work it describes -- this is what makes "I booked this fill" and "I have got past this fill" impossible to disagree.

Fill journal. A record of every half of every execution the consumer has booked, keyed by (trade_generation, exec_id, side, liquidity). Its purpose is to make booking safe to retry: if the same message arrives twice -- which a venue resend guarantees will happen -- the second attempt collides with the existing row and books nothing.

Halves and pairing. Bitnomial reports one execution as two messages, one for each side of the trade (liquidity is add for the maker, remove for the taker). The finished trade row can only be written once both have arrived. A half that never gets its counterpart is eventually written against a placeholder account so the books still balance -- see "lone halves" below.

Epoch. A stretch of the venue's sequence numbering. When Bitnomial resets the session, its numbering restarts at 1, so "sequence number 5" is ambiguous without knowing which stretch it belongs to. An epoch identifies the stretch.

Trade generation. The bookkeeping equivalent: a marker that separates everything booked before a reset from everything after, so old and new execution ids cannot collide in the journal.

Witness ledger. A record of sequence numbers the session saw but never handed up for recording. Without it, a gap in the log looks identical to lost data. With it, the recorder can tell "nothing was there" from "something was there and I lost it" -- and only the second one is an emergency.

Source liveness. A durable statement about whether the drop-copy feed is actually delivering, kept separately from whether the socket is connected. This is the direct answer to the incident's healthy: true.

The invariants

These are the properties every PR in the stack should be read against. They are what make the failure class unreachable rather than merely unlikely.

1. Nothing is interpreted before it is written down. The recorder appends raw bytes. If a message cannot even be parsed, it is still recorded, marked unreadable, and the consumer stops on it. A message we cannot understand must never be a message we silently skip.

2. Work and progress commit together, or neither. Booking a fill and advancing the watermark past it are one transaction. This is the direct fix for the torn write in A-3790, where trades landed in one store and the money that went with them never committed to the other.

3. Gaps fail closed, loudly, forever. A consumer reads a page of rows and refuses it unless the first row is exactly watermark + 1. It does not skip ahead, does not retry past it, and no restart clears it. A stuck pipeline is a bad morning; a pipeline that quietly skipped position 140,663 is a wrong balance nobody finds for a month.

flowchart TD
    A["Consumer reads next page"] --> B{"Is the first row<br/>watermark + 1?"}
    B -->|yes| C["Apply rows, advance<br/>watermark, same transaction"]
    B -->|no| D["Stop. Watermark unchanged."]
    C --> A
    D --> E["Stays stopped across restarts.<br/>Needs a human."]

4. Booking is exactly-once under replay. A venue resend after an outage re-delivers thousands of messages we already have. Each collides with its journal row and books nothing. Measured on a 570,000-message resend: zero double-books.

5. No query may get slower as history accumulates. This one is easy to violate by accident and expensive to discover. The first version of the lone-half sweep asked "which halves have no counterpart?" by searching the whole journal every time. That cost grows with everything that ever happened, even though the number of halves actually waiting stays tiny -- so it passed every test and then degraded in production. It is now a column on the row plus an index over only the rows still waiting: 190.7 ms down to 0.014 ms.

6. Being connected is not being current, and being caught up is not being current. A session can heartbeat perfectly while delivering nothing. A consumer sitting at the end of a log that stopped growing is genuinely caught up and correct to say so. Both are true at once during an outage, which is precisely how the incident looked healthy. They are separate facts and health requires both.

Which layer solves what

FIX is a mature protocol with real error correction, and a mature engine implements it well. Both are load-bearing here and neither is being replaced. The reason this pipeline exists anyway is that the guarantees stop at a specific line, and the fills we lost were on the far side of it.

The FIX session layer gives us:

A FIX engine (quickfix and its equivalents) adds:

Neither solves any of the following, because they are about our books rather than the wire:

The dividing line is worth stating plainly, because it is the whole argument for this work: the protocol guarantees we notice a gap. It guarantees nothing about our balances. Everything in the third list is unchanged by which engine, binding or language we run, which is why swapping any of those does not move it.

Worked example: an outage, before and after

The venue drops the session for two hours, then comes back and resends.

Before. The consumer crashes. On restart it asks ClickHouse for its position, gets a stale answer, rewinds the session's sequence number to match, and desynchronises. Live messages arriving during the confusion are not written anywhere. The venue answers the resend request with a gap fill and 32 messages are gone. Health is green the whole time.

After. The recorder loses its session and retries. Nothing is booked during the outage, because nothing arrives -- but nothing is lost either, because the log's last position is durable and the recorder resumes from it rather than from a guess. Source liveness goes not-current, so health goes red and the order-entry gate closes. When the venue returns and resends, every re-delivered message collides with its journal row and books nothing, and any genuinely new message appends after the last position.

How much the venue can give back is bounded, and the bound is a calendar rather than a duration. Bitnomial keeps at most a week of resend history, and it is discarded at the venue's weekly Friday maintenance reset rather than aging out message by message. So what the session can actually replay is everything since the last reset -- minutes on a Friday evening, nearly seven days on a Monday.

Two things follow, and they are the ones that matter at 3am. An outage that spans the reset is not recoverable from FIX at all, however brief it was. And the same outage is cheap early in the week and expensive late in it, so the exposure is judged by the calendar, not by how long the session was down. Once that history is gone the only route back is Bitnomial's clearing records and the daily reconciliation built on them. (CME's equivalent window, for scale, is a flat 48 hours.)

If the venue gap-fills a range we were owed, that is a hole we cannot close: the consumer stops at it and stays stopped, and an operator has to accept the loss explicitly and on the record (btnl.clear_source_liveness, which writes who accepted what to an append-only table). The one thing that cannot happen is the pipeline continuing as if the range had been delivered.

What is deliberately not built

Reviewing the stack

Slice What to look for
1 schema The tables and their constraints. Nothing reads them yet.
2-3 A FIX latency guard that silently ate messages; shared helpers.
4 commit core The single-transaction commit that invariant 2 rests on.
5-6 The recorder, and the loop that feeds consumers.
7 FIX source Session ownership, epochs, the witness ledger. The largest slice.
8, 11 End-to-end tests against a real Postgres and a real FIX venue.
9 trade consumer Booking, pairing, the lone-half sweep, exactly-once.
10 venue mock A venue that can reset, gap-fill and resend, so the above is testable.
12-14 The two services, source liveness, deploy and operator docs.
15-16 The load harness and the capacity evidence.

Evidence

Capacity and load-test results are in docs/internal/operations/btnl-pg-pipeline-capacity.md and btnl-pg-pipeline-load-tests.md. Operating instructions are in docs/internal/operations/bitnomial-pg-pipeline-overview.mdx and the deploy runbook. The performance question that prompted the measurement is A-4623.