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
6 changes: 3 additions & 3 deletions .repository-projection.json
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,11 @@
"projection": "deixic-code",
"projectionSchemaVersion": 1,
"sourceRepository": "dx-corp/mono",
"sourceSha": "0dda4b314d332ad6f95f189d10b066d1a080ee34",
"sourceSha": "373dd517beb6babba8f831ede8501edef7cd0520",
"destinationRepository": "dx-corp/code",
"priorProjectedBase": "8991ebcc446033952c75e0858f46245ab8b0580a",
"priorProjectedBase": "d30b1b01dee74cc2a51d546423e51e80ed66f31f",
"definitionDigest": "82936441c776e3e8edb5d215a75007ec9714a233f489d460075d79d5ef5ba32f",
"toolDigest": "c244d99199a7ae3eb8ff644a99462163c23b0bb6a83ef50af01efbdca0b81d04",
"contentDigest": "93b1fddc5de9161c81eb617fc50c86e6f9bb877a46f1a1206a46d8732b6b39f8",
"contentDigest": "67278d7f71d18a40a13eb757faff25c35ae2c1e05c6f963b74ac7d6da00953b9",
"publicationEligible": true
}
331 changes: 4 additions & 327 deletions packages/local-host-rs/src/hosted_runner/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,10 @@ use crate::hosted_runner::rendezvous_protocol::RendezvousMode;

#[cfg(unix)]
mod deferred_rejection;
#[cfg(unix)]
mod replay_ownership;
#[cfg(unix)]
use replay_ownership::governed_response_for_ack_test;
mod env_aliases;
mod transport_flush;

Expand Down Expand Up @@ -11281,105 +11285,6 @@ async fn queued_response_restarts_pending_and_consumes_once_after_child_exit() {
.shutdown();
}

#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn correlated_protocol_rejection_rolls_back_ownership_and_allows_retry() {
let workspace = tempdir().expect("workspace");
let fixtures = tempdir().expect("fixtures");
let log_path = fixtures.path().join("rejected-responses.log");
let script = create_reject_then_accept_script(fixtures.path(), &log_path, None);
let supervisor = connected_supervisor_for_script(&script).await;
let executor = Arc::new(AgentSupervisorHostedRunnerMessageExecutor::new(Arc::clone(
&supervisor,
)));
let handle = start_hosted_runner_with_message_executor(
test_config(workspace.path().to_path_buf()),
executor.clone(),
)
.await
.expect("hosted runner");
let client = reqwest::Client::new();
let (capability, subscription_id) =
attach_thread_controller(&client, &handle.base_url(), "conn_rejection").await;
let headers = response_headers(
"conn_rejection",
&subscription_id,
&capability,
"rejected-key",
);
let response = ToAgentMessage::ToolResponse {
call_id: "retry-call".to_string(),
tool_execution_id: Some("retry-execution".to_string()),
approved: true,
result: None,
};

let error = match handle_message(
handle.shared.clone(),
"sess_test",
headers.clone(),
response.clone(),
)
.await
{
Err(error) => error,
Ok(_) => panic!("correlated protocol rejection must not return success"),
};
assert_eq!(error.code, HostedRunnerErrorCode::RuntimeFailed);
assert!(error.message.contains("not awaiting a decision"));
{
let state = handle
.shared
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
assert!(
!state
.pending_response_idempotency
.contains_key("rejected-key")
);
}
assert!(
!executor
.queued_responses
.lock()
.expect("queued responses")
.contains_key("rejected-key")
);
assert!(
!load_executor_response_ledger(workspace.path(), "sess_test")
.expect("response ledger")
.iter()
.any(|(key, _)| key == "rejected-key")
);

handle_message(
handle.shared.clone(),
"sess_test",
headers.clone(),
response.clone(),
)
.await
.expect("corrected retry dispatches");
let replay = handle_message(handle.shared.clone(), "sess_test", headers, response)
.await
.expect("accepted retry replays");
let ResponseBody::Json { body, .. } = replay else {
panic!("accepted retry replay must return JSON");
};
assert_eq!(body["replayed"], true);
assert_eq!(
std::fs::read_to_string(&log_path)
.expect("rejected response log")
.lines()
.count(),
2,
"one rejected dispatch and one corrected dispatch are expected"
);
handle.shutdown().await;
supervisor.lock().expect("supervisor").shutdown();
}

#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn rejected_response_releases_ownership_when_ledger_cleanup_fails() {
Expand Down Expand Up @@ -11585,234 +11490,6 @@ async fn unmatched_governed_ack_does_not_finalize_or_redispatch() {
supervisor.lock().expect("supervisor").shutdown();
}

