use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use starweaver_core::{AgentId, Metadata, RunId, SessionId, TraceContext};
use crate::AgentStreamRecord;
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub enum DisplayMessageKind {
#[serde(rename = "RUN_QUEUED")]
RunQueued,
#[serde(rename = "RUN_STARTED")]
RunStarted,
#[serde(rename = "TEXT_MESSAGE_START")]
AssistantTextStart,
#[serde(rename = "TEXT_MESSAGE_CONTENT")]
AssistantTextDelta,
#[serde(rename = "TEXT_MESSAGE_END")]
AssistantTextEnd,
#[serde(rename = "TOOL_CALL_START")]
ToolCallStart,
#[serde(rename = "TOOL_CALL_ARGS")]
ToolCallDelta,
#[serde(rename = "TOOL_CALL_END")]
ToolCallEnd,
#[serde(rename = "TOOL_CALL_RESULT")]
ToolResult,
#[serde(rename = "TOOLS_UNAVAILABLE")]
ToolsUnavailable,
#[serde(rename = "TOOL_SEARCH_LOADED")]
ToolSearchLoaded,
#[serde(rename = "TOOL_SEARCH_INITIALIZED")]
ToolSearchInitialized,
#[serde(rename = "TOOL_SEARCH_REFRESHED")]
ToolSearchRefreshed,
#[serde(rename = "TOOL_SEARCH_INVALIDATED")]
ToolSearchInvalidated,
#[serde(rename = "TOOL_SEARCH_FAILED")]
ToolSearchFailed,
#[serde(rename = "TOOL_SEARCH_NO_MATCH")]
ToolSearchNoMatch,
#[serde(rename = "TOOLSET_INITIALIZED")]
ToolsetInitialized,
#[serde(rename = "TOOLSET_UNAVAILABLE")]
ToolsetUnavailable,
#[serde(rename = "TOOLSET_FAILED")]
ToolsetFailed,
#[serde(rename = "TOOLSET_REFRESHED")]
ToolsetRefreshed,
#[serde(rename = "TOOLSET_CLOSED")]
ToolsetClosed,
#[serde(rename = "APPROVAL_REQUESTED")]
ApprovalRequested,
#[serde(rename = "DEFERRED_REQUESTED")]
DeferredRequested,
#[serde(rename = "APPROVAL_RESOLVED")]
ApprovalResolved,
#[serde(rename = "HITL_RESOLVED")]
HitlResolved,
#[serde(rename = "HITL_DIAGNOSTIC")]
HitlDiagnostic,
#[serde(rename = "CHECKPOINT")]
Checkpoint,
#[serde(rename = "SKILLS_SCANNED")]
SkillsScanned,
#[serde(rename = "SKILL_ACTIVATED")]
SkillActivated,
#[serde(rename = "SKILLS_RELOADED")]
SkillsReloaded,
#[serde(rename = "SUBAGENT_STARTED")]
SubagentStarted,
#[serde(rename = "SUBAGENT_COMPLETED")]
SubagentCompleted,
#[serde(rename = "SUBAGENT_FAILED")]
SubagentFailed,
#[serde(rename = "COMPACTION_STARTED")]
CompactionStarted,
#[serde(rename = "COMPACTION_COMPLETED")]
CompactionCompleted,
#[serde(rename = "COMPACTION_FAILED")]
CompactionFailed,
#[serde(rename = "HANDOFF_STARTED")]
HandoffStarted,
#[serde(rename = "HANDOFF_COMPLETED")]
HandoffCompleted,
#[serde(rename = "HANDOFF_FAILED")]
HandoffFailed,
#[serde(rename = "STEERING_SUBMITTED")]
SteeringSubmitted,
#[serde(rename = "STEERING_RECEIVED")]
SteeringReceived,
#[serde(rename = "GOAL_ITERATION")]
GoalIteration,
#[serde(rename = "GOAL_COMPLETED")]
GoalCompleted,
#[serde(rename = "TASK_SNAPSHOT")]
TaskSnapshot,
#[serde(rename = "TASK_EVENT")]
TaskEvent,
#[serde(rename = "NOTE_EVENT")]
NoteEvent,
#[serde(rename = "FILE_EVENT")]
FileEvent,
#[serde(rename = "MEDIA_EVENT")]
MediaEvent,
#[serde(rename = "HOST_EVENT", alias = "HOST_OPERATION")]
HostEvent,
#[serde(rename = "RUN_FINISHED")]
RunCompleted,
#[serde(rename = "RUN_ERROR")]
RunFailed,
#[serde(rename = "RUN_CANCELLED")]
RunCancelled,
}
impl DisplayMessageKind {
#[must_use]
pub const fn is_terminal(self) -> bool {
matches!(
self,
Self::RunCompleted | Self::RunFailed | Self::RunCancelled
)
}
}
#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum DisplayVisibility {
#[default]
Public,
Diagnostic,
Internal,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct DisplayMessage {
#[serde(default = "default_display_schema")]
pub schema: String,
pub sequence: usize,
pub session_id: SessionId,
pub run_id: RunId,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub agent_id: Option<AgentId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub agent_name: Option<String>,
pub timestamp: DateTime<Utc>,
#[serde(default, skip_serializing_if = "TraceContext::is_empty")]
pub trace_context: TraceContext,
#[serde(rename = "type")]
pub kind: DisplayMessageKind,
#[serde(default, skip_serializing_if = "Value::is_null")]
pub payload: Value,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub preview: Option<String>,
#[serde(default)]
pub visibility: DisplayVisibility,
#[serde(default, skip_serializing_if = "Metadata::is_empty")]
pub metadata: Metadata,
}
impl DisplayMessage {
pub const SCHEMA: &'static str = "starweaver.display.v1";
#[must_use]
pub fn new(
sequence: usize,
session_id: SessionId,
run_id: RunId,
kind: DisplayMessageKind,
) -> Self {
Self {
schema: Self::SCHEMA.to_string(),
sequence,
session_id,
run_id,
agent_id: None,
agent_name: None,
timestamp: Utc::now(),
trace_context: TraceContext::default(),
kind,
payload: Value::Null,
preview: None,
visibility: DisplayVisibility::Public,
metadata: Metadata::default(),
}
}
#[must_use]
pub fn with_payload(mut self, payload: Value) -> Self {
self.payload = payload;
self
}
#[must_use]
pub fn with_preview(mut self, preview: impl Into<String>) -> Self {
self.preview = Some(preview.into());
self
}
#[must_use]
pub fn with_trace_context(mut self, trace_context: TraceContext) -> Self {
self.trace_context = trace_context;
self
}
#[must_use]
pub const fn is_terminal(&self) -> bool {
self.kind.is_terminal()
}
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 DisplayProjectionContext {
pub session_id: SessionId,
pub run_id: RunId,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub agent_id: Option<AgentId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub agent_name: Option<String>,
#[serde(default, skip_serializing_if = "TraceContext::is_empty")]
pub trace_context: TraceContext,
}
impl DisplayProjectionContext {
#[must_use]
pub fn new(session_id: SessionId, run_id: RunId) -> Self {
Self {
session_id,
run_id,
agent_id: None,
agent_name: None,
trace_context: TraceContext::default(),
}
}
}
#[async_trait]
pub trait DisplayMessageProjector: Send + Sync {
async fn project(
&self,
context: &DisplayProjectionContext,
record: &AgentStreamRecord,
) -> Vec<DisplayMessage>;
}
pub(super) fn default_display_schema() -> String {
DisplayMessage::SCHEMA.to_string()
}