Lakestream
Verification

Protocol specifications

The leaderless log protocol and task claiming are specified and model-checked in TLA+ and Fizzbee, on top of the coordination-delegated pattern.

The Lakestream Storage Spec's append and compaction protocols correspond to the coordination protocols specified in a separate repository, openlakestream/leaderless-log-protocol, which holds canonical Markdown specifications plus TLA+ and Fizzbee models of each one. The Markdown specification is the source of truth; the two models are independent implementations of it.

The repository organizes the material in three layers.

The coordination-delegated pattern

Layer 0 is the underlying architectural pattern. Worker nodes carry no consensus logic: they are stateless with respect to coordination and delegate ordering, mutual exclusion, state transitions, and failure detection to an external linearizable store. Workers never communicate with each other directly; all shared state — counters, locks, fences, cursors — lives in that store. This separates the coordination plane from the data plane, as opposed to embedded consensus such as Raft or Paxos, where every node participates in the consensus protocol.

The pattern names the primitives the store must provide:

PrimitiveRole
CompareAndSet(key, expected, new)State transitions, such as moving a log from open to fenced
AtomicIncrement(key, delta)Total ordering, such as offset assignment
ConditionalCreate(key, value)Mutual exclusion, such as lock acquisition
EphemeralRecord(key, value, session)Failure detection: the record is deleted when the client session expires

Supporting data primitives are Get, Put, Delete, ConditionalDelete, RangeDelete, and CeilingGet. Not every protocol uses every primitive.

The Layer 0 document is explicit about the cost. The coordination store becomes an availability dependency for coordination decisions, every coordinated operation costs at least one round trip to it, crash detection is bounded by its session timeout, and it must absorb the aggregate coordination load of all workers. The pattern suits systems where coordination is infrequent relative to the data path — one coordination call per batch of records rather than per record — and the document states plainly that embedded consensus is the better choice when coordination is the data path.

The leaderless log protocol

Layer 1 specifies a distributed append-only log with concurrent writers, compaction, and readers, where any writer can append without being elected leader. It models three interacting sub-protocols: the writer append path, the compaction index update, and the reader path.

  • Monotonic offsets. A writer obtains its offset from the coordination store's AtomicIncrement on a sequence counter. A batch write increments by the batch size in one call.
  • Log index. The index is held as linearizable key-value entries updated through atomic CAS. Keys are sparse: an entry covering N records is stored at the end offset of its range, and readers locate the covering entry for an arbitrary offset with CeilingGet.
  • Compaction. Replacing a range of WAL index entries is a non-atomic three-step CAS update: write the COMPACTED entry at the range end, RangeDelete the superseded WAL entries below it, then advance the compaction cursor. Ordering the write before the delete is what keeps every committed offset continuously readable, even if the compactor crashes between steps. Compaction operates on whole entries and cannot split a multi-record entry at an arbitrary boundary.
  • Fencing. A fenced log rejects further offset assignment, which prevents a stale writer from appending after a takeover.

Task claiming

Layer 2 specifies distributed task claiming: multiple workers independently scan for tasks, claim one, execute it, and release it. Mutual exclusion comes from ephemeral locks created with ConditionalCreate and released with ConditionalDelete, so at most one worker executes a given task at a time. When a worker crashes, its coordination-store session eventually expires and its locks are removed without an explicit failure detector. Tasks progress through INIT → PREPARED → COMMITTED, or move to a dead-letter state after a configured maximum number of failures.

Layers 1 and 2 are independent and composable. A system may use either alone, or both — for instance a log with background compaction scheduled as claimed tasks.

Verification status

Both protocols are model-checked in TLA+ (with the TLC model checker) and in Fizzbee, against the properties the repository's results tables report.

ProtocolSafety properties checkedLiveness properties checked
Leaderless logMonotonicOffsets, FencedRejectsAppends, CompactionPreservesData, NoPhantomEntries, CursorConsistency, NoOverlappingRanges, NoReaderError, SequentialCompactionSafetyAppendProgress, CompactionCompletes, ReaderEventuallySucceeds
Task claimingMutualExclusion, NoDoubleExecution, DLQOnlyAfterMaxFailures, LockConsistency, NoOrphanExecutionTaskCompletion, CrashRecovery, NoStarvation

