klieo-ops 0.2.0

Operational layer above klieo-core: supervisor, governor, gates, escalation, worklog, handoff.
Documentation
//! Public conformance fixtures consumed by third-party impl crates as
//! dev-dependencies. Each `run_<primitive>_conformance` asserts the
//! documented contract end-to-end.

#![allow(missing_docs)]

use crate::supervisor::{Status, Supervisor, SupervisorEvent};
use crate::types::{AgentId, AgentMeta, KillReason, KillTrigger, RuntimeState};
use futures::StreamExt;
use std::time::Duration;

/// Conformance: register, heartbeat, status report, kill-switch.
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,
    };

    // 1. heartbeat before register MUST fail with UnknownAgent.
    let err = s.heartbeat(agent.clone()).await.expect_err("must reject");
    assert!(matches!(
        err,
        crate::supervisor::SupervisorError::UnknownAgent(_)
    ));

    // 2. register, then heartbeat, then watch sees the events.
    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");

    // 3. runtime state starts at Running.
    assert_eq!(s.runtime_state().await, RuntimeState::Running);

    // 4. trip kill switch → state is Halted.
    s.trip_kill_switch(
        KillReason("integration test".into()),
        KillTrigger::Programmatic,
    )
    .await
    .expect("kill ok");
    assert_eq!(s.runtime_state().await, RuntimeState::Halted);
}

/// Conformance: acquire-release round-trip + saturation detection.
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());

    // 1. Acquire then drop — should not deadlock or starve.
    {
        let _p: Permit = g
            .acquire_llm(provider.clone(), 100)
            .await
            .expect("first acquire succeeds");
    }

    // 2. Acquire to saturation; exact count depends on refill timing.
    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);
}