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::{
self, ContentBlock, Plan, PlanEntry, PlanEntryPriority, PlanEntryStatus, SessionUpdate,
ToolCall, ToolCallContent, ToolCallStatus, ToolCallUpdate, ToolCallUpdateFields,
};
const WRITE_TODOS: &str = "write_todos";
#[derive(Default)]
pub struct Translator {
current_message_streamed: bool,
seen: HashSet<String>,
replay_input_messages: bool,
}
impl Translator {
pub fn new() -> Self {
Self::default()
}
pub fn for_replay() -> Self {
Self {
replay_input_messages: true,
..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::InputMessage(data) => {
if !self.replay_input_messages {
return Vec::new();
}
if data.message.role != MessageRole::User {
return Vec::new();
}
match data.message.text().map(str::trim) {
Some(text) if !text.is_empty() => {
vec![SessionUpdate::UserMessageChunk(protocol::text_chunk(text))]
}
_ => Vec::new(),
}
}
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(protocol::text_chunk(
&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(protocol::text_chunk(text))]
}
_ => Vec::new(),
}
}
EventData::ReasonThinkingDelta(data) => {
if data.delta.is_empty() {
return Vec::new();
}
vec![SessionUpdate::AgentThoughtChunk(protocol::text_chunk(
&data.delta,
))]
}
EventData::ReasonItem(data) => data
.summary
.iter()
.filter_map(|segment| {
let trimmed = segment.trim();
(!trimmed.is_empty())
.then(|| SessionUpdate::AgentThoughtChunk(protocol::text_chunk(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(Plan::new(entries))])
.unwrap_or_default();
}
let title = data
.narration
.as_deref()
.or(data.display_name.as_deref())
.unwrap_or(name)
.to_string();
vec![SessionUpdate::ToolCall(
ToolCall::new(data.tool_call.id.clone(), title)
.status(ToolCallStatus::InProgress)
.raw_input(non_null(data.tool_call.arguments.clone())),
)]
}
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(Plan::new(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(protocol::Content::new(block))])
.unwrap_or_default();
vec![SessionUpdate::ToolCallUpdate(ToolCallUpdate::new(
data.tool_call_id.clone(),
ToolCallUpdateFields::new().status(status).content(content),
))]
}
_ => Vec::new(),
}
}
}
fn tool_result_content(data: &ToolCompletedData) -> Option<ContentBlock> {
let summary = crate::transcript::summarize_tool_result(data);
let trimmed = summary.trim();
if trimmed.is_empty() {
None
} else {
Some(protocol::text_block(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") => PlanEntryStatus::Completed,
Some("in_progress") => PlanEntryStatus::InProgress,
_ => PlanEntryStatus::Pending,
};
Some(PlanEntry::new(content, PlanEntryPriority::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(protocol::text_chunk(
"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,
error_disclosure: 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,
error_disclosure: None,
},
)));
assert_eq!(
updates,
vec![SessionUpdate::AgentMessageChunk(protocol::text_chunk(
"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(protocol::text_chunk(
"pondering"
))]
);
}
#[test]
fn tool_started_uses_in_progress_status_without_kind() {
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(
protocol::ToolCall::new("call_1", "Listing files")
.status(ToolCallStatus::InProgress)
.raw_input(json!({ "command": "ls" })),
)]
);
let serialized = serde_json::to_value(&updates[0]).unwrap();
assert_eq!(serialized["status"], "in_progress");
assert!(
serialized.get("kind").is_none(),
"autonomous tools must not advertise approval-looking categories: {serialized}"
);
}
#[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(update) => {
assert_eq!(update.tool_call_id.to_string(), "call_1");
assert_eq!(update.fields.status, Some(ToolCallStatus::Failed));
assert_eq!(
update.fields.content,
Some(vec![protocol::content("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(Plan::new(vec![
PlanEntry::new(
"first",
PlanEntryPriority::Medium,
PlanEntryStatus::Completed,
),
PlanEntry::new(
"second",
PlanEntryPriority::Medium,
PlanEntryStatus::InProgress,
),
PlanEntry::new("third", PlanEntryPriority::Medium, PlanEntryStatus::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");
}
#[test]
fn replay_mode_emits_user_messages() {
let mut t = Translator::for_replay();
let updates = t.on_event(&event(EventData::InputMessage(
everruns_core::events::InputMessageData::new(Message::user("prior prompt")),
)));
assert_eq!(
updates,
vec![SessionUpdate::UserMessageChunk(protocol::text_chunk(
"prior prompt"
))]
);
}
#[test]
fn live_mode_suppresses_user_messages() {
let mut t = Translator::new();
let updates = t.on_event(&event(EventData::InputMessage(
everruns_core::events::InputMessageData::new(Message::user("current prompt")),
)));
assert!(updates.is_empty());
}
}