Skip to content

Publish Streams

Not everything an application observes belongs in history. This chapter is the facility for the data that does not: a typed feed with a bounded window, a publish that never waits for anything, and loss reported as an exact row count rather than prevented. You leave able to run a high-rate source beside the journal without either one distorting the other.

The schematic now has a simulation. It samples the voltage on a net, slowly in this example and thousands of times a second in a real instrument. The latest samples belong on screen and sometimes in an analysis file. They do not belong in the journal.

A journal entry is durable truth: replay must meet it again, in order, because state or the application's story depends on it. A stream row is an observation: useful now, allowed to fall out of its window, and never allowed to make its source wait. Harmos puts the two beside one another and gives them exactly one shared coordinate — the journal position the row witnessed when it was published.

That makes Streams the answer to chapter 1's high-rate-data rule. Transactions say what the schematic is. NetVoltage says what the simulator just saw.

The whole row definition is one plain Rust struct:

streams/net_voltage.rs
#[harmos::stream(id = "net_voltage", window = 4)]
#[derive(Clone, Copy, Debug, Deserialize, PartialEq, Eq, Serialize)]
pub struct NetVoltage {
/// Which net was sampled.
pub net: NetId,
/// When the simulator sampled it, in microseconds on its own clock.
pub sampled_at_us: u64,
/// The measured potential, in millivolts.
pub millivolts: i32,
}

The three arguments are the whole declaration.

  • id is the permanent artifact name. It is filesystem-safe and a valid Python identifier, because net_voltage.jsonl and a net_voltage analysis column must need no escaping convention.
  • version is the stored row shape. Historical rows follow the same flat, direct-to-current conversion grammar chapter 10 uses for records.
  • window is a row count, never a duration. Four means exactly the four newest rows stay answerable from memory, however quickly or irregularly they arrived.

The timestamp in that struct is deliberate. A stream's position stamp answers which journal state did this row witness? It does not answer when did the instrument sample? Irregular feeds carry their own time as an ordinary payload field. The rule from commit still holds: metadata is what the runtime witnesses; the payload is what the application declares. sampled_at_us is the simulator's declaration, not runtime metadata wearing a different name.

The declaration above is the one compiled by minischematic-app; there is no second teaching row with a parallel vocabulary. A stream row stays in the application crate even when a plugin publishes into it: the window it is retained in and the grouping it is summarized by are the host's policy, not a contract a plugin compiles. What the two halves share is the declared id and schema — chapter 12 shows where that agreement is asserted.

Streams have their own kind-root beside transactions and records:

streams/mod.rs
mod net_voltage;
pub use net_voltage::NetVoltage;
#[harmos::streams]
pub enum Streams {
NetVoltage(NetVoltage),
}

Rows never pass through that enum. It is a declaration catalog: one variant adds the definition to the table boot freezes and generates the Stream<SchematicEditor> bound that makes every verb typed. Assembly gets one new line and no per-definition registration:

Runtime::builder(origin, ())
.register(Transactions::catalog())
.register(Records::catalog())
.register(Streams::catalog())
.storage(storage)
.runtime()
.await

That is the same catalog discipline as chapter 3. Add one vertical stream file and one variant; forget the variant and publish(NetVoltage { .. }) has no Stream<SchematicEditor> implementation, so the omission fails at compile time.

Terminal window
mise run run:minischematic

Stop 3 of the transcript is the complete live path. A watcher appears, demand turns capture on, two rows publish, an as_of query joins them to the committed move, four more rows overrun the four-row window, and the stalled watcher gets the exact loss before continuing:

3 one feed, alive while somebody needs it live or recorded
demand 2 -> recording + watcher started capture
as_of(4) -> 2 rows stamped at the committed move
watch -> Gap { missed: 2 }, then the four retained rows
row @4 · 2900 µs · 3310 mV
row @4 · 4400 µs · 3305 mV
row @4 · 7100 µs · 3295 mV
row @4 · 9800 µs · 3300 mV
demand 1 -> recording keeps capture requested

Read the repeated @4 correctly. Four is a journal Position: the move had committed, the journal then sat idle, and every row witnessed that same unchanged state. A position is not a row number and it is not a clock.

