Market data publisher latency statistics

As built: 2026-09-23, commit 185f18d7f

Summary

The market data publisher records two latency histograms, one counter, and one gauge for the EP3 market data feed. They measure:

The design follows Order gateway latency statistics. The histograms share bucket boundaries and the _window naming convention from rs/sdk-internal/latency-series. The publisher aggregates directly into fixed-size BoundedHistogram instances and renders them alongside its other Prometheus metrics. The collector exports these over OTLP to the observability ClickHouse. Nothing is written to the exchange database.

Series

series type labels start end
md_ep3_to_recv_ms histogram none EP3 transact time of the update EP3 monitor task takes the update off the subscription stream
md_recv_to_ws_ms histogram event: book_l1, book_l2, book_l3, ticker same arrival session writer's socket write of the derived event returns; one sample per session that writes it
md_ep3_updates_total counter symbol
md_symbol_subscriptions gauge symbol

Neither histogram has a symbol label. With 155 buckets, a symbol label multiplies the exported series by the number of instruments. The per-symbol dimension is on the counter and the gauge.

md_ep3_to_recv_ms

This series compares two wall clocks: EP3's and the publisher's. Clock skew moves the whole distribution. A negative difference is recorded as zero.

The first update for each symbol on a subscription is the book snapshot. Its transact time is the time of the book's last change, which can be hours old. The monitor task records no sample for it. This applies to the first subscription and to every resubscription after an EP3 reconnect.

An update with no transact time produces no sample.

md_recv_to_ws_ms

The monitor task records a monotonic instant when the update arrives, before it updates the cache. The instant travels with the update on the cache broadcast. A session handler that derives a book update or a ticker from the update sends it to the session's outbound queue as a FeedEvent that carries the instant. After the session writer's socket write returns, the writer records the sample under the event's label.

The span therefore includes the cache update, the Redis write, the broadcast, the per-session event build, the wait in the session's outbound queue, and the socket write.

Events with no EP3 arrival produce no sample: heartbeats, trades, candles, subscribe snapshots, and tickers that an estimated funding refresh triggers.

The publisher closes a session whose outbound queue is full, and counts the closure in md_pub_slow_consumer_evictions_total. The events that session had not written produce no samples.

md_ep3_updates_total

The monitor task increments the counter for every subscribed symbol's update, including snapshots. Counter handles and snapshot flags are initialized from the instrument catalogue for each EP3 subscription. The per-update counter path does not allocate or build a label set; unexpected symbols are rejected before they can create metrics state.

message_rate remains the total rate across symbols.

md_symbol_subscriptions

ConnectionManager increments the gauge when a session first subscribes to an available symbol, and decrements it on unsubscribe. When the manager is dropped at the end of a session, it decrements the gauge for every symbol still subscribed. Unknown/unavailable symbols return an error without creating a subscription or metric label. Candle subscriptions are not counted.

ws_clients remains the number of connected sessions.

Memory bounds and load behavior

There are exactly five latency accumulators: ingress and the four outbound event kinds. Each has 16 fixed shards containing 156 exclusive bucket counters (including overflow), a sum, and a ring of 61 one-second min/max aggregates. Storage is under 240 KiB total for the five accumulators, independent of message rate, subscriber count, uptime, and scrape availability. Per-symbol metrics are bounded by the instrument catalogue.

Every observed latency updates a bucket and its window extrema. There is no sampling, raw-sample queue, full-buffer overwrite, or try-lock-and-drop path. Values beyond the largest finite bound still contribute to +Inf, count, sum, and maximum. Histogram quantiles retain the existing bucket-resolution approximation.

Writers select a shard once per thread and use a short lock to update its aggregates; a scrape copies each shard under that lock and formats the copy outside it. This avoids per-event metric registration and label allocation. More writers than shards can contend, but contention does not discard observations.

The _window summary exports only quantiles 0 and 1. Its extrema retain observations for 60–61 seconds: expiry rounds up to a one-second boundary, so a recent spike is never evicted early. Its count and sum remain cumulative. No upkeep task or scrape is needed to bound memory.

Limits

Test coverage

Unit tests cover the _window names, clock-skew clamping, rejection of arbitrary symbol labels, and subscription gauge cleanup on unsubscribe, disconnect, reconnect, cancellation, and task abort. The shared crate tests bucket boundaries and overflow, window rotation and idle expiry, concurrent writers and scrapes, preservation of tail spikes, and Prometheus rendering.

The integration suite drives EP3 updates through L1, L2, and L3 subscriptions over the real WebSocket stack. It asserts that both series render as histograms with the expected labels, that both _window series render as summaries, and that the counter and the gauge carry a symbol label.

Code map

file contents
rs/sdk-internal/latency-series/src/lib.rs LatencySeries, the bucket set, Prometheus builder
rs/sdk-internal/latency-series/src/bounded.rs fixed-memory histogram shards, min/max ring, Prometheus rendering
rs/marketdata-publisher/src/in_situ_latency_measurement.rs the series constants, event labels, per-symbol counter and gauge
rs/marketdata-publisher/src/tasks/monitor_market_data.rs arrival instant, snapshot exclusion, md_ep3_to_recv_ms
rs/marketdata-publisher/src/handlers.rs attaches the arrival instant to derived events
rs/marketdata-publisher/src/lib.rs FeedEvent; the session writer records md_recv_to_ws_ms
rs/marketdata-publisher/src/connection_manager.rs md_symbol_subscriptions