From 76a8e78384b21dfbc39b4703c91ce8450e440fca Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Wed, 16 Sep 2026 20:25:55 +0000 Subject: [PATCH 1/3] fix(code-index): restore nested admission fix lost in integration --- crates/tracedecay-code-index/src/chunks.rs | 45 +++++++++++++--------- 1 file changed, 26 insertions(+), 19 deletions(-) diff --git a/crates/tracedecay-code-index/src/chunks.rs b/crates/tracedecay-code-index/src/chunks.rs index 214717d7bd..71a45ebdf1 100644 --- a/crates/tracedecay-code-index/src/chunks.rs +++ b/crates/tracedecay-code-index/src/chunks.rs @@ -229,15 +229,19 @@ where if chunks.len() < PARALLEL_CHUNK_THRESHOLD { return chunks.iter().try_for_each(&operation); } - let failure = chunks - .par_iter() - .enumerate() - .filter_map(|(index, chunk)| { - admit(&mut || operation(chunk)) - .err() - .map(|error| (index, error)) - }) - .min_by_key(|(index, _)| *index); + // Leaves are admitted one unit at a time on whichever worker runs them, + // so the caller's own unit must not be held across the join. + let failure = crate::parallelism::with_yielded_background_cpu_permits(|| { + chunks + .par_iter() + .enumerate() + .filter_map(|(index, chunk)| { + admit(&mut || operation(chunk)) + .err() + .map(|error| (index, error)) + }) + .min_by_key(|(index, _)| *index) + }); match failure { Some((_, error)) => Err(error), None => Ok(()), @@ -349,10 +353,12 @@ impl ExactExtractionAuthorityV1 { if chunks.len() < PARALLEL_CHUNK_THRESHOLD { return chunks.into_iter().map(|chunk| self.admit(chunk)).collect(); } - let admitted = chunks - .into_par_iter() - .map(|chunk| crate::parallelism::with_background_cpu_permit(|| self.admit(chunk))) - .collect::>(); + let admitted = crate::parallelism::with_yielded_background_cpu_permits(|| { + chunks + .into_par_iter() + .map(|chunk| crate::parallelism::with_background_cpu_permit(|| self.admit(chunk))) + .collect::>() + }); admitted.into_iter().collect() } @@ -2585,13 +2591,14 @@ mod tests { use super::*; use crate::extract::ExtractionCoverageV1; use tracedecay_domain::{ - BoundedSanitizedText, ChunkerRevision, CodeGenerationId, CodeSearchChunkAnchorV1, - CodeSearchChunkGrainV1, CodeSearchChunkId, ContentDigest, FileOccurrenceId, - GrammarRevision, LanguageDescriptorRevision, LanguageId, ManifestDigest, PolicyRevisionId, - ProjectId, SanitizationReceiptId, SanitizedCodeFileV1, SanitizedCodeSnapshotV1, - SanitizerRevision, SensitivityDecision, SensitivityLevelV1, SnapshotFileDispositionV1, - SourceSpan, SymbolOccurrenceId, UtcMicros, ValidatedCodeFileV1, + BoundedSanitizedText, ChunkerRevision, CodeGenerationId, CodeIndexWorkerSelectionV1, + CodeSearchChunkAnchorV1, CodeSearchChunkGrainV1, CodeSearchChunkId, ContentDigest, + FileOccurrenceId, GrammarRevision, LanguageDescriptorRevision, LanguageId, ManifestDigest, + PolicyRevisionId, ProjectId, SanitizationReceiptId, SanitizedCodeFileV1, + SanitizedCodeSnapshotV1, SanitizerRevision, SensitivityDecision, SensitivityLevelV1, + SnapshotFileDispositionV1, SourceSpan, SymbolOccurrenceId, UtcMicros, ValidatedCodeFileV1, }; + use tracedecay_runtime_core::resident_memory::DEFAULT_PROCESS_RESIDENT_MEMORY_LIMIT_V1; use crate::extract::{ ExtractionCancellation, LanguageExtractor as CanonicalLanguageExtractor, NeverCancelled, From f2d2c1c86ae6ccf0eb4d10f345dce8eb69ccfeb2 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Wed, 16 Sep 2026 20:36:42 +0000 Subject: [PATCH 2/3] test(code-index): isolate admission and pool failure probes --- crates/tracedecay-code-index/src/chunks.rs | 84 +++++++++---------- .../tracedecay-code-index/src/incremental.rs | 26 ++++++ .../tracedecay-code-index/src/parallelism.rs | 6 +- .../code_index_suite/chunk_incremental.rs | 34 -------- .../src/background_cpu.rs | 72 ++-------------- 5 files changed, 75 insertions(+), 147 deletions(-) diff --git a/crates/tracedecay-code-index/src/chunks.rs b/crates/tracedecay-code-index/src/chunks.rs index 71a45ebdf1..3a03bda6f6 100644 --- a/crates/tracedecay-code-index/src/chunks.rs +++ b/crates/tracedecay-code-index/src/chunks.rs @@ -2745,13 +2745,27 @@ mod tests { /// out across the pool while a full-width request (the lexical sorter's /// admission) is already queued at the FIFO head. Stolen leaves must not /// wait behind that head on a unit their own parent holds. - /// - /// The installed authority is process-global, so every ordering signal - /// comes from inside the request that takes the queue position - /// (`with_permits_placed`), never from shared counters that sibling - /// tests also move. #[test] fn nested_chunk_fan_out_does_not_wedge_behind_a_full_width_head_waiter() { + // The installed authority is global. Run this scenario alone so its + // queue counters cannot be advanced by another test's admissions. + if std::env::var_os("TRACEDECAY_NESTED_ADMISSION_CHILD").is_none() { + let output = std::process::Command::new(std::env::current_exe().expect("test binary")) + .arg("--exact") + .arg(std::thread::current().name().expect("named libtest thread")) + .arg("--nocapture") + .env("TRACEDECAY_NESTED_ADMISSION_CHILD", "1") + .output() + .expect("run isolated admission test"); + assert!( + output.status.success(), + "{}\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + assert!(String::from_utf8_lossy(&output.stdout).contains("1 passed")); + return; + } let installed = crate::parallelism::install_worker_plan( CodeIndexWorkerSelectionV1::Automatic {}, DEFAULT_PROCESS_RESIDENT_MEMORY_LIMIT_V1.get(), @@ -2760,10 +2774,6 @@ mod tests { let authority = installed.background_cpu; let width = authority.width().get(); if width < 2 { - eprintln!( - "skipping nested admission wedge: needs a second pool worker to steal onto, \ - installed width is {width}" - ); return; } let chunks = std::iter::repeat_n( @@ -2771,49 +2781,36 @@ mod tests { PARALLEL_CHUNK_THRESHOLD * width, ) .collect::>(); - let chunks_len = chunks.len(); let holder_admitted = Arc::new(AtomicBool::new(false)); - let head_placed = Arc::new(AtomicBool::new(false)); - let leaves_placed = Arc::new(AtomicUsize::new(0)); + let head_queued = Arc::new(AtomicBool::new(false)); let leaves_entered = Arc::new(AtomicUsize::new(0)); let (finished, finishes) = mpsc::channel::<&'static str>(); let holder = { let authority = Arc::clone(&authority); let holder_admitted = Arc::clone(&holder_admitted); - let head_placed = Arc::clone(&head_placed); - let leaves_placed = Arc::clone(&leaves_placed); + let head_queued = Arc::clone(&head_queued); let leaves_entered = Arc::clone(&leaves_entered); let finished = finished.clone(); std::thread::spawn(move || { let outcome = crate::parallelism::install(|| { crate::parallelism::with_background_cpu_permit(|| { holder_admitted.store(true, Ordering::SeqCst); - wait_until("full-width head request placed in the queue", || { - head_placed.load(Ordering::SeqCst) + wait_until("full-width head request queued", || { + head_queued.load(Ordering::SeqCst) }); try_for_each_chunk_ordered( - // Production admits each leaf through - // `with_background_cpu_permit` on this same - // authority; the placed hook only adds the signal. - |unit| { - authority.with_permits_placed( - 1, - || { - leaves_placed.fetch_add(1, Ordering::SeqCst); - }, - unit, - ) - }, + |unit| crate::parallelism::with_background_cpu_permit(unit), &chunks, |_| { if leaves_entered.fetch_add(1, Ordering::SeqCst) == 0 { // Keep the first leaf busy until a sibling - // leaf has taken its own queue position, - // which is behind the head. - wait_until("a sibling leaf placed its request", || { - leaves_placed.load(Ordering::SeqCst) >= 2 + // either ran (admission progressed) or is + // queued behind the head (the wedge). + wait_until("a sibling leaf ran or queued", || { + leaves_entered.load(Ordering::SeqCst) >= 2 + || authority.waiting_work_units() > width }); } Ok(()) @@ -2826,40 +2823,35 @@ mod tests { }) }; wait_until("holder admitted", || holder_admitted.load(Ordering::SeqCst)); + assert_eq!(authority.active_units(), 1); let head = { - let authority = Arc::clone(&authority); - let head_placed = Arc::clone(&head_placed); let finished = finished.clone(); std::thread::spawn(move || { - authority.with_permits_placed( - width, - || head_placed.store(true, Ordering::SeqCst), - || {}, - ); + crate::parallelism::with_background_cpu_permits(width, || {}); finished.send("head").expect("test thread is waiting"); }) }; + wait_until("head request waiting for the full width", || { + authority.waiting_work_units() >= width + }); + head_queued.store(true, Ordering::SeqCst); for _ in 0..2 { assert!( finishes.recv_timeout(ADMISSION_STEP_DEADLINE).is_ok(), "nested chunk fan-out wedged behind the full-width head waiter for \ {ADMISSION_STEP_DEADLINE:?}: active_units={} waiting_work_units={} \ - leaves_placed={} leaves_entered={}", + leaves_entered={}", authority.active_units(), authority.waiting_work_units(), - leaves_placed.load(Ordering::SeqCst), leaves_entered.load(Ordering::SeqCst), ); } holder.join().expect("holder thread"); head.join().expect("head thread"); - assert_eq!( - leaves_entered.load(Ordering::SeqCst), - chunks_len, - "every leaf must run once the fan-out completes" - ); + assert_eq!(authority.active_units(), 0); + assert_eq!(authority.waiting_work_units(), 0); } const RUST_SOURCE: &str = "//! Module documentation.\n\nuse std::collections::HashMap;\n\n/// Doc comment.\npub fn alpha(x: u32) -> u32 {\n x + 1\n}\n\npub struct Holder {\n map: HashMap,\n}\n\nimpl Holder {\n pub fn get(&self, key: u32) -> Option {\n self.map.get(&key).copied()\n }\n}\n\n// A trailing free-floating comment.\n"; diff --git a/crates/tracedecay-code-index/src/incremental.rs b/crates/tracedecay-code-index/src/incremental.rs index f3e906a117..ad2f0d59f9 100644 --- a/crates/tracedecay-code-index/src/incremental.rs +++ b/crates/tracedecay-code-index/src/incremental.rs @@ -606,3 +606,29 @@ fn placeholder_digest() -> ManifestDigest { ManifestDigest::new(format!("sha256:{}", "0".repeat(64))) .expect("a zeroed sha256 digest is canonical") } + +#[cfg(test)] +mod tests { + use super::{ChunkIncrementErrorV1, GenerationChunkManifestV1}; + use crate::parallelism::{CodeIndexParallelismErrorV1, force_install_failure_for_test}; + use tracedecay_domain::CodeGenerationId; + + #[test] + fn pool_failure_remains_a_typed_parallelism_error() { + struct ClearForce; + impl Drop for ClearForce { + fn drop(&mut self) { + force_install_failure_for_test(false); + } + } + let _clear = ClearForce; + force_install_failure_for_test(true); + let generation = CodeGenerationId::new("generation.pool-failure").expect("generation"); + assert!(matches!( + GenerationChunkManifestV1::new(generation, vec![]), + Err(ChunkIncrementErrorV1::Parallelism( + CodeIndexParallelismErrorV1::PoolBuild { .. } + )) + )); + } +} diff --git a/crates/tracedecay-code-index/src/parallelism.rs b/crates/tracedecay-code-index/src/parallelism.rs index 2c70664fc3..ddfe655b36 100644 --- a/crates/tracedecay-code-index/src/parallelism.rs +++ b/crates/tracedecay-code-index/src/parallelism.rs @@ -549,6 +549,7 @@ impl std::error::Error for CodeIndexParallelismErrorV1 {} /// 0 means "use the configured host width". static FORCED_WORKERS: AtomicUsize = AtomicUsize::new(0); +#[cfg(test)] thread_local! { /// Test-only: force [`install`] on this thread to return /// [`CodeIndexParallelismErrorV1::PoolBuild`]. Thread-scoped so a fault @@ -589,8 +590,8 @@ pub fn clear_forced_indexing_workers_for_test() { /// Force [`install`] to fail so callers can assert operational pool errors stay /// typed as parallelism failures instead of identity corruption. -#[doc(hidden)] -pub fn force_install_failure_for_test(force: bool) { +#[cfg(test)] +pub(crate) fn force_install_failure_for_test(force: bool) { FORCE_INSTALL_FAILURE.with(|flag| flag.set(force)); } @@ -641,6 +642,7 @@ where R: Send, { hotpath::gauge!("code_index_worker_count").set(indexing_workers()); + #[cfg(test)] if FORCE_INSTALL_FAILURE.with(std::cell::Cell::get) { return Err(CodeIndexParallelismErrorV1::PoolBuild { message: "forced code-index worker pool failure for test".to_owned(), diff --git a/crates/tracedecay-code-index/tests/code_index_suite/chunk_incremental.rs b/crates/tracedecay-code-index/tests/code_index_suite/chunk_incremental.rs index f6055a66f2..68ea89d660 100644 --- a/crates/tracedecay-code-index/tests/code_index_suite/chunk_incremental.rs +++ b/crates/tracedecay-code-index/tests/code_index_suite/chunk_incremental.rs @@ -353,37 +353,3 @@ fn duplicate_file_occurrences_are_rejected_before_manifest_flattening() { ))) ); } - -#[test] -fn standalone_pool_failure_is_parallelism_not_identity_validation() { - struct ClearForce; - impl Drop for ClearForce { - fn drop(&mut self) { - tracedecay_code_index::parallelism::force_install_failure_for_test(false); - } - } - let _clear = ClearForce; - tracedecay_code_index::parallelism::force_install_failure_for_test(true); - - let expected_generation = generation(2); - let file = baseline_file(&expected_generation, "file.ok", "src/lib.rs"); - let err = GenerationChunkManifestV1::new(expected_generation, vec![file]).unwrap_err(); - - match err { - ChunkIncrementErrorV1::Parallelism( - tracedecay_code_index::parallelism::CodeIndexParallelismErrorV1::PoolBuild { message }, - ) => { - assert!( - message.contains("forced"), - "expected forced pool failure, got {message}" - ); - } - ChunkIncrementErrorV1::NonCanonical(cause) => { - panic!( - "operational pool failure must not be NonCanonical ({})", - cause.reason_code() - ); - } - other => panic!("unexpected increment error: {other}"), - } -} diff --git a/crates/tracedecay-runtime-core/src/background_cpu.rs b/crates/tracedecay-runtime-core/src/background_cpu.rs index e878a1aa5e..9fc48fdac4 100644 --- a/crates/tracedecay-runtime-core/src/background_cpu.rs +++ b/crates/tracedecay-runtime-core/src/background_cpu.rs @@ -69,7 +69,7 @@ struct YieldedBackgroundCpuV1<'a> { impl Drop for YieldedBackgroundCpuV1<'_> { fn drop(&mut self) { - self.authority.admit_units(self.units, || {}); + self.authority.admit_units(self.units); BACKGROUND_CPU_UNITS.with(|units| units.set(self.units)); BACKGROUND_CPU_DEPTH.with(|depth| depth.set(self.depth)); } @@ -127,7 +127,7 @@ impl ProcessBackgroundCpuV1 { /// 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, || {}) + self.acquire_units(1) } /// Acquire one CPU unit only when no earlier waiter exists and capacity is @@ -178,25 +178,10 @@ impl ProcessBackgroundCpuV1 { self: &Arc, requested_units: usize, operation: impl FnOnce() -> R, - ) -> R { - self.with_permits_placed(requested_units, || {}, operation) - } - - /// [`Self::with_permits`] that calls `placed` once this request holds its - /// FIFO position, or ran on a sufficient admission the calling thread - /// already held. Callers that must order other work strictly after this - /// request's queue position (a request that cannot fit until the current - /// holders release) get that guarantee without reading shared counters. - pub fn with_permits_placed( - self: &Arc, - requested_units: usize, - placed: impl FnOnce(), - 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 { - placed(); return operation(); } if active_units > 0 { @@ -209,12 +194,12 @@ impl ProcessBackgroundCpuV1 { units: active_units, depth, }; - let _permit = self.acquire_units(units, placed); + 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, placed); + let _permit = self.acquire_units(units); let _scope = BackgroundCpuScopeV1::enter(); BACKGROUND_CPU_UNITS.with(|active| active.set(units)); operation() @@ -242,19 +227,15 @@ impl ProcessBackgroundCpuV1 { operation() } - fn acquire_units( - self: &Arc, - units: usize, - placed: impl FnOnce(), - ) -> BackgroundCpuPermitV1 { - self.admit_units(units, placed); + fn acquire_units(self: &Arc, units: usize) -> BackgroundCpuPermitV1 { + self.admit_units(units); BackgroundCpuPermitV1 { authority: Arc::clone(self), units, } } - fn admit_units(&self, units: usize, placed: impl FnOnce()) { + fn admit_units(&self, units: usize) { let waiter = Arc::new(BackgroundCpuWaiterV1 { units }); let mut state = self .state @@ -262,7 +243,6 @@ impl ProcessBackgroundCpuV1 { .unwrap_or_else(std::sync::PoisonError::into_inner); state.waiters.push_back(Arc::clone(&waiter)); record_state(&state, self.width); - placed(); loop { let is_front = state .waiters @@ -451,44 +431,6 @@ mod tests { assert_eq!(authority.active_units(), 0); } - #[test] - fn placed_fires_once_a_blocked_request_holds_its_queue_position() { - let authority = Arc::new(ProcessBackgroundCpuV1::new( - NonZeroUsize::new(2).expect("nonzero width"), - )); - let held = authority.acquire(); - let placed = Arc::new(AtomicBool::new(false)); - let admitted = Arc::new(AtomicBool::new(false)); - let waiter = { - let authority = Arc::clone(&authority); - let placed = Arc::clone(&placed); - let admitted = Arc::clone(&admitted); - std::thread::spawn(move || { - authority.with_permits_placed( - 2, - || placed.store(true, Ordering::SeqCst), - || admitted.store(true, Ordering::SeqCst), - ); - }) - }; - while !placed.load(Ordering::SeqCst) { - std::thread::sleep(Duration::from_millis(1)); - } - assert!(!admitted.load(Ordering::SeqCst)); - assert_eq!(authority.waiting_work_units(), 2); - assert_eq!(authority.active_units(), 1); - drop(held); - waiter.join().expect("full-width waiter"); - assert!(admitted.load(Ordering::SeqCst)); - assert_eq!(authority.active_units(), 0); - - let nested = AtomicBool::new(false); - authority.with_permits(2, || { - authority.with_permits_placed(1, || nested.store(true, Ordering::SeqCst), || {}); - }); - assert!(nested.load(Ordering::SeqCst)); - } - #[test] fn panic_and_cancellation_drop_release_every_unit() { let authority = Arc::new(ProcessBackgroundCpuV1::new( From 92e44cc09ab997c686d5a5a41a0fd664582f704e Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Wed, 16 Sep 2026 22:16:32 +0000 Subject: [PATCH 3/3] fix(code-index): deduplicate admission test imports --- crates/tracedecay-code-index/src/chunks.rs | 2 -- 1 file changed, 2 deletions(-) diff --git a/crates/tracedecay-code-index/src/chunks.rs b/crates/tracedecay-code-index/src/chunks.rs index 3a03bda6f6..4f2190a0dc 100644 --- a/crates/tracedecay-code-index/src/chunks.rs +++ b/crates/tracedecay-code-index/src/chunks.rs @@ -2606,8 +2606,6 @@ mod tests { }; use crate::intake::{CodeIndexIntake, SanitizedCodeIntake}; use crate::languages::{LanguageRegistry, StaticLanguageRegistry}; - use tracedecay_domain::configuration::CodeIndexWorkerSelectionV1; - use tracedecay_runtime_core::resident_memory::DEFAULT_PROCESS_RESIDENT_MEMORY_LIMIT_V1; struct AlwaysCancelled;