RFC: Consuming EP3 executions directly from Redpanda

Date: 2026-08-28, revised 2026-09-10 Status: In review. Production uses the gRPC drop copy. Paired dev captures confirm every native field mapping. In production, EP3's drop-copy delivery accounts for about 69% of mean acknowledgement latency and most of its tail.

The problem

For most of 2026 we have been chasing order-acknowledgement latency on EP3. A client places an order, the InsertOrder RPC returns almost immediately -- and then the client waits for the ack. The wait is long enough that market makers have reduced or ceased trading over it.

The wait is structural. InsertOrder's response carries only an exchange order id (rs/order-gateway/src/rest_service.rs:641 uses it for nothing else; the synchronous-return fields are only populated when max_block_time is set on the request, and nothing in the codebase sets it). The actual acknowledgement -- the Accepted state transition, the OrderAcked WebSocket event, the og_ws_order_to_ep3_ack_ms measurement -- fires from the EP3 drop-copy stream, when an execution with ExecutionType::New arrives at rs/order-gateway/src/lib.rs:1541.

And the drop copy is not a direct line. Connamara confirmed (December 2025) that EP3's drop-copy service is itself a consumer of EP3's internal Kafka bus, a Redpanda cluster. We have since confirmed that directly: EP3's own drop-copy resume token is a plain vector of Kafka offsets, twelve pairs of a 4-byte partition and an 8-byte offset, one per partition of the events topic. The drop copy is a Kafka consumer with a serialisation layer bolted on, and we repeat that last hop four times over, because order-gateway, trade-engine2, risk-engine and marketdata-publisher each hold their own subscription.

The acknowledgement a client waits on is not the reply to InsertOrder. It is a message that was sitting on EP3's internal Kafka bus milliseconds after the match, and everything between that bus and our order-gateway is overhead we can remove.

Connamara supports reading the bus directly and documents the contract at https://connamara-ep3.readme.io/docs/consuming-directly-from-kafka. A probe on the branch michael/ep3-kafka-direct consumed both feeds side by side and put the expected p99 ack time below 10 ms. That figure was never reproduced, and it is not what this RFC rests on -- see "Where the time actually goes" below, which measures the same claim from production data.

The proposal in one sentence

Consume the events topic from EP3's Redpanda brokers in a shared ax-ep3 module that reassembles, deduplicates, and decodes the raw records behind the same subscription contract the gRPC drop copy serves today -- and switch consumers over one at a time, by per-service configuration, starting and possibly ending with order-gateway.

flowchart LR
    ME[EP3 matching engine] --> RP[("Redpanda events topic<br/>12 partitions, port 9092")]
    RP --> DC[EP3 drop-copy service]
    DC -- "gRPC CreateDropCopySubscription" --> OG[order-gateway]
    DC -- gRPC --> TE[trade-engine2]
    DC -- gRPC --> RE[risk-engine]
    DC -- gRPC --> MD[marketdata-publisher]
    RP -. "planned: Kafka fetch,<br/>ax-ep3 kafka module" .-> OG
    OG --> WS[WebSocket clients:<br/>OrderAcked, fills]

"Kafka" names the protocol, not a different server: Redpanda is a drop-in implementation of the Kafka wire protocol, which is why an ordinary Kafka client library talks to it natively.

Vocabulary

events topic. EP3's internal firehose: executions, order events, requests and market data all pass through it. Twelve partitions on the dev exchange. It is strictly read-only for us, and it carries far more than we need, so we filter client-side.

TransactionEvent. The envelope in every record value (connamara.ep3.wire.v1beta1, vendored at rs/ep3/api/protos/connamara/ep3/wire/v1beta1/wire.proto). It holds a trace_route and repeated google.protobuf.Any payloads; the type_url says what each one is.

The two execution-bearing shapes. This is the trap. Executions arrive as books.v1beta1.BookEvent and as orders.v1beta1.OrderEvent, which carry the same orders, executions, transact_time and cancel_reject fields. On a sample of live dev order flow, order events were the more common of the two. A consumer that decodes only book events looks perfectly healthy and silently drops about half the executions. The topic also carries orders.v1beta1.FindAndCancelOrdersRequest and trades.v1beta1.TradeCaptureReport; neither carries executions, and both are ignored.

