From 8fca3ee63f7ac5e55bbaa84528f3807e9f98aadc Mon Sep 17 00:00:00 2001 From: Alexy Khrabrov Date: Sun, 20 Sep 2026 00:50:53 -0700 Subject: [PATCH] Make memory accounting lock-free docs/lock-free.md, L1 to L4. live_bytes and peak_bytes are atomics admitted by the same compare-exchange as work units; only cancellation wakers keep the mutex. Co-Authored-By: Claude Fable 5.1 --- CHANGELOG.md | 13 ++ crates/grust-procedures/src/resources.rs | 112 ++++++++------- crates/grust-procedures/tests/contracts.rs | 2 + .../tests/contracts/memory.rs | 134 ++++++++++++++++++ docs/book/chapters/generalized-algorithms.md | 25 ++-- docs/lock-free.md | 44 +++++- 6 files changed, 268 insertions(+), 62 deletions(-) create mode 100644 crates/grust-procedures/tests/contracts/memory.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index 416ac071..c553639d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,19 @@ reconstructed from Git history, release commits, and the shipped docs. ## Unreleased +- **Memory accounting no longer takes a lock.** `ExecutionContext` holds + accounted bytes and their high-water mark as atomics and admits a memory + charge through the same compare-exchange as a work charge, so a byte limit + stays exact under concurrent charges; a release is one atomic subtraction. + Only cancellation wakers remain behind the mutex. The reference Cypher + executor charges memory per copied value, and on the full-path `reduce` query + this removes 11% at 1,024 nodes and 13% at 4,096 on one laptop; the fused + `UNWIND` form and the direct kernels do not move. Three behaviours change: + `usage()` reads its figures in sequence rather than as one snapshot, so exact + totals should be read after execution (`peak_bytes` is never reported below + `live_bytes`); the peak can trail a charge by an instant but never misses a + completed one; and a poisoned lock can no longer fail a memory charge, only + waker registration. See `docs/lock-free.md`. - The sequential kernels — `dijkstra`, `shortestPaths`, `bfs`, `dfs`, `topologicalSort`, `scc`, `wcc` — and the heap the path kernels share charge work through a `WorkMeter` instead of one compare-exchange on the shared diff --git a/crates/grust-procedures/src/resources.rs b/crates/grust-procedures/src/resources.rs index 8e0e8344..ab5e278a 100644 --- a/crates/grust-procedures/src/resources.rs +++ b/crates/grust-procedures/src/resources.rs @@ -35,8 +35,6 @@ pub struct ResourceUsage { #[derive(Debug, Default)] struct State { - live_bytes: usize, - peak_bytes: usize, waiters: Vec>, } @@ -47,9 +45,10 @@ struct State { /// the clock, so a caller that needs an exact poll has one. const DEADLINE_SAMPLE_UNITS: usize = 1024; -/// Work and cancellation are lock-free because kernels charge per unit of work, -/// once per visited entry or reconstructed path step. Memory reservations and -/// wakers stay behind the mutex: they are rare and need multi-field atomicity. +/// Work, memory and cancellation are lock-free: kernels charge work once per +/// visited entry, and the reference Cypher executor charges memory once per +/// copied value. Only the cancellation wakers stay behind the mutex; they are +/// rare, and they are the one place that needs several fields to move together. #[derive(Debug)] struct Shared { limits: ExecutionLimits, @@ -62,6 +61,14 @@ struct Shared { /// admitting a block holds it shared; creating, dropping and refusing hold /// it exclusively. See [`WorkMeter`]'s admission for why. grants: RwLock>>, + /// Accounted memory now, and its high-water mark. Atomics, like + /// `work_units`: the reference Cypher executor charges the logical bytes of + /// every copied value, so a streaming query charges memory per element and + /// a lock here was a tenth of its profile. Nothing waits for memory to be + /// freed, so a release needs no waker and no lock either. + live_bytes: AtomicUsize, + peak_bytes: AtomicUsize, + /// Cancellation wakers only. state: Mutex, } @@ -104,6 +111,8 @@ impl ExecutionContext { cancelled: AtomicBool::new(false), charges_since_deadline_read: AtomicUsize::new(0), grants: RwLock::new(Vec::new()), + live_bytes: AtomicUsize::new(0), + peak_bytes: AtomicUsize::new(0), state: Mutex::new(State::default()), }))) } @@ -333,14 +342,10 @@ impl ExecutionContext { /// Check admission without reserving or changing measured peak usage. /// The actual owner must still reserve before allocating. pub fn check_memory_available(&self, bytes: usize) -> Result<()> { - let state = self - .0 - .state - .lock() - .map_err(|_| ProcedureError::ResourceStatePoisoned)?; self.check_state(DeadlineCheck::Exact)?; - state + self.0 .live_bytes + .load(Ordering::Relaxed) .checked_add(bytes) .filter(|next| *next <= self.0.limits.memory_bytes) .ok_or(ProcedureError::BudgetExceeded { @@ -373,35 +378,57 @@ impl ExecutionContext { } fn charge_memory(&self, bytes: usize, deadline: DeadlineCheck) -> Result<()> { - let mut state = self - .0 - .state - .lock() - .map_err(|_| ProcedureError::ResourceStatePoisoned)?; self.check_state(deadline)?; - let next = state - .live_bytes - .checked_add(bytes) - .filter(|next| *next <= self.0.limits.memory_bytes) - .ok_or(ProcedureError::BudgetExceeded { - resource: "memory", - limit: self.0.limits.memory_bytes, - })?; - state.live_bytes = next; - state.peak_bytes = state.peak_bytes.max(next); - Ok(()) + // Admit exactly, as `admit_work` does: the exchange recomputes admission + // against the value it actually replaces, so concurrent charges cannot + // slip past the limit between them. + let mut current = self.0.live_bytes.load(Ordering::Relaxed); + loop { + let next = current + .checked_add(bytes) + .filter(|next| *next <= self.0.limits.memory_bytes) + .ok_or(ProcedureError::BudgetExceeded { + resource: "memory", + limit: self.0.limits.memory_bytes, + })?; + match self.0.live_bytes.compare_exchange_weak( + current, + next, + Ordering::Relaxed, + Ordering::Relaxed, + ) { + Ok(_) => { + // The peak can trail `live_bytes` for the instant between + // these two lines; it never misses a completed charge. + self.0.peak_bytes.fetch_max(next, Ordering::Relaxed); + return Ok(()); + } + Err(observed) => current = observed, + } + } + } + + /// Release bytes a reservation or account admitted. Only their `Drop` + /// calls this, with what they hold, so the counter cannot underflow. + fn release_memory(&self, bytes: usize) { + if bytes != 0 { + self.0.live_bytes.fetch_sub(bytes, Ordering::Relaxed); + } } /// Current accounted usage, including batches retained by a consumer. + /// + /// The three figures are read one after another, not as one snapshot: a + /// reader racing a charge can see a peak newer than the live figure beside + /// it. Live is read first and the peak only grows, so `peak_bytes >= + /// live_bytes` holds for every reader. Read after execution for exact + /// totals. pub fn usage(&self) -> Result { - let state = self - .0 - .state - .lock() - .map_err(|_| ProcedureError::ResourceStatePoisoned)?; + let live_bytes = self.0.live_bytes.load(Ordering::Relaxed); + let peak_bytes = self.0.peak_bytes.load(Ordering::Relaxed).max(live_bytes); Ok(ResourceUsage { - live_bytes: state.live_bytes, - peak_bytes: state.peak_bytes, + live_bytes, + peak_bytes, work_units: self.0.work_units.load(Ordering::Relaxed), }) } @@ -573,13 +600,7 @@ impl MemoryAccount { impl Drop for MemoryAccount { fn drop(&mut self) { - let mut state = self - .context - .0 - .state - .lock() - .unwrap_or_else(|poison| poison.into_inner()); - state.live_bytes -= self.bytes; + self.context.release_memory(self.bytes); } } @@ -591,14 +612,7 @@ struct Reservation { impl Drop for Reservation { fn drop(&mut self) { - // Recover only for cleanup; operational calls still report poisoning. - let mut state = self - .context - .0 - .state - .lock() - .unwrap_or_else(|poison| poison.into_inner()); - state.live_bytes -= self.bytes; + self.context.release_memory(self.bytes); } } diff --git a/crates/grust-procedures/tests/contracts.rs b/crates/grust-procedures/tests/contracts.rs index 0f999e87..a3bb13b8 100644 --- a/crates/grust-procedures/tests/contracts.rs +++ b/crates/grust-procedures/tests/contracts.rs @@ -12,6 +12,8 @@ mod builtins; mod cache; #[path = "contracts/cancellation.rs"] mod cancellation; +#[path = "contracts/memory.rs"] +mod memory; #[path = "contracts/parallel.rs"] mod parallel; diff --git a/crates/grust-procedures/tests/contracts/memory.rs b/crates/grust-procedures/tests/contracts/memory.rs new file mode 100644 index 00000000..632d8a3e --- /dev/null +++ b/crates/grust-procedures/tests/contracts/memory.rs @@ -0,0 +1,134 @@ +//! Memory accounting without a lock: admission stays exact under contention, +//! releases return everything, and the peak is a true high-water mark. + +use std::sync::atomic::{AtomicUsize, Ordering}; + +use super::*; + +const WORKERS: usize = 16; + +fn limits(memory_bytes: usize) -> ExecutionLimits { + ExecutionLimits { + memory_bytes, + work_units: usize::MAX, + batch_rows: 8, + deadline: None, + } +} + +#[test] +fn concurrent_reservations_never_admit_more_than_the_limit() { + const CHUNK: usize = 1000; + const FITS: usize = 50; + for _ in 0..50 { + let execution = ExecutionContext::new(limits(FITS * CHUNK + CHUNK / 2)).expect("limits"); + let admitted = AtomicUsize::new(0); + let refused = AtomicUsize::new(0); + let start = std::sync::Barrier::new(WORKERS); + let held = std::sync::Mutex::new(Vec::new()); + std::thread::scope(|scope| { + for _ in 0..WORKERS { + scope.spawn(|| { + start.wait(); + for _ in 0..FITS { + match execution.reserve(CHUNK) { + Ok(reservation) => { + admitted.fetch_add(1, Ordering::Relaxed); + held.lock().expect("held").push(reservation); + } + Err(ProcedureError::BudgetExceeded { + resource: "memory", .. + }) => { + refused.fetch_add(1, Ordering::Relaxed); + } + Err(other) => panic!("unexpected failure: {other}"), + } + } + }); + } + }); + // Every chunk that fits is admitted, not one more, whatever the race. + assert_eq!(admitted.load(Ordering::Relaxed), FITS); + assert_eq!(refused.load(Ordering::Relaxed), WORKERS * FITS - FITS); + let usage = execution.usage().expect("usage"); + assert_eq!(usage.live_bytes, FITS * CHUNK); + assert_eq!(usage.peak_bytes, FITS * CHUNK); + drop(held); + let usage = execution.usage().expect("usage"); + assert_eq!(usage.live_bytes, 0, "every release came back"); + assert_eq!(usage.peak_bytes, FITS * CHUNK, "the peak does not fall"); + } +} + +#[test] +fn accounts_and_reservations_dropped_on_many_threads_return_to_zero() { + let execution = ExecutionContext::new(limits(usize::MAX)).expect("limits"); + std::thread::scope(|scope| { + for worker in 0..WORKERS { + let execution = execution.clone(); + scope.spawn(move || { + for round in 0..2000 { + let mut account = execution.memory_account(); + account.charge(1 + (worker + round) % 97).expect("charge"); + let reservation = execution.reserve(64).expect("reserve"); + account.charge(3).expect("charge"); + drop(reservation); + } + }); + } + }); + let usage = execution.usage().expect("usage"); + assert_eq!(usage.live_bytes, 0); + assert!(usage.peak_bytes >= 64 + 4 && usage.peak_bytes <= WORKERS * (64 + 97 + 3)); +} + +#[test] +fn the_peak_is_the_true_maximum_of_a_sequence_and_never_below_live() { + let execution = ExecutionContext::new(limits(10_000)).expect("limits"); + let first = execution.reserve(4000).expect("first"); + let second = execution.reserve(5000).expect("second"); + assert_eq!(execution.usage().expect("usage").peak_bytes, 9000); + drop(first); + let usage = execution.usage().expect("usage"); + assert_eq!((usage.live_bytes, usage.peak_bytes), (5000, 9000)); + // A refused charge moves neither figure. + assert!(execution.reserve(6000).is_err()); + assert!(execution.check_memory_available(5001).is_err()); + assert!(execution.check_memory_available(5000).is_ok()); + let usage = execution.usage().expect("usage"); + assert_eq!((usage.live_bytes, usage.peak_bytes), (5000, 9000)); + let third = execution.reserve(5000).expect("exactly fits"); + assert_eq!(execution.usage().expect("usage").peak_bytes, 10_000); + drop((second, third)); + assert_eq!(execution.usage().expect("usage").live_bytes, 0); + + // A reader racing charges always sees a peak at or above the live figure. + let execution = ExecutionContext::new(limits(usize::MAX)).expect("limits"); + let done = std::sync::atomic::AtomicBool::new(false); + std::thread::scope(|scope| { + scope.spawn(|| { + for round in 0..200_000usize { + let _held = execution.reserve(1 + round % 4096).expect("reserve"); + } + done.store(true, Ordering::Release); + }); + while !done.load(Ordering::Acquire) { + let usage = execution.usage().expect("usage"); + assert!(usage.peak_bytes >= usage.live_bytes); + } + }); +} + +#[test] +fn overflowing_the_counter_is_a_budget_failure_not_a_wrap() { + let execution = ExecutionContext::new(limits(usize::MAX)).expect("limits"); + let _held = execution.reserve(usize::MAX - 10).expect("fits"); + assert!(matches!( + execution.reserve(11), + Err(ProcedureError::BudgetExceeded { + resource: "memory", + .. + }) + )); + assert!(execution.reserve(10).is_ok()); +} diff --git a/docs/book/chapters/generalized-algorithms.md b/docs/book/chapters/generalized-algorithms.md index 350add41..2ef1187d 100644 --- a/docs/book/chapters/generalized-algorithms.md +++ b/docs/book/chapters/generalized-algorithms.md @@ -218,16 +218,21 @@ Two properties of that meter are therefore part of the algorithm contract. cancellation flag are atomics. Each charge admits through a compare-exchange that recomputes admission against the value it actually replaces, so a concurrent charge cannot overshoot the limit between load and store, and an -exhausted budget still fails exactly at its limit. Memory reservations, peak -accounting and cancellation wakers keep a mutex: they are rare and need several -fields to move together. - -That is true of kernels, which reserve memory per batch. It is not true of a -materializing Cypher consumer, which charges the logical bytes of every value it -copies and so takes that mutex once or more per row. Those per-copy charges, -`charge_cumulative_memory` and `MemoryAccount::charge`, sample the deadline as -work charges do; `reserve` still reads the clock. Moving the byte counters to -atomics is planned in `docs/lock-free.md` and is not part of this release. +exhausted budget still fails exactly at its limit. Accounted memory and its +high-water mark are atomics admitted the same way, so a byte limit is as exact +as a work limit. Only the cancellation wakers keep a mutex: they are rare, and +they are the one place several fields must move together. + +Memory had to follow work off the lock because of who charges it. A kernel +reserves memory per batch, but a materializing Cypher consumer charges the +logical bytes of every value it copies, once or more per row, through +`charge_cumulative_memory` and `MemoryAccount::charge`. Those sample the +deadline as work charges do; `reserve` still reads the clock. Two things follow +from dropping the lock. `usage()` reads its figures one after another rather +than as one snapshot, so read it after execution for exact totals; while an +execution runs, `peak_bytes` is never below `live_bytes`. And a poisoned lock no +longer fails a memory charge, because there is no lock to poison; only waker +registration can still report it. **The deadline is sampled; everything else is not.** An execution that sets no deadline pays nothing for deadline enforcement — neither a clock read nor a diff --git a/docs/lock-free.md b/docs/lock-free.md index 12e4a8ff..6b74ca11 100644 --- a/docs/lock-free.md +++ b/docs/lock-free.md @@ -1,8 +1,8 @@ # Lock-free memory accounting — plan -Status: **PLANNED for the release after Tadpole 0.21.0.** Not started. Recorded -2026-09-18. Nothing here is implemented; every number below is an estimate or a -measurement of the code as it stands at `230ae31`. +Status: **DONE, 2026-09-20,** on `work/lock-free-memory`, as designed below. The +record of what was measured is at the end, under "Outcome". The sections before +it are the plan as written on 2026-09-18 and are left as they were. ## Why @@ -199,3 +199,41 @@ The remaining gap needs one of two larger changes, each its own goal: the more general fix. Neither is part of this plan. + +## Outcome (2026-09-20) + +L1 as designed: `live_bytes` and `peak_bytes` are atomics on `Shared`, admitted +by the same compare-exchange loop as work; releases are a `fetch_sub`; `State` +holds only the wakers. One refinement: `usage()` reports +`max(peak_bytes, live_bytes)`, which makes `peak >= live` hold for every reader +without depending on load order. + +L2: `crates/grust-procedures/tests/contracts/memory.rs`. Sixteen threads race +for a limit that fits exactly fifty chunks and exactly fifty are admitted, in +fifty rounds; accounts and reservations dropped across threads return +`live_bytes` to zero; the peak is the true maximum of a sequence, is unmoved by +a refused charge, and is never seen below `live_bytes` by a reader racing two +hundred thousand charges; a counter overflow is a budget failure. `loom` was not +used, so the interleavings are tested, not model-checked. 1,297 tests pass +across the crates that read `usage()`. + +L0 and L3, one macOS laptop, release, whole-process wall time, best of seven, +full-path Dijkstra on a chain through Cypher: + +| query | nodes | mutex | lock-free | +| --- | --- | --- | --- | +| `reduce` fold | 1,024 | 301 ms | 267 ms | +| `reduce` fold | 4,096 | 4,389 ms | 3,815 ms | +| `UNWIND` aggregate | 1,024 | 84 ms | 82 ms | +| `UNWIND` aggregate | 4,096 | 862 ms | 859 ms | +| direct kernel | 1,024 | 37 ms | 36 ms | +| direct kernel | 4,096 | 108 ms | 107 ms | + +The `reduce` form gains 11% and 13%, which is what the profile attributed to the +lock; the lock-free figure was measured before and after the mutex figure and +reproduced. The fused `UNWIND` form and the direct kernel, which charge memory +rarely, do not move, so the regression this plan worried about did not appear +here. `reduce` is still several times the `UNWIND` form, as "What this will not +deliver" said it would be. **The paired benchmark run on the algorithms harness +has not happened**; it is requested in `codex-to-codex.md` and these single-host +figures should not be quoted as that result.