Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

lazily Cell Model Specification

Normative source of truth for cell kinds and the multi-write merge mechanism.

This chapter defines the cell-kind model that the Wire Protocol serves. It is upstream of every transport: IPC, FFI, signaling, and the distributed plane all carry cells whose convergence semantics are fixed here. The Distributed: CRDT Cell Plane section specifies merge: crdt — the first multi-write merge mechanism defined below.

The single axis: writer count

A lazily graph is a set of cells. The only axis that determines how a cell’s value converges is how many writers can concurrently produce a value for itnot whether those writers are local or remote, in one process or many.

KindConcurrent writersConvergenceMerge
Single-writer (local / direct)exactly onedirect reactive push/pullnone
Multi-writepotentially manymerge: <mechanism> ingressmechanism-defined

The kind is a static property of what the cell represents, chosen at definition time. There is no dynamic per-write mode switching: a multi-write cell stays multi-write even when only one writer is currently live, and a single-writer cell never becomes multi-write because a value happens to arrive over a wire.

Single-writer cells (local / direct)

A single-writer cell has exactly one writer:

  • this runtime’s own derivation graph (a derived cell is always single-writer — see Derived cells), or
  • a single owning runtime that one-way mirrors the cell to other runtimes via a delta projection (today’s lazily-rs → lazily-kt #lazilystatesync push is exactly this shape).

Propagation is direct reactive push/pull over the IPC or FFI channels. A mirror is still single-writer: the receiving side observes, it does not write back. No merge step exists, and none is permitted — a single-writer cell that receives a concurrent remote write is a conformance error, not a merge.

A cell mirrored one-way to N readers is single-writer. Locality does not make a cell multi-write; a second concurrent writer does.

Multi-write cells

A multi-write cell admits concurrent writes from multiple replicas. New values arrive as remote ops merged through the cell’s declared mechanism; the merged result is fed into the reactive graph as an ordinary cell update, after which propagation is identical to a single-writer cell.

A multi-write cell carries a merge: <mechanism> attribute. The mechanism is a parameter of the cell, not a separate cell kind:

Cell      = SingleWriter
          | MultiWrite { merge: MergeMechanism }

This framing is deliberate. An implementation MUST model multi-write as one cell category parameterized by a merge mechanism, and MUST NOT hardcode a single crdt cell kind. New mechanisms slot in without a new kind, and every mechanism shares the ingress boundary, the merge-unit granularity, and the downstream propagation rules below.

Merge mechanisms

crdt is the first mechanism this spec defines and the only one with a normative wire schema today (distributed.json). The mechanism slot is open by construction so later mechanisms slot in alongside it. Every mechanism MUST be deterministic — replicas applying the same op set MUST reach the same value regardless of arrival order or lag.

mergeStatusConvergence strategy
crdtnormative (first)Conflict-free replicated data type (yrs/Yjs-family registers); converges without coordination. See Distributed: CRDT Cell Plane.
lwwreservedLast-writer-wins by HLC/Lamport timestamp.
otreservedOperational transform (server-ordered op rebase).
leasereservedLease/lock-serialized single-live-writer; degenerate-concurrency convergence.
customreservedApplication-supplied deterministic merge function.

crdt is chosen as the first mechanism because it converges without coordination — no central ordering authority, no lock round-trip — which is the property the editor-as-replica use case needs. Reserved mechanisms are named to fix the extension shape; an implementation MAY reject any mechanism it does not implement, but MUST reject it explicitly (capability negotiation) rather than silently treating it as crdt.

A multi-write cell with a single live writer degenerates to a near-free merge (no concurrent ops to reconcile) under every mechanism. This is why the kind is chosen statically by representation: a shared cell that is usually edited by one replica still declares merge: so that the moment a second writer attaches, convergence is already guaranteed.

Merge is an ingress operation on root cells only

The merge mechanism is an ingress step at the boundary where remote ops enter a replica. It applies only to root (input) cells.

Derived cells are never multi-write

A derived cell is a deterministic function of its inputs. It has exactly one writer — the derivation — and therefore is always single-writer. Replicas converge on a derived cell because they converge on its roots, never by merging the derived value itself. An implementation MUST NOT replicate or merge a derived cell directly.

Propagate guard (computed / computed_ripple_when / slot)

A derived cell carries a propagate guard: after it recomputes, the guard decides whether the recompute is forwarded to its dependents. The guard governs propagation, never computation — a derived cell is always recomputed when it is dirty and read (laziness is the only compute gate), so a reader MUST always observe a value current with the inputs. The guard only decides whether dependents are invalidated. This is the glitch-suppression that lets a diamond re-converge without re-running consumers whose input did not meaningfully change.

Three constructors expose the guard; all always compute, and differ only in the guard predicate changed(old, new) (propagate iff true):

  • computed(f) — the default. Guards by the value’s natural equality: changed = (old ≠ new). An equal recompute MUST NOT invalidate dependents. This is the #lzcellkernel guarded computed; equality semantics are the host language’s (PartialEq / == / equals), so for reference-typed values (e.g. a fresh array each recompute) suppression is host-defined and MAY not fire — this is permitted (an unsuppressed cascade is a missed optimization, never incorrect).
  • computed_ripple_when(f, changed) — the same guard with an explicit, pure changed(old, new) predicate: a cheaper/custom equality (dedup a large value by a version/hash field; epsilon compare; hysteresis; monotonic gate; or “propagate every N” when the counter lives in the value). changed MUST be a pure function of (old, new); value-carried state is a permitted input, external mutable state is NOT (it would key off recompute/read frequency and break determinism). computed(f)computed_ripple_when(f, ≠).
  • slot(f)pass-through: no guard (changed ≡ true). Every recompute propagates. The escape for values with no cheap equality (!PartialEq, non- comparable collections) that a binding is fine re-firing. This is a derived constructor and is distinct from the internal storage-sense Slot that a computed node occupies.

The guard is proved in lazily-formal as recomputeSlot_equal_preserves_dependents (equal recompute leaves dependents’ dirty flags untouched) and its specialization recomputeSlot_ripple_when_false_preserves_dependents for a custom changed.

Dependency tracking (the fortified compute view)

A derived cell discovers its dependencies dynamically: on each recompute it runs its compute function and records every cell that function reads. The recorded set is re-bound every recompute (not accumulated), so a conditional read (cond ? a : b) drops the branch it did not take — an implementation MUST NOT retain a dependency a recompute did not read.

The identity that a read must attribute to — which node is being recomputed — is carried into the compute function as a value, through a per-recompute compute view (lazily-rs: Compute, the sole implementor besides Context of the ComputeOps operations subset). It is NOT ambient (thread-local / module global). This is normative because ambient state is clobbered across suspension: an async compute that reads a dependency after an await would attribute it to whatever else ran on the executor. A value threaded through the closure survives suspension (it is captured), so it is the only mechanism that tracks correctly post-await — and the only one that works where no ambient carrier exists (browser JS has no AsyncLocalStorage). Bindings whose runtime does provide a suspension-surviving ambient carrier (Python contextvars, Dart Zone, Node AsyncLocalStorage) MAY use it; all others MUST thread the value.

The compute view SHOULD be fortified so misattribution is prevented by construction, not convention:

  • Sole tracking surface — a tracked read is available ONLY through the compute view; reading through the owning context registers no edge (the explicit untracked escape). A normal read therefore cannot silently miss tracking.
  • Non-escapable — the view MUST NOT outlive the recompute (lazily-rs binds it by lifetime and makes it !Send), so it cannot be stored and later replayed to register an edge against the wrong node.
  • Generation-stamped — a read against a node disposed/recycled mid-recompute is detected, never misattributed.

Edge-attribution invariant (normative): every dependency edge registered during the recompute of node n has n as its dependent. Because the node is a value parameter of the compute view, this holds by construction. Proved in lazily-formal as registerReads_dependent_is_recomputing_node.

Effects stay single-writer

Effects (irreversible external actions — send email, charge card, fire webhook) are not multi-write. State convergence does not authorize an effect to fire on every replica. Effects MUST be gated behind a single-writer authority (a designated peer or small consensus group) that decides when the effect fires, at-most-once. See Single-writer effect authority.

                 remote ops
                     │
                     ▼
        ┌──────────────────────────┐
        │  merge: <mechanism>      │   ← ingress, ROOT cells only
        │  (crdt | lww | …)        │
        └──────────────────────────┘
                     │  merged value as ordinary cell update
                     ▼
        root cell ──► derived cells ──► effects
                    (single-writer)   (single-writer authority)
                     └─ direct reactive propagation, identical for all kinds ─┘

Cell = merge unit

The cell is the unit of merge. Each multi-write cell converges independently; a merge mechanism operates within one cell’s value and MUST NOT move content across cell boundaries.

This is normative because the cell graph already supplies the natural merge boundaries: making the cell the merge unit makes cross-cell contamination impossible by construction. A whole-document merge that splices one logical region’s content into another (for example an agent’s console output bleeding into a queue region) is a conformance violation — each region is a distinct cell and converges on its own.

An implementation MUST scope every merge to a single cell. Coarser-grained merge (whole snapshot, whole document) is permitted only as an optimization that is observably equivalent to per-cell merge; if it can produce a result no sequence of per-cell merges could, it is non-conforming.

The per-cell merge is characterized as an associative fold by the merge algebra (#relaycell): a MergeCell<T, M> is a cell whose write folds under a MergePolicy M, and a plain Cell is MergeCell<KeepLatest>. The merge mechanisms below are the semilattice policies (CrdtJoin<C>) of that algebra; associativity is the invariant that lets a bounded relay flush at any watermark and still converge (Phase 2+).

Liveness vs mechanism

Whether a cell is multi-write (merge: present) is static. How many writers are live against it at a given moment is dynamic, governed by an attach/detach authority state machine (e.g. an editor plugin attaching or detaching). The authority SM governs only liveness — it never changes a cell’s kind or mechanism:

  • No live writer → the cell still declares its mechanism, but with zero concurrent ops the merge is inert and a durable replica MAY be ephemeral (rebuilt from a checkpoint on demand).
  • One or more live writers → ops flow and the declared mechanism reconciles them.

Mechanism is a property of the cell’s meaning; liveness is a property of the current session. Conforming implementations MUST keep these independent.

Cross-process liveness as a CRDT cell. When session liveness itself must cross a process boundary — “editor pid X has doc Y open”, “pid X holds the owner lease” — it is modeled as an ordinary multi-write cell on the CRDT plane, not as out-of-band state: an OR-set for open-set membership (observed-remove, so a re-open wins over a lagging close) and an LWW register for the per-pid alive flag / lease. The derived “is this doc live” aggregate is then a plain reactive memo over that liveness keyed map (the #lzfamilysync derived-aggregate contract), and an OS process-exit event is just the highest-stamp write to alive[pid]. This keeps liveness on the same convergent, idempotent, frontier-resumable substrate as every other replicated cell. Normative semantics: protocol.md § Liveness cells.

Conformance summary

An implementation conforms to the cell model when:

  1. Every cell is classified single-writer or multi-write, statically, by representation.
  2. Multi-write cells carry a merge: <mechanism> attribute; multi-write is not modeled as a hardcoded crdt cell kind.
  3. crdt is implemented as the first mechanism; unimplemented reserved mechanisms are rejected explicitly, never silently aliased.
  4. Merge is an ingress step on root cells only; derived cells and effects are never merged.
  5. Every merge is scoped to a single cell (cell = merge unit); no cross-cell content movement.
  6. Downstream reactive propagation is the same direct mechanism regardless of cell kind.
  7. The attach/detach authority governs writer liveness only, never a cell’s kind or mechanism.
  8. A keyed cell collection (ReactiveMap — the SourceMap/ComputedMap specializations) is implemented — entries are ordinary cells, a dedicated membership cell tracks the key set, and the value / set-membership / order reactivity-independence, stable-handle, and atomic-move invariants below hold. Collections are required of every binding, not optional.
  9. An ordered keyed tree (SourceTree) is implemented, inheriting the per-cell merge and atomic-move guarantees node-by-node (required of every binding).
  10. Keyed reconciliation emits the move-minimized {insert, remove, move, update} op set (LIS over prior indices preserved), and a stable entry is not invalidated by a sibling reorder (required of every binding).
  11. A reactive queue (QueueCell) is implemented — a FIFO collection whose reactive shell invalidates by reader kind (head / length / empty / full / closed), backed by a pluggable QueueStorage backend. The shell / storage split, closure observable contract, bounded-queue backpressure, and ordering guarantees below hold (required of every binding).
  12. Materialization is eager by default — a ComputedMap’s derived entries are pre-minted over the keyset; lazy materialization (get_or_insert_with mint-on-access) is opt-in, keyed, and observationally transparent — identical read values, allocation deferred only (see Materialization). It is a behavior, not a mode flag. Lazy evaluation (bounded-viewport recompute) is provided either way and is never conflated with lazy materialization.

Materialization (a caller-provided recipe)

Cell kind (above) fixes how a cell converges. Materialization is an orthogonal axis: it fixes when a derived cell’s backing node is allocated — not what it computes, not how it converges, not how it merges. It trades memory and first-touch latency against cold full-scan cost, and it MUST NOT be observable through the value of any cell.

Why a behavior, not a mode. Materialization was first pinned as a bespoke ReactiveFamily type carrying an eager/lazy mode — a reaction to one binding (lazily-zig) implementing lazy materialization in the spreadsheet benchmark and thereby diverging from the others. Standardizing a type with a mode invites re-divergence: each binding builds it slightly differently, and most never need it. What must agree is the observable behavior the benchmark measures — transparency (a lazy read equals an eager read) and deferral (an unread lazy entry costs nothing) — not any type or flag. So materialization is normative as a behavior of the keyed primitive: it is simply what ComputedMap does — get_or_insert_with mints a derived slot on first access (lazy); a pre-mint loop over the keyset is eager. There is no materialization mode and no family type — the family types (ReactiveFamily / CellFamily) are removed; ComputedMap (a ReactiveMap specialization) is the vehicle.

The materialization recipe

Materialization is caller-provided: a keyed collection (a SourceMap or any keyed address space) plus a per-key factory whose return type is the materialization choice. Nothing new is required beyond the cell/slot/signal primitives a binding already has:

  • Eager entry — the factory yields an input cell or an eager signal (a memo-slot + puller effect): the node is allocated/pulled up front; a read is a direct node access.
  • Lazy entry — the factory yields a lazy slot: the node is allocated on its first observe, addressed by key. A never-observed lazy entry is never allocated.

So the caller provides materialization through two levers it already owns — whether it observes a key (unread lazy entries stay unallocated) and what the factory returns (slot ⇒ lazy, signal ⇒ eager) — with no mode flag and no per-read toggle. This is the same “lazy by default, eager when asked (via signal)” model lazily already uses at the single-cell level, applied per key. Entry kind is the pinned axis:

  • Cell entries (H = Source) are input nodes — always materialized; an input has no derivation to defer. Minting an input on first get is a collection concern, not materialization.
  • Slot entries (H = Computed) are derived — the ones deferral governs: an eager factory allocates them up front, a lazy factory defers each to first observe.

Entry kind is orthogonal to the materialization choice (proved in lazily-formal’s Materialization module as cell_entries_materialized_in_every_mode / slot_entries_deferred_under_lazy): lazy defers only slot entries, never cell entries.

Normative rules on the recipe:

  • Eager is the default. Absent an explicit lazy opt-in, derived entries are eager — a read is a direct node access and a full recompute pays only compute. A binding MUST make eager the default.
  • Lazy is an explicit opt-in overlay on the eager core, addressed by key, never the default and never a per-read toggle on an eager handle. The first observe of key k constructs the same node the eager build would have, then caches it — a keyed overlay, not a second graph engine. A binding that offers lazy MUST expose it as an explicit opt-in (e.g. a keyed factory / keyed-context constructor).

Observational transparency (normative)

For every node and every read, the observed value MUST be identical under either mode. Materialization mode is not observable on the value axis — it changes allocation timing and memory, never results:

observe(build(eager, spec), id) = observe(build(lazy, spec), id) = spec.val(id)   ∀ id

This is proved in lazily-formal’s Materialization module (observe_canonical, eager_lazy_observationally_equivalent). An implementation MUST preserve these consequences:

  1. Same values. A lazy read returns the value an eager read would (observe_canonical).
  2. No churn from allocation. Materializing one node MUST NOT change any other node’s observed value (materialize_preserves_observe).
  3. Deferral, not de-allocation. Lazy materialization only grows the materialized set; a materialized node is never silently dropped, and the lazy set is a subset of the eager set (materialize_present_monotone, lazy_present_subset_eager).
  4. Reactivity is orthogonal. Lazy evaluation — leaving off-viewport derived cells dirty and never recomputing them (the microsecond bounded-viewport read) — is required of both modes and is independent of materialization. Eager materialization still evaluates lazily; lazy materialization additionally defers allocation. An implementation MUST NOT conflate the two.

Execution-context flavors (thread-safe / async)

ComputedMap runs against a context, and the context is a third axis orthogonal to both entry kind and materialization: it fixes where and how the graph executes, not what it computes. The materialization laws hold over each context a binding provides — the ReactiveMap line has one flavor per context, and each carries the context-specific law below:

  • Single-threaded (ComputedMap, over the base Context) — the reference semantics above.
  • Thread-safe (ThreadSafeComputedMap, over a lock-backed context) — a Send + Sync map that can live in a cross-thread owner (e.g. a hub behind a global mutex, where an Rc-based map cannot go). It carries the same materialization laws, plus materialization confluence: the present set and every observed value are independent of the order in which keys are materialized. This is what makes lock-serialized concurrent materialization safe — any order the lock admits yields the same observable map. Proved in lazily-formal’s Materialization module (materialize_present_comm / materialize_observe_comm).
  • Async (AsyncComputedMap, over an async context) — derived (slot) entries resolve asynchronously, so a non-blocking read returns an optional value (None while pending, Some(v) once resolved). Observational transparency weakens to eventual transparency: once a node resolves, its observed value is the canonical value — identical to what the synchronous ComputedMap observes. Input cells are resolved at build. Proved in lazily-formal’s AsyncMaterialization module (eventual_transparency, async_resolved_matches_sync; a pending read is never a stale value, observe_pending_is_none).

A binding SHOULD keep the per-key factory uniform across flavors, so the flavors differ only in execution context — a derived async slot wraps that factory in a ready computation. The factory SHOULD receive the entry’s own tracking view (Fn(&Compute, &K) -> V, not Fn(&K) -> V): without it a derived entry cannot read another reactive and have that read register a dependency edge, which silently reduces ComputedMap to a cache.

A flavor MUST preserve the entry-kind and materialization laws above, MUST expose the Core surface defined under “Core surface vs. binding extensions” below, and adds the context-specific guarantee (confluence for thread-safe, eventual transparency for async). Ordering and atomic move are not exempt: they are Core, and they bind every flavor the binding advertises. A binding that advertises a thread-safe or async map but ships Core only on its single-threaded map is non-conforming on that advertised flavor. A context declared absent is staged out of that peer group; a runner-local lock, dictionary, or synchronous stand-in MUST NOT be used to make the missing flavor appear covered.

When to opt into lazy

Lazy pays off only for sparsely-touched large keyed address spaces — e.g. a 10,000,000-cell spreadsheet where a session reads ~1% of the derived cells: it lowers peak memory and makes “open” cost O(inputs) rather than O(derived cells). It costs a keyed-cache lookup per read instead of a handle dereference, and a cold full scan pays allocation and compute together (eager_materializes_all vs lazy_defers_slots). Handle-based graphs that read most of what they build SHOULD stay eager. The choice is a per-context construction decision, not a per-cell or per-read one.

Keyed cell collections

A keyed cell collection is a composition of cells, not a new cell kind. It maps keys K to per-entry reactive nodes and adds a dedicated membership cell tracking the set of keys. There is one keyed primitive, generic over the entry’s handle kind:

  • ReactiveMap<K, V, H> — a mutable reactive keyed dict: reactive membership + order, get_or_insert_with (mint-on-access), remove, move. H is the entry handle kind. Its two specializations are the concrete types a binding exposes:
    • SourceMap<K, V> = ReactiveMap<K, V, Source>input-cell entries. Adds set(key, value) (an input is settable). Minting is eager-by-value.
    • ComputedMap<K, V> = ReactiveMap<K, V, Computed>derived-slot entries. get_or_insert_with(key, factory) mints a slot on first access (lazy materialization); a slot’s value is derived, so ComputedMap has no set. Eager materialization is a pre-mint loop over the keyset; lazy is mint-on-access — there is no eager/lazy mode flag.

Deprecated spellings. SourceMap and ComputedMap were previously CellMap and SlotMap, with ThreadSafeCellMap / ThreadSafeSlotMap and AsyncCellMap / AsyncSlotMap as their per-context variants. The rename finishes the v2 kernel migration: the node kinds became Source and Computed, and the map names now say which kind they hold instead of naming a vocabulary the kernel no longer uses. A binding SHOULD keep the old names as deprecated aliases of the new ones rather than removing them, so a caller that has not migrated still compiles. Runners MUST accept the old model spellings in a fixture ("CellMap", "SlotMap") alongside the new ones — the same dual-accept the corpus already uses for the signal/eager and dispose_signal/lazy op names. The fixture FILE names keep their historical spelling; a file name is not a type name, and renaming them would invalidate every binding’s replay ledger at once for no semantic gain.

set(key, value) is therefore cell-only (lives on the SourceMap specialization); the shared surface — get_or_insert_with / remove / move / membership / order — lives on the generic ReactiveMap. There are no family types: the “keyed materialized family” is ComputedMap + the mint recipe, and the “auto-mint keyed default” is get_or_insert_with — neither needs a separate type (see § Materialization).

Required. The keyed cell collections layer is normative for every lazily binding — it is not an optional lazily-rs extension. A conforming binding MUST implement ReactiveMap (at least its SourceMap specialization; ComputedMap where the binding supports derived slots), the ordered keyed tree (SourceTree), and keyed reconciliation, and MUST validate against the canonical fixtures in conformance/collections/. The single-writer / multi-write classification, merge: mechanism, and ingress rules below are exactly those defined above — the collection adds no new merge unit.

It conforms to the cell model when:

  1. Each entry is an ordinary cell — its single-writer / multi-write classification, merge: mechanism, and ingress rules are exactly those above. The collection adds no new merge unit; cell = merge unit still holds per entry.
  2. Value, set-membership, and order reactivity are independent: writing one entry’s value MUST NOT invalidate membership or order readers; adding/removing a key MUST NOT invalidate readers of unrelated entry values; and a pure reorder (atomic move) MUST NOT invalidate set-membership readers (len / contains) — only order readers (keys).
  3. A key resolves to a stable handle for the key’s lifetime; membership and order changes are signalled by their dedicated cells, never by mutating sibling entries.
  4. Atomic ordered move (move_to / move_before / move_after): reordering a key MUST keep the entry’s same cell handle, dependents, and lineage (not remove + re-mint) and bump only the order signal once.

Core surface vs. binding extensions

The clauses above are laws. This section says which methods a binding must expose for those laws to be observable, and — just as importantly — which methods are not spec surface at all.

The distinction exists because the family drifted without it. A survey of all nine bindings’ single-threaded maps found method counts from 11 to 32, and the drift splits cleanly in two: methods present in 8 or 9 bindings are laws with a fixture behind them and one binding that never implemented them, while methods present in 7 or fewer are conveniences with no fixture and no law — is_empty is len() == 0, len_untracked is a tracking-discipline escape hatch, reconcile belongs to a higher layer, insert is set under another name. Without this split, “kt is missing seven required methods” and “cs also ships a reconcile helper” are indistinguishable: both read as divergence, and the coverage matrix marks both green.

Core — REQUIRED

A conforming binding MUST expose all of the following, on every flavor it ships (single-threaded, thread-safe, async), spelled in the binding’s own naming convention:

groupmethods
constructionconstruct bound to a context
materializationget_or_insert_with (lazy mint-on-access); materialize_all (eager pre-mint)
kind specializationset (SourceMap only); materialize_all (ComputedMap only)
entry readobserve(key) -> Maybe<V>pure and non-minting; handle(key) -> Maybe<H>
reactive readskeys (order-tracked), len and contains_key (membership-tracked)
materialization planepresent_keys, present_count, is_present — non-reactive
orderposition, move_to, move_before, move_after
membershipremove
introspectionentry_kind

Two Core entries are load-bearing in ways that are easy to miss:

  • observe MUST NOT mint. A read that takes a factory and allocates is a different operation with a different contract; it cannot express “is this key absent” and it cannot be called from a reader without a write side effect.
  • handle is not optional. Clause 3’s stable handle and clause 4’s “same cell handle” are unassertable without a way to observe entry identity. A map whose entries are plain cached values rather than reactive nodes cannot satisfy this, and MUST be recorded as divergent rather than marked green.

The kind specializations bind per entry kind, not per flavor: set is SourceMap-only on all three flavors, materialize_all is ComputedMap-only on all three. Ordering, by contrast, binds every flavor — it touches no entry handle and awaits nothing, so it is neither thread-coloured nor async-coloured.

Extended — OPTIONAL, non-normative

The following ship in some bindings and are explicitly not spec surface. A binding MAY offer them, MAY name them differently, and MAY omit them entirely; their absence is not a conformance gap and MUST NOT be scored in the coverage matrix.

entry / entry_with (eager value-minting convenience), is_empty, len_untracked, get_or_insert_handle, reconcile, insert, Try*-prefixed naming variants, and any binding-local accessors for the membership/order signals themselves.

Ordered keyed tree

An ordered keyed tree (SourceTree) is a further composition: each node is (stable id, value cell, ordered keyed child collection). It conforms when per-node value reactivity holds (editing a node invalidates only that node’s readers), per-level membership/order reactivity holds (a sibling subtree or descendant change MUST NOT invalidate an unrelated level’s child readers), and child reorder inherits the atomic-move guarantee. The tree is still a composition of cells — not a new cell kind — so per-cell merge applies node-by-node.

This is the runtime substrate for stable keyed/wire addressing of collection entries (see the protocol spec’s node-key addressing) and for keyed reconciliation of document trees (minimal {insert, remove, move, update} ops per item → per-cell CRDT merge).

Keyed reconciliation

Reconciling a level diffs two keyed sequences by stable key, not position, emitting the minimal {insert, remove, move, update} op set. It conforms when reordering is move-minimized (keys already in relative order — the longest-increasing-subsequence over their prior indices — MUST NOT move; only the remainder emit move), and when applied to the reactive collection a stable entry (unchanged value, in the LIS) MUST NOT have its value cell invalidated by a sibling reorder. Applying this minimal op set per-cell is the enabling step for per-cell CRDT merge of a document tree — it replaces whole-subtree replacement with proportional-to-the-diff work.

Memoized semantic tree

The syntactic tree holds input cells; a semantic tree (unresolved prompts, drainable heads, summaries) is a layer of memoized computeds derived from it — one memo slot per node folding (node value, child derived values). It conforms when the derivation is incremental and glitch-free: editing one node recomputes only its ancestor chain (a sibling subtree’s derived value stays cached), and a node edit that does not change the folded result MUST NOT re-run a downstream consumer (memo equality guard). Semantics are derived, not materialized eagerly — cost is proportional to the diff, not the document.

Manufactured identity for text

Markdown has no inherent node ids, so reconciliation keys are manufactured from text in three layers: in-band anchors (exact, survive a body rewrite), content-derived hashes of normalized text (survive reflow/reorder, change on edit), and alignment by similarity (word-LCS ratio) to distinguish an edit (key inherited from the matched predecessor → an update) from a genuine insert. A true rewrite legitimately reads as insert+remove. This is why the controlled skeleton uses in-band markers: stable identity is the linchpin that keeps keyed reconciliation from degrading to whole-document replacement over unstable text.

Free-text CRDT + re-parse

For anchorless prose under concurrent edits the merge unit drops to characters: a Fugue/RGA-style character CRDT (each char an element with a unique id + left origin; deletes tombstoned) whose order is a pure function of the element set, so merge is commutative, associative, idempotent and concurrent same-point inserts converge with both preserved. The structural tree is then a projection of the merged text (re-parse → manufactured-identity keys → reconcile), not the merge unit. Honest floor: a true rewrite is a replace — no character identity survives it. The anchored layer keeps per-node lineage; the free-text layer’s guarantee is “merge the text, re-derive the tree.”

Delta sync (#lztextsync)

Whole-replica merge requires transporting the entire element set; a conformant binding also exposes delta synchronization so replicas converge by exchanging only what a partner lacks. Three operations, over the same element set:

  • version_vector(){peer → counter} — the greatest [OpId] counter this replica holds per originating peer, taken over both insert ids and tombstone (delete) ids. It is the compact frontier a replica publishes; an op (c, p) is unknown to a partner iff c > their_vv[p] (absent peer = 0).
  • delta_since(their_vv)[TextOp] — the ops this replica holds that their_vv has not observed: elements whose insert id is newer, plus elements whose tombstone id is newer (a fresh deletion of an already-shared element). Each TextOp is the transport form of one element — { id, ch, origin, deleted }. A whole-state snapshot is delta_since(∅).
  • apply_delta(ops) — applies a delta op list with the same algebra as merge: a new id adds its element (preserving that id); an incoming tombstone is merged sticky-minimally (concurrent deletes keep the smaller delete id); the local Lamport counter advances past every observed id. It is therefore commutative, associative, and idempotent, and re-applying a delta is a no-op.

Identity preservation is the load-bearing property: rebuilding a replica by apply_delta-ing a snapshot onto a fresh buffer keeps every character’s OpId, so a later concurrent edit merges without duplication — unlike re-parsing the text, which would mint fresh ids and double content on the next merge. This is what lets a canonical replica fork per-member replicas from an encoded snapshot and keep them converged by bidirectional delta_since/apply_delta.

The three operations and their convergence/idempotence/identity invariants are pinned by the compute fixture conformance/collections/textcrdt_delta_sync.json (#lztextsync): version-vector shape, bidirectional exchange, whole-snapshot fork identity, and no-op re-apply.

Move-aware sequence order

Sibling order under concurrency is a separate composition above per-cell value merge: a move-aware sequence CRDT (fractional-index positions tiebroken by peer). It conforms when a move is a single LWW reassignment of an element’s position — not delete + reinsert — so two concurrent moves of the same element converge to the later one without duplication, and a concurrent move + value-edit of one element both apply (position and value are independent registers). Removal is an LWW tombstone. This is the order layer beneath keyed reconciliation; it lives only at the multi-writer boundary, leaving the single-producer Snapshot/Delta mirror unchanged.

Tombstone garbage collection

Tombstones (both the sequence-CRDT LWW flag and the character-CRDT sticky delete, which carries the delete’s own id) accumulate without bound — the standard set-CRDT memory-bloat cost. Conformant GC is causal-stability-gated: a tombstone is collectable only once every replica has observed the deletion (the version-vector frontier supplied by the distributed plane, never a single replica’s clock). The sequence layer drops a stable tombstone directly (observationally inert: order/contains already skip it, re-merge re-adopts it as a tombstone, a genuine resurrection wins by LWW). The character layer is conservative: it collects a stable deleted element only when nothing references it as a left origin, so removal never orphans a survivor; interior tombstones are reclaimed bottom-up. Bloat is bounded to the multi-writer plane — the single-producer Snapshot/Delta mirror accrues none.

Scheduling (#lzspecgcdefer). The safety contract fixes which tombstones are collectable; it does not mandate when collection runs. A binding MAY defer GC to idle, batch it under memory pressure, or reclaim incrementally per anti-entropy round, provided no collectable tombstone is re-examined after reclaim and unbounded accrual is surfaced via instrumentation. Deferral changes only memory footprint, never observable values.

Reactive queues

A reactive queue (QueueCell) is a FIFO collection composed of cells — not a new cell kind — that adds queue semantics (push to tail, pop from head) to the reactive graph. Like the keyed collections above, it adds no new merge unit; each element’s value is an ordinary cell subject to the same single-writer / multi-write classification.

The distinguishing property of a reactive queue is that invalidation is scoped to reader kind, not to individual positions: a push invalidates length/empty/full readers (and the tail signal); a pop invalidates head/length/empty/full readers. The head reader observes the current head value — after a pop, the head reader sees the next element (or empty), not a stale value. There is no random-access queue[N] reader; per-position reactivity is the domain of SourceMap, not QueueCell.

QueueCell — SPSC primitive with MPSC usage rule

QueueCell is specified as a single-producer, single-consumer (SPSC) primitive: one writer owns the tail, one reader owns the head. The producer is the natural FIFO sequencer (push order = delivery order).

MPSC (multi-producer, single-consumer) is a usage rule on the same primitive, not a separate type. Multiple producers push to the same tail inside a batch(); the batch boundary serializes the pushes into a deterministic order. A conforming implementation MUST document the MPSC usage rule and MUST NOT introduce a separate MPSCQueueCell type.

Naming discipline. The cardinality of producers/consumers is not a type parameter. SPSCQueueCell would imply MPSCQueueCell / SPMCQueueCell / MPMCQueueCell siblings — but those shapes differ in semantics (invalidation model, handoff exclusivity), not cardinality. See § Future queue primitives for the genuinely distinct primitives (TopicCell, WorkQueueCell).

The queue family — two axes (semantics defines the primitive)

QueueCell, TopicCell, and WorkQueueCell are one family of reactive cursor-stream cells, separated by two orthogonal axes. Only the first defines the primitive; the second is a usage tier every primitive shares.

Axis 1 — consumer delivery semantics (the primitive axis). Where does each pushed element go?

Delivery semanticsEach element goes toConsumptionPrimitive
singlethe one consumerdestructive popQueueCell
competingexactly one of N consumersdestructive, exclusive handoffWorkQueueCell
broadcastevery subscribernon-destructive cursor readTopicCell

Axis 2 — topology / ordering (a usage tier, not a type). Producer and consumer counts — fan-in and fan-out — never change which primitive you have; they only set the ordering guarantee and its cost. Opt into exactly the tier you need:

Ordering tierGuaranteeMechanismCost
per-producer FIFO (default)each producer’s substream in push order; interleave arbitrarynone — ≡ multiplexing N single-producer channels at the consumerfree
agreed total orderone global order across producersa single leader sequencerone hop, single point of failure
agreed total order + HAtotal order surviving node death / partitionconsensus (Raft/Paxos)quorum latency

Why not name by cardinality or by fan-in/fan-out. Both are the same mistake — naming a primitive by topology instead of by delivery semantics:

  • SP*/MP* mixes two axes: QueueCell already covers SPSC and MPSC (both single-consumer; MPSC is the multi-producer usage). WorkQueueCell and TopicCell are each usable single- or multi-producer, so SPMC/MPMC cannot tell them apart.
  • FanOutCell/FanInCell fails identically. Fan-in (many producers → one consumer) is just QueueCell used multi-producer — the ordering tier, not a distinct primitive; a FanInCell collapses back into QueueCell. Fan-out (one → many consumers) is ambiguous between broadcast and competing — TopicCell and WorkQueueCell both fan out, and one FanOutCell name cannot say whether an element goes to all or to one. And fan-in and fan-out are not mutually exclusive (a topic or work queue may be many-producer and many-consumer at once), so they are orthogonal descriptors of wiring, not primitive identities.

So fan-in and fan-out are real and useful — but as the topology axes describing how a primitive is wired, not as names: fan-in = the multi-producer ordering tier (Axis 2); fan-out = multi-consumer, which subdivides into broadcast (TopicCell) vs competing (WorkQueueCell) by Axis 1. A primitive is always named by its Axis-1 delivery semantics, which carries meaning a topology name cannot.

Naming rationale (suffix). The shared family suffix is Cell; Queue appears only where consumption is destructive and exclusiveQueueCell and WorkQueueCell. A TopicCell is a non-destructive broadcast log (reading removes nothing; every subscriber reads every element), so it carries no QueueTopicQueueCell would conflate the deliberately contrasting messaging terms queue (consumed once) and topic (delivered to all). Parity is the Cell suffix + this family section, not a forced Queue infix.

What multi-producer actually buys (fan-in): a single merged fan-in aggregation stream, one shared backpressure / capacity bound, and — only at the total-order tier — an agreed order. Consensus is the price of agreed total order under partition only, never of multi-producer itself; per-producer FIFO needs no sequencer (it is multiplexed single-producer channels). Full cross-replica retention (every replica holds each item until all consumers ack + compaction) is a replication/HA cost, orthogonal to producer count — a non-replicated queue retains each item once.

Distribution cost differs by primitive (Axis 1 decides it):

PropertyQueueCellWorkQueueCellTopicCell
Each element →the one consumerexactly one of Nevery subscriber
Consumptiondestructive popdestructive, exclusive handoffnon-destructive cursor read
Distribution costleader-election HA (fence one consumer)assignment consensus (quorum)per-subscriber cursors, no consensus¹
Slow consumerfills queue → backpressurerouted to others; fine until all saturatedgrows its retention; evict on lease expiry
Consumer deathqueue stalls until failoverin-flight reassigned; no stallthat subscription grows; others fine
Orderingtotal FIFOassignment-FIFO; processing unorderedper-subscriber FIFO¹
Deliveryexactly-once local / effectively-once distributedat-least-once + idempotency = effectively-onceat-least-once per subscriber
Consensus?only for HAyes (assignment)only for agreed-order broadcast¹

¹ A topic is consensus-free only at the per-producer-FIFO tier; total-order broadcast (all subscribers agreeing on one order) is atomic broadcast ≡ consensus — the same cost as an ordered WorkQueueCell. Quorum-intersection safety for WorkQueueCell assignment is proven ReliableSync.majorities_intersect. Detailed semantics follow in § Future queue primitives.

Reactive shell vs storage backend

A QueueCell factors into two layers:

  ┌─────────────────────────────────────────────────────────┐
  │                   Reactive shell                         │
  │  head version cell │ tail version cell │ closed cell     │
  │  len / is_empty / is_full / head — reactive reads       │
  │  invalidation scoped by reader kind                     │
  └────────────────────────┬────────────────────────────────┘
                           │ QueueStorage trait
                           │  try_push(v) → Result<(), Full|Closed>
                           │  try_pop() → Result<T, Empty|Closed>
                           │  len() → usize
                           │  capacity() → Option<usize>
                           │  is_closed() → bool
   ┌───────────────────────┼───────────────────────────────┐
   │                       │                               │
   ▼                       ▼                               ▼
 VecDequeStorage      RaftQueueStorage              KafkaStorage
 (local default)      (embedded consensus;          (external broker;
                       per distributed-queue PRD)    via adapter)

The reactive shell owns the version cells and invalidation logic; it is storage-agnostic and is what the formal model (QueueCell.lean) pins. The reader-kind plane is pinned separately by QueueReaderKinds.lean: the invalidation set is proven to be a function of the transition alone (the bound, the pre-op length, emptiness, and the closed flag) and never of the elements, which is what licenses running the identical rule on all three flavors. It also proves the atomicity requirement above, and — as the queue-family instance of the thread-safe confluence obligation — that invalidation is order-independent and idempotent, because it consults membership in a set rather than a sequence. QueueFamilyReaderKinds.lean carries the same obligation to the other two primitives: the WorkQueueCell rule invalidates a reader kind iff that kind’s value moved (so it neither over- nor under-invalidates) and reads only the before/after counts; a TopicCell publish changes the observed suffix of exactly the connected subscribers; a cursor read depends on that subscriber’s own record alone, which is why the independence law holds structurally rather than by discipline; and safe GC provably changes no reader’s value, which is what licenses gc taking no context on any flavor. The storage backend owns the actual FIFO data structure and is pluggable via the QueueStorage trait (Rust) / concept (C++) / interface (Py/JS/etc.).

An implementation MUST split the shell from the storage. The shell MUST NOT assume a specific storage type (VecDeque, ring buffer, broker client). A binding MAY ship multiple backends; the default MUST be an unbounded VecDeque-backed storage.

Storage backend contract

Minimal required contract. A QueueStorage backend MUST implement exactly try_push / try_pop / len / is_closed / close. peek and capacity are optional capabilities with a default of “absent” (None): a backend that satisfies only the five required methods — a raw channel, a consuming stream, a Go channel — is fully conforming. It simply has no head reader (no peek) and no is_full reader (unbounded, capacity() → None), exactly as an unbounded backend has always had no is_full. A backend that can cheaply inspect its head MAY expose peek to gain a reactive head; a backend that is bounded MUST expose capacity() → Some(n). head was never in the MUST-reactive set (see “Named observables” below), so removing peek from the required contract removes no required reader.

Footnote — LookaheadShim. A caller that wants a head reader over a non-peekable backend MAY opt into a shell-level lookahead shim that prefetches (early-pops) one element into a one-slot buffer. This is SPSC-local only — early-popping is incorrect for competing-consumer or consensus backends, where an element must not be committed to one consumer before assignment. The shim is not part of the core contract.

A conforming backend MUST also satisfy:

  1. FIFO order: try_pop returns elements in the order they were try_push-ed. A backend that reorders or silently drops elements is non-conforming.
  2. Cardinality compatibility: the backend’s native producer/consumer shape MUST be a superset of the shell’s required shape. (SPSC shell = any backend; MPSC usage requires a backend that accepts multi-writer pushes.)
  3. Bounded contract (optional): a bounded backend exposes capacity() → Some(n) and try_push returns Full when at capacity. The overflow policy (block / drop-oldest / drop-newest / reject) is a backend property — the shell’s observable contract only distinguishes Full from Empty/Closed.
  4. Position identity: invalidation is phrased over reader kind (head/len/empty/full), not over storage indices. A ring-buffer backend whose slot index wraps MUST NOT cause spurious invalidations; the shell layers its own logical reader-kind derivations above the storage.

Closure and lifecycle

Closure is an observable contract, not a mechanism:

  1. try_pop on a closed, non-empty queue returns the next element (drain continues).
  2. try_pop on a closed, empty queue returns Closed — a signal distinct from Empty.
  3. try_push on a closed queue is an error, regardless of capacity.
  4. Close is idempotent (closing an already-closed queue is a no-op) and terminal (once closed, a queue cannot be reopened).

The mechanism (a dedicated closed cell, a flag in storage, a sentinel value) is a binding-level choice. The formal model pins closure as a monotonic flag: Closed_then_stays_Closed.

Bounded queue and reactive backpressure

When the storage backend is bounded (capacity() → Some(n)), the reactive shell exposes is_full as a reactive read. A consumer’s pop that transitions the queue from full to not-full MUST invalidate is_full readers (true → false), enabling push-side effects to react to capacity recovery without polling. This is the backpressure signal: a producer observes is_full and backs off; a consumer’s pop invalidates the producer’s is_full subscription and the producer resumes.

An implementation MUST expose is_full as a reactive cell when the backend is bounded. The unbounded default (capacity() → None) has no is_full reader to invalidate.

Ordering guarantee

ShapeGuarantee
SPSCTotal FIFO — pop order exactly matches push order. The producer is the single sequencer.
MPSCPer-producer FIFO — messages from each producer arrive in that producer’s push order. Inter-producer interleaving is deterministic within a batch() but implementation-defined across batches; under distribution it converges.

A consumer MUST NOT assume total-FIFO across multiple producers. If total order across producers is required, route all pushes through a single producer or use a consensus-backed storage backend (per the distributed-queue PRD).

Wire and snapshot shape

The QueueCell shell’s version cells have no own IPC schema — the head/tail/closed counters are trivial and not independently serialized. A queue reconciles by two complementary wire forms, chosen by plane:

  • Snapshot plane → storage-snapshot form. Full queue state is the storage backend’s snapshot form. The reference VecDequeStorage backend serializes as a JSON array (element order = FIFO order) for conformance fixtures; bindings MAY choose a more efficient binary encoding (bincode, postcard) for production. Cross-backend interop (e.g. VecDeque-backed on one peer, RaftQueueStorage on another) requires explicit storage-format agreement; the shell does not mandate a canonical storage snapshot.
  • Delta plane → op-log form. Incremental change is the ordered shell op-log — QueuePush / QueuePop / QueueClose — carried in a Delta like any other DeltaOp (protocol.md § QueueCell op-log delta form). These are shell ops (storage-agnostic append/remove-head/close), so they need no storage-format agreement; the op-log is the form reliable-sync fuses under backpressure (a queue cannot state-supersede coalesce — order, multiplicity, and the receiver’s pop position forbid it). The op-log delta form is normative here; a distributed storage backend (per the distributed-queue PRD) remains a separate v1 non-goal.

Distribution

Distribution of a QueueCell is a storage-backend property, not a shell property. A QueueCell is distributable iff its storage backend provides a distributed synchronization mechanism. The shell itself is sync-mechanism-agnostic.

v1 does not specify any distributed backend. The Native Distributed Queue PRD covers the future consensus-based RaftQueueStorage backend (Phase 1+) and the positioning relative to external brokers (Kafka, RabbitMQ, Redis Streams, SQS) via the QueueStorage adapter. CRDT-based distribution is explicitly out of scope for queues — destructive pop requires agreement, not merge (see the PRD’s § “Background: Why Consensus, Not CRDT”).

Threading, permissions, and instrumentation

  • Threading contract: QueueCell inherits the Context threading model. MPSC on the same Context uses batch(); cross-thread MPSC requires a ThreadSafeContext-bound flavor; cross-process requires a distributed storage backend.

    This sentence was aspirational prose for the whole v1 line: it named ThreadSafeContext as the way to reach cross-thread MPSC while no binding had a constructor that accepted one, so the only documented path to the guarantee was unimplementable. It is now normative — a thread-safe flavor is Core (see § “Core surface vs. binding extensions (queue family)”) and a binding that ships the family on a single-threaded context only is non-conforming on the other two flavors, not merely unfinished. Recording a capability as absent is fine; documenting a route to it that does not exist is not.

  • Permissions: over the distributed plane, push and pop are distinct capabilities under PeerPermissions (distributed.json). A peer MAY be granted push-only, pop-only, or both.

  • Atomicity: push and pop are individually atomic. Multi-op transactions (e.g., “enqueue N items then close”) MUST use batch() so a concurrent observer never sees a partial state.

  • Instrumentation: a binding SHOULD expose depth / push-count / pop-count metrics via the standard instrumentation surface (parallel to effect_queue_pushes / max_effect_queue_depth).

  • Named observables: is_empty and len are reactive reads dual to is_full. All three (is_empty / len / is_full) MUST be reactive when their respective conditions can change.

  • Demand-driven derivation (#lzspecdemanddriven): a reader-kind MUST be observable-consistent — a get returns the value consistent with all preceding ops — and a binding SHOULD defer its derivation until it has a subscriber (the measured collapse is ~32× per op — 327 ns → ~10 ns — per docs/relaycell-backpressure-analysis.md §5). A reader-kind is a derived value (a Slot), not an eagerly-written cell; an op with no subscriber to a given reader-kind SHOULD only mark it stale (O(1)) and derive it lazily on the next get, provided the derived value is consistent with all preceding ops. This preserves the observable contract (conformance fixtures read the values and MUST stay green) while an unsubscribed QueueCell collapses toward raw-storage cost — the reactive shell is charged only along a path an effect actually observes. A binding that eagerly derives reader-kinds remains conformant (values stay green by construction) but forgoes the write-cost win. See docs/relaycell-backpressure-analysis.md §5 (demand-driven reader-kinds) and §4.0 (the merge cost law).

Queue family extensions

QueueCell covers SPSC and MPSC. TopicCell is the broadcast member of the family; WorkQueueCell remains the competing-consumer extension. They differ in invalidation model and handoff semantics, not merely producer/consumer cardinality:

TopicCell (broadcast)

A broadcast topic: every subscriber receives every pushed element. Invalidation is “all subscribers,” not “head reader.” Each subscriber holds its own cursor; the topic retains an element until all durable cursors pass it (or a TTL expires). Reading is non-destructive — a subscriber’s advance never removes the element for others.

Relationship to QueueCell: not a multi-consumer queue. A QueueCell consumer destructively pops; a TopicCell subscriber reads by cursor and removes nothing. Different invalidation models in kind.

Normative state and operations

A TopicCell<T> state is the tuple (base_offset, elements, subscriptions):

  • base_offset: u64 is the absolute offset of elements[0]; end_offset is base_offset + elements.len.
  • elements: [T] is the retained append log in per-producer FIFO order.
  • subscriptions: Map<SubscriberId, Subscription> is keyed by stable subscriber identity. Each record contains an absolute cursor (the next offset to read), durability (durable or ephemeral), and connected. Every stored cursor MUST remain in base_offset..=end_offset. An ephemeral record MUST be connected: disconnecting it removes the record instead of persisting an offline ephemeral cursor.

The following operations and observables are required:

  • subscribe(id, durable) creates a new cursor at end_offset. Reconnecting an existing durable id MUST reuse its persisted cursor; it MUST NOT silently jump to the tail. Disconnecting an ephemeral subscription removes its record, so a later subscription with that id starts at the then-current tail.
  • publish(value) appends exactly one element and leaves every cursor unchanged. Every connected subscriber whose cursor is now behind end_offset is invalidated independently. An offline durable subscriber is not scheduled, but the element remains readable after it reconnects.
  • read(id) observes the element at that subscriber’s cursor without changing the log or any cursor. Reads from an unknown or disconnected subscription are unavailable. advance(id) moves only a connected subscription’s cursor by one; it is a no-op for unknown/disconnected subscriptions and at end_offset. Therefore base_offset ≤ cursor ≤ end_offset is preserved by every non-TTL operation, and an offline durable cursor is frozen until reconnect.
  • The retention frontier is the minimum cursor across durable subscriptions, or end_offset when there are none. Safe GC removes only offsets below that frontier, advances base_offset, and leaves absolute subscriber cursors unchanged. It MUST preserve every durable subscriber’s future read stream. Ephemeral subscriptions never hold the frontier.
  • GC is observational maintenance and invalidates no subscriber when it stays below the frontier. TTL-forced truncation beyond a durable cursor is a separate gap/resync policy and is not part of the v1 semantic core.

The canonical replay fixtures are topiccell_broadcast_cursor_isolation.json, topiccell_durable_replay_gc.json, and topiccell_ephemeral_lifecycle.json, with connected-session boundary behavior pinned by topiccell_offline_tail_bounds.json. The executable universal reference is lazily-formal/LazilyFormal/TopicCell.lean.

Distribution is cheap — no assignment consensus. Every subscriber gets every element, so there is no “who gets this” decision to arbitrate. A distributed topic is N independent per-subscriber cursor-queues — the per-peer DurableOutbox fan-out already pinned for reliable sync. At-least-once fan-out + idempotent apply per subscriber = effectively-once, no quorum. Caveat: this holds only at the per-producer-FIFO tier; total-order broadcast (all subscribers agree on one order) is atomic broadcast ≡ consensus.

Durable vs ephemeral subscriptions. A durable subscription persists its cursor and replays elements missed while offline; an ephemeral subscription sees only elements published while connected (fire-and-forget). Durable subscriptions drive retention.

Backpressure is per-subscriber, and the answer is the state-vs-event split. Each subscriber’s delivery buffer is itself a bounded QueueCell, so a slow subscriber’s is_full fires locally; what happens then is a per-subscription policy, and the right one depends on message semantics — the same dichotomy as outbox coalescing:

  • State topic (broadcasting a value) — old elements are worthless once superseded, so a lagging subscriber conflates to latest (drop intermediates, keep the newest): the LWW/last-value coalesce applied to a subscriber. Memory-bounded and effect-lossless — the laggard simply gets current state, and the producer feels nothing.
  • Event/log topic (each element a distinct event) — no conflation is possible (order and multiplicity are meaning), so a lagging subscriber resolves overflow one of three ways:
    • Drop (lossy, isolated)drop-oldest / drop-newest for that subscription; fast subscribers and the producer are unaffected. This is the broadcast default — coupling defeats fan-out.
    • Couple + backpressure (lossless, slowest-paces) — the subscriber withholds ack, retention grows, and the producer throttles to the slowest durable subscriber. Opt-in, for “all-must-receive-losslessly” topics; the operator accepts one slow subscriber pacing the whole topic.
    • Evict (bounded) — on sustained lag past a liveness lease, drop the whole subscription (§ Partition & eviction); it full-resyncs (durable) or resumes from now (ephemeral) on return.

So a TopicCell producer feels backpressure only if a subscription opts into coupling; by default a slow subscriber conflates (state) or drops/evicts (event), and failure is isolated — a dead subscriber grows only its retention, never stalling other subscribers or the producer. This is the opposite of QueueCell, where the single consumer is the backpressure path.

Consumer groups = a topic of work queues. The “Kafka consumer group” shape is a composition, not a fourth primitive: WorkQueueCell semantics within a group (competing) and TopicCell semantics across groups (each group gets the full stream). Model it as a TopicCell whose subscribers are WorkQueueCells.

Status: the local semantic core, conformance fixtures, and Lean proofs ship in v0.31.0. The consensus-backed distributed storage integration remains in the distributed-queue PRD Phase 3.

WorkQueueCell (competing consumers)

A work queue: N consumers compete for elements from a shared FIFO; each element is delivered to exactly one consumer (exclusive handoff).

Why pure CRDT cannot do this. A queue pop is not idempotent-commutative: two consumers concurrently popping the same head both survive a CRDT merge → duplicate delivery, and there is no “un-pop.” Exclusive handoff therefore needs a single serialization point for the assignment decision — a designated leader assigning each element to one consumer, or a consensus-committed assignment log. This is the queue’s CP nature made concrete.

Safety via quorum intersection. With consensus-committed assignment, “element X → consumer W” is a majority-committed log entry. Because any two majorities of an n-voter set intersect in ≥1 voter (ReliableSync.majorities_intersect / majorities_overcount), two conflicting assignments of the same element can never both commit — no double-delivery, ever — and a minority (no quorum) cannot commit an assignment, so it blocks rather than risking a duplicate.

Three populations — do not conflate.

  • Workers/consumers — any count; clients of the queue, not voters. Scale freely.
  • Replication peers (per-peer outbox) — any count; transport, not voting.
  • Voting replicas (order + commit assignments) — want an odd count: 2f+1 replicas tolerate f failures; odd maximizes tolerance-per-node and keeps quorum unambiguous. Make an even data set odd with a witness/arbiter (votes, holds no data).

Partition behavior. Only a partition holding a majority makes progress; the rest block (safe, no double-pop). A 2–2 split of 4 voters gives neither side a majority → both stall (the even-group hazard, fixed by a witness). A 2–2–1 split of 5 voters leaves no side with 3 → all block until partitions heal enough for some side to reach a majority. Odd counts protect two-way splits; they never guarantee a majority under multi-way fragmentation. Halting is correct — safety over liveness.

Delivery & lifecycle (#lzworkqueue):

  • Assignment IDs + ack/nack. Each handoff carries a delivery ID; the worker acks (done, remove) or nacks (requeue). Exactly-once commit is the consensus assignment; exactly-once effect needs an idempotent worker or a causal receipt on completion; delivery is at-least-once (redelivery on failure) — so effectively-once = at-least-once + idempotency.
  • Visibility-timeout / lease. An assigned-but-unacked element is leased with a TTL; on expiry it reassigns (worker presumed dead) — the per-item analog of the liveness lease.
  • Dead-letter queue + poison detection. An element exceeding a max-redelivery count (repeatedly crashing its worker) routes to a DLQ instead of redelivering forever.
  • Fairness / dedup. Assignment policy (round-robin, weighted, pull-based) + producer/consumer dedup keys are policy extensions; the portable core uses pull-based FIFO assignment.

Portable local-authority contract. Every binding ships the same in-process WorkQueueCell state machine. The owning instance is the serialization point; hosts MUST serialize calls using the binding’s ordinary context/threading boundary. A local instance does not claim to be a distributed consensus implementation. A cross-process or HA backend MUST put the claim decision behind a leader or consensus-committed assignment log while preserving these operations and outcomes:

  • push(value) -> item_id appends a fresh, monotonically increasing item to the pending FIFO.
  • claim(worker, now) -> delivery? removes the oldest pending item and creates a fresh, monotonically increasing delivery_id. The returned lease contains the stable item_id, value, worker, 1-based attempt, and deadline = now + visibility_timeout. Empty claim is a no-op.
  • ack(worker, delivery_id) -> bool removes only that worker’s matching in-flight delivery. Unknown, stale, or wrong-worker acknowledgements are no-ops. Re-acking is therefore idempotent.
  • nack(worker, delivery_id) -> bool validates ownership, then either appends the item to the pending tail or moves it to the dead-letter list when the configured max_deliveries has been reached.
  • reap_expired(now) -> count expires leases whose deadline < now (the deadline itself remains live), in ascending delivery-ID order. Each expired item is requeued at the pending tail or dead-lettered under the same attempt limit. A redelivery receives a new delivery ID while keeping its item ID and value.

visibility_timeout is a positive logical-clock duration and max_deliveries >= 1. Requeue-at-tail means retries do not block newer work. The portable reactive reader kinds are pending_len, is_empty (pending only), in_flight_len, and dead_letter_len; a mutation invalidates only readers whose values may have changed, and these reader-kinds SHOULD be demand-driven (#lzspecdemanddriven, see § “Demand-driven derivation” above). The canonical fixtures are conformance/collections/workqueue_competing_delivery.json and workqueue_lease_deadletter.json.

Assignment-FIFO ≠ processing-FIFO. Competing consumers trade ordering for parallelism: even if elements are assigned FIFO, N workers process concurrently, so completion order is unordered. A slow worker holding element 1 does not block element 2 (it routes to another worker) — the point of the primitive — but a consumer MUST NOT assume total processing order. For total order, use one consumer (QueueCell) or route through a single worker.

Pull beats push for balancing. A pull model (workers request when ready) is naturally load-balancing and backpressure-friendly — a fast worker pulls more, a saturated worker stops pulling. A push model needs the assigner to track worker capacity. Pull-based competing consumers are the simpler, self-throttling default.

Backpressure & failure — most resilient of the three. A single slow worker does not fill the queue (its items route to others); the queue backpressures only when all workers are saturated and depth grows. A worker death does not stall the queue — its in-flight (unacked) items reassign on lease expiry; throughput degrades, delivery continues. Opposite of QueueCell, whose single consumer dying stalls until failover.

Status: the portable local-authority state machine, reactive reader kinds, Lean safety model, and cross-language conformance are shipped. Distributed/HA assignment still requires the consensus adapter from the distributed-queue PRD Phase 2; the local shell must not be mistaken for that adapter.

Core surface vs. binding extensions (queue family)

The clauses above are laws. This section says which methods a binding must expose on QueueCell / TopicCell / WorkQueueCell for those laws to be observable, and which methods are not spec surface at all. It mirrors § “Core surface vs. binding extensions” for ReactiveMap, and exists for the same reason: without the split, “this binding ships no ordering on two flavors” and “that binding also ships a convenience helper” both read as divergence and both score green.

The split is drawn from a survey of the eight bindings that ship the family (lazily-cs does not yet). Normalising for naming convention, nine methods are present in 8 of 8 — unanimous, therefore law — and the rest are conveniences.

Core — REQUIRED

A conforming binding MUST expose all of the following on every flavor it ships (single-threaded, thread-safe, async), spelled in the binding’s own naming convention:

primitivegroupmethods
QueueCellconstructionconstruct bound to a context; bounded and unbounded forms
mutationtry_push, try_pop, close
reader kindshead, len, is_empty, is_full, is_closedall reactive
boundcapacity
TopicCellsubscriptionsubscribe, reconnect, disconnect
broadcastpublish
cursor readread (one element at the cursor), read_stream (the cursor’s tail), advance
retentioncursor/GC observables sufficient to assert the retention law
WorkQueueCellmutationpush, claim, ack, nack, reap_expired
reader kindspending_len, in_flight_len, dead_letter_lenall reactive

Three Core entries are load-bearing in ways that are easy to miss:

  • is_empty is Core here, unlike on ReactiveMap. On the map it is Extended, because it is just len() == 0. For a queue it is a named reactive observable that § “Threading, permissions, and instrumentation” already requires — “is_empty and len are reactive reads dual to is_full; all three MUST be reactive when their respective conditions can change”. The reader-kind independence law is stated over that triple, so dropping one makes the law unassertable.
  • A reader kind MUST be reactive, and polling a counter is not reactive. A version counter the caller must poll registers no dependency edge, so no reader can be invalidated by it and the independence law cannot be observed. A binding whose reader kinds are polled counters, or which derives them by diffing a snapshot, MUST be recorded as divergent rather than marked green.
  • try_pop MUST distinguish empty from closed. The closure lifecycle law needs Closed to be observably different from Empty; a pop returning a bare optional collapses them.

Ordering-style reasoning applies to the whole family: none of the Core surface touches an entry handle and none of it awaits, so it is neither thread-coloured nor async-coloured and binds every flavor. Confirmed empirically — across all eight bindings the queue-family sources contain no async fn, suspend fun, Future, Task, Promise, or await anywhere. A flavor is not permitted to colour pop.

Reader-kind derivation — one design, pinned

A conforming binding MUST implement reader kinds as memoized derivation with explicit invalidation: each reader kind is a derived node with no graph dependencies, and a successful op clears exactly those whose value provably changed, in one atomic frontier walk.

Three designs were in the field. This one is pinned because:

  • it is what § “Demand-driven derivation” already specifies, and the only one that realises its cost win — a version-cell design charges a graph write on every op even with no subscriber, which is precisely what that clause exists to avoid;
  • it is the shape the whole family converged on for ReactiveMap, where all nine bindings now hold a graph-agnostic ordering/bookkeeping core beside per-flavor version cells minted on each flavor’s own graph. Forking the queue family onto a different reader-kind design would split the family’s own structure for no gain.

A binding MAY keep a version-cell implementation and remain conformant only if reads still register a dependency edge and the invalidation set is exact; it forgoes the demand-driven win. Polled counters and snapshot diffing are not conformant under any reading.

Per-flavor obligation

A flavor MUST preserve the reader-kind independence law, the closure lifecycle, and the single-frontier-walk atomicity requirement — one op invalidates all of its changed reader kinds together, so a subscriber never observes len bumped while is_full has not yet flipped. On top of that it adds exactly its context-specific guarantee: confluence for thread-safe, eventual transparency for async.

Two flavor questions are settled here rather than left to each binding:

  • TopicCell fan-out needs nothing extra under a thread-safe context. Each subscriber owns its cursor and a cursor only advances, so concurrent subscriber reads commute and the retained-until-all-cursors-pass rule is a monotone function of the cursor set. Broadcast is confluent by construction.
  • WorkQueueCell lease expiry stays caller-driven on every flavor. reap_expired takes the current time as a parameter; a flavor MUST NOT substitute an owned timer. Time as an input is what makes lease expiry deterministic and fixture-replayable, and an async flavor that grew its own clock would be unreplayable for no gain.

Extended — OPTIONAL, non-normative

The following ship in some bindings and are explicitly not spec surface. A binding MAY offer them, MAY name them differently, and MAY omit them entirely; their absence is not a conformance gap and MUST NOT be scored in the coverage matrix.

push / pop as panicking or asserting variants of try_push / try_pop, peek, elements (whole-buffer snapshot), base_offset, snapshot / from_snapshot, restart, gc as a caller-visible entry point when retention is already automatic, reader_handle accessors for the reader-kind nodes themselves, and any binding-local accessors for the version cells.