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.
Two Questions, Two Verbs
Section titled “Two Questions, Two Verbs”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.
Reading the State
Section titled “Reading the State”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.
Why a Closure and Not a Getter
Section titled “Why a Closure and Not a Getter”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
.awaitinside 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.
What the Position Stamp Is For
Section titled “What the Position Stamp Is For”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 isat.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
queryorwatch, because reads and history share one coordinate system.
Reading History
Section titled “Reading History”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;}When Memory Runs Out
Section titled “When Memory Runs Out”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.
Common Mistakes
Section titled “Common Mistakes”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
commitfuture is still on its way back to you. Your receipt is your copy of the confirmation, not the moment of the event.
A Live Status Bar
Section titled “A Live Status Bar”// `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.
Test Your Knowledge
Section titled “Test Your Knowledge”
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.
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?
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?
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.