Skip to content

DDIR server: tail <export> as corgi, a columnar subscription - #914

Open
frankmcsherry wants to merge 1 commit into
master-nextfrom
ddir-columnar-tail
Open

frankmcsherry wants to merge 1 commit into
master-nextfrom
ddir-columnar-tail

Conversation

@frankmcsherry

Copy link
Copy Markdown
Member

tail <export> as corgi streams a columnar export as binary frames instead of text lines. It is the subscription a remote client would use. Plain tail and tail … as rows are unchanged.

Frames. Each frame is a header line <id> frame <kind> <nbytes>\n followed by exactly nbytes of payload. On TCP and stdin they follow each other on the stream; on WebSocket, each frame is one binary message holding both. Integers are little-endian.

Kind When Payload
schema once, before any data text: export=<name> encoding=corgi-container/1 time=u64 diff=i64 data=worker:u64-le,container
data per batch, per worker u64 worker, then one batch in the container format workers already exchange between processes (corgi/bytes.rs)
progress when the subscription's frontier advances u64 upper: every update before upper has been sent, from every worker

How it works.

  • Encoding. Each worker encodes its share of each imported batch straight from the export's Corgi chunks, with no rows built (Published::import_encoded, backend::corgi::export_containers), and sends it to worker 0. Workers hold disjoint keys, so a tick's data frames together are its changes.
  • Replay. The frames before the ok replay the export's compacted contents, as text tail does.
  • Row traces. A row trace has no columnar form, because an untyped variant has no shape, so as corgi on one is refused with err. Text tail reads any trace. Since DDIR server: columnar exports and export taps #913, corgi exports are columnar by default.
  • Slow clients. A client more than 256 MiB behind gets err client too slow and end, and that subscription stops sending.

Measured in one run at 4 workers, tailing the tour's scored export (tour with inspect removed, 50k nodes / 100k edges, 500 edges churned per tick). Client time includes receiving and decoding every row.

Text tail as corgi
Replay, 99,997 rows 9.83 MB, 66.2 ms 5.60 MB (−43%), 28.9 ms
Ten churn ticks, 10,000 rows 0.99 MB, 26.1 ms 0.57 MB, 23.2 ms

Tests. New server/tests/columnar_tail.rs, at 1 and 4 workers:

  • a subscription started mid-stream decodes frames that match text tail update for update, across retractions, lists and sums;
  • the decoded frames sum to peek after each tick;
  • as corgi on a row trace is refused.

cargo test --release -p interactive -p ddir-server: 141 passed, 6 ignored (the 5 ignored before this change, plus the measurement above). The WebSocket binary path has no test.

Open choices:

  • whether schema should carry declared DDIR types (Int vs F64, variant names), which are erased today;
  • progress only once per tick;
  • dropping a stopped slow subscription's dataflow; today it stays until stop or the end of the session, since dropping it on worker 0 alone isn't safe;
  • an Arrow adapter;
  • per-worker streams instead of gathering to worker 0.

🤖 Generated with Claude Code

`tail <export> as corgi` streams a columnar export as binary frames instead
of text lines: a header line `<id> frame <kind> <nbytes>` then the payload (on
WebSocket, one binary message holding both). Kinds: `schema` once (text),
`data` = u64 worker | one encoded batch in the inter-process container format
(`corgi/bytes.rs`), and `progress` = u64 upper once the subscription's
frontier passes it, so a client knows when a time is complete.

Every worker encodes its share of each imported batch without building rows
(`Published::import_encoded`, `backend::corgi::export_containers`) and sends
it to worker 0. Workers hold disjoint keys, so a tick's data frames together
are its consolidated changes. A row trace has no columnar form (an untyped
variant has no shape), so `as corgi` on one is refused; text tail reads any
trace. A client more than 256 MiB behind gets `err` and `end` and the
subscription stops sending. The session writer carries `Out::{Line, Frame}`.

Tests: decoded frames equal text `tail` update for update and `peek` after
each tick, at 1 and 4 workers, with retractions, lists and sums, starting mid
stream; `as corgi` on a row trace is refused. An ignored measurement on the
tour's `scored` export (100k rows, 500 edges churned per tick): replay 5.6 MB
vs 9.8 MB of text, 29 vs 66 ms to receive and decode.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant