Skip to content

Serve Beneath the Window

Chapter 6 left one refusal unfinished. A stream's window retains a row count, eviction is how it stays bounded, and a query that reached beneath it answered BeneathWindow — honestly, but with nowhere to go next.

Chapter 13 already solved the same problem for entries. A read beneath the journal's tail is served by an attached recovery source, through the supervisor that owns it, and the verb never changed. This chapter is that answer mirrored onto the row feeds: implement the series port, attach the adapter, and the range the ring no longer holds becomes a place a read can reach. The workspace ships no implementation of that port — this chapter is the contract one answers, and every verb above it is unchanged whether you have written one or not.

Nothing above it moves. publish is still synchronous and infallible, query is still the same verb with the same iterator, and a runtime with no adapter attached still answers BeneathWindow::Unserved — which is not an error state but the honest description of a runtime whose window is its whole history.

A served tier needs two things a bounded ring never did: somewhere to place a row, and somewhere to partition it. Both are declarations of the row, and both name a field the row already has.

streams/sim_sample.rs
/// One simulation run, as the application names it.
#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)]
pub struct RunId(pub u64);
impl harmos::Track for RunId {
fn label(&self) -> String {
format!("run-{}", self.0)
}
}
/// One probe within a run.
#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)]
pub struct ProbeId(pub String);
impl harmos::Track for ProbeId {
fn label(&self) -> String {
self.0.clone()
}
}
#[harmos::stream(
id = "sim_sample",
window = 4096,
track = (run, probe),
axis = step,
)]
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
pub struct SimSample {
/// The run this sample belongs to.
pub run: RunId,
/// The probe it was measured at.
pub probe: ProbeId,
/// The solver step it was taken at.
pub step: u64,
/// The measured potential, in millivolts.
pub millivolts: i32,
/// The solver residual at that step.
pub residual: f64,
}

axis = step names the one u64 monotone coordinate these rows are placed on. It is not Position. A position orders rows against the journal, resumes a recording, and accounts for gaps; an axis is the domain coordinate a chart is drawn against and a range is asked in. A row carries both, and conflating them is how a display ends up plotting a bookkeeping number. Physical time is neither — if the samples carry a clock, it rides as an ordinary field beside the axis.

track = (run, probe) declares the partition as a tuple. An adapter partitions by it, floors are per-partition, and retention can hang off a run's retirement. It is worth being precise about why net would be a wrong component here: a track is something that ends — an execution, a cycle, a run, and the probes within one — and a net never does. A track is not a label space, and a declaration reaches at most four components deep, so its cardinality stays something you can reason about instead of a tag system whose cost nobody can predict.

Naming the fields is the whole declaration. Nothing reads a field name for meaning, and each field's own declared type becomes that component's type, so the row stays the one place its shape is written down. A track one component deep is written track = (probe) and is still a tuple — one selection, one spelling. Both declarations are optional: a stream with neither is the common case and behaves exactly as chapter 6 described.

There is no series adapter to reach for. Harmos publishes the port — Series<A> — and no implementation of it, because where beneath-window rows live is an application's decision exactly as a journal copy's format is. An application that wants this tier writes that adapter itself, against a contract that is already closed.

A Series<A> is first an ordinary StreamConsumer<A>: it names the streams it captures, keeps its own cursor, and is fed beside the entry order exactly as chapter 6's recording is. Four more verbs are what make it a served tier rather than a recording:

VerbAskedAnswers
rowsa position range of one stream, narrowed to a trackthe rows in publish order, erased
bucketsa declared-axis range and a point budgetBuckets, never rows
aggregatea declared-axis range and one named statisticone Aggregated per measured field
sealthe barrierthe Position through which rows are durable

So the shape you write is two impls, and the first of them is the recording contract you already know:

use harmos::{Batch, Cursor, Definition, Series, StreamConsumer};
pub struct SchematicSeries { /* whatever your store is */ }
impl StreamConsumer<SchematicEditor> for SchematicSeries {
fn resume(&self) -> Cursor { … }
fn selections(&self) -> Vec<Definition> { … }
async fn drain(&mut self, batch: Batch<SchematicEditor>) -> Cursor { … }
async fn flush(&mut self) { … }
}
impl Series<SchematicEditor> for SchematicSeries {
// The four verbs above, and nothing else.
}

Attaching the adapter you wrote is a single line:

use harmos::Storage;
// `SchematicSeries` is the application's own `Series<SchematicEditor>`.
let storage = Storage::new().series(SchematicSeries::open("var/series")?);

Storage::series designates one adapter and only 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. The designation is a stream attachment and nothing more: the adapter is folded exactly as attach_streams folds one, serving is what it adds, and it feeds no durability frontier.

Two things the contract deliberately does not carry. How coarse the stored summaries get and how many rungs the ladder has is a choice only something that stores rungs can make; so is how long each rung survives. Both belong where the adapter is built, at assembly, where the operator's answer lives — not on the row, where the application's answer lives, and not on a port that would then have to mean the same thing to an adapter that keeps no rungs at all.

