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
114 changes: 89 additions & 25 deletions ARCHITECTURE.md

Large diffs are not rendered by default.

44 changes: 39 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -195,6 +195,8 @@ match watcher.watch_all_from(&cursor, tx.clone()).await {
}
```

The full replay only re-delivers what NATS still holds. On a bucket with `max_age` or `discard: old`, values that aged out are gone from it; [`watch_applied`](#cursor-expiry) restores those from an artifact.

`watch_prefix_from()` works the same way for prefix-filtered streams, and
`watch_prefixes_from()` resumes the union of several prefixes on one
multi-filter consumer.
Expand Down Expand Up @@ -298,7 +300,7 @@ let (resume, store) = AppendLogSnapshot::load(Path::new("/var/lib/svc/state.snap

let final_cursor = watch_applied(
watcher, WatchScope::All, Some(resume),
Some(reader), // arms the cursor-expired stale-key resync; None to skip
repair, // cursor-expiry repair (see "Cursor expiry"); None to skip
Some(store), None, // store; export-request channel
BatchConfig::default(),
parse, apply, on_applied, shutdown,
Expand All @@ -318,8 +320,8 @@ let final_cursor = watch_applied(
watcher,
WatchScope::All, // or Prefix("node.".into()) / Prefixes(vec![...])
Some(resume), // Option<WatchCursor> — resume here, or None
Some(reader), // Option<Arc<dyn KvReader>> — arms the
// cursor-expired stale-key resync, or None
repair, // ExpiryRepair — how an expired cursor is
// repaired (see below), or None
Some(store), // any SnapshotStore (e.g. AppendLogSnapshot), or None
None, // Option<mpsc::Receiver<ExportRequest>> — live exports
BatchConfig::default(), // 10ms window, 100 updates per batch
Expand All @@ -332,14 +334,46 @@ let final_cursor = watch_applied(

A batch closes when `window` elapses or it hits `max` updates, whichever comes first. Then, in order: `apply(batch)` runs to completion, the cursor advances to the batch's highest revision, the batch + cursor are folded into the `store` atomically (on a blocking task), and `on_applied` fires.

With a `store`, resume from the cursor the store reports (`load`/`open`), not one you persisted in `on_applied`: across a transient store failure, `on_applied` reports a cursor ahead of what the store has made durable.

Persist the cursor on receipt instead and a crash between receive and apply loses data: the cursor reads "caught up to rev N" while rev N sits in an unapplied buffer, and the next resume starts past it. `watch_applied` checkpoints at the applied cursor, so a persisted cursor always means every update up to it has been applied.

- `parse` returning `None` (corrupt bytes, irrelevant key) still advances the cursor — nothing to apply means nothing to skip.
- On `CursorExpired`, it falls back to a full watch automatically. With a `reader` wired, it first diffs the fold against the bucket's live keys and applies synthetic deletes for keys that vanished during the gap (their delete markers were evicted with the cursor) — the one case the fallback re-list can't cover.
- On `CursorExpired` it repairs the fold, then keeps watching. See [Cursor expiry](#cursor-expiry).
- It returns the final applied cursor on shutdown or stream close.

`apply` runs inline. If it panics, the panic aborts the watch.

### Cursor expiry

A cursor expires when NATS has evicted the revisions after it: at resume, or mid-watch when retention overruns a live All-scope watch (the floor guard). Both take the same repair, and the right repair depends on what the bucket's retention evicts.

**Buckets that keep current values** (`discard: new`, no `max_age`). Only old history and delete markers get evicted, so a key missing from the bucket was deleted. `ExpiryRepair::Relist(reader)` lists the bucket's live keys, deletes the fold's in-scope keys that are gone, then re-lists every live value.

**Buckets that evict current values** (`max_age`, or `discard: old` under `max_bytes`). A key missing from the bucket may just have aged out, and a write made while the node was offline may have aged out too. NATS can't tell those apart from a delete, so the key list can't be trusted. `ExpiryRepair::Restore { reader, restore }` replaces the in-scope fold with the newest published artifact (which holds every write minus every real delete) and resumes from its cursor. It never moves a key to an older revision and never deletes a key the bucket lists as live, so even an artifact exported while its exporter was catching up produces no transient phantom deletes or regressions. `Relist` is refused on these buckets: the watch fails rather than delete valid keys.

`ExpiryRepair::Auto { reader, restore }` reads the bucket's retention at expiry and picks for you:

```rust
use slipstream::{ArtifactRestore, ExpiryRepair};

let repair = ExpiryRepair::Auto {
reader: bucket.reader(),
restore: Arc::new(ArtifactRestore::new(
transport, // Arc<dyn ArtifactTransport>, same key the fleet exports to
"routes/latest",
"/var/lib/svc/scratch", // same filesystem as the fold
|artifact, dest| FjallSnapshot::import(artifact, dest, config.clone()),
)),
};
```

A restore only accepts an artifact that is ahead of the local fold, covers the watch's scope (recorded in the artifact by the exporting `watch_applied`), and whose cursor NATS still retains. Otherwise there is no safe recovery: the watch fails and logs the reason at `error`. **Run exports well inside `max_age`** (or the `discard: old` turnover), or an expired node has nothing fresh to restore from.

A node with no cursor on a bucket that has already evicted current values seeds from the artifact too, since a re-list alone would miss whatever aged out.

`Some(reader)` from older code still compiles and means `Relist`. `None` (or no `store`) falls back to the re-list alone and warns that deleted keys may persist.

## NATS mapping

| Concept | NATS primitive |
Expand Down Expand Up @@ -371,7 +405,7 @@ if caps.global_ordering { /* VersionToken is globally ordered across keys */ }
| `NotConnected` | Operation before `connect()` | Call `connect()` |
| `AlreadyExists` | `create()` on a live key | Read current state, decide |
| `RevisionMismatch` | CAS conflict on `update()` / `delete_with_version()` | Re-read, retry |
| `CursorExpired` | `watch_*_from()` cursor compacted by NATS | Fall back to `watch_all()` |
| `CursorExpired` | `watch_*_from()` cursor compacted by NATS (at resume, or mid-stream when retention overruns an All-scope watch) | `watch_applied` repairs it; raw callers fall back to `watch_all()` on buckets that keep current values |
| `WatchError` | NATS stream dropped | Re-subscribe |

## Credentials
Expand Down
9 changes: 5 additions & 4 deletions mise.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,11 @@ yamlfmt = "0.12.1"
dprint = "0.50.2"
"ubi:nats-io/nats-server" = "2.14.1"
# Real-S3 transport tests (tests/transport_s3.rs): single static binaries, no
# Docker — same pattern as nats-server. mc creates the test bucket. (ubi can't
# install minio: its releases ship extension-less raw binaries that ubi's
# asset filter rejects, hence the asdf-backend plugin.)
"asdf:mise-plugins/mise-minio" = "2025-09-07T16-13-09Z"
# Docker — same pattern as nats-server. mc creates the test bucket. MinIO no
# longer publishes community server binaries (dl.min.io answers 410 Gone), so
# the same pinned release is built from source with `go install` (~45 s).
go = "1.25"
"go:github.com/minio/minio" = "RELEASE.2025-09-07T16-13-09Z"
"aqua:minio/mc" = "RELEASE.2025-08-13T08-35-41Z"

[tasks.format]
Expand Down
Loading
Loading