Skip to content

The Series Facility

Beneath a stream's window is a place a read can reach rather than a wall. This is the design of that tier: what it adds to the stream surface, what a storage adapter must answer, and which decisions are closed.

The core surface is delivered: the stream verbs, the Series port beneath them, the aggregate vocabulary the stream surface and the guest capability wire share, and the guide chapter that teaches all of it. What is not delivered is an implementation of the port — beneath-window storage is an application's decision the way a journal copy's format is, so this page describes a contract and the obligations it places on whoever answers it.

Make beneath-window stream history a served read instead of a refusal, with bounded display reads, durable tiered storage, and a data lifecycle — without adding a third subsystem.

  1. No third subsystem. The journal already serves beneath-tail reads through registered storage. Streams gain the mirrored seam: a series storage adapter — itself a stream consumer, and written by the application — serves beneath-window ranges. One read surface, tier-transparent, honest floor evidence either way.
  2. Two coordinates, never conflated. Position orders, resumes, and accounts for gaps. A stream may declare one u64 monotone axis field for domain reads and display grids; physical time rides as an ordinary field. Nothing orders by float.
  3. The budget flows end-to-end. Every display read carries a point budget from the caller to the adapter. Keel wrote a decimation tier no read could reach because every call site passed usize::MAX; that failure mode is structurally excluded here, because the budgeted verb has no unbudgeted form.
  4. Raw is truth; buckets are display. Raw reads and budgeted reads are separate verbs with separate return types, so exporting decimated data by accident is a type error rather than a code-review catch.
  5. The producer never pays. publish stays synchronous and infallible into the bounded window. Sealing, decimation, and retention run on a Resident in the work layer. A barrier's cost is its syncs rather than its rows, and that cost belongs to a supervised service, never the publishing thread.
  6. Loss and lifecycle are observable. Floors advance visibly, gaps carry counts, sealing returns the durable position, and retention's effect is reported as a moving floor rather than left to an operator's cron job.
  7. No application policy in the mechanism. No name-based heuristics (keel shipped field.contains("discontinuity") in its generic crate), no application vocabulary. Time alignment, correlation, and domain statistics stay application-side.
  • Track — one typed component of a partition tuple, and the durable name it carries. Prefix — the compile-checked relation between a selection and a declaration: the whole tuple, or any shorter one. Tracked — the row declaration that names the application-typed partition tuple a stream carries (an execution and the probes within it; a cycle, a run). Storage partitions by it, floors are per-partition, and retention can hang off a partition's retirement. Untracked streams stay the default and the common case.

    A track is not a RequestKey (idempotency) and not a Focus (history predicate). The census forbidden list carries StreamKey, Channel, and Tag so the decision cannot be reopened by accident: a declared tuple at most four components deep, never a general label space.

  • Axial — the row declaration that names the one u64 monotone domain coordinate its rows are placed on. Only an axial stream can be surveyed.

  • survey — the budgeted display read. Answers buckets, never rows.

  • seal — the explicit durability barrier; answers the position through which rows are durable. Dokime lost 1,200 rows to the absence of this verb in keel, where it landed late as flush().

  • Bucket — what one numeric field did across one bucket of the axis. See Buckets Are a Fold.

Both coordinate declarations are optional, and both are declarations of the row rather than of its catalog membership — a track is a partition dimension and an axis is a domain coordinate, and neither depends on which application publishes the row. #[harmos::stream] generates them from a field the row already declares:

#[harmos::stream(
id = "capture",
window = 4096,
track = (run, probe),
axis = ordinal,
)]
#[derive(Clone, Deserialize, Serialize)]
pub struct Capture {
pub run: Run,
pub probe: Probe,
pub ordinal: u64,
pub value: i64,
}

The declared fields' own types become Tracked::Track — here (Run, Probe) — so the row stays the one place its shape is written down. Naming the fields is the declaration; nothing reads a field name for meaning. Depth is part of the declaration, not a convention a caller upholds: a track one component deep is track = (probe), and the attribute refuses a bare name, an empty tuple, a repeated field, a field the axis already claimed, and a fifth component.

