diff --git a/docs/plans/tally-design.md b/docs/plans/tally-design.md index ac4a607..f34521a 100644 --- a/docs/plans/tally-design.md +++ b/docs/plans/tally-design.md @@ -123,6 +123,41 @@ the host reading it. where the documents say one thing and the numbers are another, and nothing in it is wrong enough to notice. +## It is fed sequentially, by a follower + +The fold reads its source from beginning to end, once, in one process — the follower's `--run` consumer. The follower is not +merely the delivery mechanism: it holds the position durably, holds the +source's retention back while the tally is behind it, supplies the registry +and the systemd lifecycle, and restarts into exactly the re-fold `safe_offset` +makes correct. + +⚠ **Parallel workers are not the first thing to reach for**, for two +reasons: + +* **It needs additive partials**, which is their THIRD use after spilling and + repair — two workers can both contribute to one bucket, and with a merge + that replaces a cell the second silently erases the first. And it must + partition by SOURCE POSITION rather than by day: a day's entries are not + contiguous in the source (late arrivals are why a citation exists at all), + and there is no per-chunk logline range to find them by + ([logline-order.md](logline-order.md)). Partitioning by chunk range falls + out for free, a worker's consumed range being the partial identity already + wanted. +* **The gain is small.** Measured over 300,000 real lines: reading records + 0.52 s, the fold 4.53 s, writing blocks 0.05 s — so blocks are 1% and the + fold is 90%. A day is ~47 s single-threaded, and the largest backfill that + can ever be asked for is the SOURCE's retention (weeks, not years, since + nothing can tally what was dropped), so a 30-day rebuild is ~24 minutes + once. + +⚠ **Nor is sharing the extraction, obvious as it looks.** The four metrics of +the measured document each run their own extract regex over every line, and +cost is spread evenly +and the metric with the SMALLEST regex (17 characters) is joint-most expensive +at 1.31 s, because it is a histogram of 22 buckets and turns one line into 22 +samples. The cost is per sample produced and per metric evaluated, not per +regex byte, so sharing the extraction wins much less than it looks. + ## Recovery is re-derivation A follower's position is precious because it shipped bytes it cannot un-ship. @@ -135,6 +170,29 @@ and not a correctness one. nothing else. Nothing consults a watermark, so a re-derivation does not depend on where the read started. +⚠ **And a record stream is forward-only, which is fine for a CRASH and not +for a repair.** The tally is fed a stream and cannot ask for bytes again, so +the two cases separate: + +* **A crash needs no seek**, because of `Roller::safe_offset` — "the oldest + source byte any OPEN bucket still depends on. A consumer may not report past + this". The position is therefore always behind every unfinished bucket, so a + restart re-sends from before that bucket's FIRST entry, re-folds it whole, + and the complete total replaces the partial one. ⚠ **The block writer's + `How::Merge` depends on this, invisibly**: advance `safe_offset` any further + and blocks are corrupted by code that never mentions them. ⚠ That dependency + is a consequence of a merge REPLACING a cell — under additive partials + ([tally-partials.md](tally-partials.md)) the position may advance freely. +* **A repair does need to go back, and a REWIND is the wrong way to do it.** + A tally is not only a consumer: `query --records --from X --to Y | tally` + is a bounded DIRECT read of the source, and it already exists. So "the regex + was wrong, recompute last week" is a one-shot pass over that window while + the follower keeps going forward — no position is moved, and the live path + is never interrupted. ⚠ What it needs is `How::Regenerate` rather than + `Merge`, which the writer does not expose: merging a recomputation into + what is there would replace cell by cell and leave any series the new + definition no longer produces standing. + ## Copying is a file sync or a bundle Blocks that are immutable once past the floor, under a manifest with a crc32 @@ -188,7 +246,12 @@ a working set that saturates rather than drifting. for a partial — and it is also the atomicity defect below; * **compaction's schedule**, and whether a query merges or refuses; * **the write-batching mechanism** — a WAL for samples in the `.sap` shape is - the candidate, since one late sample otherwise rewrites a whole day block. + the candidate, since one late sample otherwise rewrites a whole day block; +* **a repair pass** — a bounded read with `How::Regenerate`, which is what + makes "fix the definition and recompute" an operation rather than a plan. + ⚠ Not a rewind of the follower: the direct read already exists, and moving + a live position to recompute history would stop the live path to fix the + past. ⚠ **A defect that exists today:** `commit` renames a block into place and THEN saves the manifest, so a rewrite at the same `(t0, generation)` leaves a window diff --git a/packaging/timberfs-completion.bash b/packaging/timberfs-completion.bash index a7aab1e..41a887a 100644 --- a/packaging/timberfs-completion.bash +++ b/packaging/timberfs-completion.bash @@ -110,6 +110,7 @@ _timberfs() { --deadline | --positions | --batch-size | --follow-from | --delete-empty | \ --look-in | --etc | --extractor | --metric | --width | --grace | \ --provision | --run | --pack | --unpack | --block-buckets | \ + --blocks | --block-flush | \ --query | --series | --since | --until) COMPREPLY=($(compgen -f -- "$cur")) return 0 diff --git a/packaging/timberfs.1 b/packaging/timberfs.1 index 8ae244b..ac03575 100644 --- a/packaging/timberfs.1 +++ b/packaging/timberfs.1 @@ -1332,6 +1332,42 @@ shows the tape as written, both lines included, which is what a tape viewer should do. See LATENESS below for when a provisional line is written. .PP +.B WRITING THE GRID (EXPERIMENTAL, A HARNESS) +.PP +.BI \-\-blocks " DIR" +writes the numbers as columnar BLOCKS into DIR instead of as tally lines +on stdout \(em the storage the design settles on, where a line is the +INTERCHANGE form and not the store. Verified against 300,000 real log +lines: the blocks render back to exactly what the line path wrote. +.IP +⚠ A HARNESS rather than an interface. Nothing in a provisioned deployment +runs this: it is how the writer is exercised against a real store until +.B \-\-run +can write blocks itself. And it names a DIRECTORY where every other +timberfs argument names a store: the manifest carries an id, but nothing +searches for block stores by it yet, so a path is the only address there +is. +.IP +.BI \-\-block\-flush " N" +is how many samples are buffered before a commit (default 50000). ⚠ It is +a WRITE\-AMPLIFICATION control and not a latency one: a block is a day, so +each commit rewrites up to a megabyte, and one real day of a busy store is +~268,000 samples. Measured on 300,000 lines, the same numbers cost 43, 5 +or 1 block writes at 1000, 10000 and 50000. The cost of buffering is +VISIBILITY \(em a sample is not in a block until it is flushed \(em and the +cost of losing the buffer to a crash is a re\-read, the position not +having moved. The durable form of it would be a write\-ahead log for +samples. +.IP +⚠ ONE store in, one directory out, which is why it is here and not on +.BR \-\-run : +a provisioned run serves a SELECTION, one sink per source store, and +where each one's blocks go is a provisioning question rather than a +writer one. +.IP +⚠ Nothing is written to stdout, the unit announcement included: a block +store carries its own definitions, and a unit is a property of one. +.PP .B PACKING THE GRID (EXPERIMENTAL, A MEASUREMENT) .PP .BI \-\-pack " DIR" diff --git a/src/main.rs b/src/main.rs index 5a97677..1c7c4f8 100644 --- a/src/main.rs +++ b/src/main.rs @@ -667,6 +667,28 @@ enum Command { /// format --fold takes #[arg(long)] observations: bool, + /// EXPERIMENTAL, and a HARNESS rather than an interface: write + /// the numbers as columnar BLOCKS into DIR instead of as tally + /// lines on stdout, which is how the block writer is exercised + /// against a real store until a provisioned `--run` can write + /// them. + /// + /// ⚠ It names a DIRECTORY where every other timberfs argument + /// names a store: the manifest carries an id, but nothing + /// searches for block stores by it, so a path is the only + /// address there is. + /// + /// ⚠ One store in, one directory out. Not on `--provision`'s + /// `--run`, which serves a SELECTION with a sink per source + /// store: where each one's blocks go is a provisioning question. + #[arg(long, value_name = "DIR", conflicts_with_all = ["fold", "provision", "run", "try_it", "check", "pack", "unpack", "query", "observations"])] + blocks: Option, + /// With --blocks: samples buffered before a commit. A block is a + /// day, so each commit rewrites up to a megabyte — this is the + /// write-amplification control, and a sample is not in a block + /// until it is flushed + #[arg(long, value_name = "N", default_value_t = timberfs::tally::DEFAULT_BLOCK_FLUSH, requires = "blocks")] + block_flush: usize, /// EXPERIMENTAL, and a measurement rather than a feature: read /// tally lines on stdin and write them as columnar BLOCKS into /// DIR — the grid of series x buckets, one file per range, @@ -1910,6 +1932,8 @@ fn main() -> anyhow::Result<()> { grain::cmd_reindex(&file)?; } Command::Tally { + blocks, + block_flush, extractors, try_it, check, @@ -1958,6 +1982,9 @@ fn main() -> anyhow::Result<()> { .map(append::parse_duration_ms) .transpose()?; tally::cmd_tally(&tally::TallyOpts { + blocks, + block_flush, + block_buckets, extractors, etc, try_it, diff --git a/src/tally.rs b/src/tally.rs index 5864a87..f2bf23e 100644 --- a/src/tally.rs +++ b/src/tally.rs @@ -1391,6 +1391,19 @@ pub struct TallyOpts { pub metrics: Vec, /// Override the extractors' own window. pub width_ms: Option, + /// Write the numbers as columnar BLOCKS into this directory instead + /// of as lines on stdout (docs/plans/tally-design.md). + /// + /// ⚠ One store in, one directory out, which is why it is on this + /// path and not on `--run`: a provisioned run serves a SELECTION, + /// one sink per source store, and where each one's blocks go is a + /// provisioning question rather than a writer one. + pub blocks: Option, + /// Samples buffered before a commit. The write-amplification + /// control: a block is a day, so a commit rewrites up to a + /// megabyte. + pub block_flush: usize, + pub block_buckets: usize, } #[derive(Clone, Copy, Debug)] @@ -1436,9 +1449,9 @@ pub fn fold_stream( ); } roller.add(&s); - let _ = emit(out, roller.drain(Drain::Sealed))?; + let _ = emit(out, None, roller.drain(Drain::Sealed))?; } - let _ = emit(out, roller.drain(Drain::Final))?; + let _ = emit(out, None, roller.drain(Drain::Final))?; Ok(()) } @@ -1675,10 +1688,12 @@ impl Run { e: &crate::records::EntryRec, axis: Axis, out: &mut impl Write, + blocks: Option<&mut crate::tally_block::Writer>, observations: bool, ) -> anyhow::Result> { let mut batch: Vec = Vec::new(); let mut meta: Option<(u64, u64)> = None; + let to_blocks = blocks.is_some(); for l in self.live.iter_mut() { let ts = match axis { Axis::Logline => e.ts, @@ -1706,7 +1721,12 @@ impl Run { // every line: the tape outlives the document that // described it, and a number whose unit is unknown is // a number nobody can act on. - if !l.announced { + // ⚠ Announced only on the LINE path. A block store + // carries the definitions, and a unit is a property + // of a definition — so writing it here would be the + // same fact in two places, one of them a marker the + // grid does not hold (docs/plans/tally-design.md). + if !l.announced && !to_blocks { l.announced = true; if let Some(u) = &l.unit { meta = span(meta, Some((ts, ts))); @@ -1733,7 +1753,7 @@ impl Run { batch.extend(l.roller.drain(Drain::Sealed)); } } - Ok(span(meta, emit(out, batch)?)) + Ok(span(meta, emit(out, blocks, batch)?)) } /// A quiet stream: let wall clock stand in for event time, write @@ -1746,7 +1766,19 @@ impl Run { /// has not reached — and every entry arriving after that is then /// displaced out of its own bucket. So the newest bucket is shown /// rather than sealed, and superseded when it is complete. - fn idle(&mut self, now_ms: u64, out: &mut impl Write) -> anyhow::Result> { + /// + /// ⚠ Into BLOCKS this rests on `Block::merge` REPLACING a cell: a + /// provisional bucket states its total so far, and the complete one + /// later states the whole total again. When merge becomes additive + /// for partials (docs/plans/tally-partials.md) this double-counts, + /// and the third identity component that note asks for is what + /// stops it. + fn idle( + &mut self, + now_ms: u64, + out: &mut impl Write, + blocks: Option<&mut crate::tally_block::Writer>, + ) -> anyhow::Result> { let mut batch: Vec = Vec::new(); for l in self.live.iter_mut() { if l.roller.watermark() == 0 { @@ -1757,12 +1789,13 @@ impl Run { l.roller.advance_to(now_ms); batch.extend(l.roller.drain(Drain::Provisional)); } - emit(out, batch) + emit(out, blocks, batch) } fn finish( &mut self, out: &mut impl Write, + blocks: Option<&mut crate::tally_block::Writer>, observations: bool, ) -> anyhow::Result> { if observations { @@ -1772,7 +1805,7 @@ impl Run { for l in self.live.iter_mut() { batch.extend(l.roller.drain(Drain::Final)); } - emit(out, batch) + emit(out, blocks, batch) } } @@ -1838,20 +1871,59 @@ pub fn cmd_tally(opts: &TallyOpts) -> anyhow::Result<()> { return Ok(()); } + let mut blocks = match &opts.blocks { + None => None, + Some(dir) => { + if opts.observations { + bail!("--observations prints width-0s observations, which are not a grid"); + } + Some(crate::tally_block::Writer::open( + dir, + width_of(&docs, opts)?, + opts.block_buckets, + opts.block_flush, + 3, + )?) + } + }; let stdin = std::io::stdin(); let mut reader = crate::records::Reader::new(stdin.lock()); let mut ended = false; while let Some(rec) = reader.next_rec()? { match rec { crate::records::Rec::Entry(e) => { - run.feed(&e, axis, &mut out, opts.observations)?; + run.feed(&e, axis, &mut out, blocks.as_mut(), opts.observations)?; + } + // ⚠ The stream says which store it came from, and a block + // store records it: a citation is an offset into ONE tape. + // A record without an id came from no store (a pipe), so + // there is nothing to attribute. + crate::records::Rec::Source(fields) => { + if let Some(w) = blocks.as_mut() { + if let Some((_, id)) = fields.iter().find(|(k, _)| k == "id") { + w.source_is(id)?; + } + } } crate::records::Rec::End(_) => ended = true, _ => {} } } - run.finish(&mut out, opts.observations)?; + run.finish(&mut out, blocks.as_mut(), opts.observations)?; out.flush()?; + if let Some(w) = &mut blocks { + w.flush()?; + crate::note!( + "timberfs: {} sample(s) into {} block write(s){}", + w.committed, + w.blocks_written, + if w.markers > 0 { + format!("; {} marker(s) not stored", w.markers) + } else { + String::new() + } + ); + } if !ended { bail!("the record stream ended without stream-end — the answer is truncated"); } @@ -1885,9 +1957,9 @@ pub fn try_text( let entries = entries_of_text(text)?; let mut tally: Vec = Vec::new(); for e in &entries { - run.feed(e, axis, &mut tally, opts.observations)?; + run.feed(e, axis, &mut tally, None, opts.observations)?; } - run.finish(&mut tally, opts.observations)?; + run.finish(&mut tally, None, opts.observations)?; Ok(Tried { entries: entries.len(), tally: String::from_utf8_lossy(&tally).into_owned(), @@ -1919,6 +1991,44 @@ fn axis_of(docs: &[(PathBuf, Extractor)]) -> anyhow::Result { Ok(first.window.axis) } +/// Samples buffered before a commit, and buckets per block. +/// +/// ⚠ The flush default is a WRITE-AMPLIFICATION choice, not a latency +/// one: a block is a day, so each commit rewrites up to a megabyte, and +/// a real day of a busy store is ~268,000 samples — so this is a +/// handful of rewrites a day rather than one per batch. A sample is not +/// in a block until it is flushed, which is the cost. +pub const DEFAULT_BLOCK_FLUSH: usize = 50_000; +/// A day at 60s, which two measurements settled +/// (docs/plans/tally-as-a-tally.md). +pub const DEFAULT_BLOCK_BUCKETS: usize = 1440; + +/// The one bucket width a run makes, which a block store must hold. +/// +/// ⚠ Refused rather than reconciled when the documents disagree: a +/// store holds one width, and combining two would be coarsening one of +/// them silently. +fn width_of(docs: &[(PathBuf, Extractor)], opts: &TallyOpts) -> anyhow::Result { + if let Some(w) = opts.width_ms { + return Ok(w); + } + let mut it = docs.iter(); + let (first_path, first) = it.next().expect("load_extractors refuses an empty set"); + for (path, doc) in it { + if doc.window.width_ms != first.window.width_ms { + bail!( + "{} buckets at {} and {} at {} — one block store holds one width, \ + so give --width to pick it", + first_path.display(), + render_width(first.window.width_ms), + path.display(), + render_width(doc.window.width_ms) + ); + } + } + Ok(first.window.width_ms) +} + /// Plain text into the SAME entries a store would yield: the real /// assembly, so a `--try` run and a live run cannot disagree about where /// one entry ends and the next begins. @@ -1958,14 +2068,27 @@ fn entries_of_text(text: &[u8]) -> anyhow::Result> /// /// Reports the BUCKET WINDOW it wrote, which is what stamps the chunk: /// a tally store's write axis is the minutes its lines are about. -fn emit(out: &mut impl Write, mut batch: Vec) -> anyhow::Result> { +fn emit( + out: &mut impl Write, + blocks: Option<&mut crate::tally_block::Writer>, + mut batch: Vec, +) -> anyhow::Result> { batch.sort_by(|a, b| (a.ts, &a.metric, &a.labels).cmp(&(b.ts, &b.metric, &b.labels))); let window = match (batch.first(), batch.last()) { (Some(f), Some(l)) => Some((f.ts, l.ts)), _ => None, }; - for s in batch { - writeln!(out, "{}", s.render())?; + // ⚠ One destination or the other, never both: the position this + // returns is what advances a consumer, and a sample counted twice + // by a downstream that reads the lines AND the blocks would be + // double-counted rather than merged. + match blocks { + Some(w) => w.take(batch)?, + None => { + for s in batch { + writeln!(out, "{}", s.render())?; + } + } } Ok(window) } @@ -2611,6 +2734,11 @@ impl TallyOpts { metrics: Vec::new(), width_ms, fold: None, + // A provisioned run writes its sinks' tapes; where each + // sink's blocks would go is the provisioning's question. + blocks: None, + block_flush: DEFAULT_BLOCK_FLUSH, + block_buckets: DEFAULT_BLOCK_BUCKETS, } } } @@ -2687,7 +2815,7 @@ pub fn cmd_run(set: &str, opts: &RunOpts) -> anyhow::Result<()> { let now = crate::store::now_ms(); for (id, sink) in sinks.iter_mut() { let mut lines: Vec = Vec::new(); - let window = sink.run.idle(now, &mut lines)?; + let window = sink.run.idle(now, &mut lines, None)?; sink.write(&lines, window)?; report(&mut reports, id, sink)?; } @@ -2731,7 +2859,7 @@ pub fn cmd_run(set: &str, opts: &RunOpts) -> anyhow::Result<()> { } let sink = sinks.get_mut(&id).expect("just inserted"); let mut lines: Vec = Vec::new(); - let window = sink.run.feed(&e, axis, &mut lines, false)?; + let window = sink.run.feed(&e, axis, &mut lines, None, false)?; sink.write(&lines, window)?; if let Some(off) = e.offset { sink.delivered_to = Some(off + e.payload.len() as u64); @@ -2747,7 +2875,7 @@ pub fn cmd_run(set: &str, opts: &RunOpts) -> anyhow::Result<()> { // no half-counted minute behind. for (id, sink) in sinks.iter_mut() { let mut lines: Vec = Vec::new(); - let window = sink.run.finish(&mut lines, false)?; + let window = sink.run.finish(&mut lines, None, false)?; sink.write(&lines, window)?; let cfg = sink.store.cfg; if let Some(f) = sink.store.files.get_mut(&sink.name) { @@ -3535,6 +3663,9 @@ mod tests { etc: PathBuf::from("/etc/timberfs"), try_it: true, check: false, + blocks: None, + block_flush: DEFAULT_BLOCK_FLUSH, + block_buckets: DEFAULT_BLOCK_BUCKETS, observations: false, metrics: Vec::new(), width_ms: None, @@ -3694,7 +3825,9 @@ mod tests { // stamp — so the window must cover it or the chunk holding it is // stamped from somewhere else entirely. let first = entries_of_text(b"2026-09-06T13:37:10.000Z hello\n").unwrap(); - let meta = run.feed(&first[0], Axis::Logline, &mut out, false).unwrap(); + let meta = run + .feed(&first[0], Axis::Logline, &mut out, None, false) + .unwrap(); assert_eq!( meta, Some(( @@ -3706,8 +3839,9 @@ mod tests { // Two buckets seal at the end: the window is the first and the // last, not the moment they were folded. let later = entries_of_text(b"2026-09-06T13:39:10.000Z hello\n").unwrap(); - run.feed(&later[0], Axis::Logline, &mut out, false).unwrap(); - let sealed = run.finish(&mut out, false).unwrap(); + run.feed(&later[0], Axis::Logline, &mut out, None, false) + .unwrap(); + let sealed = run.finish(&mut out, None, false).unwrap(); assert_eq!( sealed, Some(( diff --git a/src/tally_block.rs b/src/tally_block.rs index 1270b9a..2d80fbc 100644 --- a/src/tally_block.rs +++ b/src/tally_block.rs @@ -710,6 +710,45 @@ impl Entry { #[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)] pub struct Manifest { pub v: u32, + /// This store's own identity, minted once and never touched after — + /// bark's rule, and for bark's reason: a path is an address and the + /// id is what the store IS, across renames, moves and copies. + /// + /// ⚠ `Option` because a manifest written before identity existed has + /// none, and minting one on read would hand an old store a new + /// identity in silence. Missing is a defect to repair deliberately, + /// which is how `timberfs identity` treats the same absence. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub id: Option, + /// When that identity was established, RFC3339. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub created: Option, + /// The SOURCE store's id, which is load-bearing rather than + /// provenance: a citation is an offset into that store's tape and + /// means nothing without knowing which tape. Absent for a store + /// packed from tally lines, which came from no store. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub source: Option, + /// The source's labels, copied when this store was created. + /// + /// ⚠ A record of what the source said THEN, never a live view: the + /// source may be relabelled, renamed or deleted long before a + /// two-year tally, and when it is deleted this is the only surviving + /// witness of what was measured. + #[serde(default, skip_serializing_if = "serde_json::Map::is_empty")] + pub labels: serde_json::Map, + /// Retention: how long, and how much. + /// + /// ⚠ Numbers where a `.bark` holds the strings an operator wrote, + /// because that file holds a DECLARATION and this holds the resolved + /// policy. And `retain_ms` is checked to be a whole number of BLOCKS + /// (`set_retain`), retention dropping whole blocks and nothing + /// finer — so a value between two of them would be a rounding + /// dressed as a setting. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub retain_ms: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub retain_bytes: Option, pub width_ms: u64, pub block_buckets: usize, /// The oldest bucket this store still claims to know about. @@ -727,6 +766,12 @@ impl Manifest { pub fn new(width_ms: u64, block_buckets: usize) -> Manifest { Manifest { v: 1, + id: crate::bark::new_uuid().ok(), + created: Some(chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true)), + source: None, + labels: serde_json::Map::new(), + retain_ms: None, + retain_bytes: None, width_ms, block_buckets, floor: 0, @@ -734,6 +779,50 @@ impl Manifest { } } + /// How long one block covers, which is retention's granularity. + pub fn block_span_ms(&self) -> u64 { + self.width_ms * self.block_buckets as u64 + } + + /// Declare how long to keep. ⚠ Refused unless it is a whole number + /// of blocks: `drop_before` removes blocks whole, so a finer value + /// would be a promise this store cannot keep. + pub fn set_retain(&mut self, ms: Option) -> anyhow::Result<()> { + if let Some(ms) = ms { + let span = self.block_span_ms(); + if span == 0 || ms == 0 || ms % span != 0 { + bail!( + "a retention of {} is not a whole number of blocks — this store's block \ + covers {}, and retention drops blocks whole", + crate::tally::render_width(ms), + crate::tally::render_width(span) + ); + } + } + self.retain_ms = ms; + Ok(()) + } + + /// Record which store the numbers came from. + /// + /// ⚠ Refused when it disagrees with what is recorded: a store's + /// source does not change, and feeding one store's records into + /// another's blocks would leave every citation pointing into the + /// wrong tape — silently, the offsets being plausible either way. + pub fn set_source(&mut self, id: &str) -> anyhow::Result { + match &self.source { + Some(had) if had == id => Ok(false), + Some(had) => bail!( + "these blocks were derived from store {had} and this run reads {id} — a \ + citation is an offset into ONE tape" + ), + None => { + self.source = Some(id.to_string()); + Ok(true) + } + } + } + /// Read it, or `None` where a store has none yet. pub fn load(dir: &std::path::Path) -> anyhow::Result> { let path = dir.join(MANIFEST); @@ -1091,6 +1180,114 @@ pub fn query( Ok((out, t)) } +/// Writes sealed samples into a store's blocks as they arrive. +/// +/// ⚠ It BUFFERS, and the buffer is the write-amplification control. A +/// block is a day, so committing one sample rewrites up to a megabyte; +/// buffering `limit` samples makes that once per `limit` instead. Two +/// costs: a sample is not in a block until it is flushed, and the +/// buffer is not durable — losing it to a crash costs a re-read, the +/// position not having moved. +pub struct Writer { + dir: std::path::PathBuf, + m: Manifest, + block_buckets: usize, + level: i32, + limit: usize, + held: Vec, + pub committed: usize, + pub blocks_written: usize, + pub markers: usize, +} + +impl Writer { + /// Open a store to write into, creating its manifest if there is + /// none. ⚠ The width is the WINDOW's, and a manifest that disagrees + /// is refused rather than adopted: a store holds one bucket width, + /// and coarsening is a different operation on a column. + pub fn open( + dir: &std::path::Path, + width_ms: u64, + block_buckets: usize, + limit: usize, + level: i32, + ) -> anyhow::Result { + if block_buckets == 0 { + bail!("a block spans at least one bucket"); + } + std::fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?; + let m = match Manifest::load(dir)? { + Some(m) => { + if m.width_ms != width_ms { + bail!( + "{} holds {} buckets and this run makes {}", + dir.display(), + crate::tally::render_width(m.width_ms), + crate::tally::render_width(width_ms) + ); + } + m + } + None => Manifest::new(width_ms, block_buckets), + }; + Ok(Writer { + dir: dir.to_path_buf(), + m, + block_buckets, + level, + limit: limit.max(1), + held: Vec::new(), + committed: 0, + blocks_written: 0, + markers: 0, + }) + } + + /// Record which store these numbers come from, refusing a second + /// one — see `Manifest::set_source`. It reaches disk with the next + /// commit, the manifest being the commit point: a run that learns a + /// source and writes no block has nothing to attribute. + pub fn source_is(&mut self, id: &str) -> anyhow::Result<()> { + self.m.set_source(id)?; + Ok(()) + } + + /// Take a batch. Markers are counted and not stored — a marker + /// states something about a RUN and the grid holds numbers + /// (docs/plans/tally-design.md). + pub fn take(&mut self, batch: Vec) -> anyhow::Result<()> { + for s in batch { + if s.is_marker() { + self.markers += 1; + } else { + self.held.push(s); + } + } + if self.held.len() >= self.limit { + self.flush()?; + } + Ok(()) + } + + /// Commit what is held. ⚠ `How::Merge`, always: a range being + /// filled is one derivation arriving in pieces, and inferring a + /// regeneration from having seen the range before is the defect that + /// silently replaced three of six real days. + pub fn flush(&mut self) -> anyhow::Result<()> { + if self.held.is_empty() { + return Ok(()); + } + let held = std::mem::take(&mut self.held); + let n = held.len(); + for b in Block::pack(&held, self.block_buckets)? { + commit(&self.dir, &mut self.m, &b, self.level, How::Merge)?; + self.blocks_written += 1; + } + self.committed += n; + Ok(()) + } +} + // ------------------------------------------------------------- the verbs /// A block's file name: sortable by time, generation visible, derivable @@ -2002,6 +2199,158 @@ mod tests { assert!(a < b, "directory order must be time order"); } + /// A store has an identity, minted once, and reloading never + /// touches it — bark's rule, for bark's reason: a path is an + /// address and the id is what the store IS. + #[test] + fn a_store_is_minted_an_identity_that_a_reload_never_changes() { + let dir = tmpdir("id1"); + let m = Manifest::new(60_000, 60); + let id = m.id.clone().expect("minted"); + let created = m.created.clone().expect("stamped"); + m.save(&dir).unwrap(); + let back = Manifest::load(&dir).unwrap().unwrap(); + assert_eq!(back.id.as_deref(), Some(id.as_str())); + assert_eq!(back.created, Some(created)); + // And a second store is a different store. + assert_ne!(Manifest::new(60_000, 60).id, Some(id)); + let _ = std::fs::remove_dir_all(&dir); + } + + /// ⚠ A manifest written before identity existed loads, and is NOT + /// handed one on the way in: minting on read would give an old + /// store a new identity in silence, where absence is a defect to + /// repair deliberately. + #[test] + fn a_manifest_without_an_identity_is_not_given_one_by_reading_it() { + let dir = tmpdir("id2"); + std::fs::create_dir_all(&dir).unwrap(); + std::fs::write( + dir.join(MANIFEST), + r#"{"v":1,"width_ms":60000,"block_buckets":60,"floor":0,"blocks":[]}"#, + ) + .unwrap(); + let m = Manifest::load(&dir).unwrap().unwrap(); + assert_eq!(m.id, None, "reading minted an identity"); + assert_eq!(m.source, None); + assert!(m.labels.is_empty()); + let _ = std::fs::remove_dir_all(&dir); + } + + /// A store's source does not change. Feeding one store's records + /// into another's blocks would leave every citation pointing into + /// the wrong tape, and the offsets are plausible either way — so it + /// is refused rather than noticed later. + #[test] + fn a_second_source_for_one_block_store_is_refused() { + let mut m = Manifest::new(60_000, 60); + assert!(m.set_source("aaaa-1111").unwrap(), "the first is recorded"); + assert!( + !m.set_source("aaaa-1111").unwrap(), + "the same one again is a no-op" + ); + let err = m.set_source("bbbb-2222").unwrap_err().to_string(); + assert!( + err.contains("aaaa-1111") && err.contains("bbbb-2222"), + "{err}" + ); + assert_eq!( + m.source.as_deref(), + Some("aaaa-1111"), + "and it did not move" + ); + } + + /// Retention drops whole blocks, so a value between two of them is + /// a promise the store cannot keep — refused rather than rounded. + #[test] + fn a_retention_finer_than_a_block_is_refused() { + let mut m = Manifest::new(60_000, 1440); + assert_eq!(m.block_span_ms(), 86_400_000, "a day of 60s buckets"); + m.set_retain(Some(730 * 86_400_000)).unwrap(); + assert_eq!(m.retain_ms, Some(730 * 86_400_000)); + let err = m.set_retain(Some(36 * 3_600_000)).unwrap_err().to_string(); + assert!(err.contains("whole number of blocks"), "{err}"); + assert_eq!( + m.retain_ms, + Some(730 * 86_400_000), + "a refusal must not have moved it" + ); + m.set_retain(None).unwrap(); + assert_eq!(m.retain_ms, None, "and it can be cleared"); + } + + /// What the writer stores is what the lines would have said. If this + /// does not hold, the block path is a different tally rather than + /// the same one stored differently — which is the whole claim. + /// ⚠ A flush limit of 2 forces several commits over one range, so + /// this also exercises the merge every commit after the first is. + #[test] + fn what_the_writer_stores_is_what_the_lines_said() { + let dir = tmpdir("w1"); + let lines = [ + "2026-09-06T13:37:00.000Z 60s m a=1 count=1 sum=10", + "2026-09-06T13:38:00.000Z 60s m a=1 count=2 sum=20", + "2026-09-06T13:38:00.000Z 60s m a=2 count=3 sum=30", + "2026-09-06T13:39:00.000Z 60s other b=1 count=4", + ]; + let mut w = Writer::open(&dir, 60_000, 60, 2, 3).unwrap(); + for l in lines { + w.take(vec![s(l)]).unwrap(); + } + w.flush().unwrap(); + assert_eq!(w.committed, 4); + assert!(w.blocks_written > 1, "the limit should have forced commits"); + + let (got, _) = query(&dir, &sel("[]"), 0, u64::MAX).unwrap(); + assert_eq!( + got.iter().map(|x| x.render()).collect::>(), + lines, + "the grid must render back to the lines it was given" + ); + let _ = std::fs::remove_dir_all(&dir); + } + + /// Markers are counted and not stored: a marker states something + /// about a RUN and the grid holds numbers + /// (docs/plans/tally-design.md). + #[test] + fn the_writer_counts_markers_and_stores_none() { + let dir = tmpdir("w2"); + let mut w = Writer::open(&dir, 60_000, 60, 100, 3).unwrap(); + w.take(vec![ + s("2026-09-06T13:37:00.000Z 60s m a=1 count=1"), + s("2026-09-06T13:37:00.000Z 60s !drop metric=m reason=unreadable count=2"), + s("2026-09-06T13:37:00.000Z 0s !meta metric=m unit=calls"), + ]) + .unwrap(); + w.flush().unwrap(); + assert_eq!((w.committed, w.markers), (1, 2)); + let (got, _) = query(&dir, &sel("[]"), 0, u64::MAX).unwrap(); + assert_eq!(got.len(), 1); + assert_eq!(got[0].metric, "m"); + let _ = std::fs::remove_dir_all(&dir); + } + + /// A store holds ONE bucket width, so reopening it with another is + /// refused rather than adopted — coarsening is a different + /// operation on a column, not something a second run may do by + /// arriving with a different window. + #[test] + fn a_writer_will_not_change_a_stores_width() { + let dir = tmpdir("w3"); + let mut w = Writer::open(&dir, 60_000, 60, 100, 3).unwrap(); + w.take(vec![s("2026-09-06T13:37:00.000Z 60s m a=1 count=1")]) + .unwrap(); + w.flush().unwrap(); + let err = match Writer::open(&dir, 300_000, 60, 100, 3) { + Err(e) => e.to_string(), + Ok(_) => panic!("a second width was accepted"), + }; + assert!(err.contains("60s") && err.contains("300s"), "{err}"); + let _ = std::fs::remove_dir_all(&dir); + } + /// Garbage in must not panic. Defence in depth rather than a /// scenario: every production path reaches `decode` through /// `read_block`, which checks the manifest's crc32 first, so this