use chrono::{DateTime, Utc};
use greentic_types::{EventEnvelope, EventId, TenantCtx};
use serde_json::json;
const FALLBACK_EVENT_ID: &str = "audit-event-invalid-id";
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Outcome {
Ok,
Error,
}
pub struct NodeAuditRecord<'a> {
pub tenant: &'a TenantCtx,
pub flow_id: &'a str,
pub node_id: &'a str,
pub component_id: &'a str,
pub operation: &'a str,
pub session_id: &'a str,
pub duration_ms: u64,
pub outcome: Outcome,
pub error: Option<&'a str>,
}
pub fn audit_subject(tenant: &str, event: &str) -> String {
format!("audit.{tenant}.flow.{event}")
}
pub fn agent_audit_subject(tenant: &str, event: &str) -> String {
format!("audit.{tenant}.agent.{event}")
}
#[allow(clippy::too_many_arguments)]
fn base_event(
tenant: &TenantCtx,
ty: &str,
source: &str,
subject: String,
correlation_id: Option<String>,
payload: serde_json::Value,
time: DateTime<Utc>,
id: String,
) -> EventEnvelope {
let event_id = EventId::new(&id).unwrap_or_else(|_| {
EventId::new(FALLBACK_EVENT_ID).expect("static fallback event id is a valid EventId")
});
EventEnvelope {
id: event_id,
topic: "runner.flow".to_string(),
r#type: ty.to_string(),
source: source.to_string(),
tenant: tenant.clone(),
subject: Some(subject),
time,
correlation_id,
payload,
metadata: Default::default(),
}
}
pub fn build_audit_event(
rec: &NodeAuditRecord<'_>,
now: DateTime<Utc>,
id: String,
) -> EventEnvelope {
let (event_suffix, outcome_label) = match rec.outcome {
Outcome::Ok => ("node_end", "ok"),
Outcome::Error => ("node_error", "error"),
};
let event_type = format!("greentic.runner.flow.{event_suffix}");
let mut payload = json!({
"node_id": rec.node_id,
"component_id": rec.component_id,
"operation": rec.operation,
"outcome": outcome_label,
"duration_ms": rec.duration_ms,
});
if let Some(error) = rec.error
&& let Some(map) = payload.as_object_mut()
{
map.insert("error".to_string(), json!(error));
}
base_event(
rec.tenant,
&event_type,
"runner",
format!("flow:{}/node:{}", rec.flow_id, rec.node_id),
Some(rec.session_id.to_string()),
payload,
now,
id,
)
}
pub fn build_agent_audit_event(
tenant: &TenantCtx,
agent_id: &str,
session_id: &str,
event: &str,
payload: serde_json::Value,
now: DateTime<Utc>,
id: String,
) -> EventEnvelope {
let event_type = format!("greentic.runner.agent.{event}");
base_event(
tenant,
&event_type,
"runner",
format!("agent:{agent_id}"),
Some(session_id.to_string()),
payload,
now,
id,
)
}
#[cfg(test)]
mod tests {
use super::*;
use greentic_types::{EnvId, TenantId};
use serde_json::Value;
fn admin_decode(body: &Value) -> (String, String, String, String, String, String) {
let tenant = body
.get("tenant")
.and_then(|t| t.get("tenant"))
.and_then(Value::as_str)
.or_else(|| body.get("tenant").and_then(Value::as_str))
.unwrap()
.to_string();
let ty = body
.get("type")
.and_then(Value::as_str)
.unwrap()
.to_string();
let source = body
.get("source")
.and_then(Value::as_str)
.unwrap()
.to_string();
let subject = body
.get("subject")
.and_then(Value::as_str)
.unwrap()
.to_string();
let time = body
.get("time")
.and_then(Value::as_str)
.unwrap()
.to_string();
let corr = body
.get("correlation_id")
.and_then(Value::as_str)
.unwrap()
.to_string();
(tenant, ty, source, subject, time, corr)
}
fn tenant_ctx() -> TenantCtx {
TenantCtx::new(
EnvId::try_from("prod").expect("valid env id"),
TenantId::try_from("t1").expect("valid tenant id"),
)
}
#[test]
fn built_event_decodes_under_admin_contract() {
let tenant = tenant_ctx();
let rec = NodeAuditRecord {
tenant: &tenant,
flow_id: "f1",
node_id: "n1",
component_id: "greentic:http",
operation: "call",
session_id: "s1",
duration_ms: 12,
outcome: Outcome::Ok,
error: None,
};
let now = chrono::DateTime::parse_from_rfc3339("2026-07-03T00:00:00Z")
.unwrap()
.with_timezone(&chrono::Utc);
let env = build_audit_event(&rec, now, "id1".to_string());
let body = serde_json::to_value(&env).unwrap();
let (tenant, ty, source, subject, time, corr) = admin_decode(&body);
assert_eq!(tenant, "t1");
assert_eq!(ty, "greentic.runner.flow.node_end");
assert_eq!(source, "runner");
assert_eq!(subject, "flow:f1/node:n1");
assert!(time.starts_with("2026-07-03T00:00:00"));
assert_eq!(corr, "s1");
assert!(
body.get("payload")
.and_then(|p| p.get("duration_ms"))
.is_some()
);
}
#[test]
fn error_outcome_sets_type_and_payload_error() {
let tenant = tenant_ctx();
let rec = NodeAuditRecord {
tenant: &tenant,
flow_id: "f1",
node_id: "n1",
component_id: "greentic:http",
operation: "call",
session_id: "s1",
duration_ms: 5,
outcome: Outcome::Error,
error: Some("boom"),
};
let now = chrono::DateTime::parse_from_rfc3339("2026-07-03T00:00:00Z")
.unwrap()
.with_timezone(&chrono::Utc);
let env = build_audit_event(&rec, now, "id2".to_string());
let body = serde_json::to_value(&env).unwrap();
let (_, ty, _, _, _, _) = admin_decode(&body);
assert_eq!(ty, "greentic.runner.flow.node_error");
let payload = body.get("payload").unwrap();
assert_eq!(
payload.get("outcome").and_then(Value::as_str),
Some("error")
);
assert_eq!(payload.get("error").and_then(Value::as_str), Some("boom"));
}
#[test]
fn ok_outcome_never_includes_error_key() {
let tenant = tenant_ctx();
let rec = NodeAuditRecord {
tenant: &tenant,
flow_id: "f1",
node_id: "n1",
component_id: "greentic:http",
operation: "call",
session_id: "s1",
duration_ms: 5,
outcome: Outcome::Ok,
error: None,
};
let now = Utc::now();
let env = build_audit_event(&rec, now, "id3".to_string());
let body = serde_json::to_value(&env).unwrap();
assert!(body.get("payload").unwrap().get("error").is_none());
}
#[test]
fn subject_is_well_formed() {
assert_eq!(audit_subject("t1", "node_end"), "audit.t1.flow.node_end");
}
#[test]
fn agent_subject_is_well_formed() {
assert_eq!(
agent_audit_subject("t1", "tool_call"),
"audit.t1.agent.tool_call"
);
}
#[test]
fn agent_event_decodes_under_admin_contract() {
let tenant = tenant_ctx();
let now = chrono::DateTime::parse_from_rfc3339("2026-07-03T00:00:00Z")
.unwrap()
.with_timezone(&chrono::Utc);
let env = build_agent_audit_event(
&tenant,
"a1",
"s1",
"tool_call",
json!({"tool": "http", "call_id": "c1"}),
now,
"id1".to_string(),
);
let body = serde_json::to_value(&env).unwrap();
let (tenant, ty, source, subject, time, corr) = admin_decode(&body);
assert_eq!(tenant, "t1");
assert_eq!(ty, "greentic.runner.agent.tool_call");
assert_eq!(source, "runner");
assert_eq!(subject, "agent:a1");
assert!(time.starts_with("2026-07-03T00:00:00"));
assert_eq!(corr, "s1");
assert_eq!(
body.get("payload")
.and_then(|p| p.get("tool"))
.and_then(Value::as_str),
Some("http")
);
}
}