Keel's answer to the same question was a string an application concatenated, and axon's migration is what proved it insufficient: one execution publishes a dynamic set of probes, and run + "/" + probe puts the separator convention, the escaping, and the parse back in the application. Here the pair is the declaration.

Levels, retention, and cadence are storage-adapter configuration at assembly, not definition concerns. Keel split its knobs across two crates that never saw each other, and one bound knob meant three things; neither mistake recurs.

The verbs:

streams.query::<Capture>(range).track((run, probe)).as_of(at).await?; // rows
streams.query::<Capture>(range).track((run,)).await?; // the run
streams.survey::<Capture>(axis).track((run, probe)).budget(400).await?; // buckets
streams.survey::<Capture>(axis).fold(Aggregate::Mean).await?; // statistics
streams.seal::<Capture>().track((run, probe)).await; // Position

.track(…) exists only on tracked streams, so selecting a partition of a stream that declares none does not compile rather than refusing at runtime. It takes the declared tuple or any prefix of it: the whole tuple is one partition, a shorter one is every partition beneath it. A tuple of the wrong component type or the wrong depth implements no Prefix of that declaration and fails to compile — there is one spelling per selection, tuples including the one-tuple (run,), because a bare component would be a second way to write the same narrowing. What a prefix cuts the answer by is the application's own equality, component by component; what it hands the served tier is the component labels, and nothing above the adapter builds a name out of them. A survey is taken two ways and no third: budget is the only way to take buckets, and fold is the only way to take a statistic. Neither has a form that omits its argument, so there is no unbudgeted display read and no defaulted aggregate.

query reads the window first and takes its floor out of the same snapshot that produced the rows. Everything above that floor is the window's answer; everything at or below it is asked of the served tier, capped at the same floor. Exactly one tier owns the boundary position and neither tier sorts — both are ordered by construction.

Reading the floor separately would be the bug: eviction between the two reads would duplicate or drop the boundary row. Positions are not unique per row, so the floor's own stamp is partially evicted by definition and belongs wholly to the tier beneath.

A window is stood on a floor at boot, not only by eviction. A restarted runtime opens empty rings while the tier beneath them still holds everything the previous run sealed. A window whose floor only ever came from its own eviction would therefore answer a full query out of an empty ring and call the durable history missing — and that answer is not a short one a caller could detect, it is a complete-looking one. So boot asks each attached series copy the position it is durable through, per declared stream, and stands that stream's ring on it before the runtime is ready. The first query after a restart splices, on the ordinary boundary and with no eviction required: at or beneath the watermark is the tier's, above it is the ring's. The watermark is read through Series::seal, which has nothing staged to take at boot and therefore costs no barrier.

Two conditions keep that floor honest, and each is a way the splice would otherwise lie. A tier that reports the origin — nothing attached, nothing durable yet, or a stream the adapter does not select — leaves the window standing on nothing, so a full query is answered from the ring rather than refused for a history no tier has. And a watermark beyond the frontier the boot recovered means the two tiers no longer share a coordinate: the rows came back further than the journal did, so a splice on that position would cut history where the state never stood. The honest degraded mode there is the window alone.

The window is one ring per stream whether or not a track is declared. Tracks partition durable storage and select reads; they never multiply the boot-frozen table, because a track is runtime instancing and the table is closed from start onward. A tracked read is cut by the application's own component equality; the labels of the selection are what narrow the read the adapter serves.

BeneathWindow says which wall a read hit, because each is a different next move:

VariantMeaningThe caller's move
UnservedNo series adapter is attachedAttach one at assembly
Floor { floor }No tier retains anything beneath floorNarrow the range
Unreadable(_)The adapter could not read what it holdsReport, retry