The repository reports the leaderless-log properties as passing in the configurations it ships, states that the task-claiming properties are expected to pass, and requires that TLA+ and Fizzbee agree on every verdict. Where they disagree, the Markdown specification decides which model is wrong.

What the models do and do not cover

The models cover the protocol core, not the Ursa implementation. They are bounded: the published configurations run two or three writers, a handful of offsets, and a fixed compaction range, so a pass is a check over that state space rather than a proof about an unbounded system. Not every property is checked in both tools: the three liveness properties need fairness assumptions, and FencedRejectsAppends uses ENABLED, so all four run in TLA+ only. The repository also notes that one liveness property holds vacuously under the current step ordering and is retained as a regression guard. The specification also lists deliberate simplifications between model and implementation, including a compaction fast path implementations may take that the model does not. Model checking a protocol is not a statement about the correctness of code that claims to implement it.

Relationship to the Storage Spec

The Lakestream Storage Spec specifies its own append, compaction, and fencing protocols. They correspond to the ones specified here, but the conventions differ in places:

  • Offsets. This repository is 1-based: the first offset is 1, an index entry is stored at its range's inclusive end offset, and a reader locates the covering entry with CeilingGet — the smallest key greater than or equal to the query. The Storage Spec is 0-based: the first offset is 0, an index record's key is its range's exclusive end offset, and a reader uses a strictly-higher lookup — the smallest key greater than the query. The two are equivalent under translation: a 0-based exclusive end offset and the corresponding 1-based inclusive end offset are the same number, and a strictly-higher lookup on offset o matches a CeilingGet on offset o + 1.
  • Offset assignment. Layer 0 lists AtomicIncrement(key, delta) as the coordination primitive for ordering. This repository's own AssignOffset action fuses that increment with the index write into one atomic step, rather than modeling them as two separate operations. The Storage Spec's writers assign offsets the same way: a fused sequenced put that increments two counters — record count and cumulative byte size — and creates the index record in a single call.
  • Fencing. This repository models fencing as reversible: FenceLog and UnfenceLog (actions 10–11) move a log between OPEN and FENCED. The Storage Spec's fence is terminal: it is installed as one step of log deletion, blocks any new write lease from being acquired, and is never removed. Leases already held are drained: the catalog waits until the log has no lease records, and a lease record also disappears when the session that created it ends. Only then does the catalog delete the log's index records; it deletes no objects. If the catalog stops waiting before then, it leaves the log's records, objects and fence in place.
  • Compaction. The Storage Spec's compaction order matches actions 5–8 here: write the compacted record at the range's end offset, range-delete the superseded records, then advance a cursor. The Storage Spec leaves the cursor's exact scheme implementation-defined.

These are differences in convention, not in the underlying pattern. As elsewhere on this page, model-checking the protocol core specified here says nothing about the correctness of the Storage Spec's text or of any implementation, including Ursa's.

How Ursa uses it

Ursa runs this protocol with Oxia as its coordination store: Oxia holds catalog metadata, the offset index, cursor state, and the sequence keys that assign offsets. The Layer 0 document lists Oxia as supporting all of the pattern's primitives natively, including atomic increment through sequence-key deltas, ephemeral records with session management, ceiling lookups, and range deletes. Stores such as etcd, ZooKeeper, and FoundationDB support subsets and would have to emulate the rest.

The concrete key layouts and versions Ursa uses are described in the Storage Spec's Oxia mapping appendix. The write fence's key is normative and described in the Storage Spec's metadata store model. The visibility rules that follow from the three-step compaction update are in its Compaction section.

The S3-Queue example

The repository includes S3-Queue, a worked example that instantiates the leaderless log protocol on S3-compatible object storage with no separate coordination store: S3 conditional writes supply the coordination primitive. It ships a system-level specification alongside a command-line implementation, and the Layer 1 document maps each protocol action onto the corresponding source file and coordination call.

It is an example of how the protocol can be instantiated, not a component of Ursa and not the only possible mapping.