DDIR server: tail <export> as corgi, a columnar subscription - #914
Open
frankmcsherry wants to merge 1 commit into
Open
frankmcsherry wants to merge 1 commit into
frankmcsherry wants to merge 1 commit into
Conversation
`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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
tail <export> as corgistreams a columnar export as binary frames instead of text lines. It is the subscription a remote client would use. Plaintailandtail … as rowsare unchanged.Frames. Each frame is a header line
<id> frame <kind> <nbytes>\nfollowed by exactlynbytesof 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.schemaexport=<name> encoding=corgi-container/1 time=u64 diff=i64 data=worker:u64-le,containerdatau64 worker, then one batch in the container format workers already exchange between processes (corgi/bytes.rs)progressu64 upper: every update beforeupperhas been sent, from every workerHow it works.
Published::import_encoded,backend::corgi::export_containers), and sends it to worker 0. Workers hold disjoint keys, so a tick'sdataframes together are its changes.okreplay the export's compacted contents, as texttaildoes.as corgion one is refused witherr. Texttailreads any trace. Since DDIR server: columnar exports and export taps #913, corgi exports are columnar by default.err client too slowandend, and that subscription stops sending.Measured in one run at 4 workers, tailing the tour's
scoredexport (tour withinspectremoved, 50k nodes / 100k edges, 500 edges churned per tick). Client time includes receiving and decoding every row.tailas corgiTests. New
server/tests/columnar_tail.rs, at 1 and 4 workers:tailupdate for update, across retractions, lists and sums;peekafter each tick;as corgion 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:
schemashould carry declared DDIR types (IntvsF64, variant names), which are erased today;progressonly once per tick;stopor the end of the session, since dropping it on worker 0 alone isn't safe;🤖 Generated with Claude Code