From f50a4e805f8f1d7bfc4f87ee018e9033c7ab78a1 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Thu, 1 Oct 2026 18:55:04 +0000 Subject: [PATCH 1/4] fix(code-index): mount a store-lock-refused first open on its release --- .../src/code_index_generations/locking.rs | 11 +- .../src/code_index_generations/tests.rs | 34 +++++ .../code_index_scheduler/registry/mount.rs | 127 +++++++++++++----- .../reconcile_failure_isolation_tests.rs | 63 ++++++++- 4 files changed, 194 insertions(+), 41 deletions(-) diff --git a/crates/tracedecay-code-index-retention/src/code_index_generations/locking.rs b/crates/tracedecay-code-index-retention/src/code_index_generations/locking.rs index 8a7a6bcc75..76880593f7 100644 --- a/crates/tracedecay-code-index-retention/src/code_index_generations/locking.rs +++ b/crates/tracedecay-code-index-retention/src/code_index_generations/locking.rs @@ -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 @@ -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) } diff --git a/crates/tracedecay-code-index-retention/src/code_index_generations/tests.rs b/crates/tracedecay-code-index-retention/src/code_index_generations/tests.rs index 6842967644..a7129eba68 100644 --- a/crates/tracedecay-code-index-retention/src/code_index_generations/tests.rs +++ b/crates/tracedecay-code-index-retention/src/code_index_generations/tests.rs @@ -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() + ); +} diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs index a8cd732846..51ff4382c2 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs @@ -40,10 +40,10 @@ use super::{ CONVERGENCE_PARK_RECONCILE_FAILURE_REMEDIATION_V1, CONVERGENCE_PARK_STORE_RELEASE_WAIT_REMEDIATION_V1, CONVERGENCE_PARK_TASK_FAILURE_REMEDIATION_V1, CodeIndexSchedulerRegistryV1, - ColdMountAdmissionV1, GraphActivationGateV1, GraphSeatGateV1, MountedCodeIndexWorktreeV1, - PendingWakeV1, PublishedTextProjectionOutcomeV1, ServingGenerationSlot, ServingSwapOutcomeV1, - TEXT_PROJECTION_DOCUMENTS_PER_PASS_V1, clear_convergence_park, - clear_graph_resident_memory_park, convergence_park_retries_on_wake, + ColdMountAdmissionV1, ColdMountReservationV1, GraphActivationGateV1, GraphSeatGateV1, + MountedCodeIndexWorktreeV1, PendingWakeV1, PublishedTextProjectionOutcomeV1, + ServingGenerationSlot, ServingSwapOutcomeV1, TEXT_PROJECTION_DOCUMENTS_PER_PASS_V1, + clear_convergence_park, clear_graph_resident_memory_park, convergence_park_retries_on_wake, is_repeated_conflict_verdict, park_convergence, publication_authority_is_terminal, retained_noop_requires_follow_up_wake, }; @@ -152,6 +152,43 @@ impl CodeIndexSchedulerRegistryV1 { ); } + /// Wait until the holder of the store lock that refused a cold open lets + /// go of it. The holder may be another process, so the kernel lock wait + /// is the signal; shutdown or retirement of this reservation ends the wait. + async fn wait_for_cold_open_store_release( + store_root: &Path, + reservation: &ColdMountReservationV1, + ) -> Result<(), CodeIndexSchedulerErrorV1> { + let mut cancellation = reservation.slot.cancellation.subscribe(); + let cancelled = || { + CodeIndexSchedulerErrorV1::Identity(if reservation.slot.is_retired() { + "code-index scheduler owner is still retiring".to_owned() + } else { + "code-index scheduler is shutting down".to_owned() + }) + }; + if reservation.slot.is_cancelled() { + return Err(cancelled()); + } + let (released_tx, released) = tokio::sync::oneshot::channel(); + let root = store_root.to_path_buf(); + std::thread::Builder::new() + .name("code-index-store-release".to_owned()) + .spawn(move || { + let _ = released_tx.send(wait_for_code_generation_store_release(&root)); + })?; + tokio::select! { + released = released => match released { + Ok(Ok(())) => Ok(()), + Ok(Err(error)) => Err(std::io::Error::other(error.to_string()).into()), + Err(_) => Err(CodeIndexSchedulerErrorV1::Identity( + "code-index store release wait ended without a result".to_owned(), + )), + }, + _ = cancellation.changed() => Err(cancelled()), + } + } + /// Spend this mount's single automatic reset on a corrupt publication. /// /// The reset is attempted once per mount so a store that is corrupt again @@ -474,40 +511,60 @@ impl CodeIndexSchedulerRegistryV1 { let scoped_store_root = super::super::scoped_code_index_store_root(&store_root, &project_root); let worker_scope_store_root = scoped_store_root.clone(); - let open_project_id = project_id.clone(); - let open_project_root = project_root.clone(); - let open_byte_pool = Arc::clone(&self.byte_pool); - let open_resident_memory = Arc::clone(&self.resident_memory); - let open_resident_owners = Arc::clone(&self.resident_owners); let progress_daemon_incarnation = self.progress_daemon_incarnation; let progress_producer_incarnation = self.mint_progress_producer_incarnation()?; let mounted_path_policy = path_policy.clone(); - let (opened, cold_mount_reservation) = tokio::task::spawn_blocking(move || { - #[cfg(test)] - Self::pause_cold_mount_open_for_test(&open_project_root); - let opened = CodeIndexWorktreeSchedulerV1::open_with_policy( - open_project_id, - &open_project_root, - scoped_store_root, - open_byte_pool, - CodeIndexHintPolicyV1::default(), - path_policy, - ); - #[cfg(test)] - Self::finish_cold_mount_open_for_test(&open_project_root); - let mut opened = opened?; - opened.bind_resident_memory(open_resident_memory); - opened.bind_resident_owners(open_resident_owners); - opened.bind_progress_incarnations( - progress_daemon_incarnation, - progress_producer_incarnation, - ); - Ok::<_, CodeIndexSchedulerErrorV1>((opened, cold_mount_reservation)) - }) - .await - .map_err(|error| { - CodeIndexSchedulerErrorV1::Identity(format!("code-index mount task failed: {error}")) - })??; + let mut cold_mount_reservation = cold_mount_reservation; + let opened = loop { + let open_project_id = project_id.clone(); + let open_project_root = project_root.clone(); + let open_store_root = scoped_store_root.clone(); + let open_path_policy = path_policy.clone(); + let open_byte_pool = Arc::clone(&self.byte_pool); + let open_resident_memory = Arc::clone(&self.resident_memory); + let open_resident_owners = Arc::clone(&self.resident_owners); + let (opened, reservation) = tokio::task::spawn_blocking(move || { + #[cfg(test)] + Self::pause_cold_mount_open_for_test(&open_project_root); + let opened = CodeIndexWorktreeSchedulerV1::open_with_policy( + open_project_id, + &open_project_root, + open_store_root, + open_byte_pool, + CodeIndexHintPolicyV1::default(), + open_path_policy, + ); + #[cfg(test)] + Self::finish_cold_mount_open_for_test(&open_project_root); + let opened = opened.map(|mut opened| { + opened.bind_resident_memory(open_resident_memory); + opened.bind_resident_owners(open_resident_owners); + opened.bind_progress_incarnations( + progress_daemon_incarnation, + progress_producer_incarnation, + ); + opened + }); + (opened, cold_mount_reservation) + }) + .await + .map_err(|error| { + CodeIndexSchedulerErrorV1::Identity(format!( + "code-index mount task failed: {error}" + )) + })?; + cold_mount_reservation = reservation; + match opened { + Err(error) if error.is_store_lock_contended() => { + Self::wait_for_cold_open_store_release( + &scoped_store_root, + &cold_mount_reservation, + ) + .await?; + } + opened => break opened?, + } + }; let repository_id = opened.identity().repository_id().clone(); let worktree_id = opened.identity().worktree_id().clone(); let reconcile_in_progress = opened.reconcile_in_progress(); diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/reconcile_failure_isolation_tests.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/reconcile_failure_isolation_tests.rs index 90a8abc115..a293535a94 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/reconcile_failure_isolation_tests.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/reconcile_failure_isolation_tests.rs @@ -17,7 +17,9 @@ use std::time::Duration; use tempfile::TempDir; use tracedecay_code_index_retention::code_index_generations::acquire_code_generation_store_lock; use tracedecay_contracts::ResolvedScope; -use tracedecay_contracts::code_index_freshness::CodeIndexStalenessStateV1; +use tracedecay_contracts::code_index_freshness::{ + CodeIndexReadinessTargetV1, CodeIndexReadinessWaitReadV1, CodeIndexStalenessStateV1, +}; use super::super::tests::OwnerSignals; use super::super::{ @@ -487,6 +489,65 @@ async fn a_transient_capacity_refusal_is_retried_without_an_external_wake() { fixture.registry.shutdown().await; } +/// A first mount refused because another owner holds the scope's +/// code-generation store lock opens once that holder lets go, with the release +/// as its only trigger, and its worktree then reaches fresh. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_first_mount_refused_by_a_held_store_lock_mounts_on_its_release() { + let root = TempDir::new().expect("fixture root"); + let project = root.path().join("project"); + fs::create_dir_all(project.join("src")).expect("create source root"); + fs::write(project.join("src/main.rs"), "fn main() {}\n").expect("write source"); + run_git_in(&project, &["init", "-q", "-b", "main"]); + run_git_in(&project, &["add", "."]); + run_git_in(&project, &["commit", "-qm", "fixture"]); + let store = root.path().join("store"); + let canonical = canonical_existing_identity(&project).expect("canonical project"); + let scope_store = super::super::scoped_code_index_store_root(&store, &canonical); + fs::create_dir_all(&scope_store).expect("create the scope store"); + let holder = + acquire_code_generation_store_lock(&scope_store).expect("hold the scope store lock"); + + let registry = CodeIndexSchedulerRegistryV1::with_background_reconcile_permits(1, 1); + registry.install_cold_mount_open_observer(&project); + let mount = tokio::spawn({ + let registry = registry.clone(); + let project = project.clone(); + async move { + registry + .mount_worktree( + tracedecay_domain::ProjectId::new("project.first-mount-store-lock") + .expect("project identity"), + &project, + store, + ) + .await + } + }); + // The first open started and finished while the lock was held. + registry.wait_for_cold_mount_open_events(&project, 2).await; + drop(holder); + + let mounted = tokio::time::timeout(SETTLE_DEADLINE, mount) + .await + .expect("the mount ends once the lock is released") + .expect("mount task"); + assert_eq!(mounted.map_err(|error| error.to_string()), Ok(true)); + let ready = registry + .wait_for_readiness(&project, CodeIndexReadinessTargetV1::Fresh, SETTLE_DEADLINE) + .await + .expect("readiness read"); + let CodeIndexReadinessWaitReadV1::Reached { reading } = ready else { + panic!("the released first mount must reach fresh: {ready:?}"); + }; + assert_eq!( + reading.staleness_state, + Some(CodeIndexStalenessStateV1::Fresh) + ); + assert_eq!(reading.parked, None); + registry.shutdown().await; +} + /// A pass refused because another owner holds this scope's code-generation /// store lock runs again exactly when the holder lets go: no pass repeats the /// refusal while the lock stays held, and the worktree converges after the From 12989e027a51b5e0bf06d22aacacb67b03bc8984 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Thu, 1 Oct 2026 20:03:09 +0000 Subject: [PATCH 2/4] fix(code-index): end every seat wait through one authority --- .../src/code_index_executor.rs | 8 +- .../src/code_index_scheduler.rs | 6 +- .../src/code_index_scheduler/queries.rs | 8 +- .../src/code_index_scheduler/query_runtime.rs | 65 +--- .../src/code_index_scheduler/registry.rs | 3 +- .../registry/convergence_park_tests.rs | 90 ++++- .../registry/ignored_dependencies.rs | 127 +++---- .../code_index_scheduler/registry/mount.rs | 2 +- .../registry/owner_signals.rs | 356 +++++++++++++----- .../registry/serving_reads.rs | 85 +++-- .../tests/deferred_mount_tests.rs | 29 +- .../code_index_scheduler/tests/reconcile.rs | 7 +- .../ignored_dependency_admission.rs | 16 +- .../project_open_owners/advisory_runtime.rs | 1 - .../advisory_runtime/deferred.rs | 219 +---------- .../compiler_diagnostics_producer.rs | 29 +- .../code_index_ignored_dependencies_test.rs | 7 +- 17 files changed, 540 insertions(+), 518 deletions(-) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_executor.rs b/crates/tracedecay-code-index-runtime/src/code_index_executor.rs index 1df59c70f1..a7258a17f0 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_executor.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_executor.rs @@ -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", diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler.rs index ab7f81ad64..0768843e7a 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler.rs @@ -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; diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/queries.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/queries.rs index 2625c05d9e..79ecb9e76e 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/queries.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/queries.rs @@ -335,6 +335,7 @@ impl CodeIndexSchedulerRegistryV1 { ) -> Result { 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); @@ -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 { @@ -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!( @@ -455,6 +456,7 @@ impl CodeIndexSchedulerRegistryV1 { ) -> Result { 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); @@ -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 diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/query_runtime.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/query_runtime.rs index f010150780..3bb2b702c9 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/query_runtime.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/query_runtime.rs @@ -21,7 +21,7 @@ use tracedecay_domain::{ }; use super::{ - CodeIndexReconcileAdmissionV1, CodeIndexSchedulerRegistryV1, + CodeIndexOwnerSignalsV1, CodeIndexReconcileAdmissionV1, CodeIndexSchedulerRegistryV1, serving::CodeTextQueryOwnerReadinessV1, }; use tracedecay_query::retrieval::exact::{ @@ -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( registry: &CodeIndexSchedulerRegistryV1, project_root: PathBuf, @@ -103,19 +98,11 @@ pub async fn retry_deferred_query_authority_until_serving( F: FnMut() -> Fut, Fut: std::future::Future, { - 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 @@ -124,38 +111,8 @@ pub async fn retry_deferred_query_authority_until_serving( { 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; } } } diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry.rs index 4df065fef7..f8f89be7ab 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry.rs @@ -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 diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/convergence_park_tests.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/convergence_park_tests.rs index 9ffb864c24..5d09c31d5d 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/convergence_park_tests.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/convergence_park_tests.rs @@ -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; @@ -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; @@ -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 @@ -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 @@ -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) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/ignored_dependencies.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/ignored_dependencies.rs index 03cb556cfd..264db8a476 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/ignored_dependencies.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/ignored_dependencies.rs @@ -1,6 +1,7 @@ //! Single-flight activation and serving publication for ignored dependencies. use std::collections::BTreeMap; +use std::convert::Infallible; use std::path::Path; use std::sync::{ Arc, Mutex, @@ -12,7 +13,8 @@ use tracedecay_code_index::production::{CodeIndexExecutionControlV1, CodeIndexPr use tracedecay_domain::canonical_sha256; use super::{ - CodeIndexOwnerSignalsV1, CodeIndexSchedulerRegistryV1, PendingWakeV1, ServingGenerationSlot, + CodeIndexSchedulerRegistryV1, CodeIndexSeatParkV1, CodeIndexSeatWaitV1, PendingWakeV1, + ServingGenerationSlot, }; use crate::code_index_scheduler::graph_activation::CodeGraphActivationAuthorityV1; use crate::code_index_scheduler::{ @@ -396,11 +398,14 @@ impl CodeIndexExecutionControlV1 for AdmissionControlBridgeV1 { } impl CodeIndexSchedulerRegistryV1 { + /// `deadline` is the request deadline `control` also reports, as the + /// instant a wait for the decoded seat ends at. pub async fn index_verified_ignored_dependency( &self, project_root: &Path, request: CodeIndexIgnoredDependencyRequestV1, control: Arc, + deadline: tokio::time::Instant, ) -> Result { let project_root = canonical_existing_identity(project_root)?; let flight_key = AdmissionFlightKeyV1::for_request(&request)?; @@ -440,7 +445,7 @@ impl CodeIndexSchedulerRegistryV1 { ) }; if graph_activation_enabled { - self.await_decoded_seat(&project_root, control.as_ref()) + self.await_decoded_seat(&project_root, control.as_ref(), deadline) .await?; } let (flight, owns_flight) = { @@ -495,79 +500,65 @@ impl CodeIndexSchedulerRegistryV1 { /// Admission builds on the decoded generation, but a ready generation /// serves from its text owner before the worker's graph tail seats it (a /// first publication), or with no seat at all until one is demanded. - /// Demand the seat and wait for it on the owner's signals, within the - /// request's budget: an unexpired request never refuses for a seat that - /// is still being installed, and an expired one reports its deadline. - /// - /// A worker parked on a failure installs no seat, so the park is the - /// answer. A park the worker re-checks on every wake answers only once a - /// pass after this demand has observed it again (each observation - /// rewrites it); any other park answers once the worker is back at a - /// wait, since only changed input or an operator remedy lifts it. + /// Demand the seat and wait for it until `deadline`: an unexpired request + /// never refuses for a seat that is still being installed, an expired one + /// reports its deadline, and a parked worker answers with its park. async fn await_decoded_seat( &self, project_root: &Path, control: &(dyn CodeIndexExecutionControlV1 + Send + Sync), + deadline: tokio::time::Instant, ) -> Result<(), CodeIndexSchedulerErrorV1> { - let mut signals = CodeIndexOwnerSignalsV1::subscribe(self, project_root).await; - let Some(activity) = self.subscribe_owner_activity(project_root).await else { - return Err(CodeIndexIgnoredDependencyRefusalV1::Cancelled.into()); - }; - let (serving_generation, convergence_park, shutting_down, pending_wake, wake) = { - let mounted = self.mounted.lock().await; - let Some(worktree) = mounted.get(project_root) else { - return Err(CodeIndexIgnoredDependencyRefusalV1::Cancelled.into()); - }; - ( - Arc::clone(&worktree.serving_generation), - Arc::clone(&worktree.convergence_park), - Arc::clone(&worktree.shutting_down), - Arc::clone(&worktree.pending_wake), - Arc::clone(&worktree.wake), - ) - }; - let seated = || { - serving_generation - .read() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .is_some() - }; - let park = || { - convergence_park - .read() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .clone() - }; - if seated() { - return Ok(()); - } - let park_at_demand = park(); - self.request_complete_generation(project_root).await; - // The demand flag flips once per mount, so this demand posts its own - // wake: a park the worker re-checks answers only after that re-check. - Self::note_wake( - &pending_wake, - &wake, - CodeIndexCadenceTriggerV1::QueryAdmission, - ); - loop { - if seated() { - return Ok(()); + let demanded = &AtomicBool::new(false); + let Ok(waited) = self + .wait_for_seat(project_root, deadline, |_| async move { + let mounted = self.mounted.lock().await; + let Some(worktree) = mounted.get(project_root) else { + return Ok(Some(CodeIndexSeatWaitV1::Cancelled)); + }; + if worktree + .serving_generation + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .is_some() + { + return Ok(Some(CodeIndexSeatWaitV1::Seated(()))); + } + if control.is_cancelled() { + return Ok(Some(CodeIndexSeatWaitV1::Cancelled)); + } + if control.is_deadline_exceeded() { + return Ok(Some(CodeIndexSeatWaitV1::Deadline)); + } + if !demanded.swap(true, Ordering::AcqRel) { + let (pending_wake, wake) = ( + Arc::clone(&worktree.pending_wake), + Arc::clone(&worktree.wake), + ); + drop(mounted); + self.request_complete_generation(project_root).await; + // The demand flag flips once per mount, so this demand + // posts its own wake: a park the worker re-checks answers + // only after that re-check. + Self::note_wake( + &pending_wake, + &wake, + CodeIndexCadenceTriggerV1::QueryAdmission, + ); + } + Ok::<_, Infallible>(None) + }) + .await; + match waited { + CodeIndexSeatWaitV1::Seated(()) => Ok(()), + CodeIndexSeatWaitV1::Parked(CodeIndexSeatParkV1::Convergence(parked)) => { + Err(CodeIndexIgnoredDependencyRefusalV1::ConvergenceParked(parked).into()) } - if let Some(parked) = park() - && activity.pass_finished() - && (!parked.retries_on_wake || park_at_demand.as_ref() != Some(&parked)) - { - return Err(CodeIndexIgnoredDependencyRefusalV1::ConvergenceParked(parked).into()); + CodeIndexSeatWaitV1::Parked(_) | CodeIndexSeatWaitV1::Cancelled => { + Err(CodeIndexIgnoredDependencyRefusalV1::Cancelled.into()) } - refuse_if_interrupted(control, &shutting_down)?; - tokio::select! { - changed = signals.changed() => { - if changed.is_err() { - return Err(CodeIndexIgnoredDependencyRefusalV1::Cancelled.into()); - } - } - () = tokio::time::sleep(Duration::from_millis(5)) => {} + CodeIndexSeatWaitV1::Deadline => { + Err(CodeIndexIgnoredDependencyRefusalV1::DeadlineExceeded.into()) } } } diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs index 51ff4382c2..4beade2f9c 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs @@ -2450,7 +2450,7 @@ impl CodeIndexSchedulerRegistryV1 { // A configuration refusal never re-attempts, so an // unparked refusal read as an indefinite `indexing`. // A memory refusal parks typed until memory is - // given back or its retry delay elapses. + // given back. if error.is_resident_memory_graph_refusal() { park_convergence( &worker_convergence_park, diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs index 1eaf51403b..af68c93243 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs @@ -1,21 +1,23 @@ -//! Change signals for one project root, and the readiness wait built on them. +//! Change signals for one project root, and the one seat wait built on them. +use std::convert::Infallible; use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, PoisonError}; use std::time::Duration; use tracedecay_runtime_core::path_safety::canonical_existing_identity; use tracedecay_contracts::code_index_freshness::{ - CodeIndexFreshnessReadFailureV1, CodeIndexReadinessTargetV1, CodeIndexReadinessV1, - CodeIndexReadinessWaitReadV1, + CodeIndexConvergenceParkedV1, CodeIndexFreshnessReadFailureV1, CodeIndexReadinessTargetV1, + CodeIndexReadinessV1, CodeIndexReadinessWaitReadV1, }; use tracedecay_contracts::ResolvedScope; use super::{ - CodeIndexCadenceTriggerV1, CodeIndexOwnerActivityV1, CodeIndexReconcileAdmissionV1, - CodeIndexSchedulerRegistryV1, unique_mounted_for_scope, + CodeIndexCadenceTriggerV1, CodeIndexGenerationPublishedV1, CodeIndexOwnerActivityV1, + CodeIndexReconcileAdmissionV1, CodeIndexSchedulerRegistryV1, unique_mounted_for_scope, }; use crate::code_index_scheduler::reconcile::FreshnessProbeVerdictV1; use crate::code_index_scheduler::{ @@ -32,20 +34,47 @@ enum CodeIndexFreshSweepRefusedV1 { SweepFailed, } -/// How a search's wait for a restart's retained text serving ended. -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -pub enum CodeIndexRetainedTextServingWaitV1 { - /// The retained generation's exact and lexical owners serve under a - /// mounted query authority. - Serving, - /// The mounted worktree has no durable publication: a first index, with - /// nothing retained to wait for. +/// Why a seat wait ended without the seat it waited for. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum CodeIndexSeatParkV1 { + /// The worker parked on a failure; the park names its cause and remedy. + Convergence(CodeIndexConvergenceParkedV1), + /// The worktree has no durable publication: a first index, with nothing + /// retained to seat. Unpublished, - /// The budget elapsed while the retained text owners were still opening. - Warming, - /// Waiting cannot make them serve: the publication authority is corrupt, - /// the publication read failed, or the registry closed. - Unreachable, + /// The durable publication pointer could not be read. + PublicationUnreadable, + /// A generation is seated, but it does not serve the waiting read. + SeatNotServable, + /// The awaited readiness is unreachable for the named reason. + Unreachable(String), +} + +impl CodeIndexSeatParkV1 { + /// The readiness-wait reason this park reports. + fn readiness_reason(self) -> String { + match self { + Self::Convergence(_) => "code_index_convergence_parked".to_owned(), + Self::Unpublished => "code_index_unpublished".to_owned(), + Self::PublicationUnreadable => "code_index_publication_unreadable".to_owned(), + Self::SeatNotServable => "code_index_seat_not_servable".to_owned(), + Self::Unreachable(reason) => reason, + } + } +} + +/// How one wait for a worktree's seat ended. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum CodeIndexSeatWaitV1 { + /// The awaited seat is in place. + Seated(T), + /// Waiting cannot install the seat. + Parked(CodeIndexSeatParkV1), + /// The registry closed, the worktree is shutting down, or the request + /// was cancelled. + Cancelled, + /// The deadline passed while the seat was still being installed. + Deadline, } /// The registry that owns these channels is gone. @@ -53,8 +82,9 @@ pub enum CodeIndexRetainedTextServingWaitV1 { pub struct CodeIndexOwnerSignalsClosedV1; /// Every signal the registry publishes for one root: registry-wide serving -/// seats and mounts, the worktree's serving-generation changes, its owner -/// passes, pending wake and worker phase, and cadence receipts. +/// seats and mounts, the root's sealed publications, the worktree's +/// serving-generation changes, its owner passes, pending wake and worker +/// phase, and cadence receipts. /// /// Subscribe before the first read. `watch::Sender::subscribe()` marks the /// current value seen, so a change between subscribe and [`Self::changed`] @@ -68,6 +98,7 @@ pub struct CodeIndexOwnerSignalsV1 { seats: tokio::sync::watch::Receiver, root_mounted: tokio::sync::watch::Receiver, receipts: tokio::sync::watch::Receiver, + publications: tokio::sync::broadcast::Receiver, serving: Option>, activity: Option, } @@ -76,15 +107,34 @@ impl CodeIndexOwnerSignalsV1 { pub async fn subscribe(registry: &CodeIndexSchedulerRegistryV1, path: &Path) -> Self { Self { registry: registry.clone(), - path: path.to_path_buf(), + path: canonical_existing_identity(path).unwrap_or_else(|_| path.to_path_buf()), seats: registry.subscribe_serving_seats(), root_mounted: registry.subscribe_root_mounted(), receipts: registry.subscribe_cadence_receipts(), + publications: registry.subscribe_generation_publications(), serving: registry.subscribe_serving_generation_changes(path).await, activity: registry.subscribe_owner_activity(path).await, } } + /// Resolves on a sealed publication of this root; a lagged receiver may + /// have dropped one, so it resolves then too. + async fn root_published( + publications: &mut tokio::sync::broadcast::Receiver, + path: &Path, + ) -> Result<(), CodeIndexOwnerSignalsClosedV1> { + loop { + match publications.recv().await { + Ok(publication) if publication.project_root == path => return Ok(()), + Ok(_) => {} + Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => return Ok(()), + Err(tokio::sync::broadcast::error::RecvError::Closed) => { + return Err(CodeIndexOwnerSignalsClosedV1); + } + } + } + } + /// Resolves once per quiescent burst of publications, so a waiter never /// re-reads, and contends with the worker, on every intra-pass update. pub async fn changed(&mut self) -> Result<(), CodeIndexOwnerSignalsClosedV1> { @@ -111,6 +161,7 @@ impl CodeIndexOwnerSignalsV1 { changed = self.receipts.changed() => { changed.map_err(|_| CodeIndexOwnerSignalsClosedV1)?; } + published = Self::root_published(&mut self.publications, &self.path) => published?, changed = async { match self.serving.as_mut() { Some(serving) => serving.changed().await, @@ -150,6 +201,7 @@ impl CodeIndexOwnerSignalsV1 { self.seats.borrow_and_update(); self.root_mounted.borrow_and_update(); self.receipts.borrow_and_update(); + while self.publications.try_recv().is_ok() {} if let Some(serving) = self.serving.as_mut() { serving.borrow_and_update(); } @@ -161,6 +213,7 @@ impl CodeIndexOwnerSignalsV1 { ] .into_iter() .any(|changed| changed.unwrap_or(false)) + || !self.publications.is_empty() || self .serving .as_ref() @@ -183,6 +236,83 @@ impl CodeIndexOwnerSignalsV1 { } impl CodeIndexSchedulerRegistryV1 { + /// The one wait for a worktree's seat, whatever the seat is: a decoded + /// generation, a current text owner, text serving under a query + /// authority, or a readiness target. + /// + /// It subscribes to the root's signals before the first `probe`, so a + /// seat installed between a probe and the wait still wakes it, and it + /// re-probes only when the registry publishes a change. `probe` is told + /// whether the worker finished its last pass, graph tail included, and + /// answers `Some` once the wait is over and `None` while the seat is still + /// being installed. Between probes the worker's own state ends the wait the same + /// way for every reader: a worktree shutting down is `Cancelled`, and a + /// parked worker installs no seat, so its park is the answer once the + /// worker is back at a wait. A park the worker re-checks on every wake + /// answers only after a pass since the wait began observed it again (each + /// observation rewrites it). An unmounted root is waited through; the + /// probe decides whether that ends the wait. Dropping the future cancels + /// the wait. + pub(crate) async fn wait_for_seat( + &self, + project_root: &Path, + deadline: tokio::time::Instant, + mut probe: impl FnMut(bool) -> Probe, + ) -> Result, E> + where + Probe: Future>, E>>, + { + let mut signals = CodeIndexOwnerSignalsV1::subscribe(self, project_root).await; + let park_at_start = self.convergence_park(project_root).await; + loop { + if let Some(ended) = probe(signals.owner_settled()).await? { + return Ok(ended); + } + if let Some(ended) = self + .seat_worker_end(project_root, &signals, park_at_start.as_ref()) + .await + { + return Ok(ended); + } + match tokio::time::timeout_at(deadline, signals.changed()).await { + Err(_) => return Ok(CodeIndexSeatWaitV1::Deadline), + Ok(Err(CodeIndexOwnerSignalsClosedV1)) => { + return Ok(CodeIndexSeatWaitV1::Cancelled); + } + Ok(Ok(())) => {} + } + } + } + + /// The worker state that ends every seat wait; `None` while the worker + /// may still install a seat. + async fn seat_worker_end( + &self, + project_root: &Path, + signals: &CodeIndexOwnerSignalsV1, + park_at_start: Option<&CodeIndexConvergenceParkedV1>, + ) -> Option> { + let canonical = canonical_existing_identity(project_root).ok()?; + let (park, shutting_down) = { + let mounted = self.mounted.lock().await; + let worktree = mounted.get(&canonical)?; + ( + worktree + .convergence_park + .read() + .unwrap_or_else(PoisonError::into_inner) + .clone(), + worktree.shutting_down.load(Ordering::Acquire), + ) + }; + if shutting_down { + return Some(CodeIndexSeatWaitV1::Cancelled); + } + let parked = park?; + (signals.owner_settled() && (!parked.retries_on_wake || park_at_start != Some(&parked))) + .then(|| CodeIndexSeatWaitV1::Parked(CodeIndexSeatParkV1::Convergence(parked))) + } + /// Sweep the source witness now and post a wake for any proven change, /// so later freshness reads describe the source as of this call. The /// sweep reads only the freshness fence, never the scheduler mutex, so a @@ -246,11 +376,12 @@ impl CodeIndexSchedulerRegistryV1 { /// re-verification wake, and the worker renews the proof without a /// rebuild when the digests still match. Waiting on the registry's /// publications for that renewal keeps a read from answering unavailable - /// in the window. Returns `None` when the scope has no mounted text owner - /// or the registry closed; the caller bounds the wait. + /// in the window. Returns `None` when the scope has no mounted text owner, + /// the wait ends without one, or `deadline` passes. pub(crate) async fn current_text_owner_for_scope( &self, scope: &ResolvedScope, + deadline: tokio::time::Instant, ) -> Option { let root = { let mounted = self.mounted.lock().await; @@ -259,13 +390,24 @@ impl CodeIndexSchedulerRegistryV1 { .0 .clone() }; - let mut signals = CodeIndexOwnerSignalsV1::subscribe(self, &root).await; - loop { - let (latest, current) = self.retained_text_owner_freshness_for_scope(scope).await?; - if current { - return Some(latest); - } - signals.changed().await.ok()?; + let Ok(waited) = self + .wait_for_seat(&root, deadline, |_| async move { + Ok::<_, Infallible>( + match self.retained_text_owner_freshness_for_scope(scope).await { + None => Some(CodeIndexSeatWaitV1::Parked( + CodeIndexSeatParkV1::Unpublished, + )), + Some((latest, true)) => Some(CodeIndexSeatWaitV1::Seated(latest)), + Some((_, false)) => None, + }, + ) + }) + .await; + match waited { + CodeIndexSeatWaitV1::Seated(latest) => Some(latest), + CodeIndexSeatWaitV1::Parked(_) + | CodeIndexSeatWaitV1::Cancelled + | CodeIndexSeatWaitV1::Deadline => None, } } @@ -295,16 +437,16 @@ impl CodeIndexSchedulerRegistryV1 { budget: Duration, ) -> Result { let deadline = tokio::time::Instant::now() + budget; - let mut signals = CodeIndexOwnerSignalsV1::subscribe(self, project_root).await; - let owner_settled = |signals: &CodeIndexOwnerSignalsV1| { - target == CodeIndexReadinessTargetV1::GraphReady || signals.owner_settled() - }; if target != CodeIndexReadinessTargetV1::Fresh && let Some(reading) = self .dashboard_freshness_read(project_root) .await? .filter(|freshness| freshness.readiness(target) == CodeIndexReadinessV1::Reached) - && owner_settled(&signals) + && (target == CodeIndexReadinessTargetV1::GraphReady + || self + .subscribe_owner_activity(project_root) + .await + .is_some_and(|activity| activity.pass_finished())) { return Ok(CodeIndexReadinessWaitReadV1::Reached { reading: Box::new(reading), @@ -336,36 +478,41 @@ impl CodeIndexSchedulerRegistryV1 { }); } } - loop { - let last = self.dashboard_freshness_read(project_root).await?; - if let Some(freshness) = last.as_ref() { - match freshness.readiness(target) { - CodeIndexReadinessV1::Reached if owner_settled(&signals) => { - return Ok(CodeIndexReadinessWaitReadV1::Reached { - reading: Box::new(freshness.clone()), - }); - } - CodeIndexReadinessV1::Reached => {} - CodeIndexReadinessV1::Unreachable { reason } => { - return Ok(CodeIndexReadinessWaitReadV1::Unreachable { reason }); + let waited = self + .wait_for_seat(project_root, deadline, |owner_settled| async move { + let Some(freshness) = self.dashboard_freshness_read(project_root).await? else { + return Ok(None); + }; + Ok(match freshness.readiness(target) { + CodeIndexReadinessV1::Reached + if owner_settled || target == CodeIndexReadinessTargetV1::GraphReady => + { + Some(CodeIndexSeatWaitV1::Seated(freshness)) } - CodeIndexReadinessV1::Pending => {} - } - } - match tokio::time::timeout_at(deadline, signals.changed()).await { - Err(_) => { - return Ok(CodeIndexReadinessWaitReadV1::TimedOut { - last: last.map(Box::new), - }); - } - Ok(Err(CodeIndexOwnerSignalsClosedV1)) => { - return Ok(CodeIndexReadinessWaitReadV1::Unreachable { - reason: "code_index_scheduler_registry_closed".to_owned(), - }); - } - Ok(Ok(())) => {} - } - } + CodeIndexReadinessV1::Reached | CodeIndexReadinessV1::Pending => None, + CodeIndexReadinessV1::Unreachable { reason } => Some( + CodeIndexSeatWaitV1::Parked(CodeIndexSeatParkV1::Unreachable(reason)), + ), + }) + }) + .await?; + Ok(match waited { + CodeIndexSeatWaitV1::Seated(reading) => CodeIndexReadinessWaitReadV1::Reached { + reading: Box::new(reading), + }, + CodeIndexSeatWaitV1::Parked(park) => CodeIndexReadinessWaitReadV1::Unreachable { + reason: park.readiness_reason(), + }, + CodeIndexSeatWaitV1::Cancelled => CodeIndexReadinessWaitReadV1::Unreachable { + reason: "code_index_scheduler_registry_closed".to_owned(), + }, + CodeIndexSeatWaitV1::Deadline => CodeIndexReadinessWaitReadV1::TimedOut { + last: self + .dashboard_freshness_read(project_root) + .await? + .map(Box::new), + }, + }) } /// Wait, for at most `budget`, until the generation a restart retained @@ -384,45 +531,48 @@ impl CodeIndexSchedulerRegistryV1 { project_root: &Path, scope: &ResolvedScope, budget: Duration, - ) -> CodeIndexRetainedTextServingWaitV1 { - let deadline = tokio::time::Instant::now() + budget; - let mut signals = CodeIndexOwnerSignalsV1::subscribe(self, project_root).await; - loop { - match self.has_active_publication(project_root).await { - Some(Ok(true)) => break, - Some(Ok(false)) => return CodeIndexRetainedTextServingWaitV1::Unpublished, - Some(Err(_)) => return CodeIndexRetainedTextServingWaitV1::Unreachable, - None => {} - } - match tokio::time::timeout_at(deadline, signals.changed()).await { - Err(_) => return CodeIndexRetainedTextServingWaitV1::Warming, - Ok(Err(CodeIndexOwnerSignalsClosedV1)) => { - return CodeIndexRetainedTextServingWaitV1::Unreachable; - } - Ok(Ok(())) => {} - } - } - let mut reconcile_requested = false; - loop { - if self.text_serves_search(project_root, scope).await { - return CodeIndexRetainedTextServingWaitV1::Serving; - } - if !reconcile_requested { - reconcile_requested = true; - if let CodeIndexReconcileAdmissionV1::PublicationAuthorityCorrupt(_) = - self.request_query_background_reconcile(scope).await - { - return CodeIndexRetainedTextServingWaitV1::Unreachable; - } - } - match tokio::time::timeout_at(deadline, signals.changed()).await { - Err(_) => return CodeIndexRetainedTextServingWaitV1::Warming, - Ok(Err(CodeIndexOwnerSignalsClosedV1)) => { - return CodeIndexRetainedTextServingWaitV1::Unreachable; - } - Ok(Ok(())) => {} - } - } + ) -> CodeIndexSeatWaitV1<()> { + let published = AtomicBool::new(false); + let reconcile_requested = AtomicBool::new(false); + let (published, reconcile_requested) = (&published, &reconcile_requested); + let Ok(waited) = self + .wait_for_seat( + project_root, + tokio::time::Instant::now() + budget, + |_| async move { + if !published.load(Ordering::Acquire) { + match self.has_active_publication(project_root).await { + Some(Ok(true)) => published.store(true, Ordering::Release), + Some(Ok(false)) => { + return Ok(Some(CodeIndexSeatWaitV1::Parked( + CodeIndexSeatParkV1::Unpublished, + ))); + } + Some(Err(_)) => { + return Ok(Some(CodeIndexSeatWaitV1::Parked( + CodeIndexSeatParkV1::PublicationUnreadable, + ))); + } + None => return Ok(None), + } + } + if self.text_serves_search(project_root, scope).await { + return Ok(Some(CodeIndexSeatWaitV1::Seated(()))); + } + if !reconcile_requested.swap(true, Ordering::AcqRel) { + if let CodeIndexReconcileAdmissionV1::PublicationAuthorityCorrupt(parked) = + self.request_query_background_reconcile(scope).await + { + return Ok(Some(CodeIndexSeatWaitV1::Parked( + CodeIndexSeatParkV1::Convergence(parked), + ))); + } + } + Ok::<_, Infallible>(None) + }, + ) + .await; + waited } /// Whether `scope`'s text owners serve exact and lexical search under a diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/serving_reads.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/serving_reads.rs index 78f74a82c2..53cceba069 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/serving_reads.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/serving_reads.rs @@ -3,6 +3,7 @@ use std::{ collections::BTreeMap, + convert::Infallible, path::{Path, PathBuf}, sync::{Arc, atomic::Ordering}, }; @@ -18,8 +19,8 @@ use super::super::{ use super::graph_cursor_retention::GraphCursorRetentionV1; use super::scope_identity::{latest_matches_scope_identity, text_matches_scope_identity}; use super::{ - CodeIndexMountedScopeV1, CodeIndexOwnerSignalsV1, CodeIndexSchedulerRegistryV1, - CodeIndexServingScopeV1, MountedCodeIndexWorktreeV1, PendingWakeClaimV1, + CodeIndexMountedScopeV1, CodeIndexSchedulerRegistryV1, CodeIndexSeatParkV1, + CodeIndexSeatWaitV1, CodeIndexServingScopeV1, MountedCodeIndexWorktreeV1, PendingWakeClaimV1, ReadyProbeServingPartsV1, dashboard_code_graph_serving, dashboard_freshness_identity, dashboard_terminal_status, dashboard_text_freshness_identity, project_graph_publication_phase, unique_mounted_for_scope, @@ -797,7 +798,7 @@ impl CodeIndexSchedulerRegistryV1 { .await .ok()?; // Parked on resident memory, a reader's pass would only repeat the - // refusal; memory given back or the retry delay wakes the worker. + // refusal; memory given back wakes the worker. if request_reconcile && admission == GenerationDecodeAdmissionV1::AwaitDecode && !memory_retry.waiting() @@ -845,12 +846,12 @@ impl CodeIndexSchedulerRegistryV1 { /// whole decoded generation. A publication seats only its text owner and /// defers the decode until a reader needs it, so this read's own demand /// is what starts that decode: it waits for the seat rather than answering - /// the demanding request unavailable. It stops waiting once no decode is - /// pending (seated, memory-refused, shutting down, or unmounted); the - /// caller's resolution deadline bounds the rest. + /// the demanding request unavailable, until `deadline`. A seat that does + /// not serve this scope, or no published text owner to decode, ends it. pub(crate) async fn latest_complete_fresh_for_scope_awaiting_seat( &self, scope: &tracedecay_contracts::ResolvedScope, + deadline: tokio::time::Instant, ) -> Option { let root = { let mounted = self.mounted.lock().await; @@ -859,43 +860,47 @@ impl CodeIndexSchedulerRegistryV1 { .0 .clone() }; - let mut signals = CodeIndexOwnerSignalsV1::subscribe(self, &root).await; - loop { - if let Some(latest) = self.latest_complete_fresh_for_scope(scope).await { - return Some(latest); - } - if !self.complete_seat_pending(&root).await { - // The seat can land between the miss above and this check. - return self.latest_complete_fresh_for_scope(scope).await; - } - signals.changed().await.ok()?; + let root = &root; + let Ok(waited) = self + .wait_for_seat(root, deadline, |_| async move { + if let Some(latest) = self.latest_complete_fresh_for_scope(scope).await { + return Ok(Some(CodeIndexSeatWaitV1::Seated(latest))); + } + let mounted = self.mounted.lock().await; + let Some(worktree) = mounted.get(root) else { + return Ok(Some(CodeIndexSeatWaitV1::Cancelled)); + }; + let seated = worktree + .serving_generation + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .is_some(); + let text_published = worktree + .text_generation + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .is_some(); + Ok::<_, Infallible>(if seated { + Some(CodeIndexSeatWaitV1::Parked( + CodeIndexSeatParkV1::SeatNotServable, + )) + } else if text_published { + None + } else { + Some(CodeIndexSeatWaitV1::Parked( + CodeIndexSeatParkV1::Unpublished, + )) + }) + }) + .await; + match waited { + CodeIndexSeatWaitV1::Seated(latest) => Some(latest), + CodeIndexSeatWaitV1::Parked(_) + | CodeIndexSeatWaitV1::Cancelled + | CodeIndexSeatWaitV1::Deadline => None, } } - /// Whether demand has asked the worker to seat the complete generation - /// of a published text owner and nothing yet stops it from doing so. - async fn complete_seat_pending(&self, project_root: &Path) -> bool { - let mounted = self.mounted.lock().await; - let Some(worktree) = mounted.get(project_root) else { - return false; - }; - worktree - .complete_generation_requested - .load(Ordering::Acquire) - && !worktree.memory_retry.waiting() - && !worktree.shutting_down.load(Ordering::Acquire) - && worktree - .serving_generation - .read() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .is_none() - && worktree - .text_generation - .read() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .is_some() - } - /// Resolve one exact scope and admit only an already-current generation. pub async fn latest_complete_ready_for_scope( &self, diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/deferred_mount_tests.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/deferred_mount_tests.rs index d9b1cfcb4f..7f0c30caa3 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/deferred_mount_tests.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/deferred_mount_tests.rs @@ -10,7 +10,10 @@ use super::{ALPHA_LIB_V1, GitFixture, wait_for_initial_generation}; use crate::code_index_scheduler::query_runtime::{ DeferredMountAttemptV1, retry_deferred_query_authority_until_serving, }; -use crate::code_index_scheduler::{CodeIndexGenerationPublishedV1, CodeIndexSchedulerRegistryV1}; +use crate::code_index_scheduler::{ + CodeIndexGenerationPublishedV1, CodeIndexOwnerSignalsV1, CodeIndexSchedulerRegistryV1, +}; +use tracedecay_runtime_core::path_safety::canonical_existing_identity; const GENERATION_PUBLICATION_CHANNEL_CAPACITY: usize = 128; @@ -283,3 +286,27 @@ async fn deferred_query_authority_wakes_after_pre_mount_subscribe_on_partitioned ); restarted.shutdown().await; } + +/// A root's owner signals wake on that root's sealed publication and not on +/// another root's, so a waiter re-reads only when its own root may have moved. +#[tokio::test] +async fn owner_signals_wake_on_their_own_roots_publication() { + let root = TempDir::new().expect("project root"); + let foreign = TempDir::new().expect("foreign project root"); + let registry = CodeIndexSchedulerRegistryV1::new(1); + let mut signals = CodeIndexOwnerSignalsV1::subscribe(®istry, root.path()).await; + + registry.push_generation_publication_for_test(synthetic_publication(foreign.path(), 1)); + assert!( + tokio::time::timeout(Duration::from_millis(200), signals.changed()) + .await + .is_err(), + "a foreign root's publication must not wake this root's waiter" + ); + let canonical = canonical_existing_identity(root.path()).expect("canonical project root"); + registry.push_generation_publication_for_test(synthetic_publication(&canonical, 2)); + assert_eq!( + tokio::time::timeout(Duration::from_secs(5), signals.changed()).await, + Ok(Ok(())) + ); +} 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 2ed1911580..d98cb7a276 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 @@ -3928,7 +3928,12 @@ async fn ignored_dependency_waits_for_global_admission_before_publication_gate() let project_root = fixture.path().to_path_buf(); let request_task = tokio::spawn(async move { request_registry - .index_verified_ignored_dependency(&project_root, request, Arc::new(ActiveControl)) + .index_verified_ignored_dependency( + &project_root, + request, + Arc::new(ActiveControl), + tokio::time::Instant::now() + Duration::from_secs(60), + ) .await }); // The publication gate stays held until after this wait, so a request diff --git a/crates/tracedecay-code-index-runtime/src/project_reads/ignored_dependency_admission.rs b/crates/tracedecay-code-index-runtime/src/project_reads/ignored_dependency_admission.rs index bd8b941f15..6b2bec4d03 100644 --- a/crates/tracedecay-code-index-runtime/src/project_reads/ignored_dependency_admission.rs +++ b/crates/tracedecay-code-index-runtime/src/project_reads/ignored_dependency_admission.rs @@ -2,6 +2,7 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; +use std::time::Duration; use tracedecay_application::code_index::{ CodeIndexIgnoredDependencyAdmissionErrorV1, CodeIndexIgnoredDependencyAdmissionFutureV1, @@ -96,9 +97,22 @@ impl CodeIndexIgnoredDependencyAdmissionPortV1 if control.is_deadline_exceeded() { return Err(CodeIndexIgnoredDependencyAdmissionErrorV1::TimedOut); } + let remaining_micros = request + .context() + .deadline() + .expires_at + .0 + .saturating_sub(now_micros().0); + let deadline = tokio::time::Instant::now() + + Duration::from_micros(u64::try_from(remaining_micros).unwrap_or(0)); match self .schedulers - .index_verified_ignored_dependency(&project_root, scheduler_request, control) + .index_verified_ignored_dependency( + &project_root, + scheduler_request, + control, + deadline, + ) .await { Ok(outcome) => Ok(outcome.generation_id), diff --git a/crates/tracedecay/src/daemon/project_open_owners/advisory_runtime.rs b/crates/tracedecay/src/daemon/project_open_owners/advisory_runtime.rs index 001eedbc90..1704fe3aa0 100644 --- a/crates/tracedecay/src/daemon/project_open_owners/advisory_runtime.rs +++ b/crates/tracedecay/src/daemon/project_open_owners/advisory_runtime.rs @@ -115,7 +115,6 @@ use tracedecay_mcp::handlers::hook_runtime::daemon_mint_hook_v2_file_id; mod deferred; mod model; -pub(super) use deferred::wait_for_generation_change; use deferred::{AdvisoryMountPublisherV1, AdvisoryMountStateV1}; pub(crate) use model::ProjectOpenDependentOwnerState; use model::advisory_monotonic_deadline; diff --git a/crates/tracedecay/src/daemon/project_open_owners/advisory_runtime/deferred.rs b/crates/tracedecay/src/daemon/project_open_owners/advisory_runtime/deferred.rs index 38801f8f4d..693f51307f 100644 --- a/crates/tracedecay/src/daemon/project_open_owners/advisory_runtime/deferred.rs +++ b/crates/tracedecay/src/daemon/project_open_owners/advisory_runtime/deferred.rs @@ -1,8 +1,8 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; -use tokio::sync::{broadcast, watch}; -use tracedecay_code_index_runtime::code_index_scheduler::CodeIndexGenerationPublishedV1; +use tokio::sync::watch; +use tracedecay_code_index_runtime::code_index_scheduler::CodeIndexOwnerSignalsV1; use super::super::{project_open_lsp_scope_grant, register_production_lsp_owner}; use super::{ @@ -147,31 +147,20 @@ pub(super) fn spawn( )); return false; } - // Project-open schedules this owner before code-index activation. Capture - // the registry-wide seat and root-mounted cursors now so a retained - // generation seated before the background task's first poll remains - // observable, and so a pre-mount subscribe can re-attach after activation. - // The exact project and scope are still revalidated by `try_mount` after - // every wake. - let mut serving_seats = invocation.code_index_schedulers.subscribe_serving_seats(); - let mut root_mounted = invocation.code_index_schedulers.subscribe_root_mounted(); owner.spawn_background_task(hotpath::future!( async move { - let mut publications = invocation - .code_index_schedulers - .subscribe_generation_publications(); - let mut serving_changes = None; + // Project-open schedules this owner before code-index activation, + // and sealing announces durable source before its text owner is + // installed. Subscribe before the first probe so a generation + // seated or installed after it still wakes this mount; the exact + // project and scope are revalidated by `try_mount` after every wake. + let mut signals = CodeIndexOwnerSignalsV1::subscribe( + &invocation.code_index_schedulers, + &project_root, + ) + .await; let mut partial_publication_retried = false; let settled = loop { - // Sealing announces durable source before its text owner is - // installed. Subscribe before probing that owner so a later - // installation can finish this mount without another edit. - if serving_changes.is_none() { - serving_changes = invocation - .code_index_schedulers - .subscribe_serving_generation_changes(&project_root) - .await; - } publisher.publish(AdvisoryMountStateV1::Mounting); match try_mount(&invocation, &project_root, &mut state).await { Attempt::Settled(settled) => break settled, @@ -189,21 +178,12 @@ pub(super) fn spawn( tracing::info!( event = "advisory_deferred_generation_unavailable", project = %project_root.display(), - serving_watch_registered = serving_changes.is_some(), state = ?awaiting, "waiting after exact complete-generation admission declined" ); } } - if !wait_for_generation_change( - &project_root, - &mut publications, - &mut serving_changes, - &mut serving_seats, - &mut root_mounted, - ) - .await - { + if signals.changed().await.is_err() { return; } partial_publication_retried = false; @@ -230,36 +210,6 @@ async fn awaiting_state( }) } -/// Waits for the next signal that `project_root` may have a new generation: -/// a sealed publication for this root, a serving swap, a retained seat, or a -/// root mount. `false` when the scheduler's publication channel closed. -pub(in crate::daemon::project_open_owners) async fn wait_for_generation_change( - project_root: &Path, - publications: &mut broadcast::Receiver, - serving_changes: &mut Option>, - serving_seats: &mut watch::Receiver, - root_mounted: &mut watch::Receiver, -) -> bool { - loop { - tokio::select! { - publication = publications.recv() => match publication { - Ok(publication) if publication.project_root == project_root => return true, - Ok(_) => {}, - Err(broadcast::error::RecvError::Lagged(_)) => return true, - Err(broadcast::error::RecvError::Closed) => return false, - }, - serving = async { - match serving_changes.as_mut() { - Some(changes) => changes.changed().await, - None => std::future::pending().await, - } - } => return serving.is_ok(), - seat = serving_seats.changed() => return seat.is_ok(), - mounted = root_mounted.changed() => return mounted.is_ok(), - } - } -} - #[derive(Clone, Copy, PartialEq, Eq)] enum Attempt { /// `Mounted` or `Failed`: nothing retries this owner after it returns. @@ -478,146 +428,3 @@ async fn classify_failure( ); attempt } - -#[cfg(test)] -mod tests { - use std::future::Future; - use std::task::{Context, Poll, Waker}; - - use tracedecay_domain::{CodeGenerationId, ContentDigest, RepositoryId}; - - use super::{CodeIndexGenerationPublishedV1, broadcast, wait_for_generation_change, watch}; - - #[tokio::test] - async fn serving_installation_wakes_after_sealed_publication_was_consumed() { - let root = tempfile::tempdir().expect("project root"); - let foreign = tempfile::tempdir().expect("foreign project root"); - let (publication_sender, mut publications) = broadcast::channel(4); - let (serving_sender, serving_receiver) = watch::channel(()); - let mut serving_changes = Some(serving_receiver); - let (_seat_sender, mut serving_seats) = watch::channel(0_u64); - let (_root_sender, mut root_mounted) = watch::channel(0_u64); - let publication = CodeIndexGenerationPublishedV1 { - project_root: root.path().to_path_buf(), - repository_id: RepositoryId::new("repository.deferred").expect("repository"), - generation_id: CodeGenerationId::new("generation.deferred").expect("generation"), - snapshot_content_identity: ContentDigest::new(format!("sha256:{}", "a".repeat(64))) - .expect("content digest"), - observation_time_micros: 1, - }; - publication_sender - .send(publication.clone()) - .expect("sealed publication"); - assert!( - wait_for_generation_change( - root.path(), - &mut publications, - &mut serving_changes, - &mut serving_seats, - &mut root_mounted, - ) - .await - ); - - let mut waiting = Box::pin(wait_for_generation_change( - root.path(), - &mut publications, - &mut serving_changes, - &mut serving_seats, - &mut root_mounted, - )); - let mut context = Context::from_waker(Waker::noop()); - assert!(matches!(waiting.as_mut().poll(&mut context), Poll::Pending)); - publication_sender - .send(CodeIndexGenerationPublishedV1 { - project_root: foreign.path().to_path_buf(), - ..publication - }) - .expect("foreign publication"); - assert!(matches!(waiting.as_mut().poll(&mut context), Poll::Pending)); - - serving_sender.send_replace(()); - assert!(matches!( - waiting.as_mut().poll(&mut context), - Poll::Ready(true) - )); - drop(waiting); - drop(serving_sender); - assert!( - !wait_for_generation_change( - root.path(), - &mut publications, - &mut serving_changes, - &mut serving_seats, - &mut root_mounted, - ) - .await - ); - } - - #[tokio::test] - async fn retained_seat_wakes_after_subscription_precedes_scheduler_enrollment() { - let root = tempfile::tempdir().expect("project root"); - let (publication_sender, mut publications) = broadcast::channel(1); - let (seat_sender, mut serving_seats) = watch::channel(0_u64); - let (_root_sender, mut root_mounted) = watch::channel(0_u64); - - // The deferred owner subscribes while no per-project scheduler exists. - // A retained generation then seats without a new-generation broadcast. - seat_sender.send_modify(|seats| *seats += 1); - assert!( - wait_for_generation_change( - root.path(), - &mut publications, - &mut None, - &mut serving_seats, - &mut root_mounted, - ) - .await - ); - drop(publication_sender); - } - - #[tokio::test] - async fn root_mount_wakes_a_wait_before_per_worktree_subscribe() { - let root = tempfile::tempdir().expect("project root"); - let (_publication_sender, mut publications) = broadcast::channel(1); - let (_seat_sender, mut serving_seats) = watch::channel(0_u64); - let (root_sender, mut root_mounted) = watch::channel(0_u64); - let mut serving_changes = None; - - let mut waiting = Box::pin(wait_for_generation_change( - root.path(), - &mut publications, - &mut serving_changes, - &mut serving_seats, - &mut root_mounted, - )); - let mut context = Context::from_waker(Waker::noop()); - assert!(matches!(waiting.as_mut().poll(&mut context), Poll::Pending)); - root_sender.send_modify(|roots| *roots += 1); - assert!(matches!( - waiting.as_mut().poll(&mut context), - Poll::Ready(true) - )); - } - - #[tokio::test] - async fn publication_channel_closure_stops_a_wait_before_scheduler_mount() { - let root = tempfile::tempdir().expect("project root"); - let (sender, mut publications) = broadcast::channel(1); - let (_seat_sender, mut serving_seats) = watch::channel(0_u64); - let (_root_sender, mut root_mounted) = watch::channel(0_u64); - drop(sender); - assert!( - !wait_for_generation_change( - root.path(), - &mut publications, - &mut None, - &mut serving_seats, - &mut root_mounted, - ) - .await - ); - } -} diff --git a/crates/tracedecay/src/daemon/project_open_owners/compiler_diagnostics_producer.rs b/crates/tracedecay/src/daemon/project_open_owners/compiler_diagnostics_producer.rs index 90f4cc2429..18fea2baaa 100644 --- a/crates/tracedecay/src/daemon/project_open_owners/compiler_diagnostics_producer.rs +++ b/crates/tracedecay/src/daemon/project_open_owners/compiler_diagnostics_producer.rs @@ -16,13 +16,13 @@ use tracedecay_application::diagnostics_producer::{ CompilerProducerRunV1, TypeScriptProducerStateV1, run_typescript_producer_v1, }; use tracedecay_application::diagnostics_store::DiagnosticsStore; +use tracedecay_code_index_runtime::code_index_scheduler::CodeIndexOwnerSignalsV1; use tracedecay_contracts::now_micros; use tracedecay_domain::CodeGenerationId; use tracedecay_lsp::{TypeScriptProject, typescript_projects}; use tracedecay_runtime_core::logging::log_daemon_event; use super::DaemonInvocationState; -use super::advisory_runtime::wait_for_generation_change; /// The project store's directory for tsc's incremental build-info, one file /// per checked tsconfig; it leaves with the store. @@ -70,23 +70,14 @@ pub(super) fn spawn_typescript_diagnostics_producer( ], ); let schedulers = invocation.code_index_schedulers.clone(); - let mut publications = schedulers.subscribe_generation_publications(); - let mut serving_seats = schedulers.subscribe_serving_seats(); - let mut root_mounted = schedulers.subscribe_root_mounted(); - let mut serving_changes = None; + let mut signals = CodeIndexOwnerSignalsV1::subscribe(&schedulers, &project_root).await; let build_info_dir = graph.store_layout().data_root.join(BUILD_INFO_DIR); let mut state = TypeScriptProducerStateV1::default(); - // One producer run per generation, whatever its outcome: the - // wake sources below are registry-wide, and a project whose - // compiler keeps failing must not re-run it on every other - // project's publication. + // One producer run per generation, whatever its outcome: the root's + // signals also wake on serving and pass changes, and a project + // whose compiler keeps failing must not re-run it on each of them. let mut attempted_for: Option = None; loop { - if serving_changes.is_none() { - serving_changes = schedulers - .subscribe_serving_generation_changes(&project_root) - .await; - } // The sealed generation, not the serving one: tsc reads only // the source the seal proved, so the check starts at the seal // rather than after the text projection and serving swap. A @@ -125,15 +116,7 @@ pub(super) fn spawn_typescript_diagnostics_producer( ], ), } - if !wait_for_generation_change( - &project_root, - &mut publications, - &mut serving_changes, - &mut serving_seats, - &mut root_mounted, - ) - .await - { + if signals.changed().await.is_err() { return; } } diff --git a/crates/tracedecay/tests/daemon_suite/code_index_ignored_dependencies_test.rs b/crates/tracedecay/tests/daemon_suite/code_index_ignored_dependencies_test.rs index 2c36bb3cb5..4641b8e894 100644 --- a/crates/tracedecay/tests/daemon_suite/code_index_ignored_dependencies_test.rs +++ b/crates/tracedecay/tests/daemon_suite/code_index_ignored_dependencies_test.rs @@ -303,7 +303,12 @@ async fn index_dependency( control: Arc, ) -> Result { registry - .index_verified_ignored_dependency(project_root, request, control) + .index_verified_ignored_dependency( + project_root, + request, + control, + tokio::time::Instant::now() + Duration::from_secs(60), + ) .await } From f6d47e934ae1b7c021200a9bf85da08f2820870b Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Thu, 1 Oct 2026 20:18:17 +0000 Subject: [PATCH 3/4] simplify(code-index): map seat parks without a reason table --- .../registry/owner_signals.rs | 22 ++++++------------- .../code_index_ignored_dependencies_test.rs | 8 ++----- 2 files changed, 9 insertions(+), 21 deletions(-) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs index af68c93243..aaaf442ca0 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs @@ -50,19 +50,6 @@ pub enum CodeIndexSeatParkV1 { Unreachable(String), } -impl CodeIndexSeatParkV1 { - /// The readiness-wait reason this park reports. - fn readiness_reason(self) -> String { - match self { - Self::Convergence(_) => "code_index_convergence_parked".to_owned(), - Self::Unpublished => "code_index_unpublished".to_owned(), - Self::PublicationUnreadable => "code_index_publication_unreadable".to_owned(), - Self::SeatNotServable => "code_index_seat_not_servable".to_owned(), - Self::Unreachable(reason) => reason, - } - } -} - /// How one wait for a worktree's seat ended. #[derive(Clone, Debug, PartialEq, Eq)] pub enum CodeIndexSeatWaitV1 { @@ -500,8 +487,13 @@ impl CodeIndexSchedulerRegistryV1 { CodeIndexSeatWaitV1::Seated(reading) => CodeIndexReadinessWaitReadV1::Reached { reading: Box::new(reading), }, - CodeIndexSeatWaitV1::Parked(park) => CodeIndexReadinessWaitReadV1::Unreachable { - reason: park.readiness_reason(), + CodeIndexSeatWaitV1::Parked(CodeIndexSeatParkV1::Unreachable(reason)) => { + CodeIndexReadinessWaitReadV1::Unreachable { reason } + } + // Besides the probe's own readiness, only the worker's park ends + // this wait. + CodeIndexSeatWaitV1::Parked(_) => CodeIndexReadinessWaitReadV1::Unreachable { + reason: "code_index_convergence_parked".to_owned(), }, CodeIndexSeatWaitV1::Cancelled => CodeIndexReadinessWaitReadV1::Unreachable { reason: "code_index_scheduler_registry_closed".to_owned(), diff --git a/crates/tracedecay/tests/daemon_suite/code_index_ignored_dependencies_test.rs b/crates/tracedecay/tests/daemon_suite/code_index_ignored_dependencies_test.rs index 4641b8e894..51384d2cb1 100644 --- a/crates/tracedecay/tests/daemon_suite/code_index_ignored_dependencies_test.rs +++ b/crates/tracedecay/tests/daemon_suite/code_index_ignored_dependencies_test.rs @@ -302,13 +302,9 @@ async fn index_dependency( request: CodeIndexIgnoredDependencyRequestV1, control: Arc, ) -> Result { + let deadline = tokio::time::Instant::now() + Duration::from_secs(60); registry - .index_verified_ignored_dependency( - project_root, - request, - control, - tokio::time::Instant::now() + Duration::from_secs(60), - ) + .index_verified_ignored_dependency(project_root, request, control, deadline) .await } From c5137f2249040b1a601fff45e44331d6d58559cb Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Thu, 1 Oct 2026 20:40:30 +0000 Subject: [PATCH 4/4] fix(code-index): wake the publication-gate wait on shutdown, not a poll --- .../src/code_index_scheduler/registry/mount.rs | 9 ++++++--- .../code_index_scheduler/registry/owner_signals.rs | 13 ++++++------- .../src/code_index_scheduler/tests/reconcile.rs | 2 +- .../code_index_ignored_dependencies_test.rs | 2 +- 4 files changed, 14 insertions(+), 12 deletions(-) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs index 4beade2f9c..a3b8bc1ea5 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/mount.rs @@ -7,7 +7,7 @@ use std::{ Arc, Mutex, OnceLock, RwLock, atomic::{AtomicBool, AtomicU64, Ordering}, }, - time::{Duration, Instant}, + time::Instant, }; use tracedecay_code_index::parallelism::collect_installed_worker_heaps_when_idle; @@ -881,6 +881,9 @@ impl CodeIndexSchedulerRegistryV1 { .await; return; }; + // Retirement and daemon shutdown set the flag, then fire this + // watch, so the gate wait below wakes on them. + let mut shutdown_observed = worker_serving_generation_changed.subscribe(); if worker_shutting_down.load(Ordering::Acquire) { tracing::info!( event = "code_index_worker_shutdown_observed", @@ -902,8 +905,8 @@ impl CodeIndexSchedulerRegistryV1 { let _build_publication = loop { tokio::select! { guard = &mut build_publication => break guard, - () = tokio::time::sleep(Duration::from_millis(5)) => { - if worker_shutting_down.load(Ordering::Acquire) { + changed = shutdown_observed.changed() => { + if changed.is_err() || worker_shutting_down.load(Ordering::Acquire) { tracing::info!( event = "code_index_worker_shutdown_observed", phase = "build_publication_lock", diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs index aaaf442ca0..65c7b0e328 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/owner_signals.rs @@ -551,14 +551,13 @@ impl CodeIndexSchedulerRegistryV1 { if self.text_serves_search(project_root, scope).await { return Ok(Some(CodeIndexSeatWaitV1::Seated(()))); } - if !reconcile_requested.swap(true, Ordering::AcqRel) { - if let CodeIndexReconcileAdmissionV1::PublicationAuthorityCorrupt(parked) = + if !reconcile_requested.swap(true, Ordering::AcqRel) + && let CodeIndexReconcileAdmissionV1::PublicationAuthorityCorrupt(parked) = self.request_query_background_reconcile(scope).await - { - return Ok(Some(CodeIndexSeatWaitV1::Parked( - CodeIndexSeatParkV1::Convergence(parked), - ))); - } + { + return Ok(Some(CodeIndexSeatWaitV1::Parked( + CodeIndexSeatParkV1::Convergence(parked), + ))); } Ok::<_, Infallible>(None) }, 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 d98cb7a276..a3d8ff6e17 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 @@ -3932,7 +3932,7 @@ async fn ignored_dependency_waits_for_global_admission_before_publication_gate() &project_root, request, Arc::new(ActiveControl), - tokio::time::Instant::now() + Duration::from_secs(60), + tokio::time::Instant::now() + Duration::from_mins(1), ) .await }); diff --git a/crates/tracedecay/tests/daemon_suite/code_index_ignored_dependencies_test.rs b/crates/tracedecay/tests/daemon_suite/code_index_ignored_dependencies_test.rs index 51384d2cb1..1466b8e97f 100644 --- a/crates/tracedecay/tests/daemon_suite/code_index_ignored_dependencies_test.rs +++ b/crates/tracedecay/tests/daemon_suite/code_index_ignored_dependencies_test.rs @@ -302,7 +302,7 @@ async fn index_dependency( request: CodeIndexIgnoredDependencyRequestV1, control: Arc, ) -> Result { - let deadline = tokio::time::Instant::now() + Duration::from_secs(60); + let deadline = tokio::time::Instant::now() + Duration::from_mins(1); registry .index_verified_ignored_dependency(project_root, request, control, deadline) .await