runtime
.streams
.publish(NetVoltage::teaching(1_000, 3_300));

publish is synchronous and infallible. It takes one short lock, gives the row its publish ordinal and current journal stamp, replaces the oldest retained row when the ring is full, wakes observers, and returns. It never waits for a watcher, a recorder, the writer task, or disk.

That asymmetry is the contract. The journal applies backpressure because losing a change would make replay wrong. A stream overwrites because stalling a sound card or ADC would make the source wrong. What it owes instead is honest loss.

Every publish also takes a cursor. The cursor is a dense publish ordinal: one row takes one number, the next row takes the next number. It measures rows, never time. Across several registered streams it orders their interleaved publishes for recording; it is not a second journal coordinate and never appears in as_of.

A producer normally commits and publishes in the same iteration, and the two verbs do not have to agree about anything except which state the row witnessed. The schematic's focused feed publishes after the edit whose position it was handed:

let streams = runtime.streams.clone();
let mut watching = runtime.streams.watch::<NetVoltage>();
let mut demand = streams.demand::<NetVoltage>();
let mut observed_demand = demand.clone();
let (reported, report) = tokio::sync::oneshot::channel();
// Application code owns the source. Harmos reports demand; this task
// decides that demand starts the simulator and its return to zero stops it.
// The recording and watcher each hold one place until their own lifetime
// ends.
let capture = tokio::spawn(async move {
assert_eq!(*demand.borrow(), 2);
streams.publish(NetVoltage::teaching(1_000, 3_300));
streams.publish(NetVoltage::teaching(1_750, 3_290));
let cut = streams
.query::<NetVoltage>(..)
.as_of(committed)
.await
.expect("the first two rows are still retained")
.count();
for (sampled_at_us, millivolts) in [
(2_900, 3_310),
(4_400, 3_305),
(7_100, 3_295),
(9_800, 3_300),
] {
streams.publish(NetVoltage::teaching(sampled_at_us, millivolts));
}
reported.send(cut).ok();
while demand.changed().await.is_ok() {
if *demand.borrow() == 0 {
break;
}
}
});

The example asks its finite question before it deliberately overruns the window:

let cut: Vec<_> = runtime
.streams
.query::<NetVoltage>(..)
.as_of(committed)
.await?
.collect();

query is assembled first and awaited second. The await takes one snapshot; the returned standard iterator holds no lock. .as_of(committed) keeps rows stamped <= committed, so every row in the answer witnessed state no newer than the change at that position.

The idle-journal case is load-bearing. A row published ten seconds later under the same position is included, because the state did not move in those ten seconds. as_of promises state consistency, not wall-clock chronology. Put chronology in the payload, as sampled_at_us does.

Once eviction has happened, a range whose start reaches beneath the retained window gets BeneathWindow. It is never shortened silently. A short answer and a complete answer look identical to a fold, so refusing is the only honest choice.

That refusal is where chapter 14 picks up. Attach a series adapter and the range the ring no longer holds becomes a place this same verb reaches, spliced at the window's floor.

let mut watching = runtime.streams.watch::<NetVoltage>();
match next(&mut watching).await {
Some(Item::Row(At { value, position })) => draw(value, position),
Some(Item::Gap { missed }) => mark_discontinuity(missed),
None => finish(),
}

watch starts at the live head and returns the standard asynchronous Stream trait. Its item is either a typed position-stamped row or Gap { missed }. The gap is not a synthetic row and lives in no window; it is this watcher's account of how far it lagged. Another watcher over the same ring may have kept up and see no gap at all.

missed counts rows. Not microseconds, not positions, not record batches: rows. In the transcript six rows publish into a four-row window before the watch is polled, so the only correct answer is two. After that one gap it continues with the four rows it can still reach, in publish order.

Starting a hardware source because a window exists would waste power forever. Starting it because a watcher appeared is application policy, so harmos reports the fact and leaves the decision to your task:

let streams = runtime.streams.clone();
let mut demand = streams.demand::<NetVoltage>();
tokio::spawn(async move {
while demand.changed().await.is_ok() {
if *demand.borrow() > 0 {
simulator.start();
} else {
simulator.stop();
}
}
});