Unserved is the behavior that existed before this facility, and it remains the honest degraded mode rather than an error state: with nothing attached, the window is the whole served history.

Series<A> is Source<A> mirrored onto the row feeds — one designation, one owner, one dispatch idiom:

pub trait Series<A: Application>: StreamConsumer<A> {
fn rows(&mut self, stream: Definition, track: &[String],
start: Bound<Position>, end: Bound<Position>)
-> impl Future<Output = Result<Vec<At<Arc<dyn Any + Send + Sync>>>, BeneathWindow>> + Send;
fn buckets(&mut self, stream: Definition, track: &[String],
start: Bound<u64>, end: Bound<u64>, budget: usize)
-> impl Future<Output = Result<Vec<Bucket>, BeneathWindow>> + Send;
fn aggregate(&mut self, stream: Definition, track: &[String],
start: Bound<u64>, end: Bound<u64>, kind: Aggregate)
-> impl Future<Output = Result<Vec<Aggregated>, BeneathWindow>> + Send;
fn seal(&mut self, stream: Definition, track: &[String])
-> impl Future<Output = Position> + Send;
}

track is the component labels of a selection, outermost first, and an empty slice narrows nothing. This is the one boundary a typed track becomes text at. How those labels are stored — under what key, sorted how — is the adapter's own, and no layer above it sees or builds that form.

Every series adapter is first an ordinary StreamConsumer: it is fed the same delivery runs as every other recording and keeps its own Cursor. What makes it a series copy is that it can hand those rows back. What the mirror does not carry over from Source is durability designation — rows are dropped honestly rather than waited for, so a series copy feeds no frontier and gates no commit.

Rows come back erased. An adapter selects each stream by its concrete row type at assembly, which is exactly what lets it restore that type without the contract naming it — the same trick delivery plays in reverse.

Answering short is never permitted. An adapter that does not hold the whole range reports its floor; one that cannot read what it holds says so, rather than reporting a floor a caller would read as retention.

Buckets come back grouped by field in the order the adapter declares its fields, and ascending by first.axis within each field. A flat Vec with an unstated order would make a fold over levels depend on adapter happenstance; series-agg folds runs of adjacent buckets, so adjacency has to mean something. Adapters owe this ordering; readers may rely on it.

An adapter produces it by construction rather than by a pass over the answer: key every fold by (field index, grid cell), and because a cell's coordinate is its axis, laying the key out in order is the ordering. Store the declaration order durably beside the levels, so the grouping a survey answers with cannot drift with the build that serves it.

Storage::series(adapter) designates it, zero or one. Naming a second demotes the first to an ordinary stream recording — still fed, still writing, simply no longer the copy reads are served from.

A read crosses a bounded lane to the supervisor that exclusively owns the adapter, waits on a single-use reply, and is answered only after ready feeds have been drained. This is the journal's historical-read seam, mirrored rather than reinvented, and it inherits that seam's property: delivery has priority over reads, so a survey can be starved under sustained publish.

How short that starvation stays is the adapter's half of the bargain: a drain that stages rows in memory and touches no file puts the supervisor back in its loop in microseconds, and a read is then never behind an fsync. A barrier — which is behind fsyncs — runs on the same supervisor, so the sealing cadence is the knob that trades a survey's tail latency against a crash's staged-row loss. An adapter that does I/O in its drain moves both numbers the wrong way, and nothing above it can compensate.

The lane's send half lives on every Streams clone and outlives the supervisor. Once the supervisor is gone nothing is served, so a failed send and a dropped reply both mean Unserved — never a hang.

A bucket carries true M4 per numeric field — first, last, min, and max, each with the axis coordinate it occurred at — plus count, sum, and sum_sq.

The extremum coordinates are what make a bucket placeable: an envelope that knows a peak's height but not where it happened cannot draw the marker.

The three accumulators are what make the rest of the statistics a fold rather than a raw scan. Mean is sum / count, and variance comes out of the same three, so two adjacent buckets combine without going back to the rows. A stored mean cannot be merged, so none is stored.

