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
131 changes: 64 additions & 67 deletions crates/tracedecay-code-index/src/chunks.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(()),
Expand Down Expand Up @@ -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::<Vec<_>>();
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::<Vec<_>>()
});
admitted.into_iter().collect()
}

Expand Down Expand Up @@ -2585,22 +2591,21 @@ 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,
TreeSitterExtractor,
};
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;

Expand Down Expand Up @@ -2738,13 +2743,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(),
Expand All @@ -2753,60 +2772,43 @@ 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(
Arc::clone(&file_chunks().chunks[0]),
PARALLEL_CHUNK_THRESHOLD * width,
)
.collect::<Vec<_>>();
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(())
Expand All @@ -2819,40 +2821,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<u32, u32>,\n}\n\nimpl Holder {\n pub fn get(&self, key: u32) -> Option<u32> {\n self.map.get(&key).copied()\n }\n}\n\n// A trailing free-floating comment.\n";
Expand Down
26 changes: 26 additions & 0 deletions crates/tracedecay-code-index/src/incremental.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 { .. }
))
));
}
}
6 changes: 4 additions & 2 deletions crates/tracedecay-code-index/src/parallelism.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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));
Comment on lines +591 to +595

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Move forced pool failures out of production code

Remove this exported fault switch from the production parallelism path: its only repository caller is the standalone_pool_failure_is_parallelism_not_identity_validation integration test, yet every real install call now checks mutable test state. Exercise the error mapping through an internal injectable executor or a test-local unit instead of shipping a test-only production port.

AGENTS.md reference: AGENTS.md:L106-L108

Useful? React with 👍 / 👎.

}

Expand Down Expand Up @@ -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(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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}"),
}
}
Loading
Loading