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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -2503,12 +2503,34 @@ fn branch_add_admits_background_publication_and_remove_retires_its_exact_artifac
"branch list must read durable admission\nstdout:\n{}\nstderr:\n{pending_stderr}",
String::from_utf8_lossy(&pending.stdout)
);
assert!(
pending_stderr
.lines()
.any(|line| line.contains("feature/new") && line.contains("indexing")),
"admitted branch must be durably visible as indexing: {pending_stderr}"
);
// The one-file publication may seal before this listing runs, so the
// admitted branch reads either as pending or as synced, and a synced
// branch serves only once the daemon switches to it. A synced listing
// must already rest on sealed provenance; sealing never reverts.
let admitted = pending_stderr
.lines()
.find(|line| line.starts_with(" feature/new "))
.unwrap_or_else(|| panic!("admitted branch must be durably listed: {pending_stderr}"));
let listed_pending = admitted.contains("indexing");
if listed_pending {
assert!(
admitted.ends_with(", exact index pending"),
"a pending branch must name its pending exact index: {admitted}"
);
} else {
assert!(
(admitted.starts_with(" feature/new [current], ")
|| admitted.starts_with(" feature/new [current, serving], "))
&& admitted.contains(" (from main), synced "),
"admitted branch must read as pending or synced: {admitted}"
);
assert!(
tracedecay_runtime_core::branch_meta::load_branch_meta(&shard_root)
.and_then(|meta| meta.branches.get("feature/new").cloned())
.is_some_and(|entry| entry.graph_source.is_some()),
"a synced listing must rest on sealed provenance: {admitted}"
);
}
let started = Instant::now();
let meta = loop {
if let Some(meta) = tracedecay_runtime_core::branch_meta::load_branch_meta(&shard_root)
Expand All @@ -2525,6 +2547,18 @@ fn branch_add_admits_background_publication_and_remove_retires_its_exact_artifac
);
std::thread::sleep(Duration::from_millis(100));
};
let mut sealed = tracedecay_command_without_daemon(home.path(), &project_root);
sealed.args(["branch", "list"]);
let sealed = run_with_timeout(sealed, cli_timeout());
let sealed_stderr = String::from_utf8_lossy(&sealed.stderr);
assert!(
sealed_stderr
.lines()
.any(|line| line.starts_with(" feature/new [current")
&& !line.contains("indexing")
&& line.contains(" (from main), synced ")),
"a sealed branch must list as synced: {sealed_stderr}"
);
let entry = meta
.branches
.get("feature/new")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -422,15 +422,22 @@ fn retained_text_projection_gate() -> &'static Mutex<BTreeMap<PathBuf, RetainedT
}

