use crate::ops_event::OpsEvent;
use klieo_core::error::MemoryError;
use klieo_core::ids::RunId;
use klieo_core::memory::{Episode, EpisodicMemory};
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");
}
}