#[cfg(feature = "test-utils")]
mod tests {
use std::time::Duration;
use alien_commands::{
runtime::{decode_params, parse_envelope},
test_utils::{
dispatcher::MockDispatcherAssertions, server::TestCommandServerAssertions, *,
},
types::*,
};
use alien_core::{MessagePayload, QueueMessage};
use chrono::Utc;
#[tokio::test]
async fn test_core_push_small_params_small_response() {
let server = TestCommandServer::new().await;
let request = test_inline_create_command("push-agent", "generate-report");
let response = server.create_command(request).await.unwrap();
assert_eq!(response.state, CommandState::Dispatched); assert!(response.storage_upload.is_none());
let mock_dispatcher = server
.mock_dispatcher()
.expect("Should have mock dispatcher");
mock_dispatcher.assert_has_dispatched().await;
let dispatched = mock_dispatcher.get_latest().await.unwrap();
assert_eq!(dispatched.envelope.command_id, response.command_id);
assert_eq!(dispatched.envelope.command, "generate-report");
assert!(matches!(
dispatched.envelope.params,
BodySpec::Inline { .. }
));
let params = decode_params(&dispatched.envelope).await.unwrap();
assert!(params.is_object());
let agent_response = test_success_response(b"report generated");
server
.submit_command_response(&response.command_id, agent_response)
.await
.unwrap();
let final_status = server
.wait_for_completion(&response.command_id, Duration::from_secs(5))
.await
.unwrap();
assert_eq!(final_status.state, CommandState::Succeeded);
let final_response = final_status.response.unwrap();
assert!(final_response.is_success());
if let CommandResponse::Success { response: body } = final_response {
assert_inline_body(&body, b"report generated");
}
}
#[tokio::test]
async fn test_core_pull_small_params_small_response() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let request = test_inline_create_command("pull-agent", "sync-data");
let create_response = server.create_command(request).await.unwrap();
assert_eq!(create_response.state, CommandState::Pending);
let lease = server
.acquire_single_lease("pull-agent")
.await
.unwrap()
.unwrap();
assert_eq!(lease.command_id, create_response.command_id);
assert!(matches!(lease.envelope.params, BodySpec::Inline { .. }));
let status = server
.get_command_status(&create_response.command_id)
.await
.unwrap();
assert_eq!(status.state, CommandState::Dispatched);
let params = decode_params(&lease.envelope).await.unwrap();
assert!(params.is_object());
let agent_response = test_json_success_response(&serde_json::json!({
"status": "synced",
"command_id": lease.command_id
}));
server
.submit_command_response(&lease.command_id, agent_response)
.await
.unwrap();
let final_status = server
.wait_for_completion(&create_response.command_id, Duration::from_secs(5))
.await
.unwrap();
assert_eq!(final_status.state, CommandState::Succeeded);
let response = final_status.response.unwrap();
assert!(response.is_success());
if let CommandResponse::Success { response: body } = response {
let body_data = body.decode_inline().unwrap();
let json: serde_json::Value = serde_json::from_slice(&body_data).unwrap();
assert_eq!(json["status"], "synced");
}
}
#[tokio::test]
async fn test_core_push_large_params_large_response() {
let server = TestCommandServer::new().await;
let large_params = vec![b'X'; 160000]; let request = test_storage_create_command("push-agent", "process-bulk", large_params.len());
let response = server.create_command(request).await.unwrap();
assert_eq!(response.state, CommandState::PendingUpload);
assert!(response.storage_upload.is_some());
let storage_upload = response.storage_upload.unwrap();
storage_upload
.put_request
.execute(Some(large_params.clone().into()))
.await
.unwrap();
let upload_complete = test_upload_complete_request(160000);
server
.upload_complete(&response.command_id, upload_complete)
.await
.unwrap();
let mock_dispatcher = server
.mock_dispatcher()
.expect("Should have mock dispatcher");
assert!(mock_dispatcher.has_dispatched().await);
let dispatched = mock_dispatcher.get_latest().await.unwrap();
assert_eq!(dispatched.envelope.command_id, response.command_id);
assert!(matches!(
dispatched.envelope.params,
BodySpec::Storage { .. }
));
let large_response_data = vec![b'R'; 160000];
dispatched
.envelope
.response_handling
.storage_upload_request
.execute(Some(large_response_data.clone().into()))
.await
.unwrap();
let agent_response = CommandResponse::success_storage(large_response_data.len() as u64);
server
.submit_command_response(&response.command_id, agent_response)
.await
.unwrap();
let final_status = server
.wait_for_completion(&response.command_id, Duration::from_secs(5))
.await
.unwrap();
assert_eq!(final_status.state, CommandState::Succeeded);
let final_response = final_status.response.unwrap();
assert!(final_response.is_success());
if let CommandResponse::Success { response: body } = final_response {
assert_storage_body(&body, Some(160000));
assert_storage_body_content(&body, &large_response_data).await;
}
assert!(server.storage_object_count().await > 1);
}
#[tokio::test]
async fn test_core_pull_large_params_large_response() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let large_params = vec![b'Y'; 160000]; let request = test_storage_create_command("pull-agent", "bulk-process", large_params.len());
let response = server.create_command(request).await.unwrap();
assert_eq!(response.state, CommandState::PendingUpload);
assert!(response.storage_upload.is_some());
let storage_upload = response.storage_upload.unwrap();
storage_upload
.put_request
.execute(Some(large_params.clone().into()))
.await
.unwrap();
let upload_complete = test_upload_complete_request(160000);
let upload_response = server
.upload_complete(&response.command_id, upload_complete)
.await
.unwrap();
assert_eq!(upload_response.state, CommandState::Pending);
let lease = server
.acquire_single_lease("pull-agent")
.await
.unwrap()
.unwrap();
assert_eq!(lease.command_id, response.command_id);
assert!(matches!(lease.envelope.params, BodySpec::Storage { .. }));
let status = server
.get_command_status(&response.command_id)
.await
.unwrap();
assert_eq!(status.state, CommandState::Dispatched);
let params_bytes = alien_commands::runtime::decode_params_bytes(&lease.envelope)
.await
.unwrap();
assert_eq!(params_bytes.len(), 160000); assert_eq!(params_bytes, vec![b'Y'; 160000]);
let large_response_data = vec![b'Z'; 160000];
lease
.envelope
.response_handling
.storage_upload_request
.execute(Some(large_response_data.clone().into()))
.await
.unwrap();
let agent_response = CommandResponse::success_storage(large_response_data.len() as u64);
server
.submit_command_response(&lease.command_id, agent_response)
.await
.unwrap();
let final_status = server
.wait_for_completion(&response.command_id, Duration::from_secs(5))
.await
.unwrap();
assert_eq!(final_status.state, CommandState::Succeeded);
let final_response = final_status.response.unwrap();
assert!(final_response.is_success());
if let CommandResponse::Success { response: body } = final_response {
assert_storage_body(&body, Some(160000));
assert_storage_body_content(&body, &large_response_data).await;
}
assert!(server.storage_object_count().await > 0);
}
#[tokio::test]
async fn test_core_push_small_params_large_response() {
let server = TestCommandServer::new().await;
let request = test_inline_create_command("push-agent", "generate-large-report");
let response = server.create_command(request).await.unwrap();
assert_eq!(response.state, CommandState::Dispatched); assert!(response.storage_upload.is_none());
let mock_dispatcher = server
.mock_dispatcher()
.expect("Should have mock dispatcher");
mock_dispatcher.assert_has_dispatched().await;
let dispatched = mock_dispatcher.get_latest().await.unwrap();
assert_eq!(dispatched.envelope.command_id, response.command_id);
assert!(matches!(
dispatched.envelope.params,
BodySpec::Inline { .. }
));
let large_response_data = vec![b'L'; 160000];
dispatched
.envelope
.response_handling
.storage_upload_request
.execute(Some(large_response_data.clone().into()))
.await
.unwrap();
let agent_response = CommandResponse::success_storage(large_response_data.len() as u64);
server
.submit_command_response(&response.command_id, agent_response)
.await
.unwrap();
let final_status = server
.wait_for_completion(&response.command_id, Duration::from_secs(5))
.await
.unwrap();
assert_eq!(final_status.state, CommandState::Succeeded);
let final_response = final_status.response.unwrap();
assert!(final_response.is_success());
if let CommandResponse::Success { response: body } = final_response {
assert_storage_body(&body, Some(160000));
assert_storage_body_content(&body, &large_response_data).await;
}
assert!(server.storage_object_count().await > 0);
}
#[tokio::test]
async fn test_core_push_large_params_small_response() {
let server = TestCommandServer::new().await;
let large_params = vec![b'X'; 160000]; let request = test_storage_create_command("push-agent", "process-data", large_params.len());
let response = server.create_command(request).await.unwrap();
assert_eq!(response.state, CommandState::PendingUpload);
assert!(response.storage_upload.is_some());
let storage_upload = response.storage_upload.unwrap();
storage_upload
.put_request
.execute(Some(large_params.clone().into()))
.await
.unwrap();
let upload_complete = test_upload_complete_request(160000);
server
.upload_complete(&response.command_id, upload_complete)
.await
.unwrap();
let mock_dispatcher = server
.mock_dispatcher()
.expect("Should have mock dispatcher");
assert!(mock_dispatcher.has_dispatched().await);
let dispatched = mock_dispatcher.get_latest().await.unwrap();
assert_eq!(dispatched.envelope.command_id, response.command_id);
assert!(matches!(
dispatched.envelope.params,
BodySpec::Storage { .. }
));
let agent_response = test_success_response(b"ok");
server
.submit_command_response(&response.command_id, agent_response)
.await
.unwrap();
let final_status = server
.wait_for_completion(&response.command_id, Duration::from_secs(5))
.await
.unwrap();
assert_eq!(final_status.state, CommandState::Succeeded);
let final_response = final_status.response.unwrap();
assert!(final_response.is_success());
if let CommandResponse::Success { response: body } = final_response {
assert_inline_body(&body, b"ok");
}
assert!(server.storage_object_count().await > 0);
}
#[tokio::test]
async fn test_core_pull_small_params_large_response() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let request = test_inline_create_command("pull-agent", "generate-large");
let create_response = server.create_command(request).await.unwrap();
assert_eq!(create_response.state, CommandState::Pending);
let lease = server
.acquire_single_lease("pull-agent")
.await
.unwrap()
.unwrap();
assert_eq!(lease.command_id, create_response.command_id);
assert!(matches!(lease.envelope.params, BodySpec::Inline { .. }));
let status = server
.get_command_status(&create_response.command_id)
.await
.unwrap();
assert_eq!(status.state, CommandState::Dispatched);
let params = decode_params(&lease.envelope).await.unwrap();
assert!(params.is_object());
let large_response_data = vec![b'M'; 160000];
lease
.envelope
.response_handling
.storage_upload_request
.execute(Some(large_response_data.clone().into()))
.await
.unwrap();
let agent_response = CommandResponse::success_storage(large_response_data.len() as u64);
server
.submit_command_response(&lease.command_id, agent_response)
.await
.unwrap();
let final_status = server
.wait_for_completion(&create_response.command_id, Duration::from_secs(5))
.await
.unwrap();
assert_eq!(final_status.state, CommandState::Succeeded);
let response = final_status.response.unwrap();
assert!(response.is_success());
if let CommandResponse::Success { response: body } = response {
assert_storage_body(&body, Some(160000));
assert_storage_body_content(&body, &large_response_data).await;
}
assert!(server.storage_object_count().await > 0);
}
#[tokio::test]
async fn test_core_pull_large_params_small_response() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let large_params = vec![b'Y'; 160000]; let request = test_storage_create_command("pull-agent", "process-bulk", large_params.len());
let response = server.create_command(request).await.unwrap();
assert_eq!(response.state, CommandState::PendingUpload);
assert!(response.storage_upload.is_some());
let storage_upload = response.storage_upload.unwrap();
storage_upload
.put_request
.execute(Some(large_params.clone().into()))
.await
.unwrap();
let upload_complete = test_upload_complete_request(160000);
let upload_response = server
.upload_complete(&response.command_id, upload_complete)
.await
.unwrap();
assert_eq!(upload_response.state, CommandState::Pending);
let lease = server
.acquire_single_lease("pull-agent")
.await
.unwrap()
.unwrap();
assert_eq!(lease.command_id, response.command_id);
assert!(matches!(lease.envelope.params, BodySpec::Storage { .. }));
let status = server
.get_command_status(&response.command_id)
.await
.unwrap();
assert_eq!(status.state, CommandState::Dispatched);
let params_bytes = alien_commands::runtime::decode_params_bytes(&lease.envelope)
.await
.unwrap();
assert_eq!(params_bytes.len(), 160000); assert_eq!(params_bytes, vec![b'Y'; 160000]);
let agent_response = test_json_success_response(&serde_json::json!({
"status": "processed",
"command_id": lease.command_id
}));
server
.submit_command_response(&lease.command_id, agent_response)
.await
.unwrap();
let final_status = server
.wait_for_completion(&response.command_id, Duration::from_secs(5))
.await
.unwrap();
assert_eq!(final_status.state, CommandState::Succeeded);
let final_response = final_status.response.unwrap();
assert!(final_response.is_success());
if let CommandResponse::Success { response: body } = final_response {
let body_data = body.decode_inline().unwrap();
let json: serde_json::Value = serde_json::from_slice(&body_data).unwrap();
assert_eq!(json["status"], "processed");
}
assert!(server.storage_object_count().await > 0);
}
#[tokio::test]
async fn test_status_and_envelope_carry_resolved_target() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let request = test_inline_create_command("target-agent", "targeted-command");
let response = server.create_command(request).await.unwrap();
let status = server
.get_command_status(&response.command_id)
.await
.unwrap();
assert_eq!(status.target, server.default_target);
let lease = server
.acquire_single_lease("target-agent")
.await
.unwrap()
.unwrap();
assert_eq!(lease.envelope.target, server.default_target);
}
#[tokio::test]
async fn test_create_with_unknown_target_rejected() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let mut request = test_inline_create_command("target-agent", "targeted-command");
request.target_resource_id = Some("no-such-resource".to_string());
let err = server.create_command(request).await.unwrap_err();
assert_eq!(err.code, "COMMAND_TARGET_NOT_FOUND");
}
#[tokio::test]
async fn test_create_shorthand_with_two_targets_ambiguous() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
server
.registry
.register_target("second-daemon", CommandTargetType::Daemon)
.await
.unwrap();
let request = test_inline_create_command("target-agent", "targeted-command");
let err = server.create_command(request).await.unwrap_err();
assert_eq!(err.code, "COMMAND_TARGET_AMBIGUOUS");
}
#[tokio::test]
async fn test_lease_scans_only_requesting_targets_prefix() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
server
.registry
.register_target("second-daemon", CommandTargetType::Daemon)
.await
.unwrap();
let second_target = CommandTarget::new("second-daemon", CommandTargetType::Daemon);
let mut request_a = test_inline_create_command("target-agent", "for-default");
request_a.target_resource_id = Some(server.default_target.resource_id.clone());
let cmd_a = server.create_command(request_a).await.unwrap();
let mut request_b = test_inline_create_command("target-agent", "for-second");
request_b.target_resource_id = Some("second-daemon".to_string());
let cmd_b = server.create_command(request_b).await.unwrap();
let default_leases = server
.acquire_lease(
"target-agent",
LeaseRequest {
deployment_id: "target-agent".to_string(),
target: server.default_target.clone(),
max_leases: 10,
lease_seconds: 60,
},
)
.await
.unwrap();
assert_eq!(default_leases.leases.len(), 1);
assert_eq!(default_leases.leases[0].command_id, cmd_a.command_id);
assert_eq!(
default_leases.leases[0].envelope.target,
server.default_target
);
let second_leases = server
.acquire_lease(
"target-agent",
LeaseRequest {
deployment_id: "target-agent".to_string(),
target: second_target.clone(),
max_leases: 10,
lease_seconds: 60,
},
)
.await
.unwrap();
assert_eq!(second_leases.leases.len(), 1);
assert_eq!(second_leases.leases[0].command_id, cmd_b.command_id);
assert_eq!(second_leases.leases[0].envelope.target, second_target);
}
#[tokio::test]
async fn test_expired_lease_redelivers_with_incremented_attempt() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let mut request = test_inline_create_command("expiry-agent", "expiry-redelivery");
request.target_resource_id = Some(server.default_target.resource_id.clone());
let created = server.create_command(request).await.unwrap();
let lease_request = LeaseRequest {
deployment_id: "expiry-agent".to_string(),
target: server.default_target.clone(),
max_leases: 1,
lease_seconds: 1,
};
let first = server
.acquire_lease("expiry-agent", lease_request.clone())
.await
.unwrap();
assert_eq!(first.leases.len(), 1);
assert_eq!(first.leases[0].command_id, created.command_id);
assert_eq!(first.leases[0].attempt, 1);
let while_held = server
.acquire_lease("expiry-agent", lease_request.clone())
.await
.unwrap();
assert!(
while_held.leases.is_empty(),
"a live lease must not be double-leased"
);
tokio::time::sleep(std::time::Duration::from_millis(2500)).await;
let second = server
.acquire_lease("expiry-agent", lease_request)
.await
.unwrap();
assert_eq!(
second.leases.len(),
1,
"an expired lease must make the command leasable again"
);
assert_eq!(second.leases[0].command_id, created.command_id);
assert_eq!(
second.leases[0].attempt, 2,
"expiry-driven redelivery must increment the attempt"
);
assert_eq!(
second.leases[0].envelope.attempt, 2,
"the redelivered envelope must carry the incremented attempt"
);
}
#[tokio::test]
async fn test_leased_envelope_manager_urls_are_path_relative() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let mut request = test_inline_create_command("relative-agent", "relative-urls");
request.target_resource_id = Some(server.default_target.resource_id.clone());
let created = server.create_command(request).await.unwrap();
let leases = server
.acquire_lease(
"relative-agent",
LeaseRequest {
deployment_id: "relative-agent".to_string(),
target: server.default_target.clone(),
max_leases: 1,
lease_seconds: 60,
},
)
.await
.unwrap();
assert_eq!(leases.leases.len(), 1);
let envelope = &leases.leases[0].envelope;
let submit = &envelope.response_handling.submit_response_url;
assert!(
submit.starts_with(&format!("{}/response?", created.command_id)),
"submit URL must be relative to the lease endpoint, got '{submit}'"
);
assert!(
submit.contains("response_token="),
"relativization must preserve the signed query, got '{submit}'"
);
}
#[tokio::test]
async fn test_lease_target_mismatch_in_pending_index_is_loud_error() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
server
.registry
.register_target("second-daemon", CommandTargetType::Daemon)
.await
.unwrap();
let mut request = test_inline_create_command("target-agent", "for-second");
request.target_resource_id = Some("second-daemon".to_string());
let cmd = server.create_command(request).await.unwrap();
let corrupt_key = format!(
"target:target-agent:{}:pending:{}:{}",
server.default_target.resource_id,
chrono::Utc::now().timestamp_nanos_opt().unwrap_or(0),
cmd.command_id
);
alien_bindings::traits::Kv::put(server.kv.as_ref(), &corrupt_key, vec![], None)
.await
.unwrap();
let result = server
.acquire_lease(
"target-agent",
LeaseRequest {
deployment_id: "target-agent".to_string(),
target: server.default_target.clone(),
max_leases: 1,
lease_seconds: 60,
},
)
.await;
let err = result.unwrap_err();
assert!(
err.message.contains(&cmd.command_id) || err.message.contains("target"),
"expected loud target-mismatch error, got: {}",
err.message
);
}
#[tokio::test]
async fn test_idempotency_scoped_per_target() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
server
.registry
.register_target("second-daemon", CommandTargetType::Daemon)
.await
.unwrap();
let make_request = |target: &str| {
let mut request = test_inline_create_command("target-agent", "idem-command");
request.target_resource_id = Some(target.to_string());
request.idempotency_key = Some("same-key".to_string());
request
};
let default_id = server.default_target.resource_id.clone();
let first = server
.create_command(make_request(&default_id))
.await
.unwrap();
let second = server
.create_command(make_request("second-daemon"))
.await
.unwrap();
assert_ne!(first.command_id, second.command_id);
let replay = server
.create_command(make_request(&default_id))
.await
.unwrap();
assert_eq!(replay.command_id, first.command_id);
}
#[tokio::test]
async fn test_worker_and_daemon_share_command_name_route_independently() {
let server = TestCommandServer::new().await;
server
.registry
.register_target("shared-daemon", CommandTargetType::Daemon)
.await
.unwrap();
let daemon_target = CommandTarget::new("shared-daemon", CommandTargetType::Daemon);
let mut worker_request = test_inline_create_command("target-agent", "shared-command");
worker_request.target_resource_id = Some(server.default_target.resource_id.clone());
let worker_response = server.create_command(worker_request).await.unwrap();
assert_eq!(worker_response.state, CommandState::Dispatched);
let mut daemon_request = test_inline_create_command("target-agent", "shared-command");
daemon_request.target_resource_id = Some("shared-daemon".to_string());
let daemon_response = server.create_command(daemon_request).await.unwrap();
assert_eq!(daemon_response.state, CommandState::Pending);
let mock_dispatcher = server
.mock_dispatcher()
.expect("Should have mock dispatcher");
mock_dispatcher.assert_dispatch_count(1).await;
let dispatched = mock_dispatcher.get_latest().await.unwrap();
assert_eq!(dispatched.envelope.command_id, worker_response.command_id);
assert_eq!(dispatched.envelope.target, server.default_target);
let daemon_leases = server
.acquire_lease(
"target-agent",
LeaseRequest {
deployment_id: "target-agent".to_string(),
target: daemon_target.clone(),
max_leases: 10,
lease_seconds: 60,
},
)
.await
.unwrap();
assert_eq!(daemon_leases.leases.len(), 1);
assert_eq!(
daemon_leases.leases[0].command_id,
daemon_response.command_id
);
assert_eq!(daemon_leases.leases[0].envelope.target, daemon_target);
let worker_leases = server
.acquire_lease(
"target-agent",
LeaseRequest {
deployment_id: "target-agent".to_string(),
target: server.default_target.clone(),
max_leases: 10,
lease_seconds: 60,
},
)
.await
.unwrap();
assert_eq!(worker_leases.leases.len(), 0);
}
#[tokio::test]
async fn test_basic_api_operations() {
let server = TestCommandServer::new().await;
let request = test_inline_create_command("api-agent", "test-command");
let response = server.create_command(request).await.unwrap();
assert_eq!(response.state, CommandState::Dispatched);
assert!(response.command_id.starts_with("cmd_"));
assert!(response.storage_upload.is_none());
let status = server
.get_command_status(&response.command_id)
.await
.unwrap();
assert_eq!(status.command_id, response.command_id);
assert_eq!(status.state, CommandState::Dispatched);
assert_eq!(status.attempt, 1);
let large_request = test_storage_create_command("storage-agent", "upload-command", 200_000);
let response = server.create_command(large_request).await.unwrap();
assert_eq!(response.state, CommandState::PendingUpload);
assert!(response.storage_upload.is_some());
let upload_complete = test_upload_complete_request(200_000);
let complete_response = server
.upload_complete(&response.command_id, upload_complete)
.await
.unwrap();
assert_eq!(complete_response.state, CommandState::Dispatched);
}
#[tokio::test]
async fn ambiguous_push_failure_keeps_dispatched_and_accepts_late_response() {
let server = TestCommandServer::new().await;
let dispatcher = server
.mock_dispatcher()
.expect("push test server uses the mock dispatcher");
dispatcher.set_should_fail(true).await;
let created = server
.create_command(test_inline_create_command(
"push-agent",
"ambiguous-dispatch",
))
.await
.expect("ambiguous acknowledgement must still return the durable command ID");
assert_eq!(created.state, CommandState::Dispatched);
dispatcher.assert_dispatch_count(0).await;
let status = server
.get_command_status(&created.command_id)
.await
.expect("command status");
assert_eq!(status.command_id, created.command_id);
assert_eq!(status.state, CommandState::Dispatched);
assert_eq!(status.attempt, 1);
server
.submit_command_response(
&created.command_id,
test_success_response(b"completed after ambiguous acknowledgement"),
)
.await
.expect("late response remains valid");
server.assert_command_succeeded(&created.command_id).await;
}
#[tokio::test]
async fn definite_push_rejection_becomes_terminal_delivery_failure() {
let server = TestCommandServer::new().await;
let dispatcher = server
.mock_dispatcher()
.expect("push test server uses the mock dispatcher");
dispatcher.set_should_reject(true).await;
let created = server
.create_command(test_inline_create_command(
"push-agent",
"definite-rejection",
))
.await
.expect("definite rejection must return the durable terminal command ID");
assert_eq!(created.state, CommandState::Failed);
let status = server
.get_command_status(&created.command_id)
.await
.expect("command status");
assert_eq!(status.state, CommandState::Failed);
let Some(CommandResponse::Error { code, message, .. }) = status.response else {
panic!("definite delivery rejection must persist an error response");
};
assert_eq!(code, "DELIVERY_FAILED");
assert_eq!(message, "Worker runtime did not accept command delivery");
let late = server
.submit_command_response(
&created.command_id,
test_success_response(b"must not replace delivery failure"),
)
.await;
assert!(
late.is_ok(),
"terminal duplicate submissions are idempotent"
);
server.assert_command_failed(&created.command_id).await;
}
#[tokio::test]
async fn test_lease_operations() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let request = test_inline_create_command("lease-agent", "lease-command");
let response = server.create_command(request).await.unwrap();
assert_eq!(response.state, CommandState::Pending);
let lease = server
.acquire_single_lease("lease-agent")
.await
.unwrap()
.unwrap();
assert_eq!(lease.command_id, response.command_id);
assert_eq!(lease.attempt, 1);
assert!(lease.lease_expires_at > Utc::now());
assert_envelope_command_id(&lease.envelope, &response.command_id);
assert_envelope_command(&lease.envelope, "lease-command");
let empty_lease_request = LeaseRequest {
deployment_id: "nonexistent-agent".to_string(),
target: server.default_target.clone(),
max_leases: 1,
lease_seconds: 60,
};
let empty_response = server
.acquire_lease("nonexistent-agent", empty_lease_request)
.await
.unwrap();
assert_eq!(empty_response.leases.len(), 0);
server
.release_lease(&lease.command_id, &lease.lease_id)
.await
.unwrap();
server
.assert_command_state(&response.command_id, CommandState::Pending)
.await;
}
#[tokio::test]
async fn max_timeout_lease_receives_response_credentials_beyond_lease_headroom() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let created = server
.create_command(test_inline_create_command("max-timeout", "run"))
.await
.unwrap();
let lease = server
.acquire_lease(
"max-timeout",
LeaseRequest {
deployment_id: "max-timeout".to_string(),
target: server.default_target.clone(),
max_leases: 1,
lease_seconds: 3660,
},
)
.await
.unwrap()
.leases
.into_iter()
.next()
.expect("max-timeout lease");
assert_eq!(lease.command_id, created.command_id);
let submit = url::Url::parse("http://manager.invalid")
.unwrap()
.join(&lease.envelope.response_handling.submit_response_url)
.unwrap();
let expires = submit
.query_pairs()
.find_map(|(key, value)| (key == "expires").then(|| value.parse::<i64>().unwrap()))
.expect("response token expiry");
let required = Utc::now() + chrono::Duration::seconds(3660);
assert!(
expires > required.timestamp(),
"response token must outlive max execution plus lease headroom"
);
assert!(
lease
.envelope
.response_handling
.storage_upload_request
.expiration
> required,
"response upload credential must outlive max execution plus lease headroom"
);
}
#[tokio::test]
async fn delayed_storage_command_lease_refreshes_expired_params_get() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let created = server
.create_command(test_storage_create_command(
"delayed-storage",
"run",
200_000,
))
.await
.unwrap();
server
.upload_complete(&created.command_id, test_upload_complete_request(200_000))
.await
.unwrap();
let mut stored = server
.command_server
.get_params(&created.command_id)
.await
.unwrap()
.expect("stored params");
let BodySpec::Storage {
storage_get_request: Some(request),
..
} = &mut stored
else {
panic!("storage params with GET request");
};
request.expiration = Utc::now() - chrono::Duration::hours(1);
server
.command_server
.store_params(&created.command_id, &stored)
.await
.unwrap();
let lease = server
.acquire_single_lease("delayed-storage")
.await
.unwrap()
.expect("delayed command lease");
let BodySpec::Storage {
storage_get_request: Some(fresh),
..
} = lease.envelope.params
else {
panic!("leased storage params with fresh GET request");
};
assert!(
fresh.expiration > Utc::now() + chrono::Duration::minutes(50),
"leasing must replace the stale upload-time params URL"
);
}
#[tokio::test]
async fn test_response_operations() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let request = test_inline_create_command("response-agent", "response-command");
let response = server.create_command(request).await.unwrap();
let lease = server
.acquire_single_lease("response-agent")
.await
.unwrap()
.unwrap();
let agent_response = test_json_success_response(&serde_json::json!({
"result": "success",
"data": [1, 2, 3]
}));
server
.submit_command_response(&lease.command_id, agent_response)
.await
.unwrap();
server.assert_command_succeeded(&response.command_id).await;
let status = server
.get_command_status(&response.command_id)
.await
.unwrap();
let final_response = status.response.unwrap();
assert!(final_response.is_success());
if let CommandResponse::Success { response: body } = final_response {
let decoded = body.decode_inline().unwrap();
let json: serde_json::Value = serde_json::from_slice(&decoded).unwrap();
assert_eq!(json["result"], "success");
}
let duplicate_response = test_success_response(b"second response");
let result = server
.submit_command_response(&lease.command_id, duplicate_response)
.await;
assert!(result.is_ok());
let status = server
.get_command_status(&response.command_id)
.await
.unwrap();
let final_response = status.response.unwrap();
if let CommandResponse::Success { response: body } = final_response {
let decoded = body.decode_inline().unwrap();
let json: serde_json::Value = serde_json::from_slice(&decoded).unwrap();
assert_eq!(json["result"], "success"); }
}
#[tokio::test]
#[tracing_test::traced_test]
async fn test_response_cleanup_failure_keeps_command_terminal() {
let server = TestCommandServer::builder()
.with_pull_mode()
.with_fault_injection()
.build()
.await;
let fault_kv = server
.fault_kv
.clone()
.expect("fault injection was requested");
let request = test_inline_create_command("cleanup-agent", "cleanup-command");
let response = server.create_command(request).await.unwrap();
let lease = server
.acquire_single_lease("cleanup-agent")
.await
.unwrap()
.unwrap();
fault_kv.arm_pending_scan_failure();
let agent_response = test_json_success_response(&serde_json::json!({ "result": "ok" }));
let submit_result = server
.submit_command_response(&lease.command_id, agent_response)
.await;
assert!(
submit_result.is_ok(),
"cleanup failures are best-effort and must not fail submit_command_response \
once the terminal state is committed"
);
assert!(
logs_contain("Failed to clean up pending index"),
"a cleanup failure must still be logged as a warning"
);
let status = server
.get_command_status(&response.command_id)
.await
.unwrap();
assert_eq!(
status.state,
CommandState::Succeeded,
"command must be terminal even though response cleanup failed"
);
let final_response = status
.response
.expect("stored response must be visible on the terminal command");
assert!(final_response.is_success());
if let CommandResponse::Success { response: body } = final_response {
let decoded = body.decode_inline().unwrap();
let json: serde_json::Value = serde_json::from_slice(&decoded).unwrap();
assert_eq!(json["result"], "ok");
}
}
#[tokio::test]
async fn test_runtime_integration() {
let envelope = test_simple_envelope("cmd_runtime_test", "runtime-command");
let envelope_json = serde_json::to_value(&envelope).unwrap();
let queue_message = QueueMessage {
id: "msg_123".to_string(),
payload: MessagePayload::Json(envelope_json),
receipt_handle: "handle_123".to_string(),
timestamp: Utc::now(),
source: "test-queue".to_string(),
attributes: std::collections::HashMap::new(),
attempt_count: Some(1),
};
let parsed = parse_envelope(&queue_message).unwrap();
assert!(parsed.is_some());
let parsed_envelope = parsed.unwrap();
assert_eq!(parsed_envelope.command_id, "cmd_runtime_test");
assert_eq!(parsed_envelope.command, "runtime-command");
let non_arc_message = QueueMessage {
id: "msg_456".to_string(),
payload: MessagePayload::Json(serde_json::json!({"regular": "message"})),
receipt_handle: "handle_456".to_string(),
timestamp: Utc::now(),
source: "test-queue".to_string(),
attributes: std::collections::HashMap::new(),
attempt_count: Some(1),
};
let parsed = parse_envelope(&non_arc_message).unwrap();
assert!(parsed.is_none());
let params_json = serde_json::json!({"key": "value", "number": 42});
let params_bytes = serde_json::to_vec(¶ms_json).unwrap();
let test_envelope = test_envelope(
"cmd_params",
"params-command",
BodySpec::inline(¶ms_bytes),
);
let decoded_params = decode_params(&test_envelope).await.unwrap();
assert_eq!(decoded_params["key"], "value");
assert_eq!(decoded_params["number"], 42);
}
#[tokio::test]
async fn test_error_handling() {
let server = TestCommandServer::new().await;
let invalid_request = CreateCommandRequest {
deployment_id: "error-agent".to_string(),
command: "".to_string(), params: BodySpec::inline(b"{}"),
deadline: None,
idempotency_key: None,
target_resource_id: None,
};
let result = server.create_command(invalid_request).await;
assert!(result.is_err());
let invalid_request = CreateCommandRequest {
deployment_id: "".to_string(), command: "test".to_string(),
params: BodySpec::inline(b"{}"),
deadline: None,
idempotency_key: None,
target_resource_id: None,
};
let result = server.create_command(invalid_request).await;
assert!(result.is_err());
let upload_complete = test_upload_complete_request(1000);
assert!(server
.upload_complete("nonexistent", upload_complete)
.await
.is_err());
assert!(server.get_command_status("nonexistent").await.is_err());
assert!(server
.submit_command_response("nonexistent", test_success_response(b"test"))
.await
.is_err());
let past_deadline = Utc::now() - chrono::Duration::minutes(1);
let expired_request = test_create_command_with_deadline(
"expired-agent",
"expired-command",
BodySpec::inline(b"{}"),
past_deadline,
);
assert!(server.create_command(expired_request).await.is_err());
}
#[tokio::test]
async fn test_error_response_handling() {
let server = TestCommandServer::builder().with_pull_mode().build().await;
let request = test_inline_create_command("error-agent", "error-command");
let response = server.create_command(request).await.unwrap();
let lease = server
.acquire_single_lease("error-agent")
.await
.unwrap()
.unwrap();
let agent_response = test_error_response("PROCESSING_FAILED", "Something went wrong");
server
.submit_command_response(&lease.command_id, agent_response)
.await
.unwrap();
let status = server
.get_command_status(&response.command_id)
.await
.unwrap();
assert_eq!(status.state, CommandState::Failed);
let final_response = status.response.unwrap();
assert!(final_response.is_error());
if let CommandResponse::Error { code, message, .. } = final_response {
assert_eq!(code, "PROCESSING_FAILED");
assert_eq!(message, "Something went wrong");
}
}
}
#[cfg(feature = "test-utils")]
#[tokio::test]
async fn concurrent_opposite_submits_keep_state_and_response_consistent() {
use alien_commands::test_utils::*;
use alien_commands::types::{CommandResponse, CommandState};
let server = TestCommandServer::builder().with_pull_mode().build().await;
let request = test_inline_create_command("pull-agent", "flaky-op");
let created = server.create_command(request).await.unwrap();
let lease = server
.acquire_single_lease("pull-agent")
.await
.unwrap()
.unwrap();
let success = test_json_success_response(&serde_json::json!({ "ok": true }));
let failure = CommandResponse::Error {
code: "HANDLER_ERROR".to_string(),
message: "boom".to_string(),
details: None,
};
let (a, b) = tokio::join!(
server.submit_command_response(&lease.command_id, success),
server.submit_command_response(&lease.command_id, failure),
);
a.unwrap();
b.unwrap();
let status = server
.get_command_status(&created.command_id)
.await
.unwrap();
assert!(status.state.is_terminal());
let response = status.response.expect("terminal command serves a response");
match (&status.state, &response) {
(CommandState::Succeeded, CommandResponse::Success { .. }) => {}
(CommandState::Failed, CommandResponse::Error { .. }) => {}
(state, response) => {
panic!("torn terminal record: state {state:?} does not match response {response:?}")
}
}
}
#[cfg(feature = "test-utils")]
#[tokio::test]
async fn deadline_reaper_expires_overdue_commands() {
use alien_commands::test_utils::*;
use alien_commands::types::CommandState;
let server = TestCommandServer::builder().with_pull_mode().build().await;
let mut request = test_inline_create_command("pull-agent", "slow-op");
request.deadline = Some(chrono::Utc::now() + chrono::Duration::seconds(2));
let created = server.create_command(request).await.unwrap();
assert_eq!(created.state, CommandState::Pending);
tokio::time::sleep(std::time::Duration::from_millis(2600)).await;
let expired = server.command_server.reap_expired_commands().await.unwrap();
assert_eq!(expired, 1, "the overdue command must be reaped");
let status = server
.get_command_status(&created.command_id)
.await
.unwrap();
assert_eq!(status.state, CommandState::Expired);
let lease = server.acquire_single_lease("pull-agent").await.unwrap();
assert!(lease.is_none(), "expired command must not be leasable");
assert_eq!(
server.command_server.reap_expired_commands().await.unwrap(),
0
);
}