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
5 changes: 5 additions & 0 deletions src/frames.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<u64> {
// 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,
Expand Down
71 changes: 71 additions & 0 deletions src/grain.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Vec<u8>> {
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.
Expand Down Expand Up @@ -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<Vec<u8>> = (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);
}
}
115 changes: 88 additions & 27 deletions src/serve.rs
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,22 @@ pub fn push_seq(out: &mut Vec<Run>, 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<Served> {
let mut handle = crate::query::open_source(input)?;
Expand All @@ -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<crate::format::ChunkRecord> = 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<usize> = 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(
Expand All @@ -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")?,
Expand All @@ -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,
Expand All @@ -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 {
Expand Down Expand Up @@ -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<u8>)> = 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::<Vec<_>>(), [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"
);
}
}
94 changes: 93 additions & 1 deletion src/ship.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -328,7 +351,6 @@ impl Shipper {
if stores.is_empty() {
return Ok((Vec::new(), stores, matched));
}
let files: Vec<PathBuf> = 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
Expand All @@ -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<Store> = 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<PathBuf> = 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))
Expand Down Expand Up @@ -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);
}
}
Loading