From 61b7300c066c57ccdee36fc597a02c3b874ed9cc Mon Sep 17 00:00:00 2001 From: "zackary.l.jackson" Date: Sun, 4 Oct 2026 15:41:36 +0000 Subject: [PATCH 1/2] fix(index): seal an edit's successor during text projection An arrival during a publication's text projection waited for that projection to finish and its generation to seat before the worker reconciled it. The publishing pass now watches pending arrivals while it joins the projection, seals the successor when admission is free, and hands that outcome to the next pass instead of resealing from source. Fixes #2798 Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../src/code_index_scheduler/registry.rs | 11 +++ .../code_index_scheduler/registry/mount.rs | 93 ++++++++++++++++++- .../registry/test_gates.rs | 41 +++++++- .../code_index_scheduler/tests/reconcile.rs | 74 +++++++++++++++ 4 files changed, 214 insertions(+), 5 deletions(-) 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..2ebf5f0ada 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 @@ -736,6 +736,10 @@ 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. 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 = 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 +903,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, @@ -1417,8 +1422,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(); 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 +2687,80 @@ 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 + // pending arrival stays for the next pass, which takes the + // sealed outcome. 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 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); + } + 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" + ), + } + } + 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 +2836,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..3255f63693 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,80 @@ 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"); + + release_projection + .send(()) + .expect("release publication projection"); + let successor = active_pointer().generation_id; + 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) { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("the superseded projection finishes and its successor projects after it"); + 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")]); From 06f9aef112c338fd1962afbf6f13de7b142efbcb Mon Sep 17 00:00:00 2001 From: "zackary.l.jackson" Date: Sun, 4 Oct 2026 16:01:49 +0000 Subject: [PATCH 2/2] fix(index): keep later arrivals pending past an early seal The early seal now claims the arrival it answers and the next pass reuses that arrival, so an edit after the seal stays pending for its own pass instead of being consumed by the presealed outcome. Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- .../code_index_scheduler/registry/mount.rs | 50 ++++++++++++++----- .../code_index_scheduler/tests/reconcile.rs | 20 ++++++-- 2 files changed, 54 insertions(+), 16 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 2ebf5f0ada..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, @@ -737,9 +737,14 @@ impl CodeIndexSchedulerRegistryV1 { // 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. 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 = None; + // 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 @@ -1334,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(); } @@ -1422,7 +1432,7 @@ 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(); + 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 { @@ -2690,8 +2700,9 @@ impl CodeIndexSchedulerRegistryV1 { 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 - // pending arrival stays for the next pass, which takes the - // sealed outcome. Admission is only tried: this pass holds + // 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; @@ -2705,6 +2716,12 @@ impl CodeIndexSchedulerRegistryV1 { 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; @@ -2735,7 +2752,7 @@ impl CodeIndexSchedulerRegistryV1 { "code_index_successor_sealed_during_text_projection", "an arrival during text projection sealed its successor" ); - sealed_successor = Some(outcome); + sealed_successor = Some((outcome, arrival, trigger)); } Ok(Ok(_)) => {} Ok(Err(error)) => tracing::debug!( @@ -2749,6 +2766,13 @@ impl CodeIndexSchedulerRegistryV1 { "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, 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 3255f63693..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 @@ -4167,11 +4167,21 @@ async fn edit_during_text_projection_seals_its_successor_before_the_projection_f }) .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"); - let successor = active_pointer().generation_id; tokio::time::timeout(Duration::from_secs(10), async { loop { let pointer = active_pointer(); @@ -4180,14 +4190,18 @@ async fn edit_during_text_projection_seals_its_successor_before_the_projection_f entry.generation_id == generation_id && entry.text_artifact().is_some() }) }; - if projected(&projecting) && projected(&successor) { + if projected(&projecting) + && projected(&successor) + && pointer.generation_id != successor + && projected(&pointer.generation_id) + { break; } tokio::time::sleep(Duration::from_millis(10)).await; } }) .await - .expect("the superseded projection finishes and its successor projects after it"); + .expect("both projections finish and the later edit seals and projects its own generation"); registry.shutdown().await; }