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
23 changes: 12 additions & 11 deletions crates/tracedecay-mcp/src/server/project_host_admission_replay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
}
}
Expand Down
22 changes: 18 additions & 4 deletions crates/tracedecay/src/daemon/broker_stream_transport_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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",
Expand Down
32 changes: 14 additions & 18 deletions crates/tracedecay/src/daemon/http_application_tests/remote_tls.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1566,25 +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");

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
Expand All @@ -1605,7 +1602,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()
Expand Down
6 changes: 4 additions & 2 deletions crates/tracedecay/src/daemon/invocation_state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
&& &current != 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
&& &current != prior
{
return current;
}
})
.await
.expect("changed code generation")
tokio::time::sleep(Duration::from_millis(10)).await;
}
}

async fn publish_code_edit(
Expand Down
24 changes: 24 additions & 0 deletions crates/tracedecay/src/daemon/tests/lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -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(
Expand All @@ -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");
}

Expand Down Expand Up @@ -1022,6 +1025,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();
Expand Down Expand Up @@ -1054,6 +1058,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(
Expand All @@ -1071,6 +1076,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");
}

Expand All @@ -1082,6 +1088,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();
Expand Down Expand Up @@ -1115,6 +1122,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(
Expand All @@ -1132,9 +1140,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() {
Expand Down
13 changes: 13 additions & 0 deletions crates/tracedecay/src/daemon/tests/replay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();

Expand Down
2 changes: 1 addition & 1 deletion crates/tracedecay/src/daemon/tests/runtime_identity.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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"]),
(
Expand Down
Loading
Loading