#[cfg(unix)]
async fn assert_unique_protocol_request_owner_across_restart(
message: ToAgentMessage,
message_type: &str,
request_id: &str,
) {
let workspace = tempdir().expect("workspace");
let fixtures = tempdir().expect("fixtures");
let first_log = fixtures.path().join(format!("{message_type}-first.log"));
let first_script = create_delayed_identity_ack_script(
fixtures.path(),
&format!("{message_type}-first.sh"),
&first_log,
message_type,
request_id,
);
let first_supervisor = connected_supervisor_for_script(&first_script).await;
let first_executor = Arc::new(AgentSupervisorHostedRunnerMessageExecutor::new(Arc::clone(
&first_supervisor,
)));
let first = start_hosted_runner_with_message_executor(
test_config(workspace.path().to_path_buf()),
first_executor,
)
.await
.expect("first hosted runner");
let client = reqwest::Client::new();
let (capability, subscription_id) =
attach_thread_controller(&client, &first.base_url(), "conn_identity_first").await;
let owner_headers = response_headers(
"conn_identity_first",
&subscription_id,
&capability,
"identity-owner-key",
);
let competing_headers = response_headers(
"conn_identity_first",
&subscription_id,
&capability,
"identity-competing-key",
);

handle_message(
first.shared.clone(),
"sess_test",
owner_headers.clone(),
message.clone(),
)
.await
.expect("owner response queues");
let conflict = match handle_message(
first.shared.clone(),
"sess_test",
competing_headers.clone(),
message.clone(),
)
.await
{
Err(error) => error,
Ok(_) => panic!("second key must not own the same protocol request"),
};
assert_eq!(conflict.code, HostedRunnerErrorCode::IdempotencyConflict);
tokio::time::sleep(Duration::from_millis(300)).await;
let replay = handle_message(
first.shared.clone(),
"sess_test",
owner_headers,
message.clone(),
)
.await
.expect("owner key replays after delayed acknowledgement");
let ResponseBody::Json { body, .. } = replay else {
panic!("owner replay must return JSON");
};
assert_eq!(body["replayed"], true);
assert_eq!(
std::fs::read_to_string(&first_log)
.expect("first child log")
.lines()
.count(),
1
);
first.shutdown().await;
first_supervisor
.lock()
.expect("first supervisor")
.shutdown();

let second_log = fixtures.path().join(format!("{message_type}-second.log"));
let second_script = create_delayed_identity_ack_script(
fixtures.path(),
&format!("{message_type}-second.sh"),
&second_log,
message_type,
request_id,
);
let second_supervisor = connected_supervisor_for_script(&second_script).await;
let second_executor = Arc::new(AgentSupervisorHostedRunnerMessageExecutor::new(Arc::clone(
&second_supervisor,
)));
let second = start_hosted_runner_with_message_executor(
test_config(workspace.path().to_path_buf()),
second_executor,
)
.await
.expect("restarted hosted runner");
let second_connection_id = if matches!(
&message,
ToAgentMessage::GovernedClientToolResult(tool_wire::GovernedClientToolResult { .. })
) {
"conn_identity_first"
} else {
"conn_identity_second"
};
let (capability, subscription_id) =
attach_thread_controller(&client, &second.base_url(), second_connection_id).await;
let replay = handle_message(
second.shared.clone(),
"sess_test",
response_headers(
second_connection_id,
&subscription_id,
&capability,
"identity-owner-key",
),
message.clone(),
)
.await
.expect("durable owner key replays");
let ResponseBody::Json { body, .. } = replay else {
panic!("durable owner replay must return JSON");
};
assert_eq!(body["replayed"], true);
let conflict = match handle_message(
second.shared.clone(),
"sess_test",
response_headers(
second_connection_id,
&subscription_id,
&capability,
"identity-competing-key",
),
message,
)
.await
{
Err(error) => error,
Ok(_) => panic!("request ownership must survive restart"),
};
assert_eq!(conflict.code, HostedRunnerErrorCode::IdempotencyConflict);
assert!(
!second_log.exists(),
"restart must not redispatch either key"
);
second.shutdown().await;
second_supervisor
.lock()
.expect("second supervisor")
.shutdown();
}

#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn delayed_tool_response_has_one_idempotency_owner_across_restart() {
assert_unique_protocol_request_owner_across_restart(
ToAgentMessage::ToolResponse {
call_id: "unique-tool-call".to_string(),
tool_execution_id: Some("unique-tool-execution".to_string()),
approved: true,
result: None,
},
"tool_response",
"unique-tool-call",
)
.await;
}

#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn delayed_server_request_response_has_one_idempotency_owner_across_restart() {
assert_unique_protocol_request_owner_across_restart(
ToAgentMessage::ServerRequestResponse {
request_id: "unique-server-request".to_string(),
request_type: ServerRequestType::UserInput,
approved: None,
result: None,
content: Some(Vec::new()),
is_error: Some(false),
decision_action: None,
reason: Some("answer".to_string()),
},
"server_request_response",
"unique-server-request",
)
.await;
}

#[cfg(unix)]
fn governed_response_for_ack_test() -> ToAgentMessage {
ToAgentMessage::GovernedClientToolResult(tool_wire::GovernedClientToolResult {
process_tool_cost_micros: None,
call_id: "unique-governed-call".to_string(),
content: Vec::new(),
is_error: false,
tool_execution_id: "unique-governed-execution".to_string(),
client_instance_id: "conn_identity_first".to_string(),
grant_id: "grant-1".to_string(),
grant_version: 1,
grant_hash: "hash".to_string(),
turn_digest: "turn-digest".to_string(),
definition_digest: "definition-digest".to_string(),
args_digest: "args-digest".to_string(),
owner_lease_epoch: 1,
idempotency_key: "identity-owner-key".to_string(),
})
}

#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn delayed_governed_result_has_one_idempotency_owner_across_restart() {
assert_unique_protocol_request_owner_across_restart(
governed_response_for_ack_test(),
"governed_client_tool_result",
"unique-governed-call",
)
.await;
}

#[cfg(unix)]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn delayed_ack_event_pump_finalizes_before_restart_without_redispatch() {
Expand Down
Loading
Loading