Delay Shard and materialized state
The Delay Shard is the unit that keeps ordering, ownership, local state, and recovery aligned.
one physical Command Topic partition
= one Shard Log source order
= one Delay Shard
= one Owner Lease and owner epoch
= one RocksDB database
= one checkpoint and recovery unit
A Worker can host several Delay Shards, but V1 does not create a cross-Shard transaction. Capacity scales by adding physical source partitions and assigning their Shards across Workers.
Three storage roles
| Storage | Answers | Why it exists |
|---|---|---|
| Command Topic / Shard Log | What happened, and in what order? | Durable replay input for Commands and System Mutations |
| RocksDB | What is the current state now? | Fast point queries, cancellation, rescheduling, due scans, retry state, and local recovery progress |
| Object Store | Where are large immutable objects? | Large payloads and complete checkpoint objects |
The Shard Log is authoritative for ordered facts. RocksDB is the durable materialized view produced by applying those facts. Object Store is neither the ordered log nor the live state database.
Why one RocksDB database per Shard?
Using the same boundary for ownership and storage makes takeover and deletion easier to reason about:
- a Shard checkpoint is a complete physical image rather than a range inside a shared Worker database;
- a failed or migrating Shard does not require splitting another Shard's files;
- disk, compaction, restore, and local deletion can be gated by the exact Shard lifecycle;
- the Shard Runtime remains the only semantic writer.
The tradeoff is operational. A Worker may open many databases, so file descriptors, cache, compaction bandwidth, disk pressure, and placement must be governed across all hosted Shards.
Single-writer application
Producer callbacks, timers, control requests, and background cleanup do not write semantic state independently. They first become ordered Shard Log mutations. The Shard Runtime then applies one source position at a time and commits the business result together with its new applied position.
This avoids relying on thread scheduling to decide whether Cancel, Reschedule, Publish Outcome, expiration, or an operator action wins.
Timeline, not an in-memory authority
The persistent time index supports long delays, point cancellation, rescheduling, and recovery without loading every future message into memory. An in-memory timer wheel can be added later as an acceleration cache, but it cannot become the authority for messages that must survive process and disk loss.
Read Shard Log and Source Position next, then checkpoints and recovery.