Skip to content

Read and Observe

Two questions pull things back out of a runtime — what is the state now? and what has happened? — and they take two different verbs. This chapter is both, the position stamp that makes their answers comparable, and the loop watch is actually built out of. The mistake it exists to prevent is answering the second question with the first.

Everything so far has pushed changes into the runtime. Now we pull things out — and there are two genuinely different questions to ask, so there are two different answers.

What is the state right now? State is one value, so the answer is one value.

What has happened? History is a sequence, so the answer is a sequence.

Harmos keeps these apart on purpose. Reading the state tells you the fold of everything committed; reading history tells you the entries that produced it. Confusing the two is the most common mistake on this page, and the last section of it says why.

let At { value: nets, position } = runtime
.journal
.read(|schematic| schematic.nets.len())
.await?;

You hand in a closure. You get back whatever the closure returned, paired with the position it was read at — the same At<T> from Chapter 1, except the value is now what you computed rather than the whole Schematic.

Two small things in that line. The read is async, because the state lives behind an asynchronous lock and a reader waits its turn. And it is fallible in exactly one way: Error::Unavailable, which means a change panicked inside apply and the state is poisoned. Reads are refused runtime-wide from that instant, because nobody may observe half-mutated state — the whole story is in Recover and Survive. There is no other reason a read can fail.

The Journal<A> handle is the one state-reading capability. It is Send + Sync and cheap to clone, so request paths, services, and UI tasks take the handle they need and call journal.read(|schematic| ...).await. Keeping the verb on the handle also means a long-lived host never needs to pass the whole boot-and-stop Runtime<A> into ordinary work.

An API returning a state reference or lock guard would be worse in a way that is worth understanding, because the reason shapes how you write every read.

The state lives behind a read-write lock. Many readers run inside it at once, and they do not block each other. The single writer, however, must wait for the readers currently inside before it can apply the next transaction. So the duration of your read is the duration of somebody else's commit latency.

A closure makes that duration a bounded, visible region of code. It also makes the dangerous version unwriteable: the signature is FnOnce(&A::State) -> R, and the borrow checker will not let a reference to the schematic escape the call. You cannot hold what you were never handed.

Three habits follow, and they are the whole discipline:

  • Keep the closure short. Compute, do not orchestrate.
  • Return owned data — a count, a clone of one symbol, a small summary struct.
  • Never .await inside it, and never do I/O in it. Both are ways of holding the lock for an unbounded time.

Long-lived observation is not a long-lived read. It is a projection or a consumer, covered in Persist.

Every read is stamped, which turns three fuzzy questions into integer comparisons:

  • "Which version of the document did the user actually see?" The number you drew from is right there beside the pixels.
  • "Does this read include the change I just committed?" Compare it with your Receipt: the test is at.position >= receipt.position. Never == — other people commit too, and your change is included by everything after it, not only by the moment it landed.
  • "Where do I start watching from?" A position from a read is a position you can hand straight to query or watch, because reads and history share one coordinate system.
pub async fn query<X>(&self, after: Position) -> impl Iterator<Item = Result<Entry<X>, Error>>;
pub fn watch<X>(&self, after: Position) -> impl Stream<Item = Result<Entry<X>, Error>>;

query is finite: an ordinary Iterator over a snapshot of history, which ends when it reaches the head. Taking that snapshot means reading the tail, so the call itself is async; walking the result afterwards is not. Use it to catch up. watch is an ordinary asynchronous Stream that does not end, and starting one costs nothing, so it is not async at all. Use it to follow.

Neither returns a bespoke harmos type. query hands back the standard library's Iterator, so map, filter, and take_while work on it directly. watch hands back futures_core::Stream — the plain trait, which is the only thing harmos depends on. The .next()-style combinators over a stream come from whichever combinator crate your application already uses; harmos does not pick one for you.

Both take the position you want to hear about entries after, and both yield the same item:

pub struct Entry<X> {
/// The entry's permanent identity.
pub position: Position,
/// What the runtime witnessed when the entry was sealed.
pub metadata: Metadata,
/// The payload the application declared.
pub value: X,
}

value is the payload your application declared. metadata is what the runtime witnessed when the entry was sealed — principal, time, correlation, causation (Chapter 4). position is the entry's permanent identity.

The type parameter says which catalog you want to hear about. watch::<Transactions>(p) yields changes; query::<Records>(p) yields facts. Transactions and records live in one total order and one position space, so a record at 41 and a transaction at 42 are adjacent in a single history — you are filtering one stream, not merging two.

A projection catching up from where it left off looks like this:

let mut cursor = last_checkpoint; // a Position your fold stored
for entry in runtime.journal.query::<Records>(cursor).await {
let entry = entry?;
index.absorb(&entry.value);
cursor = entry.position;
}

