use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use crate::{DisplayMessage, DisplayMessageKind};
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct AguiEvent {
#[serde(rename = "type")]
pub event_type: String,
pub id: String,
pub sequence: usize,
pub session_id: String,
pub run_id: String,
#[serde(default, skip_serializing_if = "Value::is_null")]
pub payload: Value,
}
impl AguiEvent {
pub fn to_jsonl_line(&self) -> serde_json::Result<String> {
serde_json::to_string(self).map(|line| format!("{line}\n"))
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct VercelDataStreamPart {
#[serde(rename = "type")]
pub part_type: String,
#[serde(default, skip_serializing_if = "Value::is_null")]
pub value: Value,
}
#[must_use]
pub fn display_to_agui_event(message: &DisplayMessage) -> AguiEvent {
AguiEvent {
event_type: display_event_type(message.kind).to_string(),
id: message.sequence.to_string(),
sequence: message.sequence,
session_id: message.session_id.as_str().to_string(),
run_id: message.run_id.as_str().to_string(),
payload: display_payload(message),
}
}
pub fn display_to_agui_jsonl(messages: &[DisplayMessage]) -> serde_json::Result<String> {
let mut out = String::new();
for message in messages {
out.push_str(&display_to_agui_event(message).to_jsonl_line()?);
}
Ok(out)
}
#[must_use]
pub fn display_to_vercel_data_stream(message: &DisplayMessage) -> Vec<VercelDataStreamPart> {
match message.kind {
DisplayMessageKind::RunStarted => vec![part(
"start",
json!({
"runId": message.run_id.as_str(),
"sessionId": message.session_id.as_str(),
}),
)],
DisplayMessageKind::AssistantTextStart => vec![part("text-start", Value::Null)],
DisplayMessageKind::AssistantTextDelta => vec![part(
"text-delta",
json!({"textDelta": text_delta(&message.payload)}),
)],
DisplayMessageKind::AssistantTextEnd => vec![part("text-end", Value::Null)],
DisplayMessageKind::ToolCallStart => vec![part("tool-call-start", message.payload.clone())],
DisplayMessageKind::ToolCallDelta => vec![part("tool-call-delta", message.payload.clone())],
DisplayMessageKind::ToolCallEnd => vec![part("tool-call-end", message.payload.clone())],
DisplayMessageKind::ToolResult => vec![part("tool-result", message.payload.clone())],
DisplayMessageKind::RunCompleted => vec![part("finish", message.payload.clone())],
DisplayMessageKind::RunFailed => vec![part("error", message.payload.clone())],
DisplayMessageKind::RunCancelled => vec![part("abort", message.payload.clone())],
DisplayMessageKind::RunQueued
| DisplayMessageKind::ToolsUnavailable
| DisplayMessageKind::ToolSearchLoaded
| DisplayMessageKind::ToolSearchInitialized
| DisplayMessageKind::ToolSearchRefreshed
| DisplayMessageKind::ToolSearchInvalidated
| DisplayMessageKind::ToolSearchFailed
| DisplayMessageKind::ToolSearchNoMatch
| DisplayMessageKind::ToolsetInitialized
| DisplayMessageKind::ToolsetUnavailable
| DisplayMessageKind::ToolsetFailed
| DisplayMessageKind::ToolsetRefreshed
| DisplayMessageKind::ToolsetClosed
| DisplayMessageKind::ApprovalRequested
| DisplayMessageKind::DeferredRequested
| DisplayMessageKind::ApprovalResolved
| DisplayMessageKind::HitlResolved
| DisplayMessageKind::HitlDiagnostic
| DisplayMessageKind::Checkpoint
| DisplayMessageKind::SkillsScanned
| DisplayMessageKind::SkillActivated
| DisplayMessageKind::SkillsReloaded
| DisplayMessageKind::SubagentStarted
| DisplayMessageKind::SubagentCompleted
| DisplayMessageKind::SubagentFailed
| DisplayMessageKind::CompactionStarted
| DisplayMessageKind::CompactionCompleted
| DisplayMessageKind::CompactionFailed
| DisplayMessageKind::HandoffStarted
| DisplayMessageKind::HandoffCompleted
| DisplayMessageKind::HandoffFailed
| DisplayMessageKind::SteeringSubmitted
| DisplayMessageKind::SteeringReceived
| DisplayMessageKind::GoalIteration
| DisplayMessageKind::GoalCompleted
| DisplayMessageKind::TaskSnapshot
| DisplayMessageKind::TaskEvent
| DisplayMessageKind::NoteEvent
| DisplayMessageKind::FileEvent
| DisplayMessageKind::MediaEvent
| DisplayMessageKind::HostEvent => vec![part(
"data",
json!({
"type": display_event_type(message.kind),
"payload": message.payload,
}),
)],
}
}
pub fn display_to_vercel_data_stream_jsonl(
messages: &[DisplayMessage],
) -> serde_json::Result<String> {
let mut out = String::new();
for message in messages {
for part in display_to_vercel_data_stream(message) {
out.push_str(&serde_json::to_string(&part)?);
out.push('\n');
}
}
Ok(out)
}
fn display_payload(message: &DisplayMessage) -> Value {
let mut payload = match &message.payload {
Value::Object(map) => map.clone(),
Value::Null => serde_json::Map::new(),
value => {
let mut map = serde_json::Map::new();
map.insert("value".to_string(), value.clone());
map
}
};
payload.insert(
"timestamp".to_string(),
Value::String(message.timestamp.to_rfc3339()),
);
if let Some(preview) = &message.preview {
payload.insert("preview".to_string(), Value::String(preview.clone()));
}
Value::Object(payload)
}
const fn display_event_type(kind: DisplayMessageKind) -> &'static str {
match kind {
DisplayMessageKind::RunQueued => "RUN_QUEUED",
DisplayMessageKind::RunStarted => "RUN_STARTED",
DisplayMessageKind::AssistantTextStart => "TEXT_MESSAGE_START",
DisplayMessageKind::AssistantTextDelta => "TEXT_MESSAGE_CONTENT",
DisplayMessageKind::AssistantTextEnd => "TEXT_MESSAGE_END",
DisplayMessageKind::ToolCallStart => "TOOL_CALL_START",
DisplayMessageKind::ToolCallDelta => "TOOL_CALL_ARGS",
DisplayMessageKind::ToolCallEnd => "TOOL_CALL_END",
DisplayMessageKind::ToolResult => "TOOL_CALL_RESULT",
DisplayMessageKind::ToolsUnavailable => "TOOLS_UNAVAILABLE",
DisplayMessageKind::ToolSearchLoaded => "TOOL_SEARCH_LOADED",
DisplayMessageKind::ToolSearchInitialized => "TOOL_SEARCH_INITIALIZED",
DisplayMessageKind::ToolSearchRefreshed => "TOOL_SEARCH_REFRESHED",
DisplayMessageKind::ToolSearchInvalidated => "TOOL_SEARCH_INVALIDATED",
DisplayMessageKind::ToolSearchFailed => "TOOL_SEARCH_FAILED",
DisplayMessageKind::ToolSearchNoMatch => "TOOL_SEARCH_NO_MATCH",
DisplayMessageKind::ToolsetInitialized => "TOOLSET_INITIALIZED",
DisplayMessageKind::ToolsetUnavailable => "TOOLSET_UNAVAILABLE",
DisplayMessageKind::ToolsetFailed => "TOOLSET_FAILED",
DisplayMessageKind::ToolsetRefreshed => "TOOLSET_REFRESHED",
DisplayMessageKind::ToolsetClosed => "TOOLSET_CLOSED",
DisplayMessageKind::ApprovalRequested => "APPROVAL_REQUESTED",
DisplayMessageKind::DeferredRequested => "DEFERRED_REQUESTED",
DisplayMessageKind::ApprovalResolved => "APPROVAL_RESOLVED",
DisplayMessageKind::HitlResolved => "HITL_RESOLVED",
DisplayMessageKind::HitlDiagnostic => "HITL_DIAGNOSTIC",
DisplayMessageKind::Checkpoint => "CHECKPOINT",
DisplayMessageKind::SkillsScanned => "SKILLS_SCANNED",
DisplayMessageKind::SkillActivated => "SKILL_ACTIVATED",
DisplayMessageKind::SkillsReloaded => "SKILLS_RELOADED",
DisplayMessageKind::SubagentStarted => "SUBAGENT_STARTED",
DisplayMessageKind::SubagentCompleted => "SUBAGENT_COMPLETED",
DisplayMessageKind::SubagentFailed => "SUBAGENT_FAILED",
DisplayMessageKind::CompactionStarted => "COMPACTION_STARTED",
DisplayMessageKind::CompactionCompleted => "COMPACTION_COMPLETED",
DisplayMessageKind::CompactionFailed => "COMPACTION_FAILED",
DisplayMessageKind::HandoffStarted => "HANDOFF_STARTED",
DisplayMessageKind::HandoffCompleted => "HANDOFF_COMPLETED",
DisplayMessageKind::HandoffFailed => "HANDOFF_FAILED",
DisplayMessageKind::SteeringSubmitted => "STEERING_SUBMITTED",
DisplayMessageKind::SteeringReceived => "STEERING_RECEIVED",
DisplayMessageKind::GoalIteration => "GOAL_ITERATION",
DisplayMessageKind::GoalCompleted => "GOAL_COMPLETED",
DisplayMessageKind::TaskSnapshot => "TASK_SNAPSHOT",
DisplayMessageKind::TaskEvent => "TASK_EVENT",
DisplayMessageKind::NoteEvent => "NOTE_EVENT",
DisplayMessageKind::FileEvent => "FILE_EVENT",
DisplayMessageKind::MediaEvent => "MEDIA_EVENT",
DisplayMessageKind::HostEvent => "HOST_EVENT",
DisplayMessageKind::RunCompleted => "RUN_FINISHED",
DisplayMessageKind::RunFailed => "RUN_ERROR",
DisplayMessageKind::RunCancelled => "RUN_CANCELLED",
}
}
fn text_delta(payload: &Value) -> String {
payload
.get("delta")
.or_else(|| payload.get("text_delta"))
.or_else(|| payload.get("text"))
.and_then(Value::as_str)
.unwrap_or_default()
.to_string()
}
fn part(part_type: impl Into<String>, value: Value) -> VercelDataStreamPart {
VercelDataStreamPart {
part_type: part_type.into(),
value,
}
}