use serde::{Deserialize, Serialize};
use time::OffsetDateTime;
use crate::journal::{LfEvent, LfEventType, LfNode};
use crate::lfd::id::LfdId;
use crate::lfd::types::{AttentionItem, Session};
use crate::provider_auth::Provider;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum Event {
Connected {
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "auth.flow_started")]
AuthFlowStarted {
provider: Provider,
verification_uri: String,
#[serde(skip_serializing_if = "Option::is_none")]
verification_uri_complete: Option<String>,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "auth.connected")]
AuthConnected {
provider: Provider,
#[serde(skip_serializing_if = "Option::is_none")]
login: Option<String>,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "auth.failed")]
AuthFailed {
provider: Provider,
error: String,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "auth.disconnected")]
AuthDisconnected {
provider: Provider,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "auth.token_refreshed")]
AuthTokenRefreshed {
provider: Provider,
#[serde(skip_serializing_if = "Option::is_none")]
login: Option<String>,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "auth.refresh_failed")]
AuthRefreshFailed {
provider: Provider,
reason: String,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "auth.refresh_required")]
AuthRefreshRequired {
provider: Provider,
reason: String,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
WaveCreated {
wave_id: LfdId,
name: String,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
WaveUpdated {
wave_id: LfdId,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
WaveDeleted {
wave_id: LfdId,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
WaveStarted {
wave_id: LfdId,
run_id: LfdId,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
WaveStopped {
wave_id: LfdId,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
WaveWaiting {
wave_id: LfdId,
run_id: LfdId,
step: String,
#[serde(skip_serializing_if = "Option::is_none")]
session_id: Option<LfdId>,
#[serde(skip_serializing_if = "Option::is_none")]
initial_user_message: Option<String>,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "run.started")]
RunStarted {
run_id: LfdId,
wave_name: String,
worktree: String,
command: Vec<String>,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "flow.started")]
FlowStarted {
run_id: LfdId,
flow: String,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "flow.completed")]
FlowCompleted {
run_id: LfdId,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "flow.errored")]
FlowErrored {
run_id: LfdId,
error: String,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "flow.escalated")]
FlowEscalated {
run_id: LfdId,
signal: String,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "step.started")]
StepStarted {
run_id: LfdId,
step: String,
index: u32,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "step.completed")]
StepCompleted {
run_id: LfdId,
step: String,
index: u32,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "step.errored")]
StepErrored {
run_id: LfdId,
step: String,
index: u32,
error: String,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "step.escalated")]
StepEscalated {
run_id: LfdId,
step: String,
index: u32,
signal: String,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "run.completed")]
RunCompleted {
run_id: LfdId,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "run.errored")]
RunErrored {
run_id: LfdId,
error: String,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
#[serde(rename = "run.escalated")]
RunEscalated {
run_id: LfdId,
signal: String,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
AttentionCreated {
item: AttentionItem,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
AttentionUpdated {
item: AttentionItem,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
AttentionResolved {
item: AttentionItem,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
SessionCreated {
session: Session,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
SessionUpdated {
session: Session,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
OutputLine {
wave_id: LfdId,
agent_id: LfdId,
text: String,
#[serde(with = "time::serde::rfc3339")]
timestamp: OffsetDateTime,
},
}
impl Event {
pub fn now() -> OffsetDateTime {
OffsetDateTime::now_utc()
}
pub fn attention_created(item: AttentionItem) -> Self {
Self::AttentionCreated {
item,
timestamp: Self::now(),
}
}
pub fn attention_updated(item: AttentionItem) -> Self {
Self::AttentionUpdated {
item,
timestamp: Self::now(),
}
}
pub fn attention_resolved(item: AttentionItem) -> Self {
Self::AttentionResolved {
item,
timestamp: Self::now(),
}
}
pub fn session_created(session: Session) -> Self {
Self::SessionCreated {
session,
timestamp: Self::now(),
}
}
pub fn session_updated(session: Session) -> Self {
Self::SessionUpdated {
session,
timestamp: Self::now(),
}
}
pub fn wave_created(wave_id: LfdId, name: String) -> Self {
Self::WaveCreated {
wave_id,
name,
timestamp: Self::now(),
}
}
pub fn auth_flow_started(
provider: Provider,
verification_uri: String,
verification_uri_complete: Option<String>,
) -> Self {
Self::AuthFlowStarted {
provider,
verification_uri,
verification_uri_complete,
timestamp: Self::now(),
}
}
pub fn auth_connected(provider: Provider, login: Option<String>) -> Self {
Self::AuthConnected {
provider,
login,
timestamp: Self::now(),
}
}
pub fn auth_failed(provider: Provider, error: String) -> Self {
Self::AuthFailed {
provider,
error,
timestamp: Self::now(),
}
}
pub fn auth_disconnected(provider: Provider) -> Self {
Self::AuthDisconnected {
provider,
timestamp: Self::now(),
}
}
pub fn auth_token_refreshed(provider: Provider, login: Option<String>) -> Self {
Self::AuthTokenRefreshed {
provider,
login,
timestamp: Self::now(),
}
}
pub fn auth_refresh_failed(provider: Provider, reason: String) -> Self {
Self::AuthRefreshFailed {
provider,
reason,
timestamp: Self::now(),
}
}
pub fn auth_refresh_required(provider: Provider, reason: String) -> Self {
Self::AuthRefreshRequired {
provider,
reason,
timestamp: Self::now(),
}
}
pub fn wave_updated(wave_id: LfdId) -> Self {
Self::WaveUpdated {
wave_id,
timestamp: Self::now(),
}
}
pub fn wave_deleted(wave_id: LfdId) -> Self {
Self::WaveDeleted {
wave_id,
timestamp: Self::now(),
}
}
pub fn wave_started(wave_id: LfdId, run_id: LfdId) -> Self {
Self::WaveStarted {
wave_id,
run_id,
timestamp: Self::now(),
}
}
pub fn wave_stopped(wave_id: LfdId) -> Self {
Self::WaveStopped {
wave_id,
timestamp: Self::now(),
}
}
pub fn wave_waiting(
wave_id: LfdId,
run_id: LfdId,
step: String,
session_id: Option<LfdId>,
initial_user_message: Option<String>,
) -> Self {
Self::WaveWaiting {
wave_id,
run_id,
step,
session_id,
initial_user_message,
timestamp: Self::now(),
}
}
}
impl From<LfEvent> for Event {
fn from(event: LfEvent) -> Self {
match (event.node, event.event) {
(LfNode::Run, LfEventType::Started) => Self::RunStarted {
run_id: event.run_id,
wave_name: event
.wave_name
.expect("run.started journal events must include wave_name"),
worktree: event
.worktree
.expect("run.started journal events must include worktree"),
command: event
.command
.expect("run.started journal events must include command"),
timestamp: event.ts,
},
(LfNode::Flow, LfEventType::Started) => Self::FlowStarted {
run_id: event.run_id,
flow: event
.flow
.expect("flow.started journal events must include flow"),
timestamp: event.ts,
},
(LfNode::Flow, LfEventType::Completed) => Self::FlowCompleted {
run_id: event.run_id,
timestamp: event.ts,
},
(LfNode::Flow, LfEventType::Errored) => Self::FlowErrored {
run_id: event.run_id,
error: event
.error
.expect("flow.errored journal events must include error"),
timestamp: event.ts,
},
(LfNode::Flow, LfEventType::Escalated) => Self::FlowEscalated {
run_id: event.run_id,
signal: event
.signal
.expect("flow.escalated journal events must include signal"),
timestamp: event.ts,
},
(LfNode::Step, LfEventType::Started) => Self::StepStarted {
run_id: event.run_id,
step: event
.step
.expect("step.started journal events must include step"),
index: event
.index
.expect("step.started journal events must include index"),
timestamp: event.ts,
},
(LfNode::Step, LfEventType::Completed) => Self::StepCompleted {
run_id: event.run_id,
step: event
.step
.expect("step.completed journal events must include step"),
index: event
.index
.expect("step.completed journal events must include index"),
timestamp: event.ts,
},
(LfNode::Step, LfEventType::Errored) => Self::StepErrored {
run_id: event.run_id,
step: event
.step
.expect("step.errored journal events must include step"),
index: event
.index
.expect("step.errored journal events must include index"),
error: event
.error
.expect("step.errored journal events must include error"),
timestamp: event.ts,
},
(LfNode::Step, LfEventType::Escalated) => Self::StepEscalated {
run_id: event.run_id,
step: event
.step
.expect("step.escalated journal events must include step"),
index: event
.index
.expect("step.escalated journal events must include index"),
signal: event
.signal
.expect("step.escalated journal events must include signal"),
timestamp: event.ts,
},
(LfNode::Run, LfEventType::Completed) => Self::RunCompleted {
run_id: event.run_id,
timestamp: event.ts,
},
(LfNode::Run, LfEventType::Errored) => Self::RunErrored {
run_id: event.run_id,
error: event
.error
.expect("run.errored journal events must include error"),
timestamp: event.ts,
},
(LfNode::Run, LfEventType::Escalated) => Self::RunEscalated {
run_id: event.run_id,
signal: event
.signal
.expect("run.escalated journal events must include signal"),
timestamp: event.ts,
},
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn test_id(s: &str) -> LfdId {
LfdId::from_raw(s)
}
#[test]
fn wave_waiting_serializes_correctly() {
let event = Event::wave_waiting(
test_id("wave-1"),
test_id("run-1"),
"implement".to_string(),
Some(test_id("session-1")),
Some("Start with user prompt".to_string()),
);
let json = serde_json::to_value(&event).unwrap();
assert_eq!(json["type"], "wave_waiting");
assert_eq!(json["wave_id"], "wave-1");
assert_eq!(json["run_id"], "run-1");
assert_eq!(json["step"], "implement");
assert_eq!(json["session_id"], "session-1");
assert_eq!(json["initial_user_message"], "Start with user prompt");
assert!(json["timestamp"].is_string());
}
#[test]
fn wave_waiting_omits_session_id_when_absent() {
let event = Event::wave_waiting(
test_id("wave-1"),
test_id("run-1"),
"implement".to_string(),
None,
None,
);
let json = serde_json::to_value(&event).unwrap();
assert!(json.get("session_id").is_none());
assert!(json.get("initial_user_message").is_none());
}
#[test]
fn event_roundtrips_through_json() {
let id = || LfdId::new();
let timestamp = Event::now();
let events = vec![
Event::wave_waiting(id(), id(), "step".to_string(), None, None),
Event::RunStarted {
run_id: id(),
wave_name: "wave".to_string(),
worktree: "/tmp/wt".to_string(),
command: vec!["lf".to_string(), "build".to_string()],
timestamp,
},
Event::FlowStarted {
run_id: id(),
flow: "build".to_string(),
timestamp,
},
Event::StepStarted {
run_id: id(),
step: "implement".to_string(),
index: 0,
timestamp,
},
Event::StepCompleted {
run_id: id(),
step: "implement".to_string(),
index: 0,
timestamp,
},
Event::StepErrored {
run_id: id(),
step: "review".to_string(),
index: 1,
error: "boom".to_string(),
timestamp,
},
Event::FlowCompleted {
run_id: id(),
timestamp,
},
Event::RunCompleted {
run_id: id(),
timestamp,
},
Event::RunErrored {
run_id: id(),
error: "boom".to_string(),
timestamp,
},
Event::auth_connected(Provider::GitHub, Some("jackdanger".to_string())),
Event::auth_token_refreshed(Provider::GitHub, Some("jackdanger".to_string())),
Event::auth_refresh_failed(Provider::GitHub, "refresh failed".to_string()),
Event::auth_refresh_required(Provider::Claude, "user must re-authenticate".to_string()),
];
for event in events {
let json = serde_json::to_string(&event).unwrap();
let parsed: Event = serde_json::from_str(&json).unwrap();
let json2 = serde_json::to_string(&parsed).unwrap();
assert_eq!(json, json2);
}
}
#[test]
fn auth_events_use_dotted_type_names() {
let event = Event::auth_connected(Provider::GitHub, Some("jackdanger".to_string()));
let json = serde_json::to_value(&event).expect("serialize");
assert_eq!(json["type"], "auth.connected");
assert_eq!(json["provider"], "github");
assert_eq!(json["login"], "jackdanger");
let refreshed = Event::auth_token_refreshed(Provider::Claude, None);
let refreshed_json = serde_json::to_value(&refreshed).expect("serialize refreshed");
assert_eq!(refreshed_json["type"], "auth.token_refreshed");
assert_eq!(refreshed_json["provider"], "claude");
assert!(refreshed_json.get("login").is_none());
let failed = Event::auth_refresh_failed(Provider::Codex, "timed out".to_string());
let failed_json = serde_json::to_value(&failed).expect("serialize failure");
assert_eq!(failed_json["type"], "auth.refresh_failed");
assert_eq!(failed_json["provider"], "codex");
assert_eq!(failed_json["reason"], "timed out");
let required =
Event::auth_refresh_required(Provider::Claude, "user must re-authenticate".into());
let required_json = serde_json::to_value(&required).expect("serialize required");
assert_eq!(required_json["type"], "auth.refresh_required");
assert_eq!(required_json["provider"], "claude");
assert_eq!(required_json["reason"], "user must re-authenticate");
}
#[test]
fn wave_event_can_be_enriched_with_extra_field() {
let event = Event::wave_updated(test_id("wave-1"));
let mut base = serde_json::to_value(&event).unwrap();
let wave_data = serde_json::json!({ "id": "wave-1", "name": "test-wave" });
if let serde_json::Value::Object(ref mut map) = base {
map.insert("wave".to_string(), wave_data);
}
let json_str = serde_json::to_string(&base).unwrap();
let parsed: serde_json::Value = serde_json::from_str(&json_str).unwrap();
assert_eq!(parsed["type"], "wave_updated");
assert_eq!(parsed["wave_id"], "wave-1");
assert!(parsed["timestamp"].is_string());
assert_eq!(parsed["wave"]["id"], "wave-1");
assert_eq!(parsed["wave"]["name"], "test-wave");
}
#[test]
fn runtime_events_use_dotted_type_names() {
let timestamp = Event::now();
let event = Event::RunStarted {
run_id: test_id("run-1"),
wave_name: "lfd".to_string(),
worktree: "/tmp/repo.lfd".to_string(),
command: vec!["lf".to_string(), "build".to_string()],
timestamp,
};
let json = serde_json::to_value(&event).expect("serialize");
assert_eq!(json["type"], "run.started");
assert_eq!(json["run_id"], "run-1");
assert_eq!(json["wave_name"], "lfd");
assert!(json.get("flow").is_none());
let flow = Event::FlowStarted {
run_id: test_id("run-1"),
flow: "build".to_string(),
timestamp,
};
let flow_json = serde_json::to_value(&flow).expect("serialize flow");
assert_eq!(flow_json["type"], "flow.started");
assert_eq!(flow_json["flow"], "build");
let step = Event::StepCompleted {
run_id: test_id("run-1"),
step: "implement".to_string(),
index: 1,
timestamp,
};
let step_json = serde_json::to_value(&step).expect("serialize step");
assert_eq!(step_json["type"], "step.completed");
assert!(step_json.get("exit_code").is_none());
}
}