The number is live watchers plus attached recordings that selected this definition. It is advisory — transitions race with publishes — and exists for idle-off, never admission control. Harmos starts and stops no application task. The example has one selected recording before the watcher appears, so demand rises from one to two for the live view and falls back to one when that view closes. The recorder keeps the source requested until runtime stop drops the last demand and the capture task ends.

Live windows are not durability. When analysis needs an artifact, attach a stream consumer and select exactly the streams it should record. Harmos ships no recorder: what it publishes is the StreamConsumer contract, and the reference adapter's JsonlStreams is a worked implementation of it you can copy (reference):

use harmos::Storage;
use jsonl_adapter::JsonlStreams;
let recording = JsonlStreams::<SchematicEditor>::open(&capture_directory)?
.select::<NetVoltage>();
let storage = Storage::new().attach_streams(recording);

The attachment counts as demand for NetVoltage, so the same capture loop stays on while recording even when no screen is watching. The adapter appends one JSON line per row to net_voltage.jsonl and makes each delivery run durable before it answers. Nothing is finalized and nothing is renamed: a line is its own frame, so every byte up to the last newline is already a complete, readable artifact. The one thing an interrupted write can leave is a final line that never landed whole, and the next open drops it.

That framing is also why the file needs no header and is never rotated. There is no stored schema for a new row shape to contradict, so a version bump keeps writing into the same capture; an application that wants historical files separated per version declares that policy itself.

The stream-consumer checkpoint has one exact reading. resume() says what the fold has taken in; drain(batch) answers the cursor it made durable. The space between those two answers is that recording's risk window. A cursor is still a dense row ordinal — never a timestamp and never a journal position — and the adapter keeps its own beside the captures rather than inferring one from their rows.

Every recorded row carries two technical members beside the fields the application declared:

  • position: the journal position witnessed at publish;
  • gap_missed: loss evidence on the first real row after a gap.

gap_missed counts rows. It is a member on the line, not file metadata and not an invented row, and it is present on every line so a reader's Option<u64> never meets a missing member. Null means no gap precedes this row; 3 means three publish ordinals were missed before it. A row type that already names position or gap_missed is refused rather than silently overwritten.

Reading a capture back is one call, and the type parameter is whatever you want a line to be:

use jsonl_adapter::streams::read_rows;
let rows: Vec<NetVoltage> = read_rows(capture_directory.join("net_voltage.jsonl"))?;

Serde ignores the members it was not asked for, so the same file answers a request for the domain row and a request for the technical members beside it. An analysis that specifically wants the loss evidence names a struct with nothing else in it and reads the very same lines:

#[derive(Deserialize)]
struct Columns {
position: u64,
gap_missed: Option<u64>,
}
let columns: Vec<Columns> = read_rows(capture_directory.join("net_voltage.jsonl"))?;
demand 0 -> capture stopped; recording flushed
position -> [4, 4, 4, 4]
gap_missed -> [Some(2), None, None, None]

The example's integration test reopens that capture through the same call and asserts both vectors, so this is executable readback rather than a format sketch.

And the analyst's entire Python side can be one line:

Terminal window
python -c 'import pandas as pd; print(pd.read_json("net_voltage.jsonl", lines=True))'

There is no live form and finalized form to tell apart. A capture a process is still appending to and a capture its process died holding read exactly the same, one complete line at a time.

The same split is the framework-embedding rule recorded in .codex/skills/harmos-tauri/SKILL.md: handles are the request-path currency; the runtime is boot/stop capability kept for the exit path. Manage cheap Journal<A> and Streams<A> clones in Tauri commands. Keep Runtime<A> out of managed request state, then drop the managed handles before calling runtime.stop().await during exit.

1. The journal has not moved since position 40. A sensor publishes rows at 12:00 and 12:05, both stamped 40. Is one stamp wrong, and can as_of(40) include both?

Neither stamp is wrong, and the cut includes both. The stamp says the row witnessed state no newer than position 40; the state was unchanged for those five minutes, so both observations satisfy the promise.

If five minutes matters to the application, it belongs in a timestamp field on the row. Turning Position into time would destroy the shared coordinate without supplying a reliable clock.

2. A watcher reports Gap { missed: 120 }. A teammate logs ‘120 ms of data lost’. What information do they actually have?

