use std::collections::HashSet;
use everruns_core::events::{Event as RuntimeEvent, EventData, ToolCompletedData};
use everruns_core::message::{ContentPart, MessageRole};
use serde_json::Value;
use super::protocol::{
ContentBlock, PlanEntry, PlanPriority, PlanStatus, SessionUpdate, ToolCallContent,
ToolCallStatus, ToolKind,
};
const WRITE_TODOS: &str = "write_todos";
#[derive(Default)]
pub struct Translator {
current_message_streamed: bool,
seen: HashSet<String>,
}
impl Translator {
pub fn new() -> Self {
Self::default()
}
pub fn on_event(&mut self, event: &RuntimeEvent) -> Vec<SessionUpdate> {
if !self.seen.insert(event.id.to_string()) {
return Vec::new();
}
match &event.data {
EventData::OutputMessageStarted(_) => {
self.current_message_streamed = false;
Vec::new()
}
EventData::OutputMessageDelta(data) => {
if data.delta.is_empty() {
return Vec::new();
}
self.current_message_streamed = true;
vec![SessionUpdate::AgentMessageChunk {
content: ContentBlock::text(&data.delta),
}]
}
EventData::OutputMessageCompleted(data) => {
if self.current_message_streamed {
return Vec::new();
}
if data.message.role != MessageRole::Agent || data.message.has_tool_calls() {
return Vec::new();
}
match data.message.text().map(str::trim) {
Some(text) if !text.is_empty() => vec![SessionUpdate::AgentMessageChunk {
content: ContentBlock::text(text),
}],
_ => Vec::new(),
}
}
EventData::ReasonThinkingDelta(data) => {
if data.delta.is_empty() {
return Vec::new();
}
vec![SessionUpdate::AgentThoughtChunk {
content: ContentBlock::text(&data.delta),
}]
}
EventData::ReasonItem(data) => data
.summary
.iter()
.filter_map(|segment| {
let trimmed = segment.trim();
(!trimmed.is_empty()).then(|| SessionUpdate::AgentThoughtChunk {
content: ContentBlock::text(trimmed),
})
})
.collect(),
EventData::ToolStarted(data) => {
let name = data.tool_call.name.as_str();
if name == WRITE_TODOS {
return plan_from_value(&data.tool_call.arguments)
.map(|entries| vec![SessionUpdate::Plan { entries }])
.unwrap_or_default();
}
let title = data
.narration
.as_deref()
.or(data.display_name.as_deref())
.unwrap_or(name)
.to_string();
vec![SessionUpdate::ToolCall {
tool_call_id: data.tool_call.id.clone(),
title,
kind: tool_kind(name),
status: ToolCallStatus::InProgress,
raw_input: non_null(data.tool_call.arguments.clone()),
content: Vec::new(),
}]
}
EventData::ToolCompleted(data) => {
if data.tool_name == WRITE_TODOS {
return result_value(data)
.as_ref()
.and_then(plan_from_value)
.map(|entries| vec![SessionUpdate::Plan { entries }])
.unwrap_or_default();
}
let status = if data.success {
ToolCallStatus::Completed
} else {
ToolCallStatus::Failed
};
let content = tool_result_content(data)
.map(|block| vec![ToolCallContent::Content { content: block }])
.unwrap_or_default();
vec![SessionUpdate::ToolCallUpdate {
tool_call_id: data.tool_call_id.clone(),
status: Some(status),
content,
raw_output: None,
}]
}
_ => Vec::new(),
}
}
}
fn tool_kind(name: &str) -> ToolKind {
match name {
"read_file" | "stat_file" | "list_directory" => ToolKind::Read,
"grep" | "grep_files" | "duckduckgo_search" => ToolKind::Search,
"edit_file" | "write_file" | "create_directory" => ToolKind::Edit,
"delete_file" => ToolKind::Delete,
"bash" => ToolKind::Execute,
"web_fetch" => ToolKind::Fetch,
_ => ToolKind::Other,
}
}
fn tool_result_content(data: &ToolCompletedData) -> Option<ContentBlock> {
let summary = crate::app::summarize_tool_result(data);
let trimmed = summary.trim();
if trimmed.is_empty() {
None
} else {
Some(ContentBlock::text(trimmed))
}
}
fn plan_from_value(value: &Value) -> Option<Vec<PlanEntry>> {
let todos = value.get("todos")?.as_array()?;
let entries = todos
.iter()
.filter_map(|todo| {
let content = todo.get("content").and_then(Value::as_str)?;
if content.trim().is_empty() {
return None;
}
let status = match todo.get("status").and_then(Value::as_str) {
Some("completed") => PlanStatus::Completed,
Some("in_progress") => PlanStatus::InProgress,
_ => PlanStatus::Pending,
};
Some(PlanEntry {
content: content.to_string(),
priority: PlanPriority::Medium,
status,
})
})
.collect::<Vec<_>>();
Some(entries)
}
fn result_value(data: &ToolCompletedData) -> Option<Value> {
let parts = data.result.as_ref()?;
for part in parts {
if let ContentPart::Text(t) = part
&& let Ok(v) = serde_json::from_str::<Value>(&t.text)
{
return Some(v);
}
}
None
}
fn non_null(value: Value) -> Option<Value> {
if value.is_null() { None } else { Some(value) }
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::Utc;
use everruns_core::events::{
Event, EventContext, OutputMessageCompletedData, OutputMessageDeltaData,
ReasonThinkingDeltaData, ToolCompletedData, ToolStartedData,
};
use everruns_core::message::Message;
use everruns_core::tool_types::ToolCall;
use everruns_core::typed_id::{EventId, SessionId, TurnId};
use serde_json::json;
fn event(data: EventData) -> Event {
Event {
id: EventId::new(),
event_type: data.event_type().to_string(),
ts: Utc::now(),
session_id: SessionId::new(),
context: EventContext::empty(),
data,
metadata: None,
tags: None,
sequence: None,
}
}
#[test]
fn streaming_deltas_become_message_chunks() {
let mut t = Translator::new();
let updates = t.on_event(&event(EventData::OutputMessageDelta(
OutputMessageDeltaData {
turn_id: TurnId::new(),
delta: "Hel".into(),
accumulated: "Hel".into(),
},
)));
assert_eq!(
updates,
vec![SessionUpdate::AgentMessageChunk {
content: ContentBlock::text("Hel")
}]
);
}
#[test]
fn completed_message_suppressed_after_streaming() {
let mut t = Translator::new();
let _ = t.on_event(&event(EventData::OutputMessageDelta(
OutputMessageDeltaData {
turn_id: TurnId::new(),
delta: "Hi".into(),
accumulated: "Hi".into(),
},
)));
let completed = t.on_event(&event(EventData::OutputMessageCompleted(
OutputMessageCompletedData {
message: Message::assistant("Hi"),
metadata: None,
usage: None,
error_code: None,
error_fields: None,
},
)));
assert!(
completed.is_empty(),
"streamed text must not be re-sent: {completed:?}"
);
}
#[test]
fn completed_message_synthesised_when_not_streamed() {
let mut t = Translator::new();
let updates = t.on_event(&event(EventData::OutputMessageCompleted(
OutputMessageCompletedData {
message: Message::assistant("full answer"),
metadata: None,
usage: None,
error_code: None,
error_fields: None,
},
)));
assert_eq!(
updates,
vec![SessionUpdate::AgentMessageChunk {
content: ContentBlock::text("full answer")
}]
);
}
#[test]
fn thinking_deltas_become_thought_chunks() {
let mut t = Translator::new();
let updates = t.on_event(&event(EventData::ReasonThinkingDelta(
ReasonThinkingDeltaData {
turn_id: TurnId::new(),
delta: "pondering".into(),
accumulated: "pondering".into(),
},
)));
assert_eq!(
updates,
vec![SessionUpdate::AgentThoughtChunk {
content: ContentBlock::text("pondering")
}]
);
}
#[test]
fn tool_started_maps_kind_and_in_progress_status() {
let mut t = Translator::new();
let updates = t.on_event(&event(EventData::ToolStarted(ToolStartedData {
tool_call: ToolCall {
id: "call_1".into(),
name: "bash".into(),
arguments: json!({ "command": "ls" }),
},
tool_call_fingerprint: None,
display_name: Some("Bash".into()),
narration: Some("Listing files".into()),
})));
assert_eq!(
updates,
vec![SessionUpdate::ToolCall {
tool_call_id: "call_1".into(),
title: "Listing files".into(),
kind: ToolKind::Execute,
status: ToolCallStatus::InProgress,
raw_input: Some(json!({ "command": "ls" })),
content: vec![],
}]
);
}
#[test]
fn tool_completed_failure_maps_to_failed_status() {
let mut t = Translator::new();
let updates = t.on_event(&event(EventData::ToolCompleted(ToolCompletedData {
tool_call_id: "call_1".into(),
tool_name: "bash".into(),
tool_call_fingerprint: None,
tool_result_fingerprint: None,
display_name: None,
success: false,
status: "error".into(),
result: None,
error: Some("boom".into()),
duration_ms: None,
capability_id: None,
capability_name: None,
narration: None,
})));
assert_eq!(updates.len(), 1);
match &updates[0] {
SessionUpdate::ToolCallUpdate {
tool_call_id,
status,
content,
..
} => {
assert_eq!(tool_call_id, "call_1");
assert_eq!(*status, Some(ToolCallStatus::Failed));
assert_eq!(
content,
&vec![ToolCallContent::Content {
content: ContentBlock::text("error: boom")
}]
);
}
other => panic!("expected tool_call_update, got {other:?}"),
}
}
#[test]
fn write_todos_started_becomes_plan() {
let mut t = Translator::new();
let updates = t.on_event(&event(EventData::ToolStarted(ToolStartedData {
tool_call: ToolCall {
id: "call_todos".into(),
name: "write_todos".into(),
arguments: json!({
"todos": [
{ "content": "first", "status": "completed" },
{ "content": "second", "status": "in_progress" },
{ "content": "third", "status": "pending" },
]
}),
},
tool_call_fingerprint: None,
display_name: None,
narration: None,
})));
assert_eq!(
updates,
vec![SessionUpdate::Plan {
entries: vec![
PlanEntry {
content: "first".into(),
priority: PlanPriority::Medium,
status: PlanStatus::Completed,
},
PlanEntry {
content: "second".into(),
priority: PlanPriority::Medium,
status: PlanStatus::InProgress,
},
PlanEntry {
content: "third".into(),
priority: PlanPriority::Medium,
status: PlanStatus::Pending,
},
]
}]
);
}
#[test]
fn duplicate_event_id_is_ignored() {
let mut t = Translator::new();
let ev = event(EventData::OutputMessageDelta(OutputMessageDeltaData {
turn_id: TurnId::new(),
delta: "x".into(),
accumulated: "x".into(),
}));
assert_eq!(t.on_event(&ev).len(), 1);
assert_eq!(t.on_event(&ev).len(), 0, "second delivery must be ignored");
}
}