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
4 changes: 3 additions & 1 deletion ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -478,7 +478,7 @@ Cursor: cur_len:u8 ++ cursor_bytes

Version bytes are stored as length-prefixed raw bytes, not a fixed `u64`. A 10-byte FDB versionstamp round-trips intact; a `u64`-only field would flatten it to 0 and break every subsequent CAS on a restored entry.

CRC covers from the type byte through the end of the record. A truncated final record (crash mid-write) is silently discarded. A CRC mismatch in the middle of the file returns `SnapshotError::Corrupted`.
CRC covers from the type byte through the end of the record. A truncated final record (crash mid-write) is silently discarded. A CRC mismatch in the middle of the file returns `SnapshotError::Corrupted`. Lengths are read before the CRC can be checked, so a corrupted length that points past EOF reads as a truncated tail and drops the records after it: the fold lands at an earlier cursor, the same shape as a tail lost to power loss. A value is at most 64 MiB (NATS's payload ceiling), enforced by the writer and the reader alike; a larger length is treated as corruption.

### State Machine

Expand Down Expand Up @@ -742,11 +742,13 @@ Production code and an exhaustive model checker running the same logic closes th
| `Timeout` on any NATS op | CLOSE_WAIT half-dead TCP parks the `await` without this guard | Call `shutdown()` + `connect()` |
| `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 |
| Value over 64 MiB written to the append log | `apply` returns `InvalidFormat` before writing anything; under `watch_applied` the batch re-queues until the watch fail-stops | Unreachable from a NATS watch (its payload ceiling is the same); a direct `apply` caller must keep values under 64 MiB or use an LSM backend |
| 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` | 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 |
| Multipart upload fails mid-stream | The upload is aborted, so S3/GCS drop its parts; if the abort itself fails, a warning names the key | Next round re-uploads; a bucket lifecycle rule for incomplete uploads is the backstop (a crashed process can't abort) |
| 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 |
| Tampered / torn artifact | `ArtifactInvalid` at import (hash + cursor gate); nothing written to destination | Fetch another artifact |
Expand Down
2 changes: 1 addition & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ rust-version = "1.92"

[package]
name = "beyond-slipstream"
version = "0.8.2"
version = "0.8.3"
edition.workspace = true
license.workspace = true
rust-version.workspace = true
Expand Down
26 changes: 23 additions & 3 deletions src/nats.rs
Original file line number Diff line number Diff line change
Expand Up @@ -298,9 +298,16 @@ fn classify_raw_create_response(payload: &[u8]) -> RawCreateOutcome {
return RawCreateOutcome::AlreadyExists;
}

// 400 "maximum number of streams reached" may also mean bucket exists
// (Synadia Cloud returns this when at stream limit but bucket exists)
if code == 400 && description.contains("maximum number of streams") {
// At the stream limit, the bucket may still exist (Synadia Cloud answers
// this way for a bucket it already has). 10027 is the server's
// JSMaximumStreamsLimitErr; the text is the fallback for replies without
// it. A false positive is harmless: StreamLimit re-checks the bucket.
if err_code == 10027
|| (code == 400
&& description
.to_ascii_lowercase()
.contains("maximum number of streams"))
{
return RawCreateOutcome::StreamLimit;
}

Expand Down Expand Up @@ -1854,6 +1861,19 @@ mod tests {
);
}

#[test]
fn raw_create_stream_limit_by_code_or_reworded_text() {
for payload in [
&br#"{"error":{"code":400,"err_code":10027,"description":"stream limit hit"}}"#[..],
br#"{"error":{"code":400,"description":"Maximum Number Of Streams reached"}}"#,
] {
assert_eq!(
classify_raw_create_response(payload),
RawCreateOutcome::StreamLimit
);
}
}

#[test]
fn raw_create_propagates_unknown_error() {
// Any other JetStream error is fatal and must surface code + description.
Expand Down
4 changes: 4 additions & 0 deletions src/protocol.rs
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,10 @@ pub fn payload_prunable(
/// resync. Machine-checked as `bootstrap never silently diverges` in
/// `tests/model.rs` (where the model's retention floor is
/// `first_sequence - 1`).
///
/// This answers only whether the window is intact; it holds at
/// `revision == u64::MAX`. The sequence to start at must come from
/// [`resume_start_sequence`], which refuses to wrap there.
pub fn resume_window_ok(revision: u64, first_sequence: u64) -> bool {
first_sequence <= revision.saturating_add(1)
}
Expand Down
55 changes: 42 additions & 13 deletions src/snapshot.rs
Original file line number Diff line number Diff line change
Expand Up @@ -526,8 +526,8 @@ impl AppendLogSnapshot {
/// manifest — checksums, backend identity, format generation — and the
/// staged copy is loaded and its cursor compared against the manifest's
/// before anything lands at `dest_path`; a bad artifact never becomes a
/// fold. A crash mid-import leaves nothing at `dest_path`; a crash after
/// the final rename leaves a valid fold (a retried import then refuses the
/// fold. A crash mid-import leaves nothing at `dest_path`; a crash or a
/// failed open after the final rename leaves a valid fold (a retried import then refuses the
/// existing destination — just [`open`](Self::open) it).
pub fn import(
artifact_dir: &Path,
Expand Down Expand Up @@ -857,10 +857,8 @@ fn parse_put(data: &[u8], stored_crc: u32) -> Result<(Record<'_>, usize), Record
data[vl_off + 3],
]) as usize;
// Lengths are read before CRC. A bit-flipped value_len of 2e9 looks like
// EOF and used to silent-drop every later record. Cap far above any
// legitimate config value; a torn last record of a real-sized value is
// still Truncated below.
const MAX_VALUE_LEN: usize = 16 * 1024 * 1024;
// EOF and used to silent-drop every later record. A torn last record of a
// real-sized value is still Truncated below.
if value_len > MAX_VALUE_LEN {
return Err(RecordError::CrcMismatch {
consumed: data.len().min(vl_off + 4),
Expand Down Expand Up @@ -984,6 +982,11 @@ fn parse_cursor(data: &[u8], stored_crc: u32) -> Result<(Record<'_>, usize), Rec
// Internal: record writing (incremental CRC, no allocations)
// ---------------------------------------------------------------------------

/// Largest value a put record may carry. NATS caps a message payload at
/// 64 MiB, so every value a watch can deliver fits. The same bound gates the
/// writer and the reader.
const MAX_VALUE_LEN: usize = 64 * 1024 * 1024;

fn write_put_record(
w: &mut impl Write,
key: &str,
Expand All @@ -1004,13 +1007,15 @@ fn write_put_record(
u16::MAX
))
})?;
let value_len = u32::try_from(value.len()).map_err(|_| {
SnapshotError::InvalidFormat(format!(
"value too long: {} bytes (max {})",
value.len(),
u32::MAX
))
})?;
// The reader treats a longer value as corruption, so writing one would
// leave a record no later `load` can get past.
if value.len() > MAX_VALUE_LEN {
return Err(SnapshotError::InvalidFormat(format!(
"value too long: {} bytes (max {MAX_VALUE_LEN})",
value.len()
)));
}
let value_len = value.len() as u32;
// The version is stored as length-prefixed raw bytes so any backend's token
// (NATS u64, FDB 10-byte versionstamp) round-trips intact. `VersionToken`
// caps inline storage at 10 bytes, so this `u8` length never truncates today;
Expand Down Expand Up @@ -1244,6 +1249,30 @@ mod tests {
assert_eq!(snap.entries["node.eu-west-1"].value, b"val2");
}

/// The writer used to accept any u32 length while the reader called
/// anything over 16 MiB corruption: one large value bricked the fold.
#[test]
fn largest_value_round_trips_and_larger_is_refused_before_writing() {
let dir = TempDir::new().unwrap();
let path = dir.path().join("test.snap");

let (_, mut snap) = AppendLogSnapshot::load(&path).unwrap();
let big = vec![7u8; MAX_VALUE_LEN];
snap.apply(&[put("big", &big, 1)], &cursor(1)).unwrap();
let err = snap
.apply(&[put("huge", &vec![0u8; MAX_VALUE_LEN + 1], 2)], &cursor(2))
.unwrap_err();
assert!(matches!(err, SnapshotError::InvalidFormat(_)), "{err}");
snap.apply(&[put("after", b"x", 3)], &cursor(3)).unwrap();
drop(snap);

let (c, snap) = AppendLogSnapshot::load(&path).unwrap();
assert_eq!(c.as_u64(), Some(3));
assert_eq!(snap.get("big").unwrap().unwrap().value.len(), MAX_VALUE_LEN);
assert!(snap.get("huge").unwrap().is_none());
assert_eq!(snap.get("after").unwrap().unwrap().value, b"x");
}

#[test]
fn multiple_batches() {
let dir = TempDir::new().unwrap();
Expand Down
2 changes: 1 addition & 1 deletion src/snapshot_rocksdb.rs
Original file line number Diff line number Diff line change
Expand Up @@ -423,7 +423,7 @@ impl RocksDbSnapshot {
let (staged_cursor, verify) = Self::open(
&stage.payload(),
RocksDbConfig {
sync: config.sync,
sync: true,
cache_size_bytes: 0,
},
)?;
Expand Down
59 changes: 39 additions & 20 deletions src/transport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -478,28 +478,47 @@ impl ArtifactTransport for ObjectStoreTransport {
.await?
.map_err(map_obj)?;
let mut wm = WriteMultipart::new_with_chunk_size(upload, CHUNK);
let mut buf = vec![0u8; CHUNK];
loop {
use tokio::io::AsyncReadExt;
let n = file.read(&mut buf).await.map_err(SnapshotError::Io)?;
if n == 0 {
break;
let streamed = async {
let mut buf = vec![0u8; CHUNK];
loop {
use tokio::io::AsyncReadExt;
let n = file.read(&mut buf).await.map_err(SnapshotError::Io)?;
if n == 0 {
break;
}
timed("part upload", wm.wait_for_capacity(MAX_CONCURRENT_PARTS))
.await?
.map_err(map_obj)?;
wm.write(&buf[..n]);
}
timed("part upload", wm.wait_for_capacity(MAX_CONCURRENT_PARTS))
.await?
.map_err(map_obj)?;
wm.write(&buf[..n]);
// Drains every in-flight part (up to MAX_CONCURRENT_PARTS × CHUNK
// bytes), so it gets a proportionally larger stall bound than a
// single-request await.
timed_by(
"multipart drain",
OP_TIMEOUT * MAX_CONCURRENT_PARTS as u32,
wm.wait_for_capacity(0),
)
.await?
.map_err(map_obj)
}
// finish() drains every in-flight part (up to MAX_CONCURRENT_PARTS ×
// CHUNK bytes) plus the completion request, so it gets a proportionally
// larger stall bound than a single-request await.
timed_by(
"multipart finish",
OP_TIMEOUT * MAX_CONCURRENT_PARTS as u32,
wm.finish(),
)
.await?
.map_err(map_obj)?;
.await;
if let Err(e) = streamed {
// S3 and GCS keep uploaded parts of an unfinished upload (billed,
// invisible to listings) until it is aborted; dropping does not.
if let Err(abort) = timed("multipart abort", wm.abort())
.await
.and_then(|r| r.map_err(map_obj))
{
warn!(key, error = %abort, "multipart abort failed; parts linger until a bucket lifecycle rule removes them");
}
return Err(e);
}
// Every part is in; finish() only flushes the (empty) buffer and sends
// the completion request, aborting the upload itself if that fails.
timed("multipart finish", wm.finish())
.await?
.map_err(map_obj)?;

// Pointer LAST, by monotonic conditional swap: its presence marks the
// payload complete, and it can never regress past a newer round.
Expand Down
Loading