use super::manager::ExtensionProcess;
use everruns_core::events::{
Event, OUTPUT_MESSAGE_DELTA, REASON_THINKING_DELTA, TOOL_OUTPUT_DELTA, TOOL_PROGRESS,
};
use everruns_core::typed_id::SessionId;
use serde_json::Value;
use std::sync::Arc;
use tokio::sync::broadcast;
fn is_forwardable(event_type: &str) -> bool {
!matches!(
event_type,
OUTPUT_MESSAGE_DELTA | REASON_THINKING_DELTA | TOOL_OUTPUT_DELTA | TOOL_PROGRESS
)
}
pub(crate) fn trace_event_params(event: &Event) -> Value {
serde_json::json!({
"event_type": event.event_type,
"id": event.id.to_string(),
"ts": event.ts.to_rfc3339(),
"session_id": event.session_id.to_string(),
"context": serde_json::to_value(&event.context).unwrap_or(Value::Null),
"data": serde_json::to_value(&event.data).unwrap_or(Value::Null),
})
}
pub(crate) struct TraceForwarder {
name: String,
process: Arc<ExtensionProcess>,
}
impl TraceForwarder {
pub(crate) fn for_capability(
capability: &super::capability::ExtensionCapability,
config: &serde_json::Value,
) -> Option<Self> {
use everruns_core::capabilities::Capability;
capability.trace_process(config).map(|process| Self {
name: capability.name().to_string(),
process,
})
}
pub(crate) fn start(self, session_id: SessionId, events: broadcast::Receiver<Event>) {
spawn_forwarder(self.name, self.process, session_id, events);
}
}
pub(crate) fn spawn_forwarder(
ext_name: String,
process: Arc<ExtensionProcess>,
session_id: SessionId,
mut events: broadcast::Receiver<Event>,
) {
tokio::spawn(async move {
loop {
match events.recv().await {
Ok(event) if event.session_id == session_id => {
if !is_forwardable(&event.event_type) {
continue;
}
if let Err(err) = process.send_trace_event(trace_event_params(&event)).await {
tracing::warn!(
target: "yolop::ext", ext = %ext_name,
"trace/event forward failed: {err}"
);
}
}
Ok(_) => {}
Err(broadcast::error::RecvError::Lagged(n)) => {
tracing::debug!(
target: "yolop::ext", ext = %ext_name,
"trace forwarder lagged, dropped {n} events"
);
}
Err(broadcast::error::RecvError::Closed) => break,
}
}
});
}
#[cfg(test)]
mod tests {
use super::*;
use everruns_core::events::{EventContext, EventRequest};
#[test]
fn deltas_are_dropped_lifecycle_is_kept() {
assert!(is_forwardable("turn.started"));
assert!(is_forwardable("tool.completed"));
assert!(is_forwardable("llm.generation"));
assert!(!is_forwardable(OUTPUT_MESSAGE_DELTA));
assert!(!is_forwardable(TOOL_OUTPUT_DELTA));
assert!(!is_forwardable(TOOL_PROGRESS));
}
#[test]
fn params_carry_type_ids_and_verbatim_payload() {
let session = SessionId::new();
let request = EventRequest {
event_type: "tool.completed".into(),
ts: chrono::Utc::now(),
session_id: session,
context: EventContext::empty(),
data: everruns_core::events::EventData::unsupported(
"tool.completed".into(),
serde_json::json!({ "tool_name": "bash" }),
),
metadata: None,
tags: None,
};
let event = request.into_event(everruns_core::typed_id::EventId::new(), 1);
let params = trace_event_params(&event);
assert_eq!(params["event_type"], "tool.completed");
assert_eq!(params["session_id"], session.to_string());
assert_eq!(params["id"], event.id.to_string());
assert!(params["ts"].as_str().unwrap().contains('T'));
assert!(params["context"].is_object(), "{params}");
}
}