Bucket::merge is the law this shape exists for, and it is associative: folding a run of rows gives the same bucket as merging any grouping of its sub-buckets, in the same order. Order is not free — first and last are the leftmost minimum and rightmost maximum of the axis, so an axis tie keeps the earlier row — which is exactly what associativity promises and commutativity would not. Merging buckets of different fields is an authoring error and panics.

Floating addition remains floating addition: the accumulators merge exactly in the sense that no statistic is approximated away, not in the sense that summation stops rounding. The unit tests state the law over exactly representable values so a test of mergeability is not a test of rounding.

This is also the structural argument for the LTTB non-goal. A reduction whose bucket boundaries depend on the requested output count, and which carries running state across those boundaries, has no merge at all: two of its buckets cannot be combined into the bucket of their union, so it can never be a stored level and never folds from the level beneath it. It belongs on the application side of a raw read, permanently.

Buckets come from the served tier alone — the window retains rows, not levels — so a survey sees what a seal made durable. Seal, then survey.

A bucket is four placed points, so a point budget is a bucket count, and every measured field answers its own run of buckets out of the same budget:

cells per field = budget / (4 × measured fields)

The adapter answers from the finest built level whose cells over the requested range fit that count. Choosing the coarsest would satisfy every budget trivially and make every finer level unreachable — which is precisely keel's decimation tier that no read could get to, and the failure principle 3 exists to exclude. The cell count is bounded from the requested span rather than from the segments, so choosing a level costs no read; a bound that overestimates picks a coarser level than strictly necessary, which spends less than the ceiling, and spending less is always allowed.

When even the coarsest built level would exceed the ceiling, adjacent cells of it are merged until they fit. That is the same associative fold that built the level, so the answer is exactly the bucket a still-coarser level would have held — and the budget is never exceeded, because the budget is a ceiling.

The unit of an answer is a whole cell, so a survey admits every cell its range touches and a boundary bucket folds rows just outside it. That is the display convention and it is deliberate: clipping a cell would mean a bucket that is no longer the fold of a grid cell, which is exactly the property the merge law and every stored level rest on. A caller who needs the range exactly reads rows.

What a measured field is, is worth naming. The adapter folds every numeric column a row declares except those it declared as coordinates: Axial::AXIS names the axis and Tracked::TRACK names every track component, and all of them are excluded by their declaration rather than by their spelling. A component carried as a number is the common case, and it is constant within its own partition — which is what an adapter stores a run of rows under — so folding it would answer a statistic of a repeated value and spend budget doing it. Both consts exist for exactly this: a declared coordinate is not a measure.

No implementation of this port ships. This section is what one owes, written down so that two adapters answering the same question answer it the same way — and so that an application writing its first one knows which of its choices are its own and which are the contract's.

A partition tuple becomes text exactly once, at the adapter's own boundary, and the encoding is the adapter's to choose. Two properties are not:

  • Encoded order is tuple order. Comparing keys byte by byte must be comparing tuples component by component, and a shorter tuple must sort immediately before everything beneath it.
  • A prefix scan is a prefix selection. The keys beginning with one tuple's encoding must be exactly the partitions whose track begins with that tuple — run_70 is not beneath run_7.

An encoding that satisfies both is why a narrowed read can be a range scan that stops at the end of its subtree rather than a filter over everything stored. The wire format page writes down one canonical form that has both properties — each component's label followed by the unit separator U+001F, including the last — so a reader in another language can order and select partitions the way the adapter that wrote them does. An adapter that invents its own owes the two laws anyway.

Which is why a label carrying a control character is refused rather than stored: it would not mis-sort, it would mis-select. Everything else an application may write is a label — dotted probe paths, hyphenated run names, text in any script.