The history harmos keeps in memory is a bounded suffix — the most recent entries, not all of them. With a recovery source attached, query serves an exact older range through that source. Without one, asking below the floor gets Error::BeyondTail.

Read that error precisely. It says I cannot answer that from memory. It does not say the entries never existed or were thrown away — where they live is the subject of Persist. Harmos refuses rather than silently returning a short answer, because a short answer to "everything after 900" is indistinguishable from the truth and would corrupt any fold that trusted it.

Chapter 13 follows the storage-backed path and its retention boundary. The fix for a lagging live fold is still almost never a bigger tail. A consumer that keeps its own checkpoint and folds forward as entries arrive never falls behind the floor in the first place, and that is what derived views are supposed to be.

Polling state instead of watching. A timer that reads state every 16 ms and diffs the result cannot see entries — it sees folds. Two renames of the same net between ticks look like one rename; a rename followed by a rename back looks like nothing happened at all. Anything that must count, log, audit, or react to what was done has to read entries.

Holding the read open. Awaiting inside the closure, or doing file I/O in it, makes the writer wait on your slowest operation. The commit that stalls is somebody else's.

Treating watch as a subscription. There is nothing to register and nothing to unregister — see below.

Under the Hood: The Tail, the Scan Frontier, the Doorbell

Section titled “Under the Hood: The Tail, the Scan Frontier, the Doorbell”

Chapter 1 introduced the runtime's single shared reality. Here is the rest of it:

struct TrackedState<A: Application> {
state: RwLock<At<A::State>>,
tail: RwLock<Tail<A>>,
applied: watch::Sender<Position>,
stored: watch::Receiver<Position>,
poisoned: AtomicBool,
}

Which end of a channel a field holds is the capability. The writer alone holds the applied sender, and the journal only ever receives stored, because it consumes durability and never claims it (Persist). poisoned is the flag reads consult before they take the lock.

The tail is the bounded in-memory suffix of the entry order, holding typed entries. It is what query and watch read, and it is the only thing they read. There is no separate "sealed entry" type standing between the writer and the readers: the entry in the tail is the entry.

The scan frontier is the internal position each iterator and stream carries as it walks. The important property is that it advances across the mixed order, even when nothing matches what you asked for. Ask for records while twenty transactions go by, and the frontier moves past all twenty. Two consequences fall out of that one behaviour: you never rescan a span you have already walked, and the position you resume from is a position in the one shared order rather than a per-type counter — which is why a record consumer and a transaction consumer can compare notes at all.

The doorbell is applied, a watch channel carrying the position of the last entry applied in memory. The writer holds the sending end; everyone else holds a receiver. journal.applied() reads its current value.

And that is enough to build watch out of, because a watcher is not a subsystem. A watcher is a position plus a loop:

loop {
read matching entries from the tail, starting after my position
advance my position past everything I scanned
when the tail has nothing more for me → await the doorbell
}

That is the whole implementation. No subscription registry, no cursor store, no coalescing machinery, nothing to unregister — drop the stream and it is gone, with no bookkeeping left behind anywhere.

Two behaviours follow directly from that shape, and both surprise people:

  • The doorbell carries no entries. It only says "history moved"; the watcher then reads from its own position. So a watcher that wakes late has not missed anything — it simply does one longer read. Notification is not delivery, and the position is the truth.
  • A watcher can see your change before you do. Chapter 4's pipeline publishes before it replies: the entry joins the tail and the doorbell rings while your commit future is still on its way back to you. Your receipt is your copy of the confirmation, not the moment of the event.
// `watch` hands back the plain `futures_core::Stream`. `next()` is not part of
// it, so it comes from whichever combinator crate this application already has.
use futures::StreamExt;
/// Keeps the editor's status bar in step with the journal, forever.
async fn run_status_bar(runtime: Runtime, ui: StatusBar) -> Result<(), Error> {
// Draw once from the current state — and remember when "now" was.
let At { value: summary, position } = runtime.journal.read(Summary::of).await?;
ui.draw(&summary);
// Then follow every change committed after that exact point.
let mut changes = runtime.journal.watch::<Transactions>(position);
while let Some(entry) = changes.next().await {
let entry = entry?;
let At { value: summary, .. } = runtime.journal.read(Summary::of).await?;
ui.draw(&summary);
ui.footer(format!(
"edit {} by {}",
entry.position,
// Witnessed, so it is an Option: an entry committed through an
// unscoped handle belongs to nobody in particular.
entry
.metadata
.principal
.as_ref()
.map_or("the editor", PrincipalId::as_str),
));
}
Ok(())
}
struct Summary { symbols: usize, nets: usize }
impl Summary {
fn of(schematic: &Schematic) -> Self {
Summary { symbols: schematic.symbols.len(), nets: schematic.nets.len() }
}
}