The two Execution types. The drop copy delivers ep3.v1beta1.Execution: an execution with the full Order snapshot embedded (clord_id, cum_qty, leaves_qty, state, price) plus its own transact_time. An event carries the leaner orders.v1beta1.Execution: an order_id string pointing at the affected order carried once in the same event, its own order_status (the order's cum, leaves and state as of this execution), and no timestamp. The two Order types share exactly one field tag -- id -- because the lean one nests everything under attributes and status. Our consumer performs the same join EP3's drop copy does, because order-gateway reads the embedded order everywhere (cum_qty/leaves_qty at lib.rs:1438, replace resolution at lib.rs:1392). Connamara confirmed on 2026-09-10 that the per-execution status is a contract: every book-side execution clones the order's live status at creation, mutated for fills immediately before the clone, and executions are appended in the order they are applied to the book.

id header. Every record carries an 8-byte big-endian signed 64-bit id, globally unique, assigned by EP3 at publish time. Delivery is at-least-once, so this is the deduplication key. Ordering is per partition only.

Chunked message. Brokers cap record size, so EP3 splits any larger message across records carrying chunkedMessageId / chunkedMessageIndex (a byte offset, not a sequence number) / chunkedMessageTotalSize. Reassembly is copying bytes into a buffer of the stated size until it is full. Not an edge case: our own drop-copy client raises its gRPC frame limit to 64 MB (rs/ep3/src/client/admin.rs:506) precisely because single messages -- mass cancels, expirations -- get that large.

Sentinel record. After the final chunk of a split message EP3 publishes one ordinary record with the same id and an empty payload. Consumers skip any TransactionEvent with an empty payload. The sentinel reuses the split message's id deliberately, so what keeps it safe is not the order of the two checks but that chunks never feed id-based dedup at all: chunked records are deduplicated on chunkedMessageId, which leaves the sentinel's id unseen and lets it reach the empty-payload check on its own.

Offset vs resume token. The drop copy's cursor is an opaque resume_token with a known failure mode, vendor-confirmed in December 2025: a token older than retention is silently clamped to the start of what remains, with no error. A Kafka offset fails loudly with OffsetOutOfRange instead. Loud beats silent: a stalled consumer is a bad morning, a consumer that silently skipped a fill is a wrong position nobody finds for a month.

What the consumer must do with every record

flowchart TD
    R[fetch record from partition] --> C{chunk headers<br/>present?}
    C -- yes --> D2{seen chunkedMessageId?}
    D2 -- yes --> SKIP[skip]
    D2 -- no --> B[copy bytes into buffer at<br/>chunkedMessageIndex]
    B --> F{buffer full?}
    F -- no --> R
    F -- yes --> DEC[decode TransactionEvent]
    C -- no --> D{seen id header?}
    D -- yes --> SKIP
    D -- no --> DEC
    DEC --> P{payload empty?}
    P -- "yes (sentinel)" --> SKIP
    P -- no --> U{payload type?}
    U -- "any other type: carries no executions" --> SKIP
    U -- "BookEvent or OrderEvent" --> SYN[synthesize drop-copy-shaped<br/>execution batch]
    SYN --> OUT[downstream, unchanged:<br/>state machine, WS, ClickHouse]

What the diagram does not show is why each step is mandatory:

The same order, before and after

A client sends a limit order on EURUSD-PERP over the order-gateway WebSocket.

Today: order-gateway calls InsertOrder; EP3 matches or rests it and publishes the resulting event to its internal Redpanda bus; EP3's drop-copy service consumes that record, joins the order snapshot into an ep3.v1beta1.Execution, waits out its own batching, and streams a gRPC response to each of our four subscriptions; order-gateway's reader receives the batch, and only now records og_dropcopy_delivery_latency_ms, transitions the order to Accepted, and emits OrderAcked. Every hop after Redpanda re-reads and re-sends a message that already existed.

Proposed: the same record lands on the events topic; order-gateway's fetch (already long-polling, so the broker responds the moment data arrives) returns it directly; the module dedups, decodes the event in whichever shape it arrived, synthesizes the identical execution batch, and hands it to the same channel the gRPC reader fills today. The drop-copy service, its batching and the gRPC hop are gone; og_ws_order_to_ep3_ack_ms is the before/after scorecard.

Where the time actually goes

The whole value of this change is capped by one quantity: how long EP3's drop copy takes to tell us about a match. Everything before it (our validation, the network to EP3, the match itself) and everything after it (our state update, our acknowledgement) is untouched by switching transport. That quantity is measured in production, per execution, by og_dropcopy_delivery_latency_ms: EP3 stamps TransactTime when it matches, order-gateway stamps arrival, and the difference is the drop copy's contribution.

Over 24 hours of production, 14.9M executions and 6.3M accepted orders:

segment p50 p90 p95 mean
EP3 matching engine 0.025 ms
match to the events topic 0 ms 1 ms 0.1 ms
EP3 drop-copy delivery 2.32 ms 24.7 ms 158.6 ms 16.8 ms
our batch processing 0.03 ms 0.16 ms 0.42 ms
our processor queue wait 0.03 ms 0.12 ms 0.80 ms
our ClickHouse permit wait 0.001 ms 0.001 ms 0.001 ms
total order to acknowledgement 9 ms 124 ms 24.3 ms

Delivery is about 69% of mean acknowledgement latency, and our entire share of that same return path stays under a millisecond at p95.

The tail is the sharper finding. Half of all executions reach us in 2.3 ms, but the slowest 5% take 158 ms -- roughly seventy times worse -- while our own handling of those same executions never exceeds a millisecond. A tail that severe on the far side of the wire, and flat on ours, is what drags acknowledgement p99 to 303 ms. Production has never acknowledged an order in under 5 ms.

This is also a test the proposal could have failed. Had delivery come back at half a millisecond, the 24 ms would have been our own code and no change of transport would have helped. It did not.

What replaces those 16.8 ms is known: the record is on the topic 0.1 ms after the match, and reading it from inside the exchange's network costs about 1 ms.

Two caveats on the numbers. Delivery compares EP3's clock to ours, so the split between the outbound and return legs carries unknown clock skew; the tail argument is unaffected, because 158 ms dwarfs any plausible skew. And delivery is measured per execution while acknowledgement is per order, so their means and shapes compare but their percentiles do not.

Exactness requirement

Direct Kafka delivers the same execution and cancel-reject fields as native gRPC, or it fails with an error. It does not exclude fields, substitute end states, deliver partial batches, or log and skip. A position that can never be served again is different: the feed reports it as unservable, and the gateway cold starts, resyncs and counts the gap (see "Recovery paths" below).

Paired captures from EP3 dev confirm every field mapping: a historical comparison of 27,598 executions for the order and execution fields, and live captures of ten commission rule shapes for the four commission fields. Synthesis fails an event whose native rendering was not observed: spread and expression commission rules, and rate strings that do not parse.

What we measured

The measurements below have different scopes. Demo samples and a retained-dev report identified candidate mappings; the paired dev captures of 2026-09-09 confirmed them. Raw captures, comparison reports and the findings write-up live outside the repository (ep3-demo-capture-20260908/evidence-20260909/DEV_KAFKA_FINDINGS.md).

Why a library module, and why order-gateway first

A shared module in ax-ep3, behind the contract consumers already use: create_drop_copy_subscription keeps its request and response types, and the transport behind it is selected per service by env (EP3_EXECUTION_FEED=dropcopy|kafka plus a broker list). Call sites keep their request and response types, but the method now returns a boxed stream rather than tonic::Streaming, so the few helpers that named that type were retyped, and order-gateway gained explicit handling for the one failure only the direct feed can raise. This is also what makes migration gradual: every service links the same library, but each process picks its own transport, so flipping a consumer is a one-variable config change and rolling back is changing it back.

Not a relay service that re-serves the gRPC API: a relay is a new deployable and a single point of failure in front of every consumer, and the vendor checklist lives in one place either way.

Order-gateway is where the latency is felt, and it is the easiest correct migration: it skips block-trade executions outright, and keeps no durable cursor -- its token lives in a local, so a restart subscribes with an empty token, tails live, and repopulates open orders from a SearchOrders snapshot (sync_initial_open_orders, filtered to OrderStateFilter::Open). That maps exactly onto "start at latest offsets, resync on OffsetOutOfRange", including the gap it leaves: the snapshot restores resting orders, never terminal ones. trade-engine2 is the opposite: post-trade, latency-insensitive, and its Postgres-atomic resume-token machinery is the most carefully tested code on this stream. It should move late or never.

rskafka is the right client: pure Rust, no librdkafka C toolchain in a workspace whose builds already wedge on quickfix-ffi, and a manual partition/offset model with no consumer-group rebalancing to reason about.

One contract, two transports

The swap is signature-preserving because everything the drop copy does server-side is either emulable in the library or deliberately improved:

Two differences are deliberate:

Tests

flowchart LR
    subgraph ep3-mock
        MEQ[mock matching engine] --> EV[BookEvent or OrderEvent]
    end
    EV -- "production synthesis" --> GDC[gRPC DropCopyAPI<br/>existing suites, unchanged]
    EV -. "kafka mirror: id headers,<br/>chunking, sentinels" .-> TRP[("Redpanda testcontainer")]
    TRP -. fetch .-> KM[ax-ep3 kafka feed<br/>under test]

Deliberately not built

Questions for Connamara

Direct broker access answered most of the original list by measurement: events has 12 partitions (as does every EP3 topic), records carry no Kafka key at all, the readable window is time-based, expirations and unsolicited cancels both ride as ordinary executions, and cancel rejects ride on the event. What remains, most load-bearing first:

  1. Answered 2026-09-10, see below. Is Execution.order_status the exact status the native drop copy uses for that execution, and which affected-order fields can change within one event? The proto includes the status, and the retained-dev sample shows per-execution semantics. Native Order.create_time equalled Order.initial_order_receive_time on all 27,598 paired dev executions, including replacements. Reconnect scenarios are not observed.
  2. Answered 2026-09-10, see below. Is BookEvent plus OrderEvent the complete set of execution-bearing payloads, including settlement and rare lifecycle types? TradeCaptureReport, FindAndCancelOrdersRequest and FirmsNotification carry none. The feed skips every type other than BookEvent and OrderEvent.
  3. Answered 2026-09-11, see below. Records carry no key, so partition assignment is entirely producer-side logic. Is that mapping deterministic and stable across restarts and redeployments? Per-order sequencing rests on a symbol staying on one partition, and without a key we cannot verify that from the wire.
  4. Answered 2026-09-11, see below. What is the configured retention on events, in time and in bytes?
  5. Answered 2026-09-10 and 2026-09-11, see below. Delivery is at least once, so a record can in principle appear twice at distinct offsets. We have never observed a repeated id header on the topic; the consumer deduplicates within a session, but a duplicate whose original arrived just before a restart would be re-emitted, since offset suppression cannot cover a new offset. The consequence is bounded to a repeated client event -- quantities are taken from the order snapshot, not accumulated. Does the drop-copy service itself deduplicate producer retries, and does the producer retry in practice?
  6. The dev Kafka listener accepts connections with no credentials at all. We would like a principal restricted to DESCRIBE and READ on events with no WRITE, so the "never produce" rule is enforced by the broker rather than by our code. What is required for production?
  7. Withdrawn 2026-09-11. Delivery from the bus to a drop-copy subscriber has a p95 of 158 ms against a p50 of 2.3 ms, measured on our side over 14.9M executions in a day. Is that expected, and is it something you would want to address independently of our reading the bus directly?

Answers from Connamara (2026-09-10)

  1. Guaranteed. Every book-side execution clones the order's live status at creation time, and for fills the status (cum, leaves, state) is mutated immediately before the clone, so the status an execution carries is the post-this-execution snapshot. create_time is a literal copy of initial_order_receive_time in the drop-copy conversion, for every execution type. Executions are appended in the order they are applied to the book.
  2. Complete. Only books.v1beta1.BookEvent and orders.v1beta1.OrderEvent carry an Execution or a CancelReject; settlement and expiry executions, mass-cancel results and instrument state changes all arrive in them. Trade busts produce no execution: a consumer that needs busts must follow TradeCaptureReport state changes. Everything else on the topic can be ignored.
  3. Repeated trade-capture reports are lifecycle transitions, not retries. The trades server publishes a fresh TradeCaptureReport for the same trade on every change (one TYPE_NEW when the trade is created from a BookEvent, then one per transition: per-side acknowledgement, an admin state update such as CLEARED, a counterparty or metadata update). Each carries a new report id and a fresh header id while the embedded executions stay identical and the trade id is the same; the count depends on the clearing workflow. A report may also be published more than once; the report id and the trade id together are the idempotency key. Whether the drop-copy service deduplicates producer retries of executions is still open.
  4. A Redpanda ACL, outside EP3. Read-only access is configured on the broker (Redpanda authorization ACLs); once ACLs are on, EP3's own connection string needs credentials (Kafka connection string format, user and password authentication).

Follow-up, 2026-09-10 (asked after these answers): can records of other messages land between two chunks of one split message? Yes. One producer writes a message's chunks in order to one partition, but other producers' records can sit between them. This is why the direct feed's resume token carries two offsets per partition, the earliest chunk still being reassembled and the position everything has been delivered through: a single offset would either re-emit the whole events that landed between the chunks or lose the message.

Answers from Connamara, second round (2026-09-11)

  1. Duplicates. The cluster runs acks=all with idempotence, so a producer retry cannot commit a duplicate record; the id header exists for exchanges that run acks=1 without idempotence. Consumer-group rebalances can redeliver, which does not apply to a consumer with no group. The consumer keeps its id ring as defence in depth; the configuration is the guard. Note: enable_idempotence: true is pinned in the demo Redpanda config only; prod and dev rely on Redpanda's default, and the producer's acks is EP3's own setting.
  2. Busts are a distinct trade state with their own report (TRADE_CAPTURE_REPORT_TYPE_BUST, TRADE_STATE_BUSTED), never a republish of the NEW report. The four identical NEW copies per trade on dev were not explained.
  3. Retention is per-topic retention_bytes. The cluster configs set 5 GiB per partition on prod, demo and dev with no time-based limit, so Redpanda's default applies on time. At the dev capture's mean record size of about 1 KB that is roughly 5 million records per partition; at prod's rate, with most traffic on two partitions, a hot partition holds well under a day. The actual topic-level values (rpk topic describe events -c) are still to be read.
  4. Partitions. The partition of requests equals the partition of events. Symbols are sticky: even when partitions are resized a symbol always lands on the same partition, in order, which the exchange needs for price-time priority. Question 3 is closed.
  5. Topic recreation. EP3 creates a brand new events topic natively, and that is what a Redpanda migration does: devops performed one during the dev upgrade, which is why dev's partition 4 restarts at offset 0, and considers another prod migration likely. Broker-side consumer-group offsets are deleted with the topic; the vendor's advice to reset the group to earliest does not apply because the consumer stores its own position. Risk: a position that survives a recreation resumes silently inside unrelated records once the new topic grows past it. An operating rule prevents this (2026-09-15): the position is in process memory, and a migration happens only with the gateway stopped. The restarted gateway cold starts on the new topic. Nothing compares the topic's identity with the position.

Question 7 (drop-copy tail latency) is withdrawn: the transport decision is made.

Roadmap

  1. Infrastructure. Done for dev: an external listener with a reachable advertised address. Production needs the same, with credentials.
  2. Latency measurement. Done, and from production rather than a test environment -- see "Where the time actually goes". Note that a consumer measured from outside the exchange's network cannot answer this question: a polling consumer has no fetch outstanding at the broker for one one-way trip after each response, so records arriving in that window wait. Measured from a laptop that gap accounted for essentially all of a 58 ms apparent deficit against the drop copy, and it shrinks to about 1 ms inside the network. Any future comparison must run in-VPC.
  3. Order-gateway against a real EP3 Kafka feed. Done on ax-dev: the current build ran on the direct feed against ep3-dev, healthy, with its own state machine, WebSocket service and ClickHouse writes.
  4. Native field mappings. Done on paired dev captures (see "What we measured"); spread and expression commission rules remain unobservable there.
  5. Shadow observation on dev, demo and production, with gRPC authoritative, until the comparison is clean. Meet the coverage and failure requirements below before proposing a separate order-gateway cutover.
  6. Recovery paths under load, the ones no environment has exercised yet: restarting mid-flow without duplicating or losing executions, and resuming from a position retention has aged past. Both need order flow on an exchange we own.
  7. Second wave, optional and separately decided: marketdata-publisher and risk-engine.

Production shadow acceptance contract

The first production use of direct Kafka is observation only. Keep the serving gateway on gRPC. Run the observer separately, with its own subscription and resource limits, so that a comparator failure or backpressure cannot stop the serving feed. The observer does not submit orders, change exchange state, or send Kafka reports to downstream business logic.

Do not use the finite ep3-shadow-compare probe to accept a production shadow run. It excludes fields, discards some synthesis errors, overwrites repeated native IDs, does not compare cancel rejects, and accepts extra Kafka executions.

The continuous observer must:

The Kafka side uses mappings from the wire only. Do not copy a native value into a Kafka report, because that removes the independence of the comparison. An input with an unobserved rendering (a spread or expression commission rule, a rate that does not parse) fails synthesis, and the observer reports the run as failed, not as coverage. The feed skips payload types that it does not decode.

Before an unattended run, test that the observer detects:

Each defect must produce a visible failed or unverified interval. A rollout decision must cite report volume and lifecycle coverage, including a replacement after earlier fills, and recovery. Time with no alerts is not sufficient. Observe production orders passively; run directed lifecycle tests in an authorized test environment. Cutover is a separate, explicit decision.

Where to look