Exactly 120 rows were lost to that watcher. Nothing in the gap says how long those rows span. At a steady 1 kHz it might be 120 ms; on an irregular feed it might be a minute; across a burst it might be less than one millisecond. Duration can only be derived from payload timestamps around the gap, never from missed itself.

3. A recorded row follows a loss of seven rows. Where is that evidence stored, and why is adding a synthetic row the wrong shape?

The first real row after the loss carries gap_missed = 7. Later rows carry null until another loss — the member is on every row, so the absence of a gap is stated rather than left to a missing field.

A synthetic row would have to invent values for every application field — net, timestamp, voltage — and an analyst could mistake it for an observation. Gap evidence describes the boundary before a real row; it is not itself a row. File-level metadata is not used.

4. No UI is watching, but a recording selected NetVoltage. Should the simulator stop when its last window closes?

No. Demand counts both live watchers and attached recordings that selected the stream, so the recording keeps the count above zero. The application task stops the simulator only when the combined count returns to zero.

Harmos still makes no control decision. It reports the count; the simulator task owns start and stop, including any debounce or hardware-specific policy around the transition.

Transport packing and statistical grouping are independent. The runtime can summarize declared numeric scalar fields while still publishing every raw row immediately:

#[harmos::stream(id = "voltage", window = 4096,
batch(samples = 1024, gaps = "report"))]
#[harmos::message]
#[derive(Clone, PartialEq, serde::Serialize, serde::Deserialize)]
struct Voltage {
#[harmos(measure)]
volts: f64,
}

message is optional for process-local streams; include it when the same row also crosses a protobuf publication boundary.

runtime.streams.batches::<Voltage>() returns a typed asynchronous stream; batches.next().await yields BatchItem::Summary(Arc<Summary>) or BatchItem::Gap { missed }. Each field has count, sum, mean, min, max, population variance, population standard deviation and RMS. Welford's online recurrence avoids subtracting large squared means; scaled RMS avoids directly squaring large values. No percentile or median scan is performed.

One group contains exactly samples received rows, independent of gRPC packet boundaries and observers. Source gaps do not fill group slots: their count accumulates in Summary::missed, in constant work even for a large gap. complete() means a full group with no source loss. It does not assert numerical validity. NaN/infinite field values are excluded and counted in invalid; finite arithmetic overflow sets overflow and leaves statistics absent. A field with no finite values also has absent statistics. Other fields in the same row remain independent. Integer scalar measures through 64 bits convert to f64, so integers beyond 2^53 may lose precision; these are floating-point statistics, not exact integer sums. Buffers, arrays, optional fields and nested messages cannot be measures. A publication containing a byte buffer remains one row, never an implicit sequence of values.

The runtime computes each group once under its stream window lock, even with no summary observers. Raw streams that omit batch allocate no accumulator or summary ring. Summary observers contribute to stream demand. Every opted-in stream retains the latest 64 summaries; a slow observer sees the exact number of summaries evicted, separately from source rows lost within each summary. An observer starts at the current summary head, including any unfinished group when it completes. Dropping an observer releases its demand without affecting computation or other observers.

Runtime close flushes one unfinished group before closing observers. A tail containing only source loss produces a partial summary with count zero, its missed count and absent statistics. An empty tail produces nothing. Observers drain retained results and end; rows published after that stream's close do not create more summaries. Summary retention is memory-only, resets at boot, and is separate from raw storage/history. batch and measure are runtime policy, not a change to the serialized row schema or its history ordinal.

Only gaps = "report" is supported, and it is the default. samples accepts 1–4294967295. At least one scalar must be marked. Tracked streams explicitly reject statistical batching for now: no partitions are silently combined and no unbounded map of partition accumulators is created. This feature reduces neither raw transport traffic nor raw retention work.

Run mise run bench:streams:batches for a separate repeated release measurement of runtime rows/second with equivalent numeric row shapes and summary grouping enabled/disabled. It includes the raw runtime ring and statistics, excludes sidecar transport and storage, and has no timing pass/fail threshold. These rows/second measurements are not the child-process useful MiB/s measurement.

Stokker Technologies markDesigned and built by Stokker Technologies