diff --git a/nodedb-crdt/src/state/document_cell.rs b/nodedb-crdt/src/state/document_cell.rs index 919e589f5..b5a5a3dc2 100644 --- a/nodedb-crdt/src/state/document_cell.rs +++ b/nodedb-crdt/src/state/document_cell.rs @@ -7,6 +7,35 @@ use std::ops::Deref; use loro::LoroDoc; +/// A measured point relating the document's operation count to its encoded +/// size, taken from one real snapshot export. +struct Calibration { + ops: usize, + bytes: usize, +} + +impl Calibration { + /// Encoded size implied by `ops`, holding the measured bytes-per-operation + /// ratio fixed. Widened to `u128` because a large document's + /// `bytes * ops` overflows `usize` on 32-bit targets long before either + /// factor does. + fn scale_to(&self, ops: usize) -> usize { + let scaled = self.bytes as u128 * ops as u128 / self.ops as u128; + usize::try_from(scaled).unwrap_or(usize::MAX) + } + + /// Whether `ops` is close enough to the measured point to interpolate + /// from it. Encoded density is stable within a document — what changes it + /// is a different *kind* of content, which arrives gradually — so the + /// ratio is re-measured once the operation count has halved or doubled. + /// + /// Bounding recalibration to a doubling makes the export cost amortise to + /// O(1) per write, rather than being paid on every write. + fn covers(&self, ops: usize) -> bool { + ops > 0 && ops <= self.ops.saturating_mul(2) && ops.saturating_mul(2) >= self.ops + } +} + /// A `LoroDoc` together with the values cached from it. /// /// Compaction does not mutate a document, it replaces one: `compact_history` @@ -26,9 +55,14 @@ use loro::LoroDoc; /// Reads go through `Deref`, so every `self.doc.…` call site is untouched. pub(in crate::state) struct DocumentCell { doc: LoroDoc, - /// Snapshot size and the oplog version it was measured at. `None` until - /// the estimate is first asked for. - memory_estimate: RefCell>, + /// Last real measurement of encoded size, and the operation count it was + /// taken at. `None` until the estimate is first asked for. + calibration: RefCell>, + /// Real snapshot exports performed to answer `estimated_bytes`. The point + /// of the calibration is that this grows logarithmically with the number + /// of writes, not linearly, which is a property worth asserting. + #[cfg(test)] + exports: std::cell::Cell, } impl DocumentCell { @@ -36,7 +70,9 @@ impl DocumentCell { pub(in crate::state) fn new(doc: LoroDoc) -> Self { Self { doc, - memory_estimate: RefCell::new(None), + calibration: RefCell::new(None), + #[cfg(test)] + exports: std::cell::Cell::new(0), } } @@ -46,33 +82,56 @@ impl DocumentCell { *self = Self::new(doc); } - /// Snapshot size in bytes, as a proxy for memory footprint. + /// Estimated encoded size in bytes, as a proxy for memory footprint. + /// + /// Loro exposes no direct memory metric, and a snapshot export — the + /// honest proxy — costs O(document). Callers put this on the write path + /// (a memory governor updated after every operation), so paying a full + /// re-encode per call means every write re-serialises the whole document: + /// ~100 ms per write on a 4 MB document, and proportionally worse above + /// that. + /// + /// So the export is used to *calibrate* rather than to answer. `len_ops` + /// is an inlined oplog counter, and it counts operations that are still in + /// an open transaction, so it tracks writes the moment they happen. The + /// answer is that counter scaled by the measured bytes-per-operation, and + /// a real export runs only when the count leaves the calibrated range. /// - /// Loro exposes no direct memory metric. A snapshot export is proportional - /// to state size, which is good enough for pressure monitoring but costs - /// O(document) — and the callers that want it are polling loops. It is - /// therefore measured once per version rather than once per call. - /// The version is a sound cache key because `oplog_vv` counts operations - /// that are still in an open transaction — a write is visible to the key - /// the moment it happens, not when it commits. `state::tests` pins that - /// property, since the cache is wrong the day it stops holding. + /// Exact whenever the document has not changed since it was measured; + /// an interpolation otherwise, which is what a pressure signal needs. pub(in crate::state) fn estimated_bytes(&self) -> usize { - let version = self.doc.oplog_vv(); - if let Some((measured_at, bytes)) = self.memory_estimate.borrow().as_ref() - && *measured_at == version - { - return *bytes; + let ops = self.doc.len_ops(); + if let Some(calibration) = self.calibration.borrow().as_ref() { + if ops == calibration.ops { + return calibration.bytes; + } + if calibration.covers(ops) { + return calibration.scale_to(ops); + } } + #[cfg(test)] + self.exports.set(self.exports.get() + 1); let Ok(snapshot) = self.doc.export(loro::ExportMode::Snapshot) else { // A failed export is not a measurement. Caching the zero would pin - // the document at "empty" until its next write moved the version. + // the document at "empty" until enough writes moved it out of + // range again. return 0; }; let bytes = snapshot.len(); - *self.memory_estimate.borrow_mut() = Some((version, bytes)); + // An empty document is not a ratio: it would divide by zero, and the + // export that measured it was trivial anyway. + if ops > 0 { + *self.calibration.borrow_mut() = Some(Calibration { ops, bytes }); + } bytes } + + /// How many real snapshot exports `estimated_bytes` has performed. + #[cfg(test)] + pub(in crate::state) fn export_count(&self) -> usize { + self.exports.get() + } } impl Deref for DocumentCell { diff --git a/nodedb-crdt/src/state/snapshot.rs b/nodedb-crdt/src/state/snapshot.rs index 237e5858a..d6dd7ffe1 100644 --- a/nodedb-crdt/src/state/snapshot.rs +++ b/nodedb-crdt/src/state/snapshot.rs @@ -147,13 +147,24 @@ impl CrdtState { /// Includes operation history, current state, and internal caches. /// Use this to decide when to trigger `compact_history()`. /// - /// The proxy is a full snapshot export, so it is measured once per version - /// rather than once per call: a document nobody is writing cannot have - /// changed size. Compaction discards the measurement along with the - /// document it described. + /// Cheap enough to call on the write path. The underlying proxy is a + /// snapshot export, which costs O(document), so it is used to calibrate a + /// bytes-per-operation ratio and the answer comes from the oplog's + /// operation counter — a real export runs only when the document has + /// halved or doubled since it was last measured. Exact for a document that + /// has not changed since then, an interpolation otherwise. + /// + /// Compaction discards the measurement along with the document it + /// described. pub fn estimated_memory_bytes(&self) -> usize { self.doc.estimated_bytes() } + + /// Real snapshot exports performed to answer `estimated_memory_bytes`. + #[cfg(test)] + pub(crate) fn export_count_for_test(&self) -> usize { + self.doc.export_count() + } } #[cfg(test)] diff --git a/nodedb-crdt/src/state/tests.rs b/nodedb-crdt/src/state/tests.rs index 3bccf4e24..87def6ccd 100644 --- a/nodedb-crdt/src/state/tests.rs +++ b/nodedb-crdt/src/state/tests.rs @@ -221,6 +221,55 @@ fn estimated_memory_follows_compaction() { ); } +#[test] +fn estimated_memory_does_not_re_encode_on_every_write() { + let state = CrdtState::new(1).unwrap(); + for i in 0..2_000 { + state + .upsert("items", &format!("i-{i}"), &[("v", LoroValue::I64(i))]) + .unwrap(); + state.estimated_memory_bytes(); + } + + // The estimate is called after every write by the memory governor. Paying + // a full document re-encode per call is what made a large store take + // ~0.5 s per record; the calibrated ratio has to make that logarithmic in + // the number of writes, not linear. + let exports = state.export_count_for_test(); + assert!( + exports < 32, + "{exports} full exports for 2000 writes — the estimate is re-encoding \ + the document per write again" + ); +} + +#[test] +fn estimated_memory_stays_close_to_the_real_encoded_size() { + let state = CrdtState::new(1).unwrap(); + for i in 0..2_000 { + state + .upsert( + "items", + &format!("i-{i}"), + &[("v", LoroValue::String("payload".repeat(4).into()))], + ) + .unwrap(); + state.estimated_memory_bytes(); + } + + // Interpolating between calibrations trades exactness for cost. It is a + // pressure signal, so it may drift — but it has to stay the same order as + // the truth, or the thresholds built on it mean nothing. + let estimate = state.estimated_memory_bytes() as f64; + let actual = state.export_snapshot().unwrap().len() as f64; + let ratio = estimate / actual; + assert!( + (0.5..=2.0).contains(&ratio), + "estimate {estimate} vs actual {actual} (ratio {ratio:.2}) — drifted \ + beyond what a pressure signal can carry" + ); +} + #[test] fn oplog_version_counts_uncommitted_operations() { let state = CrdtState::new(1).unwrap();