A track is the one piece of this the contract does name, and it names it as text: every verb takes the component labels of a caller's selection, outermost first, and the encoding those labels are stored under is the adapter's own. Nothing above it builds or sees one.

Sealing Is the One Thing You Have to Ask For

Section titled “Sealing Is the One Thing You Have to Ask For”

The producer never pays. publish stays synchronous and infallible, and the adapter's side of the bargain is that taking a delivery run stages the rows in memory and touches no file at all. That keeps a publish off the write path and keeps 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 rather than to a publishing thread. A resident owns the cadence and asks for the barrier through the same handle every other caller uses (chapter 11):

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

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 keeps a barrier from racing a drain.

Make a barrier store-wide, and a cadence tick costs one commit however many feeds it answers for: one commit point covering every stream and every track staged at the moment it is taken. Sealing them in turn instead pays the commit tail once per stream and buys nothing, because rows left staged are rows a crash still loses. Either way the cadence is the knob worth thinking about. It is the size of your risk window — rows staged when a process dies are gone, honestly and by design — and it is also what trades a display read's tail latency against that loss.

A caller can always take a barrier itself:

let durable = runtime.streams.seal::<SimSample>().track((run, probe)).await;

seal answers the Position through which those rows are durable, and with no adapter attached it answers Position::ORIGIN — nothing is durable, said plainly, rather than a success that kept nothing. It exists because of the asymmetry chapter 6 introduced: entries are waited for, so a caller never has to ask; rows are dropped honestly, so a caller who needs them kept must. The cadence is what means nobody has to.

The verb did not change:

let samples: Vec<_> = runtime
.streams
.query::<SimSample>(..)
.track((run, probe))
.as_of(committed)
.await?
.collect();

The answer path did. The window is read first and answers for everything above its floor; what remains is asked of the adapter, capped at that same floor. The seam is invisible in the answer, exactly one tier owns the boundary position, and neither tier sorts — both are ordered by construction.

A restart does not have to wait for an eviction to get that seam back. Boot asks the attached adapter how far it is durable and stands each empty ring on that watermark, so the first query after a reboot already splices: beneath the watermark is what the previous run sealed, above it is what this one has published. There is nothing to reopen and nothing to merge on your side — recovering a feed is the same query it was before the process died.

.track(…) exists only on tracked streams. Selecting a partition of a stream that declares none is not a runtime refusal; it is a call that does not compile.

It takes the declared tuple or any prefix of it, which is the second half of the same guarantee:

// One probe of one run.
.track((RunId(7), ProbeId("V(out)".to_owned())))
// Everything run 7 published, across every probe.
.track((RunId(7),))
// Neither of these compiles: wrong component type, wrong depth.
// .track((ProbeId("V(out)".to_owned()),))
// .track((RunId(7), ProbeId("V(out)".to_owned()), RunId(8)))

A prefix is checked where it is written, and the answer is cut by your own components' equality rather than by a name anybody spelled. There is no string to concatenate and no separator to agree on: what a partition is called on disk is the storage adapter's business, and it is the only layer that knows.

A chart of a million samples does not want a million rows. It wants a few hundred points that keep the shape. That is a different question with a different answer type:

let envelope = runtime
.streams
.survey::<SimSample>(0..=10_000)
.track((run, probe))
.budget(400)
.await?;

survey answers Buckets, never rows. Keeping them separate verbs with separate return types is deliberate: exporting decimated data by accident becomes a type error rather than something code review has to catch.

It takes only an axial stream. This is the tier a display read is answered from, and a stream with no coordinate to place rows on has nothing to bucket — the bound is on the verb, so asking is a compile error rather than an empty answer nobody noticed.

Every bucket carries true M4 for one field — first, last, min, and max, each with the axis coordinate it occurred at — plus count, sum, and sum_sq. The 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 everything else a fold; a mean cannot be merged, so none is stored, and Bucket::merge combines two adjacent buckets into the bucket of their union exactly.

budget is the only way to take that answer, so this verb has no unbudgeted form — which is precisely what stops a stored summary from being a rung no read can reach. A bucket is four placed points, so a point budget is a bucket count, and every measured field answers its own run out of the same budget:

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

SimSample measures two fields, not five. No declared coordinate is a measure: the axis places the rows, and every track component is one repeated value within its own partition, so folding any of them would spend budget on a statistic of a constant. A budget of 400 therefore buys fifty cells of millivolts and fifty of residual.

The adapter answers from the finest stored rung that fits beneath that ceiling, never the coarsest — a budget is a ceiling, and spending less than one is always allowed. Buckets come from the served tier alone, because the window retains rows and not summaries. Seal, then survey.

One convention is worth knowing before you draw an axis label: a survey admits every whole cell its range touches, so a boundary bucket may fold a few rows just outside the range you asked for. Clipping it would produce a bucket that is no longer the fold of a grid cell, which is the property every stored rung rests on. A caller who needs the range exactly reads rows.

