diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/serving.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/serving.rs index 8f0fab57d9..e15fdbd7d7 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/serving.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/serving.rs @@ -701,6 +701,9 @@ pub(super) struct CodeTextArtifactBuildV1 { pub(super) struct CodeTextProjectionStateV1 { slot: Mutex, ready: Condvar, + /// Whether the slot holds `Idle`, written with every slot transition so + /// a read probe never has to take the slot lock. + slot_idle: AtomicBool, } pub(super) enum CodeTextProjectionSlotV1 { @@ -720,9 +723,27 @@ impl CodeTextProjectionStateV1 { Self { slot: Mutex::new(CodeTextProjectionSlotV1::Idle), ready: Condvar::new(), + slot_idle: AtomicBool::new(true), } } + fn replace_slot( + &self, + slot: &mut CodeTextProjectionSlotV1, + next: CodeTextProjectionSlotV1, + ) -> CodeTextProjectionSlotV1 { + let previous = std::mem::replace(slot, next); + self.slot_idle.store( + matches!(slot, CodeTextProjectionSlotV1::Idle), + Ordering::Release, + ); + previous + } + + fn slot_is_idle(&self) -> bool { + self.slot_idle.load(Ordering::Acquire) + } + pub(super) fn lock_slot(&self) -> MutexGuard<'_, CodeTextProjectionSlotV1> { { let _span = tracing::trace_span!("query.artifact.head_open.lock_wait").entered(); @@ -754,7 +775,8 @@ impl<'a> TextHeadOpenClaimV1<'a> { ) -> MutexGuard<'a, CodeTextProjectionSlotV1> { self.armed = false; let mut slot = self.state.lock_slot(); - *slot = CodeTextProjectionSlotV1::Building(build); + self.state + .replace_slot(&mut slot, CodeTextProjectionSlotV1::Building(build)); self.state.ready.notify_all(); slot } @@ -767,7 +789,8 @@ impl Drop for TextHeadOpenClaimV1<'_> { } let mut slot = self.state.lock_slot(); if matches!(&*slot, CodeTextProjectionSlotV1::HeadOpening) { - *slot = CodeTextProjectionSlotV1::Idle; + self.state + .replace_slot(&mut slot, CodeTextProjectionSlotV1::Idle); } drop(slot); self.state.ready.notify_all(); @@ -1845,16 +1868,10 @@ impl LatestCodeTextGenerationV1 { if !self.query_owners_are_ready() { return true; } - // A read probe must not queue behind an advance: a wake holds the - // slot lock for its whole bounded slice. A contended slot is by - // definition work in progress. - match self.text_projection_build.slot.try_lock() { - Ok(slot) => !matches!(&*slot, CodeTextProjectionSlotV1::Idle), - Err(std::sync::TryLockError::WouldBlock) => true, - Err(std::sync::TryLockError::Poisoned(poisoned)) => { - !matches!(&*poisoned.into_inner(), CodeTextProjectionSlotV1::Idle) - } - } + // A read probe must not queue behind an advance, which holds the + // slot lock for its whole bounded slice. Nor may it read a held lock + // as work: another probe or an already-ready advance holds it too. + !self.text_projection_build.slot_is_idle() } pub(super) fn clone_index_status( @@ -3257,7 +3274,8 @@ impl LatestCodeTextGenerationV1 { if self.query_owners_are_ready() { return Ok(true); } - *slot = CodeTextProjectionSlotV1::HeadOpening; + self.text_projection_build + .replace_slot(&mut slot, CodeTextProjectionSlotV1::HeadOpening); drop(slot); let mut claim = TextHeadOpenClaimV1::new(&self.text_projection_build); let outcome = { @@ -3522,10 +3540,12 @@ impl LatestCodeTextGenerationV1 { // same discipline as the durable-head reopen. On failure the claim // restores `Idle` and the durable staging file resumes on a later // wake. - let CodeTextProjectionSlotV1::Building(finished) = - std::mem::replace(&mut *slot, CodeTextProjectionSlotV1::HeadOpening) + let CodeTextProjectionSlotV1::Building(finished) = self + .text_projection_build + .replace_slot(&mut slot, CodeTextProjectionSlotV1::HeadOpening) else { - *slot = CodeTextProjectionSlotV1::Idle; + self.text_projection_build + .replace_slot(&mut slot, CodeTextProjectionSlotV1::Idle); return Err(RetrievalPortError::Contract( "code-index text artifact build state vanished during publication".to_owned(), )); diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/reconcile.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/reconcile.rs index b9467bc02e..d2c2f2a022 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/reconcile.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/reconcile.rs @@ -3789,6 +3789,14 @@ async fn ready_wait_ends_only_after_the_graph_tail_seats_the_generation() { let registry = CodeIndexSchedulerRegistryV1::with_background_reconcile_permits(1, 1); let canonical_root = canonical_existing_identity(fixture.path()).expect("canonical fixture"); let (swap_entered, release_swap) = registry.pause_next_serving_swap(canonical_root); + // Hold admission until the complete-generation demand is posted, so it + // coalesces into the first pass. A demand arriving after that pass claims + // its wake queues a follow-up, which reads as `verifying` at the gate. + let admission = registry + .background_reconcile_admission() + .acquire_owned() + .await + .expect("hold worker before the complete-generation demand"); registry .mount_worktree( test_project_id(), @@ -3798,6 +3806,7 @@ async fn ready_wait_ends_only_after_the_graph_tail_seats_the_generation() { .await .expect("mount worktree"); assert!(registry.request_complete_generation(fixture.path()).await); + drop(admission); tokio::time::timeout(SERVING_SEAT_FAILURE_CEILING, swap_entered) .await .expect("publication did not reach its serving swap") diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/serving.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/serving.rs index 68bdfd21b2..59dffd8ced 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/serving.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/serving.rs @@ -885,6 +885,33 @@ fn text_artifact_publication_serializes_pointer_attachment_with_retention() { ); } +/// A held slot lock is not projection work: concurrent probes and a no-op +/// advance over ready owners take it too. Reading the contention as work +/// stamped a phantom continuation that kept freshness `verifying`. +#[test] +fn a_held_idle_projection_slot_is_not_text_projection_work() { + let fixture = GitFixture::new(ALPHA_LIB_V1); + let store = TempDir::new().expect("store root"); + let mut scheduler = scheduler( + &fixture, + store.path().to_path_buf(), + Arc::new(SharedCodeIndexBytePoolV1::default()), + ); + published(scheduler.reconcile_now().expect("publish generation")); + let latest = scheduler.latest_complete().expect("latest generation"); + while !latest.query_owners_are_ready() { + latest.advance_text_serving(1).expect("advance text build"); + } + assert!(!latest.text_projection_needs_work()); + + let held_slot = latest.text_projection_build.lock_slot(); + assert!( + !latest.text_projection_needs_work(), + "another holder of an idle slot is not projection work" + ); + drop(held_slot); +} + /// The artifact seals its clone index with its lexical rows, so the owners /// that serve search serve clone lookups at once: no work is left behind /// the first seal, and the clone status is ready (stale only when the source diff --git a/crates/tracedecay-code-index-runtime/src/git_watch.rs b/crates/tracedecay-code-index-runtime/src/git_watch.rs index 8675a4ca05..939245c42d 100644 --- a/crates/tracedecay-code-index-runtime/src/git_watch.rs +++ b/crates/tracedecay-code-index-runtime/src/git_watch.rs @@ -749,6 +749,8 @@ async fn debounce_loop( // If an operation is in flight, do not fire yet, wait for the next // event (marker removal wakes us) or a short recheck tick. if operation_state == OperationState::InFlight { + #[cfg(test)] + state.operation_held.notify_one(); tokio::select! { biased; () = cancellation.cancelled() => return DebounceExit::Cancelled, diff --git a/crates/tracedecay-code-index-runtime/src/git_watch/state.rs b/crates/tracedecay-code-index-runtime/src/git_watch/state.rs index 31dd7c7c7b..ef3d379ad5 100644 --- a/crates/tracedecay-code-index-runtime/src/git_watch/state.rs +++ b/crates/tracedecay-code-index-runtime/src/git_watch/state.rs @@ -151,6 +151,8 @@ pub struct WatchState { #[cfg(test)] pub plan_drained: Notify, #[cfg(test)] + pub operation_held: Notify, + #[cfg(test)] pub operation_scan_probe: OperationScanProbe, #[cfg(test)] pub retirement_probe: RetirementRaceProbe, @@ -206,6 +208,8 @@ impl WatchState { #[cfg(test)] plan_drained: Notify::new(), #[cfg(test)] + operation_held: Notify::new(), + #[cfg(test)] operation_scan_probe: OperationScanProbe::default(), #[cfg(test)] retirement_probe: RetirementRaceProbe::default(), diff --git a/crates/tracedecay-code-index-runtime/src/git_watch/tests.rs b/crates/tracedecay-code-index-runtime/src/git_watch/tests.rs index 1241910e83..000133a944 100644 --- a/crates/tracedecay-code-index-runtime/src/git_watch/tests.rs +++ b/crates/tracedecay-code-index-runtime/src/git_watch/tests.rs @@ -3,6 +3,7 @@ use super::identity::{ }; use super::*; +use futures_util::FutureExt; use notify::event::EventAttributes; use std::process::Command; use tokio::sync::Notify; @@ -591,7 +592,9 @@ async fn callback_failure_requests_conservative_reconciliation() { ); } -#[tokio::test(start_paused = true)] +// Real time: the live notify thread wakes this runtime, and a paused clock +// would auto-advance past the test's own deadlines while that wake is in flight. +#[tokio::test] async fn linked_worktree_operation_holds_real_debounce_until_marker_clears() { let (_container, primary, linked) = linked_worktree_fixture(); let watcher = GitWatcher::new(fast_watch_config()); @@ -616,11 +619,14 @@ async fn linked_worktree_operation_holds_real_debounce_until_marker_clears() { attrs: EventAttributes::default(), }; classify_and_mark(&state, &event); - for _ in 0..8 { - tokio::task::yield_now().await; - } - tokio::time::advance(Duration::from_millis(max_delay_ms + 1)).await; - tokio::task::yield_now().await; + tokio::time::timeout(TEST_READY_TIMEOUT, state.operation_held.notified()) + .await + .expect("debounce must observe the linked-worktree operation marker"); + tokio::time::sleep(Duration::from_millis(2 * max_delay_ms)).await; + let _ = state.operation_held.notified().now_or_never(); + tokio::time::timeout(TEST_READY_TIMEOUT, state.operation_held.notified()) + .await + .expect("the operation must keep holding past the hard deadline"); assert_eq!( state.drained_plans.load(Ordering::Relaxed), 0, @@ -636,11 +642,12 @@ async fn linked_worktree_operation_holds_real_debounce_until_marker_clears() { attrs: EventAttributes::default(), }, ); - tokio::time::advance(Duration::from_secs(1)).await; tokio::time::timeout(TEST_READY_TIMEOUT, state.plan_drained.notified()) .await .expect("clearing the linked-worktree marker must release the debounce"); - assert_eq!(state.drained_plans.load(Ordering::Relaxed), 1); + // The live watcher may still deliver its own copies of the marker writes + // after this drain; those are new evidence, so only the release is exact. + assert!(state.drained_plans.load(Ordering::Relaxed) >= 1); watcher.shutdown().await; }