greentic-runner-host 1.1.6

Host runtime shim for Greentic runner: config, pack loading, activity handling
Documentation
//! `StepObserver` implementation that emits a best-effort audit event per
//! agentic-worker step (tool call/result), published to
//! `audit.<tenant>.agent.<event>` and ingested by the same admin subscriber
//! B-2 wired up for per-flow-node events. See
//! `docs/superpowers/specs/2026-07-03-agent-audit-emitter-design.md`.
//!
//! Gated behind the `agentic-worker` feature: `greentic_aw_runtime::StepObserver`
//! is only available when that (default-on) feature pulls the crate in.

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;

/// Emits `audit.<tenant>.agent.<event>` events for one agentic-worker
/// invocation's tool calls/results.
///
/// Best-effort throughout — [`AuditSink::emit`] never blocks, errors, or
/// panics, so this observer can never affect the agent loop's behavior.
/// `wants_streaming` is always `false`: audit only cares about tool steps,
/// not token-level streaming.
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());
    }
}