Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -462,6 +462,15 @@ fn published_text_projection_gate() -> &'static Mutex<BTreeMap<PathBuf, WorkerSt
GATE.get_or_init(|| Mutex::new(BTreeMap::new()))
}

/// Holds a publication's text projection after its first advance opened the
/// build, while the pass that published it waits to join it.
#[cfg(test)]
fn opened_published_text_projection_gate() -> &'static Mutex<BTreeMap<PathBuf, WorkerStepGateV1>> {
static GATE: std::sync::OnceLock<Mutex<BTreeMap<PathBuf, WorkerStepGateV1>>> =
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<BTreeMap<PathBuf, WorkerStepGateV1>> {
Expand Down Expand Up @@ -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)) => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<tracedecay_domain::CodeGenerationId> = 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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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();
}
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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() {
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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")]);
Expand Down
Loading