use async_trait::async_trait;
use serde::Serialize;
use crate::provision::{ProvisionOutcome, ProvisionResult};
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ProvisionEvent {
Succeeded { runner_name: String },
Failed {
runner_name: String,
error: String,
diagnostics: serde_json::Map<String, serde_json::Value>,
},
#[cfg_attr(not(target_os = "linux"), allow(dead_code))]
AtCapacity {
runner_name: String,
code: String,
message: String,
retry_after_secs: u64,
},
}
impl ProvisionEvent {
pub fn is_success(&self) -> bool {
matches!(self, ProvisionEvent::Succeeded { .. })
}
}
impl From<ProvisionResult> for ProvisionEvent {
fn from(pr: ProvisionResult) -> Self {
match pr.outcome {
ProvisionOutcome::Success => ProvisionEvent::Succeeded {
runner_name: pr.runner_name,
},
ProvisionOutcome::Failed { error, diagnostics } => ProvisionEvent::Failed {
runner_name: pr.runner_name,
error,
diagnostics,
},
ProvisionOutcome::HostFull {
code,
message,
retry_after_secs,
} => ProvisionEvent::AtCapacity {
runner_name: pr.runner_name,
code,
message,
retry_after_secs,
},
}
}
}
#[derive(Debug, Clone, Serialize, PartialEq)]
pub struct AgentEvent {
pub runner_name: String,
pub kind: EventKind,
pub severity: Severity,
pub title: String,
pub message: String,
#[serde(skip_serializing_if = "serde_json::Map::is_empty")]
pub metadata: serde_json::Map<String, serde_json::Value>,
}
#[derive(Debug, Clone, Copy, Serialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum EventKind {
#[allow(dead_code)]
ProvisionSucceeded,
ProvisionFailed,
HostAtCapacity,
}
#[derive(Debug, Clone, Copy, Serialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum Severity {
#[allow(dead_code)]
Info,
Warning,
Error,
}
pub fn to_agent_event(ev: &ProvisionEvent, attempt: u32) -> Option<AgentEvent> {
match ev {
ProvisionEvent::Succeeded { .. } => None,
ProvisionEvent::Failed {
runner_name,
error,
diagnostics,
} => {
let mut metadata = diagnostics.clone();
metadata.insert("attempt".into(), serde_json::Value::from(attempt));
metadata.insert("error".into(), serde_json::Value::from(error.clone()));
Some(AgentEvent {
runner_name: runner_name.clone(),
kind: EventKind::ProvisionFailed,
severity: Severity::Error,
title: "Provision failed".into(),
message: error.clone(),
metadata,
})
}
ProvisionEvent::AtCapacity {
runner_name,
code,
message,
retry_after_secs,
} => {
let mut metadata = serde_json::Map::new();
metadata.insert("code".into(), serde_json::Value::from(code.clone()));
metadata.insert(
"retry_after_secs".into(),
serde_json::Value::from(*retry_after_secs),
);
Some(AgentEvent {
runner_name: runner_name.clone(),
kind: EventKind::HostAtCapacity,
severity: Severity::Warning,
title: "Host at capacity".into(),
message: message.clone(),
metadata,
})
}
}
}
#[async_trait]
pub trait ProvisionReporter: Send + Sync {
async fn report(&self, event: ProvisionEvent);
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Mutex;
pub struct RecordingReporter {
events: Mutex<Vec<ProvisionEvent>>,
}
impl RecordingReporter {
pub fn new() -> Self {
Self {
events: Mutex::new(Vec::new()),
}
}
pub fn events(&self) -> Vec<ProvisionEvent> {
self.events.lock().unwrap().clone()
}
}
#[async_trait]
impl ProvisionReporter for RecordingReporter {
async fn report(&self, event: ProvisionEvent) {
self.events.lock().unwrap().push(event);
}
}
fn pr_success(name: &str) -> ProvisionResult {
ProvisionResult {
runner_name: name.into(),
executor_kind: None,
outcome: ProvisionOutcome::Success,
}
}
fn pr_failed(name: &str, err: &str) -> ProvisionResult {
ProvisionResult {
runner_name: name.into(),
executor_kind: None,
outcome: ProvisionOutcome::failed(err),
}
}
fn pr_host_full(name: &str, code: &str) -> ProvisionResult {
ProvisionResult {
runner_name: name.into(),
executor_kind: None,
outcome: ProvisionOutcome::HostFull {
code: code.into(),
message: format!("{code} test"),
retry_after_secs: 10,
},
}
}
#[test]
fn lifts_success_outcome_to_succeeded_event() {
let ev: ProvisionEvent = pr_success("r1").into();
assert!(ev.is_success());
}
#[test]
fn lifts_failed_outcome_to_failed_event() {
let ev: ProvisionEvent = pr_failed("r1", "boom").into();
assert!(matches!(ev, ProvisionEvent::Failed { .. }));
}
#[test]
fn lifts_host_full_outcome_to_at_capacity_event() {
let ev: ProvisionEvent = pr_host_full("r1", "CPU_EXHAUSTED").into();
assert!(matches!(ev, ProvisionEvent::AtCapacity { .. }));
}
#[test]
fn success_emits_no_wire_event() {
let ev = ProvisionEvent::Succeeded {
runner_name: "r1".into(),
};
assert!(to_agent_event(&ev, 0).is_none());
}
#[test]
fn failed_emits_provision_failed_event_with_attempt_and_error() {
let ev = ProvisionEvent::Failed {
runner_name: "r1".into(),
error: "boom".into(),
diagnostics: serde_json::Map::new(),
};
let agent_event = to_agent_event(&ev, 3).expect("should emit");
assert_eq!(agent_event.kind, EventKind::ProvisionFailed);
assert_eq!(agent_event.severity, Severity::Error);
assert_eq!(agent_event.message, "boom");
assert_eq!(
agent_event.metadata.get("attempt"),
Some(&serde_json::Value::from(3))
);
assert_eq!(
agent_event.metadata.get("error"),
Some(&serde_json::Value::from("boom"))
);
}
#[test]
fn failed_event_carries_caller_supplied_diagnostics_into_metadata() {
let mut diag = serde_json::Map::new();
diag.insert("elapsed_ms".into(), serde_json::Value::from(37386));
diag.insert("partial_stdout".into(), serde_json::Value::from("18066\n"));
diag.insert(
"vm_diagnostics".into(),
serde_json::Value::from("=== ps ===\n..."),
);
let ev = ProvisionEvent::Failed {
runner_name: "r1".into(),
error: "ssh exec timed out".into(),
diagnostics: diag,
};
let agent_event = to_agent_event(&ev, 1).expect("should emit");
assert_eq!(
agent_event.metadata.get("elapsed_ms"),
Some(&serde_json::Value::from(37386))
);
assert_eq!(
agent_event.metadata.get("partial_stdout"),
Some(&serde_json::Value::from("18066\n"))
);
assert!(agent_event
.metadata
.get("vm_diagnostics")
.and_then(|v| v.as_str())
.unwrap()
.contains("=== ps ==="));
}
#[test]
fn canonical_metadata_keys_override_caller_diagnostics() {
let mut diag = serde_json::Map::new();
diag.insert("attempt".into(), serde_json::Value::from(99));
diag.insert("error".into(), serde_json::Value::from("wrong"));
let ev = ProvisionEvent::Failed {
runner_name: "r1".into(),
error: "real error".into(),
diagnostics: diag,
};
let agent_event = to_agent_event(&ev, 2).expect("should emit");
assert_eq!(
agent_event.metadata.get("attempt"),
Some(&serde_json::Value::from(2))
);
assert_eq!(
agent_event.metadata.get("error"),
Some(&serde_json::Value::from("real error"))
);
}
#[test]
fn at_capacity_emits_host_at_capacity_event_with_code_and_retry_after() {
let ev = ProvisionEvent::AtCapacity {
runner_name: "r1".into(),
code: "CPU_EXHAUSTED".into(),
message: "CPU exhausted: ...".into(),
retry_after_secs: 10,
};
let agent_event = to_agent_event(&ev, 0).expect("should emit");
assert_eq!(agent_event.kind, EventKind::HostAtCapacity);
assert_eq!(agent_event.severity, Severity::Warning);
assert_eq!(agent_event.title, "Host at capacity");
assert_eq!(
agent_event.metadata.get("code"),
Some(&serde_json::Value::from("CPU_EXHAUSTED"))
);
assert_eq!(
agent_event.metadata.get("retry_after_secs"),
Some(&serde_json::Value::from(10u64))
);
}
#[test]
fn serializes_kind_and_severity_as_snake_lowercase() {
let ev = ProvisionEvent::AtCapacity {
runner_name: "r1".into(),
code: "CPU_EXHAUSTED".into(),
message: "x".into(),
retry_after_secs: 5,
};
let agent_event = to_agent_event(&ev, 0).unwrap();
let json = serde_json::to_value(&agent_event).unwrap();
assert_eq!(json["kind"], "host_at_capacity");
assert_eq!(json["severity"], "warning");
assert_eq!(json["title"], "Host at capacity");
assert_eq!(json["metadata"]["code"], "CPU_EXHAUSTED");
assert_eq!(json["metadata"]["retry_after_secs"], 5);
}
#[tokio::test]
async fn recording_reporter_captures_events_in_order() {
let r = RecordingReporter::new();
r.report(pr_success("r1").into()).await;
r.report(pr_failed("r2", "x").into()).await;
let evs = r.events();
assert_eq!(evs.len(), 2);
assert!(matches!(evs[0], ProvisionEvent::Succeeded { .. }));
assert!(matches!(evs[1], ProvisionEvent::Failed { .. }));
}
}