roder-core 0.1.6

Agentic software development tools and SDKs for Roder.
Documentation
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;
        }
    }
}

/// Bus-backed [`AgentSwarmProgressSink`]: the `agent_swarm` tool publishes a
/// live tick through this, which the runtime turns into an `AgentSwarmProgress`
/// event (and persists it like other thread events).
#[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:?}"),
        }
    }
}