From cf72244adee9c4ec3f2049179f30aa4089c424a8 Mon Sep 17 00:00:00 2001 From: Torstein Tauno Svendsen Date: Sat, 3 Oct 2026 00:58:04 +0200 Subject: [PATCH] An idle store is not read on every poll A follower polls every store it ships, once a second, and most are idle on most polls. Each poll parsed the store's whole .rings, and serve() also loaded its whole .grain, to find that there was nothing to send. The shipper leaves a store out of the read when its cursor sits exactly at the end of the flushed chunks, no write-ahead entries are pending and no .sap.seal handoff is in flight: the read also follows the live edge, so a chunk-level test alone would miss a wal store. The stores returned are the ones read, because the parse pairs them with the buffer. The frames sender skips serve() when the store has no chunk at or after the resume seq, decided from the index header and its last record. serve() no longer copies the store's selection twice before looking at max_chunks, and reads only the grain pages it will send instead of loading the whole grain: grain::read_pages walks the length prefixes with a bounded buffer. Three stores of 300,000 chunks, idle under feed: 48 MB/s read and 128 MB peak RSS before, 0 MB/s and 7 MB after. Co-Authored-By: Claude Sonnet 5.5 --- src/frames.rs | 5 +++ src/grain.rs | 71 +++++++++++++++++++++++++++++++ src/serve.rs | 115 ++++++++++++++++++++++++++++++++++++++------------ src/ship.rs | 94 ++++++++++++++++++++++++++++++++++++++++- 4 files changed, 257 insertions(+), 28 deletions(-) diff --git a/src/frames.rs b/src/frames.rs index 78288ed..32d73de 100644 --- a/src/frames.rs +++ b/src/frames.rs @@ -924,6 +924,11 @@ fn open_send_stream( /// Serve one store's turn: from where its stream stands, up to /// `CHUNKS_PER_TURN`. Returns how many chunks went on the wire. fn ship_turn(w: &mut impl Write, s: &mut Stream, opts: &SendOpts) -> anyhow::Result { + // An idle store is the common case on every poll: a header and one + // record say so, where serving reads the whole index to find nothing. + if crate::serve::nothing_from(&s.path, s.resume) { + return Ok(0); + } let mut body = Vec::new(); let served = crate::serve::serve( &s.path, diff --git a/src/grain.rs b/src/grain.rs index 5fc0d2f..33b3577 100644 --- a/src/grain.rs +++ b/src/grain.rs @@ -320,6 +320,47 @@ impl Grain { } } +/// The filters of records `from..from + count`, read by walking the length +/// prefixes with a bounded buffer: memory is the pages asked for, not the +/// grain. Shorter than `count` where the grain does not reach, which is +/// "missing means scan" for the rest. +pub fn read_pages(path: &Path, from: usize, count: usize) -> Vec> { + let mut pages = Vec::new(); + let Ok(f) = File::open(path) else { + return pages; + }; + let Ok(Some(s)) = stamp(&f) else { + return pages; + }; + let mut r = BufReader::with_capacity(WALK_BUF, &f); + if r.seek(SeekFrom::Start(s.first)).is_err() { + return pages; + } + let (mut off, mut idx) = (s.first, 0usize); + while idx < from.saturating_add(count) && off + 4 <= s.len { + let mut p = [0u8; 4]; + if r.read_exact(&mut p).is_err() { + break; + } + let len = u32::from_le_bytes(p) as u64; + if off + 4 + len > s.len { + break; + } + if idx >= from { + let mut page = vec![0u8; len as usize]; + if r.read_exact(&mut page).is_err() { + break; + } + pages.push(page); + } else if r.seek_relative(len as i64).is_err() { + break; + } + off += 4 + len; + idx += 1; + } + pages +} + /// The grain's size and how many chunks it covers, without loading it: /// the commit when it still describes the file, otherwise one walk. Opens /// nothing for writing, so a caller that may only read the store can use it. @@ -1159,4 +1200,34 @@ mod tests { .unwrap(); assert_eq!(coverage(d.path(), "a.log"), None); } + + #[test] + fn pages_read_by_range_are_the_loaded_ones() { + let d = TempDir::new(); + write_grain(d.path(), "a.log", 9); + let path = format::grain_path(d.path(), "a.log"); + let loaded = load(&path).unwrap(); + for (from, count) in [(0, 9), (0, 3), (2, 4), (7, 2), (8, 50), (0, 0)] { + let got = read_pages(&path, from, count); + let want: Vec> = (from..(from + count).min(9)) + .map(|i| loaded.page(i).unwrap().to_vec()) + .collect(); + assert_eq!(got, want, "pages {from}..{}", from + count); + } + assert!(read_pages(&path, 40, 3).is_empty(), "past the end"); + assert!(read_pages(&d.path().join("none"), 0, 3).is_empty()); + } + + #[test] + fn pages_stop_at_a_torn_tail() { + let d = TempDir::new(); + write_grain(d.path(), "a.log", 4); + let path = format::grain_path(d.path(), "a.log"); + let g = OpenOptions::new().append(true).open(&path).unwrap(); + (&g).write_all(&90u32.to_le_bytes()).unwrap(); + (&g).write_all(&[7u8; 5]).unwrap(); + drop(g); + assert_eq!(read_pages(&path, 0, 10).len(), 4); + assert_eq!(read_pages(&path, 3, 10).len(), 1); + } } diff --git a/src/serve.rs b/src/serve.rs index 4be8af3..4eb1a39 100644 --- a/src/serve.rs +++ b/src/serve.rs @@ -115,6 +115,22 @@ pub fn push_seq(out: &mut Vec, s: u64) { } } +/// True when `input` provably holds no chunk at or after `first_seq`: its +/// last record is older, or it has none. Answers false for anything it +/// cannot read, so the request is served and says why. +pub fn nothing_from(input: &Path, first_seq: u64) -> bool { + if crate::query::is_bundle(input) { + return false; + } + let Ok((dir, name)) = crate::query::resolve_backing(input) else { + return false; + }; + matches!( + crate::format::read_index_last(&crate::format::rings_path(&dir, &name)), + Ok((_, last)) if last.is_none_or(|c| c.seq < first_seq) + ) +} + /// 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)?; @@ -137,23 +153,12 @@ pub fn serve(input: &Path, req: &Request, out: &mut impl Write) -> anyhow::Resul // Which chunks the request selects. Done before stream-open so its // declared range is what is actually coming, not what was asked for. - let selected: Vec = handle - .records - .iter() - .filter(|c| { - c.seq >= req.first_seq && (req.last_seq == frame::OPEN_ENDED || c.seq <= req.last_seq) - }) - .copied() - .collect(); - let positions: Vec = handle - .records - .iter() - .enumerate() - .filter(|(_, c)| { - c.seq >= req.first_seq && (req.last_seq == frame::OPEN_ENDED || c.seq <= req.last_seq) - }) - .map(|(i, _)| i) - .collect(); + // By position, not copied out: a store is far larger than a turn. + let wanted = |c: &crate::format::ChunkRecord| { + c.seq >= req.first_seq && (req.last_seq == frame::OPEN_ENDED || c.seq <= req.last_seq) + }; + let first = handle.records.iter().position(wanted); + let last = handle.records.iter().rposition(wanted); let mut buf = Vec::new(); frame::encode( @@ -162,8 +167,12 @@ pub fn serve(input: &Path, req: &Request, out: &mut impl Write) -> anyhow::Resul frame: Frame::StreamOpen { origin_id: origin_of(&bark), sender_id: id_of(&bark), - first_seq: selected.first().map(|c| c.seq).unwrap_or(req.first_seq), - last_seq: selected.last().map(|c| c.seq).unwrap_or(frame::OPEN_ENDED), + first_seq: first + .map(|i| handle.records[i].seq) + .unwrap_or(req.first_seq), + last_seq: last + .map(|i| handle.records[i].seq) + .unwrap_or(frame::OPEN_ENDED), mode: req.mode, provenance: serde_json::to_vec(&travels(&bark, input)) .context("serializing what the store is")?, @@ -181,7 +190,7 @@ pub fn serve(input: &Path, req: &Request, out: &mut impl Write) -> anyhow::Resul &Framed { stream: req.stream, frame: Frame::Coverage { - runs: runs_of(selected.iter().map(|c| c.seq)), + runs: runs_of(handle.records.iter().filter(|c| wanted(c)).map(|c| c.seq)), }, }, &mut buf, @@ -191,22 +200,33 @@ pub fn serve(input: &Path, req: &Request, out: &mut impl Write) -> anyhow::Resul return Ok(stats); } - let grain = if req.sidecars { - grain_path(input).and_then(|p| crate::grain::load(&p).ok()) - } else { - None + let Some((first, last)) = first.zip(last) else { + return Ok(stats); + }; + // Only the pages this request can send: a grain is read by position, + // and the chunks it ships are the ones from `first` on. + let span = req + .max_chunks + .map_or(last - first + 1, |n| (n as usize).min(last - first + 1)); + let pages = match req.sidecars.then(|| grain_path(input)).flatten() { + Some(p) => crate::grain::read_pages(&p, first, span), + None => Vec::new(), }; - for (c, pos) in selected.into_iter().zip(positions) { + for pos in first..=last { + let c = handle.records[pos]; + if !wanted(&c) { + continue; + } if req.max_chunks.is_some_and(|n| stats.chunks >= n) { break; } stats.last_examined = Some(c.seq); let mut sidecars = Vec::new(); - if let Some(page) = grain.as_ref().and_then(|g| g.page(pos)) { + if let Some(page) = pos.checked_sub(first).and_then(|i| pages.get(i)) { sidecars.push(Sidecar { kind: Sidecar::tag(GRAIN_TAG), - bytes: page.to_vec(), + bytes: page.clone(), }); } let comp = if req.mode == Mode::Frames { @@ -664,4 +684,45 @@ mod tests { std::fs::write(&gpath, &bytes[..15]).unwrap(); assert_eq!(grain_header(&store), None, "shorter than a header"); } + + #[test] + fn a_resume_mid_store_sends_the_pages_of_its_own_chunks() { + let d = TempDir::new(); + let p = a_store(d.path(), "mid", 6); + let (dir, name) = crate::query::resolve_backing(&p).unwrap(); + crate::grain::extend_grain(&dir, &name).unwrap(); + let g = crate::grain::load(&crate::format::grain_path(&dir, &name)).unwrap(); + + let mut req = Request::everything(Mode::Index); + req.first_seq = 2; + req.max_chunks = Some(3); + let mut buf = Vec::new(); + let served = serve(&p, &req, &mut buf).unwrap(); + assert_eq!(served.chunks, 3); + let got: Vec<(u64, Vec)> = frames(&buf) + .iter() + .filter_map(|f| match &f.frame { + Frame::Chunk { seq, sidecars, .. } => Some((*seq, sidecars[0].bytes.clone())), + _ => None, + }) + .collect(); + assert_eq!(got.iter().map(|c| c.0).collect::>(), [2, 3, 4]); + for (seq, page) in got { + assert_eq!(page, g.page(seq as usize).unwrap(), "chunk {seq}"); + } + } + + #[test] + fn nothing_from_is_true_only_past_the_last_chunk() { + let d = TempDir::new(); + let p = a_store(d.path(), "idle", 3); + assert!(!nothing_from(&p, 0)); + assert!(!nothing_from(&p, 2), "the last chunk is still to send"); + assert!(nothing_from(&p, 3)); + assert!(nothing_from(&p, 100)); + assert!( + !nothing_from(&d.path().join("absent.log"), 3), + "unreadable is served, so the reason is said" + ); + } } diff --git a/src/ship.rs b/src/ship.rs index ec82536..3f529a8 100644 --- a/src/ship.rs +++ b/src/ship.rs @@ -106,6 +106,29 @@ fn tape_end(dir: &Path, name: &str) -> u64 { dropped + last.map(|c| c.uncomp_end()).unwrap_or(0) } +/// True when reading `name` from tape offset `at` provably returns nothing: +/// the cursor sits exactly at the end of the flushed chunks, and the live +/// edge, which the read also follows, has no write-ahead entries pending or +/// mid-handoff. Anything it cannot read answers false, so the read happens +/// and says why. +fn nothing_to_read(dir: &Path, name: &str, at: u64) -> bool { + use crate::format; + if format::sap_seal_path(dir, name).exists() { + return false; + } + if let Ok(m) = std::fs::metadata(format::sap_path(dir, name)) { + if m.len() > crate::sap::HEADER_LEN { + return false; + } + } + let Ok((_, last)) = format::read_index_last(&format::rings_path(dir, name)) else { + return false; + }; + let end = + crate::query::dropped_bytes_of(&dir.join(name)) + last.map(|c| c.uncomp_end()).unwrap_or(0); + end == at +} + /// How many stores a batch may span, and how many entries it may hold, by /// default. The entry cap is the destination's batch size; the store cap /// bounds a poll's syscalls, since every store in the selection is opened @@ -328,7 +351,6 @@ impl Shipper { if stores.is_empty() { return Ok((Vec::new(), stores, matched)); } - let files: Vec = stores.iter().map(|(_, p, _)| p.clone()).collect(); let mut cursor = self.positions.cursor(); for (id, _, _) in &stores { // Never backwards: the recorded position is the floor, so a @@ -338,6 +360,22 @@ impl Shipper { cursor.insert(id.clone(), at); } } + // The read opens every store it is given, so an idle one is left + // out of it, and out of the answer: most stores are idle on most + // polls. The stores returned are the ones read, because the parse + // pairs them with the buffer. + let stores: Vec = stores + .into_iter() + .filter(|(id, path, _)| { + let at = cursor.get(id).copied(); + let (dir, name) = (path.parent(), path.file_name().and_then(|n| n.to_str())); + !matches!((at, dir, name), (Some(at), Some(d), Some(n)) if nothing_to_read(d, n, at)) + }) + .collect(); + if stores.is_empty() { + return Ok((Vec::new(), stores, matched)); + } + let files: Vec = stores.iter().map(|(_, p, _)| p.clone()).collect(); let mut buf = Vec::new(); crate::query::read_forward(&mut buf, &files, &cursor, self.batch_entries)?; Ok((buf, stores, matched)) @@ -990,4 +1028,58 @@ mod tests { assert!(bodies(&b).iter().all(|l| l.contains("web"))); let _ = std::fs::remove_dir_all(&root); } + + #[test] + fn a_store_is_idle_only_at_exactly_its_tape_end_with_nothing_live() { + let root = forest("idle", &[("a", "x", 3)]); + let (dir, name) = (root.join("a"), "a.log"); + let end = tape_end(&dir, name); + assert!(end > 0); + assert!(nothing_to_read(&dir, name, end)); + assert!(!nothing_to_read(&dir, name, end - 1), "behind it: to read"); + assert!(!nothing_to_read(&dir, name, end + 1), "beyond it: say so"); + + // The live edge is read too: entries pending in the write-ahead + // segment, or a segment mid-handoff, are not idle. + let sap = crate::format::sap_path(&dir, name); + std::fs::write(&sap, vec![0u8; crate::sap::HEADER_LEN as usize]).unwrap(); + assert!(nothing_to_read(&dir, name, end), "a header alone is empty"); + std::fs::write(&sap, vec![0u8; crate::sap::HEADER_LEN as usize + 1]).unwrap(); + assert!(!nothing_to_read(&dir, name, end), "pending entries"); + std::fs::remove_file(&sap).unwrap(); + std::fs::write(crate::format::sap_seal_path(&dir, name), b"x").unwrap(); + assert!(!nothing_to_read(&dir, name, end), "a handoff in flight"); + + assert!( + !nothing_to_read(&dir, "absent.log", 0), + "unreadable is read, so the reason is said" + ); + let _ = std::fs::remove_dir_all(&root); + } + + #[test] + fn a_caught_up_store_is_left_out_of_the_read_and_wakes_when_it_grows() { + let root = forest("wake", &[("a", "x", 3), ("b", "x", 3)]); + let mut sh = shipper(&root, "service=x", 100); + let first = sh.poll().unwrap(); + assert_eq!(first.slices.len(), 2); + sh.accept(&first).unwrap(); + + let idle = sh.poll().unwrap(); + assert!(idle.slices.is_empty(), "both caught up: nothing is read"); + assert_eq!(idle.matched, 2, "still matched, just idle"); + + append(&root, "b", 2); + let woke = sh.poll().unwrap(); + let ids: Vec<_> = woke.slices.iter().map(|s| s.path.clone()).collect(); + assert_eq!( + ids, + vec![root.join("b").join("b.log")], + "only the grown one" + ); + assert_eq!(woke.slices[0].entries.len(), 2); + sh.accept(&woke).unwrap(); + assert!(sh.poll().unwrap().slices.is_empty()); + let _ = std::fs::remove_dir_all(&root); + } }