struct SleepingGuard {
label: String,
}
impl Guard for SleepingGuard {
fn name(&self) -> &str {
&self.label
}
fn evaluate(&self, _ctx: &GuardContext) -> Result<GuardDecision, KernelError> {
std::thread::sleep(Duration::from_secs(2));
Ok(GuardDecision {
verdict: Verdict::Allow,
evidence: Vec::new(),
})
}
}
struct RecordingGuard {
label: String,
ran: Arc<AtomicU64>,
}
impl Guard for RecordingGuard {
fn name(&self) -> &str {
&self.label
}
fn evaluate(&self, _ctx: &GuardContext) -> Result<GuardDecision, KernelError> {
self.ran.fetch_add(1, Ordering::SeqCst);
Ok(GuardDecision {
verdict: Verdict::Allow,
evidence: Vec::new(),
})
}
}
struct HangingToolServer {
id: String,
tools: Vec<String>,
invocations: Arc<AtomicU64>,
}
#[async_trait::async_trait]
impl ToolServerConnection for HangingToolServer {
fn server_id(&self) -> &str {
&self.id
}
fn tool_names(&self) -> Vec<String> {
self.tools.clone()
}
async fn invoke(
&self,
_tool_name: &str,
_arguments: serde_json::Value,
_nested_flow_bridge: Option<&mut dyn NestedFlowBridge>,
) -> Result<serde_json::Value, KernelError> {
self.invocations.fetch_add(1, Ordering::SeqCst);
std::future::pending::<()>().await;
unreachable!("hanging tool server never returns")
}
}
struct BlockingToolServer {
id: String,
tools: Vec<String>,
invocations: Arc<AtomicU64>,
}
#[async_trait::async_trait]
impl ToolServerConnection for BlockingToolServer {
fn server_id(&self) -> &str {
&self.id
}
fn tool_names(&self) -> Vec<String> {
self.tools.clone()
}
async fn invoke(
&self,
_tool_name: &str,
_arguments: serde_json::Value,
_nested_flow_bridge: Option<&mut dyn NestedFlowBridge>,
) -> Result<serde_json::Value, KernelError> {
self.invocations.fetch_add(1, Ordering::SeqCst);
std::thread::sleep(Duration::from_secs(2));
Ok(serde_json::json!({ "ok": true }))
}
}
struct WedgedLivenessStore;
impl ReceiptStore for WedgedLivenessStore {
fn append_chio_receipt(&self, _receipt: &ChioReceipt) -> Result<(), ReceiptStoreError> {
Ok(())
}
fn append_child_receipt(
&self,
_receipt: &ChildRequestReceipt,
) -> Result<(), ReceiptStoreError> {
Ok(())
}
fn writer_liveness(&self, _stall_threshold: std::time::Duration) -> ReceiptWriterLiveness {
ReceiptWriterLiveness::Wedged
}
}
struct FirstReceiptAppendTimesOutStore {
calls: Arc<AtomicU64>,
unbounded_calls: Arc<AtomicU64>,
first_entered: mpsc::Sender<()>,
}
impl ReceiptStore for FirstReceiptAppendTimesOutStore {
fn append_chio_receipt(&self, _receipt: &ChioReceipt) -> Result<(), ReceiptStoreError> {
self.unbounded_calls.fetch_add(1, Ordering::SeqCst);
Err(ReceiptStoreError::Conflict(
"unbounded receipt append must not run".to_string(),
))
}
fn append_chio_receipt_with_timeout(
&self,
_receipt: &ChioReceipt,
budget: std::time::Duration,
) -> Result<Option<u64>, ReceiptStoreError> {
let call = self.calls.fetch_add(1, Ordering::SeqCst);
if call == 0 {
let _ = self.first_entered.send(());
std::thread::sleep(budget);
return Err(ReceiptStoreError::Timeout {
operation: "test receipt append".to_string(),
timeout_ms: budget.as_millis().min(u128::from(u64::MAX)) as u64,
});
}
Ok(Some(call + 1))
}
fn append_child_receipt(
&self,
_receipt: &ChildRequestReceipt,
) -> Result<(), ReceiptStoreError> {
Ok(())
}
fn writer_liveness(&self, _stall_threshold: std::time::Duration) -> ReceiptWriterLiveness {
ReceiptWriterLiveness::Healthy
}
}
#[test]
fn receipt_append_timeout_releases_write_lock_within_budget() {
let calls = Arc::new(AtomicU64::new(0));
let unbounded_calls = Arc::new(AtomicU64::new(0));
let (entered_tx, entered_rx) = mpsc::channel();
let mut config = make_config();
config.checkpoint_batch_size = 0;
config.deadlines.receipt_append_budget_ms = MIN_RECEIPT_APPEND_BUDGET_MS;
let keypair = config.keypair.clone();
let mut kernel = make_kernel(config);
kernel
.set_receipt_store(Box::new(FirstReceiptAppendTimesOutStore {
calls: Arc::clone(&calls),
unbounded_calls: Arc::clone(&unbounded_calls),
first_entered: entered_tx,
}))
.expect("install timeout store");
let kernel = Arc::new(kernel);
let first_receipt = make_signed_receipt(&keypair, "timeout-first");
let first_id = first_receipt.id.clone();
let second_receipt = make_signed_receipt(&keypair, "timeout-second");
let second_id = second_receipt.id.clone();
let first_kernel = Arc::clone(&kernel);
let started = std::time::Instant::now();
let first = thread::spawn(move || first_kernel.record_chio_receipt(&first_receipt));
entered_rx
.recv_timeout(Duration::from_secs(1))
.expect("bounded receipt append did not start");
let (second_tx, second_rx) = mpsc::channel();
let second_kernel = Arc::clone(&kernel);
let second = thread::spawn(move || {
let _ = second_tx.send(second_kernel.record_chio_receipt(&second_receipt));
});
let first_result = first.join().expect("first receipt thread panicked");
assert!(
started.elapsed() < Duration::from_secs(1),
"receipt append timeout did not return near its configured budget"
);
assert!(matches!(
first_result,
Err(KernelError::ReceiptPersistence(
ReceiptStoreError::Timeout {
timeout_ms: MIN_RECEIPT_APPEND_BUDGET_MS,
..
}
))
));
assert!(matches!(
second_rx.recv_timeout(Duration::from_secs(1)),
Ok(Ok(()))
));
second.join().expect("second receipt thread panicked");
assert_eq!(calls.load(Ordering::SeqCst), 2);
assert_eq!(unbounded_calls.load(Ordering::SeqCst), 0);
let receipt_ids: Vec<String> = kernel
.receipt_log()
.iter()
.map(|receipt| receipt.id.clone())
.collect();
assert!(!receipt_ids.contains(&first_id));
assert!(receipt_ids.contains(&second_id));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn guard_pipeline_budget_denies_hung_guard_and_frees_worker(
) -> Result<(), Box<dyn std::error::Error>> {
let mut config = make_config();
config.deadlines.guard_pipeline_budget_ms = 200;
let mut kernel = make_kernel(config);
kernel.add_guard(Box::new(SleepingGuard {
label: "sleeping".to_string(),
}));
kernel.register_tool_server(Box::new(EchoServer::new("srv-hpd", vec!["noop"])));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-hpd", "noop")]),
300,
);
let request = make_request("req-hpd-guard", &cap, "noop", "srv-hpd");
let kernel = Arc::new(kernel);
let start = std::time::Instant::now();
let response = kernel.evaluate_tool_call(&request).await?;
let elapsed = start.elapsed();
assert_eq!(response.verdict, Verdict::Deny);
assert!(
elapsed < Duration::from_secs(1),
"deadline should fire near 200ms, well before the 2s guard sleep, took {elapsed:?}"
);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn per_guard_budget_bounds_single_guard_not_pipeline(
) -> Result<(), Box<dyn std::error::Error>> {
let fast_ran = Arc::new(AtomicU64::new(0));
let mut config = make_config();
config
.deadlines
.per_guard_budget_ms
.insert("slow".to_string(), 200);
let mut kernel = make_kernel(config);
kernel.add_guard(Box::new(RecordingGuard {
label: "fast".to_string(),
ran: Arc::clone(&fast_ran),
}));
kernel.add_guard(Box::new(SleepingGuard {
label: "slow".to_string(),
}));
kernel.register_tool_server(Box::new(EchoServer::new("srv-pg", vec!["noop"])));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-pg", "noop")]),
300,
);
let request = make_request("req-pg", &cap, "noop", "srv-pg");
let kernel = Arc::new(kernel);
let start = std::time::Instant::now();
let response = kernel.evaluate_tool_call(&request).await?;
assert_eq!(response.verdict, Verdict::Deny);
assert!(start.elapsed() < Duration::from_secs(1));
assert_eq!(fast_ran.load(Ordering::SeqCst), 1);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn pipeline_budget_bounds_the_per_guard_loop() -> Result<(), Box<dyn std::error::Error>> {
let mut config = make_config();
config.deadlines.guard_pipeline_budget_ms = 300;
config
.deadlines
.per_guard_budget_ms
.insert("slow".to_string(), 5_000);
let mut kernel = make_kernel(config);
kernel.add_guard(Box::new(SleepingGuard {
label: "slow".to_string(),
}));
kernel.register_tool_server(Box::new(EchoServer::new("srv-pipeline", vec!["noop"])));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-pipeline", "noop")]),
300,
);
let request = make_request("req-pipeline", &cap, "noop", "srv-pipeline");
let kernel = Arc::new(kernel);
let start = std::time::Instant::now();
let response = kernel.evaluate_tool_call(&request).await?;
let elapsed = start.elapsed();
assert_eq!(response.verdict, Verdict::Deny);
assert!(
elapsed < Duration::from_secs(1),
"the pipeline deadline must fire near 300ms, well before the 2s guard sleep, took {elapsed:?}"
);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn dispatch_budget_expiry_runs_full_unwind_and_emits_cancelled_receipt(
) -> Result<(), Box<dyn std::error::Error>> {
let invocations = Arc::new(AtomicU64::new(0));
let mut config = make_config();
config.deadlines.dispatch_budget_ms = 200;
let mut kernel = make_kernel(config);
kernel.register_tool_server(Box::new(HangingToolServer {
id: "srv-hang".to_string(),
tools: vec!["noop".to_string()],
invocations: Arc::clone(&invocations),
}));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-hang", "noop")]),
300,
);
let request = make_request("req-dispatch-deadline", &cap, "noop", "srv-hang");
let kernel = Arc::new(kernel);
let start = std::time::Instant::now();
let response = kernel.evaluate_tool_call(&request).await?;
assert!(
start.elapsed() < Duration::from_secs(1),
"deadline must fire near 200ms"
);
assert_eq!(invocations.load(Ordering::SeqCst), 1, "dispatch did start");
assert_eq!(response.verdict, Verdict::Deny);
assert_eq!(
kernel.receipt_log().len(),
1,
"one Cancelled receipt persisted"
);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn wedged_writer_watchdog_denies_before_side_effect() -> Result<(), Box<dyn std::error::Error>>
{
let invocations = Arc::new(AtomicU64::new(0));
let mut kernel = make_kernel(make_config());
kernel.set_receipt_store(Box::new(WedgedLivenessStore))?;
kernel.register_tool_server(Box::new(SideEffectServer::new(
"srv-wedged",
vec!["noop"],
Arc::clone(&invocations),
)));
kernel.refresh_receipt_writer_liveness_for_test();
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-wedged", "noop")]),
300,
);
let request = make_request("req-wedged", &cap, "noop", "srv-wedged");
let kernel = Arc::new(kernel);
let response = kernel.evaluate_tool_call(&request).await?;
assert_eq!(response.verdict, Verdict::Deny);
assert_eq!(
invocations.load(Ordering::SeqCst),
0,
"no tool side effect may occur while the writer is wedged"
);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn wedged_writer_denies_before_dispatch_without_a_running_watchdog(
) -> Result<(), Box<dyn std::error::Error>> {
let invocations = Arc::new(AtomicU64::new(0));
let mut kernel = make_kernel(make_config());
kernel.set_receipt_store(Box::new(WedgedLivenessStore))?;
kernel.register_tool_server(Box::new(SideEffectServer::new(
"srv-no-watchdog",
vec!["noop"],
Arc::clone(&invocations),
)));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-no-watchdog", "noop")]),
300,
);
let request = make_request("req-no-watchdog", &cap, "noop", "srv-no-watchdog");
let kernel = Arc::new(kernel);
let response = kernel.evaluate_tool_call(&request).await?;
assert_eq!(response.verdict, Verdict::Deny);
assert_eq!(
invocations.load(Ordering::SeqCst),
0,
"no tool side effect may occur while the writer is wedged"
);
Ok(())
}
struct SnapshotCountingWedgedStore {
snapshot_writes: Arc<AtomicU64>,
}
impl ReceiptStore for SnapshotCountingWedgedStore {
fn append_chio_receipt(&self, _receipt: &ChioReceipt) -> Result<(), ReceiptStoreError> {
Ok(())
}
fn append_child_receipt(
&self,
_receipt: &ChildRequestReceipt,
) -> Result<(), ReceiptStoreError> {
Ok(())
}
fn record_capability_snapshot(
&self,
_token: &CapabilityToken,
_parent_capability_id: Option<&str>,
) -> Result<(), ReceiptStoreError> {
self.snapshot_writes.fetch_add(1, Ordering::SeqCst);
Ok(())
}
fn writer_liveness(&self, _stall_threshold: std::time::Duration) -> ReceiptWriterLiveness {
ReceiptWriterLiveness::Wedged
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn wedged_writer_denies_before_evaluation_capability_snapshot(
) -> Result<(), Box<dyn std::error::Error>> {
let invocations = Arc::new(AtomicU64::new(0));
let snapshot_writes = Arc::new(AtomicU64::new(0));
let mut kernel = make_kernel(make_config());
kernel.set_receipt_store(Box::new(SnapshotCountingWedgedStore {
snapshot_writes: Arc::clone(&snapshot_writes),
}))?;
kernel.register_tool_server(Box::new(SideEffectServer::new(
"srv-snapshot",
vec!["noop"],
Arc::clone(&invocations),
)));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-snapshot", "noop")]),
300,
);
kernel.refresh_receipt_writer_liveness_for_test();
let snapshots_before_dispatch = snapshot_writes.load(Ordering::SeqCst);
let request = make_request("req-snapshot", &cap, "noop", "srv-snapshot");
let kernel = Arc::new(kernel);
let response = kernel.evaluate_tool_call(&request).await?;
assert_eq!(response.verdict, Verdict::Deny);
assert_eq!(
snapshot_writes.load(Ordering::SeqCst),
snapshots_before_dispatch,
"evaluation must deny before entering the capability snapshot write"
);
assert_eq!(
invocations.load(Ordering::SeqCst),
0,
"no tool side effect may occur while the writer is wedged"
);
Ok(())
}
struct SnapshotBudgetStore {
bounded_budget_ms: Arc<AtomicU64>,
unbounded_calls: Arc<AtomicU64>,
}
impl ReceiptStore for SnapshotBudgetStore {
fn append_chio_receipt(&self, _receipt: &ChioReceipt) -> Result<(), ReceiptStoreError> {
Ok(())
}
fn append_child_receipt(
&self,
_receipt: &ChildRequestReceipt,
) -> Result<(), ReceiptStoreError> {
Ok(())
}
fn record_capability_snapshot(
&self,
_token: &CapabilityToken,
_parent_capability_id: Option<&str>,
) -> Result<(), ReceiptStoreError> {
self.unbounded_calls.fetch_add(1, Ordering::SeqCst);
Ok(())
}
fn record_capability_snapshot_with_timeout(
&self,
_token: &CapabilityToken,
_parent_capability_id: Option<&str>,
budget: std::time::Duration,
) -> Result<(), ReceiptStoreError> {
self.bounded_budget_ms
.store(budget.as_millis() as u64, Ordering::SeqCst);
Ok(())
}
fn writer_liveness(&self, _stall_threshold: std::time::Duration) -> ReceiptWriterLiveness {
ReceiptWriterLiveness::Healthy
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn healthy_writer_records_capability_snapshot_through_the_bounded_path(
) -> Result<(), Box<dyn std::error::Error>> {
let bounded_budget_ms = Arc::new(AtomicU64::new(0));
let unbounded_calls = Arc::new(AtomicU64::new(0));
let config = make_config();
let expected_budget_ms =
u64::try_from(config.deadlines.receipt_append_budget().as_millis()).unwrap_or(u64::MAX);
let mut kernel = make_kernel(config);
kernel.set_receipt_store(Box::new(SnapshotBudgetStore {
bounded_budget_ms: Arc::clone(&bounded_budget_ms),
unbounded_calls: Arc::clone(&unbounded_calls),
}))?;
kernel.register_tool_server(Box::new(EchoServer::new("srv-snap-budget", vec!["noop"])));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-snap-budget", "noop")]),
300,
);
bounded_budget_ms.store(0, Ordering::SeqCst);
unbounded_calls.store(0, Ordering::SeqCst);
let request = make_request("req-snap-budget", &cap, "noop", "srv-snap-budget");
let kernel = Arc::new(kernel);
let response = kernel.evaluate_tool_call(&request).await?;
assert_eq!(response.verdict, Verdict::Allow);
assert_eq!(
bounded_budget_ms.load(Ordering::SeqCst),
expected_budget_ms,
"the observed-capability snapshot must use the bounded writer path with the append budget"
);
assert_eq!(
unbounded_calls.load(Ordering::SeqCst),
0,
"the hot-path snapshot must not use the unbounded writer path"
);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn phase_dispatch_requires_the_full_evaluation_pipeline(
) -> Result<(), Box<dyn std::error::Error>> {
use crate::kernel::evaluator::{BlockingToolEvaluator, ToolEvaluator};
let invocations = Arc::new(AtomicU64::new(0));
let mut kernel = make_kernel(make_config());
kernel.register_tool_server(Box::new(HangingToolServer {
id: "srv-phase-hang".to_string(),
tools: vec!["noop".to_string()],
invocations: Arc::clone(&invocations),
}));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-phase-hang", "noop")]),
300,
);
let request = make_request("req-phase-dispatch", &cap, "noop", "srv-phase-hang");
let kernel = Arc::new(kernel);
let result = BlockingToolEvaluator
.dispatch(&kernel, &request, false)
.await;
assert_eq!(invocations.load(Ordering::SeqCst), 0);
assert!(matches!(
result,
Err(KernelError::DirectDispatchUnavailable)
));
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn dispatch_budget_bounds_a_connection_that_blocks_before_awaiting(
) -> Result<(), Box<dyn std::error::Error>> {
let invocations = Arc::new(AtomicU64::new(0));
let mut config = make_config();
config.deadlines.dispatch_budget_ms = 200;
let mut kernel = make_kernel(config);
kernel.register_tool_server(Box::new(BlockingToolServer {
id: "srv-blocking".to_string(),
tools: vec!["noop".to_string()],
invocations: Arc::clone(&invocations),
}));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-blocking", "noop")]),
300,
);
let request = make_request("req-blocking-dispatch", &cap, "noop", "srv-blocking");
let kernel = Arc::new(kernel);
let start = std::time::Instant::now();
let response = kernel.evaluate_tool_call(&request).await?;
let elapsed = start.elapsed();
assert!(
elapsed < Duration::from_secs(1),
"a connection that blocks before awaiting must be bounded near the 200ms dispatch budget, took {elapsed:?}"
);
assert_eq!(
invocations.load(Ordering::SeqCst),
1,
"dispatch did start the blocking connection"
);
assert_eq!(response.verdict, Verdict::Deny);
Ok(())
}
struct ThreadRecordingGuard {
label: String,
thread: Arc<std::sync::Mutex<Option<std::thread::ThreadId>>>,
}
impl Guard for ThreadRecordingGuard {
fn name(&self) -> &str {
&self.label
}
fn evaluate(&self, _ctx: &GuardContext) -> Result<GuardDecision, KernelError> {
*self.thread.lock().expect("record guard thread") = Some(std::thread::current().id());
Ok(GuardDecision {
verdict: Verdict::Allow,
evidence: Vec::new(),
})
}
}
#[test]
fn always_offload_moves_guards_off_the_async_worker_without_a_timer(
) -> Result<(), Box<dyn std::error::Error>> {
let mut config = make_config();
config.deadlines.always_offload_guards = true;
let mut kernel = make_kernel(config);
let guard_thread = Arc::new(std::sync::Mutex::new(None));
kernel.add_guard(Box::new(ThreadRecordingGuard {
label: "recording".to_string(),
thread: Arc::clone(&guard_thread),
}));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-offload", "noop")]),
300,
);
let request = make_request("req-offload", &cap, "noop", "srv-offload");
let scope = make_scope(vec![make_grant("srv-offload", "noop")]);
let runtime = tokio::runtime::Builder::new_current_thread().build()?;
runtime.block_on(async {
assert!(
!super::dispatch::dispatch_timer_available(),
"the test runtime must be timerless"
);
let worker = std::thread::current().id();
let outcome = kernel
.run_guards_within_budget(&request, &scope, None, None)
.await;
assert!(outcome.is_ok(), "the recording guard allows");
let guard = guard_thread
.lock()
.expect("read guard thread")
.expect("guard ran");
assert_ne!(
guard, worker,
"always_offload must move the guard off the async worker even without a timer"
);
});
Ok(())
}
#[test]
fn always_offload_moves_guards_off_the_worker_without_a_timer_even_with_a_budget(
) -> Result<(), Box<dyn std::error::Error>> {
let mut config = make_config();
config.deadlines.always_offload_guards = true;
config.deadlines.guard_pipeline_budget_ms = 200;
let mut kernel = make_kernel(config);
let guard_thread = Arc::new(std::sync::Mutex::new(None));
kernel.add_guard(Box::new(ThreadRecordingGuard {
label: "recording".to_string(),
thread: Arc::clone(&guard_thread),
}));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-offload", "noop")]),
300,
);
let request = make_request("req-offload-budget", &cap, "noop", "srv-offload");
let scope = make_scope(vec![make_grant("srv-offload", "noop")]);
let runtime = tokio::runtime::Builder::new_current_thread().build()?;
runtime.block_on(async {
assert!(
!super::dispatch::dispatch_timer_available(),
"the test runtime must be timerless"
);
let worker = std::thread::current().id();
let outcome = kernel
.run_guards_within_budget(&request, &scope, None, None)
.await;
assert!(outcome.is_ok(), "the recording guard allows");
let guard = guard_thread
.lock()
.expect("read guard thread")
.expect("guard ran");
assert_ne!(
guard, worker,
"a configured budget must not defeat always_offload in a timerless runtime"
);
});
Ok(())
}
#[test]
fn always_offload_runs_guards_inline_without_a_tokio_runtime(
) -> Result<(), Box<dyn std::error::Error>> {
let ran = Arc::new(AtomicU64::new(0));
let mut config = make_config();
config.deadlines.always_offload_guards = true;
let mut kernel = make_kernel(config);
kernel.add_guard(Box::new(RecordingGuard {
label: "recording".to_string(),
ran: Arc::clone(&ran),
}));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-offload", "noop")]),
300,
);
let request = make_request("req-offload-no-runtime", &cap, "noop", "srv-offload");
let scope = make_scope(vec![make_grant("srv-offload", "noop")]);
let outcome =
futures::executor::block_on(kernel.run_guards_within_budget(&request, &scope, None, None));
assert!(
outcome.is_ok(),
"guards must run inline without a runtime instead of panicking in spawn_blocking"
);
assert_eq!(
ran.load(Ordering::SeqCst),
1,
"the guard must still execute on the inline fallback"
);
Ok(())
}
#[test]
fn nested_dispatch_isolates_a_synchronously_blocking_call_from_the_async_pool(
) -> Result<(), Box<dyn std::error::Error>> {
let budget = Duration::from_millis(50);
let block = Duration::from_millis(400);
const HEARTBEAT_MS: u64 = 5;
async fn count_heartbeats_while<F, Fut>(make_call: F) -> u64
where
F: FnOnce() -> Fut,
Fut: std::future::Future<Output = ()> + Send + 'static,
{
let ticks = Arc::new(AtomicU64::new(0));
let ticks_beat = Arc::clone(&ticks);
let beat = tokio::spawn(async move {
loop {
ticks_beat.fetch_add(1, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(HEARTBEAT_MS)).await;
}
});
let blocking = tokio::spawn(make_call());
let _ = blocking.await;
beat.abort();
ticks.load(Ordering::SeqCst)
}
let runtime = tokio::runtime::Builder::new_multi_thread()
.worker_threads(1)
.enable_all()
.build()?;
let inline_ticks = runtime.block_on(count_heartbeats_while(|| async move {
let _ = tokio::time::timeout(budget, async {
std::thread::sleep(block);
Ok::<_, KernelError>(ToolServerOutput::Value(serde_json::json!({ "ok": true })))
})
.await;
}));
let helper_ticks = runtime.block_on(count_heartbeats_while(|| async move {
let call = async {
std::thread::sleep(block);
Ok::<_, KernelError>(ToolServerOutput::Value(serde_json::json!({ "ok": true })))
};
let output =
crate::kernel::dispatch::dispatch_nested_call_within_budget(call, budget).await;
assert!(
matches!(output, Ok(ToolServerOutput::Value(_))),
"the blocking nested call completes through the helper"
);
}));
let expected_live = block.as_millis() as u64 / HEARTBEAT_MS;
assert!(
inline_ticks <= 2,
"the inline timeout pins the sole worker, starving the heartbeat (ticks={inline_ticks})"
);
assert!(
helper_ticks >= expected_live / 4,
"block_in_place must keep the async pool alive while the nested call blocks (ticks={helper_ticks}, expected ~{expected_live})"
);
Ok(())
}
#[test]
fn timed_out_dispatch_aborts_a_queued_blocking_task_before_it_runs_the_tool(
) -> Result<(), Box<dyn std::error::Error>> {
let invocations = Arc::new(AtomicU64::new(0));
let mut config = make_config();
config.deadlines.dispatch_budget_ms = 100;
let mut kernel = make_kernel(config);
kernel.register_tool_server(Box::new(SideEffectServer::new(
"srv-abort",
vec!["noop"],
Arc::clone(&invocations),
)));
let agent_kp = make_keypair();
let cap = make_capability(
&kernel,
&agent_kp,
make_scope(vec![make_grant("srv-abort", "noop")]),
300,
);
let request = make_request("req-abort", &cap, "noop", "srv-abort");
let kernel = Arc::new(kernel);
let runtime = tokio::runtime::Builder::new_multi_thread()
.worker_threads(2)
.max_blocking_threads(1)
.enable_all()
.build()?;
runtime.block_on(async {
let blocker_started = Arc::new(AtomicBool::new(false));
let started = Arc::clone(&blocker_started);
let _blocker = tokio::task::spawn_blocking(move || {
started.store(true, Ordering::SeqCst);
std::thread::sleep(Duration::from_millis(800));
});
while !blocker_started.load(Ordering::SeqCst) {
tokio::time::sleep(Duration::from_millis(5)).await;
}
let start = std::time::Instant::now();
let response = kernel
.evaluate_tool_call(&request)
.await
.unwrap_or_else(|e| panic!("evaluate should return a deny response, not error: {e}"));
assert!(
start.elapsed() < Duration::from_millis(600),
"the dispatch deadline must fire near 100ms, well before the 800ms blocker frees the pool"
);
assert_eq!(response.verdict, Verdict::Deny);
assert_eq!(
invocations.load(Ordering::SeqCst),
0,
"the queued dispatch must not have started before the deadline fired"
);
tokio::time::sleep(Duration::from_millis(1_000)).await;
assert_eq!(
invocations.load(Ordering::SeqCst),
0,
"a timed-out dispatch must be aborted, not left to run the tool after the deadline"
);
});
Ok(())
}
#[test]
fn watchdog_does_not_start_without_a_timer() -> Result<(), Box<dyn std::error::Error>> {
let kernel = Arc::new(make_kernel(make_config()));
let runtime = tokio::runtime::Builder::new_current_thread().build()?;
runtime.block_on(async {
assert!(
!super::dispatch::dispatch_timer_available(),
"the test runtime must be timerless"
);
kernel.spawn_receipt_writer_watchdog();
assert!(
!kernel.receipt_writer_watchdog_is_running(),
"no watchdog poll task may start without a time driver"
);
});
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn watchdog_does_not_keep_the_kernel_alive() {
let mut config = make_config();
config.deadlines.receipt_writer_poll_ms = 60_000;
let kernel = Arc::new(make_kernel(config));
let weak = Arc::downgrade(&kernel);
kernel.spawn_receipt_writer_watchdog();
tokio::time::sleep(Duration::from_millis(50)).await;
drop(kernel);
tokio::time::sleep(Duration::from_millis(50)).await;
assert!(
weak.upgrade().is_none(),
"the watchdog task must not keep the kernel alive after its last external Arc is dropped"
);
}