Within a partition, rows must come back in publish order, and an unnarrowed read must come back in publish order across partitions. An adapter that stores an ordinal beside each row gets the second from a linear merge of runs that are each ordered by construction, which is what lets it partition by track — per-track floors and track retirement both need that — without a read path that sorts.

Whatever an adapter's commit point is, the claim is the record. A watermark, a floor, and a level ladder are derived from what was durably recorded rather than advanced beside it, so a crash between writing a file and recording it cannot leave a claim the data does not support. An adapter that advances a claim first has a boot that lies.

A torn record is a recovered write, not a bricked boot: read the trusted prefix, stop at the first entry that fails its own check, keep the failed tail as evidence rather than deleting it, and repair before claiming anything new. Heal before claim. Every intermediate a crash leaves — a temporary that never claimed its name, a file that claimed one and never got recorded — is kept as evidence rather than adopted on a guess, and the next name is claimed above everything the store has ever held, live or retired, so a claim never lands on a name a crash left evidence under.

A barrier may cover more than it was asked about, and that is not a lie: seal answers only for the stream and track named, and what it made durable being more than that costs nothing to report honestly. Sealing only what a caller named would pay the commit tail once per stream a cadence covers and buy nothing — rows left staged are rows a crash still loses.

A stored digest earns its place by being verified on read. A re-read at seal proves nothing the write already knew; a check on the way out is what stands between a flipped byte and a statistic nobody can trace.

Stored summaries sit on a global axis grid anchored at zero — never per-segment chunks, because chunk-aligned buckets shift with the size of the thing that rotated them and two of them cannot be merged into the bucket of their union. Level zero folds the raw rows; every level above it folds the level beneath by Bucket::merge, so a stored level equals a scan of its rows in any grouping. That is the property every budgeted answer rests on, and it is the one an adapter cannot get wrong quietly.

How long each tier survives is the adapter's own declaration, made where the adapter is built. What the contract fixes is that floor advance is the observable outcome, and that a floor is derived from retirement evidence and from nothing else. Raw rows and folded levels retire on their own lifetimes and answer their own floors, so a survey is never refused on a floor its levels still hold. An unnarrowed read takes the highest floor of any track, because a range is answered whole only where every partition of it is.

Retiring a whole partition is the one tracks exist for. A track is an execution, a cycle, a run — something that ends — and the only evidence of its end that reaches a storage adapter is that its rows stopped. A duration-based sweep must never retire a track's newest live segment, because a sealed position is derived from what is stored and a track with nothing live has no honest position to report; retiring a track outright is the exception, and it is the exception on purpose.

publish stays synchronous and infallible; a drain should stage rows and touch no file, because that is what keeps the supervisor back in its loop in microseconds and a read from ever queueing behind an fsync. Turning staged rows into stored ones is a barrier, and a barrier is priced work — so it belongs to a supervised service:

let mut ticker = tokio::time::interval(Duration::from_millis(250));
while context.until_stopping(ticker.tick()).await.is_some() {
context.streams().seal::<Capture>().await;
}

The resident owns the cadence; the work runs where the adapter is owned, across the same bounded lane every read crosses. There is no second handle and no second writer, which is what the core's dispatch seam requires and what keeps a barrier from racing a drain. A tick should close one barrier however many streams it names: sealing them in turn costs a commit tail apiece, and under sustained publish every crossing after the first finds rows drained in behind it, so the syncs multiply with the streams instead of amortizing across them. A caller may still take a barrier itself with streams.seal(); the cadence is what means nobody has to. A clean stop closes the last barrier through StreamConsumer::flush, which is the only durability hook a stop offers — seal is never called on the way out.

  • An interior gap is not detectable at this seam. A barrier can record the rows the delivery cursor crossed beside the rows the adapter took, as durable evidence, but it cannot convert a shortfall into a floor: the cursor counts every stream's rows, so a shortfall proves loss only for an adapter that selected them all. A consumer gets no per-stream miss count from the StreamConsumer contract — Advanced::missed reaches watchers only — so rows a window evicted before delivery reached the adapter leave a hole it cannot see. Keeping a drain free of I/O is what makes that window small.
  • Rows staged when a process dies are gone, honestly and by design. The sealing cadence is the size of that risk window, and it is declared.

