use chrono::Utc;
use greentic_aw_runtime::StepObserver;
use greentic_types::TenantCtx;
use serde_json::{Value, json};
use super::audit_event::{agent_audit_subject, build_agent_audit_event};
use super::audit_sink::AuditSink;
use super::recorder::generate_audit_event_id;
pub struct AgentAuditObserver {
sink: AuditSink,
tenant: TenantCtx,
agent_id: String,
session_id: String,
}
impl AgentAuditObserver {
pub fn new(sink: AuditSink, tenant: TenantCtx, agent_id: String, session_id: String) -> Self {
Self {
sink,
tenant,
agent_id,
session_id,
}
}
fn emit(&self, event: &str, payload: Value) {
let envelope = build_agent_audit_event(
&self.tenant,
&self.agent_id,
&self.session_id,
event,
payload,
Utc::now(),
generate_audit_event_id(),
);
self.sink.emit(
agent_audit_subject(self.tenant.tenant.as_str(), event),
&envelope,
);
}
}
impl StepObserver for AgentAuditObserver {
fn wants_streaming(&self) -> bool {
false
}
fn on_token_delta(&self, _chunk: &str) {}
fn on_tool_call(&self, name: &str, call_id: &str) {
self.emit(
"tool_call",
json!({
"agent_id": self.agent_id,
"tool": name,
"call_id": call_id,
}),
);
}
fn on_tool_result(&self, name: &str, call_id: &str, result: &Value) {
self.emit(
"tool_result",
json!({
"agent_id": self.agent_id,
"tool": name,
"call_id": call_id,
"result": result,
}),
);
}
}
#[cfg(test)]
mod tests {
use super::*;
use greentic_types::{EnvId, TenantId};
use tokio::sync::mpsc;
fn tenant_ctx() -> TenantCtx {
TenantCtx::new(
EnvId::try_from("prod").expect("valid env id"),
TenantId::try_from("t1").expect("valid tenant id"),
)
}
fn observer_with_channel() -> (AgentAuditObserver, mpsc::Receiver<(String, Vec<u8>)>) {
let (tx, rx) = mpsc::channel(16);
let sink = AuditSink::from_sender(tx);
let observer =
AgentAuditObserver::new(sink, tenant_ctx(), "a1".to_string(), "s1".to_string());
(observer, rx)
}
#[test]
fn wants_streaming_is_false() {
let (observer, _rx) = observer_with_channel();
assert!(!observer.wants_streaming());
}
#[tokio::test]
async fn on_tool_call_enqueues_one_tool_call_event() {
let (observer, mut rx) = observer_with_channel();
observer.on_tool_call("http", "c1");
let (subject, bytes) = rx.try_recv().expect("event enqueued");
assert_eq!(subject, "audit.t1.agent.tool_call");
let value: Value = serde_json::from_slice(&bytes).expect("valid JSON");
assert!(
value
.get("type")
.and_then(Value::as_str)
.expect("type is a string")
.ends_with("agent.tool_call")
);
assert_eq!(
value.get("payload").and_then(|p| p.get("tool")),
Some(&json!("http"))
);
assert_eq!(
value.get("payload").and_then(|p| p.get("call_id")),
Some(&json!("c1"))
);
assert!(rx.try_recv().is_err(), "exactly one event enqueued");
}
#[tokio::test]
async fn on_tool_result_enqueues_one_tool_result_event_with_result_payload() {
let (observer, mut rx) = observer_with_channel();
observer.on_tool_result("http", "c1", &json!({"ok": true}));
let (subject, bytes) = rx.try_recv().expect("event enqueued");
assert_eq!(subject, "audit.t1.agent.tool_result");
let value: Value = serde_json::from_slice(&bytes).expect("valid JSON");
assert!(
value
.get("type")
.and_then(Value::as_str)
.expect("type is a string")
.ends_with("agent.tool_result")
);
assert_eq!(
value.get("payload").and_then(|p| p.get("result")),
Some(&json!({"ok": true}))
);
assert!(rx.try_recv().is_err(), "exactly one event enqueued");
}
#[tokio::test]
async fn on_token_delta_enqueues_nothing() {
let (observer, mut rx) = observer_with_channel();
observer.on_token_delta("chunk");
assert!(rx.try_recv().is_err());
}
}