#![allow(missing_docs)]
use crate::supervisor::{Status, Supervisor, SupervisorEvent};
use crate::types::{AgentId, AgentMeta, KillReason, KillTrigger, RuntimeState};
use futures::StreamExt;
use std::time::Duration;
pub async fn run_supervisor_conformance<S: Supervisor>(s: S) {
let agent = AgentId("test-agent".into());
let meta = AgentMeta {
id: agent.clone(),
role: "test".into(),
version: "0.0.0".into(),
identity_pubkey: [0u8; 32],
expected_step_p99: None,
};
let err = s.heartbeat(agent.clone()).await.expect_err("must reject");
assert!(matches!(
err,
crate::supervisor::SupervisorError::UnknownAgent(_)
));
let mut stream = s.watch().await;
s.register(agent.clone(), meta).await.expect("register ok");
s.heartbeat(agent.clone()).await.expect("heartbeat ok");
s.report(agent.clone(), Status::Idle)
.await
.expect("report ok");
let mut seen_heartbeat = false;
for _ in 0..4u8 {
let next = tokio::time::timeout(Duration::from_millis(200), stream.next()).await;
if let Ok(Some(SupervisorEvent::Heartbeat(_))) = next {
seen_heartbeat = true;
}
if seen_heartbeat {
break;
}
}
assert!(seen_heartbeat, "Heartbeat event must be emitted");
assert_eq!(s.runtime_state().await, RuntimeState::Running);
s.trip_kill_switch(
KillReason("integration test".into()),
KillTrigger::Programmatic,
)
.await
.expect("kill ok");
assert_eq!(s.runtime_state().await, RuntimeState::Halted);
}
pub async fn run_governor_conformance<G: crate::governor::Governor>(g: G) {
use crate::governor::Permit;
use crate::types::ProviderId;
let provider = ProviderId("anthropic".into());
{
let _p: Permit = g
.acquire_llm(provider.clone(), 100)
.await
.expect("first acquire succeeds");
}
let mut held: Vec<Permit> = Vec::new();
for _ in 0..6u8 {
match g.acquire_llm(provider.clone(), 100).await {
Ok(p) => held.push(p),
Err(_) => break,
}
}
drop(held);
}
#[cfg(feature = "escalation")]
use crate::escalation::{
Escalation, EscalationFilter, EscalationTicket, Resolution, ResolutionOutcome, Severity,
};
#[cfg(feature = "escalation")]
pub async fn run_escalation_conformance<E: Escalation>(esc: E) {
let id_high = esc
.raise(EscalationTicket {
tenant: Some(crate::types::TenantId("tenant_a".into())),
severity: Severity::High,
reason: "high ticket".into(),
provenance: None,
})
.await
.expect("raise high");
let id_low = esc
.raise(EscalationTicket {
tenant: Some(crate::types::TenantId("tenant_b".into())),
severity: Severity::Low,
reason: "low ticket".into(),
provenance: None,
})
.await
.expect("raise low");
let _id_crit = esc
.raise(EscalationTicket {
tenant: None,
severity: Severity::Critical,
reason: "critical ticket".into(),
provenance: None,
})
.await
.expect("raise critical");
let all = esc.list(EscalationFilter::default()).await;
assert_eq!(all.len(), 3, "default filter returns all raised tickets");
let at_least_high = esc
.list(EscalationFilter {
min_severity: Some(Severity::High),
..Default::default()
})
.await;
assert_eq!(at_least_high.len(), 2, "min_severity=High drops Low ticket");
let tenant_a = esc
.list(EscalationFilter {
tenant: Some(crate::types::TenantId("tenant_a".into())),
..Default::default()
})
.await;
assert_eq!(tenant_a.len(), 1, "tenant filter narrows to one");
esc.resolve(
id_high,
Resolution {
outcome: ResolutionOutcome::Approved,
reason: Some("approved".into()),
approvers: vec![],
signatures: vec![],
},
)
.await
.expect("resolve approved");
let err = esc
.resolve(
crate::escalation::EscalationId("does_not_exist".into()),
Resolution {
outcome: ResolutionOutcome::Denied,
reason: None,
approvers: vec![],
signatures: vec![],
},
)
.await
.expect_err("unknown ticket must error");
assert!(matches!(
err,
crate::escalation::EscalationError::UnknownTicket(_)
));
esc.resolve(
id_low,
Resolution {
outcome: ResolutionOutcome::Halted,
reason: Some("kill-switch tripped".into()),
approvers: vec![],
signatures: vec![],
},
)
.await
.expect("resolve halted");
assert_eq!(Severity::Low.escalate_one_level(), Severity::Medium);
assert_eq!(Severity::High.escalate_one_level(), Severity::Critical);
assert_eq!(Severity::Critical.escalate_one_level(), Severity::Critical);
}
#[cfg(feature = "worklog")]
use crate::worklog::{WorkId, WorkItem, WorkLog, WorkLogError, WorkStatus};
#[cfg(feature = "worklog")]
fn _fresh_item(title: &str) -> WorkItem {
WorkItem {
id: WorkId(String::new()),
tenant: None,
title: title.into(),
payload: serde_json::json!({}),
status: WorkStatus::Planned,
depends_on: vec![],
last_transition_at: String::new(),
}
}
#[cfg(feature = "worklog")]
pub async fn run_worklog_conformance<W: WorkLog>(wl: W) {
let solo = wl.plan(_fresh_item("solo")).await.expect("plan solo");
assert_eq!(
wl.get(solo.clone()).await.unwrap().status,
WorkStatus::Ready
);
let parent = wl.plan(_fresh_item("parent")).await.expect("plan parent");
let child_item = WorkItem {
depends_on: vec![parent.clone()],
.._fresh_item("child")
};
let child = wl.plan(child_item).await.expect("plan child");
assert_eq!(
wl.get(child.clone()).await.unwrap().status,
WorkStatus::Planned
);
wl.transition(parent.clone(), WorkStatus::Done)
.await
.expect("transition parent Done");
assert_eq!(
wl.get(child.clone()).await.unwrap().status,
WorkStatus::Ready
);
let still_planned = WorkItem {
depends_on: vec![child.clone()],
.._fresh_item("blocked")
};
let blocked = wl.plan(still_planned).await.expect("plan blocked");
let err = wl
.dispatch(blocked.clone())
.await
.expect_err("planned cannot dispatch");
assert!(matches!(err, WorkLogError::Internal(_)));
wl.dispatch(solo.clone()).await.expect("dispatch solo");
assert_eq!(
wl.get(solo.clone()).await.unwrap().status,
WorkStatus::InProgress
);
let a = wl.plan(_fresh_item("cyc_a")).await.expect("plan cyc_a");
let b = wl.plan(_fresh_item("cyc_b")).await.expect("plan cyc_b");
wl.depend(b.clone(), a.clone())
.await
.expect("b depends on a");
let err = wl
.depend(a.clone(), b.clone())
.await
.expect_err("cycle must be rejected");
assert!(matches!(err, WorkLogError::CycleDetected { .. }));
wl.transition(solo.clone(), WorkStatus::Done)
.await
.expect("solo Done");
let err = wl
.transition(solo.clone(), WorkStatus::Ready)
.await
.expect_err("cannot un-Done");
assert!(matches!(err, WorkLogError::Internal(_)));
let err = wl
.transition(WorkId("does_not_exist".into()), WorkStatus::Done)
.await
.expect_err("missing id");
assert!(matches!(err, WorkLogError::UnknownWorkItem(_)));
for i in 0..3 {
wl.plan(_fresh_item(&format!("extra_{i}")))
.await
.expect("plan extras");
}
let bounded = wl.ready(2).await;
assert!(bounded.len() <= 2, "ready(2) must return at most 2");
let subgraph = wl.dag(child.clone()).await;
assert!(subgraph.items.iter().any(|i| i.id == parent));
assert!(subgraph.items.iter().any(|i| i.id == child));
assert!(subgraph.items.iter().any(|i| i.id == blocked));
}
#[cfg(feature = "handoff")]
use crate::handoff::{Handoff, HandoffError, HandoffState};
#[cfg(feature = "handoff")]
use crate::signer::SoftwareSigner;
#[cfg(feature = "handoff")]
use crate::types::{AgentId as HoAgentId, TenantId as HoTenantId};
#[cfg(feature = "handoff")]
pub async fn run_handoff_conformance<H: Handoff>(h: H) {
let signer = SoftwareSigner::from_bytes([11u8; 32]);
let state = HandoffState {
from_run: "rn_conformance".into(),
merkle_root: "abad1dea".into(),
at_seq: Some(7),
redacted_state: b"opaque-cbor-bytes".to_vec(),
source_identity: HoAgentId("source-agent".into()),
tenant: Some(HoTenantId("T1".into())),
};
let envelope = h
.package(state.clone(), &signer, chrono::Duration::hours(1))
.await
.expect("package");
let proof = h.verify(&envelope).await.expect("verify");
assert_eq!(proof.from_run, "rn_conformance");
assert_eq!(proof.source_identity.0, "source-agent");
assert_eq!(proof.tenant.as_ref().map(|t| t.0.as_str()), Some("T1"));
assert_eq!(proof.redacted_state, b"opaque-cbor-bytes".to_vec());
let mut tampered = envelope.clone();
tampered.redacted_state[0] ^= 0xff;
let err = h.verify(&tampered).await.expect_err("tamper must reject");
assert!(matches!(err, HandoffError::SignatureInvalid(_)));
let short_state = HandoffState { ..state.clone() };
let short = h
.package(short_state, &signer, chrono::Duration::seconds(0))
.await
.expect("package zero-ttl");
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
let err = h.verify(&short).await.expect_err("expired must reject");
assert!(matches!(err, HandoffError::Expired { .. }));
let delivered = h
.deliver(envelope.clone(), HoAgentId("receiver".into()))
.await
.expect("deliver");
assert_eq!(delivered, "rn_conformance");
}
#[cfg(feature = "gates")]
use crate::approver_registry::{ApproverId, ApproverRegistry, StaticApproverRegistry};
#[cfg(feature = "gates")]
use crate::gates::{ApprovalError, ApprovalOutcome, FourEyesGate, Gate, GateDecision, GateRequest};
#[cfg(feature = "gates")]
use ed25519_dalek::{Signer, SigningKey};
#[cfg(feature = "gates")]
use klieo_core::KvStore;
#[cfg(feature = "gates")]
use std::collections::HashMap;
#[cfg(feature = "gates")]
use std::sync::Arc;
#[cfg(feature = "gates")]
pub fn build_test_four_eyes(
kv: Arc<dyn KvStore>,
approver_seeds: &[(&str, u8)],
quorum: u8,
dual_control_tool_names: &[&str],
) -> (FourEyesGate, HashMap<String, SigningKey>) {
let mut keys = HashMap::new();
let mut signers = HashMap::new();
for (id, seed) in approver_seeds {
let sk = SigningKey::from_bytes(&[*seed; 32]);
keys.insert(ApproverId((*id).to_string()), sk.verifying_key());
signers.insert((*id).to_string(), sk);
}
let registry: Arc<dyn ApproverRegistry> = Arc::new(StaticApproverRegistry::from_map(keys));
let classifier = FourEyesGate::dual_control_tools(dual_control_tool_names);
let gate = FourEyesGate::new(kv, registry, quorum, classifier);
(gate, signers)
}
#[cfg(feature = "gates")]
pub async fn run_four_eyes_conformance(kv: Arc<dyn KvStore>) {
let (gate, _signers) =
build_test_four_eyes(kv.clone(), &[("alice", 1), ("bob", 2)], 2, &["payout"]);
let decision = gate
.evaluate(GateRequest::new("payout", serde_json::json!({})))
.await;
assert!(matches!(
decision,
GateDecision::RequireApproval { quorum: 2, .. }
));
let decision = gate
.evaluate(GateRequest::new("read_only", serde_json::json!({})))
.await;
assert!(matches!(decision, GateDecision::Allow));
let (gate, signers) =
build_test_four_eyes(kv.clone(), &[("alice", 1), ("bob", 2)], 2, &["payout"]);
let GateDecision::RequireApproval { ticket, .. } = gate
.evaluate(GateRequest::new("payout", serde_json::json!({})))
.await
else {
panic!("expected RequireApproval");
};
let payload = b"conformance-payload";
gate.submit_approval(
&ticket,
ApproverId("alice".into()),
signers["alice"].sign(payload),
payload.to_vec(),
)
.expect("alice");
gate.submit_approval(
&ticket,
ApproverId("bob".into()),
signers["bob"].sign(payload),
payload.to_vec(),
)
.expect("bob");
let outcome = gate
.wait_for_approval(ticket, Duration::from_secs(1))
.await
.expect("approval ok");
assert!(matches!(outcome, ApprovalOutcome::Allow));
let (gate, _) = build_test_four_eyes(kv.clone(), &[("alice", 1), ("bob", 2)], 2, &["payout"]);
let GateDecision::RequireApproval { ticket, .. } = gate
.evaluate(GateRequest::new("payout", serde_json::json!({})))
.await
else {
panic!();
};
let err = gate
.wait_for_approval(ticket, Duration::from_millis(350))
.await
.expect_err("must timeout");
assert!(matches!(err, ApprovalError::TimedOut { .. }));
let (gate, _) = build_test_four_eyes(kv.clone(), &[("alice", 1), ("bob", 2)], 2, &["payout"]);
let gate = Arc::new(gate);
let GateDecision::RequireApproval { ticket, .. } = gate
.evaluate(GateRequest::new("payout", serde_json::json!({})))
.await
else {
panic!();
};
let g2 = gate.clone();
let ticket_copy = ticket.clone();
let waiter = tokio::spawn(async move {
g2.wait_for_approval(ticket_copy, Duration::from_secs(2))
.await
});
tokio::time::sleep(Duration::from_millis(50)).await;
gate.deny_approval(&ticket, "operator denied for conformance test");
let result = waiter.await.unwrap();
assert!(matches!(result, Err(ApprovalError::Denied(_))));
let (gate, signers) = build_test_four_eyes(kv.clone(), &[("alice", 1)], 2, &["payout"]);
let GateDecision::RequireApproval { ticket, .. } = gate
.evaluate(GateRequest::new("payout", serde_json::json!({})))
.await
else {
panic!();
};
let payload = b"conformance";
gate.submit_approval(
&ticket,
ApproverId("alice".into()),
signers["alice"].sign(payload),
payload.to_vec(),
)
.expect("alice");
let bob = SigningKey::from_bytes(&[99u8; 32]);
let err = gate
.submit_approval(
&ticket,
ApproverId("bob".into()),
bob.sign(payload),
payload.to_vec(),
)
.expect_err("bob not registered");
assert!(matches!(err, ApprovalError::VerificationFailed(_)));
let (gate, signers) = build_test_four_eyes(kv.clone(), &[("alice", 1)], 2, &["payout"]);
let GateDecision::RequireApproval { ticket, .. } = gate
.evaluate(GateRequest::new("payout", serde_json::json!({})))
.await
else {
panic!();
};
let payload = b"conformance";
gate.submit_approval(
&ticket,
ApproverId("alice".into()),
signers["alice"].sign(payload),
payload.to_vec(),
)
.expect("alice 1");
gate.submit_approval(
&ticket,
ApproverId("alice".into()),
signers["alice"].sign(payload),
payload.to_vec(),
)
.expect("alice 2 — accepted but should not count for quorum");
let err = gate
.wait_for_approval(ticket, Duration::from_millis(350))
.await
.expect_err("must NOT reach quorum from duplicate identity");
assert!(matches!(err, ApprovalError::TimedOut { .. }));
}