#[cfg(test)]
struct PublishedTextProjectionGateV1 {
struct WorkerStepGateV1 {
entered: tokio::sync::oneshot::Sender<()>,
release: tokio::sync::oneshot::Receiver<()>,
}

#[cfg(test)]
fn published_text_projection_gate()
-> &'static Mutex<BTreeMap<PathBuf, PublishedTextProjectionGateV1>> {
static GATE: std::sync::OnceLock<Mutex<BTreeMap<PathBuf, PublishedTextProjectionGateV1>>> =
fn 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>> {
static GATE: std::sync::OnceLock<Mutex<BTreeMap<PathBuf, WorkerStepGateV1>>> =
std::sync::OnceLock::new();
GATE.get_or_init(|| Mutex::new(BTreeMap::new()))
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2464,6 +2464,10 @@ impl CodeIndexSchedulerRegistryV1 {
// nested guard and publishes the witness before either
// lifetime becomes idle.
}
#[cfg(test)]
if matches!(&result, Ok((Ok(_), Some(_), _))) {
Self::wait_for_serving_swap_gate(&worker_project_root).await;
}
if let Ok((Ok(_), Some(latest), _)) = &result {
let scheduler = Arc::clone(&worker_scheduler);
let serving_generation = Arc::clone(&worker_serving_generation);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,14 @@ impl CodeIndexOwnerSignalsV1 {
Ok(())
}

/// Whether the worktree's worker finished its last pass, graph tail
/// included, and is back at a wait.
fn owner_settled(&self) -> bool {
self.activity
.as_ref()
.is_some_and(CodeIndexOwnerActivityV1::pass_finished)
}

/// Consume publications until a scheduler turn passes without one.
async fn settle_burst(&mut self) {
loop {
Expand Down Expand Up @@ -255,6 +263,11 @@ impl CodeIndexSchedulerRegistryV1 {
/// only where a stat moved, without the scheduler mutex. `ready` and
/// `graph_ready` accept a reading that already satisfies them, the answer
/// a plain status read gives; a pending `ready` still sweeps first.
/// `fresh` and `ready` are reached only once the worker has also finished
/// the pass behind that reading: its graph tail seats the decoded
/// generation and binds the source proof to that seat under a counted
/// step, so a wait that ended before the tail was followed by reads
/// reporting `verifying` for the generation it had just reported fresh.
/// An unmounted root is waited through: a mount that lands inside the
/// budget reconciles the source as of that mount. Dropping the future
/// abandons the wait; a wake the sweep posted is ordinary demand.
Expand All @@ -265,17 +278,21 @@ impl CodeIndexSchedulerRegistryV1 {
budget: Duration,
) -> Result<CodeIndexReadinessWaitReadV1, CodeIndexFreshnessReadFailureV1> {
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)
{
return Ok(CodeIndexReadinessWaitReadV1::Reached {
reading: Box::new(reading),
});
}
let mut signals = CodeIndexOwnerSignalsV1::subscribe(self, project_root).await;
// The caller's budget bounds the sweep, and an unproven source cannot
// be reported as reached.
let sweep_source = target != CodeIndexReadinessTargetV1::GraphReady;
Expand Down Expand Up @@ -306,11 +323,12 @@ impl CodeIndexSchedulerRegistryV1 {
let last = self.dashboard_freshness_read(project_root).await?;
if let Some(freshness) = last.as_ref() {
match freshness.readiness(target) {
CodeIndexReadinessV1::Reached => {
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 });
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,10 @@ use super::super::{
use super::{
CodeIndexSchedulerRegistryV1, ColdMountOpenEventV1, ColdMountOpenTestControlV1,
ColdMountPostCheckTestControlV1, PendingWakeDropGateTestV1, PendingWakeV1,
PublishedTextProjectionGateV1, QueryAdmissionTestControlV1, ServingGenerationInstallationV1,
ServingGenerationRollbackOutcomeV1, cold_mount_admission_barriers, cold_mount_open_controls,
cold_mount_post_check_controls, published_text_projection_gate, query_admission_controls,
unique_mounted_for_scope, wait_notified_if_unset,
QueryAdmissionTestControlV1, ServingGenerationInstallationV1,
ServingGenerationRollbackOutcomeV1, WorkerStepGateV1, cold_mount_admission_barriers,
cold_mount_open_controls, cold_mount_post_check_controls, published_text_projection_gate,
query_admission_controls, serving_swap_gate, unique_mounted_for_scope, wait_notified_if_unset,
};
use tracedecay_runtime_core::path_safety::canonical_existing_identity;

Expand All @@ -38,10 +38,7 @@ impl CodeIndexSchedulerRegistryV1 {
.unwrap_or_else(std::sync::PoisonError::into_inner);
assert!(
gates
.insert(
project_root.clone(),
PublishedTextProjectionGateV1 { entered, release },
)
.insert(project_root.clone(), WorkerStepGateV1 { entered, release })
.is_none(),
"one published text projection gate per worktree: {}",
project_root.display()
Expand All @@ -61,6 +58,39 @@ impl CodeIndexSchedulerRegistryV1 {
}
}

/// Hold the next graph tail of the worker for `project_root` right before
/// it seats the decoded generation. The first receiver resolves once the
/// worker waits there; sending on the returned sender releases it.
#[cfg(test)]
pub fn pause_next_serving_swap(
&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 replaced = serving_swap_gate()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(project_root, WorkerStepGateV1 { entered, release });
assert!(replaced.is_none(), "one serving swap gate per worktree");
(entered_observed, released)
}

#[cfg(test)]
pub(super) async fn wait_for_serving_swap_gate(project_root: &Path) {
let gate = serving_swap_gate()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(project_root);
if let Some(gate) = gate {
let _ = gate.entered.send(());
let _ = gate.release.await;
}
}

/// Test-only observation of an exact mounted worktree's active owner pass.
#[cfg(test)]
pub async fn reconcile_in_progress_for_test(&self, project_root: &Path) -> bool {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,10 @@ use std::{

use tempfile::TempDir;
use tracedecay_application::diagnostics_publication::CodeIndexPublicationIdentityPortV1;
use tracedecay_contracts::code_index_freshness::{
CodeIndexReadinessTargetV1, CodeIndexReadinessV1, CodeIndexReadinessWaitReadV1,
CodeIndexStalenessStateV1,
};
use tracedecay_contracts::{
CallableCodeOperationKind, CallableCodeQueryPort, CodeQueryScope, Deadline,
ExactOccurrenceRequest, ResolvedScope, RetrievalPortContext, RetrievalPortOutcome,
Expand Down Expand Up @@ -3085,6 +3089,108 @@ async fn sealed_publication_identity_answers_before_the_generation_seats() {
registry.shutdown().await;
}

/// The text owner serves a first publication before the worker's graph tail
/// seats its decoded generation, and that seat step reads as `verifying`.
/// A `ready` wait must not end in that gap, or the next read contradicts it.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn ready_wait_ends_only_after_the_graph_tail_seats_the_generation() {
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 (swap_entered, release_swap) = registry.pause_next_serving_swap(canonical_root);
registry
.mount_worktree(
test_project_id(),
fixture.path(),
store.path().to_path_buf(),
)
.await
.expect("mount worktree");
assert!(registry.request_complete_generation(fixture.path()).await);
tokio::time::timeout(Duration::from_secs(10), swap_entered)
.await
.expect("publication did not reach its serving swap")
.expect("serving swap gate stays armed");

let before_seat = registry
.dashboard_freshness_read(fixture.path())
.await
.expect("freshness read")
.expect("mounted worktree");
assert_eq!(
(
before_seat.staleness_state,
before_seat.readiness(CodeIndexReadinessTargetV1::Ready),
),
(
Some(CodeIndexStalenessStateV1::Fresh),
CodeIndexReadinessV1::Reached
),
"the text owner already serves the generation: {before_seat:?}"
);
assert!(
registry
.latest_complete_serving_for_test(fixture.path())
.await
.is_none(),
"the decoded generation is not seated yet"
);
let held = registry
.wait_for_readiness(
fixture.path(),
CodeIndexReadinessTargetV1::Ready,
Duration::ZERO,
)
.await
.expect("readiness wait");
assert!(
matches!(held, CodeIndexReadinessWaitReadV1::TimedOut { .. }),
"ready must not be reached while the graph tail is held: {held:?}"
);

release_swap.send(()).expect("release serving swap");
let reached = registry
.wait_for_readiness(
fixture.path(),
CodeIndexReadinessTargetV1::Ready,
Duration::from_secs(10),
)
.await
.expect("readiness wait");
let CodeIndexReadinessWaitReadV1::Reached { reading } = reached else {
panic!("ready after the graph tail: {reached:?}");
};
assert_eq!(
reading.staleness_state,
Some(CodeIndexStalenessStateV1::Fresh)
);
assert_eq!(
registry
.latest_complete_serving_for_test(fixture.path())
.await
.map(|seat| seat
.generation()
.manifest()
.generation_id
.as_str()
.to_owned()),
before_seat.latest_generation_id,
"ready implies the advertised generation is seated"
);
assert_eq!(
registry
.dashboard_freshness_read(fixture.path())
.await
.expect("freshness read")
.expect("mounted worktree")
.staleness_state,
Some(CodeIndexStalenessStateV1::Fresh),
"a read after ready still reads fresh"
);
registry.shutdown().await;
}

/// A publication can finish source capture long before its text artifact is
/// ready. The serving swap must reverify source evidence that landed during
/// that projection, otherwise the exact active generation seats without a
Expand Down
Loading
Loading