From e7298c698fbb67ed0eaa17ffe411e7965c543f4c Mon Sep 17 00:00:00 2001 From: Torstein Tauno Svendsen Date: Fri, 2 Oct 2026 22:13:19 +0200 Subject: [PATCH 1/2] The grain is extended and rebased without reading it whole extend_grain and rebase_head read the entire .grain on every call just to learn how many records it held and where the last one ended; append_page read and rewrote it per received page. At a few hundred MB that was a transient allocation of the same size on every chunk flush. A .grain.commit beside the grain now records both numbers after each completed write, bound to the grain's inode and the bytes at its two ends, so a grain replaced or rebased by anyone else is walked once instead of trusted. Whatever lies past the committed end is a torn write and is cut before the next append. The commit is renamed into place without an fsync: a lost or older one describes a prefix of the file and costs re-deriving the difference. extend_grain: seek to the committed end and append. rebase_head: walk only the k dropped records, carry the commit through both the COLLAPSE_RANGE and the rewrite path. append_page: append in place. .rings is still read whole by extend_grain, and grain::load still copies every filter; neither is changed here. Co-Authored-By: Claude Sonnet 5.5 --- src/format.rs | 8 + src/grain.rs | 620 +++++++++++++++++++++++++++++++++++++++++++++----- src/rotate.rs | 1 + src/store.rs | 5 + 4 files changed, 572 insertions(+), 62 deletions(-) diff --git a/src/format.rs b/src/format.rs index ee866f4..0361f5e 100644 --- a/src/format.rs +++ b/src/format.rs @@ -52,6 +52,9 @@ pub enum RingsVersion { pub const TRUNK_EXT: &str = "trunk"; pub const RINGS_EXT: &str = "rings"; pub const GRAIN_EXT: &str = "grain"; +/// Where the grain's records end and how many there are, as of its last +/// completed write (grain.rs). A hint bound to the file it describes. +pub const GRAIN_COMMIT_EXT: &str = "grain.commit"; pub const BARK_EXT: &str = "bark"; /// The collapse-head seqlock counter (store.rs): even means idle, odd /// means a collapse is in flight. Missing reads as 0 (never collapsed). @@ -80,6 +83,7 @@ pub fn every_path(dir: &Path, name: &str) -> Vec { RINGS_EXT, TRUNK_EXT, GRAIN_EXT, + GRAIN_COMMIT_EXT, BARK_EXT, SEQ_EXT, TRIM_EXT, @@ -104,6 +108,10 @@ pub fn grain_path(dir: &Path, name: &str) -> PathBuf { dir.join(format!("{name}.{GRAIN_EXT}")) } +pub fn grain_commit_path(dir: &Path, name: &str) -> PathBuf { + dir.join(format!("{name}.{GRAIN_COMMIT_EXT}")) +} + pub fn bark_path(dir: &Path, name: &str) -> PathBuf { dir.join(format!("{name}.{BARK_EXT}")) } diff --git a/src/grain.rs b/src/grain.rs index 2debcbd..5ef41e8 100644 --- a/src/grain.rs +++ b/src/grain.rs @@ -24,10 +24,19 @@ //! A record carries no chunk id: its POSITION is the chunk index. That is //! what makes appending a chunk cost one appended record, and what makes a //! retention head-drop hostile — see `rebase_head`. +//! +//! The file states neither how many records it holds nor where the last +//! one ends, so every write would have to read it whole to find out. The +//! `.grain.commit` beside it records both after each completed write, +//! bound to the grain's inode and the bytes at its two ends: a grain +//! replaced or rebased by anyone else no longer matches and is walked +//! once instead. Whatever lies past the committed end is a torn write and +//! is cut before the next append. use std::collections::HashSet; use std::fs::{self, File, OpenOptions}; -use std::os::unix::fs::FileExt; +use std::io::{self, BufReader, BufWriter, Read, Seek, SeekFrom, Write}; +use std::os::unix::fs::{FileExt, MetadataExt}; use std::path::Path; use anyhow::{bail, Context}; @@ -123,34 +132,142 @@ fn header_bytes(first_rec: usize) -> [u8; HEADER_LEN] { /// understand — in which case the caller scans (readers) or rebuilds /// (writers), never guesses. fn first_record_offset(buf: &[u8]) -> Option { - if buf.len() < HEADER_LEN { + header_first_rec(buf, buf.len() as u64).map(|o| o as usize) +} + +fn header_first_rec(h: &[u8], file_len: u64) -> Option { + if h.len() < HEADER_LEN { return None; } - match &buf[..8] { - m if m == GRAIN_MAGIC => Some(HEADER_LEN), + match &h[..8] { + m if m == GRAIN_MAGIC => Some(HEADER_LEN as u64), m if m == GRAIN_MAGIC_V2 => { - let off = u32::from_le_bytes(buf[12..16].try_into().unwrap()) as usize; - (HEADER_LEN..=buf.len()).contains(&off).then_some(off) + let off = u32::from_le_bytes(h[12..16].try_into().unwrap()) as u64; + (HEADER_LEN as u64..=file_len).contains(&off).then_some(off) } _ => None, } } -/// Walk `n` records from `off`, returning where the next one starts. -/// None when the file ends first — a partial tail (crash debris) or a -/// grain that simply doesn't reach that far. -fn skip_records(buf: &[u8], mut off: usize, n: usize) -> Option { - for _ in 0..n { - if off + 4 > buf.len() { - return None; - } - let len = u32::from_le_bytes(buf[off..off + 4].try_into().unwrap()) as usize; - if off + 4 + len > buf.len() { - return None; +const LEAD_LEN: usize = 16; +const WALK_BUF: usize = 256 * 1024; + +/// What a commit is bound to: the file itself, its header, and the bytes +/// where its first record starts. A grain replaced or head-dropped by +/// anyone else no longer matches, whatever its length. +struct Stamp { + ino: u64, + header: [u8; HEADER_LEN], + lead: [u8; LEAD_LEN], + first: u64, + len: u64, +} + +fn stamp(f: &File) -> io::Result> { + let md = f.metadata()?; + let len = md.len(); + if len < HEADER_LEN as u64 { + return Ok(None); + } + let mut header = [0u8; HEADER_LEN]; + f.read_exact_at(&mut header, 0)?; + let Some(first) = header_first_rec(&header, len) else { + return Ok(None); + }; + let mut lead = [0u8; LEAD_LEN]; + let n = (len - first).min(LEAD_LEN as u64) as usize; + f.read_exact_at(&mut lead[..n], first)?; + Ok(Some(Stamp { + ino: md.ino(), + header, + lead, + first, + len, + })) +} + +/// How many whole records the grain holds and where the last one ends. +#[derive(Clone, Copy, Debug, PartialEq)] +struct Tail { + count: u64, + end: u64, +} + +/// Walk up to `max` length-prefixed records from `first`, stopping at the +/// last whole one: a partial tail is where a crash left off. +fn walk(f: &File, first: u64, len: u64, max: u64) -> io::Result { + let mut r = BufReader::with_capacity(WALK_BUF, f); + r.seek(SeekFrom::Start(first))?; + let mut tail = Tail { + count: 0, + end: first, + }; + while tail.count < max && tail.end + 4 <= len { + let mut p = [0u8; 4]; + r.read_exact(&mut p)?; + let body = u32::from_le_bytes(p) as u64; + if tail.end + 4 + body > len { + break; } - off += 4 + len; + r.seek_relative(body as i64)?; + tail.end += 4 + body; + tail.count += 1; + } + Ok(tail) +} + +fn hex(b: &[u8]) -> String { + b.iter().map(|x| format!("{x:02x}")).collect() +} + +fn unhex(s: &str) -> Option<[u8; N]> { + if s.len() != N * 2 { + return None; + } + let mut out = [0u8; N]; + for (i, o) in out.iter_mut().enumerate() { + *o = u8::from_str_radix(s.get(i * 2..i * 2 + 2)?, 16).ok()?; } - Some(off) + Some(out) +} + +/// The recorded tail, if it still describes THIS file. A commit that is +/// stale because it lags the file is safe to trust (the caller truncates +/// to it and re-derives the rest); one bound to a different file is not. +fn load_commit(dir: &Path, name: &str, s: &Stamp) -> Option { + let text = fs::read_to_string(format::grain_commit_path(dir, name)).ok()?; + let mut it = text.split_whitespace(); + let ino: u64 = it.next()?.parse().ok()?; + let header = unhex::(it.next()?)?; + let lead = unhex::(it.next()?)?; + let count: u64 = it.next()?.parse().ok()?; + let end: u64 = it.next()?.parse().ok()?; + (ino == s.ino && header == s.header && lead == s.lead && end >= s.first && end <= s.len) + .then_some(Tail { count, end }) +} + +/// Atomic by rename and deliberately not fsynced: a lost or older commit +/// describes a prefix of the file, which costs re-deriving the difference +/// and never correctness. +fn save_commit(dir: &Path, name: &str, f: &File, t: Tail) -> io::Result<()> { + let path = format::grain_commit_path(dir, name); + let Some(s) = stamp(f)? else { + let _ = fs::remove_file(&path); + return Ok(()); + }; + let tmp = dir.join(format!("{name}.{}.tmp", format::GRAIN_COMMIT_EXT)); + fs::write( + &tmp, + format!( + "{} {} {} {} {}\n", + s.ino, + hex(&s.header), + hex(&s.lead), + t.count, + t.end + ), + )?; + fs::rename(&tmp, &path) } fn build_filter(tokens: &HashSet<&[u8]>) -> Vec { @@ -256,38 +373,29 @@ pub fn cmd_reindex(file: &Path) -> anyhow::Result<()> { /// caller holds the writer locks. pub fn extend_grain(dir: &Path, name: &str) -> anyhow::Result<()> { let gpath = format::grain_path(dir, name); - let existing = match fs::read(&gpath) { - Ok(b) => b, - Err(_) => return build_grain(dir, name), + let Ok(out) = OpenOptions::new().read(true).write(true).open(&gpath) else { + return build_grain(dir, name); }; - let Some(first) = first_record_offset(&existing) else { + let Some((tail, known)) = locate(dir, name, &out)? else { return build_grain(dir, name); }; - let mut off = first; - let mut covered = 0usize; - loop { - if off + 4 > existing.len() { - break; - } - let len = u32::from_le_bytes(existing[off..off + 4].try_into().unwrap()) as usize; - if off + 4 + len > existing.len() { - break; // partial tail (crash): overwritten below - } - off += 4 + len; - covered += 1; - } let records = format::read_index(&format::rings_path(dir, name))?; + let covered = tail.count as usize; if covered > records.len() { return build_grain(dir, name); } + if out.metadata()?.len() > tail.end { + out.set_len(tail.end)?; + } if covered == records.len() { + if !known { + save_commit(dir, name, &out, tail)?; + } return Ok(()); } let trunk = File::open(format::trunk_path(dir, name)) .with_context(|| format!("opening {}", format::trunk_path(dir, name).display()))?; - let out = OpenOptions::new().write(true).open(&gpath)?; - out.set_len(off as u64)?; - let mut woff = off as u64; + let mut woff = tail.end; for c in &records[covered..] { let mut comp = vec![0u8; c.comp_len as usize]; trunk.read_exact_at(&mut comp, c.comp_start)?; @@ -301,9 +409,31 @@ pub fn extend_grain(dir: &Path, name: &str) -> anyhow::Result<()> { woff += filter.len() as u64; } out.sync_all()?; + save_commit( + dir, + name, + &out, + Tail { + count: records.len() as u64, + end: woff, + }, + )?; Ok(()) } +/// Where the records end: the commit when it still describes this file, +/// otherwise one walk. The flag says which. None when the file is not a +/// grain. +fn locate(dir: &Path, name: &str, f: &File) -> io::Result> { + let Some(s) = stamp(f)? else { + return Ok(None); + }; + Ok(Some(match load_commit(dir, name, &s) { + Some(t) => (t, true), + None => (walk(f, s.first, s.len, u64::MAX)?, false), + })) +} + /// Does a grain header from elsewhere describe the tokenizer THIS build /// uses? Parameters live in bytes 8..12 (case folding, MIN_TOKEN, /// MAX_TOKEN, K) and a page built under different ones, read under ours, @@ -325,16 +455,53 @@ pub fn header_matches(bytes: &[u8]) -> bool { /// extended is removed and the next `extend_grain` rebuilds it. pub fn append_page(dir: &Path, name: &str, page: &[u8]) -> anyhow::Result<()> { let gpath = format::grain_path(dir, name); - let mut buf = match fs::read(&gpath) { - Ok(b) if first_record_offset(&b).is_some() => b, - _ => header_bytes(HEADER_LEN).to_vec(), + let open = || OpenOptions::new().read(true).write(true).open(&gpath); + let (f, tail) = match open() { + Ok(f) => match locate(dir, name, &f)? { + Some((tail, _)) => (f, tail), + None => (fresh_grain(dir, name, &open)?, EMPTY), + }, + Err(e) if e.kind() == io::ErrorKind::NotFound => (fresh_grain(dir, name, &open)?, EMPTY), + Err(e) => return Err(e).with_context(|| format!("opening {}", gpath.display())), }; - buf.extend_from_slice(&(page.len() as u32).to_le_bytes()); - buf.extend_from_slice(page); + let mut record = Vec::with_capacity(4 + page.len()); + record.extend_from_slice(&(page.len() as u32).to_le_bytes()); + record.extend_from_slice(page); + if f.metadata()?.len() > tail.end { + f.set_len(tail.end)?; + } + f.write_all_at(&record, tail.end) + .with_context(|| format!("appending to {}", gpath.display()))?; + save_commit( + dir, + name, + &f, + Tail { + count: tail.count + 1, + end: tail.end + record.len() as u64, + }, + )?; + Ok(()) +} + +const EMPTY: Tail = Tail { + count: 0, + end: HEADER_LEN as u64, +}; + +/// A header-only grain installed by rename, so a reader sees the old file +/// or the new one and never a file with half a header. +fn fresh_grain( + dir: &Path, + name: &str, + open: &dyn Fn() -> io::Result, +) -> anyhow::Result { + let gpath = format::grain_path(dir, name); let tmp = gpath.with_extension("grain.tmp"); - fs::write(&tmp, &buf).with_context(|| format!("writing {}", tmp.display()))?; + fs::write(&tmp, header_bytes(HEADER_LEN)) + .with_context(|| format!("writing {}", tmp.display()))?; fs::rename(&tmp, &gpath).with_context(|| format!("renaming onto {}", gpath.display()))?; - Ok(()) + Ok(open()?) } /// Drop the first `k` chunks' filters, after retention has cut the same @@ -354,29 +521,38 @@ pub fn append_page(dir: &Path, name: &str, page: &[u8]) -> anyhow::Result<()> { /// keeps a seqlock window open around both. Best-effort by contract: a /// grain that is missing, unreadable or shorter than the drop is simply /// removed, and the next extend rebuilds it. -pub fn rebase_head(dir: &Path, name: &str, k: usize) -> std::io::Result<()> { +pub fn rebase_head(dir: &Path, name: &str, k: usize) -> io::Result<()> { let gpath = format::grain_path(dir, name); if k == 0 { return Ok(()); } - let buf = match fs::read(&gpath) { - Ok(b) => b, - Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()), + let f = match OpenOptions::new().read(true).write(true).open(&gpath) { + Ok(f) => f, + Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(()), Err(e) => return Err(e), }; - let drop_to = first_record_offset(&buf).and_then(|first| skip_records(&buf, first, k)); - let Some(survivors) = drop_to else { + let k64 = k as u64; + let s = stamp(&f)?; + let survivors = match &s { + Some(s) => { + let t = walk(&f, s.first, s.len, k64)?; + (t.count == k64).then_some(t.end) + } + None => None, + }; + let (Some(s), Some(survivors)) = (s, survivors) else { // Not a grain we understand, or it covered fewer chunks than were // dropped: nothing left worth rebasing. let _ = fs::remove_file(&gpath); + let _ = fs::remove_file(format::grain_commit_path(dir, name)); return Ok(()); }; + let known = load_commit(dir, name, &s).filter(|t| t.count >= k64 && survivors <= t.end); - let f = OpenOptions::new().read(true).write(true).open(&gpath)?; // Keep room to re-stamp the header over dead bytes: the cut can never // reach past `survivors - HEADER_LEN`. let bsize = crate::store::fstatvfs_bsize(&f)?; - let aligned = ((survivors - HEADER_LEN) as u64 / bsize) * bsize; + let aligned = ((survivors - HEADER_LEN as u64) / bsize) * bsize; if aligned > 0 { let rc = unsafe { libc::fallocate( @@ -387,11 +563,21 @@ pub fn rebase_head(dir: &Path, name: &str, k: usize) -> std::io::Result<()> { ) }; if rc == 0 { - let first_rec = survivors - aligned as usize; + let first_rec = (survivors - aligned) as usize; f.write_all_at(&header_bytes(first_rec), 0)?; - return f.sync_all(); + f.sync_all()?; + recommit( + dir, + name, + &f, + known.map(|t| Tail { + count: t.count - k64, + end: t.end - aligned, + }), + ); + return Ok(()); } - let e = std::io::Error::last_os_error(); + let e = io::Error::last_os_error(); match e.raw_os_error() { // No COLLAPSE_RANGE here (tmpfs, btrfs, NFS, older ext4/xfs): // rewrite instead, below. @@ -403,17 +589,46 @@ pub fn rebase_head(dir: &Path, name: &str, k: usize) -> std::io::Result<()> { // renamed, so a reader sees the whole old file or the whole new one. // A store on a filesystem without COLLAPSE_RANGE therefore never // upgrades its magic at all. - drop(f); let tmp = dir.join(format!("{name}.{}.tmp", format::GRAIN_EXT)); let out = OpenOptions::new() + .read(true) .write(true) .create(true) .truncate(true) .open(&tmp)?; - out.write_all_at(&header_bytes(HEADER_LEN), 0)?; - out.write_all_at(&buf[survivors..], HEADER_LEN as u64)?; + let src_end = known.map_or(s.len, |t| t.end); + { + let mut w = BufWriter::with_capacity(WALK_BUF, &out); + w.write_all(&header_bytes(HEADER_LEN))?; + let mut r = BufReader::with_capacity(WALK_BUF, &f); + r.seek(SeekFrom::Start(survivors))?; + io::copy(&mut r.take(src_end - survivors), &mut w)?; + w.flush()?; + } out.sync_all()?; - fs::rename(&tmp, &gpath) + fs::rename(&tmp, &gpath)?; + recommit( + dir, + name, + &out, + known.map(|t| Tail { + count: t.count - k64, + end: HEADER_LEN as u64 + (t.end - survivors), + }), + ); + Ok(()) +} + +/// Carry a commit across a rebase, or drop it so the next extend walks. +fn recommit(dir: &Path, name: &str, f: &File, t: Option) { + match t { + Some(t) => { + let _ = save_commit(dir, name, f, t); + } + None => { + let _ = fs::remove_file(format::grain_commit_path(dir, name)); + } + } } /// The grain build itself; the caller holds the writer locks. @@ -424,6 +639,7 @@ pub fn build_grain(dir: &Path, name: &str) -> anyhow::Result<()> { .with_context(|| format!("opening {}", format::trunk_path(dir, name).display()))?; let tmp = dir.join(format!("{name}.{}.tmp", format::GRAIN_EXT)); let out = OpenOptions::new() + .read(true) .write(true) .create(true) .truncate(true) @@ -462,6 +678,15 @@ pub fn build_grain(dir: &Path, name: &str) -> anyhow::Result<()> { format::grain_path(dir, name).display() ) })?; + let _ = save_commit( + dir, + name, + &out, + Tail { + count: records.len() as u64, + end: off, + }, + ); crate::note!( "timberfs: indexed {} chunk(s), {} distinct tokens ({} avg/chunk), grain is {} bytes \ ({} bytes/chunk avg)", @@ -588,7 +813,9 @@ mod tests { } else { assert_eq!(first, HEADER_LEN, "the rewrite path resets to V1"); } - assert_eq!(skip_records(&buf, first, 6), Some(buf.len())); + let g = fs::File::open(format::grain_path(d.path(), "a.log")).unwrap(); + let t = walk(&g, first as u64, buf.len() as u64, u64::MAX).unwrap(); + assert_eq!((t.count, t.end), (6, buf.len() as u64)); } #[test] @@ -608,4 +835,273 @@ mod tests { rebase_head(d.path(), "missing.log", 3).unwrap(); assert!(!format::grain_path(d.path(), "missing.log").exists()); } + + fn cfg() -> crate::store::Config { + crate::store::Config { + chunk_size: 1 << 20, + level: 1, + flush_age_ms: u64::MAX, + } + } + + fn a_store(dir: &Path, chunks: usize) -> crate::store::Store { + let mut st = crate::store::Store { + dir: dir.to_path_buf(), + cfg: cfg(), + files: std::collections::BTreeMap::new(), + }; + st.create("a.log").unwrap(); + add_chunks(&mut st, 0, chunks); + st + } + + fn add_chunks(st: &mut crate::store::Store, from: usize, to: usize) { + let f = st.files.get_mut("a.log").unwrap(); + for i in from..to { + f.append_windowed( + format!("line {i} marker{i:05} padding\n").as_bytes(), + 1_000 + i as u64, + 1_000 + i as u64, + &cfg(), + ) + .unwrap(); + f.flush_chunk(&cfg()).unwrap(); + } + } + + fn marker(i: usize) -> Vec> { + vec![format!("marker{i:05}").into_bytes()] + } + + fn grain_bytes(d: &Path) -> Vec { + fs::read(format::grain_path(d, "a.log")).unwrap() + } + + fn committed(d: &Path) -> Option { + let f = File::open(format::grain_path(d, "a.log")).ok()?; + load_commit(d, "a.log", &stamp(&f).ok()??) + } + + fn assert_indexes_every_chunk(d: &Path, chunks: usize) { + let g = load(&format::grain_path(d, "a.log")).unwrap(); + assert_eq!(g.chunk_count(), chunks); + for i in 0..chunks { + assert!(g.may_contain_all(i, &marker(i)), "chunk {i} lost its token"); + } + } + + #[test] + fn extending_in_steps_equals_building_at_once() { + let d = TempDir::new(); + let mut st = a_store(d.path(), 3); + extend_grain(d.path(), "a.log").unwrap(); + assert_eq!(committed(d.path()).map(|t| t.count), Some(3)); + add_chunks(&mut st, 3, 8); + extend_grain(d.path(), "a.log").unwrap(); + let stepped = grain_bytes(d.path()); + + build_grain(d.path(), "a.log").unwrap(); + assert_eq!(stepped, grain_bytes(d.path())); + assert_indexes_every_chunk(d.path(), 8); + let t = committed(d.path()).expect("a build leaves a commit"); + assert_eq!((t.count, t.end), (8, stepped.len() as u64)); + } + + #[test] + fn the_commit_governs_where_extending_resumes() { + // A complete-looking record past the committed end is what a walk + // would count as covered; the commit says it is not. + let d = TempDir::new(); + let mut st = a_store(d.path(), 3); + extend_grain(d.path(), "a.log").unwrap(); + let clean = grain_bytes(d.path()); + let g = OpenOptions::new() + .append(true) + .open(format::grain_path(d.path(), "a.log")) + .unwrap(); + (&g).write_all(&8u32.to_le_bytes()).unwrap(); + (&g).write_all(&[0u8; 8]).unwrap(); + drop(g); + + add_chunks(&mut st, 3, 4); + extend_grain(d.path(), "a.log").unwrap(); + assert_indexes_every_chunk(d.path(), 4); + assert_eq!(&grain_bytes(d.path())[..clean.len()], &clean[..]); + } + + #[test] + fn a_torn_tail_is_cut_before_extending() { + let d = TempDir::new(); + let mut st = a_store(d.path(), 3); + extend_grain(d.path(), "a.log").unwrap(); + let g = OpenOptions::new() + .append(true) + .open(format::grain_path(d.path(), "a.log")) + .unwrap(); + (&g).write_all(&100u32.to_le_bytes()).unwrap(); + (&g).write_all(&[1, 2, 3]).unwrap(); + drop(g); + // Without the commit, the walk has to find the same boundary. + fs::remove_file(format::grain_commit_path(d.path(), "a.log")).unwrap(); + + add_chunks(&mut st, 3, 5); + extend_grain(d.path(), "a.log").unwrap(); + assert_indexes_every_chunk(d.path(), 5); + let stepped = grain_bytes(d.path()); + build_grain(d.path(), "a.log").unwrap(); + assert_eq!(stepped, grain_bytes(d.path())); + } + + #[test] + fn a_commit_for_another_file_is_not_believed() { + let d = TempDir::new(); + let _st = a_store(d.path(), 3); + extend_grain(d.path(), "a.log").unwrap(); + assert_eq!(committed(d.path()).map(|t| t.count), Some(3)); + + // Someone else replaces the grain with one covering fewer chunks. + let other = TempDir::new(); + let _ = a_store(other.path(), 2); + build_grain(other.path(), "a.log").unwrap(); + fs::remove_file(format::grain_commit_path(other.path(), "a.log")).unwrap(); + fs::copy( + format::grain_path(other.path(), "a.log"), + d.path().join("swap"), + ) + .unwrap(); + fs::rename(d.path().join("swap"), format::grain_path(d.path(), "a.log")).unwrap(); + + extend_grain(d.path(), "a.log").unwrap(); + assert_indexes_every_chunk(d.path(), 3); + } + + #[test] + fn a_lagging_commit_re_derives_what_it_does_not_cover() { + // The commit rename can be lost across a power cut while the + // records it described survive. + let d = TempDir::new(); + let mut st = a_store(d.path(), 3); + extend_grain(d.path(), "a.log").unwrap(); + let old = fs::read(format::grain_commit_path(d.path(), "a.log")).unwrap(); + add_chunks(&mut st, 3, 6); + extend_grain(d.path(), "a.log").unwrap(); + fs::write(format::grain_commit_path(d.path(), "a.log"), old).unwrap(); + + extend_grain(d.path(), "a.log").unwrap(); + assert_indexes_every_chunk(d.path(), 6); + } + + #[test] + fn rebasing_carries_the_commit() { + let d = TempDir::new(); + let _st = a_store(d.path(), 12); + extend_grain(d.path(), "a.log").unwrap(); + + rebase_head(d.path(), "a.log", 5).unwrap(); + + let t = committed(d.path()).expect("the rebase kept the commit usable"); + assert_eq!(t.count, 7); + assert_eq!(t.end, grain_bytes(d.path()).len() as u64); + let g = load(&format::grain_path(d.path(), "a.log")).unwrap(); + assert_eq!(g.chunk_count(), 7); + for i in 5..12 { + assert!(g.may_contain_all(i - 5, &marker(i))); + } + } + + #[test] + fn rebasing_without_a_commit_still_works_and_leaves_none() { + let d = TempDir::new(); + let _st = a_store(d.path(), 12); + extend_grain(d.path(), "a.log").unwrap(); + fs::remove_file(format::grain_commit_path(d.path(), "a.log")).unwrap(); + + rebase_head(d.path(), "a.log", 4).unwrap(); + assert!(committed(d.path()).is_none()); + assert_eq!( + load(&format::grain_path(d.path(), "a.log")) + .unwrap() + .chunk_count(), + 8 + ); + } + + #[test] + fn adopted_pages_build_the_same_grain_and_commit() { + let src = TempDir::new(); + let _st = a_store(src.path(), 6); + build_grain(src.path(), "a.log").unwrap(); + let g = load(&format::grain_path(src.path(), "a.log")).unwrap(); + + let dst = TempDir::new(); + for i in 0..6 { + append_page(dst.path(), "a.log", g.page(i).unwrap()).unwrap(); + } + assert_eq!(grain_bytes(dst.path()), grain_bytes(src.path())); + let t = committed(dst.path()).expect("adoption commits"); + assert_eq!((t.count, t.end), (6, grain_bytes(dst.path()).len() as u64)); + } + + #[test] + fn adopting_a_page_after_a_torn_tail_does_not_keep_the_debris() { + let src = TempDir::new(); + let _st = a_store(src.path(), 3); + build_grain(src.path(), "a.log").unwrap(); + let g = load(&format::grain_path(src.path(), "a.log")).unwrap(); + + let dst = TempDir::new(); + append_page(dst.path(), "a.log", g.page(0).unwrap()).unwrap(); + let f = OpenOptions::new() + .append(true) + .open(format::grain_path(dst.path(), "a.log")) + .unwrap(); + (&f).write_all(&50u32.to_le_bytes()).unwrap(); + drop(f); + append_page(dst.path(), "a.log", g.page(1).unwrap()).unwrap(); + append_page(dst.path(), "a.log", g.page(2).unwrap()).unwrap(); + + assert_eq!(grain_bytes(dst.path()), grain_bytes(src.path())); + } + + #[test] + fn rebasing_a_wide_grain_carries_the_commit_through_a_collapse() { + // Filters big enough that the dropped prefix spans a filesystem + // block, so COLLAPSE_RANGE is the branch taken wherever it exists. + let d = TempDir::new(); + let mut st = crate::store::Store { + dir: d.path().to_path_buf(), + cfg: cfg(), + files: std::collections::BTreeMap::new(), + }; + st.create("a.log").unwrap(); + let f = st.files.get_mut("a.log").unwrap(); + for i in 0..10usize { + let wide: String = (0..1500).map(|t| format!("w{i:02}t{t:05} ")).collect(); + f.append_windowed( + format!("{wide}marker{i:05}\n").as_bytes(), + 1_000 + i as u64, + 1_000 + i as u64, + &cfg(), + ) + .unwrap(); + f.flush_chunk(&cfg()).unwrap(); + } + extend_grain(d.path(), "a.log").unwrap(); + + rebase_head(d.path(), "a.log", 6).unwrap(); + + let buf = grain_bytes(d.path()); + eprintln!( + "rebased grain magic: {}", + String::from_utf8_lossy(&buf[..8]) + ); + let t = committed(d.path()).expect("the rebase kept the commit usable"); + assert_eq!((t.count, t.end), (4, buf.len() as u64)); + let g = load(&format::grain_path(d.path(), "a.log")).unwrap(); + assert_eq!(g.chunk_count(), 4); + for i in 6..10 { + assert!(g.may_contain_all(i - 6, &marker(i))); + } + assert!(!g.may_contain_all(0, &marker(0))); + } } diff --git a/src/rotate.rs b/src/rotate.rs index 00a972b..4e92393 100644 --- a/src/rotate.rs +++ b/src/rotate.rs @@ -792,6 +792,7 @@ mod tests { format::TRUNK_EXT, format::RINGS_EXT, format::GRAIN_EXT, + format::GRAIN_COMMIT_EXT, format::BARK_EXT, format::SEQ_EXT, format::TRIM_EXT, diff --git a/src/store.rs b/src/store.rs index 82a1da2..d924e5b 100644 --- a/src/store.rs +++ b/src/store.rs @@ -1566,6 +1566,7 @@ impl Store { let _ = fs::remove_file(format::trunk_path(&self.dir, name)); let _ = fs::remove_file(format::rings_path(&self.dir, name)); let _ = fs::remove_file(format::grain_path(&self.dir, name)); + let _ = fs::remove_file(format::grain_commit_path(&self.dir, name)); let _ = fs::remove_file(format::bark_path(&self.dir, name)); let _ = fs::remove_file(format::sap_path(&self.dir, name)); let _ = fs::remove_file(format::sap_seal_path(&self.dir, name)); @@ -1600,6 +1601,10 @@ impl Store { format::grain_path(&self.dir, old), format::grain_path(&self.dir, new), ); + let _ = fs::rename( + format::grain_commit_path(&self.dir, old), + format::grain_commit_path(&self.dir, new), + ); let _ = fs::rename( format::bark_path(&self.dir, old), format::bark_path(&self.dir, new), From 08da32df1c5e623e7a42779a21222a7b1bc786c1 Mon Sep 17 00:00:00 2001 From: Torstein Tauno Svendsen Date: Fri, 2 Oct 2026 22:18:13 +0200 Subject: [PATCH 2/2] Extending the grain reads only the new end of the rings extend_grain parsed all of .rings on every call to learn the chunk count and take the records past what the grain covers. The records are a fixed stride after a declared header, so the count is the file's length and the tail is a single read. read_index_tail returns both. Its header validation is the one parse_index_versioned already did, factored into index_layout so the two cannot disagree about what a rings file is. On a store with an 89 MB grain, appending one chunk peaked at 96 MB before this series and 9 MB after. Co-Authored-By: Claude Sonnet 5.5 --- src/format.rs | 136 ++++++++++++++++++++++++++++++++++++++++++++------ src/grain.rs | 10 ++-- 2 files changed, 126 insertions(+), 20 deletions(-) diff --git a/src/format.rs b/src/format.rs index 0361f5e..315df5d 100644 --- a/src/format.rs +++ b/src/format.rs @@ -477,34 +477,29 @@ pub fn read_header_next_seq(f: &File) -> io::Result { } } -/// Parse rings content wherever it came from (a file, a bundle member). -pub fn parse_index_bytes(buf: &[u8]) -> io::Result> { - parse_index_versioned(buf).map(|(recs, _)| recs) -} - -/// As `parse_index_bytes`, plus which layout it found — what a WRITER needs, -/// since it must migrate a v1 file before appending rather than mixing two -/// record strides in one file. -pub fn parse_index_versioned(buf: &[u8]) -> io::Result<(Vec, RingsVersion)> { +/// The layout a `.rings` file declares, from its first bytes and its total +/// length: which version, how long its header is, and the stride of its +/// records. `prefix` need only reach the fixed header fields. +fn index_layout(prefix: &[u8], total_len: u64) -> io::Result<(RingsVersion, usize, usize)> { let bad = || { io::Error::new( io::ErrorKind::InvalidData, "not a timberfs index (bad magic)", ) }; - if buf.len() < RINGS_HEADER_LEN_V1 as usize { + if prefix.len() < RINGS_HEADER_LEN_V1 as usize { return Err(bad()); } - let (version, header, rec_len) = match &buf[..8] { + let (version, header, rec_len) = match &prefix[..8] { // The FILE's header length, not ours: a longer header from a later // version must shift record offsets for this reader too, which is // the whole point of declaring it. m if m == RINGS_MAGIC => { - if buf.len() < RINGS_HEADER_MIN as usize { + if prefix.len() < RINGS_HEADER_MIN as usize { return Err(bad()); } - let declared = u64::from_le_bytes(buf[8..16].try_into().unwrap()); - if declared < RINGS_HEADER_MIN || declared > buf.len() as u64 { + let declared = u64::from_le_bytes(prefix[8..16].try_into().unwrap()); + if declared < RINGS_HEADER_MIN || declared > total_len { return Err(io::Error::new( io::ErrorKind::InvalidData, format!( @@ -516,7 +511,7 @@ pub fn parse_index_versioned(buf: &[u8]) -> io::Result<(Vec, RingsV // older reader has to stop rather than answer from a layout it // does not know. Optional additions live in the reserved bytes // and set nothing, precisely so this stays rare. - let incompat = u64::from_le_bytes(buf[16..24].try_into().unwrap()); + let incompat = u64::from_le_bytes(prefix[16..24].try_into().unwrap()); if incompat != 0 { return Err(io::Error::new( io::ErrorKind::InvalidData, @@ -535,6 +530,50 @@ pub fn parse_index_versioned(buf: &[u8]) -> io::Result<(Vec, RingsV ), _ => return Err(bad()), }; + Ok((version, header, rec_len)) +} + +/// How many records a `.rings` file holds, and the ones from `from` on, +/// reading only those: the records are a fixed stride after a declared +/// header, so the count is the file's length and the tail is one seek. For +/// a caller that wants the end of a long index and not all of it. +pub fn read_index_tail(path: &Path, from: usize) -> io::Result<(usize, Vec)> { + let wrap = + |e: io::Error| io::Error::new(e.kind(), format!("reading index {}: {e}", path.display())); + let f = File::open(path) + .map_err(|e| io::Error::new(e.kind(), format!("opening index {}: {e}", path.display())))?; + let total = f.metadata().map_err(wrap)?.len(); + let mut prefix = vec![0u8; total.min(RINGS_HEADER_LEN) as usize]; + f.read_exact_at(&mut prefix, 0).map_err(wrap)?; + let (version, header, rec_len) = index_layout(&prefix, total).map_err(wrap)?; + let n = ((total - header as u64) / rec_len as u64) as usize; + if from >= n { + return Ok((n, Vec::new())); + } + let mut buf = vec![0u8; (n - from) * rec_len]; + f.read_exact_at(&mut buf, header as u64 + (from * rec_len) as u64) + .map_err(wrap)?; + let recs = buf + .chunks_exact(rec_len) + .enumerate() + .map(|(j, b)| match version { + RingsVersion::V2 => ChunkRecord::from_bytes(b), + RingsVersion::V1 => ChunkRecord::from_bytes_v1(b, (from + j) as u64), + }) + .collect(); + Ok((n, recs)) +} + +/// Parse rings content wherever it came from (a file, a bundle member). +pub fn parse_index_bytes(buf: &[u8]) -> io::Result> { + parse_index_versioned(buf).map(|(recs, _)| recs) +} + +/// As `parse_index_bytes`, plus which layout it found — what a WRITER needs, +/// since it must migrate a v1 file before appending rather than mixing two +/// record strides in one file. +pub fn parse_index_versioned(buf: &[u8]) -> io::Result<(Vec, RingsVersion)> { + let (version, header, rec_len) = index_layout(buf, buf.len() as u64)?; let n = (buf.len() - header) / rec_len; let mut out = Vec::with_capacity(n); for i in 0..n { @@ -697,4 +736,71 @@ mod tests { let back = ChunkRecord::from_bytes_v1(v1, 99); assert_eq!((back.seq, back.uncomp_start), (99, 70)); } + + fn write_image(name: &str, buf: &[u8]) -> PathBuf { + let p = std::env::temp_dir().join(format!("timberfs-tail-{}-{name}", std::process::id())); + std::fs::write(&p, buf).unwrap(); + p + } + + #[test] + fn the_tail_is_the_full_parse_from_that_point() { + for (label, buf) in [ + ("v2", image(RINGS_HEADER_LEN, 0, 9)), + ("later", image(128, 0, 9)), + ] { + let p = write_image(label, &buf); + let all = read_index(&p).unwrap(); + for from in [0, 1, 4, 8, 9, 10, 100] { + let (n, tail) = read_index_tail(&p, from).unwrap(); + assert_eq!(n, all.len(), "{label} from {from}"); + let want = &all[from.min(all.len())..]; + assert_eq!( + tail.iter().map(|c| c.seq).collect::>(), + want.iter().map(|c| c.seq).collect::>(), + "{label} from {from}" + ); + } + let _ = std::fs::remove_file(p); + } + } + + #[test] + fn the_tail_numbers_a_v1_index_the_way_the_full_parse_does() { + let mut buf = RINGS_MAGIC_V1.to_vec(); + for i in 0..5u64 { + buf.extend_from_slice(&rec(i).to_bytes()[..RECORD_LEN_V1]); + } + let p = write_image("v1", &buf); + let all = read_index(&p).unwrap(); + let (n, tail) = read_index_tail(&p, 2).unwrap(); + assert_eq!((n, tail.len()), (5, 3)); + assert_eq!( + tail.iter().map(|c| c.seq).collect::>(), + all[2..].iter().map(|c| c.seq).collect::>() + ); + let _ = std::fs::remove_file(p); + } + + #[test] + fn a_partial_trailing_record_is_not_counted() { + let mut buf = image(RINGS_HEADER_LEN, 0, 4); + buf.extend_from_slice(&[7u8; 20]); + let p = write_image("partial", &buf); + let (n, tail) = read_index_tail(&p, 3).unwrap(); + assert_eq!((n, tail.len()), (4, 1)); + let _ = std::fs::remove_file(p); + } + + #[test] + fn the_tail_refuses_what_the_full_parse_refuses() { + let p = write_image("incompat", &image(RINGS_HEADER_LEN, 0b10, 2)); + let e = read_index_tail(&p, 0).unwrap_err(); + assert!(e.to_string().contains("unsupported features"), "{e}"); + let _ = std::fs::remove_file(p); + + let p = write_image("magic", b"NOTRINGS and then some bytes"); + assert!(read_index_tail(&p, 0).is_err()); + let _ = std::fs::remove_file(p); + } } diff --git a/src/grain.rs b/src/grain.rs index 5ef41e8..f29489a 100644 --- a/src/grain.rs +++ b/src/grain.rs @@ -379,15 +379,15 @@ pub fn extend_grain(dir: &Path, name: &str) -> anyhow::Result<()> { let Some((tail, known)) = locate(dir, name, &out)? else { return build_grain(dir, name); }; - let records = format::read_index(&format::rings_path(dir, name))?; let covered = tail.count as usize; - if covered > records.len() { + let (chunks, new) = format::read_index_tail(&format::rings_path(dir, name), covered)?; + if covered > chunks { return build_grain(dir, name); } if out.metadata()?.len() > tail.end { out.set_len(tail.end)?; } - if covered == records.len() { + if covered == chunks { if !known { save_commit(dir, name, &out, tail)?; } @@ -396,7 +396,7 @@ pub fn extend_grain(dir: &Path, name: &str) -> anyhow::Result<()> { let trunk = File::open(format::trunk_path(dir, name)) .with_context(|| format!("opening {}", format::trunk_path(dir, name).display()))?; let mut woff = tail.end; - for c in &records[covered..] { + for c in &new { let mut comp = vec![0u8; c.comp_len as usize]; trunk.read_exact_at(&mut comp, c.comp_start)?; let data = zstd::stream::decode_all(&comp[..]) @@ -414,7 +414,7 @@ pub fn extend_grain(dir: &Path, name: &str) -> anyhow::Result<()> { name, &out, Tail { - count: records.len() as u64, + count: chunks as u64, end: woff, }, )?;