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
13 changes: 13 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
112 changes: 63 additions & 49 deletions crates/grust-procedures/src/resources.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,8 +35,6 @@ pub struct ResourceUsage {

#[derive(Debug, Default)]
struct State {
live_bytes: usize,
peak_bytes: usize,
waiters: Vec<Option<std::task::Waker>>,
}

Expand All @@ -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,
Expand All @@ -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<Vec<Arc<AtomicUsize>>>,
/// 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<State>,
}

Expand Down Expand Up @@ -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()),
})))
}
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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<ResourceUsage> {
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),
})
}
Expand Down Expand Up @@ -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);
}
}

Expand All @@ -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);
}
}

Expand Down
2 changes: 2 additions & 0 deletions crates/grust-procedures/tests/contracts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
134 changes: 134 additions & 0 deletions crates/grust-procedures/tests/contracts/memory.rs
Original file line number Diff line number Diff line change
@@ -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());
}
25 changes: 15 additions & 10 deletions docs/book/chapters/generalized-algorithms.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading