use std::collections::{BTreeMap, HashSet, VecDeque};
use std::sync::Arc;
use crate::events::{
OutputMessageCompletedData, OutputMessageDeltaData, ReasonItemData, SessionTaskEventData,
TaskMessageEventData, TurnFailedData,
};
use crate::session_task::{
SessionTask, SessionTaskState, TASK_KIND_SUBAGENT, TaskMessageDirection, TaskMessagePart,
};
use crate::{ContentPart, RuntimeMessage};
use everruns_contracts::execution_phase::ExecutionPhase;
use serde::de::DeserializeOwned;
use serde_json::Value;
use crate::ag_ui::{
ActivitySnapshotEvent, BaseEvent, Event, Interrupt, Metadata, ReasoningMessageContentEvent,
ReasoningMessageEndEvent, ReasoningMessageStartEvent, ReasoningSpanEvent, RunErrorEvent,
RunFinishedEvent, RunFinishedOutcome, SubagentErrorEvent, SubagentFinishedEvent,
SubagentFinishedOutcome, SubagentStartedEvent, TextMessageContentEvent, TextMessageEndEvent,
TextMessageStartEvent, TokenUsage, ToolCall, ToolCallArgsEvent, ToolCallEndEvent,
ToolCallStartEvent,
};
pub const METADATA_KEY: &str = "everruns";
pub const SUBAGENT_ACTIVITY_TYPE: &str = "everruns.subagent";
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct TurnFailure {
pub message: String,
pub code: Option<String>,
}
pub type ErrorProjection = Arc<dyn Fn(&TurnFailure) -> RunErrorEvent + Send + Sync>;
#[derive(Clone)]
pub struct ProjectionPolicy {
pub reasoning_visible: bool,
pub tool_activity_text: Option<String>,
pub usage_visible: bool,
pub subagents_visible: bool,
pub model_visible: bool,
pub session_id: Option<String>,
pub error: ErrorProjection,
}
impl Default for ProjectionPolicy {
fn default() -> Self {
Self {
reasoning_visible: true,
tool_activity_text: None,
usage_visible: true,
subagents_visible: true,
model_visible: true,
session_id: None,
error: Arc::new(|failure: &TurnFailure| RunErrorEvent {
code: failure.code.clone(),
..RunErrorEvent::new(failure.message.clone())
}),
}
}
}
impl std::fmt::Debug for ProjectionPolicy {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ProjectionPolicy")
.field("reasoning_visible", &self.reasoning_visible)
.field("tool_activity_text", &self.tool_activity_text)
.field("usage_visible", &self.usage_visible)
.field("subagents_visible", &self.subagents_visible)
.field("model_visible", &self.model_visible)
.field("session_id", &self.session_id)
.finish_non_exhaustive()
}
}
#[derive(Debug)]
pub struct Projector {
policy: ProjectionPolicy,
thread_id: String,
run_id: String,
queue: VecDeque<Event>,
assistant_message_id: Option<String>,
assistant_content_started: bool,
assistant_emitted_delta: bool,
tool_call_message_id: Option<String>,
reasoning_span: Option<String>,
reasoning_message: Option<String>,
span_opened_by_tools: bool,
thinking: bool,
active_tools: usize,
tool_activity_shown: bool,
next_reasoning_id: usize,
usage: BTreeMap<(Option<String>, Option<String>), TokenUsage>,
turn_id: Option<String>,
model: Option<String>,
open_subagents: Vec<(String, bool)>,
seen_subagents: HashSet<String>,
finished: bool,
}
impl Projector {
pub fn new(
thread_id: impl Into<String>,
run_id: impl Into<String>,
policy: ProjectionPolicy,
) -> Self {
Self {
policy,
thread_id: thread_id.into(),
run_id: run_id.into(),
queue: VecDeque::new(),
assistant_message_id: None,
assistant_content_started: false,
assistant_emitted_delta: false,
tool_call_message_id: None,
reasoning_span: None,
reasoning_message: None,
span_opened_by_tools: false,
thinking: false,
active_tools: 0,
tool_activity_shown: false,
next_reasoning_id: 0,
usage: BTreeMap::new(),
turn_id: None,
model: None,
open_subagents: Vec::new(),
seen_subagents: HashSet::new(),
finished: false,
}
}
pub fn observe_turn(&mut self, turn_id: impl Into<String>) {
self.turn_id = Some(turn_id.into());
}
pub fn run_metadata(&self) -> Option<Metadata> {
let mut everruns = serde_json::Map::new();
let mut put = |key: &str, value: Option<&String>| {
if let Some(value) = value {
everruns.insert(key.to_string(), Value::String(value.clone()));
}
};
put("sessionId", self.policy.session_id.as_ref());
put("turnId", self.turn_id.as_ref());
put(
"model",
self.model.as_ref().filter(|_| self.policy.model_visible),
);
(!everruns.is_empty())
.then(|| Metadata::from_iter([(METADATA_KEY.to_string(), Value::Object(everruns))]))
}
pub fn is_finished(&self) -> bool {
self.finished
}
pub fn drain(&mut self) -> impl Iterator<Item = Event> + '_ {
self.queue.drain(..)
}
pub fn pop(&mut self) -> Option<Event> {
self.queue.pop_front()
}
pub fn fail(&mut self, error: RunErrorEvent) {
if self.finished {
return;
}
self.close_all();
self.close_subagents();
let mut error = error;
if error.usage.is_none() {
error.usage = self.usage_entries();
}
if error.base.metadata.is_none() {
error.base.metadata = self.run_metadata();
}
self.queue.push_back(Event::RunError(error));
self.finished = true;
}
pub fn interrupt(&mut self, interrupts: Vec<Interrupt>) {
self.park(Vec::new(), interrupts);
}
pub fn park(&mut self, tool_calls: Vec<ToolCall>, interrupts: Vec<Interrupt>) {
if self.finished || (tool_calls.is_empty() && interrupts.is_empty()) {
return;
}
self.close_all();
let parent = self.tool_call_message_id.clone();
let mut pending = Vec::with_capacity(tool_calls.len());
for call in tool_calls {
let mut start = ToolCallStartEvent::new(call.id.clone(), call.function.name);
start.parent_message_id = parent.clone();
self.queue.push_back(Event::ToolCallStart(start));
if !call.function.arguments.is_empty() {
self.queue
.push_back(Event::ToolCallArgs(ToolCallArgsEvent::new(
call.id.clone(),
call.function.arguments,
)));
}
self.queue
.push_back(Event::ToolCallEnd(ToolCallEndEvent::new(call.id.clone())));
pending.push(call.id);
}
self.finish(Some(if interrupts.is_empty() {
RunFinishedOutcome::Success {
pending_tool_call_ids: Some(pending),
}
} else {
RunFinishedOutcome::Interrupt { interrupts }
}));
}
pub fn project(&mut self, event_type: &str, data: &Value) {
if self.finished {
return;
}
match event_type {
"output.message.delta" => {
if let Some(data) = parse::<OutputMessageDeltaData>(data) {
let message_id =
self.ensure_assistant_message(data.message_id.uuid().to_string());
self.open_assistant_text(&message_id);
self.queue
.push_back(Event::TextMessageContent(TextMessageContentEvent::new(
message_id, data.delta,
)));
self.assistant_emitted_delta = true;
}
}
"output.message.completed" => {
if let Some(data) = parse::<OutputMessageCompletedData>(data) {
self.output_completed(&data.message);
}
}
"reason.thinking.started" if self.policy.reasoning_visible => {
if self.reasoning_span.is_none() {
self.open_span(false);
}
self.span_opened_by_tools = false;
self.close_reasoning_message();
self.thinking = true;
}
"reason.thinking.delta" if self.policy.reasoning_visible => {
if !self.thinking {
self.project("reason.thinking.started", &Value::Null);
}
if let Some(delta) = data.get("delta").and_then(Value::as_str) {
let message_id = self.ensure_reasoning_message();
self.queue.push_back(Event::ReasoningMessageContent(
ReasoningMessageContentEvent::new(message_id, delta),
));
}
}
"reason.thinking.completed" if self.policy.reasoning_visible => {
self.close_reasoning_message();
self.close_span();
self.thinking = false;
self.span_opened_by_tools = false;
}
"reason.item" if self.policy.reasoning_visible => {
if let Some(data) = parse::<ReasonItemData>(data) {
self.reasoning_summary(&data.summary);
}
}
"tool.started" => {
self.active_tools += 1;
if let Some(text) = self.policy.tool_activity_text.clone() {
self.tool_activity(&text);
}
}
"tool.completed" => {
self.active_tools = self.active_tools.saturating_sub(1);
if self.active_tools == 0 && self.tool_activity_shown {
if self.span_opened_by_tools {
self.close_reasoning_message();
self.close_span();
}
self.tool_activity_shown = false;
self.span_opened_by_tools = false;
}
}
"llm.generation" => self.record_usage(data),
"task.created" | "task.updated" if self.policy.subagents_visible => {
if let Some(data) = parse::<SessionTaskEventData>(data) {
self.subagent_task(&data.task);
}
}
"task.message.received" if self.policy.subagents_visible => {
if let Some(data) = parse::<TaskMessageEventData>(data) {
self.subagent_message(&data);
}
}
"turn.completed" | "session.idled" => self.finish(None),
"turn.cancelled" => self.finish(Some(RunFinishedOutcome::Cancelled)),
"turn.failed" => {
let failure = parse::<TurnFailedData>(data)
.map(|data| TurnFailure {
message: data.error,
code: data.error_code,
})
.unwrap_or_default();
let error = (self.policy.error)(&failure);
self.fail(error);
}
_ => {}
}
}
fn finish(&mut self, outcome: Option<RunFinishedOutcome>) {
self.close_all();
self.close_subagents();
self.queue.push_back(Event::RunFinished(RunFinishedEvent {
base: BaseEvent {
metadata: self.run_metadata(),
..BaseEvent::default()
},
outcome,
usage: self.usage_entries(),
..RunFinishedEvent::new(self.thread_id.clone(), self.run_id.clone())
}));
self.finished = true;
}
fn record_usage(&mut self, data: &Value) {
let metadata = &data["metadata"];
if let Some(model) = ["response_model", "model"]
.iter()
.find_map(|key| metadata.get(*key).and_then(Value::as_str))
{
self.model = Some(model.to_string());
}
let Some(usage) = metadata.get("usage").filter(|usage| usage.is_object()) else {
return;
};
let count = |key: &str| usage.get(key).and_then(Value::as_u64);
let text = |key: &str| {
metadata
.get(key)
.and_then(Value::as_str)
.map(str::to_string)
};
let model = text("response_model").or_else(|| text("model"));
let entry = self
.usage
.entry((text("provider"), model.clone()))
.or_insert_with(|| TokenUsage {
provider: text("provider"),
model,
..TokenUsage::default()
});
let cache_read = count("cache_read_tokens");
let cache_write = count("cache_creation_tokens");
let input = count("input_tokens")
.map(|input| input + cache_read.unwrap_or_default() + cache_write.unwrap_or_default());
let add = |total: &mut Option<u64>, value: Option<u64>| {
if let Some(value) = value {
*total = Some(total.unwrap_or_default().saturating_add(value));
}
};
add(&mut entry.input_tokens, input);
add(&mut entry.output_tokens, count("output_tokens"));
add(&mut entry.cached_input_tokens, cache_read);
add(&mut entry.cache_write_input_tokens, cache_write);
entry.total_tokens = match (entry.input_tokens, entry.output_tokens) {
(Some(input), Some(output)) => Some(input.saturating_add(output)),
_ => None,
};
}
fn usage_entries(&self) -> Option<Vec<TokenUsage>> {
(self.policy.usage_visible && !self.usage.is_empty())
.then(|| self.usage.values().cloned().collect())
}
fn output_completed(&mut self, message: &RuntimeMessage) {
if !is_terminal_public_output(message, self.assistant_emitted_delta) {
self.close_assistant_text();
self.tool_call_message_id = Some(message.id.uuid().to_string());
return;
}
let message_id = self.ensure_assistant_message(message.id.uuid().to_string());
self.open_assistant_text(&message_id);
let text = public_text(&message.content);
if !self.assistant_emitted_delta && !text.is_empty() {
self.queue
.push_back(Event::TextMessageContent(TextMessageContentEvent::new(
message_id.clone(),
text,
)));
}
self.close_assistant_text();
self.finish(None);
}
fn ensure_assistant_message(&mut self, message_id: String) -> String {
if self.assistant_message_id.as_deref() != Some(message_id.as_str()) {
self.close_assistant_text();
self.assistant_message_id = Some(message_id.clone());
}
message_id
}
fn open_assistant_text(&mut self, message_id: &str) {
if !self.assistant_content_started {
self.queue
.push_back(Event::TextMessageStart(TextMessageStartEvent::assistant(
message_id,
)));
self.assistant_content_started = true;
}
}
fn close_assistant_text(&mut self) {
if self.assistant_content_started
&& let Some(message_id) = self.assistant_message_id.clone()
{
self.queue
.push_back(Event::TextMessageEnd(TextMessageEndEvent::new(message_id)));
}
self.assistant_message_id = None;
self.assistant_content_started = false;
self.assistant_emitted_delta = false;
}
fn next_id(&mut self, kind: &str) -> String {
self.next_reasoning_id += 1;
format!("{}-{kind}-{}", self.run_id, self.next_reasoning_id)
}
fn open_span(&mut self, by_tools: bool) {
let id = self.next_id("reasoning");
self.queue
.push_back(Event::ReasoningStart(ReasoningSpanEvent::new(id.clone())));
self.reasoning_span = Some(id);
self.span_opened_by_tools = by_tools;
}
fn close_span(&mut self) {
if let Some(id) = self.reasoning_span.take() {
self.queue
.push_back(Event::ReasoningEnd(ReasoningSpanEvent::new(id)));
}
}
fn ensure_reasoning_message(&mut self) -> String {
if let Some(id) = &self.reasoning_message {
return id.clone();
}
let id = self.next_id("reasoning-message");
self.queue.push_back(Event::ReasoningMessageStart(
ReasoningMessageStartEvent::new(id.clone()),
));
self.reasoning_message = Some(id.clone());
id
}
fn close_reasoning_message(&mut self) {
if let Some(id) = self.reasoning_message.take() {
self.queue
.push_back(Event::ReasoningMessageEnd(ReasoningMessageEndEvent::new(
id,
)));
}
}
fn reasoning_text(&mut self, text: &str) {
if self.thinking && self.reasoning_message.is_some() {
let id = self.ensure_reasoning_message();
self.queue.push_back(Event::ReasoningMessageContent(
ReasoningMessageContentEvent::new(id, format!("\n{text}")),
));
return;
}
let opened_span = self.reasoning_span.is_none();
if opened_span {
self.open_span(false);
}
self.close_reasoning_message();
let id = self.ensure_reasoning_message();
self.queue.push_back(Event::ReasoningMessageContent(
ReasoningMessageContentEvent::new(id, text),
));
self.close_reasoning_message();
if opened_span {
self.close_span();
}
}
fn reasoning_summary(&mut self, summary: &[String]) {
let text = summary
.iter()
.map(|segment| segment.trim())
.filter(|segment| !segment.is_empty())
.collect::<Vec<_>>()
.join("\n");
if !text.is_empty() {
self.reasoning_text(&text);
}
}
fn tool_activity(&mut self, text: &str) {
if !self.tool_activity_shown {
if self.reasoning_span.is_none() {
self.open_span(true);
}
self.tool_activity_shown = true;
}
if self.thinking && self.reasoning_message.is_some() {
let id = self.ensure_reasoning_message();
self.queue.push_back(Event::ReasoningMessageContent(
ReasoningMessageContentEvent::new(id, format!("\n{text}")),
));
return;
}
self.close_reasoning_message();
let id = self.ensure_reasoning_message();
self.queue.push_back(Event::ReasoningMessageContent(
ReasoningMessageContentEvent::new(id, text),
));
self.close_reasoning_message();
}
fn subagent_task(&mut self, task: &SessionTask) {
if task.kind != TASK_KIND_SUBAGENT {
return;
}
if !self.seen_subagents.contains(&task.id) {
if task.state.is_terminal() {
return;
}
self.seen_subagents.insert(task.id.clone());
let background = task.spec.get("mode").and_then(Value::as_str) == Some("background");
self.open_subagents.push((task.id.clone(), background));
self.queue
.push_back(Event::SubagentStarted(SubagentStartedEvent {
subagent_run_id: task.id.clone(),
name: task.display_name.clone(),
..SubagentStartedEvent::default()
}));
}
if !task.state.is_terminal() {
return;
}
let Some(position) = self
.open_subagents
.iter()
.position(|(id, _)| *id == task.id)
else {
return;
};
self.open_subagents.remove(position);
let id = task.id.clone();
match task.state {
SessionTaskState::Succeeded => {
if let Some(summary) = task.summary.as_deref().filter(|s| !s.trim().is_empty()) {
self.subagent_text(&id, &format!("{id}-summary"), summary);
}
self.queue
.push_back(Event::SubagentFinished(SubagentFinishedEvent {
subagent_run_id: id,
..SubagentFinishedEvent::default()
}));
}
SessionTaskState::Canceled => {
self.queue
.push_back(Event::SubagentError(SubagentErrorEvent {
subagent_run_id: id,
message: "Subagent canceled".to_string(),
code: Some("cancelled".to_string()),
..SubagentErrorEvent::default()
}));
}
_ => {
let failure = task
.error
.as_ref()
.map(|error| TurnFailure {
message: error.message.clone(),
code: Some(error.kind.clone()),
})
.unwrap_or_default();
let error = (self.policy.error)(&failure);
self.queue
.push_back(Event::SubagentError(SubagentErrorEvent {
subagent_run_id: id,
message: error.message,
code: error.code,
..SubagentErrorEvent::default()
}));
}
}
}
fn subagent_message(&mut self, data: &TaskMessageEventData) {
let message = &data.message;
if message.direction != TaskMessageDirection::Outbound
|| !self
.open_subagents
.iter()
.any(|(id, _)| *id == data.task_id)
{
return;
}
let text = message
.content
.iter()
.filter_map(|part| match part {
TaskMessagePart::Text { text } => Some(text.as_str()),
TaskMessagePart::Data { .. } => None,
})
.filter(|text| !text.is_empty())
.collect::<Vec<_>>()
.join("\n");
if !text.is_empty() {
self.subagent_text(&data.task_id, &message.id, &text);
}
for (index, part) in message.content.iter().enumerate() {
if let TaskMessagePart::Data { data: value } = part {
let content = serde_json::Map::from_iter([("data".to_string(), value.clone())]);
self.queue
.push_back(Event::ActivitySnapshot(ActivitySnapshotEvent {
message_id: format!("{}-{index}", message.id),
activity_type: SUBAGENT_ACTIVITY_TYPE.to_string(),
content,
subagent_run_id: Some(data.task_id.clone()),
..ActivitySnapshotEvent::default()
}));
}
}
}
fn subagent_text(&mut self, subagent_run_id: &str, message_id: &str, text: &str) {
let attributed = Some(subagent_run_id.to_string());
self.queue
.push_back(Event::TextMessageStart(TextMessageStartEvent {
subagent_run_id: attributed.clone(),
..TextMessageStartEvent::assistant(message_id)
}));
self.queue
.push_back(Event::TextMessageContent(TextMessageContentEvent {
subagent_run_id: attributed.clone(),
..TextMessageContentEvent::new(message_id, text)
}));
self.queue
.push_back(Event::TextMessageEnd(TextMessageEndEvent {
subagent_run_id: attributed,
..TextMessageEndEvent::new(message_id)
}));
}
fn close_subagents(&mut self) {
for (id, background) in std::mem::take(&mut self.open_subagents) {
if background {
let content = serde_json::Map::from_iter([(
"state".to_string(),
Value::String("running".to_string()),
)]);
self.queue
.push_back(Event::ActivitySnapshot(ActivitySnapshotEvent {
message_id: format!("{id}-background"),
activity_type: SUBAGENT_ACTIVITY_TYPE.to_string(),
content,
subagent_run_id: Some(id.clone()),
..ActivitySnapshotEvent::default()
}));
}
self.queue
.push_back(Event::SubagentFinished(SubagentFinishedEvent {
subagent_run_id: id,
outcome: Some(SubagentFinishedOutcome::Suspended {
interrupt_ids: None,
}),
..SubagentFinishedEvent::default()
}));
}
}
fn close_all(&mut self) {
self.close_assistant_text();
self.close_reasoning_message();
self.close_span();
self.thinking = false;
self.tool_activity_shown = false;
self.span_opened_by_tools = false;
}
}
fn parse<T: DeserializeOwned>(data: &Value) -> Option<T> {
serde_json::from_value(data.clone()).ok()
}
pub fn is_terminal_public_output(message: &RuntimeMessage, emitted_delta: bool) -> bool {
if matches!(message.phase, Some(ExecutionPhase::Commentary)) {
return false;
}
if message
.content
.iter()
.any(|part| matches!(part, ContentPart::ToolCall(_) | ContentPart::ToolResult(_)))
{
return false;
}
emitted_delta || !public_text(&message.content).is_empty()
}
pub fn public_text(parts: &[ContentPart]) -> String {
parts
.iter()
.filter_map(public_part_text)
.filter(|part| !part.is_empty())
.collect::<Vec<_>>()
.join("\n")
}
fn public_part_text(part: &ContentPart) -> Option<String> {
match part {
ContentPart::Text(text) => Some(text.text.clone()),
ContentPart::Image(image) => {
Some(image.url.clone().unwrap_or_else(|| "[Image]".to_string()))
}
ContentPart::ImageFile(image) => Some(format!(
"[Image file: {}]",
image.filename.as_deref().unwrap_or("unnamed")
)),
ContentPart::File(file) => Some(format!(
"[File: {}]",
file.filename.as_deref().unwrap_or("unnamed")
)),
_ => None,
}
}