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); + } }