Look at what the position does between the two calls. The read returns the position it reflects; the watch begins after that same position. Entries committed in the gap between the two lines are not lost and not counted twice, because both verbs speak the same coordinates and the gap is measured in them. There is no race to reason about, and no "prime the cache then subscribe" window to get wrong.

1. A colleague wants to delete the closure: fn read_guard(&self) -> StateGuard<'_, Schematic>, so the renderer can walk a 200-sheet schematic directly instead of copying summaries out. Name the two distinct costs, one of which is not about performance.

Cost one — you have handed out the writer's throttle. Readers and the writer share one read-write lock, and the writer must wait for the readers currently inside before it can apply. With a closure, the read is a bounded region you can see on one screen. With a guard, the read lasts as long as the renderer holds the value — across a frame, across an .await, across a file dialog if somebody is careless. Commit latency for every user of the runtime becomes a function of the slowest renderer.

Cost two — the position stamp stops meaning anything. At<R> promises "this value is history up to N". A guard held over time is a live view of a state that keeps changing under it, so there is no single N it reflects. You would be handing out a value that cannot be labelled — and the position is what makes reads comparable with receipts and usable as a watch start. You would lose the coordinate system, not just some milliseconds.

Both costs are prevented by the same signature: FnOnce(&A::State) -> R cannot leak a reference, so the borrow checker refuses to compile the dangerous version. This is the pattern from Chapter 3's mandatory origin again — do not validate the bad usage, make it unconstructable.

2. After a slow deploy, a projection restarts, asks to resume from the position it stored last week, and gets Error::BeyondTail. A teammate concludes the journal dropped those entries and proposes raising the tail size until it stops happening. Both halves of that are wrong. Why?

The entries were not dropped. The tail is a bounded in-memory window over history, not history itself. BeyondTail says harmos cannot answer that question from memory — it is a refusal, deliberately chosen over silently returning a short answer, because a truncated "everything after 900" is indistinguishable from a complete one and would corrupt any fold that believed it. What exists on disk is a separate matter, and Persist is where it is decided.

A bigger tail does not fix the class of bug. It buys a longer outage before the same failure, and it costs memory on every runtime forever. The real defect is a derived view that stopped folding and expected history to wait for it. A checkpointed consumer that folds forward as entries arrive keeps its own resume position and never asks about a position below the floor, because it is never that far behind. The right fix is a consumer with a checkpoint, not a larger window.

3. Your commit resolves with Receipt { position: 900, .. }. The very next line reads state and gets back At { position: 903, .. }. Meanwhile your logs show a watcher printed entry 900 before your commit returned. Reconstruct all three numbers — and is your change in the state you just read?

Yes, your change is in it, and the reason is not that 903 is close to 900. State is the fold of the whole order, so a state at position 903 includes everything at or below 903 — your entry at 900 among them. This is exactly why the "did it land?" test is at.position >= receipt.position and never ==. Waiting for a read that reports exactly 900 would be waiting for a moment nobody promised would ever be observable.

903 means three more entries were sealed between your seal and your read. Other tasks, other users, other services commit into the same single order; positions are assigned in channel-arrival order and nothing reserves you a quiet moment afterwards.

The watcher printed first because publish precedes reply. The writer joins the entry to the tail and rings the applied doorbell, and only then sends your receipt back down its reply channel. A watcher woken by that doorbell can read, format, and log entry 900 while your future is still being polled. Your receipt is your copy of the confirmation — history does not wait for you to read your mail.

4. A teammate builds the “nets renamed today” counter by reading state every 16 ms and diffing the net names against the previous read. It passes every test they wrote. What class of fact can this design never observe, and what is the shape that can?

It can never see anything that does not survive to the next tick. A read returns the fold, not the entries. Between two ticks, VCC renamed to VDD and then renamed to VBUS is one difference, so the counter reads 1 instead of 2. VCC renamed to VDD and back to VCC is no difference, so the counter reads 0 instead of 2. The tests pass because a test commits one change and waits — which is precisely the case where folds and entries agree.

It also cannot see who did it, when, under which correlation, or whether the rename was itself an undo, because none of that is in the state. It is in the entry's metadata.

The right shape is a fold over entries: watch::<Transactions>(from), matching on the entries themselves, keeping a resume Position. That is a projection, and it is one of the checkpointed consumers of Persist. As a bonus it costs less than the polling loop — the doorbell wakes it when there is something to do, rather than 62 times a second when there is not.

Stokker Technologies markDesigned and built by Stokker Technologies