use super::*;
use chio_core_types::capability::governance::{
GovernedTransactionIntent, ThresholdApprovalProposal,
};
use chio_core_types::capability::supplemental_authorization::OpaqueSupplementalAuthorization;
use chio_kernel::budget_store::BudgetStore;
use chio_kernel::dpop::{DpopConfig, DpopNonceStore, DpopProof};
use chio_kernel::execution_nonce::{
ExecutionNonceConfig, InMemoryExecutionNonceStore, SignedExecutionNonce,
};
use chio_kernel::{
ChioKernel, KernelConfig, KernelError, ToolCallRequest, ToolInvocationCost,
ToolServerConnection, DEFAULT_CHECKPOINT_BATCH_SIZE, DEFAULT_MAX_STREAM_DURATION_SECS,
DEFAULT_MAX_STREAM_TOTAL_BYTES,
};
#[path = "mediated/budget_configuration.rs"]
mod budget_configuration;
pub(crate) use budget_configuration::build_budget_store;
pub(crate) fn load_revocation_db_ids(
config: &ProtectConfig,
) -> Result<std::collections::HashSet<String>, ProtectError> {
let Some(path) = config.revocation_db.as_deref() else {
return Ok(std::collections::HashSet::new());
};
let store = chio_store_sqlite::SqliteRevocationStore::open(path).map_err(|error| {
ProtectError::Config(format!("cannot open revocation-db `{path}`: {error}"))
})?;
const PAGE_SIZE: usize = 1024;
let mut ids = std::collections::HashSet::new();
let mut cursor: Option<(i64, String)> = None;
loop {
let (after_revoked_at, after_capability_id) = match &cursor {
Some((revoked_at, capability_id)) => (Some(*revoked_at), Some(capability_id.as_str())),
None => (None, None),
};
let page = store
.list_revocations_after(PAGE_SIZE, after_revoked_at, after_capability_id)
.map_err(|error| {
ProtectError::Config(format!("cannot read revocation-db `{path}`: {error}"))
})?;
let page_len = page.len();
for record in page {
cursor = Some((record.revoked_at, record.capability_id.clone()));
ids.insert(record.capability_id);
}
if page_len < PAGE_SIZE {
break;
}
}
Ok(ids)
}
pub(crate) fn build_mediation_kernel(
signer: &Keypair,
budget_store: Arc<dyn BudgetStore>,
trusted_capability_issuers: &[PublicKey],
tool_servers: Vec<Box<dyn ToolServerConnection>>,
payment_adapter: Option<Box<dyn chio_kernel::PaymentAdapter>>,
durable_admission: Option<DurableAdmissionStores>,
) -> Result<ChioKernel, ProtectError> {
let mut ca_public_keys = vec![signer.public_key()];
for issuer in trusted_capability_issuers {
if !ca_public_keys.contains(issuer) {
ca_public_keys.push(issuer.clone());
}
}
let mut kernel = ChioKernel::new(KernelConfig {
keypair: signer.clone(),
ca_public_keys,
max_delegation_depth: 5,
policy_hash: "chio_api_protect_mediation_v1".to_string(),
allow_sampling: false,
allow_sampling_tool_use: false,
allow_elicitation: false,
max_stream_duration_secs: DEFAULT_MAX_STREAM_DURATION_SECS,
max_stream_total_bytes: DEFAULT_MAX_STREAM_TOTAL_BYTES,
require_web3_evidence: false,
allow_ephemeral_receipt_log: true,
allow_ephemeral_revocation_store: true,
checkpoint_batch_size: DEFAULT_CHECKPOINT_BATCH_SIZE,
retention_config: None,
memory_budget: chio_kernel::MemoryBudgetConfig::defaults(),
deadlines: chio_kernel::HotPathDeadlineConfig::default(),
});
kernel.set_budget_store_handle(budget_store);
if let Some(durable) = durable_admission {
kernel
.set_durable_admission_store(durable.store, durable.outcome_store, durable.fence)
.map_err(|error| {
ProtectError::Config(format!(
"failed to install durable admission stores on the mediation kernel: {error}"
))
})?;
}
let nonce_cfg = ExecutionNonceConfig {
require_nonce: true,
..ExecutionNonceConfig::default()
};
kernel.set_execution_nonce_store(
nonce_cfg,
Box::new(InMemoryExecutionNonceStore::from_config(
&ExecutionNonceConfig::default(),
)),
);
kernel.set_dpop_store(
DpopNonceStore::new(
DpopConfig::default().nonce_store_capacity,
std::time::Duration::from_secs(DpopConfig::default().proof_ttl_secs),
),
DpopConfig::default(),
);
if let Some(payment_adapter) = payment_adapter {
kernel.set_payment_adapter(payment_adapter);
}
for server in tool_servers {
kernel.register_tool_server(server);
}
kernel
.arm_restart_reserved_hold_gate()
.map_err(|error| ProtectError::Config(error.to_string()))?;
Ok(kernel)
}
#[derive(Debug, serde::Deserialize)]
pub(crate) struct SidecarEvaluateToolCallMediatedRequest {
capability: chio_core_types::capability::token::CapabilityToken,
tool_server: String,
tool_name: String,
#[serde(default)]
parameters: serde_json::Value,
#[serde(default)]
agent_id: Option<String>,
#[serde(default)]
request_id: Option<String>,
#[serde(default)]
governed_intent: Option<GovernedTransactionIntent>,
#[serde(default)]
approval_token: Option<GovernedApprovalToken>,
#[serde(default)]
approval_tokens: Vec<GovernedApprovalToken>,
#[serde(default)]
threshold_approval_proposal: Option<ThresholdApprovalProposal>,
#[serde(default)]
dpop_proof: Option<DpopProof>,
#[serde(default)]
supplemental_authorization: Option<OpaqueSupplementalAuthorization>,
#[serde(default)]
execution_nonce: Option<SignedExecutionNonce>,
}
pub(crate) async fn sidecar_evaluate_tool_call_mediated_handler(
State(state): State<Arc<ProxyState>>,
request: Request<Body>,
) -> Response {
let (_parts, body) = request.into_parts();
let body_bytes = match axum::body::to_bytes(body, 1024 * 1024).await {
Ok(bytes) => bytes,
Err(error) => {
warn!("failed to read mediated evaluate body: {error}");
return sidecar_bad_request("failed to read evaluate body").into_response();
}
};
let parsed: SidecarEvaluateToolCallMediatedRequest = match serde_json::from_slice(&body_bytes) {
Ok(parsed) => parsed,
Err(error) => {
return sidecar_bad_request(&format!("invalid mediated payload: {error}"))
.into_response();
}
};
if parsed.supplemental_authorization.is_some() {
return sidecar_bad_request(
"supplemental_authorization is unavailable: no supplemental verifier is configured",
)
.into_response();
}
if !parsed.approval_tokens.is_empty() || parsed.threshold_approval_proposal.is_some() {
return sidecar_bad_request(
"threshold approvals are unavailable on the mediated endpoint: no threshold policy resolver is configured",
)
.into_response();
}
if parsed.execution_nonce.is_some() {
return sidecar_bad_request(
"/v1/evaluate issues execution nonces; it does not accept a presented nonce. \
Present the minted nonce to the tool server, not to this endpoint",
)
.into_response();
}
let Some(mediation_kernel) = state.mediation_kernel.as_ref() else {
return internal_json_error_response(
"chio_mediation_unavailable",
"mediated tool-call route requires a configured budget store (--control-url or --budget-db)",
);
};
if !state.mediation_hold_capable {
return internal_json_error_response(
"chio_mediation_requires_local_budget_store",
"mediated authorization requires a hold-capable local budget store (--budget-db); \
a remote control-plane budget store (--control-url) cannot persist a reserved hold",
);
}
if state
.sidecar_control_token
.as_deref()
.map(str::trim)
.is_none_or(str::is_empty)
{
return internal_json_error_response(
"chio_mediation_requires_reconcile_token",
"mediated authorization requires a configured sidecar-control token so the reserved \
hold can be settled on /v1/reconcile; without one the reservation could only expire \
and forfeit budget",
);
}
const MAX_MEDIATED_DELEGATION_CHAIN: usize = 32;
if parsed.capability.delegation_chain.len() > MAX_MEDIATED_DELEGATION_CHAIN {
return sidecar_bad_request("capability delegation chain is too long").into_response();
}
let mut revoked = state.capability_is_revoked(&parsed.capability.id).await;
if !revoked {
for ancestor in &parsed.capability.delegation_chain {
if state.capability_is_revoked(&ancestor.capability_id).await {
revoked = true;
break;
}
}
}
if revoked {
return (
StatusCode::FORBIDDEN,
axum::Json(serde_json::json!({
"error": "chio_capability_revoked",
"message": "capability has been revoked",
})),
)
.into_response();
}
let agent_id = parsed
.agent_id
.unwrap_or_else(|| parsed.capability.subject.to_hex());
let request_id = parsed
.request_id
.unwrap_or_else(|| uuid::Uuid::now_v7().to_string());
let Some(budget_store) = state.budget_store.as_ref() else {
return internal_json_error_response(
"chio_mediation_requires_local_budget_store",
"mediated authorization requires a configured hold-capable budget store",
);
};
const MAX_MEDIATED_SCOPE_GRANTS: usize = 64;
if parsed.capability.scope.grants.len() > MAX_MEDIATED_SCOPE_GRANTS {
return sidecar_bad_request("capability scope carries too many grants").into_response();
}
match budget_store.request_id_has_reserved_hold(&request_id) {
Ok(Some(true)) => {
return (
StatusCode::CONFLICT,
axum::Json(serde_json::json!({
"error": "chio_request_id_reused",
"message": "request_id already backs a reservation; choose a fresh request_id",
})),
)
.into_response();
}
Ok(Some(false)) => {}
Ok(None) => {
for grant_index in 0..parsed.capability.scope.grants.len() {
let hold_id = format!(
"budget-hold:{}:{}:{}",
request_id, parsed.capability.id, grant_index
);
match budget_store.get_budget_hold(&hold_id) {
Ok(Some(_)) => {
return (
StatusCode::CONFLICT,
axum::Json(serde_json::json!({
"error": "chio_request_id_reused",
"message": "request_id already backs a reservation; choose a fresh request_id",
})),
)
.into_response();
}
Ok(None) => {}
Err(error) => {
warn!("durable hold lookup failed: {error}");
return internal_json_error_response(
"chio_mediation_failed",
&error.to_string(),
);
}
}
}
}
Err(error) => {
warn!("durable hold lookup failed: {error}");
return internal_json_error_response("chio_mediation_failed", &error.to_string());
}
}
let now_unix = chrono::Utc::now().timestamp();
if !state
.minted_request_ids
.lock()
.await
.claim(&request_id, now_unix)
{
return (
StatusCode::CONFLICT,
axum::Json(serde_json::json!({
"error": "chio_request_id_reused",
"message":
"request_id has already been used for a reservation; choose a fresh request_id",
})),
)
.into_response();
}
let kernel_request = ToolCallRequest {
request_id: request_id.clone(),
capability: parsed.capability,
tool_name: parsed.tool_name,
server_id: parsed.tool_server,
agent_id,
arguments: parsed.parameters,
dpop_proof: parsed.dpop_proof,
execution_nonce: None,
governed_intent: parsed.governed_intent,
approval_token: parsed.approval_token,
approval_tokens: Vec::new(),
threshold_approval_proposal: None,
supplemental_authorization: None,
model_metadata: None,
federated_origin_kernel_id: None,
};
let response = {
let kernel = mediation_kernel.lock().await;
match kernel.authorize_tool_call_reserving_blocking_with_metadata(&kernel_request, None) {
Ok(response) => response,
Err(error) => {
state.minted_request_ids.lock().await.release(&request_id);
warn!("mediated authorization error: {error}");
return internal_json_error_response("chio_mediation_failed", &error.to_string());
}
}
};
if !matches!(response.verdict, chio_kernel::Verdict::Allow) {
state.minted_request_ids.lock().await.release(&request_id);
}
if let Err(error) = record_tool_receipt(&state, &response.receipt).await {
if !matches!(
(&response.verdict, response.execution_nonce.as_deref()),
(chio_kernel::Verdict::Allow, Some(_))
) {
warn!("failed to persist mediated receipt: {error}");
return internal_json_error_response(
"chio_receipt_persistence_failed",
&error.to_string(),
);
}
warn!(
"mediated reserve receipt persistence failed; returning minted nonce to caller: {error}"
);
}
let status_str = match &response.verdict {
chio_kernel::Verdict::Allow => "authorized",
chio_kernel::Verdict::Deny => "deny",
chio_kernel::Verdict::PendingApproval => "pending_approval",
};
(
StatusCode::OK,
axum::Json(serde_json::json!({
"status": status_str,
"receipt": response.receipt,
"execution_nonce": response.execution_nonce,
})),
)
.into_response()
}
#[derive(Debug, serde::Deserialize)]
pub(crate) struct SidecarReconcileRequest {
execution_nonce: SignedExecutionNonce,
#[serde(default)]
arguments: serde_json::Value,
realized_cost: ToolInvocationCost,
}
pub(crate) async fn sidecar_reconcile_handler(
State(state): State<Arc<ProxyState>>,
request: Request<Body>,
) -> Response {
let (_parts, body) = request.into_parts();
let body_bytes = match axum::body::to_bytes(body, 1024 * 1024).await {
Ok(bytes) => bytes,
Err(error) => {
warn!("failed to read reconcile body: {error}");
return sidecar_bad_request("failed to read reconcile body").into_response();
}
};
let parsed: SidecarReconcileRequest = match serde_json::from_slice(&body_bytes) {
Ok(parsed) => parsed,
Err(error) => {
return sidecar_bad_request(&format!("invalid reconcile payload: {error}"))
.into_response();
}
};
let Some(mediation_kernel) = state.mediation_kernel.as_ref() else {
return internal_json_error_response(
"chio_mediation_unavailable",
"reconcile route requires a configured budget store (--control-url or --budget-db)",
);
};
if !state.mediation_hold_capable {
return internal_json_error_response(
"chio_mediation_requires_local_budget_store",
"mediated reconcile requires a hold-capable local budget store (--budget-db); \
a remote control-plane budget store (--control-url) cannot resolve a reserved hold",
);
}
let reconciled = {
let kernel = mediation_kernel.lock().await;
kernel.reconcile_reserved_authorization_by_nonce(
&parsed.execution_nonce,
&parsed.arguments,
&parsed.realized_cost,
)
};
let reconciled = match reconciled {
Ok(response) => response,
Err(error) => {
warn!("reconcile rejected: {error}");
return (
StatusCode::BAD_REQUEST,
axum::Json(serde_json::json!({
"error": "chio_reconcile_rejected",
"message": error.to_string(),
})),
)
.into_response();
}
};
if let Err(error) = record_tool_receipt(&state, &reconciled.receipt).await {
warn!("reconcile settled but receipt persistence failed; returning authoritative receipt to caller: {error}");
}
(
StatusCode::OK,
axum::Json(serde_json::json!({
"status": "reconciled",
"receipt": reconciled.receipt,
})),
)
.into_response()
}
pub(crate) async fn reap_expired_reserved_holds_once(
state: &Arc<ProxyState>,
now_unix_secs: i64,
) -> Result<usize, KernelError> {
let Some(mediation_kernel) = state.mediation_kernel.as_ref() else {
return Ok(0);
};
let kernel = mediation_kernel.lock().await;
kernel.reap_expired_reserved_budget_holds(now_unix_secs)
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used)]
mod tests {
use super::*;
use chio_kernel::budget_store::{BudgetStore, InMemoryBudgetStore};
use chio_test_support::prelude::*;
use tower::ServiceExt;
#[path = "mediated_boundary_tests.rs"]
mod boundary_tests;
fn issuing_kernel(
signer: &Keypair,
budget: Arc<dyn BudgetStore>,
trusted_capability_issuers: &[PublicKey],
) -> Arc<ChioKernel> {
Arc::new(
build_mediation_kernel(
signer,
budget,
trusted_capability_issuers,
Vec::new(),
None,
None,
)
.test_unwrap(),
)
}
fn issue_cost_bearing_capability(
kernel: &Arc<ChioKernel>,
agent: &Keypair,
server: &str,
tool: &str,
max_per: u64,
max_total: u64,
currency: &str,
) -> CapabilityToken {
use chio_core_types::capability::scope::MonetaryAmount;
let grant = ToolGrant {
server_id: server.to_string(),
tool_name: tool.to_string(),
operations: vec![Operation::Invoke],
constraints: vec![],
max_invocations: None,
max_cost_per_invocation: Some(MonetaryAmount {
units: max_per,
currency: currency.to_string(),
}),
max_total_cost: Some(MonetaryAmount {
units: max_total,
currency: currency.to_string(),
}),
dpop_required: None,
};
let scope = ChioScope {
grants: vec![grant],
..ChioScope::default()
};
kernel
.issue_capability(&agent.public_key(), scope, 3600)
.test_unwrap()
}
fn issue_invocation_capability(
kernel: &Arc<ChioKernel>,
agent: &Keypair,
server: &str,
tool: &str,
max_invocations: u32,
) -> CapabilityToken {
let grant = ToolGrant {
server_id: server.to_string(),
tool_name: tool.to_string(),
operations: vec![Operation::Invoke],
constraints: vec![],
max_invocations: Some(max_invocations),
max_cost_per_invocation: None,
max_total_cost: None,
dpop_required: None,
};
let scope = ChioScope {
grants: vec![grant],
..ChioScope::default()
};
kernel
.issue_capability(&agent.public_key(), scope, 3600)
.test_unwrap()
}
const MEDIATED_CONTROL_TOKEN: &str = "tool-server-control-token";
fn mediated_test_state(
signer: Keypair,
budget: Arc<dyn BudgetStore>,
trusted_capability_issuers: Vec<PublicKey>,
) -> Arc<ProxyState> {
mediated_test_state_with_control_token(
signer,
budget,
trusted_capability_issuers,
Some(MEDIATED_CONTROL_TOKEN.to_string()),
)
}
fn mediated_test_state_with_control_token(
signer: Keypair,
budget: Arc<dyn BudgetStore>,
trusted_capability_issuers: Vec<PublicKey>,
sidecar_control_token: Option<String>,
) -> Arc<ProxyState> {
mediated_test_state_inner(
signer,
budget,
trusted_capability_issuers,
sidecar_control_token,
None,
true,
)
}
fn mediated_test_state_inner(
signer: Keypair,
budget: Arc<dyn BudgetStore>,
trusted_capability_issuers: Vec<PublicKey>,
sidecar_control_token: Option<String>,
receipt_store: Option<SqliteReceiptStore>,
hold_capable: bool,
) -> Arc<ProxyState> {
mediated_test_state_core(
signer,
budget,
trusted_capability_issuers,
sidecar_control_token,
receipt_store,
hold_capable,
None,
None,
)
}
#[allow(clippy::too_many_arguments)]
fn mediated_test_state_core(
signer: Keypair,
budget: Arc<dyn BudgetStore>,
trusted_capability_issuers: Vec<PublicKey>,
sidecar_control_token: Option<String>,
receipt_store: Option<SqliteReceiptStore>,
hold_capable: bool,
payment_adapter: Option<Box<dyn chio_kernel::PaymentAdapter>>,
revocation_store: Option<Arc<dyn chio_kernel::RevocationStore>>,
) -> Arc<ProxyState> {
let approval_store: Arc<dyn ApprovalStore> = Arc::new(InMemoryApprovalStore::new());
let signer_public_key = signer.public_key();
let mut trusted_capability_issuers = trusted_capability_issuers;
if !trusted_capability_issuers.contains(&signer_public_key) {
trusted_capability_issuers.push(signer_public_key.clone());
}
let trusted_receipt_signers = vec![signer_public_key];
let evaluator = RequestEvaluator::new_ephemeral_with_approval_store(
Vec::new(),
signer.clone(),
"test-policy".to_string(),
Arc::clone(&approval_store),
);
let egress_contract = default_upstream_egress_contract("http://127.0.0.1:1").test_unwrap();
let http_client = client_builder_with_contract(&egress_contract)
.build()
.test_unwrap();
let mediation_kernel = Mutex::new(
build_mediation_kernel(
&signer,
Arc::clone(&budget),
&trusted_capability_issuers,
Vec::new(),
payment_adapter,
None,
)
.test_unwrap(),
);
Arc::new(ProxyState {
evaluator,
signer_keypair: signer,
upstream: "http://127.0.0.1:1".to_string(),
http_client,
egress_contract,
approval_admin: ApprovalAdmin::new(approval_store),
receipt_log: Mutex::new(ReceiptLog {
receipts: Vec::new(),
}),
tool_receipt_log: Mutex::new(ToolReceiptLog {
receipts: Vec::new(),
}),
receipt_store: receipt_store.map(Mutex::new),
revocation_store,
revoked_capability_ids: Mutex::new(std::collections::HashSet::new()),
trusted_capability_issuers,
trusted_receipt_signers,
sidecar_control_token,
budget_store: Some(budget),
mediation_hold_capable: hold_capable,
mediation_kernel: Some(mediation_kernel),
minted_request_ids: Mutex::new(MintedRequestIdWindow::new(
chio_kernel::DEFAULT_EXECUTION_NONCE_TTL_SECS,
)),
reaper_handle: Mutex::new(None),
allow_advisory: false,
receipt_backend: "ephemeral",
revocation_backend: "ephemeral",
})
}
fn with_loopback_peer(request: axum::http::Request<Body>) -> axum::http::Request<Body> {
use axum::extract::ConnectInfo;
let mut request = request;
request
.extensions_mut()
.insert(ConnectInfo(std::net::SocketAddr::from((
[127, 0, 0, 1],
4100,
))));
request
}
use chio_core_types::capability::governance::{
GovernedApprovalDecision, GovernedApprovalToken, GovernedApprovalTokenBody,
GovernedTransactionIntent, MeteredBillingContext, MeteredBillingQuote,
MeteredSettlementMode,
};
use chio_core_types::capability::scope::{Constraint, MonetaryAmount};
use chio_core_types::receipt::authoritative_spend::is_authoritative_spend_receipt;
use chio_kernel::dpop::{DpopProof, DpopProofBody, DPOP_SCHEMA};
fn issue_governed_capability(
kernel: &Arc<ChioKernel>,
agent: &Keypair,
server: &str,
tool: &str,
max_per: u64,
currency: &str,
approval_threshold_units: u64,
) -> CapabilityToken {
let grant = ToolGrant {
server_id: server.to_string(),
tool_name: tool.to_string(),
operations: vec![Operation::Invoke],
constraints: vec![
Constraint::GovernedIntentRequired,
Constraint::RequireApprovalAbove {
threshold_units: approval_threshold_units,
},
],
max_invocations: None,
max_cost_per_invocation: Some(MonetaryAmount {
units: max_per,
currency: currency.to_string(),
}),
max_total_cost: Some(MonetaryAmount {
units: max_per,
currency: currency.to_string(),
}),
dpop_required: None,
};
let scope = ChioScope {
grants: vec![grant],
..ChioScope::default()
};
kernel
.issue_capability(&agent.public_key(), scope, 3600)
.test_unwrap()
}
fn issue_dpop_capability(
kernel: &Arc<ChioKernel>,
agent: &Keypair,
server: &str,
tool: &str,
max_per: u64,
currency: &str,
) -> CapabilityToken {
issue_dpop_capability_with_total(kernel, agent, server, tool, max_per, max_per, currency)
}
fn issue_dpop_capability_with_total(
kernel: &Arc<ChioKernel>,
agent: &Keypair,
server: &str,
tool: &str,
max_per: u64,
max_total: u64,
currency: &str,
) -> CapabilityToken {
let grant = ToolGrant {
server_id: server.to_string(),
tool_name: tool.to_string(),
operations: vec![Operation::Invoke],
constraints: vec![],
max_invocations: None,
max_cost_per_invocation: Some(MonetaryAmount {
units: max_per,
currency: currency.to_string(),
}),
max_total_cost: Some(MonetaryAmount {
units: max_total,
currency: currency.to_string(),
}),
dpop_required: Some(true),
};
let scope = ChioScope {
grants: vec![grant],
..ChioScope::default()
};
kernel
.issue_capability(&agent.public_key(), scope, 3600)
.test_unwrap()
}
fn governed_intent(
id: &str,
server: &str,
tool: &str,
units: u64,
currency: &str,
) -> GovernedTransactionIntent {
GovernedTransactionIntent {
id: id.to_string(),
server_id: server.to_string(),
tool_name: tool.to_string(),
purpose: "invoice-settlement".to_string(),
max_amount: Some(MonetaryAmount {
units,
currency: currency.to_string(),
}),
commerce: None,
metered_billing: None,
runtime_attestation: None,
call_chain: None,
autonomy: None,
context: None,
body: Default::default(),
}
}
fn governed_mustprepay_intent(
id: &str,
server: &str,
tool: &str,
units: u64,
currency: &str,
) -> GovernedTransactionIntent {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.test_unwrap()
.as_secs();
let mut intent = governed_intent(id, server, tool, units, currency);
intent.metered_billing = Some(MeteredBillingContext {
settlement_mode: MeteredSettlementMode::MustPrepay,
quote: MeteredBillingQuote {
quote_id: format!("quote-{id}"),
provider: "billing.chio".to_string(),
billing_unit: "call".to_string(),
quoted_units: 1,
quoted_cost: MonetaryAmount {
units,
currency: currency.to_string(),
},
issued_at: now.saturating_sub(5),
expires_at: Some(now + 300),
},
max_billed_units: Some(2),
verified_outcome: None,
});
intent
}
fn governed_approval_token(
approver: &Keypair,
subject: &PublicKey,
intent: &GovernedTransactionIntent,
request_id: &str,
) -> GovernedApprovalToken {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.test_unwrap()
.as_secs();
GovernedApprovalToken::sign(
GovernedApprovalTokenBody {
id: format!("approval-{request_id}"),
approver: approver.public_key(),
subject: subject.clone(),
governed_intent_hash: intent.binding_hash().test_unwrap(),
request_id: request_id.to_string(),
threshold_proposal_hash: None,
issued_at: now.saturating_sub(1),
expires_at: now + 300,
decision: GovernedApprovalDecision::Approved,
},
approver,
)
.test_unwrap()
}
fn dpop_proof_for(
agent: &Keypair,
cap: &CapabilityToken,
server: &str,
tool: &str,
parameters: &serde_json::Value,
) -> DpopProof {
let args_bytes = chio_core_types::canonical::canonical_json_bytes(parameters).test_unwrap();
let action_hash = chio_core_types::crypto::sha256_hex(&args_bytes);
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.test_unwrap()
.as_secs();
DpopProof::sign(
DpopProofBody {
schema: DPOP_SCHEMA.to_string(),
capability_id: cap.id.clone(),
tool_server: server.to_string(),
tool_name: tool.to_string(),
action_hash,
nonce: uuid::Uuid::now_v7().to_string(),
issued_at: now,
agent_key: agent.public_key(),
},
agent,
)
.test_unwrap()
}
async fn post_evaluate(
state: Arc<ProxyState>,
body: &serde_json::Value,
) -> (StatusCode, serde_json::Value) {
post_json(state, "/v1/evaluate", body).await
}
async fn post_reconcile(
state: Arc<ProxyState>,
body: &serde_json::Value,
) -> (StatusCode, serde_json::Value) {
post_json_with_bearer(state, "/v1/reconcile", body, Some(MEDIATED_CONTROL_TOKEN)).await
}
async fn post_json(
state: Arc<ProxyState>,
uri: &str,
body: &serde_json::Value,
) -> (StatusCode, serde_json::Value) {
post_json_with_bearer(state, uri, body, None).await
}
async fn post_json_with_bearer(
state: Arc<ProxyState>,
uri: &str,
body: &serde_json::Value,
bearer: Option<&str>,
) -> (StatusCode, serde_json::Value) {
let mut builder = Request::builder()
.method("POST")
.uri(uri)
.header("content-type", "application/json");
if let Some(bearer) = bearer {
builder = builder.header("authorization", format!("Bearer {bearer}"));
}
let request = with_loopback_peer(
builder
.body(Body::from(serde_json::to_vec(body).unwrap()))
.unwrap(),
);
let response = build_app(state).oneshot(request).await.unwrap();
let status = response.status();
let bytes = axum::body::to_bytes(response.into_body(), 1 << 20)
.await
.unwrap();
let json: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
(status, json)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_authorization_reserves_hold_and_mints_non_authoritative_receipt() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 1000, "USD");
let cap_id = cap.id.clone();
let signer_pub = signer.public_key();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let body = serde_json::json!({
"capability": cap,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" }
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(json["status"], "authorized");
assert!(
json["execution_nonce"].is_object(),
"authorization must mint an execution nonce object"
);
assert_eq!(
json["receipt"]["decision"]["verdict"], "incomplete",
"a reserved authorization is not a completed-spend decision"
);
let budget_authority = &json["receipt"]["metadata"]["budget_authority"];
assert_eq!(budget_authority["authorize"]["exposure_units"], 100);
assert!(
budget_authority.get("terminal").is_none(),
"a reserved (not reconciled) hold must carry no terminal disposition"
);
let receipt: ChioReceipt = serde_json::from_value(json["receipt"].clone()).unwrap();
let nonce: SignedExecutionNonce =
serde_json::from_value(json["execution_nonce"].clone()).unwrap();
assert!(
is_authoritative_spend_receipt(&receipt, &[signer_pub], &nonce).is_err(),
"a reserved authorization receipt must not be an authoritative spend"
);
let usage = budget.get_usage(&cap_id, 0).unwrap();
let usage = usage.expect("the reserved hold must be recorded in the budget store");
assert_eq!(
usage.committed_cost_units().unwrap(),
100,
"the pre-execution hold must remain reserved (open), not reversed"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_durable_reuse_guard_rejects_reused_request_id_across_capabilities() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap_a =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 1000, "USD");
let cap_b =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 1000, "USD");
assert_ne!(
cap_a.id, cap_b.id,
"the two capabilities must carry distinct ids"
);
let state_before = mediated_test_state(signer.clone(), Arc::clone(&budget), Vec::new());
let body_a = serde_json::json!({
"capability": cap_a,
"tool_server": "cost-srv",
"tool_name": "compute",
"request_id": "shared-request-id",
"parameters": { "invoice": "inv-1" }
});
let (status_a, json_a) = post_evaluate(Arc::clone(&state_before), &body_a).await;
assert_eq!(status_a, StatusCode::OK, "{json_a}");
assert_eq!(json_a["status"], "authorized");
let state_after = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let body_b = serde_json::json!({
"capability": cap_b,
"tool_server": "cost-srv",
"tool_name": "compute",
"request_id": "shared-request-id",
"parameters": { "invoice": "inv-2" }
});
let (status_b, json_b) = post_evaluate(Arc::clone(&state_after), &body_b).await;
assert_eq!(status_b, StatusCode::CONFLICT, "{json_b}");
assert_eq!(json_b["error"], "chio_request_id_reused");
let body_fresh = serde_json::json!({
"capability": cap_b,
"tool_server": "cost-srv",
"tool_name": "compute",
"request_id": "fresh-request-id",
"parameters": { "invoice": "inv-3" }
});
let (status_fresh, json_fresh) = post_evaluate(Arc::clone(&state_after), &body_fresh).await;
assert_eq!(status_fresh, StatusCode::OK, "{json_fresh}");
assert_eq!(json_fresh["status"], "authorized");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_reserved_hold_blocks_oversubscription() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 100, "USD");
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let body = serde_json::json!({
"capability": cap,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" }
});
let (_, first) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(first["status"], "authorized");
let (_, second) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(
second["status"], "deny",
"the reserved hold must block a second authorization past max_total_cost"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_revoked_capability_is_rejected() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 1000, "USD");
let cap_id = cap.id.clone();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
state
.revoked_capability_ids
.lock()
.await
.insert(cap_id.clone());
let body = serde_json::json!({
"capability": cap,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" }
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::FORBIDDEN);
assert_ne!(json["status"], "authorized");
assert_eq!(json["error"], "chio_capability_revoked");
let usage = budget.get_usage(&cap_id, 0).unwrap();
assert!(usage.is_none() || usage.unwrap().committed_cost_units().unwrap() == 0);
}
fn delegated_child_capability(
issuer: &Keypair,
delegator: &Keypair,
subject: &Keypair,
ancestor_id: &str,
server: &str,
tool: &str,
) -> CapabilityToken {
use chio_core_types::capability::attenuation::{DelegationLink, DelegationLinkBody};
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.test_unwrap()
.as_secs();
let grant = ToolGrant {
server_id: server.to_string(),
tool_name: tool.to_string(),
operations: vec![Operation::Invoke],
constraints: vec![],
max_invocations: None,
max_cost_per_invocation: Some(MonetaryAmount {
units: 100,
currency: "USD".to_string(),
}),
max_total_cost: Some(MonetaryAmount {
units: 1000,
currency: "USD".to_string(),
}),
dpop_required: None,
};
let scope = ChioScope {
grants: vec![grant],
..ChioScope::default()
};
let link = DelegationLink::sign(
DelegationLinkBody {
capability_id: ancestor_id.to_string(),
delegator: delegator.public_key(),
delegatee: subject.public_key(),
attenuations: vec![],
timestamp: now,
scope_hash: None,
aggregate_budget: None,
cumulative_approval: None,
},
delegator,
)
.test_unwrap();
let body = CapabilityTokenBody {
id: uuid::Uuid::now_v7().to_string(),
issuer: issuer.public_key(),
subject: subject.public_key(),
scope,
issued_at: now,
expires_at: now + 3600,
delegation_chain: vec![link],
aggregate_invocation_budget: None,
};
CapabilityToken::sign(body, issuer).test_unwrap()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_revoked_delegation_ancestor_rejects_delegated_child() {
let signer = Keypair::generate();
let root = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let ancestor_id = "cap-root-revoked".to_string();
let child =
delegated_child_capability(&signer, &root, &agent, &ancestor_id, "cost-srv", "compute");
let child_id = child.id.clone();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
state
.revoked_capability_ids
.lock()
.await
.insert(ancestor_id.clone());
assert!(
!state
.revoked_capability_ids
.lock()
.await
.contains(&child_id),
"the presented leaf id must not itself be revoked"
);
let body = serde_json::json!({
"capability": child,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" }
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::FORBIDDEN);
assert_ne!(json["status"], "authorized");
assert_eq!(json["error"], "chio_capability_revoked");
let usage = budget.get_usage(&child_id, 0).unwrap();
assert!(usage.is_none() || usage.unwrap().committed_cost_units().unwrap() == 0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_route_honors_a_durable_only_revocation() {
let signer = Keypair::generate();
let root = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let child =
delegated_child_capability(&signer, &root, &agent, "cap-root", "cost-srv", "compute");
let child_id = child.id.clone();
let durable: Arc<dyn chio_kernel::RevocationStore> =
Arc::new(chio_kernel::InMemoryRevocationStore::new());
chio_kernel::RevocationStore::revoke(durable.as_ref(), &child_id).test_unwrap();
let state = mediated_test_state_core(
signer,
Arc::clone(&budget),
Vec::new(),
Some(MEDIATED_CONTROL_TOKEN.to_string()),
None,
true,
None,
Some(Arc::clone(&durable)),
);
assert!(
!state
.revoked_capability_ids
.lock()
.await
.contains(&child_id),
"the in-memory release set must not carry the durable-only revocation"
);
let body = serde_json::json!({
"capability": child,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" }
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::FORBIDDEN);
assert_eq!(json["error"], "chio_capability_revoked");
let usage = budget.get_usage(&child_id, 0).unwrap();
assert!(usage.is_none() || usage.unwrap().committed_cost_units().unwrap() == 0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_authorization_admits_caller_named_server_id() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap = issue_cost_bearing_capability(
&kernel,
&agent,
"arbitrary-srv",
"invoke",
100,
1000,
"USD",
);
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let body = serde_json::json!({
"capability": cap,
"tool_server": "arbitrary-srv",
"tool_name": "invoke",
"parameters": {}
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(json["status"], "authorized");
assert!(json["execution_nonce"].is_object());
assert_eq!(json["receipt"]["decision"]["verdict"], "incomplete");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_trusts_configured_external_capability_issuers() {
let signer = Keypair::generate();
let external_signer = Keypair::generate();
let agent = Keypair::generate();
let issuer_budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let issuer = issuing_kernel(&external_signer, issuer_budget, &[]);
let cap =
issue_cost_bearing_capability(&issuer, &agent, "cost-srv", "compute", 100, 1000, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let trusting_budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let trusting_state = mediated_test_state(
signer.clone(),
trusting_budget,
vec![external_signer.public_key()],
);
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": {}
});
let (_, json) = post_evaluate(trusting_state, &body).await;
assert_eq!(json["status"], "authorized");
assert!(json["execution_nonce"].is_object());
let untrusting_budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let untrusting_state = mediated_test_state(signer, untrusting_budget, Vec::new());
let (_, json) = post_evaluate(untrusting_state, &body).await;
assert_eq!(json["status"], "deny");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_deny_leaves_committed_cost_zero() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 40, "USD");
let cap_id = cap.id.clone();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let body = serde_json::json!({ "capability": cap, "tool_server": "cost-srv",
"tool_name": "compute", "parameters": {} });
let (_, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(json["status"], "deny");
assert_ne!(json["receipt"]["decision"]["verdict"], "allow");
let usage = budget.get_usage(&cap_id, 0).unwrap();
assert!(usage.is_none() || usage.unwrap().committed_cost_units().unwrap() == 0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_governed_capability_requires_intent_and_approval() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap = issue_governed_capability(&kernel, &agent, "cost-srv", "compute", 100, "USD", 50);
let cap_value = serde_json::to_value(&cap).unwrap();
let approver = signer.clone();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let bare_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" }
});
let (_, denied) = post_evaluate(Arc::clone(&state), &bare_body).await;
assert_eq!(
denied["status"], "deny",
"a governed grant without a forwarded intent must be denied"
);
let request_id = "req-governed-1";
let intent = governed_intent("intent-gov-1", "cost-srv", "compute", 100, "USD");
let approval = governed_approval_token(&approver, &agent.public_key(), &intent, request_id);
let governed_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" },
"request_id": request_id,
"governed_intent": intent,
"approval_token": approval
});
let (_, authorized) = post_evaluate(Arc::clone(&state), &governed_body).await;
assert_eq!(
authorized["status"], "authorized",
"a governed grant with a valid intent and approval must be authorized"
);
assert!(authorized["execution_nonce"].is_object());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_governed_mustprepay_authorizes_with_payment_adapter() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap = issue_governed_capability(&kernel, &agent, "cost-srv", "compute", 100, "USD", 50);
let cap_value = serde_json::to_value(&cap).unwrap();
let approver = signer.clone();
let state = mediated_test_state_core(
signer,
Arc::clone(&budget),
Vec::new(),
Some(MEDIATED_CONTROL_TOKEN.to_string()),
None,
true,
Some(Box::new(chio_kernel::payment::SimPaymentAdapter::new())),
None,
);
let request_id = "req-mustprepay-adapter";
let intent =
governed_mustprepay_intent("intent-prepay-1", "cost-srv", "compute", 100, "USD");
let approval = governed_approval_token(&approver, &agent.public_key(), &intent, request_id);
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" },
"request_id": request_id,
"governed_intent": intent,
"approval_token": approval,
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(
json["status"], "authorized",
"a configured payment adapter must let an approved governed MustPrepay authorize"
);
assert!(
json["execution_nonce"].is_object(),
"an authorized MustPrepay reservation must mint an execution nonce"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_governed_mustprepay_denied_without_payment_adapter() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap = issue_governed_capability(&kernel, &agent, "cost-srv", "compute", 100, "USD", 50);
let cap_id = cap.id.clone();
let cap_value = serde_json::to_value(&cap).unwrap();
let approver = signer.clone();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let request_id = "req-mustprepay-no-adapter";
let intent =
governed_mustprepay_intent("intent-prepay-2", "cost-srv", "compute", 100, "USD");
let approval = governed_approval_token(&approver, &agent.public_key(), &intent, request_id);
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" },
"request_id": request_id,
"governed_intent": intent,
"approval_token": approval,
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(
json["status"], "deny",
"governed MustPrepay must deny fail-closed without a configured payment adapter"
);
assert!(
json["execution_nonce"].is_null(),
"a denied MustPrepay must not mint a reserved nonce"
);
let usage = budget.get_usage(&cap_id, 0).unwrap();
assert!(usage.is_none() || usage.unwrap().committed_cost_units().unwrap() == 0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_dpop_capability_requires_valid_proof() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap = issue_dpop_capability(&kernel, &agent, "cost-srv", "compute", 100, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let params = serde_json::json!({ "invoice": "inv-1" });
let bare_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params.clone()
});
let (_, denied) = post_evaluate(Arc::clone(&state), &bare_body).await;
assert_eq!(
denied["status"], "deny",
"a dpop_required grant without a proof must be denied"
);
let proof = dpop_proof_for(&agent, &cap, "cost-srv", "compute", ¶ms);
let dpop_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"dpop_proof": proof
});
let (_, authorized) = post_evaluate(Arc::clone(&state), &dpop_body).await;
assert_eq!(
authorized["status"], "authorized",
"a valid DPoP proof must authorize the dpop_required grant"
);
assert!(authorized["execution_nonce"].is_object());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_dpop_proof_replay_is_rejected_across_requests() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap = issue_dpop_capability_with_total(
&kernel, &agent, "cost-srv", "compute", 100, 1000, "USD",
);
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let params = serde_json::json!({ "invoice": "inv-1" });
let proof = dpop_proof_for(&agent, &cap, "cost-srv", "compute", ¶ms);
let first_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "dpop-req-1",
"dpop_proof": proof,
});
let (_, first) = post_evaluate(Arc::clone(&state), &first_body).await;
assert_eq!(
first["status"], "authorized",
"the first presentation of a valid DPoP proof must authorize"
);
let second_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "dpop-req-2",
"dpop_proof": proof,
});
let (_, second) = post_evaluate(Arc::clone(&state), &second_body).await;
assert_eq!(
second["status"], "deny",
"replaying the DPoP proof must be rejected by the shared kernel's persistent nonce store"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_reused_request_id_is_conflict() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 1000, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" },
"request_id": "fixed-req-id",
});
let (status, first) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(first["status"], "authorized");
let (status, second) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::CONFLICT);
assert_ne!(second["status"], "authorized");
assert_eq!(second["error"], "chio_request_id_reused");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_authorization_requires_hold_capable_budget_store() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 1000, "USD");
let cap_id = cap.id.clone();
let state = mediated_test_state_inner(
signer,
Arc::clone(&budget),
Vec::new(),
Some(MEDIATED_CONTROL_TOKEN.to_string()),
None,
false,
);
let body = serde_json::json!({
"capability": cap,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" }
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR);
assert_eq!(json["error"], "chio_mediation_requires_local_budget_store");
assert_ne!(json["status"], "authorized");
assert!(
json["execution_nonce"].is_null(),
"a fail-closed rejection must not mint a reserved nonce"
);
let usage = budget.get_usage(&cap_id, 0).unwrap();
assert!(usage.is_none() || usage.unwrap().committed_cost_units().unwrap() == 0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_authorization_requires_reconcile_control_token() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 1000, "USD");
let cap_id = cap.id.clone();
let state =
mediated_test_state_with_control_token(signer, Arc::clone(&budget), Vec::new(), None);
let body = serde_json::json!({
"capability": cap,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" }
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR);
assert_eq!(json["error"], "chio_mediation_requires_reconcile_token");
assert_ne!(json["status"], "authorized");
assert!(
json["execution_nonce"].is_null(),
"a fail-closed rejection must not mint a reserved nonce"
);
let usage = budget.get_usage(&cap_id, 0).unwrap();
assert!(usage.is_none() || usage.unwrap().committed_cost_units().unwrap() == 0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_authorization_rejects_blank_reconcile_control_token() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 1000, "USD");
let cap_id = cap.id.clone();
let state = mediated_test_state_with_control_token(
signer,
Arc::clone(&budget),
Vec::new(),
Some(" ".to_string()),
);
let body = serde_json::json!({
"capability": cap,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" }
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR);
assert_eq!(json["error"], "chio_mediation_requires_reconcile_token");
assert_ne!(json["status"], "authorized");
assert!(
json["execution_nonce"].is_null(),
"a fail-closed rejection must not mint a reserved nonce"
);
let usage = budget.get_usage(&cap_id, 0).unwrap();
assert!(usage.is_none() || usage.unwrap().committed_cost_units().unwrap() == 0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn reconcile_requires_hold_capable_budget_store() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 150, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let params = serde_json::json!({ "invoice": "inv-1" });
let hold_capable = mediated_test_state(signer.clone(), Arc::clone(&budget), Vec::new());
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "recon-remote",
});
let (_, authorized) = post_evaluate(hold_capable, &body).await;
let nonce_json = authorized["execution_nonce"].clone();
assert!(nonce_json.is_object());
let remote = mediated_test_state_inner(
signer,
budget,
Vec::new(),
Some(MEDIATED_CONTROL_TOKEN.to_string()),
None,
false,
);
let reconcile_body = serde_json::json!({
"execution_nonce": nonce_json,
"arguments": params,
"realized_cost": { "units": 30, "currency": "USD" },
});
let (status, json) = post_reconcile(remote, &reconcile_body).await;
assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR);
assert_eq!(json["error"], "chio_mediation_requires_local_budget_store");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_durable_hold_rejects_request_id_reuse_across_restart() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 1000, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let before = mediated_test_state(signer.clone(), Arc::clone(&budget), Vec::new());
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" },
"request_id": "restart-req",
});
let (status, authorized) = post_evaluate(Arc::clone(&before), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(authorized["status"], "authorized");
let after = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let (status, replay) = post_evaluate(Arc::clone(&after), &body).await;
assert_eq!(status, StatusCode::CONFLICT);
assert_ne!(replay["status"], "authorized");
assert!(
replay["execution_nonce"].is_null(),
"no second nonce is minted against the open hold"
);
assert_eq!(replay["error"], "chio_request_id_reused");
let fresh_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-2" },
"request_id": "restart-req-fresh",
});
let (status, fresh) = post_evaluate(after, &fresh_body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(fresh["status"], "authorized");
assert!(fresh["execution_nonce"].is_object());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn reconcile_settles_reserved_hold_and_frees_budget() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let signer_pub = signer.public_key();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 150, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let params = serde_json::json!({ "invoice": "inv-1" });
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "recon-reserve",
});
let (_, authorized) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(authorized["status"], "authorized");
let nonce_json = authorized["execution_nonce"].clone();
assert!(nonce_json.is_object());
let blocked_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "recon-blocked",
});
let (_, blocked) = post_evaluate(Arc::clone(&state), &blocked_body).await;
assert_eq!(blocked["status"], "deny");
let reconcile_body = serde_json::json!({
"execution_nonce": nonce_json.clone(),
"arguments": params,
"realized_cost": { "units": 30, "currency": "USD" },
});
let (status, reconciled) = post_reconcile(Arc::clone(&state), &reconcile_body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(reconciled["status"], "reconciled");
let receipt: ChioReceipt = serde_json::from_value(reconciled["receipt"].clone()).unwrap();
let nonce: SignedExecutionNonce = serde_json::from_value(nonce_json).unwrap();
assert_eq!(
is_authoritative_spend_receipt(&receipt, &[signer_pub], &nonce),
Ok(())
);
let after_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "recon-after",
});
let (_, after) = post_evaluate(Arc::clone(&state), &after_body).await;
assert_eq!(
after["status"], "authorized",
"the budget freed by reconcile must admit a new authorization"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn reconcile_rejects_replayed_nonce() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 150, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let params = serde_json::json!({ "invoice": "inv-1" });
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "recon-replay",
});
let (_, authorized) = post_evaluate(Arc::clone(&state), &body).await;
let nonce_json = authorized["execution_nonce"].clone();
let reconcile_body = serde_json::json!({
"execution_nonce": nonce_json,
"arguments": params,
"realized_cost": { "units": 30, "currency": "USD" },
});
let (status, _) = post_reconcile(Arc::clone(&state), &reconcile_body).await;
assert_eq!(status, StatusCode::OK);
let (status, replay) = post_reconcile(Arc::clone(&state), &reconcile_body).await;
assert_eq!(status, StatusCode::BAD_REQUEST);
assert_eq!(replay["error"], "chio_reconcile_rejected");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn reconcile_rejects_argument_mismatch() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 150, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let params = serde_json::json!({ "invoice": "inv-1" });
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "recon-mismatch",
});
let (_, authorized) = post_evaluate(Arc::clone(&state), &body).await;
let nonce_json = authorized["execution_nonce"].clone();
let reconcile_body = serde_json::json!({
"execution_nonce": nonce_json,
"arguments": { "invoice": "tampered" },
"realized_cost": { "units": 30, "currency": "USD" },
});
let (status, rejected) = post_reconcile(Arc::clone(&state), &reconcile_body).await;
assert_eq!(status, StatusCode::BAD_REQUEST);
assert_eq!(rejected["error"], "chio_reconcile_rejected");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn reaper_forfeits_expired_hold_at_worst_case() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 100, "USD");
let cap_id = cap.id.clone();
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let params = serde_json::json!({ "invoice": "inv-1" });
let reserve_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "reap-reserve",
});
let (_, authorized) = post_evaluate(Arc::clone(&state), &reserve_body).await;
assert_eq!(authorized["status"], "authorized");
let blocked_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "reap-blocked",
});
let (_, blocked) = post_evaluate(Arc::clone(&state), &blocked_body).await;
assert_eq!(blocked["status"], "deny");
let settled = reap_expired_reserved_holds_once(&state, i64::MAX)
.await
.unwrap();
assert_eq!(settled, 1, "the expired reserved hold must be settled");
let after_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "reap-after",
});
let (_, after) = post_evaluate(Arc::clone(&state), &after_body).await;
assert_eq!(
after["status"], "deny",
"a forfeited reserved hold must keep the grant committed, not free it"
);
let usage = budget.get_usage(&cap_id, 0).unwrap();
let usage = usage.expect("the forfeited hold must remain recorded in the budget store");
assert_eq!(
usage.committed_cost_units().unwrap(),
100,
"reaping settles the reserved worst-case as realized spend (fail-closed)"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_authorization_works_with_both_control_url_and_budget_db() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let dir =
std::env::temp_dir().join(format!("chio-budget-both-e2e-{}", uuid::Uuid::now_v7()));
std::fs::create_dir_all(&dir).unwrap();
let db = dir.join("budget.sqlite");
let config = ProtectConfig {
upstream: "http://127.0.0.1:1".to_string(),
spec_content: Some("{}".to_string()),
spec_path: None,
listen_addr: "127.0.0.1:0".to_string(),
receipt_db: None,
allow_ephemeral_receipts: true,
sidecar_control_token: None,
signer_seed_hex: None,
trusted_capability_issuers: Vec::new(),
control_url: Some("http://127.0.0.1:1".to_string()),
control_token: Some("token".to_string()),
budget_db: Some(db.to_string_lossy().to_string()),
revocation_db: None,
require_nonce: false,
allow_advisory: false,
upstream_request_timeout: crate::DEFAULT_UPSTREAM_REQUEST_TIMEOUT,
};
let configured = build_budget_store(&config)
.unwrap()
.expect("a budget store must be built");
assert!(configured.hold_capable);
let budget = configured.store;
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 150, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state_inner(
signer,
Arc::clone(&budget),
Vec::new(),
Some(MEDIATED_CONTROL_TOKEN.to_string()),
None,
configured.hold_capable,
);
let params = serde_json::json!({ "invoice": "inv-both" });
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "both-configured",
});
let (status, authorized) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(
authorized["status"], "authorized",
"both configured must authorize via the local hold-capable store"
);
let nonce_json = authorized["execution_nonce"].clone();
assert!(nonce_json.is_object());
let reconcile_body = serde_json::json!({
"execution_nonce": nonce_json,
"arguments": params,
"realized_cost": { "units": 30, "currency": "USD" },
});
let (status, reconciled) = post_reconcile(state, &reconcile_body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(reconciled["status"], "reconciled");
}
fn open_temp_receipt_store() -> (std::path::PathBuf, SqliteReceiptStore) {
let dir = std::env::temp_dir().join(format!("chio-receipt-{}", uuid::Uuid::now_v7()));
std::fs::create_dir_all(&dir).unwrap();
let db = dir.join("receipts.sqlite");
let store = SqliteReceiptStore::open(&db.to_string_lossy()).unwrap();
(db, store)
}
fn failing_receipt_store() -> SqliteReceiptStore {
let (db, store) = open_temp_receipt_store();
let dropper = rusqlite::Connection::open(&db).unwrap();
dropper.execute("DROP TABLE tool_receipts", []).unwrap();
drop(dropper);
store
}
#[derive(Clone, Default)]
struct RecordingPaymentAdapter {
inner: chio_kernel::payment::SimPaymentAdapter,
captures: Arc<std::sync::atomic::AtomicUsize>,
releases: Arc<std::sync::atomic::AtomicUsize>,
refunds: Arc<std::sync::atomic::AtomicUsize>,
}
impl chio_kernel::PaymentAdapter for RecordingPaymentAdapter {
fn authorize(
&self,
request: &chio_kernel::PaymentAuthorizeRequest,
) -> Result<chio_kernel::PaymentAuthorization, chio_kernel::PaymentError> {
let authorization = self.inner.authorize(request)?;
if authorization.state.is_final() {
self.captures
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
}
Ok(authorization)
}
fn capture(
&self,
authorization_id: &str,
amount_units: u64,
currency: &str,
reference: &str,
) -> Result<chio_kernel::PaymentResult, chio_kernel::PaymentError> {
self.captures
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
self.inner
.capture(authorization_id, amount_units, currency, reference)
}
fn release(
&self,
authorization_id: &str,
reference: &str,
) -> Result<chio_kernel::PaymentResult, chio_kernel::PaymentError> {
self.releases
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
self.inner.release(authorization_id, reference)
}
fn refund(
&self,
transaction_id: &str,
amount_units: u64,
currency: &str,
reference: &str,
) -> Result<chio_kernel::PaymentResult, chio_kernel::PaymentError> {
self.refunds
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
self.inner
.refund(transaction_id, amount_units, currency, reference)
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_receipt_persistence_failure_returns_nonce_and_keeps_reservation() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 100, "USD");
let cap_id = cap.id.clone();
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state_inner(
signer,
Arc::clone(&budget),
Vec::new(),
Some(MEDIATED_CONTROL_TOKEN.to_string()),
Some(failing_receipt_store()),
true,
);
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" },
"request_id": "persist-fail",
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(json["status"], "authorized");
assert!(
json["execution_nonce"].is_object(),
"a persistence failure after a successful reserve must still return the nonce"
);
let hold_id = json["execution_nonce"]["nonce"]["reserved_hold_id"]
.as_str()
.expect("the returned nonce must name its reserved hold");
let hold = budget.get_budget_hold(hold_id).unwrap();
assert!(
hold.map(|hold| hold.disposition.is_open()).unwrap_or(false),
"a returned reservation must keep its reserved hold open"
);
let usage = budget.get_usage(&cap_id, 0).unwrap();
let usage = usage.expect("the reserved hold must remain recorded in the budget store");
assert_eq!(
usage.committed_cost_units().unwrap(),
100,
"the returned reservation must keep its reserved budget committed"
);
assert_eq!(
state.minted_request_ids.lock().await.len(),
1,
"a returned reservation must retain its request-id claim"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_receipt_persistence_success_keeps_reservation() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 100, "USD");
let cap_id = cap.id.clone();
let cap_value = serde_json::to_value(&cap).unwrap();
let (_db, receipt_store) = open_temp_receipt_store();
let state = mediated_test_state_inner(
signer,
Arc::clone(&budget),
Vec::new(),
Some(MEDIATED_CONTROL_TOKEN.to_string()),
Some(receipt_store),
true,
);
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" },
"request_id": "persist-ok",
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(json["status"], "authorized");
assert!(json["execution_nonce"].is_object());
let usage = budget.get_usage(&cap_id, 0).unwrap();
let usage = usage.expect("the reserved hold must remain recorded in the budget store");
assert_eq!(
usage.committed_cost_units().unwrap(),
100,
"a persisted authorization must keep its reserved hold"
);
assert_eq!(
state.minted_request_ids.lock().await.len(),
1,
"a persisted authorization must retain its request-id claim"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_invocation_receipt_persistence_failure_returns_nonce_and_keeps_reservation() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap = issue_invocation_capability(&kernel, &agent, "cost-srv", "compute", 1);
let cap_id = cap.id.clone();
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state_inner(
signer,
Arc::clone(&budget),
Vec::new(),
Some(MEDIATED_CONTROL_TOKEN.to_string()),
Some(failing_receipt_store()),
true,
);
let params = serde_json::json!({ "invoice": "inv-1" });
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "invoke-persist-fail",
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(json["status"], "authorized");
let nonce_json = json["execution_nonce"].clone();
assert!(
nonce_json.is_object(),
"an invocation reserve whose receipt fails to persist must still return the nonce"
);
let usage = budget.get_usage(&cap_id, 0).unwrap();
assert_eq!(
usage.map(|usage| usage.invocation_count).unwrap_or(0),
1,
"the returned invocation reservation stays consumed against the grant"
);
let hold_id = nonce_json["nonce"]["reserved_hold_id"]
.as_str()
.expect("the returned nonce must name its reserved hold");
let hold = budget.get_budget_hold(hold_id).unwrap();
assert!(
hold.map(|hold| hold.disposition.is_open()).unwrap_or(false),
"the returned invocation reservation must keep its reserved hold open"
);
let reconcile_body = serde_json::json!({
"execution_nonce": nonce_json,
"arguments": params,
"realized_cost": { "units": 0, "currency": "USD" },
});
let (recon_status, reconciled) = post_reconcile(state, &reconcile_body).await;
assert_eq!(recon_status, StatusCode::OK);
assert_eq!(reconciled["status"], "reconciled");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_mustprepay_receipt_persistence_failure_returns_nonce_without_refund() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap = issue_governed_capability(&kernel, &agent, "cost-srv", "compute", 100, "USD", 50);
let cap_value = serde_json::to_value(&cap).unwrap();
let approver = signer.clone();
let captures = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let refunds = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let adapter = RecordingPaymentAdapter {
inner: chio_kernel::payment::SimPaymentAdapter::new(),
captures: Arc::clone(&captures),
releases: Arc::new(std::sync::atomic::AtomicUsize::new(0)),
refunds: Arc::clone(&refunds),
};
let state = mediated_test_state_core(
signer,
Arc::clone(&budget),
Vec::new(),
Some(MEDIATED_CONTROL_TOKEN.to_string()),
Some(failing_receipt_store()),
true,
Some(Box::new(adapter)),
None,
);
let request_id = "req-mustprepay-persist-fail";
let intent =
governed_mustprepay_intent("intent-prepay-persist", "cost-srv", "compute", 100, "USD");
let approval = governed_approval_token(&approver, &agent.public_key(), &intent, request_id);
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-1" },
"request_id": request_id,
"governed_intent": intent,
"approval_token": approval,
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(
status,
StatusCode::OK,
"a captured MustPrepay reserve whose receipt fails to persist must not 500"
);
assert_eq!(json["status"], "authorized");
assert!(
json["execution_nonce"].is_object(),
"the caller must receive the nonce the captured prepayment backs"
);
assert_eq!(
captures.load(std::sync::atomic::Ordering::SeqCst),
1,
"the MustPrepay quote must be captured exactly once"
);
assert_eq!(
refunds.load(std::sync::atomic::Ordering::SeqCst),
0,
"the captured prepayment must not be refunded: it backs the returned nonce"
);
let hold_id = json["execution_nonce"]["nonce"]["reserved_hold_id"]
.as_str()
.expect("the returned nonce must name its reserved hold");
let hold = budget.get_budget_hold(hold_id).unwrap();
assert!(
hold.map(|hold| hold.disposition.is_open()).unwrap_or(false),
"the captured MustPrepay reservation must stay open, backing the returned nonce"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_invocation_reconcile_keeps_invocation_consumed() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap = issue_invocation_capability(&kernel, &agent, "cost-srv", "compute", 1);
let cap_id = cap.id.clone();
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let params = serde_json::json!({ "invoice": "inv-1" });
let reserve_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "invoke-reconcile",
});
let (_, authorized) = post_evaluate(Arc::clone(&state), &reserve_body).await;
assert_eq!(authorized["status"], "authorized");
let nonce_json = authorized["execution_nonce"].clone();
let reconcile_body = serde_json::json!({
"execution_nonce": nonce_json,
"arguments": params,
"realized_cost": { "units": 0, "currency": "USD" },
});
let (status, reconciled) = post_reconcile(Arc::clone(&state), &reconcile_body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(reconciled["status"], "reconciled");
let usage = budget.get_usage(&cap_id, 0).unwrap();
assert_eq!(
usage.map(|usage| usage.invocation_count).unwrap_or(0),
1,
"a legitimate reconcile must keep the invocation consumed, not refund it"
);
let after_body = serde_json::json!({
"capability": serde_json::to_value(&cap).unwrap(),
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "invoke-reconcile-after",
});
let (_, after) = post_evaluate(state, &after_body).await;
assert_eq!(
after["status"], "deny",
"a reconciled invocation stays consumed against max_invocations"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn reaper_forfeits_expired_invocation_reserve() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap = issue_invocation_capability(&kernel, &agent, "cost-srv", "compute", 1);
let cap_id = cap.id.clone();
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let params = serde_json::json!({ "invoice": "inv-1" });
let reserve_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "invoke-reap",
});
let (_, authorized) = post_evaluate(Arc::clone(&state), &reserve_body).await;
assert_eq!(authorized["status"], "authorized");
let settled = reap_expired_reserved_holds_once(&state, i64::MAX)
.await
.unwrap();
assert_eq!(
settled, 1,
"the expired invocation reservation must be settled"
);
let usage = budget.get_usage(&cap_id, 0).unwrap();
assert_eq!(
usage.map(|usage| usage.invocation_count).unwrap_or(0),
1,
"reaping an abandoned invocation reservation forfeits it (stays consumed)"
);
let after_body = serde_json::json!({
"capability": serde_json::to_value(&cap).unwrap(),
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "invoke-reap-after",
});
let (_, after) = post_evaluate(state, &after_body).await;
assert_eq!(
after["status"], "deny",
"a forfeited invocation reservation must keep the grant committed"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_open_invocation_reserve_blocks_oversubscription() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap = issue_invocation_capability(&kernel, &agent, "cost-srv", "compute", 1);
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let params = serde_json::json!({ "invoice": "inv-1" });
let first_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "invoke-open-1",
});
let (_, first) = post_evaluate(Arc::clone(&state), &first_body).await;
assert_eq!(first["status"], "authorized");
let second_body = serde_json::json!({
"capability": serde_json::to_value(&cap).unwrap(),
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "invoke-open-2",
});
let (_, second) = post_evaluate(state, &second_body).await;
assert_eq!(
second["status"], "deny",
"an open invocation reservation must block a second reserve past max_invocations"
);
}
#[test]
fn build_budget_store_local_sqlite_when_no_control_url() {
let dir = std::env::temp_dir().join(format!("chio-budget-{}", uuid::Uuid::now_v7()));
std::fs::create_dir_all(&dir).unwrap();
let db = dir.join("budget.sqlite");
let config = ProtectConfig {
upstream: "http://127.0.0.1:1".to_string(),
spec_content: Some("{}".to_string()),
spec_path: None,
listen_addr: "127.0.0.1:0".to_string(),
receipt_db: None,
allow_ephemeral_receipts: true,
sidecar_control_token: None,
signer_seed_hex: None,
trusted_capability_issuers: Vec::new(),
control_url: None,
control_token: None,
budget_db: Some(db.to_string_lossy().to_string()),
revocation_db: None,
require_nonce: false,
allow_advisory: false,
upstream_request_timeout: crate::DEFAULT_UPSTREAM_REQUEST_TIMEOUT,
};
let configured = build_budget_store(&config).unwrap();
let configured = configured.expect("local sqlite budget store must be built");
assert!(
configured.hold_capable,
"the local sqlite budget store implements the hold APIs and must be hold-capable"
);
}
#[test]
fn build_budget_store_remote_is_not_hold_capable() {
let config = ProtectConfig {
upstream: "http://127.0.0.1:1".to_string(),
spec_content: Some("{}".to_string()),
spec_path: None,
listen_addr: "127.0.0.1:0".to_string(),
receipt_db: None,
allow_ephemeral_receipts: true,
sidecar_control_token: None,
signer_seed_hex: None,
trusted_capability_issuers: Vec::new(),
control_url: Some("http://127.0.0.1:1".to_string()),
control_token: Some("token".to_string()),
budget_db: None,
revocation_db: None,
require_nonce: false,
allow_advisory: false,
upstream_request_timeout: crate::DEFAULT_UPSTREAM_REQUEST_TIMEOUT,
};
let configured = build_budget_store(&config).unwrap();
let configured = configured.expect("remote budget store must be built");
assert!(
!configured.hold_capable,
"the remote control-plane budget store does not implement the hold APIs"
);
}
#[test]
fn build_budget_store_prefers_local_hold_capable_when_both_configured() {
let dir = std::env::temp_dir().join(format!("chio-budget-both-{}", uuid::Uuid::now_v7()));
std::fs::create_dir_all(&dir).unwrap();
let db = dir.join("budget.sqlite");
let config = ProtectConfig {
upstream: "http://127.0.0.1:1".to_string(),
spec_content: Some("{}".to_string()),
spec_path: None,
listen_addr: "127.0.0.1:0".to_string(),
receipt_db: None,
allow_ephemeral_receipts: true,
sidecar_control_token: None,
signer_seed_hex: None,
trusted_capability_issuers: Vec::new(),
control_url: Some("http://127.0.0.1:1".to_string()),
control_token: Some("token".to_string()),
budget_db: Some(db.to_string_lossy().to_string()),
revocation_db: None,
require_nonce: false,
allow_advisory: false,
upstream_request_timeout: crate::DEFAULT_UPSTREAM_REQUEST_TIMEOUT,
};
let configured = build_budget_store(&config).unwrap();
let configured = configured.expect("a budget store must be built when both are configured");
assert!(
configured.hold_capable,
"with both configured, mediation must use the hold-capable local store"
);
assert_eq!(
configured.store.count_open_holds().unwrap(),
0,
"the preferred store must be the functional local hold-capable store"
);
}
fn revocation_db_config(revocation_db: Option<String>) -> ProtectConfig {
ProtectConfig {
upstream: "http://127.0.0.1:1".to_string(),
spec_content: Some("{}".to_string()),
spec_path: None,
listen_addr: "127.0.0.1:0".to_string(),
receipt_db: None,
allow_ephemeral_receipts: true,
sidecar_control_token: None,
signer_seed_hex: None,
trusted_capability_issuers: Vec::new(),
control_url: None,
control_token: None,
budget_db: None,
revocation_db,
require_nonce: false,
allow_advisory: false,
upstream_request_timeout: crate::DEFAULT_UPSTREAM_REQUEST_TIMEOUT,
}
}
#[test]
fn load_revocation_db_ids_is_empty_without_configured_store() {
let ids = load_revocation_db_ids(&revocation_db_config(None)).unwrap();
assert!(ids.is_empty());
}
#[test]
fn load_revocation_db_ids_reads_operator_revocations() {
let dir = std::env::temp_dir().join(format!("chio-revocation-{}", uuid::Uuid::now_v7()));
std::fs::create_dir_all(&dir).unwrap();
let db = dir.join("revocations.sqlite3");
let db_path = db.to_string_lossy().to_string();
let store = chio_store_sqlite::SqliteRevocationStore::open(&db).unwrap();
assert!(chio_kernel::RevocationStore::revoke(&store, "cap-operator-revoked").unwrap());
drop(store);
let ids = load_revocation_db_ids(&revocation_db_config(Some(db_path))).unwrap();
assert!(
ids.contains("cap-operator-revoked"),
"durable operator revocation must be loaded into the enforced set"
);
}
#[test]
fn load_revocation_db_ids_fails_closed_on_unreadable_store() {
let dir =
std::env::temp_dir().join(format!("chio-revocation-bad-{}", uuid::Uuid::now_v7()));
std::fs::create_dir_all(&dir).unwrap();
let db = dir.join("not-a-db.sqlite3");
std::fs::write(&db, b"this is not a sqlite database").unwrap();
let result = load_revocation_db_ids(&revocation_db_config(Some(
db.to_string_lossy().to_string(),
)));
assert!(
result.is_err(),
"an unreadable revocation-db must fail closed rather than start with no revocations"
);
}
#[test]
fn mediation_kernel_installs_budget_store_and_strict_nonce_config() {
let signer = Keypair::generate();
let budget: Arc<dyn BudgetStore> =
Arc::new(chio_kernel::budget_store::InMemoryBudgetStore::new());
let kernel =
build_mediation_kernel(&signer, Arc::clone(&budget), &[], Vec::new(), None, None)
.unwrap();
assert!(
kernel.execution_nonce_required(),
"mediation kernel must always run execution-nonce strict mode"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn reconcile_requires_sidecar_control_token() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 150, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let params = serde_json::json!({ "invoice": "inv-1" });
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "recon-auth-gate",
});
let (_, authorized) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(authorized["status"], "authorized");
let nonce_json = authorized["execution_nonce"].clone();
let reconcile_body = serde_json::json!({
"execution_nonce": nonce_json,
"arguments": params,
"realized_cost": { "units": 30, "currency": "USD" },
});
let (status, denied) =
post_json(Arc::clone(&state), "/v1/reconcile", &reconcile_body).await;
assert_eq!(status, StatusCode::FORBIDDEN);
assert_eq!(denied["error"], "chio_control_forbidden");
let (status, reconciled) = post_reconcile(Arc::clone(&state), &reconcile_body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(reconciled["status"], "reconciled");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn reconcile_without_configured_control_token_is_rejected() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 150, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let params = serde_json::json!({ "invoice": "inv-1" });
let with_token = mediated_test_state(signer.clone(), Arc::clone(&budget), Vec::new());
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "recon-no-token",
});
let (_, authorized) = post_evaluate(with_token, &body).await;
let nonce_json = authorized["execution_nonce"].clone();
assert!(nonce_json.is_object());
let state = mediated_test_state_with_control_token(signer, budget, Vec::new(), None);
let reconcile_body = serde_json::json!({
"execution_nonce": nonce_json,
"arguments": params,
"realized_cost": { "units": 0, "currency": "USD" },
});
let (status, denied) = post_reconcile(state, &reconcile_body).await;
assert_eq!(status, StatusCode::FORBIDDEN);
assert_eq!(denied["error"], "chio_control_forbidden");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_authorization_needs_no_tool_server_registration() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
for index in 0..8 {
let server = format!("arbitrary-srv-{index}");
let cap =
issue_cost_bearing_capability(&kernel, &agent, &server, "invoke", 100, 1000, "USD");
let body = serde_json::json!({
"capability": cap,
"tool_server": server,
"tool_name": "invoke",
"parameters": {},
"request_id": format!("noreg-{index}"),
});
let (status, json) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(
json["status"], "authorized",
"an arbitrary caller server id must authorize without any registration"
);
assert!(json["execution_nonce"].is_object());
}
}
#[test]
fn minted_request_id_window_bounds_reuse_and_expiry() {
let mut window = MintedRequestIdWindow::new(30);
assert!(window.claim("req-a", 1_000));
assert!(!window.claim("req-a", 1_000));
assert_eq!(window.len(), 1);
window.release("req-a");
assert_eq!(window.len(), 0);
assert!(window.claim("req-a", 1_000));
assert!(window.claim("req-b", 1_010));
assert_eq!(window.len(), 2);
assert!(window.claim("req-c", 1_031));
assert_eq!(window.len(), 2);
assert!(window.claim("req-a", 1_031));
assert_eq!(window.len(), 3);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_denied_authorization_does_not_burn_request_id() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 40, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": {},
"request_id": "denied-then-retry",
});
let (status, first) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(first["status"], "deny");
let (status, second) = post_evaluate(Arc::clone(&state), &body).await;
assert_ne!(
status,
StatusCode::CONFLICT,
"a denied authorization must release its request-id claim"
);
assert_eq!(second["status"], "deny");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn reserved_hold_reaper_handle_is_retained_not_detached() {
let signer = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let state = mediated_test_state(signer, budget, Vec::new());
assert!(state.reaper_handle.lock().await.is_none());
spawn_reserved_hold_reaper(&state).await;
{
let guard = state.reaper_handle.lock().await;
let handle = guard
.as_ref()
.expect("the reaper handle must be retained on the state");
assert!(
!handle.is_finished(),
"the retained reaper handle must reference a live, abortable task"
);
}
let handle = state.reaper_handle.lock().await.take().unwrap();
handle.abort();
assert!(
handle.await.is_err(),
"aborting the retained handle must cancel the reaper task"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn reconcile_returns_authoritative_receipt_when_persistence_fails() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let signer_pub = signer.public_key();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 150, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let (db, receipt_store) = open_temp_receipt_store();
let state = mediated_test_state_inner(
signer,
Arc::clone(&budget),
Vec::new(),
Some(MEDIATED_CONTROL_TOKEN.to_string()),
Some(receipt_store),
true,
);
let params = serde_json::json!({ "invoice": "inv-1" });
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "recon-persist-fail",
});
let (status, authorized) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(authorized["status"], "authorized");
let nonce_json = authorized["execution_nonce"].clone();
assert!(nonce_json.is_object());
let dropper = rusqlite::Connection::open(&db).unwrap();
dropper.execute("DROP TABLE tool_receipts", []).unwrap();
drop(dropper);
let reconcile_body = serde_json::json!({
"execution_nonce": nonce_json.clone(),
"arguments": params,
"realized_cost": { "units": 30, "currency": "USD" },
});
let (status, reconciled) = post_reconcile(Arc::clone(&state), &reconcile_body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(reconciled["status"], "reconciled");
let receipt: ChioReceipt = serde_json::from_value(reconciled["receipt"].clone()).unwrap();
let nonce: SignedExecutionNonce = serde_json::from_value(nonce_json).unwrap();
assert_eq!(
is_authoritative_spend_receipt(&receipt, &[signer_pub], &nonce),
Ok(()),
"a persistence failure after settlement must still return the authoritative receipt"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn reconcile_still_fails_closed_on_replayed_nonce_when_persistence_fails() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 150, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let (db, receipt_store) = open_temp_receipt_store();
let state = mediated_test_state_inner(
signer,
Arc::clone(&budget),
Vec::new(),
Some(MEDIATED_CONTROL_TOKEN.to_string()),
Some(receipt_store),
true,
);
let params = serde_json::json!({ "invoice": "inv-1" });
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "recon-replay-persist-fail",
});
let (_, authorized) = post_evaluate(Arc::clone(&state), &body).await;
let nonce_json = authorized["execution_nonce"].clone();
let reconcile_body = serde_json::json!({
"execution_nonce": nonce_json,
"arguments": params,
"realized_cost": { "units": 30, "currency": "USD" },
});
let (status, _) = post_reconcile(Arc::clone(&state), &reconcile_body).await;
assert_eq!(status, StatusCode::OK);
let dropper = rusqlite::Connection::open(&db).unwrap();
dropper.execute("DROP TABLE tool_receipts", []).unwrap();
drop(dropper);
let (status, replay) = post_reconcile(Arc::clone(&state), &reconcile_body).await;
assert_eq!(status, StatusCode::BAD_REQUEST);
assert_eq!(replay["error"], "chio_reconcile_rejected");
assert!(replay.get("receipt").is_none());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn mediated_durable_hold_rejects_request_id_reuse_after_settle() {
let signer = Keypair::generate();
let agent = Keypair::generate();
let budget: Arc<dyn BudgetStore> = Arc::new(InMemoryBudgetStore::new());
let kernel = issuing_kernel(&signer, Arc::clone(&budget), &[]);
let cap =
issue_cost_bearing_capability(&kernel, &agent, "cost-srv", "compute", 100, 1000, "USD");
let cap_value = serde_json::to_value(&cap).unwrap();
let state = mediated_test_state(signer.clone(), Arc::clone(&budget), Vec::new());
let params = serde_json::json!({ "invoice": "inv-1" });
let body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": params,
"request_id": "settled-reuse",
});
let (_, authorized) = post_evaluate(Arc::clone(&state), &body).await;
assert_eq!(authorized["status"], "authorized");
let nonce_json = authorized["execution_nonce"].clone();
let reconcile_body = serde_json::json!({
"execution_nonce": nonce_json,
"arguments": params,
"realized_cost": { "units": 30, "currency": "USD" },
});
let (status, _) = post_reconcile(Arc::clone(&state), &reconcile_body).await;
assert_eq!(status, StatusCode::OK);
let after = mediated_test_state(signer, Arc::clone(&budget), Vec::new());
let (status, replay) = post_evaluate(Arc::clone(&after), &body).await;
assert_eq!(status, StatusCode::CONFLICT);
assert_eq!(replay["error"], "chio_request_id_reused");
assert_ne!(replay["status"], "authorized");
assert!(
replay["execution_nonce"].is_null(),
"a reused settled request_id must not mint a second nonce"
);
let fresh_body = serde_json::json!({
"capability": cap_value,
"tool_server": "cost-srv",
"tool_name": "compute",
"parameters": { "invoice": "inv-2" },
"request_id": "settled-reuse-fresh",
});
let (status, fresh) = post_evaluate(after, &fresh_body).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(fresh["status"], "authorized");
assert!(fresh["execution_nonce"].is_object());
}
}