Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -701,6 +701,9 @@ pub(super) struct CodeTextArtifactBuildV1 {
pub(super) struct CodeTextProjectionStateV1 {
slot: Mutex<CodeTextProjectionSlotV1>,
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 {
Expand All @@ -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();
Expand Down Expand Up @@ -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
}
Expand All @@ -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();
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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 = {
Expand Down Expand Up @@ -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(),
));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand All @@ -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")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions crates/tracedecay-code-index-runtime/src/git_watch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
4 changes: 4 additions & 0 deletions crates/tracedecay-code-index-runtime/src/git_watch/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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(),
Expand Down
23 changes: 15 additions & 8 deletions crates/tracedecay-code-index-runtime/src/git_watch/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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());
Expand All @@ -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,
Expand All @@ -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;
}

Expand Down
Loading