Sometimes the question is not a shape but a number. The same assembled survey takes a second answer:

use harmos::Aggregate;
let mean = runtime
.streams
.survey::<SimSample>(0..=10_000)
.track((run, probe))
.fold(Aggregate::Mean)
.await?;
let tail = runtime
.streams
.survey::<SimSample>(0..=10_000)
.track((run,))
.fold(Aggregate::Percentile(95))
.await?;

fold answers one Aggregated value per measured field. There is no budget, because there is nothing to spend one on, and there is no default kind, because naming the aggregate is how the answer is taken at all.

The vocabulary is one enum with one parse table, and the whole of it is Min, Max, Mean, Sum, Count, and Percentile(rank) for a rank in 0..=100. The median is Percentile(50) and has no second spelling — one statistic, one value, one name, which is what keeps a facade and a wire speaking the same language instead of drifting into two.

The two halves are priced differently, and the difference is real rather than documentary:

KindReadsAnswers to
Min, Max, Mean, Sum, Countthe finest stored rungthe summary floor
Percentile(rank)the rows in the rangethe raw floor

The first five come out of the accumulators a stored bucket already carries, so they cost a summary read and no rows at all. A percentile has no mergeable accumulator — it needs the values themselves — so it is a scan and is priced as one. Because they read different tiers, they can answer differently: once retention has reclaimed a range's raw rows, a mean still comes back from the summaries that outlived them and a p95 is refused. Each says which.

The same vocabulary can cross a guest capability boundary. A guest names the aggregate as a string because the component wire carries strings, and the host parses it with this table: a name outside the vocabulary refuses with unsupported and the name it could not honor. Nothing falls back to a default kind.

Retention Is the Adapter's Declaration, Reported as a Floor

Section titled “Retention Is the Adapter's Declaration, Reported as a Floor”

A store that grows without limit is a bug on a delay. Nothing in harmos retires a stored row, because nothing in harmos stores one: how long each tier survives is declared where the adapter is built, beside the rungs it decides to keep. What the contract fixes is not the policy but how its effect is reported — and that is the part a caller writes code against.

The observable outcome is that floors advance, and an adapter is held to saying so. Answering short is never permitted: one that no longer holds the whole range answers Floor with the oldest position it does hold, and one that cannot read what it holds answers Unreadable instead. Raw rows and folded levels are allowed to answer different floors, which is why a survey is never refused on a floor its summaries still hold, and why the two halves of the aggregate vocabulary can disagree about the same range.

Retiring a whole partition is the one tracks exist for. A run ends, and the only evidence of that ending which ever reaches a storage adapter is that its rows stopped — so that is what an adapter can measure, and how long a silence has to last to mean it is its own declaration. It measures each stored partition, which is the whole tuple: (run, probe) retires probe by probe as each one falls silent, and a run is gone once its last probe is.

BeneathWindow names three different walls, because each one is a different next move:

VariantMeaningYour move
Unservedno series adapter is attachedimplement the port and attach one at assembly
Floor { floor }no tier retains anything beneath floornarrow the range
Unreadable(_)the adapter could not read what it holdsreport it, retry

The third is the one worth reading carefully. A segment whose bytes no longer match the digest the manifest recorded is Unreadable, never a floor — because a floor reads as retention, and a caller told a range was retired stops looking for data that is still sitting there. Answering short is never permitted: an adapter that does not hold the whole range reports its floor, and one that cannot read what it holds says so.

1. You published ten thousand samples, then surveyed immediately and got nothing back. What is missing?

A barrier. Summaries come from the served tier alone, and a drain only stages rows — nothing is stored until a barrier commits it. Either wait for your sealing resident's cadence or take one with seal. Seal, then survey.

2. Your rows carry a microsecond timestamp. Should it be the axis?

Only if it is monotone and you actually read and display on it. An axis is a u64 domain coordinate for exact grid arithmetic, so a solver step or a sample ordinal is usually the better answer, and the clock rides beside it as an ordinary field. Never make Position the axis: it orders rows against the journal and says nothing about your domain.

3. A range answers a Mean and refuses a Percentile(95). Is the adapter broken?

No — that is retention being visible. A mean folds from stored summaries and a percentile has to read rows, so once the adapter has reclaimed the raw rows the two questions stand on different tiers. The refusal carries the raw floor; narrow the range or ask for a statistic the summaries can answer.

4. You want a per-net breakdown. Should net be a track component?

No. A track is a partition that ends, which is what makes per-track floors and track retirement meaningful, and a net never ends. The tuple is bounded at four components on purpose — cardinality you can reason about instead of a tag space whose cost nobody can predict. Cut by net in the application, over the rows a query answers.

5. Your track is (run, probe) and you want one run's whole envelope. Do you survey each probe and merge?

No — pass the prefix. .track((run,)) folds every partition beneath it into the same buckets, out of the same stored levels, in one read. Merging on your side would fold the same cells twice and get the same answer more slowly.

Stokker Technologies markDesigned and built by Stokker Technologies