From 0a1b757bb28d065f0fcf4f14d11e75a7c2c880b7 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 17:53:55 +0000 Subject: [PATCH 1/4] test(daemon): wait on readiness in load-sensitive lib tests Replace test-invented wall-clock bounds and races in daemon and MCP server lib tests with the readiness signal each one is waiting for: armed idle notifications, the production store shutdown before an in-process restart, a fake daemon that stays live, explicit socket shutdown, paused clocks, and the cancellation signal itself. Also read the routed refusal from structuredContent after #2651. --- .../server/project_host_admission_replay.rs | 23 +++--- .../daemon/broker_stream_transport_tests.rs | 22 ++++-- .../http_application_tests/remote_tls.rs | 5 +- .../tracedecay/src/daemon/invocation_state.rs | 6 +- .../tracedecay/src/daemon/tests/lifecycle.rs | 21 ++++++ crates/tracedecay/src/daemon/tests/replay.rs | 13 ++++ .../src/daemon/tests/runtime_identity.rs | 2 +- .../mcp/server/cancel_candidate_journey.rs | 71 +++++-------------- .../tracedecay/src/mcp/server/connection.rs | 7 +- .../src/mcp/server/host_admission_tests.rs | 14 +--- crates/tracedecay/src/mcp/server/routing.rs | 9 +-- 11 files changed, 99 insertions(+), 94 deletions(-) diff --git a/crates/tracedecay-mcp/src/server/project_host_admission_replay.rs b/crates/tracedecay-mcp/src/server/project_host_admission_replay.rs index 851e75b05a..44a836a231 100644 --- a/crates/tracedecay-mcp/src/server/project_host_admission_replay.rs +++ b/crates/tracedecay-mcp/src/server/project_host_admission_replay.rs @@ -109,26 +109,27 @@ impl ProjectHostAdmissionReplayWorker { #[cfg(any(test, feature = "test-transport"))] #[hotpath::skip] - pub async fn wait_idle(&self, timeout: Duration) -> bool { - let deadline = tokio::time::Instant::now() + timeout; + pub async fn wait_idle(&self) { loop { + // `notify_waiters` stores no permit, so both waits are armed before + // the state is read; an idle transition in between still wakes us. + let idle = self.idle.notified(); + let cancelled = self.cancel_notify.notified(); + tokio::pin!(idle, cancelled); + idle.as_mut().enable(); + cancelled.as_mut().enable(); if self.cancel.load(Ordering::Acquire) { - return true; + return; } if !self.busy.load(Ordering::Acquire) && !self.dirty.load(Ordering::Acquire) && !self.broker.has_pending_replay().await { - return true; - } - let remaining = deadline.saturating_duration_since(tokio::time::Instant::now()); - if remaining.is_zero() { - return false; + return; } tokio::select! { - () = self.idle.notified() => {} - () = self.cancel_notify.notified() => return true, - () = tokio::time::sleep(remaining) => return false, + () = idle => {} + () = cancelled => return, } } } diff --git a/crates/tracedecay/src/daemon/broker_stream_transport_tests.rs b/crates/tracedecay/src/daemon/broker_stream_transport_tests.rs index be586be79e..8ce714592f 100644 --- a/crates/tracedecay/src/daemon/broker_stream_transport_tests.rs +++ b/crates/tracedecay/src/daemon/broker_stream_transport_tests.rs @@ -23,6 +23,22 @@ struct DeliverySettlementFixture { _runtime: tracedecay_global_db::tests::harness::RegisteredGlobalDbTestRuntime, } +/// Closes the client socket itself rather than this process's descriptor: a +/// child forked concurrently by another test inherits every open descriptor +/// until it execs, and a plain drop leaves the peer alive through that window. +fn close_client_socket( + reader: tokio::net::unix::OwnedReadHalf, + writer: tokio::net::unix::OwnedWriteHalf, +) { + reader + .reunite(writer) + .expect("client socket halves") + .into_std() + .expect("client socket") + .shutdown(std::net::Shutdown::Both) + .expect("shut down client socket"); +} + async fn delivery_settlement_fixture() -> DeliverySettlementFixture { let profile = tempfile::tempdir().expect("profile"); let project = tempfile::tempdir().expect("project"); @@ -149,8 +165,7 @@ async fn rmcp_receive_waits_for_full_close_after_request_half_close() { "rmcp receive must not treat a request-half close as full peer loss while a response is owed" ); - drop(client_writer); - drop(client_reader); + close_client_socket(client_reader, client_writer); assert!( tokio::time::timeout(std::time::Duration::from_secs(1), &mut receive) .await @@ -533,8 +548,7 @@ async fn rmcp_peer_disconnect_mid_delivery_settles_dropped_rather_than_unknown() ); // The client is gone before the daemon can write its response. - drop(client_reader); - drop(client_writer); + close_client_socket(client_reader, client_writer); let response = serde_json::from_value(serde_json::json!({ "jsonrpc": "2.0", diff --git a/crates/tracedecay/src/daemon/http_application_tests/remote_tls.rs b/crates/tracedecay/src/daemon/http_application_tests/remote_tls.rs index a8404d9bda..b16d09b558 100644 --- a/crates/tracedecay/src/daemon/http_application_tests/remote_tls.rs +++ b/crates/tracedecay/src/daemon/http_application_tests/remote_tls.rs @@ -1566,6 +1566,10 @@ async fn remote_tls_listener_expires_saturated_non_reading_responses() { .expect("bind Remote Brain TLS egress service"); let endpoint = service.remote_tls_endpoint().expect("TLS endpoint"); + // The 5 s idle deadlines must not race 128 handshakes on a loaded host: + // paused time advances only once every task waits, so the bounds below + // catch a stall rather than host scheduling. + tokio::time::pause(); let mut peer_tasks = Vec::with_capacity(128); for _ in 0..128 { let certificate = certificate.clone(); @@ -1605,7 +1609,6 @@ async fn remote_tls_listener_expires_saturated_non_reading_responses() { .expect("every large response must reach real TLS backpressure"); assert_eq!(service.remote_tls_available_admissions(), Some(0)); - tokio::time::pause(); for _ in 0..4 { if service .remote_tls_egress_snapshot() diff --git a/crates/tracedecay/src/daemon/invocation_state.rs b/crates/tracedecay/src/daemon/invocation_state.rs index 299133e71f..a39ac4daf0 100644 --- a/crates/tracedecay/src/daemon/invocation_state.rs +++ b/crates/tracedecay/src/daemon/invocation_state.rs @@ -1279,12 +1279,14 @@ mod shutdown_tests { assert_eq!(receipt.unfinished(), &["invocation"]); } - #[tokio::test] + /// Paused time advances only through timers, so the bound below measures + /// grace the shutdown waited out, not host scheduling. + #[tokio::test(start_paused = true)] async fn cancel_admissions_then_empty_shutdown_is_prompt() { let state = DaemonInvocationState::default(); state.cancel_admissions(); state.cancel_admissions(); - let started = std::time::Instant::now(); + let started = tokio::time::Instant::now(); assert!( state.shutdown().await.is_clean(), "empty invocation shutdown must expire cleanly" diff --git a/crates/tracedecay/src/daemon/tests/lifecycle.rs b/crates/tracedecay/src/daemon/tests/lifecycle.rs index b33bc278cd..51f147e65f 100644 --- a/crates/tracedecay/src/daemon/tests/lifecycle.rs +++ b/crates/tracedecay/src/daemon/tests/lifecycle.rs @@ -1022,6 +1022,7 @@ async fn one_shot_tool_call_allows_long_response_while_daemon_stays_live() { let authority = seed_socket_authority(&socket); let token = authority.auth_token().to_string(); let listener = tokio::net::UnixListener::bind(&socket).expect("bind daemon socket"); + let (client_done, client_done_rx) = tokio::sync::oneshot::channel(); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.expect("accept tool call"); let (reader, mut writer) = stream.into_split(); @@ -1054,6 +1055,7 @@ async fn one_shot_tool_call_allows_long_response_while_daemon_stays_live() { .await .expect("write response"); writer.write_all(b"\n").await.expect("write newline"); + accept_liveness_probes_until(&listener, client_done_rx).await; }); let result = tokio::time::timeout( @@ -1071,6 +1073,7 @@ async fn one_shot_tool_call_allows_long_response_while_daemon_stays_live() { .expect("healthy long-running request timed out") .expect("healthy long-running request must complete"); assert_eq!(result["status"], json!("ok")); + client_done.send(()).expect("fake daemon awaits the client"); server.await.expect("fake daemon task"); } @@ -1082,6 +1085,7 @@ async fn one_shot_tool_call_preserves_response_split_across_liveness_poll() { let authority = seed_socket_authority(&socket); let token = authority.auth_token().to_string(); let listener = tokio::net::UnixListener::bind(&socket).expect("bind daemon socket"); + let (client_done, client_done_rx) = tokio::sync::oneshot::channel(); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.expect("accept tool call"); let (reader, mut writer) = stream.into_split(); @@ -1115,6 +1119,7 @@ async fn one_shot_tool_call_preserves_response_split_across_liveness_poll() { .write_all(&response[split..]) .await .expect("write response suffix"); + accept_liveness_probes_until(&listener, client_done_rx).await; }); let result = tokio::time::timeout( @@ -1132,9 +1137,25 @@ async fn one_shot_tool_call_preserves_response_split_across_liveness_poll() { .expect("split-frame response timed out") .expect("split-frame response must reassemble across liveness polls"); assert_eq!(result["status"], json!("split-across-poll")); + client_done.send(()).expect("fake daemon awaits the client"); server.await.expect("fake daemon task"); } +/// A live daemon keeps accepting after it answers, so the client's liveness +/// probes keep connecting until the client has read the whole response. +#[cfg(unix)] +async fn accept_liveness_probes_until( + listener: &tokio::net::UnixListener, + mut client_done: tokio::sync::oneshot::Receiver<()>, +) { + loop { + tokio::select! { + _ = &mut client_done => return, + probe = listener.accept() => drop(probe.expect("accept liveness probe")), + } + } +} + #[cfg(unix)] #[tokio::test] async fn persistent_idle_client_closes_on_draining_without_timeout() { diff --git a/crates/tracedecay/src/daemon/tests/replay.rs b/crates/tracedecay/src/daemon/tests/replay.rs index 0376c6894b..6c43089056 100644 --- a/crates/tracedecay/src/daemon/tests/replay.rs +++ b/crates/tracedecay/src/daemon/tests/replay.rs @@ -129,6 +129,19 @@ async fn client_identity_startup_replays_retained_profile_receipts() { // stopped, so restart replay remains the acceptance path under test. drop(broker); drop(user_db); + // A restart follows the first daemon's store shutdown: its detached + // session-runtime workers still write the profile database until joined. + first_admin + .prepare_memory_graph_reconciliation_shutdown() + .await + .unwrap() + .shutdown() + .await + .unwrap(); + first_admin + .close_retained_graph_runtimes_for_shutdown() + .await + .unwrap(); drop(first_admin); std::fs::remove_file(&automation_root).unwrap(); diff --git a/crates/tracedecay/src/daemon/tests/runtime_identity.rs b/crates/tracedecay/src/daemon/tests/runtime_identity.rs index c7415c03f8..00ac8b71b7 100644 --- a/crates/tracedecay/src/daemon/tests/runtime_identity.rs +++ b/crates/tracedecay/src/daemon/tests/runtime_identity.rs @@ -293,7 +293,7 @@ async fn concurrent_same_identity_worktrees_keep_exact_server_and_scheduler_bind serde_json::json!(true), "a linked worktree without the watch opt-in must not serve a file listing: {routed_result}" ); - let problem = &routed_result["problem"]; + let problem = &routed_result["structuredContent"]["problem"]; assert_eq!( (&problem["kind"], &problem["code"], &problem["message"]), ( diff --git a/crates/tracedecay/src/mcp/server/cancel_candidate_journey.rs b/crates/tracedecay/src/mcp/server/cancel_candidate_journey.rs index 5395ac840f..cff7dd26d7 100644 --- a/crates/tracedecay/src/mcp/server/cancel_candidate_journey.rs +++ b/crates/tracedecay/src/mcp/server/cancel_candidate_journey.rs @@ -75,7 +75,6 @@ struct PausingAdmission { scan_checkpoints: Arc, pause_at: Arc, scan_paused: Arc, - resume: Arc>>, server: Arc>>>, observed_signal: Arc>>, } @@ -111,26 +110,11 @@ impl CodeIndexMcpReadAdmissionV1 for PausingAdmission { .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(signal); } self.scan_paused.notify_one(); - let resume = self - .resume - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - let deadline = std::time::Instant::now() + Duration::from_secs(5); - loop { - if signal - .as_ref() - .is_some_and(tracedecay_contracts::CancellationSignal::is_cancelled) - { - break; - } - match resume.try_recv() { - Ok(()) | Err(std::sync::mpsc::TryRecvError::Disconnected) => break, - Err(std::sync::mpsc::TryRecvError::Empty) - if std::time::Instant::now() < deadline => - { - std::thread::yield_now(); - } - Err(std::sync::mpsc::TryRecvError::Empty) => break, + // The scan holds this checkpoint until the transport cancel lands, + // however late the canceller runs. + if let Some(signal) = signal { + while !signal.is_cancelled() { + std::thread::yield_now(); } } } @@ -173,7 +157,7 @@ async fn cancelled_large_candidate_search_stops_before_the_next_batch() { authorization_revision: AuthorizationRevision::new("authorization.cancel-journey.fixture") .expect("authorization revision"), }; - let (resume_tx, admission) = pausing_admission(authority.clone()); + let admission = pausing_admission(authority.clone()); let executor = code_index_search_executor( corpus.registry.clone(), ProjectId::new("project.cancel-candidate-journey").expect("corpus project"), @@ -190,7 +174,7 @@ async fn cancelled_large_candidate_search_stops_before_the_next_batch() { .pause_at .store(NEXT_BATCH_CHECKPOINT, Ordering::SeqCst); let pause_at = admission.pause_at.load(Ordering::SeqCst); - drive_rmcp(&held.server, &admission, &resume_tx, pause_at).await; + drive_rmcp(&held.server, &admission, pause_at).await; held.server.shutdown().await; corpus.registry.shutdown().await; @@ -264,21 +248,16 @@ async fn mount_candidate_corpus() -> MountedCorpus { } } -fn pausing_admission( - authority: CodeIndexSearchAuthorityV1, -) -> (std::sync::mpsc::Sender<()>, PausingAdmission) { - let (resume_tx, resume_rx) = std::sync::mpsc::channel(); - let admission = PausingAdmission { +fn pausing_admission(authority: CodeIndexSearchAuthorityV1) -> PausingAdmission { + PausingAdmission { authority, runtime_thread: std::thread::current().id(), scan_checkpoints: Arc::new(AtomicUsize::new(0)), pause_at: Arc::new(AtomicUsize::new(usize::MAX)), scan_paused: Arc::new(tokio::sync::Notify::new()), - resume: Arc::new(StdMutex::new(resume_rx)), server: Arc::new(StdMutex::new(None)), observed_signal: Arc::new(StdMutex::new(None)), - }; - (resume_tx, admission) + } } async fn open_search_server( @@ -326,12 +305,7 @@ async fn open_search_server( } } -async fn drive_rmcp( - server: &Arc, - admission: &PausingAdmission, - _resume_tx: &std::sync::mpsc::Sender<()>, - pause_at: usize, -) { +async fn drive_rmcp(server: &Arc, admission: &PausingAdmission, pause_at: usize) { let adapter = tracedecay_mcp::server::RmcpConnectionAdapter::new( super::connection::ProductionMcpConnectionContext::with_activity(Arc::clone(server), None), false, @@ -487,22 +461,13 @@ fn tools_call_request_id(messages: &StdMutex>) -> rmcp::model::Reques } async fn wait_until_observed_signal_cancels(admission: &PausingAdmission, label: &str) { - tokio::time::timeout(Duration::from_secs(5), async { - loop { - let cancelled = admission - .observed_signal - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .as_ref() - .is_some_and(tracedecay_contracts::CancellationSignal::is_cancelled); - if cancelled { - return; - } - tokio::task::yield_now().await; - } - }) - .await - .unwrap_or_else(|_| panic!("{label}: the transport cancel did not reach the in-flight scan")); + let signal = admission + .observed_signal + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone() + .unwrap_or_else(|| panic!("{label}: the paused scan held no cancellation signal")); + signal.cancelled().await; } async fn wait_for_batch_pause(admission: &PausingAdmission, label: &str) { diff --git a/crates/tracedecay/src/mcp/server/connection.rs b/crates/tracedecay/src/mcp/server/connection.rs index 903633e178..676e3673f4 100644 --- a/crates/tracedecay/src/mcp/server/connection.rs +++ b/crates/tracedecay/src/mcp/server/connection.rs @@ -546,16 +546,15 @@ impl McpServer { #[cfg(test)] #[hotpath::skip] - pub(crate) async fn wait_project_host_admission_replay_idle(&self, timeout: Duration) -> bool { + pub(crate) async fn wait_project_host_admission_replay_idle(&self) { let worker = self .project_host_admission_replay .lock() .await .as_ref() .map(|task| Arc::clone(task.worker())); - match worker { - Some(worker) => worker.wait_idle(timeout).await, - None => true, + if let Some(worker) = worker { + worker.wait_idle().await; } } diff --git a/crates/tracedecay/src/mcp/server/host_admission_tests.rs b/crates/tracedecay/src/mcp/server/host_admission_tests.rs index d61cc4e076..6d5331a47f 100644 --- a/crates/tracedecay/src/mcp/server/host_admission_tests.rs +++ b/crates/tracedecay/src/mcp/server/host_admission_tests.rs @@ -1286,12 +1286,7 @@ async fn durable_route_survives_unavailable_effect_for_same_connection_retry() { | HostAdmissionStatus::Committed | HostAdmissionStatus::ExactDuplicate )); - assert!( - server - .wait_project_host_admission_replay_idle(Duration::from_secs(5)) - .await, - "owned project replay worker should settle the retained admission" - ); + server.wait_project_host_admission_replay_idle().await; assert_eq!(broker.pending_count().await, 0); server.shutdown().await; } @@ -1585,12 +1580,7 @@ async fn owned_project_replay_worker_continues_past_one_bounded_batch() { ) .await; - assert!( - server - .wait_project_host_admission_replay_idle(Duration::from_secs(5)) - .await, - "owned worker must drain a 65-record startup backlog across bounded passes" - ); + server.wait_project_host_admission_replay_idle().await; assert_eq!(broker.pending_count().await, 0); assert!( server.project_host_admission_replay_pass_count().await >= 2, diff --git a/crates/tracedecay/src/mcp/server/routing.rs b/crates/tracedecay/src/mcp/server/routing.rs index 935f447f81..c3a4a601fd 100644 --- a/crates/tracedecay/src/mcp/server/routing.rs +++ b/crates/tracedecay/src/mcp/server/routing.rs @@ -679,19 +679,16 @@ mod tests { params_roots.push(json!({"uri": uri.as_str(), "name": format!("slow-{index}")})); } let params = json!({"roots": params_roots}); - let started = std::time::Instant::now(); let route = resolve_initialize_roots_project_route(Some(¶ms), Some(registry), None).await; - let elapsed = started.elapsed(); for index in 0..3 { tracedecay_runtime_core::git_repository::reset_repository_discovery_for_test( &projects.path().join(format!("slow-{index}")), ); } - assert!( - elapsed < std::time::Duration::from_secs(3), - "three {probe:?} probes must not stack past one 2s budget, took {elapsed:?}" - ); + // Each probe fits one 2 s budget, so per-root budgets would resolve all + // three unregistered roots and answer NotFound. Only a budget shared + // across roots runs out, on the second probe. let Some(crate::mcp::project_route::WorkspaceProjectRoute::Failed(failure)) = route else { panic!("shared budget must defer before every slow root resolves"); }; From 0be3314d47a7ba8c996b49957a30caa551fab263 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 18:28:41 +0000 Subject: [PATCH 2/4] test(daemon): connect TLS egress peers before pausing time Connect the 128 non-reading peers one at a time so none waits out the 5 s request-read deadlines behind the others' handshakes, join the handler barrier on its own readiness, and pause time at the release so only the test's sleeps age the write-idle deadlines. --- .../http_application_tests/remote_tls.rs | 35 ++++++++----------- 1 file changed, 14 insertions(+), 21 deletions(-) diff --git a/crates/tracedecay/src/daemon/http_application_tests/remote_tls.rs b/crates/tracedecay/src/daemon/http_application_tests/remote_tls.rs index b16d09b558..4f92f6ee68 100644 --- a/crates/tracedecay/src/daemon/http_application_tests/remote_tls.rs +++ b/crates/tracedecay/src/daemon/http_application_tests/remote_tls.rs @@ -1566,29 +1566,22 @@ async fn remote_tls_listener_expires_saturated_non_reading_responses() { .expect("bind Remote Brain TLS egress service"); let endpoint = service.remote_tls_endpoint().expect("TLS endpoint"); - // The 5 s idle deadlines must not race 128 handshakes on a loaded host: - // paused time advances only once every task waits, so the bounds below - // catch a stall rather than host scheduling. - tokio::time::pause(); - let mut peer_tasks = Vec::with_capacity(128); - for _ in 0..128 { - let certificate = certificate.clone(); - peer_tasks.push(tokio::spawn(async move { - let mut peer = remote_tls_connect(endpoint, &certificate).await; - peer.write_all(b"GET /remote/egress HTTP/1.1\r\nHost: localhost\r\n\r\n") - .await - .expect("request large TLS response"); - peer.flush().await.expect("flush large TLS request"); - peer - })); - } - tokio::time::timeout(std::time::Duration::from_secs(2), handler_barrier.wait()) - .await - .expect("every large-response handler must reach the egress barrier"); + // One peer at a time finishes its handshake and request right after its + // accept, so no connection waits out the 5 s request-read deadlines + // behind the other 127 handshakes on a loaded host. let mut non_reading_peers = Vec::with_capacity(128); - for task in peer_tasks { - non_reading_peers.push(task.await.expect("join non-reading TLS peer")); + for _ in 0..128 { + let mut peer = remote_tls_connect(endpoint, &certificate).await; + peer.write_all(b"GET /remote/egress HTTP/1.1\r\nHost: localhost\r\n\r\n") + .await + .expect("request large TLS response"); + peer.flush().await.expect("flush large TLS request"); + non_reading_peers.push(peer); } + // Egress starts when this test joins the handler barrier. Pausing at the + // release leaves the 5 s write-idle deadlines to the sleeps below. + handler_barrier.wait().await; + tokio::time::pause(); tokio::time::timeout(std::time::Duration::from_secs(2), async { loop { let snapshot = service From dd6c9534ad12993cd96bba93519f760c35d4bcff Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 19:10:56 +0000 Subject: [PATCH 3/4] test(daemon): await the reindexed generation without a wall clock --- .../generation_retention_test.rs | 20 +++++++++---------- 1 file changed, 9 insertions(+), 11 deletions(-) diff --git a/crates/tracedecay/src/daemon/production_harness/generation_retention_test.rs b/crates/tracedecay/src/daemon/production_harness/generation_retention_test.rs index e21dfe3e03..ecf470330f 100644 --- a/crates/tracedecay/src/daemon/production_harness/generation_retention_test.rs +++ b/crates/tracedecay/src/daemon/production_harness/generation_retention_test.rs @@ -39,18 +39,16 @@ async fn wait_for_changed_generation( project_root: &Path, prior: &CodeGenerationId, ) -> CodeGenerationId { - tokio::time::timeout(Duration::from_secs(20), async { - loop { - if let Some(current) = schedulers.latest_generation_id(project_root).await - && ¤t != prior - { - return current; - } - tokio::time::sleep(Duration::from_millis(10)).await; + // Background reindexing has no deadline; the published generation id is + // the readiness signal. + loop { + if let Some(current) = schedulers.latest_generation_id(project_root).await + && ¤t != prior + { + return current; } - }) - .await - .expect("changed code generation") + tokio::time::sleep(Duration::from_millis(10)).await; + } } async fn publish_code_edit( From b20fb9cee89d3b368b13cdc502f3c852da628ff6 Mon Sep 17 00:00:00 2001 From: ScriptedAlchemy Date: Tue, 29 Sep 2026 19:20:24 +0000 Subject: [PATCH 4/4] test(daemon): keep the saturated fake daemon live for the client --- crates/tracedecay/src/daemon/tests/lifecycle.rs | 3 +++ 1 file changed, 3 insertions(+) diff --git a/crates/tracedecay/src/daemon/tests/lifecycle.rs b/crates/tracedecay/src/daemon/tests/lifecycle.rs index 51f147e65f..93de5f24d9 100644 --- a/crates/tracedecay/src/daemon/tests/lifecycle.rs +++ b/crates/tracedecay/src/daemon/tests/lifecycle.rs @@ -655,6 +655,7 @@ async fn one_shot_tool_call_receives_a_matching_saturation_response() { let socket = temp.path().join("daemon.sock"); let _authority = seed_socket_authority(&socket); let listener = tokio::net::UnixListener::bind(&socket).expect("bind daemon socket"); + let (client_done, client_done_rx) = tokio::sync::oneshot::channel(); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.expect("accept tool call"); super::super::reject_saturated_daemon_client( @@ -666,6 +667,7 @@ async fn one_shot_tool_call_receives_a_matching_saturation_response() { }, ) .await; + accept_liveness_probes_until(&listener, client_done_rx).await; }); let error = tokio::time::timeout( @@ -687,6 +689,7 @@ async fn one_shot_tool_call_receives_a_matching_saturation_response() { message.contains("daemon client capacity reached"), "expected a matching saturation response, got: {message}" ); + client_done.send(()).expect("fake daemon awaits the client"); server.await.expect("saturation server task"); }