use agent_context_contract::{
ContextDerivationKind, ContextSnapshotItem, ContextSnapshotItemKind, ContextSnapshotStage,
ConversationRecord, ConversationScope, CreateContextDerivationRequest,
CreateContextSnapshotRequest, ModelInvocationState, PutModelInvocationRequest,
};
use agent_deploy_contract::RollbackDeploymentRequest;
use agent_infra_sdk::{
CallOptions, ClientOptions, CredentialsProvider, GatewayExchangeCredentials, InfraClient,
InfraClientError, StaticCredentials,
};
use agent_knowledge_contract::{
KnowledgeQueryRequest, KnowledgeRequirement, ResolveKnowledgeResourcesRequest,
};
use std::collections::BTreeMap;
use std::sync::Arc;
use std::time::Duration;
use wiremock::matchers::{body_json, header, header_regex, method, path, query_param};
use wiremock::{Mock, MockServer, ResponseTemplate};
fn client(server: &MockServer) -> InfraClient {
InfraClient::new(server.uri()).expect("build Gateway-backed client")
}
#[tokio::test]
async fn runtime_bootstrap_uses_the_public_gateway_path_without_bearer_auth() {
use agent_infra_sdk::runtime_identity_contract::{
CredentialCapability, IssueRuntimeCredentialRequest,
};
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path(
agent_infra_sdk::gateway_contract::RUNTIME_CREDENTIAL_BOOTSTRAP_PATH,
))
.and(body_json(serde_json::json!({
"workloadProof": "bootstrap-proof",
"audience": "agent-infra",
"capabilities": ["gateway:delegation:manage"]
})))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"accessToken": "runtime-token",
"tokenType": "Bearer",
"expiresInSeconds": 120,
"keyId": "runtime-key"
})))
.expect(1)
.mount(&server)
.await;
let client = InfraClient::builder(server.uri())
.credentials(Arc::new(StaticCredentials::new("must-not-be-sent").unwrap()))
.build()
.unwrap();
let issued = client
.gateway()
.bootstrap_runtime_credential(&IssueRuntimeCredentialRequest {
workload_proof: "bootstrap-proof".into(),
proof_audience: None,
audience: "agent-infra".into(),
capabilities: vec![
CredentialCapability::new("gateway:delegation:manage").unwrap(),
],
})
.await
.unwrap();
assert_eq!(issued.access_token, "runtime-token");
let requests = server.received_requests().await.unwrap();
assert!(!requests[0].headers.contains_key("authorization"));
}
#[tokio::test]
async fn unified_admin_sdk_uses_the_gateway_admin_namespace() {
use agent_infra_sdk::agent_contract::ListAgentsRequest;
use agent_infra_sdk::context_contract::{ConversationScope, SearchMessagesRequest};
use agent_infra_sdk::runtime_identity_contract::{
ListRuntimeIdentitiesRequest, RuntimeIdentityId, TransitionIdentityRequest,
};
let server = MockServer::start().await;
for (verb, endpoint, response) in [
("GET", "/admin/api/agents", serde_json::json!({"items":[]})),
(
"POST",
"/admin/api/context/messages/search",
serde_json::json!({"items":[]}),
),
(
"POST",
"/admin/api/runtime-identities:list",
serde_json::json!({"items":[]}),
),
] {
let mut mock = Mock::given(method(verb)).and(path(endpoint));
if endpoint == "/admin/api/agents" {
mock = mock.and(query_param("limit", "25"));
}
mock.and(header("authorization", "Bearer admin-token"))
.respond_with(ResponseTemplate::new(200).set_body_json(response))
.expect(1)
.mount(&server)
.await;
}
let client = InfraClient::builder(server.uri())
.credentials(Arc::new(StaticCredentials::new("admin-token").unwrap()))
.build()
.unwrap();
client
.admin()
.list_agents(&ListAgentsRequest {
limit: Some(25),
..Default::default()
})
.await
.unwrap();
client
.admin()
.search_context(&SearchMessagesRequest {
scope: ConversationScope::conversation("conversation-admin"),
query: Some("needle".into()),
role: None,
limit: Some(25),
})
.await
.unwrap();
client
.admin()
.list_runtime_identities(&ListRuntimeIdentitiesRequest {
limit: 25,
..Default::default()
})
.await
.unwrap();
let invalid_runtime = client
.admin()
.transition_runtime_identity(
&RuntimeIdentityId::new("runtime-admin").unwrap(),
"destroy",
&TransitionIdentityRequest {
expected_revision: 1,
idempotency_key: "invalid-action".into(),
reason: None,
},
)
.await
.unwrap_err();
assert!(matches!(invalid_runtime, InfraClientError::Protocol { .. }));
}
#[tokio::test]
async fn unified_admin_sdk_covers_mutating_resource_contracts() {
use agent_infra_sdk::agent_contract::{
AgentId, CreateAgentRequest, ResourceRelations, UpdateAgentRequest,
};
use agent_infra_sdk::runtime_identity_contract::{
RuntimeIdentityId, TransitionIdentityRequest,
};
let server = MockServer::start().await;
let agent = serde_json::json!({
"id":"agt_admin","ownerUserId":"user-admin","name":"Admin Agent",
"description":"managed","status":"active","revision":1,"relations":{},
"createdAt":"2026-08-13T00:00:00Z","updatedAt":"2026-08-13T00:00:00Z"
});
Mock::given(method("POST"))
.and(path("/admin/api/agents"))
.respond_with(ResponseTemplate::new(201).set_body_json(agent.clone()))
.expect(1)
.mount(&server)
.await;
Mock::given(method("PATCH"))
.and(path("/admin/api/agents/agt_admin"))
.respond_with(ResponseTemplate::new(200).set_body_json(agent))
.expect(1)
.mount(&server)
.await;
let runtime = serde_json::json!({
"id":"runtime-admin","scope":{"tenantId":"tenant-admin","projectId":"project-admin"},
"kind":"code","status":"pending","revision":1,"grantRevision":0,
"credentialGeneration":0,"metadata":{},"relations":{},
"createdAt":"2026-08-13T00:00:00Z","updatedAt":"2026-08-13T00:00:00Z"
});
Mock::given(method("POST"))
.and(path("/admin/api/runtime-identities/runtime-admin:activate"))
.respond_with(ResponseTemplate::new(200).set_body_json(runtime))
.expect(1)
.mount(&server)
.await;
let client = client(&server);
client
.admin()
.create_agent(&CreateAgentRequest {
name: "Admin Agent".into(),
description: "managed".into(),
relations: ResourceRelations::default(),
idempotency_key: "create-admin-agent".into(),
})
.await
.unwrap();
client
.admin()
.update_agent(
&AgentId::new("agt_admin").unwrap(),
&UpdateAgentRequest {
expected_revision: 1,
name: Some("Admin Agent".into()),
description: None,
status: None,
relations: None,
},
)
.await
.unwrap();
client
.admin()
.transition_runtime_identity(
&RuntimeIdentityId::new("runtime-admin").unwrap(),
"activate",
&TransitionIdentityRequest {
expected_revision: 1,
idempotency_key: "activate-admin-runtime".into(),
reason: None,
},
)
.await
.unwrap();
}
#[tokio::test]
async fn agent_registry_sdk_uses_gateway_paths_and_preserves_sky_token() {
use agent_infra_sdk::agent_contract::{
AgentId, CreateAgentRequest, ListAgentsRequest, ResourceRelations, UpdateAgentRequest,
};
let server = MockServer::start().await;
let created = serde_json::json!({
"id": "agt_sdk",
"ownerUserId": "user-sdk",
"name": "SDK Agent",
"description": "created",
"status": "active",
"revision": 1,
"relations": {"runtimeIds":["runtime-sdk"], "userIds":["user-sdk"]},
"createdAt": "2026-08-11T00:00:00Z",
"updatedAt": "2026-08-11T00:00:00Z"
});
Mock::given(method("POST"))
.and(path("/v1/agents"))
.and(header("authorization", "Bearer sky-sdk-token"))
.and(body_json(serde_json::json!({
"name":"SDK Agent", "description":"created",
"relations":{"runtimeIds":["runtime-sdk"]},
"idempotencyKey":"create-sdk-agent"
})))
.respond_with(ResponseTemplate::new(201).set_body_json(created.clone()))
.expect(1)
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/v1/agents/agt_sdk"))
.and(header("authorization", "Bearer sky-sdk-token"))
.respond_with(ResponseTemplate::new(200).set_body_json(created.clone()))
.expect(1)
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/v1/agents"))
.and(query_param("relatedRuntimeId", "runtime-sdk"))
.and(query_param("limit", "10"))
.and(header("authorization", "Bearer sky-sdk-token"))
.respond_with(
ResponseTemplate::new(200)
.set_body_json(serde_json::json!({"items":[created.clone()]})),
)
.expect(1)
.mount(&server)
.await;
let updated = serde_json::json!({
"id": "agt_sdk", "ownerUserId": "user-sdk", "name": "SDK Agent",
"description": "updated", "status": "archived", "revision": 2,
"relations": {"runtimeIds":["runtime-sdk"], "userIds":["user-sdk"]},
"createdAt": "2026-08-11T00:00:00Z", "updatedAt": "2026-08-11T00:01:00Z"
});
Mock::given(method("PATCH"))
.and(path("/v1/agents/agt_sdk"))
.and(header("authorization", "Bearer sky-sdk-token"))
.and(body_json(serde_json::json!({
"expectedRevision":1, "description":"updated", "status":"archived"
})))
.respond_with(ResponseTemplate::new(200).set_body_json(updated))
.expect(1)
.mount(&server)
.await;
let client = InfraClient::builder(server.uri())
.credentials(Arc::new(StaticCredentials::new("sky-sdk-token").unwrap()))
.build()
.unwrap();
let request = CreateAgentRequest {
name: "SDK Agent".into(),
description: "created".into(),
relations: ResourceRelations {
runtime_ids: vec!["runtime-sdk".into()],
..Default::default()
},
idempotency_key: "create-sdk-agent".into(),
};
let agent = client.agents().create(&request).await.unwrap();
assert_eq!(
client.agents().get(&agent.id).await.unwrap().id,
AgentId::new("agt_sdk").unwrap()
);
assert_eq!(
client
.agents()
.list(&ListAgentsRequest {
related_runtime_id: Some("runtime-sdk".into()),
limit: Some(10),
..Default::default()
})
.await
.unwrap()
.items
.len(),
1
);
let updated = client
.agents()
.update(
&agent.id,
&UpdateAgentRequest {
expected_revision: 1,
name: None,
description: Some("updated".into()),
status: Some(agent_infra_sdk::agent_contract::AgentStatus::Archived),
relations: None,
},
)
.await
.unwrap();
assert_eq!(updated.revision, 2);
}
#[tokio::test]
async fn context_export_chunks_follow_the_bounded_cursor_contract() {
let server = MockServer::start().await;
let chunk_path = "/internal/v1/context/export-artifacts/chunk";
Mock::given(method("POST"))
.and(path(chunk_path))
.and(body_json(serde_json::json!({
"artifactId":"artifact-1", "offset":0, "limit":4
})))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"artifactId":"artifact-1", "offset":0, "dataBase64":"YWJjZA==",
"nextOffset":4, "eof":false, "digest":"sha256:digest"
})))
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(chunk_path))
.and(body_json(serde_json::json!({
"artifactId":"artifact-1", "offset":4, "limit":4
})))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"artifactId":"artifact-1", "offset":4, "dataBase64":"ZWY=",
"nextOffset":6, "eof":true, "digest":"sha256:digest"
})))
.expect(1)
.mount(&server)
.await;
let client = InfraClient::builder(server.uri()).build().unwrap();
let first = client
.context()
.export_chunk("artifact-1", 0, 4)
.await
.unwrap();
let terminal = client
.context()
.next_export_chunk(&first, 4)
.await
.unwrap()
.unwrap();
assert!(terminal.eof);
assert!(
client
.context()
.next_export_chunk(&terminal, 4)
.await
.unwrap()
.is_none()
);
}
#[tokio::test]
async fn context_records_and_knowledge_sdk_cover_every_new_http_contract() {
let server = MockServer::start().await;
let snapshot = serde_json::json!({
"id": "snapshot-sdk",
"scope": {"tenantId": "tenant-sdk", "projectId": "project-sdk", "subject": "runtime-sdk"},
"conversation": {"conversationId": "conversation-sdk"},
"stage": "compressed",
"purpose": "model_request",
"tokenCount": 4,
"items": [{
"id": "item-sdk", "ordinal": 0, "kind": "message", "role": "user",
"content": "exact request", "metadata": {}
}],
"createdAtMs": 1
});
Mock::given(method("POST"))
.and(path(agent_context_contract::SNAPSHOTS_PATH))
.respond_with(ResponseTemplate::new(201).set_body_json(snapshot.clone()))
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(agent_context_contract::SNAPSHOT_GET_PATH))
.and(body_json(serde_json::json!({"snapshotId": "snapshot-sdk"})))
.respond_with(ResponseTemplate::new(200).set_body_json(snapshot))
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(agent_context_contract::DERIVATIONS_PATH))
.respond_with(ResponseTemplate::new(201).set_body_json(serde_json::json!({
"id": "derivation-sdk",
"scope": {"tenantId": "tenant-sdk", "projectId": "project-sdk", "subject": "runtime-sdk"},
"sourceSnapshotId": "snapshot-source", "targetSnapshotId": "snapshot-sdk",
"kind": "compress", "strategy": "summary", "strategyVersion": "v1",
"inputTokenCount": 10, "outputTokenCount": 4,
"inputItemCount": 2, "outputItemCount": 1, "metadata": {}, "createdAtMs": 2
})))
.expect(1)
.mount(&server)
.await;
let invocation = serde_json::json!({
"id": "invocation-sdk",
"scope": {"tenantId": "tenant-sdk", "projectId": "project-sdk", "subject": "runtime-sdk"},
"conversationId": "conversation-sdk", "requestSnapshotId": "snapshot-sdk",
"model": "model-sdk", "provider": "provider-sdk", "state": "prepared",
"usage": {}, "createdAtMs": 3, "completedAtMs": null
});
Mock::given(method("POST"))
.and(path(agent_context_contract::MODEL_INVOCATIONS_PATH))
.respond_with(ResponseTemplate::new(200).set_body_json(invocation.clone()))
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(agent_context_contract::MODEL_INVOCATION_GET_PATH))
.and(body_json(
serde_json::json!({"invocationId": "invocation-sdk"}),
))
.respond_with(ResponseTemplate::new(200).set_body_json(invocation))
.expect(1)
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path(agent_knowledge_contract::KNOWLEDGE_CAPABILITIES_PATH))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"provider": "knowledge-sdk", "queryModes": ["hybrid"],
"resourceResolution": true, "maxHits": 20
})))
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(agent_knowledge_contract::KNOWLEDGE_QUERY_PATH))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"hits": [{
"provider": "knowledge-sdk", "resourceRef": "opaque-sdk",
"excerpt": "evidence", "score": 0.9, "metadata": {}
}],
"degraded": false
})))
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(agent_knowledge_contract::KNOWLEDGE_RESOLVE_PATH))
.and(body_json(
serde_json::json!({"resourceRefs": ["opaque-sdk"]}),
))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"resources": [{
"provider": "knowledge-sdk", "resourceRef": "opaque-sdk",
"title": "SDK resource", "metadata": {}
}]
})))
.expect(1)
.mount(&server)
.await;
let client = InfraClient::builder(server.uri()).build().unwrap();
let context = client.context();
let snapshot_request = CreateContextSnapshotRequest {
id: "snapshot-sdk".into(),
conversation: ConversationScope::conversation("conversation-sdk"),
stage: ContextSnapshotStage::Compressed,
purpose: "model_request".into(),
token_count: 4,
items: vec![ContextSnapshotItem {
id: "item-sdk".into(),
ordinal: 0,
kind: ContextSnapshotItemKind::Message,
role: Some("user".into()),
content: serde_json::json!("exact request"),
source_record_id: None,
knowledge_evidence: None,
metadata: serde_json::json!({}),
}],
};
assert_eq!(
context.create_snapshot(&snapshot_request).await.unwrap().id,
"snapshot-sdk"
);
assert_eq!(
context
.get_snapshot("snapshot-sdk")
.await
.unwrap()
.items
.len(),
1
);
assert_eq!(
context
.create_derivation(&CreateContextDerivationRequest {
id: "derivation-sdk".into(),
source_snapshot_id: "snapshot-source".into(),
target_snapshot_id: "snapshot-sdk".into(),
kind: ContextDerivationKind::Compress,
strategy: "summary".into(),
strategy_version: "v1".into(),
input_token_count: 10,
output_token_count: 4,
input_item_count: 2,
output_item_count: 1,
metadata: serde_json::json!({}),
})
.await
.unwrap()
.id,
"derivation-sdk"
);
let invocation_request = PutModelInvocationRequest {
id: "invocation-sdk".into(),
conversation_id: "conversation-sdk".into(),
request_snapshot_id: "snapshot-sdk".into(),
model: "model-sdk".into(),
provider: "provider-sdk".into(),
state: ModelInvocationState::Prepared,
response: None,
error_code: None,
usage: serde_json::json!({}),
completed_at_ms: None,
};
assert_eq!(
context
.put_model_invocation(&invocation_request)
.await
.unwrap()
.id,
"invocation-sdk"
);
assert_eq!(
context
.get_model_invocation("invocation-sdk")
.await
.unwrap()
.model,
"model-sdk"
);
assert_eq!(
context.knowledge_capabilities().await.unwrap().provider,
"knowledge-sdk"
);
assert_eq!(
context
.query_knowledge(&KnowledgeQueryRequest {
query: "evidence".into(),
requirement: KnowledgeRequirement::Required,
provider: None,
scope_ref: None,
top_k: 5,
filters: serde_json::json!({}),
})
.await
.unwrap()
.hits[0]
.resource_ref,
"opaque-sdk"
);
assert_eq!(
context
.resolve_knowledge_resources(&ResolveKnowledgeResourcesRequest {
resource_refs: vec!["opaque-sdk".into()],
})
.await
.unwrap()
.resources[0]
.title,
"SDK resource"
);
}
#[tokio::test]
async fn deploy_rollback_returns_a_durable_operation_handle() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path(
"/internal/v1/deploy/deployments/deployment-1/rollback",
))
.and(header("idempotency-key", "rollback-key"))
.respond_with(ResponseTemplate::new(202).set_body_json(serde_json::json!({
"id": "operation-rollback",
"deploymentId": "deployment-1",
"state": "queued",
"phase": "queued",
"attempt": 0,
"leaseOwner": null,
"leaseDeadline": null,
"fencingToken": 0,
"providerExternalId": "",
"version": 1,
"createdAt": "2026-08-09T00:00:00Z",
"startedAt": null,
"finishedAt": null
})))
.expect(1)
.mount(&server)
.await;
let client = InfraClient::builder(server.uri()).build().unwrap();
let handle = client
.deploy()
.rollback(
"deployment-1",
&RollbackDeploymentRequest {
target_deployment_id: "deployment-0".into(),
environment: "production".into(),
},
"rollback-key",
)
.await
.unwrap();
assert_eq!(handle.id(), "operation-rollback");
}
#[tokio::test]
async fn gateway_exchange_credentials_are_audience_bound_cached_and_redacted() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/v1/gateway/credentials:exchange"))
.and(header("authorization", "Bearer external-workload-token"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"accessToken": "short-lived-context-token",
"tokenType": "Bearer",
"expiresInSeconds": 120
})))
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/internal/v1/context/messages/recent"))
.and(header("authorization", "Bearer short-lived-context-token"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({"items": []})))
.expect(2)
.mount(&server)
.await;
let exchange = GatewayExchangeCredentials::new(
server.uri(),
Arc::new(StaticCredentials::new("external-workload-token").unwrap()),
BTreeMap::from([(
"agent-context-infra".to_string(),
vec!["messages:read".to_string()],
)]),
)
.unwrap();
assert!(!format!("{exchange:?}").contains("external-workload-token"));
let client = InfraClient::builder(server.uri())
.credentials(Arc::new(exchange))
.build()
.unwrap();
for _ in 0..2 {
assert!(
client
.context()
.recent_messages("conversation", Some(1), &[])
.await
.unwrap()
.is_empty()
);
}
}
#[tokio::test]
async fn runtime_identity_client_uses_the_gateway_admin_audience() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/v1/gateway/credentials:exchange"))
.and(header("authorization", "Bearer external-workload-token"))
.and(body_json(serde_json::json!({
"audience": "agent-runtime-identity-admin",
"capabilities": ["runtime-identities:read"]
})))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"accessToken": "short-lived-runtime-token",
"tokenType": "Bearer",
"expiresInSeconds": 120
})))
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/internal/v1/runtime-identities:list"))
.and(header("authorization", "Bearer short-lived-runtime-token"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!([])))
.expect(1)
.mount(&server)
.await;
let exchange = GatewayExchangeCredentials::new(
server.uri(),
Arc::new(StaticCredentials::new("external-workload-token").unwrap()),
BTreeMap::from([(
"agent-runtime-identity-admin".to_string(),
vec!["runtime-identities:read".to_string()],
)]),
)
.unwrap();
let client = InfraClient::builder(server.uri())
.credentials(Arc::new(exchange))
.build()
.unwrap();
assert!(
client
.runtime_identity()
.list(&agent_runtime_identity_contract::ListRuntimeIdentitiesRequest::default())
.await
.unwrap()
.is_empty()
);
}
#[tokio::test]
async fn gateway_exchange_refreshes_independent_audiences_concurrently() {
async fn exchange(
axum::extract::State(barrier): axum::extract::State<Arc<tokio::sync::Barrier>>,
) -> axum::Json<serde_json::Value> {
barrier.wait().await;
axum::Json(serde_json::json!({
"accessToken": "audience-token",
"tokenType": "Bearer",
"expiresInSeconds": 120
}))
}
let barrier = Arc::new(tokio::sync::Barrier::new(2));
let app = axum::Router::new()
.route(
"/v1/gateway/credentials:exchange",
axum::routing::post(exchange),
)
.with_state(barrier);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
let provider = GatewayExchangeCredentials::new(
format!("http://{address}"),
Arc::new(StaticCredentials::new("external-workload-token").unwrap()),
BTreeMap::from([
(
"agent-context-infra".to_string(),
vec!["messages:read".to_string()],
),
(
"agent-deploy".to_string(),
vec!["deployments:read".to_string()],
),
]),
)
.unwrap();
let (context, deploy) = tokio::time::timeout(Duration::from_secs(1), async {
tokio::join!(
provider.credential("agent-context-infra"),
provider.credential("agent-deploy")
)
})
.await
.expect("audience refreshes must not share a global lock");
assert!(context.is_ok());
assert!(deploy.is_ok());
server.abort();
}
#[tokio::test]
async fn gateway_exchange_retries_only_bounded_transient_failures() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/v1/gateway/credentials:exchange"))
.respond_with(ResponseTemplate::new(503))
.up_to_n_times(1)
.with_priority(1)
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/v1/gateway/credentials:exchange"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"accessToken": "recovered-token",
"tokenType": "Bearer",
"expiresInSeconds": 120
})))
.with_priority(2)
.expect(1)
.mount(&server)
.await;
let provider = GatewayExchangeCredentials::new(
server.uri(),
Arc::new(StaticCredentials::new("external-workload-token").unwrap()),
BTreeMap::from([(
"agent-context-infra".to_string(),
vec!["messages:read".to_string()],
)]),
)
.unwrap();
assert!(provider.credential("agent-context-infra").await.is_ok());
}
#[tokio::test]
async fn sends_service_credential_and_preserves_base_path() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/gateway/internal/v1/computers"))
.and(header("authorization", "Bearer workspace-secret"))
.and(header(
"user-agent",
concat!("agent-infra-sdk/", env!("CARGO_PKG_VERSION")),
))
.and(header_regex(
"traceparent",
"^00-[0-9a-f]{32}-[0-9a-f]{16}-01$",
))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!([])))
.expect(1)
.mount(&server)
.await;
let client = InfraClient::builder(format!("{}/gateway/", server.uri()))
.credentials(Arc::new(
StaticCredentials::new("workspace-secret").expect("valid credential"),
))
.build()
.expect("build client");
assert!(
client
.environment()
.computers(None)
.await
.expect("computers")
.is_empty()
);
assert!(!format!("{client:?}").contains("workspace-secret"));
}
#[tokio::test]
async fn maps_non_success_status_to_a_typed_error() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/internal/v1/computers"))
.respond_with(
ResponseTemplate::new(503)
.set_body_json(serde_json::json!({"message": "workspace unavailable"})),
)
.mount(&server)
.await;
let client = client(&server);
let error = client
.environment()
.computers(None)
.await
.expect_err("503 must not be accepted");
assert!(matches!(
error,
InfraClientError::HttpStatus {
status: 503,
ref message,
..
} if message == "workspace unavailable"
));
}
#[tokio::test]
async fn preserves_structured_error_and_request_id() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/internal/v1/computers"))
.respond_with(
ResponseTemplate::new(409)
.insert_header("x-request-id", "req-header")
.set_body_json(serde_json::json!({
"code": "VERSION_CONFLICT",
"message": "stale revision",
"retryable": false,
"requestId": "req-body"
})),
)
.mount(&server)
.await;
let client = client(&server);
let error = client
.environment()
.computers(None)
.await
.expect_err("conflict must be returned");
assert_eq!(error.code(), "VERSION_CONFLICT");
assert!(!error.retryable());
assert_eq!(error.request_id(), Some("req-body"));
}
#[tokio::test]
async fn retries_only_an_idempotent_environment_call() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/internal/v1/computers"))
.respond_with(
ResponseTemplate::new(503)
.insert_header("retry-after", "0")
.set_body_json(serde_json::json!({
"code": "UNAVAILABLE",
"message": "try again",
"retryable": true
})),
)
.up_to_n_times(1)
.with_priority(1)
.expect(1)
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/internal/v1/computers"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!([])))
.with_priority(2)
.expect(1)
.mount(&server)
.await;
let client = client(&server);
assert!(client.environment().computers(None).await.unwrap().is_empty());
}
#[tokio::test]
async fn retries_context_append_because_caller_message_ids_deduplicate_replay() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/internal/v1/context/messages/batch"))
.respond_with(
ResponseTemplate::new(503)
.insert_header("retry-after", "0")
.set_body_json(serde_json::json!({
"code": "UNAVAILABLE",
"retryable": true
})),
)
.up_to_n_times(1)
.with_priority(1)
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/internal/v1/context/messages/batch"))
.respond_with(
ResponseTemplate::new(200)
.set_body_json(serde_json::json!({"appended": 1, "skipped": 0})),
)
.with_priority(2)
.expect(1)
.mount(&server)
.await;
let client = InfraClient::builder(server.uri()).build().unwrap();
let response = client
.context()
.append_messages(vec![ConversationRecord {
id: "message-stable-id".into(),
scope: ConversationScope::conversation("conversation-1"),
role: "user".into(),
content: serde_json::json!("hello"),
name: None,
tool_call_id: None,
sequence: 0,
created_at_ms: 0,
metadata: serde_json::Value::Null,
}])
.await
.unwrap();
assert_eq!(response.appended, 1);
}
#[tokio::test]
async fn sends_idempotency_key_from_per_call_options() {
let server = MockServer::start().await;
Mock::given(method("PUT"))
.and(path("/internal/v1/workspaces/ws-sdk/files/state.txt"))
.and(header("idempotency-key", "write-operation-42"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"ok": true,
"revision": "fnv1a64:abc"
})))
.expect(1)
.mount(&server)
.await;
let client = client(&server);
let revision = client
.workspace()
.put_workspace_file(
"ws-sdk",
"state.txt",
&agent_infra_sdk::workspace_contract::PutWorkspaceFileRequest {
content_base64: "Yw==".into(),
if_match: None,
},
CallOptions::default().idempotency_key("write-operation-42"),
)
.await
.unwrap();
assert_eq!(revision.as_deref(), Some("fnv1a64:abc"));
}
#[tokio::test]
async fn gateway_builder_constructs_every_enabled_domain_client() {
let server = MockServer::start().await;
let client = InfraClient::builder(server.uri())
.build()
.expect("build Gateway-backed client");
let _ = (
client.admin(),
client.agents(),
client.deploy(),
client.evaluate(),
client.context(),
client.workspace(),
client.environment(),
client.gateway(),
client.runtime_identity(),
);
}
#[tokio::test]
async fn rejects_success_status_with_negative_acknowledgement() {
let server = MockServer::start().await;
Mock::given(method("PUT"))
.and(path("/internal/v1/workspaces/ws-sdk/files/result.txt"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({"ok": false})))
.mount(&server)
.await;
let client = client(&server);
let error = client
.workspace()
.put_workspace_file(
"ws-sdk",
"result.txt",
&agent_infra_sdk::workspace_contract::PutWorkspaceFileRequest {
content_base64: "Yw==".into(),
if_match: None,
},
CallOptions::default(),
)
.await
.expect_err("ok=false must not be accepted");
assert!(matches!(error, InfraClientError::Protocol { .. }));
}
#[tokio::test]
async fn rejects_missing_required_response_fields_instead_of_defaulting() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/internal/v1/runtime-identities:list"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({})))
.mount(&server)
.await;
let client = InfraClient::builder(server.uri())
.build()
.expect("build runtime identity client");
let error = client
.runtime_identity()
.list(&agent_runtime_identity_contract::ListRuntimeIdentitiesRequest::default())
.await
.expect_err("missing instances must fail decoding");
assert!(matches!(error, InfraClientError::Decode { .. }));
}
#[tokio::test]
async fn rejects_responses_above_the_configured_limit() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/internal/v1/computers"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"content": "a response that is intentionally too large"
})))
.mount(&server)
.await;
let client = InfraClient::builder(server.uri())
.max_response_bytes(16)
.build()
.expect("build client");
let error = client
.environment()
.computers(None)
.await
.expect_err("oversized response must fail");
assert!(matches!(
error,
InfraClientError::ResponseTooLarge { limit: 16, .. }
));
}
#[tokio::test]
async fn enforces_the_shared_request_timeout() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/internal/v1/computers"))
.respond_with(
ResponseTemplate::new(200)
.set_delay(Duration::from_millis(150))
.set_body_json(serde_json::json!([])),
)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/internal/v1/context/messages/recent"))
.respond_with(
ResponseTemplate::new(200)
.set_delay(Duration::from_millis(150))
.set_body_json(serde_json::json!({"items": []})),
)
.mount(&server)
.await;
let client = InfraClient::builder(server.uri())
.request_timeout(Duration::from_millis(25))
.build()
.expect("build client");
let error = client
.environment()
.computers(None)
.await
.expect_err("slow response must time out");
assert!(
matches!(&error, InfraClientError::DeadlineExceeded { .. })
|| matches!(
&error,
InfraClientError::Request { source, .. } if source.is_timeout()
)
);
let context_error = client
.context()
.recent_messages("slow", Some(1), &[])
.await
.expect_err("the same timeout must apply to the context facade");
assert!(
matches!(&context_error, InfraClientError::DeadlineExceeded { .. })
|| matches!(
&context_error,
InfraClientError::Request { source, .. } if source.is_timeout()
)
);
}
#[test]
fn rejects_unbounded_retry_configuration() {
let mut options = ClientOptions::default();
options.retry.max_attempts = 1_000_000;
let error = InfraClient::builder("http://127.0.0.1:1")
.options(options)
.build()
.expect_err("retry budget must have a hard bound");
assert_eq!(error.code(), "INVALID_OPTIONS");
}
#[tokio::test]
async fn gateway_ready_reports_ok_when_gateway_is_ready() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/health/ready"))
.respond_with(
ResponseTemplate::new(200)
.set_body_json(serde_json::json!({"status": "ok", "service": "infra-api-gateway"})),
)
.mount(&server)
.await;
let client = InfraClient::builder(server.uri()).build().unwrap();
let status = client.gateway().ready().await.unwrap();
assert!(status.is_ready());
}
#[tokio::test]
async fn gateway_ready_treats_503_as_not_ready_rather_than_error() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/health/ready"))
.respond_with(ResponseTemplate::new(503).set_body_json(
serde_json::json!({"status": "not_ready", "service": "infra-api-gateway"}),
))
.mount(&server)
.await;
let client = InfraClient::builder(server.uri()).build().unwrap();
let status = client.gateway().ready().await.unwrap();
assert!(!status.is_ready());
}
#[tokio::test]
async fn gateway_wait_until_ready_times_out_while_not_ready() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/health/ready"))
.respond_with(ResponseTemplate::new(503).set_body_json(
serde_json::json!({"status": "not_ready", "service": "infra-api-gateway"}),
))
.mount(&server)
.await;
let client = InfraClient::builder(server.uri()).build().unwrap();
let error = client
.gateway()
.wait_until_ready(Duration::from_millis(150))
.await
.expect_err("a gateway that never becomes ready must time out");
assert!(matches!(error, InfraClientError::DeadlineExceeded { .. }));
}
#[tokio::test]
async fn managed_workspace_optional_read_and_remove_use_resource_file_contract() {
let server = MockServer::start().await;
let file_path = "/internal/v1/workspaces/ws-sdk/files/src/lib.rs";
Mock::given(method("GET"))
.and(path(file_path))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"contentBase64": "aGVsbG8="
})))
.expect(1)
.mount(&server)
.await;
Mock::given(method("DELETE"))
.and(path(file_path))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"ok": true
})))
.expect(1)
.mount(&server)
.await;
let client = client(&server);
let file = client
.workspace()
.read_workspace_file_if_exists_with_options("ws-sdk", "src/lib.rs", CallOptions::default())
.await
.unwrap()
.expect("file should exist");
assert_eq!(file.content_base64, "aGVsbG8=");
client
.workspace()
.remove_workspace_file_with_options("ws-sdk", "src/lib.rs", CallOptions::default())
.await
.unwrap();
}
#[tokio::test]
async fn managed_workspace_optional_read_maps_not_found_to_none() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/internal/v1/workspaces/ws-sdk/files/missing.txt"))
.respond_with(ResponseTemplate::new(404).set_body_json(serde_json::json!({
"code": "WORKSPACE_FILE_NOT_FOUND",
"message": "file not found",
"retryable": false
})))
.expect(1)
.mount(&server)
.await;
let result = client(&server)
.workspace()
.read_workspace_file_if_exists_with_options("ws-sdk", "missing.txt", CallOptions::default())
.await
.unwrap();
assert!(result.is_none());
}