use std::sync::Arc;
use roder_api::events::{
RoderEvent, SubagentTraceCompleted, SubagentTraceCreated, SubagentTraceDeltaEvent,
SubagentTraceFailed, SubagentTraceStatusChanged,
};
use roder_api::subagents::{AgentSwarmProgress, AgentSwarmProgressSink, AgentSwarmProgressSnapshot};
use roder_api::thread::ThreadStore;
use roder_api::trace::{
ParentTurnRef, SubagentTraceDelta, SubagentTraceId, SubagentTraceSink, SubagentTraceStatus,
SubagentTraceSummary,
};
use time::OffsetDateTime;
use crate::bus::EventBus;
#[derive(Clone)]
pub(crate) struct RuntimeSubagentTraceSink {
bus: EventBus,
thread_store: Option<Arc<dyn ThreadStore>>,
}
impl RuntimeSubagentTraceSink {
pub(crate) fn new(bus: EventBus, thread_store: Option<Arc<dyn ThreadStore>>) -> Self {
Self { bus, thread_store }
}
async fn emit(&self, event: RoderEvent) {
let envelope = self.bus.emit(event);
if let (Some(store), Some(thread_id)) = (&self.thread_store, envelope.thread_id.as_ref()) {
let _ = store.append_event(thread_id, &envelope).await;
}
}
}
#[derive(Clone)]
pub(crate) struct RuntimeAgentSwarmProgressSink {
bus: EventBus,
thread_store: Option<Arc<dyn ThreadStore>>,
}
impl RuntimeAgentSwarmProgressSink {
pub(crate) fn new(bus: EventBus, thread_store: Option<Arc<dyn ThreadStore>>) -> Self {
Self { bus, thread_store }
}
}
#[async_trait::async_trait]
impl AgentSwarmProgressSink for RuntimeAgentSwarmProgressSink {
async fn emit_progress(
&self,
thread_id: &str,
turn_id: &str,
tool_id: &str,
snapshot: AgentSwarmProgressSnapshot,
) {
let envelope = self.bus.emit(RoderEvent::AgentSwarmProgress(AgentSwarmProgress {
thread_id: thread_id.to_string(),
turn_id: turn_id.to_string(),
tool_id: tool_id.to_string(),
snapshot,
timestamp: OffsetDateTime::now_utc(),
}));
if let (Some(store), Some(thread_id)) = (&self.thread_store, envelope.thread_id.as_ref()) {
let _ = store.append_event(thread_id, &envelope).await;
}
}
}
#[async_trait::async_trait]
impl SubagentTraceSink for RuntimeSubagentTraceSink {
async fn trace_created(&self, summary: SubagentTraceSummary) {
self.emit(RoderEvent::SubagentTraceCreated(SubagentTraceCreated {
summary,
timestamp: OffsetDateTime::now_utc(),
}))
.await;
}
async fn trace_delta(&self, delta: SubagentTraceDelta) {
self.emit(RoderEvent::SubagentTraceDelta(SubagentTraceDeltaEvent {
delta,
timestamp: OffsetDateTime::now_utc(),
}))
.await;
}
async fn trace_status_changed(
&self,
trace_id: SubagentTraceId,
parent: ParentTurnRef,
status: SubagentTraceStatus,
detail: Option<String>,
) {
self.emit(RoderEvent::SubagentTraceStatusChanged(
SubagentTraceStatusChanged {
trace_id,
parent,
status,
detail,
timestamp: OffsetDateTime::now_utc(),
},
))
.await;
}
async fn trace_completed(&self, summary: SubagentTraceSummary) {
self.emit(RoderEvent::SubagentTraceCompleted(SubagentTraceCompleted {
summary,
timestamp: OffsetDateTime::now_utc(),
}))
.await;
}
async fn trace_failed(&self, summary: SubagentTraceSummary, error: String) {
self.emit(RoderEvent::SubagentTraceFailed(SubagentTraceFailed {
summary,
error,
timestamp: OffsetDateTime::now_utc(),
}))
.await;
}
}
#[cfg(test)]
mod tests {
use super::*;
use roder_api::subagents::{SubagentExitReason, SubagentLane};
use roder_api::trace::{SubagentDestination, SubagentDestinationKind};
#[tokio::test]
async fn subagent_trace_sink_emits_parent_turn_envelope() {
let bus = EventBus::new(16);
let sink = RuntimeSubagentTraceSink::new(bus.clone(), None);
let mut events = bus.subscribe();
sink.trace_created(SubagentTraceSummary {
trace_id: "trace-1".to_string(),
parent: ParentTurnRef {
thread_id: "parent-thread".to_string(),
turn_id: "parent-turn".to_string(),
},
child_thread_id: "child-thread".to_string(),
child_turn_id: "child-turn".to_string(),
title: "Inspect".to_string(),
role: "explore".to_string(),
model: Some("mock".to_string()),
lane: Some(SubagentLane::Scout),
status: SubagentTraceStatus::Queued,
elapsed_ms: 0,
usage: None,
destination: Some(SubagentDestination {
kind: SubagentDestinationKind::InProcess,
label: "in-process".to_string(),
path: None,
provider_id: None,
destination_id: None,
}),
latest_activity: Some("queued".to_string()),
error_summary: None,
exit_reason: Some(SubagentExitReason::Completed),
})
.await;
let envelope = events.recv().await.unwrap();
assert_eq!(envelope.kind, "turn/subagentTraceCreated");
assert_eq!(envelope.thread_id.as_deref(), Some("parent-thread"));
assert_eq!(envelope.turn_id.as_deref(), Some("parent-turn"));
match envelope.event {
RoderEvent::SubagentTraceCreated(event) => {
assert_eq!(event.summary.lane, Some(SubagentLane::Scout));
assert_eq!(
event.summary.exit_reason,
Some(SubagentExitReason::Completed)
);
}
other => panic!("unexpected event: {other:?}"),
}
}
#[tokio::test]
async fn swarm_progress_sink_emits_progress_event() {
let bus = EventBus::new(16);
let sink = RuntimeAgentSwarmProgressSink::new(bus.clone(), None);
let mut events = bus.subscribe();
sink.emit_progress(
"thread-1",
"turn-1",
"swarm-1",
AgentSwarmProgressSnapshot {
total: 3,
completed: 1,
failed: 1,
aborted: 0,
},
)
.await;
let envelope = events.recv().await.unwrap();
assert_eq!(envelope.kind, "agent_swarm.progress");
assert_eq!(envelope.thread_id.as_deref(), Some("thread-1"));
assert_eq!(envelope.turn_id.as_deref(), Some("turn-1"));
match envelope.event {
RoderEvent::AgentSwarmProgress(event) => {
assert_eq!(event.tool_id, "swarm-1");
assert_eq!(event.snapshot.total, 3);
assert_eq!(event.snapshot.resolved(), 2);
}
other => panic!("unexpected event: {other:?}"),
}
}
}