Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 9 additions & 8 deletions ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -123,7 +123,7 @@ Every non-`_from` watch is a **state-sync** stream (NATS `DeliverPolicy::LastPer
| `VersionToken` | Opaque per-key version (NATS: 8-byte u64; FDB: 10-byte versionstamp) | Not a wall-clock timestamp; not globally ordered |
| `KvEntry` | One key + value + version from a read | Not a watch event; immutable once returned |
| `KvUpdate` | One watch event: `Put`, `Delete`, or `Purge` | Not a read result; carries deletes too |
| `Snapshot` | Deduplicated KV state + cursor persisted to disk | Not the source of truth; a cache of NATS |
| `Snapshot` | Deduplicated KV state + cursor persisted to disk; on a bounded log, a replica of record | Not a cache: NATS no longer holds what it evicted |
| `SnapshotWriter` | Append-only log of `KvUpdate`s; no in-memory state beyond a counter | Not the in-memory cache itself |
| `SnapshotStore` | Trait: the durable-fold contract — atomic `apply(batch, cursor)`, `load`, `get`, `range` | Not a serving index; stops at fold + cursor + query |
| `AppendLogSnapshot` | Default `SnapshotStore`: append-only log + in-RAM fold (pure-Rust) | Not for folds larger than RAM |
Expand Down Expand Up @@ -417,9 +417,9 @@ fn range(prefix) -> Vec<KvEntry>; // ordered prefix scan

Three invariants bind every implementation:

- **Pure function of the log.** Delete the store, replay every update with revision `> cursor`, and the state is identical. The store caches the fold; NATS is the source of truth.
- **Fold + tail = truth.** The fold at cursor C, plus a resume from C, is every write minus every real delete. A NATS cursor means every *retained* message at or below it is applied — not "the truth at C": a key whose latest write is after C arrives with the resume.
- **Cursor-after-apply.** `apply` makes data and cursor durable together, so the cursor never names a revision whose data is absent — one transaction on a transactional backend, data-then-cursor on the append log (a torn write leaves data *ahead* of the cursor, which replay re-folds, never skips).
- **Snapshot is a cache.** A tail lost to power loss (under a no-sync durability mode) is rebuilt by resuming the watch from the recovered cursor.
- **Replica of record.** NATS is a bounded log and keeps only the retained tail, so a fold can't be rebuilt from NATS alone. A tail lost to power loss (under a no-sync durability mode) leaves a consistent fold at an earlier cursor, rebuilt by resuming from the recovered cursor while it is inside NATS's retention, and by the cursor-expiry repair (an artifact restore) once it isn't. A lost or corrupt fold is rebuilt by importing an artifact, never by deleting it and re-listing NATS (complete only on a bucket that has never evicted a current value).

`watch_applied` is generic over `SnapshotStore`: on each flush, after `apply` returns, it hands the raw batch + post-apply cursor to `store.apply(...)` on a blocking task. The trait stops at fold + cursor + query; serving structures built from the fold (routing rings, hashrings) live in the consumer, which reads them out via `get`/`range`.

Expand Down Expand Up @@ -473,7 +473,7 @@ write_update() compact() [blocking: replay → dedup

`compact()` flushes the BufWriter first so un-checkpointed records survive. It reads the current file, replays it, writes to a same-directory tempfile (same filesystem = atomic rename, no `EXDEV`), `sync_all`s, then renames.

`checkpoint()` writes only a cursor record and calls `BufWriter::flush()` — a `write(2)` into the page cache. This survives a process crash but NOT a power loss. The only `fsync` is in `compact()`. The snapshot is a cache; a lost tail is rebuilt from a NATS scan + watch replay.
`checkpoint()` writes only a cursor record and calls `BufWriter::flush()` — a `write(2)` into the page cache. This survives a process crash but NOT a power loss. The only `fsync` is in `compact()`. A tail lost to power loss leaves a consistent fold at an earlier cursor, rebuilt by resuming from the recovered cursor while it is inside NATS's retention, and by the cursor-expiry repair (an artifact restore) once it isn't.

### Export Lease (`export_lease.rs`)

Expand Down Expand Up @@ -610,7 +610,7 @@ async-nats ≤0.46 has a race: the server can deliver the first batch of push-co

### Why checkpoint() does not fsync?

Checkpoints are frequent (every N watch events). An fsync per checkpoint would add milliseconds of disk-sync latency to the hot watch path. Since the snapshot is a cache backed by NATS, a tail lost to power loss is rebuilt from a NATS replay — not a correctness failure. The only `fsync` is in `compact_to_file()`, where it guarantees the new compact file is durable before the atomic rename replaces the old one.
Checkpoints are frequent (every N watch events). An fsync per checkpoint would add milliseconds of disk-sync latency to the hot watch path. A tail lost to power loss leaves a consistent fold at an earlier cursor (data and cursor commit together), rebuilt by resuming from the recovered cursor while it is inside NATS's retention, and by the cursor-expiry repair (an artifact restore) once it isn't — not a correctness failure. The only `fsync` is in `compact_to_file()`, where it guarantees the new compact file is durable before the atomic rename replaces the old one.

### Why raw backend-dir artifacts instead of a logical re-encode?

Expand Down Expand Up @@ -700,6 +700,7 @@ Production code and an exhaustive model checker running the same logic closes th
| `CursorExpired`, bucket keeps current values | `watch_applied` lists live keys, diffs fold, applies synthetic deletes, then falls back to state-sync re-list | Automatic; raw callers use `stale_keys()` manually |
| `CursorExpired`, bucket evicts current values | `watch_applied` restores the in-scope fold from the newest artifact and resumes from its cursor | Automatic with `ExpiryRepair::Restore`/`Auto` |
| Fold holds data but no cursor (torn first checkpoint) | Repaired like an expiry at revision 0 before watching | Automatic with a repair armed |
| Every message after a node's cursor SUPERSEDED (every key rewritten while it was down), evicting bucket | The conservative resume check sees the cursor as expired though nothing was evicted; the restore then needs an artifact ahead of the node, and without one the watch fail-stops (safe, unavailable) — pinned by `tests/eviction.rs::supersession_only_restart_fail_stops_without_a_newer_artifact` | Restarts once the next export round publishes; export more often on small hot buckets. Telling supersession from eviction needs the cursor's message timestamp against `max_age` (not stored today) |
| A write ages out before ANY node folds it | Every node that needed it fail-stops (no fold anywhere holds it; `tests/model_fleet.rs` proves this is the only way the fleet fail-stops) | Data loss by retention; widen retention or export more often |
| Key-listing resync on a bucket that evicts current values | Refused: `KvError::WatchError` ("evicts current values"); fold untouched | Wire `ExpiryRepair::Restore`/`Auto` with an artifact source |
| Restore: no artifact ahead, stale artifact, or one that doesn't cover the scope | `KvError::WatchError`, reason logged at `error`; fold untouched | Fix the export pipeline (run exports well inside `max_age` / `discard: old` turnover); restart |
Expand All @@ -708,9 +709,9 @@ Production code and an exhaustive model checker running the same logic closes th
| `RevisionMismatch` on CAS | Concurrent writer won the race | Re-read with `entry()`, resolve, retry |
| `AlreadyExists` on `create()` | Key already present; caller's create was not exclusive | Read live value, decide whether to proceed |
| Snapshot truncated tail | `load()` discards partial final record; earlier records intact | Resume from recovered cursor; tail re-folded |
| Snapshot mid-file CRC mismatch | `SnapshotError::Corrupted` | Delete snapshot, full NATS replay |
| Snapshot wrong format version | `SnapshotError::InvalidFormat` | Delete snapshot, full NATS replay |
| `compact()` I/O error | Writer poisoned; subsequent writes return `Io` | Delete snapshot, rebuild from NATS |
| Snapshot mid-file CRC mismatch | `SnapshotError::Corrupted` | Import the latest artifact (a NATS re-list loses everything evicted; complete only on a bucket that never evicted a current value) |
| Snapshot wrong format version | `SnapshotError::InvalidFormat` | Same: import the latest artifact |
| `compact()` I/O error | Writer poisoned; subsequent writes return `Io` | Reopen (the old file is intact until the atomic rename); if unreadable, import the latest artifact |
| Synadia Cloud stream limit | Raw API path treats as non-fatal; verifies bucket with `get_key_value` | Non-fatal if bucket exists |
| Crash between payload upload and pointer swap | Old pointer remains fully consistent; payload orphaned | Next export round publishes new pointer; stale payload pruned after grace |
| Slow exporter after newer round published | `pointer_publish_allowed` returns false → `SupersededByNewer` | Lease abandoned, local artifact deleted; payload orphaned until prune |
Expand Down
8 changes: 5 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@ beyond-slipstream = { version = "0.7", features = ["transport"] } # export/impor
| `VersionToken` | Opaque version — NATS: u64 revision; FDB: 10-byte versionstamp |
| `KvEntry` | One key + value + version from a read |
| `KvUpdate` | One watch event: `Put`, `Delete`, or `Purge` |
| `Snapshot` | Deduplicated KV state + cursor at a point in time. Disk cache, not source of truth |
| `Snapshot` | Deduplicated KV state + cursor at a point in time. On a bounded log, a replica of record: NATS no longer holds what it evicted |
| `SnapshotWriter` | Append-only log of `KvUpdate`s; survives restarts without a full NATS scan |
| `SnapshotStore` | Trait: the durable-fold contract — `apply` (data + cursor, atomically), `load`, `get`, `range` |
| `AppendLogSnapshot` | Default `SnapshotStore`: the append-only log + an in-RAM fold (pure-Rust, small state) |
Expand Down Expand Up @@ -216,6 +216,8 @@ if let Some(snap) = snapshot::load(Path::new("/var/lib/svc/state.snap"))? {
}
watcher.watch_all_from(&snap.cursor, tx).await?;
} else {
// Complete only if the bucket has never evicted a current value; on a
// bounded log, seed a fresh node from an artifact instead.
watcher.watch_all(tx).await?;
}
```
Expand Down Expand Up @@ -244,7 +246,7 @@ while let Some(update) = rx.recv().await {

This loop has a trap: `current_cursor` must track what `cache.apply()` has consumed, not what `rx.recv()` delivered. Get it wrong and a crash skips updates on resume. [`watch_applied`](#applied-watch) runs this loop for you with that invariant enforced.

The snapshot is a cache. Delete it and the service falls back to full replay on next start.
On a bounded log the snapshot is a replica of record, not a cache. Deleting it and re-listing NATS loses everything NATS has evicted: rebuild a lost or corrupt fold by importing the latest artifact. Re-listing is complete only on a bucket that never evicts current values.

### File format

Expand Down Expand Up @@ -276,7 +278,7 @@ pub trait SnapshotStore: Sized + Send {
}
```

Every backend keeps the same invariants: the fold is a pure function of the log (delete the store, replay from the cursor, get identical state), the cursor never names a revision whose data isn't durable (cursor-after-apply), and the store is a cache — a tail lost to power loss is rebuilt by resuming the watch.
Every backend keeps the same invariants. **Fold + tail = truth**: the fold at cursor C, plus a resume from C, is every write minus every real delete (a cursor means every *retained* message at or below it is applied). **Cursor-after-apply**: the cursor never names a revision whose data isn't durable. **Replica of record**: on a bounded log the fold can't be rebuilt from NATS alone; a tail lost to power loss is rebuilt by resuming from the recovered cursor while it is inside NATS's retention, and by the cursor-expiry repair (an artifact restore) once it isn't, and a lost or corrupt fold is rebuilt by importing an artifact, never by re-listing NATS.

| Backend | When | Notes |
| ------- | ---- | ----- |
Expand Down
69 changes: 65 additions & 4 deletions src/applied.rs
Original file line number Diff line number Diff line change
Expand Up @@ -82,8 +82,8 @@ use crate::snapshot::{SnapshotError, SnapshotStore};
/// embedded cursor is exactly the applied cursor at the moment of export.
///
/// The export result (or error) comes back on `reply`; an export failure is
/// reported there and the watch keeps running (the snapshot is a cache — a
/// failed artifact is the requester's problem, not the fold's).
/// reported there and the watch keeps running (a failed export leaves the fold
/// untouched; the next export round retries).
pub struct ExportRequest {
/// Where the artifact directory will be created. Must not exist (or be an
/// empty directory); same filesystem as the fold for cheap hardlinks.
Expand Down Expand Up @@ -328,10 +328,29 @@ where
// Repair channel, armed only when there is a repair to run AND a store to
// run it against (both repairs edit the fold). The watch task sends the
// repair here and waits for the ack/reply before starting its next watch.
let mode: ExpiryRepair<S> = repair.into();
// A restore replaces the fold, so asking for one without a fold is a
// contradiction: on a bounded log such a consumer would silently hold
// NATS's retained view (missing everything evicted) while believing it
// repairs. Refuse it up front; `ExpiryRepair::None` is the explicit way
// to accept the retained view.
if store.is_none()
&& matches!(
mode,
ExpiryRepair::Restore { .. } | ExpiryRepair::Auto { .. }
)
{
return Err(KvError::WatchError(
"ExpiryRepair::Restore / ::Auto repair the fold, so they need a store: a consumer \
without one cannot repair a bounded log. Pass a SnapshotStore, or \
ExpiryRepair::None to accept NATS's retained view"
.into(),
));
}
let (repair_handle, mut repairs): (
Option<RepairHandle<S>>,
Option<mpsc::Receiver<RepairRequest<S>>>,
) = match repair.into() {
) = match mode {
ExpiryRepair::None => (None, None),
mode if store.is_some() => {
let (tx, rx) = mpsc::channel(1);
Expand Down Expand Up @@ -963,7 +982,23 @@ async fn run_watch<S: Send + 'static>(
}
seed
}
None => false,
None => {
// Nothing armed: the re-list stands as the fold. Say so when
// the bucket has already evicted current values — the copy
// then lacks whatever aged out, and nothing will repair it.
if watcher
.retention()
.await?
.is_some_and(|r| !listing_is_truth(&r))
{
warn!(
"no resume cursor and no repair armed, on a bucket that has already \
evicted current values: the re-list lacks whatever aged out. Wire a \
store and ExpiryRepair::Auto to seed from an artifact"
);
}
false
}
};

loop {
Expand Down Expand Up @@ -3583,4 +3618,30 @@ mod tests {
"the restored fold must equal the artifact's"
);
}

/// A restore needs a fold to restore into: `Restore`/`Auto` without a
/// store is refused up front instead of silently running on NATS's
/// retained view.
#[tokio::test]
async fn restore_without_a_store_is_refused() {
let artifact = MockRestore::new(vec![], 10, Some(&[""]));
let watcher = Arc::new(ScriptedWatcher::new(evicting(9), vec![]));
let (_sd_tx, sd_rx) = watch::channel(false);
let err = watch_applied(
watcher,
WatchScope::All,
None,
restore_mode::<AppendLogSnapshot>(artifact),
None::<AppendLogSnapshot>,
None,
BatchConfig::default(),
parse_put,
|_b: Vec<Vec<u8>>| {},
|_| {},
sd_rx,
)
.await
.expect_err("Restore without a store must be refused");
assert!(err.to_string().contains("need a store"), "{err}");
}
}
10 changes: 9 additions & 1 deletion src/kv.rs
Original file line number Diff line number Diff line change
Expand Up @@ -216,7 +216,11 @@ pub struct KvEntry {
pub enum KvUpdate {
/// Key was created or updated.
Put(KvEntry),
/// Key was deleted.
/// Key was deleted — by a delete marker, or (NATS) by a
/// [`delete_with_version`](KvWriter::delete_with_version) tombstone, which
/// watches report as a delete with the tombstone's revision: the same rule
/// `get`/`scan`/`keys` apply. Only [`entry`](KvReader::entry) exposes the raw
/// tombstone.
Delete { key: String, version: VersionToken },
/// Key was purged (NATS-specific: all history removed).
/// Stores without purge semantics should map this to Delete.
Expand Down Expand Up @@ -393,6 +397,10 @@ pub trait KvWriter: Send + Sync {
/// CAS-gated delete: delete only if current version matches `expected`.
/// Returns `RevisionMismatch` on conflict.
/// Writes an empty value (logical delete) so concurrent writers get a conflict.
/// Every read path treats it as a delete — `get`/`scan`/`keys` hide it, and
/// watches (hence folds) receive a [`KvUpdate::Delete`] — except
/// [`entry`](KvReader::entry), which exposes the tombstone and its version for
/// CAS callers.
async fn delete_with_version(
&self,
key: &str,
Expand Down
Loading
Loading