From 398a8b18344e0257864592b7d5c581b25dc5383f Mon Sep 17 00:00:00 2001 From: Torstein Tauno Svendsen Date: Sat, 3 Oct 2026 00:33:29 +0200 Subject: [PATCH] Small reads stop reading whole files Five places read a whole file to use a few bytes of it: - migrate_rings read the entire .rings on every writer open to compare an 8-byte magic. It now reads the magic, and the file only when it is v1. - grain_header read the entire .grain on every ship turn to keep its first 16 bytes. - tape_end and the rotate summary parsed the entire .rings for the last record. read_index_last reads the header and that one record. - Session::coverage walked every chunk of the destination on every ack, which is every received chunk. The session now keeps the runs as chunks arrive; nothing in a receive session removes chunks, so the list stays exact. The records sink re-parsed .bark about four times a second to look for a retention or wal change. It now uses the stat-gated LivePolicy that append already uses, so a steady state is one stat per tick. A manifest that stops parsing keeps the last good policy with one warning, as it does for append, where the sink used to skip enforcement silently. Co-Authored-By: Claude Sonnet 5.5 --- src/format.rs | 54 +++++++++++++++++++++++++++++++++++++++++++ src/receive.rs | 14 ++++++++++-- src/rotate.rs | 10 ++++---- src/serve.rs | 54 +++++++++++++++++++++++++++++++++++++------ src/ship.rs | 7 +++--- src/sink.rs | 62 ++++++++++++++++++++++++++++---------------------- src/store.rs | 26 +++++++++++++++++---- 7 files changed, 179 insertions(+), 48 deletions(-) diff --git a/src/format.rs b/src/format.rs index 315df5d..59a10e5 100644 --- a/src/format.rs +++ b/src/format.rs @@ -538,6 +538,22 @@ fn index_layout(prefix: &[u8], total_len: u64) -> io::Result<(RingsVersion, usiz /// 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)> { + index_from(path, |_| from) +} + +/// How many records a `.rings` file holds and the last of them, from one +/// header read and one record read: for a caller that wants where the tape +/// ends and not the index. +pub fn read_index_last(path: &Path) -> io::Result<(usize, Option)> { + let (n, mut recs) = index_from(path, |n| n.saturating_sub(1))?; + Ok((n, recs.pop())) +} + +/// The records from the index `start(count)` on, and the count. +fn index_from( + path: &Path, + start: impl FnOnce(usize) -> 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) @@ -547,6 +563,7 @@ pub fn read_index_tail(path: &Path, from: usize) -> io::Result<(usize, Vec= n { return Ok((n, Vec::new())); } @@ -803,4 +820,41 @@ mod tests { assert!(read_index_tail(&p, 0).is_err()); let _ = std::fs::remove_file(p); } + + #[test] + fn the_last_record_is_what_a_full_parse_ends_with() { + for (label, buf) in [ + ("last-v2", image(RINGS_HEADER_LEN, 0, 6)), + ("last-later", image(128, 0, 6)), + ("last-empty", image(RINGS_HEADER_LEN, 0, 0)), + ] { + let p = write_image(label, &buf); + let all = read_index(&p).unwrap(); + let (n, last) = read_index_last(&p).unwrap(); + assert_eq!(n, all.len(), "{label}"); + assert_eq!(last.map(|c| c.seq), all.last().map(|c| c.seq), "{label}"); + let _ = std::fs::remove_file(p); + } + } + + #[test] + fn the_last_record_ignores_a_partial_one_and_numbers_v1_by_position() { + let mut buf = image(RINGS_HEADER_LEN, 0, 3); + buf.extend_from_slice(&[9u8; 17]); + let p = write_image("last-partial", &buf); + assert_eq!(read_index_last(&p).unwrap().1.map(|c| c.seq), Some(2)); + let _ = std::fs::remove_file(p); + + let mut v1 = RINGS_MAGIC_V1.to_vec(); + for i in 0..4u64 { + v1.extend_from_slice(&rec(i).to_bytes()[..RECORD_LEN_V1]); + } + let p = write_image("last-v1", &v1); + let all = read_index(&p).unwrap(); + assert_eq!( + read_index_last(&p).unwrap().1.map(|c| c.seq), + all.last().map(|c| c.seq) + ); + let _ = std::fs::remove_file(p); + } } diff --git a/src/receive.rs b/src/receive.rs index 35074b8..e831deb 100644 --- a/src/receive.rs +++ b/src/receive.rs @@ -188,6 +188,7 @@ pub struct Session { cfg: crate::store::Config, adopt_pages: bool, out: Received, + runs: Vec, _dir_lock: std::fs::File, _file_lock: std::fs::File, } @@ -307,7 +308,14 @@ impl Session { } } + let runs = crate::serve::runs_of( + st.files + .get(&name) + .into_iter() + .flat_map(|f| f.chunks.iter().map(|c| c.seq)), + ); Ok(Session { + runs, out: Received { store: dest.to_path_buf(), created: !existed, @@ -340,8 +348,7 @@ impl Session { /// What this destination holds now — the ack, and a coverage answer. pub fn coverage(&self) -> Vec { - let file = self.st.files.get(&self.name).expect("created in open"); - crate::serve::runs_of(file.chunks.iter().map(|c| c.seq)) + self.runs.clone() } /// Apply one frame. Returns false for a frame that ends the stream @@ -382,6 +389,9 @@ impl Session { .with_context(|| format!("appending chunk {seq} to {}", self.name))?; self.out.chunks += 1; self.out.comp_bytes += comp_len; + if let Some(c) = self.st.files.get(&self.name).and_then(|f| f.chunks.last()) { + crate::serve::push_seq(&mut self.runs, c.seq); + } for s in &sidecars { if self.adopt_pages && s.kind == crate::frame::Sidecar::tag(crate::serve::GRAIN_TAG) diff --git a/src/rotate.rs b/src/rotate.rs index 4e92393..3d73187 100644 --- a/src/rotate.rs +++ b/src/rotate.rs @@ -680,15 +680,15 @@ pub fn cmd_rotate( } })?; println!("rotated through the live mount on {}", mp.display()); - let after = format::read_index(&rings)?; - println!(" source keeps {} chunk(s)", after.len()); + let (kept, _) = format::read_index_last(&rings)?; + println!(" source keeps {kept} chunk(s)"); if let Some(t) = &target_name { - let ti = format::read_index(&format::rings_path(&dir, t))?; + let (chunks, last) = format::read_index_last(&format::rings_path(&dir, t))?; println!( " {} now has {} chunk(s), {} on disk", t, - ti.len(), - human_bytes(ti.last().map(|c| c.comp_end()).unwrap_or(0)) + chunks, + human_bytes(last.map(|c| c.comp_end()).unwrap_or(0)) ); } } diff --git a/src/serve.rs b/src/serve.rs index 6af5549..4be8af3 100644 --- a/src/serve.rs +++ b/src/serve.rs @@ -101,15 +101,20 @@ pub struct LastSent { pub fn runs_of(seqs: impl IntoIterator) -> Vec { let mut out: Vec = Vec::new(); for s in seqs { - match out.last_mut() { - Some(last) if s == last.end + 1 => last.end = s, - Some(last) if s <= last.end => {} // duplicate or out of order - _ => out.push(Run { start: s, end: s }), - } + push_seq(&mut out, s); } out } +/// One step of `runs_of`, for a caller that keeps the runs as chunks arrive. +pub fn push_seq(out: &mut Vec, s: u64) { + match out.last_mut() { + Some(last) if s == last.end + 1 => last.end = s, + Some(last) if s <= last.end => {} // duplicate or out of order + _ => out.push(Run { start: s, end: s }), + } +} + /// Write `input`'s answer to `req` as frames. pub fn serve(input: &Path, req: &Request, out: &mut impl Write) -> anyhow::Result { let mut handle = crate::query::open_source(input)?; @@ -308,8 +313,10 @@ fn grain_path(input: &Path) -> Option { /// The `.grain`'s 16-byte parameter header, which is store-level and so /// rides `stream-open` rather than every chunk. fn grain_header(input: &Path) -> Option> { - let bytes = std::fs::read(grain_path(input)?).ok()?; - (bytes.len() >= 16).then(|| bytes[..16].to_vec()) + let f = std::fs::File::open(grain_path(input)?).ok()?; + let mut header = [0u8; 16]; + std::os::unix::fs::FileExt::read_exact_at(&f, &mut header, 0).ok()?; + Some(header.to_vec()) } #[cfg(test)] @@ -624,4 +631,37 @@ mod tests { other => panic!("{other:?}"), } } + + #[test] + fn keeping_runs_as_chunks_arrive_is_the_fold_over_all_of_them() { + let seqs = [0u64, 1, 2, 2, 5, 6, 4, 9, 10, 11, 11, 30]; + let mut kept = Vec::new(); + for s in seqs { + push_seq(&mut kept, s); + } + assert_eq!(kept, runs_of(seqs)); + assert_eq!( + kept, + vec![ + Run { start: 0, end: 2 }, + Run { start: 5, end: 6 }, + Run { start: 9, end: 11 }, + Run { start: 30, end: 30 }, + ] + ); + } + + #[test] + fn the_grain_header_is_its_first_sixteen_bytes_and_nothing_else_is() { + let d = TempDir::new(); + let store = a_store(d.path(), "g", 2); + assert_eq!(grain_header(&store), None, "no grain, no header"); + let (dir, name) = crate::query::resolve_backing(&store).unwrap(); + let gpath = crate::format::grain_path(&dir, &name); + let bytes: Vec = (0..64u8).collect(); + std::fs::write(&gpath, &bytes).unwrap(); + assert_eq!(grain_header(&store), Some(bytes[..16].to_vec())); + std::fs::write(&gpath, &bytes[..15]).unwrap(); + assert_eq!(grain_header(&store), None, "shorter than a header"); + } } diff --git a/src/ship.rs b/src/ship.rs index 2290dac..ec82536 100644 --- a/src/ship.rs +++ b/src/ship.rs @@ -100,9 +100,10 @@ fn instant(s: Option<&str>) -> Option> { /// writer's own lock. fn tape_end(dir: &Path, name: &str) -> u64 { let dropped = crate::query::dropped_bytes_of(&dir.join(name)); - let chunks = - crate::format::read_index(&crate::format::rings_path(dir, name)).unwrap_or_default(); - dropped + chunks.last().map(|c| c.uncomp_end()).unwrap_or(0) + let last = crate::format::read_index_last(&crate::format::rings_path(dir, name)) + .ok() + .and_then(|(_, last)| last); + dropped + last.map(|c| c.uncomp_end()).unwrap_or(0) } /// How many stores a batch may span, and how many entries it may hold, by diff --git a/src/sink.rs b/src/sink.rs index 7bb20e9..018b5b5 100644 --- a/src/sink.rs +++ b/src/sink.rs @@ -169,10 +169,18 @@ pub fn cmd_records_sink( } else { None }; - // Chunks already folded into the declared grain. Extending re-reads - // the whole grain file, so we only do it when the flushed-chunk set - // actually changed — not every idle second. + // Chunks already folded into the declared grain: extend only when + // the flushed-chunk set actually changed, not every idle second. let mut indexed_chunks: usize = 0; + let mut live = crate::append::LivePolicy { + dir: dir.clone(), + name: name.clone(), + last: crate::bark::Retention::default(), + fields: Default::default(), + warned: false, + stamp: None, + reparsed: false, + }; Some(thread::spawn(move || { while !stop.load(Ordering::Relaxed) { thread::sleep(Duration::from_millis(1000)); @@ -218,33 +226,33 @@ pub fn cmd_records_sink( } st.lock().unwrap().flush_aged(); st.lock().unwrap().sap_sync_all(); - st.lock().unwrap().sync_wal_declarations(); - match crate::bark::declared_retention(&dir, &name) { - Ok(policy) if policy.is_some() => { - let fields = crate::follower::subject_of(&dir, &name); - let next_seq = st.lock().unwrap().next_seq(&name).unwrap_or(0); - let held = crate::follower::TickInterest::default() - .floor(&policy, &fields, next_seq); - match st.lock().unwrap().enforce_retention( - &name, - policy.max_age_ms, - policy.max_comp_bytes, - held.floor, - ) { - Err(e) => { - eprintln!("timberfs: {name}: background retention failed: {e}") - } - Ok(Some(stats)) => { - if let Some(record) = - crate::follower::override_record(&name, &policy, &stats, &held) - { - eprintln!("{record}"); - } + let policy = live.refresh(); + if live.reparsed { + st.lock().unwrap().sync_wal_declarations(); + } + if policy.is_some() { + let fields = live.fields.clone(); + let next_seq = st.lock().unwrap().next_seq(&name).unwrap_or(0); + let held = + crate::follower::TickInterest::default().floor(&policy, &fields, next_seq); + match st.lock().unwrap().enforce_retention( + &name, + policy.max_age_ms, + policy.max_comp_bytes, + held.floor, + ) { + Err(e) => { + eprintln!("timberfs: {name}: background retention failed: {e}") + } + Ok(Some(stats)) => { + if let Some(record) = + crate::follower::override_record(&name, &policy, &stats, &held) + { + eprintln!("{record}"); } - Ok(None) => {} } + Ok(None) => {} } - _ => {} } // Keep the declared index current while streaming: extend the // grain whenever the flushed-chunk set changed (a flush added diff --git a/src/store.rs b/src/store.rs index d924e5b..475ab79 100644 --- a/src/store.rs +++ b/src/store.rs @@ -173,16 +173,21 @@ fn parse_trim_marker(text: &str) -> Option<(u64, u64, u64)> { /// all — they synthesize the same numbers when parsing v1. fn migrate_rings(dir: &Path, name: &str) -> io::Result<()> { let p = format::rings_path(dir, name); - let buf = match fs::read(&p) { - Ok(b) => b, + let f = match fs::File::open(&p) { + Ok(f) => f, Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(()), Err(e) => return Err(e), }; // Absent, empty, already v2, or something else entirely: not ours to // touch. A bad magic is left for `open` to report as it always has. - if buf.len() < 8 || &buf[..8] != format::RINGS_MAGIC_V1 { - return Ok(()); + let mut magic = [0u8; 8]; + match f.read_exact_at(&mut magic, 0) { + Ok(()) if &magic == format::RINGS_MAGIC_V1 => {} + Ok(()) => return Ok(()), + Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => return Ok(()), + Err(e) => return Err(e), } + let buf = fs::read(&p)?; let (records, _) = format::parse_index_versioned(&buf)?; let next_seq = records.last().map(|c| c.seq + 1).unwrap_or(0); let mut idx = @@ -3270,4 +3275,17 @@ mod tests { "the sap's logical base must track buffer_start after a head trim" ); } + + #[test] + fn migrating_leaves_alone_whatever_is_not_a_v1_index() { + let dir = TempDir::new(); + let p = format::rings_path(dir.path(), "app"); + migrate_rings(dir.path(), "app").unwrap(); + assert!(!p.exists(), "a missing index stays missing"); + for body in [&b""[..], &b"short"[..], &b"NOTRINGS and then some more"[..]] { + fs::write(&p, body).unwrap(); + migrate_rings(dir.path(), "app").unwrap(); + assert_eq!(fs::read(&p).unwrap(), body, "{body:?} is not ours to touch"); + } + } }