One enum, one parse table, one validity rule, one implementation of the arithmetic — reached by the stream surface and by the guest capability wire alike. Keel carried three vocabularies of this, its wire answered eight of the sixteen kinds its facade did, and a request it could not parse silently became count. None of those is representable here.

Aggregate is that enum, and it lives beside Bucket because the fold law is what makes half of it cheap:

KindNamesReadsAnswers to
Level foldsmin, max, mean, sum, countthe finest stored levelthe folded floor
Raw foldsp0 … p100the rows the range holdsthe raw floor

The first five come out of the three accumulators a stored level already carries, so they cost a level read and no rows at all. A percentile has no mergeable accumulator — that is the same structural fact that keeps LTTB out — so it reads the values themselves and is priced as a scan. folds_levels is where that split is written down, and it is not advice: the two halves answer to different floors, so a range whose raw rows retention has reclaimed still answers a mean out of the levels that outlived them and refuses a p95. Each says which.

Three properties hold the vocabulary to one language:

  • One name per value. name and parse are inverses. The median is p50 and has no second spelling, because a second spelling is how three vocabularies start.
  • Nothing defaults. parse answers nothing for a name outside the table — an unknown word, or a rank that is not a percentage — so the boundary that received it refuses and says which name it could not honor. A rank outside 0..=100 built in Rust is an authoring error and panics, exactly as merging two fields' buckets does; a value that arrived over a wire is a refusal, because a wire is not an author.
  • One arithmetic. folded reads a statistic off a merged bucket and scanned reads one off raw values. Both live on the vocabulary, so the facade path and the wire path cannot drift into two answers for one name.

A level fold covers the whole cells its range touches, exactly as a survey does and for the same reason — it is folded from the same stored cells. A percentile covers the range exactly, because it read the coordinates. A caller who needs a cell-exact fold of the first kind reads rows.

  • mise run check and mise run test green at every package boundary; the census test updated in the same commit as any surface change.
  • The two track-key laws stated over a generated set of tuples rather than over examples: sorting keys equals sorting tuples, and a key begins with another exactly when its tuple begins with that tuple. A label that would break either is refused where it is encoded.
  • Crash interleavings for the seal path, each produced by the real two-phase write rather than planted before a boot (axon's "crash test" planted a file and never crashed anything): post-write/pre-rename, post-rename/pre-record, and a record torn mid-entry. Each reopens, quarantines what it found without deleting it, serves the barriers beneath it, and takes the next one.

The last two are bars an adapter clears in its own suite. The core's bar is the surface above them: that the splice, the refusal vocabulary, the budget, and the aggregate arithmetic hold whether an adapter is attached or not.

One cost statement is worth carrying into an implementation rather than rediscovering. A barrier is priced in syncs, not in rows. Packing levels into fewer files does not move it, because the count that prices a barrier is fsyncs rather than segments; what group commit buys is the commit tail paid once per tick instead of once per stream, which grows with the assembly and leaves the per-segment data sync — the larger half — untouched. Group commit changes the slope, not the intercept. A read across the seam is CPU-bound in typed deserialization for as long as the serving contract answers concrete row types, so mapping a segment rather than copying it changes the shape of a read — no buffer proportional to the segment — rather than its rate.

Dynamic stream registration (the table stays boot-frozen; tracks carry runtime instancing), general tags or labels (cardinality and cost-model honesty; a declared tuple at most four components deep, never an open label space), LTTB or any viewport-dependent reduction (see Buckets Are a Fold for why this one is structural), cross-stream alignment and correlation, redb or KV backends, and migration of keel's on-disk stream data (consumers repin pre-1.0).

Stokker Technologies markDesigned and built by Stokker Technologies