From 656d35423d3e0cae84ac56f209c741551fd3a20d Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Mon, 24 Aug 2026 01:31:35 +0000 Subject: [PATCH 1/2] fix(runtime-core): commit the missing background CPU module perf(observation) landed an import of tracedecay_runtime_core::background_cpu, but neither the module file nor its declaration was added, so the pushed integration branch does not compile: every branch cut from it fails on an unresolved import. Co-Authored-By: Claude Opus 5 (1M context) --- .../src/background_cpu.rs | 461 ++++++++++++++++++ crates/tracedecay-runtime-core/src/lib.rs | 1 + 2 files changed, 462 insertions(+) create mode 100644 crates/tracedecay-runtime-core/src/background_cpu.rs diff --git a/crates/tracedecay-runtime-core/src/background_cpu.rs b/crates/tracedecay-runtime-core/src/background_cpu.rs new file mode 100644 index 0000000000..75366049a4 --- /dev/null +++ b/crates/tracedecay-runtime-core/src/background_cpu.rs @@ -0,0 +1,461 @@ +//! Process-wide weighted admission for background CPU work. +//! +//! The authority counts active CPU units rather than owning an executor. Code +//! indexing, semantic native threads, and session preparation can therefore +//! use their existing execution substrates while sharing one hard process +//! ceiling. FIFO waiter order prevents a continuously busy class from starving +//! another class, and RAII releases capacity on success, cancellation, or +//! unwind. + +use std::cell::Cell; +use std::collections::VecDeque; +use std::fmt; +use std::num::NonZeroUsize; +use std::sync::{Arc, Condvar, Mutex, OnceLock}; + +#[derive(Debug)] +struct BackgroundCpuWaiterV1 { + units: usize, +} + +#[derive(Default)] +struct BackgroundCpuStateV1 { + active_units: usize, + waiters: VecDeque>, +} + +/// One process-wide background CPU budget shared across subsystems. +pub struct ProcessBackgroundCpuV1 { + width: NonZeroUsize, + state: Mutex, + available: Condvar, +} + +impl fmt::Debug for ProcessBackgroundCpuV1 { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ProcessBackgroundCpuV1") + .field("width", &self.width) + .field("active_units", &self.active_units()) + .finish_non_exhaustive() + } +} + +thread_local! { + static BACKGROUND_CPU_DEPTH: Cell = const { Cell::new(0) }; + static BACKGROUND_CPU_UNITS: Cell = const { Cell::new(0) }; +} + +struct BackgroundCpuScopeV1; + +impl BackgroundCpuScopeV1 { + fn enter() -> Self { + BACKGROUND_CPU_DEPTH.with(|depth| depth.set(depth.get().saturating_add(1))); + Self + } +} + +struct YieldedBackgroundCpuV1<'a> { + authority: &'a Arc, + units: usize, + depth: usize, +} + +impl Drop for YieldedBackgroundCpuV1<'_> { + fn drop(&mut self) { + self.authority.admit_units(self.units); + BACKGROUND_CPU_UNITS.with(|units| units.set(self.units)); + BACKGROUND_CPU_DEPTH.with(|depth| depth.set(self.depth)); + } +} + +impl Drop for BackgroundCpuScopeV1 { + fn drop(&mut self) { + BACKGROUND_CPU_DEPTH.with(|depth| { + let remaining = depth.get().saturating_sub(1); + depth.set(remaining); + if remaining == 0 { + BACKGROUND_CPU_UNITS.with(|units| units.set(0)); + } + }); + } +} + +impl ProcessBackgroundCpuV1 { + fn new(width: NonZeroUsize) -> Self { + Self { + width, + state: Mutex::new(BackgroundCpuStateV1::default()), + available: Condvar::new(), + } + } + + #[must_use] + pub const fn width(&self) -> NonZeroUsize { + self.width + } + + #[must_use] + pub fn active_units(&self) -> usize { + self.state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .active_units + } + + #[must_use] + pub fn waiting_work_units(&self) -> usize { + let state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + waiting_units(&state) + } + + /// Acquire one CPU unit, waiting in FIFO order when the process budget is + /// full. The returned guard must remain alive for the active work unit. + pub fn acquire(self: &Arc) -> BackgroundCpuPermitV1 { + self.acquire_units(1) + } + + /// Acquire one CPU unit only when no earlier waiter exists and capacity is + /// immediately available. + pub fn try_acquire(self: &Arc) -> Option { + let mut state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + if !state.waiters.is_empty() || state.active_units >= self.width.get() { + return None; + } + state.active_units += 1; + record_state(&state, self.width); + Some(BackgroundCpuPermitV1 { + authority: Arc::clone(self), + units: 1, + }) + } + + /// Run one active work unit under the process budget. Nested work on the + /// same thread reuses a sufficient parent admission instead of waiting on + /// itself. + pub fn with_permit(self: &Arc, operation: impl FnOnce() -> R) -> R { + self.with_permits(1, operation) + } + + /// Run a weighted work unit, clamped to the entire process width. Semantic + /// inference uses its native intra-op thread count as the weight; ordinary + /// index/session preparation uses one. + pub fn with_permits( + self: &Arc, + requested_units: usize, + operation: impl FnOnce() -> R, + ) -> R { + let units = requested_units.max(1).min(self.width.get()); + let active_units = BACKGROUND_CPU_UNITS.with(Cell::get); + if active_units >= units { + return operation(); + } + if active_units > 0 { + let depth = BACKGROUND_CPU_DEPTH.with(Cell::get); + self.release(active_units); + BACKGROUND_CPU_UNITS.with(|active| active.set(0)); + BACKGROUND_CPU_DEPTH.with(|active| active.set(0)); + let _restore = YieldedBackgroundCpuV1 { + authority: self, + units: active_units, + depth, + }; + let _permit = self.acquire_units(units); + let _scope = BackgroundCpuScopeV1::enter(); + BACKGROUND_CPU_UNITS.with(|active| active.set(units)); + return operation(); + } + let _permit = self.acquire_units(units); + let _scope = BackgroundCpuScopeV1::enter(); + BACKGROUND_CPU_UNITS.with(|active| active.set(units)); + operation() + } + + /// Temporarily yield the caller's active units while a nested executor + /// fans out independently admitted leaf work. This prevents a parent + /// Rayon worker from holding capacity while it waits for child workers, + /// including a full-width weighted child. Capacity is reacquired before + /// the parent resumes, including during unwind. + pub fn with_yielded_permits(self: &Arc, operation: impl FnOnce() -> R) -> R { + let units = BACKGROUND_CPU_UNITS.with(Cell::get); + if units == 0 { + return operation(); + } + let depth = BACKGROUND_CPU_DEPTH.with(Cell::get); + self.release(units); + BACKGROUND_CPU_UNITS.with(|active| active.set(0)); + BACKGROUND_CPU_DEPTH.with(|active| active.set(0)); + let _restore = YieldedBackgroundCpuV1 { + authority: self, + units, + depth, + }; + operation() + } + + fn acquire_units(self: &Arc, units: usize) -> BackgroundCpuPermitV1 { + self.admit_units(units); + BackgroundCpuPermitV1 { + authority: Arc::clone(self), + units, + } + } + + fn admit_units(&self, units: usize) { + let waiter = Arc::new(BackgroundCpuWaiterV1 { units }); + let mut state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + state.waiters.push_back(Arc::clone(&waiter)); + record_state(&state, self.width); + loop { + let is_front = state + .waiters + .front() + .is_some_and(|front| Arc::ptr_eq(front, &waiter)); + if is_front && state.active_units.saturating_add(waiter.units) <= self.width.get() { + state.waiters.pop_front(); + state.active_units += waiter.units; + record_state(&state, self.width); + self.available.notify_all(); + return; + } + state = self + .available + .wait(state) + .unwrap_or_else(std::sync::PoisonError::into_inner); + } + } + + fn release(&self, units: usize) { + let mut state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + debug_assert!(state.active_units >= units); + state.active_units = state.active_units.saturating_sub(units); + record_state(&state, self.width); + self.available.notify_all(); + } +} + +fn record_state(state: &BackgroundCpuStateV1, width: NonZeroUsize) { + hotpath::gauge!("runtime_core.background_cpu.width").set(width.get()); + hotpath::gauge!("runtime_core.background_cpu.active_units").set(state.active_units); + hotpath::gauge!("runtime_core.background_cpu.waiting_work_units").set(waiting_units(state)); +} + +fn waiting_units(state: &BackgroundCpuStateV1) -> usize { + state + .waiters + .iter() + .fold(0usize, |total, waiter| total.saturating_add(waiter.units)) +} + +/// RAII ownership of active CPU capacity. Dropping it is cancellation-safe and +/// releases the exact acquired weight. +pub struct BackgroundCpuPermitV1 { + authority: Arc, + units: usize, +} + +impl fmt::Debug for BackgroundCpuPermitV1 { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("BackgroundCpuPermitV1") + .field("units", &self.units) + .finish_non_exhaustive() + } +} + +impl Drop for BackgroundCpuPermitV1 { + fn drop(&mut self) { + self.authority.release(self.units); + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq, thiserror::Error)] +pub enum BackgroundCpuInstallErrorV1 { + #[error( + "background CPU authority is already installed at width {installed_width}, not requested width {requested_width}" + )] + ConflictingWidth { + installed_width: usize, + requested_width: usize, + }, + #[error("background CPU authority installation did not settle")] + InstallationDidNotSettle, +} + +static PROCESS_BACKGROUND_CPU: OnceLock> = OnceLock::new(); + +/// Install or idempotently reuse the one process background CPU authority. +pub fn install_process_background_cpu( + width: NonZeroUsize, +) -> Result, BackgroundCpuInstallErrorV1> { + if let Some(installed) = PROCESS_BACKGROUND_CPU.get() { + return compare_installed_width(installed, width); + } + let requested = Arc::new(ProcessBackgroundCpuV1::new(width)); + match PROCESS_BACKGROUND_CPU.set(Arc::clone(&requested)) { + Ok(()) => Ok(requested), + Err(_) => PROCESS_BACKGROUND_CPU.get().map_or_else( + || Err(BackgroundCpuInstallErrorV1::InstallationDidNotSettle), + |installed| compare_installed_width(installed, width), + ), + } +} + +fn compare_installed_width( + installed: &Arc, + requested: NonZeroUsize, +) -> Result, BackgroundCpuInstallErrorV1> { + if installed.width == requested { + Ok(Arc::clone(installed)) + } else { + Err(BackgroundCpuInstallErrorV1::ConflictingWidth { + installed_width: installed.width.get(), + requested_width: requested.get(), + }) + } +} + +/// Installed process authority, or `None` before daemon worker-plan admission. +#[must_use] +pub fn process_background_cpu() -> Option> { + PROCESS_BACKGROUND_CPU.get().map(Arc::clone) +} + +#[cfg(test)] +mod tests { + use std::num::NonZeroUsize; + use std::sync::{ + Arc, Barrier, + atomic::{AtomicUsize, Ordering}, + }; + use std::time::Duration; + + use super::*; + + #[test] + fn combined_classes_never_exceed_width_and_both_progress() { + let authority = Arc::new(ProcessBackgroundCpuV1::new( + NonZeroUsize::new(4).expect("nonzero width"), + )); + let active = Arc::new(AtomicUsize::new(0)); + let maximum = Arc::new(AtomicUsize::new(0)); + let index_completed = Arc::new(AtomicUsize::new(0)); + let session_completed = Arc::new(AtomicUsize::new(0)); + let start = Arc::new(Barrier::new(17)); + let mut workers = Vec::new(); + for ordinal in 0..16 { + let authority = Arc::clone(&authority); + let active = Arc::clone(&active); + let maximum = Arc::clone(&maximum); + let index_completed = Arc::clone(&index_completed); + let session_completed = Arc::clone(&session_completed); + let start = Arc::clone(&start); + workers.push(std::thread::spawn(move || { + start.wait(); + authority.with_permit(|| { + let current = active.fetch_add(1, Ordering::SeqCst) + 1; + maximum.fetch_max(current, Ordering::SeqCst); + std::thread::sleep(Duration::from_millis(5)); + active.fetch_sub(1, Ordering::SeqCst); + if ordinal % 2 == 0 { + index_completed.fetch_add(1, Ordering::SeqCst); + } else { + session_completed.fetch_add(1, Ordering::SeqCst); + } + }); + })); + } + start.wait(); + for worker in workers { + worker.join().expect("background worker"); + } + + assert!(maximum.load(Ordering::SeqCst) <= 4); + assert_eq!(index_completed.load(Ordering::SeqCst), 8); + assert_eq!(session_completed.load(Ordering::SeqCst), 8); + } + + #[test] + fn waiting_demand_sums_weighted_work_units() { + let state = BackgroundCpuStateV1 { + active_units: 4, + waiters: VecDeque::from([ + Arc::new(BackgroundCpuWaiterV1 { units: 4 }), + Arc::new(BackgroundCpuWaiterV1 { units: 1 }), + ]), + }; + + assert_eq!(waiting_units(&state), 5); + } + + #[test] + fn weighted_units_and_nested_work_share_one_width() { + let authority = Arc::new(ProcessBackgroundCpuV1::new( + NonZeroUsize::new(4).expect("nonzero width"), + )); + authority.with_permits(4, || { + assert_eq!(authority.active_units(), 4); + authority.with_permit(|| assert_eq!(authority.active_units(), 4)); + }); + authority.with_permit(|| { + assert_eq!(authority.active_units(), 1); + authority.with_permits(4, || assert_eq!(authority.active_units(), 4)); + assert_eq!(authority.active_units(), 1); + }); + assert_eq!(authority.active_units(), 0); + } + + #[test] + fn nested_executor_yields_parent_units_and_reacquires_them() { + let authority = Arc::new(ProcessBackgroundCpuV1::new( + NonZeroUsize::new(4).expect("nonzero width"), + )); + authority.with_permit(|| { + assert_eq!(authority.active_units(), 1); + authority.with_yielded_permits(|| { + assert_eq!(authority.active_units(), 0); + authority.with_permits(4, || assert_eq!(authority.active_units(), 4)); + }); + assert_eq!(authority.active_units(), 1); + }); + assert_eq!(authority.active_units(), 0); + } + + #[test] + fn panic_and_cancellation_drop_release_every_unit() { + let authority = Arc::new(ProcessBackgroundCpuV1::new( + NonZeroUsize::new(2).expect("nonzero width"), + )); + let panic = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + authority.with_permits(2, || panic!("injected background panic")); + })); + assert!(panic.is_err()); + assert_eq!(authority.active_units(), 0); + + let nested_panic = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + authority.with_permit(|| { + authority.with_yielded_permits(|| panic!("injected nested executor panic")); + }); + })); + assert!(nested_panic.is_err()); + assert_eq!(authority.active_units(), 0); + + let cancelled = authority.acquire(); + assert_eq!(authority.active_units(), 1); + drop(cancelled); + assert_eq!(authority.active_units(), 0); + assert!(authority.try_acquire().is_some()); + } +} diff --git a/crates/tracedecay-runtime-core/src/lib.rs b/crates/tracedecay-runtime-core/src/lib.rs index 7a9d0ef453..7a176ffb26 100644 --- a/crates/tracedecay-runtime-core/src/lib.rs +++ b/crates/tracedecay-runtime-core/src/lib.rs @@ -76,6 +76,7 @@ #![allow(rustdoc::broken_intra_doc_links)] #![allow(rustdoc::private_intra_doc_links)] +pub mod background_cpu; pub mod branch; pub mod branch_meta; pub mod cancellation; From d4ce5ecfe2ddc3d2d20abf94e55d86762c1b6a00 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Mon, 24 Aug 2026 02:01:48 +0000 Subject: [PATCH 2/2] test(writer): pin the batch-to-frame coalescing ratio MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The session-ingest profile measured 26,382 `submit_authorized` calls settling as 26,380 durable transactions — an effective batch size of 1.0, so the writer paid one BEGIN IMMEDIATE/COMMIT pair per frame. That amplification is already addressed: `persist_observations` now submits one request per scan batch, and the writer worker coalesces compatible queued arrivals. Both are covered above this crate, in the stores that decide how many requests to submit. What was uncovered is `build_batches` itself — the pure function that groups queued requests into transactions. It can regress to one batch per request while every existing assertion stays green, because those assertions constrain submission counts, not grouping. These assert on a count, not on elapsed time: the timings this work came from swing double digits run to run on identical input. Co-Authored-By: Claude Opus 5 (1M context) --- .../src/writer/worker/mod.rs | 201 ++++++++++++++++++ 1 file changed, 201 insertions(+) diff --git a/crates/tracedecay-rusqlite-runtime/src/writer/worker/mod.rs b/crates/tracedecay-rusqlite-runtime/src/writer/worker/mod.rs index 9851d34f51..07649c05b1 100644 --- a/crates/tracedecay-rusqlite-runtime/src/writer/worker/mod.rs +++ b/crates/tracedecay-rusqlite-runtime/src/writer/worker/mod.rs @@ -937,3 +937,204 @@ mod auxiliary_scheduling_tests { assert!(connection.is_autocommit()); } } + +/// Regression cover for writer transaction amplification. +/// +/// The measured session-ingest profile showed 26,382 `submit_authorized` calls +/// settling as 26,380 durable transactions — an effective batch size of 1.0, +/// so the writer paid one `BEGIN IMMEDIATE`/`COMMIT` pair per frame. The +/// coalescing itself lives in [`build_batches`], and nothing in this crate +/// pinned its ratio: every existing assertion sits above the writer, in the +/// stores that decide how many requests to submit. Those stores can stay fixed +/// while this function silently regresses to one batch per request. +/// +/// These assert on a count — batches produced for a known request count — +/// because the timings this work came from swing double digits run to run on +/// identical input, so an elapsed-time assertion here would prove nothing. +#[cfg(test)] +mod batch_coalescing_tests { + use std::sync::Arc; + + use tracedecay_store::{ + AdmissionConfigV1, BatchBudgetV1, RuntimeCancellationIdentityV1, RuntimeDeadlineV1, + RuntimeInterruptionV1, RuntimeRequestProbeV1, RuntimeSubmitRequestV1, + }; + + use crate::{ + admission::{Admission, Capacity, Limits}, + test_support::{metadata, request}, + writer::UnrestrictedRuntimeWriteAuthority, + }; + + use super::{AcceptedRequest, ExecutionBatch, build_batches}; + + /// A probe that never interrupts, so batch shape is decided only by the + /// budget and by `requires_isolated_commit`. + struct BatchProbe { + cancellation: RuntimeCancellationIdentityV1, + deadline: RuntimeDeadlineV1, + isolated: bool, + } + + impl BatchProbe { + fn new(request: &RuntimeSubmitRequestV1, isolated: bool) -> Self { + Self { + cancellation: request.control().cancellation.clone(), + deadline: request.control().deadline.clone(), + isolated, + } + } + } + + impl RuntimeRequestProbeV1 for BatchProbe { + fn cancellation_identity(&self) -> &RuntimeCancellationIdentityV1 { + &self.cancellation + } + + fn deadline_identity(&self) -> &RuntimeDeadlineV1 { + &self.deadline + } + + fn interruption(&self) -> Option { + None + } + + fn try_begin_commit(&self) -> bool { + true + } + + fn requires_isolated_commit(&self) -> bool { + self.isolated + } + } + + fn config(max_operations: u32) -> AdmissionConfigV1 { + let defaults = AdmissionConfigV1::default(); + AdmissionConfigV1 { + foreground_batch: BatchBudgetV1 { + max_operations, + max_bytes: u64::MAX, + ..defaults.foreground_batch + }, + background_batch: BatchBudgetV1 { + max_operations, + max_bytes: u64::MAX, + ..defaults.background_batch + }, + ..defaults + } + } + + /// Builds `count` mutually compatible foreground requests, the shape one + /// scan batch submits. + fn compatible_requests(count: usize, isolate_last: bool) -> Vec { + let admission = Admission::new( + Limits::new( + Capacity { + operations: u32::try_from(count).unwrap() + 8, + bytes: u64::MAX, + }, + Capacity { + operations: 1, + bytes: u64::MAX, + }, + u64::MAX, + u64::MAX, + ) + .unwrap(), + ); + (0..count) + .map(|index| { + let submit = Arc::new(request(metadata( + &format!("operation.batch.ratio.{index}"), + &format!("key.batch.ratio.{index}"), + 'a', + ))); + let permit = admission.reserve(&submit.envelope().metadata).unwrap(); + let isolated = isolate_last && index + 1 == count; + let (reply, _response) = tokio::sync::oneshot::channel(); + AcceptedRequest::new( + Arc::clone(&submit), + Arc::new(BatchProbe::new(&submit, isolated)), + Arc::new(UnrestrictedRuntimeWriteAuthority), + reply, + permit, + ) + }) + .collect() + } + + fn frames(batches: &[ExecutionBatch]) -> usize { + batches.iter().map(|batch| batch.items.len()).sum() + } + + #[test] + fn one_scan_batch_of_compatible_writes_opens_one_transaction() { + const REQUESTS: usize = 256; + + let batches = build_batches(compatible_requests(REQUESTS, false), &config(512)); + + assert_eq!( + frames(&batches), + REQUESTS, + "coalescing must not drop or duplicate a request" + ); + assert_eq!( + batches.len(), + 1, + "a budget-fitting scan batch must open one durable transaction: \ + batches={} frames={REQUESTS} ratio={:.4} batches/frame (target <= 1.0)", + batches.len(), + batches.len() as f64 / REQUESTS as f64, + ); + } + + #[test] + fn batch_count_is_bounded_by_the_budget_not_by_the_request_count() { + const REQUESTS: usize = 256; + const MAX_OPERATIONS: u32 = 64; + const EXPECTED: usize = REQUESTS / MAX_OPERATIONS as usize; + + let batches = build_batches( + compatible_requests(REQUESTS, false), + &config(MAX_OPERATIONS), + ); + + assert_eq!(frames(&batches), REQUESTS); + assert_eq!( + batches.len(), + EXPECTED, + "batch count must track the operation budget, not the request count: \ + batches={} frames={REQUESTS} budget={MAX_OPERATIONS}", + batches.len(), + ); + assert!( + batches + .iter() + .all(|batch| batch.items.len() <= MAX_OPERATIONS as usize), + "no batch may exceed its operation budget" + ); + } + + /// Widening the transaction window must not silently swallow a request that + /// asked to commit alone. + #[test] + fn an_isolated_commit_still_gets_its_own_transaction() { + const REQUESTS: usize = 8; + + let batches = build_batches(compatible_requests(REQUESTS, true), &config(512)); + + assert_eq!(frames(&batches), REQUESTS); + assert_eq!( + batches.len(), + 2, + "an isolated-commit request must split the batch: batches={}", + batches.len(), + ); + assert_eq!( + batches.last().expect("isolated batch").items.len(), + 1, + "the isolated-commit request must not share its transaction" + ); + } +}