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
Original file line number Diff line number Diff line change
Expand Up @@ -95,11 +95,12 @@ pub fn try_acquire_code_generation_store_read_lock(
}
}

/// Block the calling thread until no exclusive holder, in this or any other
/// process, holds the generation store, then let go at once.
/// Block the calling thread until no holder, in this or any other process,
/// holds the generation store, then let go at once.
///
/// The kernel wakes this wait on the holder's release itself, so a caller
/// refused by a held lock retries exactly when the lock becomes free. On
/// The kernel wakes this wait on the holders' release itself, so a caller
/// refused by a held lock retries exactly when the lock becomes free. The
/// wait is exclusive because a refused writer must also outlast readers. On
/// Windows it first waits out a retention pass fencing the scope; a scope a
/// pending retention transaction collected stays busy until that transaction
/// resolves, so it is reported as
Expand All @@ -112,7 +113,7 @@ pub fn wait_for_code_generation_store_release(
wait_for_generation_scope_fence_release(store_root)?;
let store_root = canonical_store_root(store_root)?;
open_lock_file(&store_root.join(STORE_LOCK_FILE))?
.lock_shared()
.lock()
.map_err(storage)
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3646,3 +3646,37 @@ fn a_pre_key_text_artifact_descriptor_reads_as_retired_and_is_replaced() {
"attaching a current artifact replaces the retired slot"
);
}

