klieo-ops 3.3.0

Operational layer above klieo-core: supervisor, governor, gates, escalation, worklog, handoff.
Documentation
//! Helper for emitting `OpsEvent` values into a `klieo-core` `EpisodicMemory`.
//!
//! This is the first production emission path for `Episode::Ops`. Agent code
//! calls `emit_ops_event` directly for now; automatic interception at the
//! `SupervisedAgent` level (tool-call / LLM-call hooks) is a separate Phase B
//! integration task.

use crate::ops_event::OpsEvent;
use klieo_core::error::MemoryError;
use klieo_core::ids::RunId;
use klieo_core::memory::{Episode, EpisodicMemory};

/// Emit a typed `OpsEvent` into `memory` as a `klieo_core::Episode::Ops`
/// payload under the provided `run_id`.
///
/// Serialises the event via [`OpsEvent::into_episode_payload`] and records
/// it with [`EpisodicMemory::record`]. Errors from either operation are
/// propagated as [`MemoryError`].
pub async fn emit_ops_event(
    memory: &dyn EpisodicMemory,
    run_id: RunId,
    event: OpsEvent,
) -> Result<(), MemoryError> {
    let payload = event
        .into_episode_payload()
        .map_err(|e| MemoryError::Serialization(format!("OpsEvent serialise failed: {e}")))?;
    memory.record(run_id, Episode::Ops(payload)).await
}

#[cfg(test)]
mod tests {
    use super::*;
    use klieo_core::memory::Episode;
    use klieo_core::test_utils::InMemoryEpisodic;
    use klieo_core::RunId;

    use crate::types::{AgentId, TenantId};

    #[tokio::test]
    async fn emit_records_ops_episode() {
        let memory = InMemoryEpisodic::default();
        let run_id = RunId::new();

        emit_ops_event(
            &memory,
            run_id,
            OpsEvent::SupervisorStateChange {
                tenant: Some(TenantId("BL_test".into())),
                agent: AgentId("test-agent".into()),
                state: "running".into(),
                reason: None,
            },
        )
        .await
        .expect("emit ok");

        let events = memory.replay(run_id).await.expect("replay ok");
        assert_eq!(events.len(), 1);
        assert!(
            matches!(&events[0], Episode::Ops(v) if v.get("kind").and_then(|k| k.as_str()) == Some("supervisor_state_change")),
            "expected Episode::Ops with kind=supervisor_state_change"
        );
    }

    #[tokio::test]
    async fn emit_gate_decision_round_trips() {
        let memory = InMemoryEpisodic::default();
        let run_id = RunId::new();

        emit_ops_event(
            &memory,
            run_id,
            OpsEvent::GateDecision {
                tenant: Some(TenantId("BL_test".into())),
                tool: "TestTool".into(),
                decision: "allow".into(),
                gate: "AllowlistGate".into(),
                policy_ref: None,
                reason: None,
            },
        )
        .await
        .expect("emit ok");

        let events = memory.replay(run_id).await.expect("replay ok");
        assert_eq!(events.len(), 1);
        let Episode::Ops(payload) = &events[0] else {
            panic!("expected Episode::Ops");
        };
        let recovered = OpsEvent::from_episode_payload(payload).expect("round-trip");
        let as_val = serde_json::to_value(&recovered).unwrap();
        assert_eq!(as_val["tool"], "TestTool");
        assert_eq!(as_val["decision"], "allow");
    }
}