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
54 changes: 54 additions & 0 deletions src/format.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<ChunkRecord>)> {
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<ChunkRecord>)> {
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<ChunkRecord>)> {
let wrap =
|e: io::Error| io::Error::new(e.kind(), format!("reading index {}: {e}", path.display()));
let f = File::open(path)
Expand All @@ -547,6 +563,7 @@ pub fn read_index_tail(path: &Path, from: usize) -> io::Result<(usize, Vec<Chunk
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;
let from = start(n);
if from >= n {
return Ok((n, Vec::new()));
}
Expand Down Expand Up @@ -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);
}
}
14 changes: 12 additions & 2 deletions src/receive.rs
Original file line number Diff line number Diff line change
Expand Up @@ -188,6 +188,7 @@ pub struct Session {
cfg: crate::store::Config,
adopt_pages: bool,
out: Received,
runs: Vec<Run>,
_dir_lock: std::fs::File,
_file_lock: std::fs::File,
}
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -340,8 +348,7 @@ impl Session {

/// What this destination holds now — the ack, and a coverage answer.
pub fn coverage(&self) -> Vec<Run> {
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
Expand Down Expand Up @@ -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)
Expand Down
10 changes: 5 additions & 5 deletions src/rotate.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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))
);
}
}
Expand Down
54 changes: 47 additions & 7 deletions src/serve.rs
Original file line number Diff line number Diff line change
Expand Up @@ -101,15 +101,20 @@ pub struct LastSent {
pub fn runs_of(seqs: impl IntoIterator<Item = u64>) -> Vec<Run> {
let mut out: Vec<Run> = 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<Run>, 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<Served> {
let mut handle = crate::query::open_source(input)?;
Expand Down Expand Up @@ -308,8 +313,10 @@ fn grain_path(input: &Path) -> Option<std::path::PathBuf> {
/// 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<Vec<u8>> {
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)]
Expand Down Expand Up @@ -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<u8> = (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");
}
}
7 changes: 4 additions & 3 deletions src/ship.rs
Original file line number Diff line number Diff line change
Expand Up @@ -100,9 +100,10 @@ fn instant(s: Option<&str>) -> Option<chrono::DateTime<chrono::FixedOffset>> {
/// 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
Expand Down
62 changes: 35 additions & 27 deletions src/sink.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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));
Expand Down Expand Up @@ -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
Expand Down
26 changes: 22 additions & 4 deletions src/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand Down Expand Up @@ -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");
}
}
}
Loading