/// A writer refused because readers hold the store must wait out those
/// readers. A release wait that returned while a reader still held the lock
/// woke the refused writer only to be refused again.
#[test]
fn store_release_wait_outlasts_a_reader() {
let store = tempfile::TempDir::new().expect("generation store");
let reader = try_acquire_code_generation_store_read_lock(store.path())
.expect("try reader lock")
.expect("the store is free for a reader");
let root = store.path().to_path_buf();
let (released_tx, released) = std::sync::mpsc::channel();
let waiter = std::thread::spawn(move || {
released_tx
.send(wait_for_code_generation_store_release(&root).is_ok())
.expect("report the release");
});
assert_eq!(
released.recv_timeout(std::time::Duration::from_millis(300)),
Err(std::sync::mpsc::RecvTimeoutError::Timeout),
"the release wait must not return while a reader holds the store"
);
drop(reader);
assert_eq!(
released.recv_timeout(std::time::Duration::from_secs(10)),
Ok(true)
);
waiter.join().expect("waiter thread");
assert!(
try_acquire_code_generation_store_lock(store.path())
.expect("try writer lock")
.is_some()
);
}
Original file line number Diff line number Diff line change
Expand Up @@ -749,16 +749,16 @@ where
return None;
}
match serving {
code_index_scheduler::CodeIndexRetainedTextServingWaitV1::Serving
| code_index_scheduler::CodeIndexRetainedTextServingWaitV1::Warming => {
code_index_scheduler::CodeIndexSeatWaitV1::Seated(())
| code_index_scheduler::CodeIndexSeatWaitV1::Deadline => {
Some(code_index_search_unavailable(
code_search::CodeIndexSearchUnavailableReasonV1::GraphWarming,
code_search::CodeIndexSearchUnavailableReasonV1::GraphWarming
.as_str(),
))
}
code_index_scheduler::CodeIndexRetainedTextServingWaitV1::Unpublished
| code_index_scheduler::CodeIndexRetainedTextServingWaitV1::Unreachable => {
code_index_scheduler::CodeIndexSeatWaitV1::Parked(_)
| code_index_scheduler::CodeIndexSeatWaitV1::Cancelled => {
Some(code_index_search_unavailable(
code_search::CodeIndexSearchUnavailableReasonV1::AuthorityUnavailable,
"query_authority_unavailable",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -214,9 +214,9 @@ pub use registry::RetainedGraphRecoveryPauseV1;
pub use registry::watch_ingress::GitStateChangeRequestV1;
pub use registry::{
CodeIndexOwnerActivityV1, CodeIndexOwnerSignalsClosedV1, CodeIndexOwnerSignalsV1,
CodeIndexRetainedTextServingWaitV1, CodeIndexWorkerPhaseV1, ScopedFeedbackDocumentIdentityV1,
ServingGenerationInstallationOutcomeV1, ServingGenerationRollbackOutcomeV1,
feedback_document_identity_from_generation,
CodeIndexSeatParkV1, CodeIndexSeatWaitV1, CodeIndexWorkerPhaseV1,
ScopedFeedbackDocumentIdentityV1, ServingGenerationInstallationOutcomeV1,
ServingGenerationRollbackOutcomeV1, feedback_document_identity_from_generation,
feedback_language_document_identity_from_generation,
};
pub(crate) use residency::ServingReadLeaseV1;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -335,6 +335,7 @@ impl CodeIndexSchedulerRegistryV1 {
) -> Result<LatestCompleteCodeIndexV1, CallableCodeCursorError> {
let wait = remaining_generation_resolution_wait(request)
.ok_or(CallableCodeCursorError::Unavailable)?;
let deadline = tokio::time::Instant::now() + wait;
let resolution = async {
if let Some(cursor) = page.cursor.as_ref() {
let expected_generation = (!is_unpinned_latest(requested)).then_some(requested);
Expand All @@ -353,7 +354,7 @@ impl CodeIndexSchedulerRegistryV1 {
// longer held, so no later retry can serve it either.
.ok_or(CallableCodeCursorError::Stale)
} else if is_unpinned_latest(requested) {
self.latest_complete_fresh_for_scope_awaiting_seat(request.scope())
self.latest_complete_fresh_for_scope_awaiting_seat(request.scope(), deadline)
.await
.ok_or(CallableCodeCursorError::Unavailable)
} else {
Expand All @@ -363,7 +364,7 @@ impl CodeIndexSchedulerRegistryV1 {
.ok_or(CallableCodeCursorError::GenerationNotHeld)
}
};
let latest = tokio::time::timeout(wait, resolution)
let latest = tokio::time::timeout_at(deadline, resolution)
.await
.map_err(|_| CallableCodeCursorError::Unavailable)??;
if !matches!(
Expand Down Expand Up @@ -455,6 +456,7 @@ impl CodeIndexSchedulerRegistryV1 {
) -> Result<LatestCodeTextGenerationV1, CallableCodeCursorError> {
let wait = remaining_generation_resolution_wait(request)
.ok_or(CallableCodeCursorError::Unavailable)?;
let deadline = tokio::time::Instant::now() + wait;
let resolution = async {
if let Some(cursor) = page.cursor.as_ref() {
let expected_generation = (!is_unpinned_latest(requested)).then_some(requested);
Expand Down Expand Up @@ -489,7 +491,7 @@ impl CodeIndexSchedulerRegistryV1 {
.map_err(|_| CallableCodeCursorError::Unavailable)?
.ok_or(CallableCodeCursorError::Stale)
} else if is_unpinned_latest(requested) {
self.current_text_owner_for_scope(request.scope())
self.current_text_owner_for_scope(request.scope(), deadline)
.await
.ok_or(CallableCodeCursorError::Unavailable)
} else if let Some(latest) = self
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ use tracedecay_domain::{
};

use super::{
CodeIndexReconcileAdmissionV1, CodeIndexSchedulerRegistryV1,
CodeIndexOwnerSignalsV1, CodeIndexReconcileAdmissionV1, CodeIndexSchedulerRegistryV1,
serving::CodeTextQueryOwnerReadinessV1,
};
use tracedecay_query::retrieval::exact::{
Expand Down Expand Up @@ -84,17 +84,12 @@ pub enum DeferredMountAttemptV1 {

/// Waits for the first retained generation of `project_root` and then
/// retries the query-authority mount. Exits when the mount reaches any terminal
/// outcome or the publication channel closes (daemon shutdown).
/// outcome or the registry closes (daemon shutdown).
///
/// The open-time mount runs before code-index activation, so the first ready
/// check usually misses. Wake sources are event-driven only:
/// - a matching `Published` broadcast on a fresh build;
/// - registry-wide root-mounted watches so a pre-activation subscribe can
/// re-attach the per-worktree serving-generation watch after mount;
/// - serving-slot / serving-generation watches for a restart `Noop` restore
/// that never rebroadcasts (partitioned recovery may leave the decoded seat
/// empty and only flip the generation watch);
/// - on `Lagged`, retry immediately. Never a standing 1 Hz ready poll.
/// check usually misses. The root's owner signals wake each retry, including
/// a restart `Noop` restore that publishes nothing and only flips the
/// serving-generation watch; there is no standing ready poll.
pub async fn retry_deferred_query_authority_until_serving<F, Fut>(
registry: &CodeIndexSchedulerRegistryV1,
project_root: PathBuf,
Expand All @@ -103,19 +98,11 @@ pub async fn retry_deferred_query_authority_until_serving<F, Fut>(
F: FnMut() -> Fut,
Fut: std::future::Future<Output = DeferredMountAttemptV1>,
{
let mut publications = registry.subscribe_generation_publications();
let mut serving_seats = registry.subscribe_serving_seats();
let mut root_mounted = registry.subscribe_root_mounted();
let mut serving_changes = None;
// Subscribe before probing so a seat that lands between subscribe and the
// ready check still wakes the retry. The mount needs only the text owner,
// so it demands no decoded generation.
let mut signals = CodeIndexOwnerSignalsV1::subscribe(registry, &project_root).await;
loop {
// Subscribe before probing so a seat that lands between subscribe and
// the ready check remains observable. The mount needs only the text
// owner, so it demands no decoded generation.
if serving_changes.is_none() {
serving_changes = registry
.subscribe_serving_generation_changes(&project_root)
.await;
}
if registry
.retained_text_owner_for_root(&project_root)
.await
Expand All @@ -124,38 +111,8 @@ pub async fn retry_deferred_query_authority_until_serving<F, Fut>(
{
return;
}
tokio::select! {
publication = publications.recv() => match publication {
// Any publication wakes the next attempt at the top of the
// loop; a foreign project's publication costs one attempt,
// which is what the previous per-root filter also paid.
Ok(_) => {}
// A lagged receiver dropped publications; one of them may have
// been this project's. Retry immediately; do not install a
// standing timer.
Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => {}
Err(tokio::sync::broadcast::error::RecvError::Closed) => return,
},
serving = async {
match serving_changes.as_mut() {
Some(changes) => changes.changed().await,
None => std::future::pending().await,
}
} => {
if serving.is_err() {
return;
}
}
seat = serving_seats.changed() => {
if seat.is_err() {
return;
}
}
mounted = root_mounted.changed() => {
if mounted.is_err() {
return;
}
}
if signals.changed().await.is_err() {
return;
}
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -468,7 +468,8 @@ mod resident_memory;
mod test_gates;
pub mod watch_ingress;
pub use owner_signals::{
CodeIndexOwnerSignalsClosedV1, CodeIndexOwnerSignalsV1, CodeIndexRetainedTextServingWaitV1,
CodeIndexOwnerSignalsClosedV1, CodeIndexOwnerSignalsV1, CodeIndexSeatParkV1,
CodeIndexSeatWaitV1,
};

/// At most two distinct worktrees may reconcile concurrently. Each reconcile
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,8 @@ use tracedecay_code_index_retention::code_index_generations::{
};

use tracedecay_contracts::code_index_freshness::{
CodeGraphServingReadinessV1, CodeIndexBuildPhaseV1,
CodeGraphServingReadinessV1, CodeIndexBuildPhaseV1, CodeIndexReadinessTargetV1,
CodeIndexReadinessWaitReadV1,
};
use tracedecay_contracts::{CallableCodeOperationKind, callable_code_operation};
use tracedecay_domain::UtcMicros;
Expand All @@ -41,7 +42,7 @@ use super::super::graph_activation::{
set_injected_publication_deadline,
};
use super::super::tests::{OwnerSignals, application_context, query_authority};
use super::{CodeIndexRetainedTextServingWaitV1, CodeIndexSchedulerRegistryV1};
use super::{CodeIndexSchedulerRegistryV1, CodeIndexSeatParkV1, CodeIndexSeatWaitV1};
use crate::project_reads::project_code_graph_projection_read_port;
use tracedecay_runtime_core::path_safety::canonical_existing_identity;

Expand Down Expand Up @@ -690,10 +691,7 @@ async fn a_text_serving_wait_answers_while_the_graph_publishes() {
)
.await
.expect("the wait answers at its own budget");
assert_eq!(
without_authority,
CodeIndexRetainedTextServingWaitV1::Warming
);
assert_eq!(without_authority, CodeIndexSeatWaitV1::Deadline);

let text = fixture
.registry
Expand All @@ -719,7 +717,7 @@ async fn a_text_serving_wait_answers_while_the_graph_publishes() {
)
.await
.expect("the wait answers while the graph publication is held");
assert_eq!(serving, CodeIndexRetainedTextServingWaitV1::Serving);
assert_eq!(serving, CodeIndexSeatWaitV1::Seated(()));
assert_eq!(
fixture
.registry
Expand All @@ -735,6 +733,84 @@ async fn a_text_serving_wait_answers_while_the_graph_publishes() {
fixture.registry.shutdown().await;
}

/// A seat wait answers a parked worker with its park instead of waiting out
/// its budget for a seat the worker will not install. The squatted artifacts
/// root parks the published generation's text owner on every pass, so the
/// search's wait for retained text serving can never be satisfied.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_seat_wait_answers_a_parked_worker_with_its_park() {
let fixture =
Fixture::mount_with_poisoned_artifacts_root("project.seat-wait-parked", |artifacts_root| {
fs::write(artifacts_root, b"squatter").expect("occupy artifacts root path");
})
.await;
fixture
.wait_for_freshness(|freshness| freshness.parked.is_some())
.await
.expect("the worker parks on the squatted artifacts root");
let scope = fixture
.registry
.serving_code_scope(&fixture.project)
.await
.expect("mounted scope");
let operation =
callable_code_operation(CallableCodeOperationKind::Callers).expect("callers operation");
let context = application_context(&operation, scope.repository_id, scope.worktree_id);

let waited = tokio::time::timeout(
CONVERGENCE_DEADLINE,
fixture.registry.wait_for_retained_text_serving(
&fixture.project,
context.scope(),
CONVERGENCE_DEADLINE * 2,
),
)
.await
.expect("a parked worker ends the wait before its budget");
let CodeIndexSeatWaitV1::Parked(CodeIndexSeatParkV1::Convergence(park)) = waited else {
panic!("the wait must answer with the worker's park: {waited:?}");
};
assert!(
park.reason.contains("code text artifacts root"),
"{}",
park.reason
);
assert!(park.retries_on_wake);
fixture.registry.shutdown().await;
}

/// A readiness wait (`status` `wait_for`) ends as soon as the registry is
/// cancelled: a worktree shutting down installs nothing more, so the wait
/// reports the closed registry instead of spending its budget.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_readiness_wait_ends_when_the_registry_is_cancelled() {
let (fixture, _held) =
Fixture::mount_with_poisoned_artifacts_root_held("project.seat-wait-cancelled", |_| {})
.await;
let registry = fixture.registry.clone();
let project = fixture.project.clone();
let wait = tokio::spawn(async move {
registry
.wait_for_readiness(
&project,
CodeIndexReadinessTargetV1::Fresh,
CONVERGENCE_DEADLINE * 2,
)
.await
});
fixture.registry.cancel();
let waited = tokio::time::timeout(CONVERGENCE_DEADLINE, wait)
.await
.expect("cancellation ends the wait before its budget")
.expect("wait task")
.expect("readiness read");
let CodeIndexReadinessWaitReadV1::Unreachable { reason } = waited else {
panic!("a cancelled registry must end the wait: {waited:?}");
};
assert_eq!(reason, "code_index_scheduler_registry_closed");
fixture.registry.shutdown().await;
}

fn run_git_in(root: &Path, args: &[&str]) {
let output = Command::new("git")
.args(args)
Expand Down
Loading
Loading