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 6b984ade20..d4fddefe53 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 @@ -462,6 +462,15 @@ fn published_text_projection_gate() -> &'static Mutex &'static Mutex> { + static GATE: std::sync::OnceLock>> = + std::sync::OnceLock::new(); + GATE.get_or_init(|| Mutex::new(BTreeMap::new())) +} + /// Holds a worker's graph tail right before it seats the decoded generation. #[cfg(test)] fn serving_swap_gate() -> &'static Mutex> { @@ -2683,6 +2692,8 @@ impl CodeIndexSchedulerRegistryV1 { && let Some(opened) = opened.take() { let _ = opened.send(()); + #[cfg(test)] + Self::wait_for_opened_published_text_projection_gate(&project_root).await; } match advance { Ok(Ok(true)) => { 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 e59a954a3b..8efa371b61 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 @@ -21,7 +21,7 @@ use tracedecay_contracts::code_index_freshness::{ use tracedecay_domain::{IndexPathPolicyV1, ProjectId}; use super::super::{ - CodeIndexCadenceTriggerV1, CodeIndexHintPolicyV1, CodeIndexNoopEvidenceV1, + CodeIndexArrivalV1, CodeIndexCadenceTriggerV1, CodeIndexHintPolicyV1, CodeIndexNoopEvidenceV1, CodeIndexReconcileOutcomeV1, CodeIndexSchedulerErrorV1, CodeIndexWorktreeSchedulerV1, DaemonCodeIndexPublicationStoreV1, LatestCodeTextGenerationV1, LatestCompleteCodeIndexV1, RetainedTextGenerationRestoreV1, @@ -736,6 +736,15 @@ impl CodeIndexSchedulerRegistryV1 { // generation that serves text without a native graph to one per // generation, so a permanently unactivatable seal cannot spin. let mut graph_seat_attempted: Option = None; + // A successor sealed while the previous pass joined its text + // projection, with the arrival it answers. The next pass takes it + // in place of its own source reconcile, so an edit never queues + // behind a projection. + let mut sealed_successor: Option<( + CodeIndexReconcileOutcomeV1, + CodeIndexArrivalV1, + CodeIndexCadenceTriggerV1, + )> = None; // A retained revision-7 graph gets one verified-head attempt // before ordinary source reconciliation owns any repair. A failed // verification falls through to one canonical replay of the same @@ -899,6 +908,7 @@ impl CodeIndexSchedulerRegistryV1 { // another actor and would reset if this loop state were lost // while the slot remained. if publication_authority_is_terminal(&worker_convergence_park) { + sealed_successor = None; let _ = Self::take_pending_arrival( &worker_pending_wake, &worker_wake, @@ -1329,11 +1339,16 @@ impl CodeIndexSchedulerRegistryV1 { .is_some_and(LatestCodeTextGenerationV1::query_owners_are_ready); // Admission is held: queue wait ends and service time begins. let started_micros = now_micros().0; - let (arrival, trigger) = Self::take_pending_arrival( - &worker_pending_wake, - &worker_wake, - CodeIndexCadenceTriggerV1::Mount, - ); + // A sealed successor answers the arrival it claimed; a later + // arrival stays pending for its own pass. + let (arrival, trigger) = match sealed_successor.as_ref() { + Some((_, arrival, trigger)) => (*arrival, *trigger), + None => Self::take_pending_arrival( + &worker_pending_wake, + &worker_wake, + CodeIndexCadenceTriggerV1::Mount, + ), + }; if text_slice_incomplete { worker_wake.notify_one(); } @@ -1417,8 +1432,12 @@ impl CodeIndexSchedulerRegistryV1 { let bind_serving_generation = Arc::clone(&worker_serving_generation); let bind_serving_source_witness = Arc::clone(&worker_serving_source_witness); let bind_source_freshness = worker_source_freshness.clone(); + let stashed_successor = sealed_successor.take().map(|(outcome, ..)| outcome); let source_result = tracing::Instrument::instrument( tokio::task::spawn_blocking(move || { + if let Some(outcome) = stashed_successor { + return Ok(outcome); + } let mut scheduler = Self::lock_scheduler_unless_shutting_down(&scheduler, &shutting_down)?; // One arrival per attempted pass, before the branch: the @@ -2678,8 +2697,94 @@ impl CodeIndexSchedulerRegistryV1 { // Join the publication's projection, then process its outcome // at the existing source-proof and serving-swap boundary. Graph // work above overlapped it; nothing seats before text is done. - if let Some(projection) = published_text_projection.take() { - published_text_projection_outcome = Some(match projection.await { + if let Some(mut projection) = published_text_projection.take() { + // An arrival during the projection seals its successor now + // instead of waiting for this projection and its seat. The + // seal claims the arrival, and the next pass answers it + // with the sealed outcome; an arrival after the seal stays + // pending for its own pass. Admission is only tried: this pass holds + // the publication gate, and waiting on admission under it + // inverts the gate order. + let mut preserve_worker_wake = false; + let joined = loop { + if !worker_shutting_down.load(Ordering::Acquire) + && sealed_successor.is_none() + && worker_pending_wake.has_pending_arrival() + && let Some(metadata) = + graph_text.as_ref().map(|text| text.metadata().clone()) + && let Ok(admission) = + Arc::clone(&worker_background_reconcile_admission) + .try_acquire_owned() + { + let (arrival, trigger) = Self::take_pending_arrival( + &worker_pending_wake, + &worker_wake, + CodeIndexCadenceTriggerV1::Mount, + ); + preserve_worker_wake = true; + let scheduler = Arc::clone(&worker_scheduler); + let shutting_down = Arc::clone(&worker_shutting_down); + let retained_text_only = !graph_activation_enabled; + let sealed = tokio::task::spawn_blocking(move || { + let mut scheduler = Self::lock_scheduler_unless_shutting_down( + &scheduler, + &shutting_down, + )?; + match scheduler.reconcile_retained_text_generation_with( + &metadata, + retained_text_only, + )? { + Some(outcome) => Ok(Some(outcome)), + None if retained_text_only => { + scheduler.reconcile_now().map(Some) + } + None => scheduler.activate_or_reconcile().map(Some), + } + }) + .await; + drop(admission); + match sealed { + Ok(Ok(Some( + outcome @ CodeIndexReconcileOutcomeV1::Published(_), + ))) => { + tracing::info!( + event = + "code_index_successor_sealed_during_text_projection", + "an arrival during text projection sealed its successor" + ); + sealed_successor = Some((outcome, arrival, trigger)); + } + Ok(Ok(_)) => {} + Ok(Err(error)) => tracing::debug!( + event = "code_index_successor_seal_during_text_projection_deferred", + error = %error, + "the successor seal waits for the next pass" + ), + Err(error) => tracing::warn!( + event = "code_index_successor_seal_during_text_projection_task_failed", + error = %error, + "the successor seal task failed; the next pass reconciles" + ), + } + if sealed_successor.is_none() { + Self::restore_pending_arrival( + &worker_pending_wake, + arrival, + trigger, + ); + } + } + tokio::select! { + outcome = &mut projection => break outcome, + () = worker_wake.notified() => { + preserve_worker_wake = true; + } + } + }; + if preserve_worker_wake { + worker_wake.notify_one(); + } + published_text_projection_outcome = Some(match joined { Ok(outcome) => outcome, Err(error) => { if let Some(text) = graph_text.as_ref() { @@ -2755,7 +2860,13 @@ impl CodeIndexSchedulerRegistryV1 { // sealed digests are swept again here: a write that // landed during the projection is observed now, // leaves the seat stale, and wakes its successor. - if let Some(text) = graph_text.as_ref() { + // + // A successor sealed during the projection already + // proved the move; re-reconciling this generation + // would race that seal. + if sealed_successor.is_none() + && let Some(text) = graph_text.as_ref() + { let proof_unmoved = worker_source_freshness.serves_verified_source( &text.metadata().snapshot().content_identity, &worker_project_root, diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/test_gates.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/test_gates.rs index 83ab0f6c22..9415934d71 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/test_gates.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/test_gates.rs @@ -18,8 +18,9 @@ use super::{ QueryAdmissionTestControlV1, ServingGenerationInstallationV1, ServingGenerationRollbackOutcomeV1, WorkerStepGateV1, cold_mount_admission_barriers, cold_mount_open_controls, cold_mount_post_check_controls, complete_seat_probe_miss_gate, - graph_decode_gate, published_text_projection_gate, query_admission_controls, serving_swap_gate, - test_gate_root, unique_mounted_for_scope, wait_notified_if_unset, + graph_decode_gate, opened_published_text_projection_gate, published_text_projection_gate, + query_admission_controls, serving_swap_gate, test_gate_root, unique_mounted_for_scope, + wait_notified_if_unset, }; use tracedecay_runtime_core::path_safety::canonical_existing_identity; @@ -50,6 +51,42 @@ impl CodeIndexSchedulerRegistryV1 { (entered_observed, released) } + /// Hold the next publication's text projection once its first advance + /// opened the build, so the publishing pass waits in its join. + #[cfg(test)] + pub async fn pause_next_opened_published_text_projection( + &self, + project_root: PathBuf, + ) -> ( + tokio::sync::oneshot::Receiver<()>, + tokio::sync::oneshot::Sender<()>, + ) { + let (entered, entered_observed) = tokio::sync::oneshot::channel(); + let (released, release) = tokio::sync::oneshot::channel(); + let previous = opened_published_text_projection_gate() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .insert( + test_gate_root(&project_root), + WorkerStepGateV1 { entered, release }, + ); + assert!( + previous.is_none(), + "one opened text projection gate per worktree: {}", + project_root.display() + ); + (entered_observed, released) + } + + #[cfg(test)] + pub(super) async fn wait_for_opened_published_text_projection_gate(project_root: &Path) { + let gate = opened_published_text_projection_gate() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .remove(&test_gate_root(project_root)); + Self::pass_worker_step_gate(gate).await; + } + #[cfg(test)] pub(super) async fn wait_for_published_text_projection_gate(project_root: &Path) { let gate = published_text_projection_gate() diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/reconcile.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/tests/reconcile.rs index b9467bc02e..f9723283b1 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 @@ -4117,6 +4117,94 @@ async fn raw_edit_during_text_projection_is_stale_after_seat_and_reconciles_with registry.shutdown().await; } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn edit_during_text_projection_seals_its_successor_before_the_projection_finishes() { + let fixture = GitFixture::new(ALPHA_LIB_V1); + let store = TempDir::new().expect("store root"); + let registry = CodeIndexSchedulerRegistryV1::with_background_reconcile_permits(1, 1); + let canonical_root = canonical_existing_identity(fixture.path()).expect("canonical fixture"); + let (projection_started, release_projection) = registry + .pause_next_opened_published_text_projection(canonical_root) + .await; + registry + .mount_worktree( + test_project_id(), + fixture.path(), + store.path().to_path_buf(), + ) + .await + .expect("mount worktree"); + tokio::time::timeout(Duration::from_secs(10), projection_started) + .await + .expect("publication did not reach text projection") + .expect("publication projection gate stays armed"); + let scheduler = registry + .scheduler_handle(fixture.path()) + .await + .expect("scheduler handle"); + let active_pointer = || { + scheduler + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .publication + .read_publication_pointer() + .expect("read publication pointer") + .expect("publication pointer") + }; + let projecting = active_pointer().generation_id; + + fixture.edit( + "src/lib.rs", + "pub fn alpha() -> u32 { 1 }\npub fn edited_during_projection() -> u32 { 2 }\n", + ); + registry + .notify_path(fixture.path(), fixture.path().join("src/lib.rs")) + .await; + tokio::time::timeout(Duration::from_secs(10), async { + while active_pointer().generation_id == projecting { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("the edit must seal its successor while the projection is still running"); + let successor = active_pointer().generation_id; + + // An edit after the early seal is not part of that successor and must + // still reach a generation of its own. + fixture.edit( + "src/lib.rs", + "pub fn alpha() -> u32 { 1 }\npub fn edited_during_projection() -> u32 { 2 }\n\ + pub fn edited_after_seal() -> u32 { 3 }\n", + ); + registry + .notify_path(fixture.path(), fixture.path().join("src/lib.rs")) + .await; + release_projection + .send(()) + .expect("release publication projection"); + tokio::time::timeout(Duration::from_secs(10), async { + loop { + let pointer = active_pointer(); + let projected = |generation_id: &str| { + pointer.generation_index.iter().any(|entry| { + entry.generation_id == generation_id && entry.text_artifact().is_some() + }) + }; + if projected(&projecting) + && projected(&successor) + && pointer.generation_id != successor + && projected(&pointer.generation_id) + { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("both projections finish and the later edit seals and projects its own generation"); + registry.shutdown().await; +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn verified_empty_source_remains_observable_while_scheduler_is_busy() { let fixture = GitFixture::new(&[("assets/blob.bin", "not source\n")]);