diff --git a/crates/tracedecay-cli/tests/core_cli_suite/cli_non_interactive_test.rs b/crates/tracedecay-cli/tests/core_cli_suite/cli_non_interactive_test.rs index 4684f0d986..ed97a812d8 100644 --- a/crates/tracedecay-cli/tests/core_cli_suite/cli_non_interactive_test.rs +++ b/crates/tracedecay-cli/tests/core_cli_suite/cli_non_interactive_test.rs @@ -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) @@ -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") 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 3acceb6a2f..84e3c26ce9 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 @@ -422,15 +422,22 @@ fn retained_text_projection_gate() -> &'static Mutex, release: tokio::sync::oneshot::Receiver<()>, } #[cfg(test)] -fn published_text_projection_gate() --> &'static Mutex> { - static GATE: std::sync::OnceLock>> = +fn published_text_projection_gate() -> &'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> { + static GATE: std::sync::OnceLock>> = std::sync::OnceLock::new(); GATE.get_or_init(|| Mutex::new(BTreeMap::new())) } 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 c89f095368..d131c09fec 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 @@ -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); 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 7730563c1f..0c9d432b5e 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 @@ -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 { @@ -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. @@ -265,17 +278,21 @@ 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) { 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; @@ -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 }); } 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 6ee364d406..a0b22dd636 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 @@ -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; @@ -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() @@ -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 { 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 369a6c7616..78450746e1 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 @@ -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, @@ -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 diff --git a/crates/tracedecay-daemon-service/src/invocation/recovery_schedule.rs b/crates/tracedecay-daemon-service/src/invocation/recovery_schedule.rs index bef8597109..f9fc4194b8 100644 --- a/crates/tracedecay-daemon-service/src/invocation/recovery_schedule.rs +++ b/crates/tracedecay-daemon-service/src/invocation/recovery_schedule.rs @@ -183,11 +183,12 @@ mod tests { use std::io; use std::io::Write; use std::sync::Arc; - use std::sync::Mutex; use std::sync::atomic::{AtomicUsize, Ordering}; + use std::sync::{Mutex, Once}; use std::time::Duration; use tracedecay_runtime_core::cancellation::CancellationToken; + use tracing::subscriber::NoSubscriber; use tracing_subscriber::fmt::MakeWriter; use super::{ @@ -427,8 +428,22 @@ mod tests { task.await.expect("failing recovery task"); } - #[test] - fn identical_failure_logs_are_byte_and_event_bounded() { + /// Runs `f` under a subscriber scoped to this thread and returns what it + /// wrote. + /// + /// tracing-core caches each callsite's interest process-wide. While only + /// one dispatcher is registered, a callsite first reached on another + /// thread takes its interest from that thread's default, so a parallel + /// test emitting the same event without a subscriber caches `never` and + /// the capture silently loses it. A registered global dispatcher makes + /// every later registration consult the registered set, which includes + /// this capture. + fn captured_tracing(f: impl FnOnce()) -> String { + static REGISTERED_GLOBAL: Once = Once::new(); + REGISTERED_GLOBAL.call_once(|| { + tracing::subscriber::set_global_default(NoSubscriber::default()) + .expect("no other global subscriber in the daemon-service tests"); + }); let bytes = Arc::new(Mutex::new(Vec::new())); let subscriber = tracing_subscriber::fmt() .with_max_level(tracing::Level::TRACE) @@ -438,18 +453,23 @@ mod tests { bytes: Arc::clone(&bytes), }) .finish(); - tracing::subscriber::with_default(subscriber, || { + tracing::subscriber::with_default(subscriber, f); + let bytes = bytes + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone(); + String::from_utf8(bytes).expect("captured tracing is UTF-8") + } + + #[test] + fn identical_failure_logs_are_byte_and_event_bounded() { + let output = captured_tracing(|| { let mut failures = RecoveryFailureTrackerV1::default(); for _ in 0..100 { failures.fail("fixture recovery", "same fixture failure".to_owned()); } failures.recover("fixture recovery"); }); - let bytes = bytes - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .clone(); - let output = String::from_utf8(bytes).expect("captured tracing is UTF-8"); assert_eq!( output.matches("durable recovery attempt failed").count(), @@ -471,27 +491,13 @@ mod tests { #[test] fn skipped_receipt_logs_are_byte_and_event_bounded() { - let bytes = Arc::new(Mutex::new(Vec::new())); - let subscriber = tracing_subscriber::fmt() - .with_max_level(tracing::Level::TRACE) - .without_time() - .with_ansi(false) - .with_writer(CapturedWriter { - bytes: Arc::clone(&bytes), - }) - .finish(); - tracing::subscriber::with_default(subscriber, || { + let output = captured_tracing(|| { let mut warnings = RecoveryWarningTrackerV1::default(); for _ in 0..100 { warnings.warn("fixture recovery", "invalid fixture receipt"); } warnings.recover("fixture recovery"); }); - let bytes = bytes - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .clone(); - let output = String::from_utf8(bytes).expect("captured tracing is UTF-8"); assert_eq!( output.matches("durable recovery receipt skipped").count(), diff --git a/crates/tracedecay/src/daemon/http_application_tests.rs b/crates/tracedecay/src/daemon/http_application_tests.rs index c174fd2a49..048879774c 100644 --- a/crates/tracedecay/src/daemon/http_application_tests.rs +++ b/crates/tracedecay/src/daemon/http_application_tests.rs @@ -1043,17 +1043,53 @@ async fn daemon_http_timed_out_cold_resolution_preserves_curate_request_identity async fn daemon_http_shutdown_releases_loopback_listener() { let (service, _) = service_with_probe().await; let endpoint = service.endpoint(); + #[cfg(target_os = "linux")] + let listener = match listening_socket_inodes(endpoint).as_slice() { + [inode] => *inode, + inodes => panic!("the service must hold one listener at {endpoint}: {inodes:?}"), + }; service.shutdown().await.expect("shutdown HTTP service"); - // Rebinding the exact address proves the listener was released. A raw - // connect probe can false-positive on a freed ephemeral port when the - // kernel self-connects (source port == destination port) or when a - // parallel test rebinds the port first. + // A freed ephemeral port goes back to the kernel pool, where any process + // may bind it, so neither rebinding nor connecting to it proves release. + // The listener's socket inode does: no other socket can take it over. + #[cfg(target_os = "linux")] + assert!( + !listening_socket_inodes(endpoint).contains(&listener), + "the daemon HTTP listener at {endpoint} must be closed on shutdown" + ); + #[cfg(not(target_os = "linux"))] tokio::net::TcpListener::bind(endpoint) .await .expect("released daemon HTTP loopback address must be rebindable"); } +/// Inodes of the IPv4 sockets listening at `endpoint`, read from the kernel +/// socket table of this network namespace. +#[cfg(target_os = "linux")] +fn listening_socket_inodes(endpoint: SocketAddr) -> Vec { + const TCP_LISTEN: &str = "0A"; + let SocketAddr::V4(endpoint) = endpoint else { + panic!("the daemon HTTP listener binds IPv4 loopback, not {endpoint}"); + }; + // The kernel prints the network-order address as a native-endian word. + let local = format!( + "{:08X}:{:04X}", + u32::from_ne_bytes(endpoint.ip().octets()), + endpoint.port() + ); + std::fs::read_to_string("/proc/net/tcp") + .expect("kernel TCP socket table") + .lines() + .skip(1) + .filter_map(|row| { + let fields = row.split_whitespace().collect::>(); + (fields[1] == local && fields[3] == TCP_LISTEN) + .then(|| fields[9].parse().expect("socket inode")) + }) + .collect() +} + #[tokio::test] async fn daemon_http_shutdown_marks_registry_inactive() { let registry = DaemonHttpApplicationRegistry::default(); diff --git a/crates/tracedecay/tests/mcp_suite/status_behavior_test.rs b/crates/tracedecay/tests/mcp_suite/status_behavior_test.rs index 337fcd19a3..857cbce82a 100644 --- a/crates/tracedecay/tests/mcp_suite/status_behavior_test.rs +++ b/crates/tracedecay/tests/mcp_suite/status_behavior_test.rs @@ -97,27 +97,29 @@ fn parse_status(text: &str) -> Value { }) } +/// The JSON status read held by `wait_for` until the sealed generation is +/// ready, the one readiness authority agent hosts wait on. async fn sealed_json_status( harness: &ProductionProjectCompositionHarnessV1, project_root: &Path, ) -> Value { - let started = Instant::now(); - let mut last = Value::Null; - while started.elapsed() < Duration::from_secs(20) { - let payload = - parse_status(&call_status(harness, project_root, json!({ "format": "json" })).await); - let freshness = &payload["code_index_freshness"]; - let graph = &freshness["worktree"]["code_graph_serving"]; - if freshness["status"] == "current" && graph["state"] == "ready" { - return payload; - } - if graph["state"] == "refused" || graph["reason"] == "activation_disabled" { - panic!("code index refused to serve: {payload}"); - } - last = payload; - tokio::time::sleep(Duration::from_millis(50)).await; - } - panic!("tracedecay_status did not report a current sealed generation: {last}"); + let payload = parse_status( + &call_status( + harness, + project_root, + json!({ + "format": "json", + "wait_for": { "state": "ready", "timeout_ms": 20_000 }, + }), + ) + .await, + ); + assert_eq!( + payload["wait"], + json!({ "outcome": "reached" }), + "tracedecay_status did not report a ready sealed generation: {payload}" + ); + payload } #[tokio::test]