1use std::sync::{Arc, Mutex, PoisonError};
4
5use serde::{Deserialize, Serialize};
6use starweaver_context::AgentEvent;
7use starweaver_core::{AgentId, ConversationId, Metadata, RunId, TaskId};
8use starweaver_model::{ModelResponse, ModelResponseStreamEvent, ToolCallPart, ToolReturnPart};
9
10use crate::{executor::AgentExecutionNode, run::RunStatus, AgentResult};
11
12#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
14#[serde(rename_all = "snake_case")]
15pub enum AgentSidebandEventCategory {
16 Run,
18 Model,
20 Tool,
22 ToolSearch,
24 Hitl,
26 Skill,
28 Task,
30 Note,
32 File,
34 Media,
36 HostOperation,
38 Subagent,
40 Message,
42 Usage,
44 Compact,
46 Steering,
48 Background,
50 Capability,
52}
53
54#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
56pub struct AgentSidebandEvent {
57 pub category: AgentSidebandEventCategory,
59 pub kind: String,
61 #[serde(default)]
63 pub payload: serde_json::Value,
64 #[serde(default, skip_serializing_if = "Metadata::is_empty")]
66 pub metadata: Metadata,
67}
68
69impl AgentSidebandEvent {
70 #[must_use]
72 pub fn from_agent_event(event: &AgentEvent) -> Option<Self> {
73 Self::category_for_kind(&event.kind).map(|category| Self {
74 category,
75 kind: event.kind.clone(),
76 payload: event.payload.clone(),
77 metadata: event.metadata.clone(),
78 })
79 }
80
81 #[must_use]
83 pub fn category_for_kind(kind: &str) -> Option<AgentSidebandEventCategory> {
84 match kind {
85 "run_start" | "run_complete" | "run_failed" | "run_waiting" | "run_cancelled" => {
86 Some(AgentSidebandEventCategory::Run)
87 }
88 "model_error_retry" | "model_stream_resume" => Some(AgentSidebandEventCategory::Model),
89 "tools_unavailable"
90 | "toolset_initialized"
91 | "toolset_unavailable"
92 | "toolset_failed"
93 | "toolset_refreshed"
94 | "toolset_closed" => Some(AgentSidebandEventCategory::Tool),
95 "tool_search_loaded" | "tool_search_initialized" | "tool_search_refreshed" => {
96 Some(AgentSidebandEventCategory::ToolSearch)
97 }
98 "hitl_resolved" => Some(AgentSidebandEventCategory::Hitl),
99 "skills_scanned" | "skill_activated" | "skills_reloaded" => {
100 Some(AgentSidebandEventCategory::Skill)
101 }
102 "task_snapshot" => Some(AgentSidebandEventCategory::Task),
103 "usage_snapshot" => Some(AgentSidebandEventCategory::Usage),
104 "compact_start" | "compact_failed" | "compact_complete" => {
105 Some(AgentSidebandEventCategory::Compact)
106 }
107 "steering_received" | "steering_submitted" => {
108 Some(AgentSidebandEventCategory::Steering)
109 }
110 "background_shell_complete" => Some(AgentSidebandEventCategory::Background),
111 "message_received" => Some(AgentSidebandEventCategory::Message),
112 "subagent_started" | "subagent_completed" | "subagent_failed" => {
113 Some(AgentSidebandEventCategory::Subagent)
114 }
115 _ if kind.starts_with("task_") => Some(AgentSidebandEventCategory::Task),
116 _ if kind.starts_with("note_") => Some(AgentSidebandEventCategory::Note),
117 _ if kind.starts_with("file_") => Some(AgentSidebandEventCategory::File),
118 _ if kind.starts_with("media_") => Some(AgentSidebandEventCategory::Media),
119 _ if kind.starts_with("host_") => Some(AgentSidebandEventCategory::HostOperation),
120 _ if kind.starts_with("tool_search_") => Some(AgentSidebandEventCategory::ToolSearch),
121 _ if kind.starts_with("tool_") || kind.starts_with("toolset_") => {
122 Some(AgentSidebandEventCategory::Tool)
123 }
124 _ if kind.starts_with("approval_")
125 || kind.starts_with("deferred_")
126 || kind.starts_with("hitl_") =>
127 {
128 Some(AgentSidebandEventCategory::Hitl)
129 }
130 _ if kind.starts_with("skill_") || kind.starts_with("skills_") => {
131 Some(AgentSidebandEventCategory::Skill)
132 }
133 _ if kind.starts_with("subagent_") => Some(AgentSidebandEventCategory::Subagent),
134 _ if kind.starts_with("message_") => Some(AgentSidebandEventCategory::Message),
135 _ if kind.starts_with("usage_") => Some(AgentSidebandEventCategory::Usage),
136 _ if kind.starts_with("compact_") => Some(AgentSidebandEventCategory::Compact),
137 _ if kind.starts_with("steering_") => Some(AgentSidebandEventCategory::Steering),
138 _ if kind.starts_with("background_") => Some(AgentSidebandEventCategory::Background),
139 _ if kind.starts_with("capability_") => Some(AgentSidebandEventCategory::Capability),
140 _ => None,
141 }
142 }
143}
144
145#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
147#[serde(tag = "kind", rename_all = "snake_case")]
148pub enum AgentStreamEvent {
149 RunStart {
151 run_id: RunId,
153 conversation_id: ConversationId,
155 },
156 NodeStart {
158 node: AgentExecutionNode,
160 step: usize,
162 status: RunStatus,
164 },
165 NodeComplete {
167 node: AgentExecutionNode,
169 step: usize,
171 status: RunStatus,
173 },
174 Custom {
176 event: AgentEvent,
178 },
179 ModelRequest {
181 step: usize,
183 },
184 ModelStream {
186 step: usize,
188 event: ModelResponseStreamEvent,
190 },
191 ModelResponse {
193 step: usize,
195 response: ModelResponse,
197 },
198 Checkpoint {
200 node: AgentExecutionNode,
202 step: usize,
204 },
205 Suspended {
207 node: AgentExecutionNode,
209 reason: String,
211 },
212 ToolCall {
214 step: usize,
216 call: ToolCallPart,
218 },
219 ToolReturn {
221 step: usize,
223 tool_return: ToolReturnPart,
225 },
226 OutputRetry {
228 retries: usize,
230 prompt: String,
232 },
233 SteeringGuard {
235 step: usize,
237 prompt: String,
239 },
240 RunComplete {
242 run_id: RunId,
244 output: String,
246 },
247 RunFailed {
249 run_id: RunId,
251 error_kind: String,
253 message: String,
255 },
256}
257
258impl AgentStreamEvent {
259 #[must_use]
261 pub fn sideband_event(&self) -> Option<AgentSidebandEvent> {
262 match self {
263 Self::Custom { event } => AgentSidebandEvent::from_agent_event(event),
264 _ => None,
265 }
266 }
267}
268
269#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
271#[serde(rename_all = "snake_case")]
272pub enum AgentStreamSourceKind {
273 Subagent,
275}
276
277#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
279pub struct AgentStreamSource {
280 pub kind: AgentStreamSourceKind,
282 pub agent_id: AgentId,
284 pub agent_name: String,
286 #[serde(default, skip_serializing_if = "Option::is_none")]
288 pub task_id: Option<TaskId>,
289 #[serde(default, skip_serializing_if = "Option::is_none")]
291 pub run_id: Option<RunId>,
292 #[serde(default, skip_serializing_if = "Option::is_none")]
294 pub parent_run_id: Option<RunId>,
295 pub source_sequence: usize,
297}
298
299impl AgentStreamSource {
300 #[must_use]
302 pub fn subagent(
303 agent_id: AgentId,
304 agent_name: impl Into<String>,
305 task_id: TaskId,
306 run_id: Option<RunId>,
307 parent_run_id: Option<RunId>,
308 source_sequence: usize,
309 ) -> Self {
310 Self {
311 kind: AgentStreamSourceKind::Subagent,
312 agent_id,
313 agent_name: agent_name.into(),
314 task_id: Some(task_id),
315 run_id,
316 parent_run_id,
317 source_sequence,
318 }
319 }
320}
321
322#[derive(Clone, Debug, Default)]
324pub struct AgentStreamSink {
325 records: Arc<Mutex<Vec<AgentStreamRecord>>>,
326}
327
328impl AgentStreamSink {
329 pub fn push(&self, record: AgentStreamRecord) {
331 self.records
332 .lock()
333 .unwrap_or_else(PoisonError::into_inner)
334 .push(record);
335 }
336
337 pub fn extend(&self, records: impl IntoIterator<Item = AgentStreamRecord>) {
339 self.records
340 .lock()
341 .unwrap_or_else(PoisonError::into_inner)
342 .extend(records);
343 }
344
345 #[must_use]
347 pub fn drain(&self) -> Vec<AgentStreamRecord> {
348 self.records
349 .lock()
350 .unwrap_or_else(PoisonError::into_inner)
351 .drain(..)
352 .collect()
353 }
354
355 #[must_use]
357 pub fn is_empty(&self) -> bool {
358 self.records
359 .lock()
360 .unwrap_or_else(PoisonError::into_inner)
361 .is_empty()
362 }
363}
364
365#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
367pub struct AgentStreamRecord {
368 pub sequence: usize,
370 #[serde(default, skip_serializing_if = "Option::is_none")]
372 pub source: Option<AgentStreamSource>,
373 pub event: AgentStreamEvent,
375}
376
377impl AgentStreamRecord {
378 #[must_use]
380 pub const fn new(sequence: usize, event: AgentStreamEvent) -> Self {
381 Self {
382 sequence,
383 source: None,
384 event,
385 }
386 }
387
388 #[must_use]
390 pub fn with_source(mut self, source: AgentStreamSource) -> Self {
391 self.source = Some(source);
392 self
393 }
394
395 #[must_use]
397 pub const fn with_sequence(mut self, sequence: usize) -> Self {
398 self.sequence = sequence;
399 self
400 }
401}
402
403#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
405pub struct AgentStreamResult {
406 pub result: AgentResult,
408 pub events: Vec<AgentStreamRecord>,
410}
411
412impl AgentStreamResult {
413 #[must_use]
415 pub fn events(&self) -> &[AgentStreamRecord] {
416 &self.events
417 }
418
419 #[must_use]
421 pub const fn result(&self) -> &AgentResult {
422 &self.result
423 }
424}
425
426pub(crate) fn push_stream_event(
427 events: &mut Option<&mut Vec<AgentStreamRecord>>,
428 event: AgentStreamEvent,
429) {
430 if let Some(events) = events.as_deref_mut() {
431 events.push(AgentStreamRecord::new(events.len(), event));
432 }
433}
434
435pub(crate) fn push_stream_record(
436 events: &mut Option<&mut Vec<AgentStreamRecord>>,
437 record: AgentStreamRecord,
438) {
439 if let Some(events) = events.as_deref_mut() {
440 events.push(record.with_sequence(events.len()));
441 }
442}