1use std::sync::{Arc, Mutex, PoisonError};
4
5use serde::{Deserialize, Serialize};
6use starweaver_core::{
7 AgentEvent, AgentExecutionNode, AgentId, ConversationId, Metadata, RunId,
8 RunLifecycle as RunStatus, TaskId,
9};
10use starweaver_model::{ModelResponse, ModelResponseStreamEvent, ToolCallPart, ToolReturnPart};
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 HostEvent,
38 Subagent,
40 Message,
42 Usage,
44 Compact,
46 Steering,
48 Goal,
50 Background,
52 Capability,
54}
55
56#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
58pub struct AgentSidebandEvent {
59 pub category: AgentSidebandEventCategory,
61 pub kind: String,
63 #[serde(default)]
65 pub payload: serde_json::Value,
66 #[serde(default, skip_serializing_if = "Metadata::is_empty")]
68 pub metadata: Metadata,
69}
70
71impl AgentSidebandEvent {
72 #[must_use]
74 pub fn from_agent_event(event: &AgentEvent) -> Option<Self> {
75 Self::category_for_kind(&event.kind).map(|category| Self {
76 category,
77 kind: event.kind.clone(),
78 payload: event.payload.clone(),
79 metadata: event.metadata.clone(),
80 })
81 }
82
83 #[must_use]
85 pub fn category_for_kind(kind: &str) -> Option<AgentSidebandEventCategory> {
86 match kind {
87 "run_start" | "run_complete" | "run_failed" | "run_waiting" | "run_cancelled" => {
88 Some(AgentSidebandEventCategory::Run)
89 }
90 "model_error_retry"
91 | "model_stream_resume"
92 | "model_transport_selected"
93 | "model_transport_fallback" => Some(AgentSidebandEventCategory::Model),
94 "tools_unavailable"
95 | "toolset_initialized"
96 | "toolset_unavailable"
97 | "toolset_failed"
98 | "toolset_refreshed"
99 | "toolset_closed" => Some(AgentSidebandEventCategory::Tool),
100 "tool_search_loaded" | "tool_search_initialized" | "tool_search_refreshed" => {
101 Some(AgentSidebandEventCategory::ToolSearch)
102 }
103 "approval_requested"
104 | "approval_resolved"
105 | "deferred_requested"
106 | "deferred_completed"
107 | "deferred_failed"
108 | "deferred_cancelled"
109 | "hitl_resolved"
110 | "hitl_decision_diagnostic" => Some(AgentSidebandEventCategory::Hitl),
111 "skills_scanned" | "skill_activated" | "skills_reloaded" => {
112 Some(AgentSidebandEventCategory::Skill)
113 }
114 "task_snapshot" => Some(AgentSidebandEventCategory::Task),
115 "usage_snapshot" => Some(AgentSidebandEventCategory::Usage),
116 "compact_start" | "compact_failed" | "compact_complete" => {
117 Some(AgentSidebandEventCategory::Compact)
118 }
119 "steering_received" | "steering_submitted" => {
120 Some(AgentSidebandEventCategory::Steering)
121 }
122 "goal_iteration" | "goal_complete" => Some(AgentSidebandEventCategory::Goal),
123 "background_shell_complete" => Some(AgentSidebandEventCategory::Background),
124 "message_received" => Some(AgentSidebandEventCategory::Message),
125 "subagent_started" | "subagent_completed" | "subagent_failed" => {
126 Some(AgentSidebandEventCategory::Subagent)
127 }
128 _ if kind.starts_with("model_") => Some(AgentSidebandEventCategory::Model),
129 _ if kind.starts_with("task_") => Some(AgentSidebandEventCategory::Task),
130 _ if kind.starts_with("note_") => Some(AgentSidebandEventCategory::Note),
131 _ if kind.starts_with("file_") => Some(AgentSidebandEventCategory::File),
132 _ if kind.starts_with("media_") => Some(AgentSidebandEventCategory::Media),
133 _ if kind.starts_with("host_") => Some(AgentSidebandEventCategory::HostEvent),
134 _ if kind.starts_with("tool_search_") => Some(AgentSidebandEventCategory::ToolSearch),
135 _ if kind.starts_with("tool_") || kind.starts_with("toolset_") => {
136 Some(AgentSidebandEventCategory::Tool)
137 }
138 _ if kind.starts_with("approval_")
139 || kind.starts_with("deferred_")
140 || kind.starts_with("hitl_") =>
141 {
142 Some(AgentSidebandEventCategory::Hitl)
143 }
144 _ if kind.starts_with("skill_") || kind.starts_with("skills_") => {
145 Some(AgentSidebandEventCategory::Skill)
146 }
147 _ if kind.starts_with("subagent_") => Some(AgentSidebandEventCategory::Subagent),
148 _ if kind.starts_with("message_") => Some(AgentSidebandEventCategory::Message),
149 _ if kind.starts_with("usage_") => Some(AgentSidebandEventCategory::Usage),
150 _ if kind.starts_with("compact_") => Some(AgentSidebandEventCategory::Compact),
151 _ if kind.starts_with("steering_") => Some(AgentSidebandEventCategory::Steering),
152 _ if kind.starts_with("goal_") => Some(AgentSidebandEventCategory::Goal),
153 _ if kind.starts_with("background_") => Some(AgentSidebandEventCategory::Background),
154 _ if kind.starts_with("capability_") => Some(AgentSidebandEventCategory::Capability),
155 _ => None,
156 }
157 }
158}
159
160#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
162#[serde(tag = "kind", rename_all = "snake_case")]
163pub enum AgentStreamEvent {
164 RunStart {
166 run_id: RunId,
168 conversation_id: ConversationId,
170 },
171 NodeStart {
173 node: AgentExecutionNode,
175 step: usize,
177 status: RunStatus,
179 },
180 NodeComplete {
182 node: AgentExecutionNode,
184 step: usize,
186 status: RunStatus,
188 },
189 Custom {
191 event: AgentEvent,
193 },
194 ModelRequest {
196 step: usize,
198 },
199 ModelStream {
201 step: usize,
203 event: ModelResponseStreamEvent,
205 },
206 ModelResponse {
208 step: usize,
210 response: ModelResponse,
212 },
213 Checkpoint {
215 node: AgentExecutionNode,
217 step: usize,
219 },
220 Suspended {
222 node: AgentExecutionNode,
224 reason: String,
226 },
227 ToolCall {
229 step: usize,
231 call: ToolCallPart,
233 },
234 ToolReturn {
236 step: usize,
238 tool_return: ToolReturnPart,
240 },
241 OutputRetry {
243 retries: usize,
245 prompt: String,
247 },
248 SteeringGuard {
250 step: usize,
252 prompt: String,
254 },
255 RunComplete {
257 run_id: RunId,
259 output: String,
261 },
262 RunCancelled {
264 run_id: RunId,
266 reason: String,
268 },
269 RunFailed {
271 run_id: RunId,
273 error_kind: String,
275 message: String,
277 },
278}
279
280impl AgentStreamEvent {
281 #[must_use]
283 pub fn sideband_event(&self) -> Option<AgentSidebandEvent> {
284 match self {
285 Self::Custom { event } => AgentSidebandEvent::from_agent_event(event),
286 _ => None,
287 }
288 }
289}
290
291#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
293#[serde(rename_all = "snake_case")]
294pub enum AgentStreamSourceKind {
295 Subagent,
297}
298
299#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
301pub struct AgentStreamSource {
302 pub kind: AgentStreamSourceKind,
304 pub agent_id: AgentId,
306 pub agent_name: String,
308 #[serde(default, skip_serializing_if = "Option::is_none")]
310 pub task_id: Option<TaskId>,
311 #[serde(default, skip_serializing_if = "Option::is_none")]
313 pub run_id: Option<RunId>,
314 #[serde(default, skip_serializing_if = "Option::is_none")]
316 pub parent_run_id: Option<RunId>,
317 pub source_sequence: usize,
319}
320
321impl AgentStreamSource {
322 #[must_use]
324 pub fn subagent(
325 agent_id: AgentId,
326 agent_name: impl Into<String>,
327 task_id: TaskId,
328 run_id: Option<RunId>,
329 parent_run_id: Option<RunId>,
330 source_sequence: usize,
331 ) -> Self {
332 Self {
333 kind: AgentStreamSourceKind::Subagent,
334 agent_id,
335 agent_name: agent_name.into(),
336 task_id: Some(task_id),
337 run_id,
338 parent_run_id,
339 source_sequence,
340 }
341 }
342}
343
344#[derive(Clone, Debug, Default)]
346pub struct AgentStreamSink {
347 records: Arc<Mutex<Vec<AgentStreamRecord>>>,
348}
349
350impl AgentStreamSink {
351 pub fn push(&self, record: AgentStreamRecord) {
353 self.records
354 .lock()
355 .unwrap_or_else(PoisonError::into_inner)
356 .push(record);
357 }
358
359 pub fn extend(&self, records: impl IntoIterator<Item = AgentStreamRecord>) {
361 self.records
362 .lock()
363 .unwrap_or_else(PoisonError::into_inner)
364 .extend(records);
365 }
366
367 #[must_use]
369 pub fn drain(&self) -> Vec<AgentStreamRecord> {
370 self.records
371 .lock()
372 .unwrap_or_else(PoisonError::into_inner)
373 .drain(..)
374 .collect()
375 }
376
377 #[must_use]
379 pub fn is_empty(&self) -> bool {
380 self.records
381 .lock()
382 .unwrap_or_else(PoisonError::into_inner)
383 .is_empty()
384 }
385}
386
387#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
389pub struct AgentStreamRecord {
390 pub sequence: usize,
392 #[serde(default, skip_serializing_if = "Option::is_none")]
394 pub source: Option<AgentStreamSource>,
395 pub event: AgentStreamEvent,
397}
398
399impl starweaver_core::VersionedRecord for AgentStreamRecord {
400 const SCHEMA: &'static str = "starweaver.runtime.stream_record";
401 const ALLOW_BARE_V0: bool = true;
402}
403
404impl AgentStreamRecord {
405 #[must_use]
407 pub const fn new(sequence: usize, event: AgentStreamEvent) -> Self {
408 Self {
409 sequence,
410 source: None,
411 event,
412 }
413 }
414
415 #[must_use]
421 pub fn is_subagent_steering_event(&self) -> bool {
422 if !matches!(
423 self.source.as_ref().map(|source| &source.kind),
424 Some(AgentStreamSourceKind::Subagent)
425 ) {
426 return false;
427 }
428 match &self.event {
429 AgentStreamEvent::SteeringGuard { .. } => true,
430 AgentStreamEvent::Custom { event } => {
431 let normalized = event.kind.to_ascii_lowercase().replace(['.', '-'], "_");
432 [
433 "steering_submitted",
434 "steer_submitted",
435 "steering_received",
436 "steer_received",
437 "steering_ack",
438 "steer_ack",
439 ]
440 .iter()
441 .any(|candidate| {
442 normalized == *candidate || normalized.ends_with(&format!("_{candidate}"))
443 })
444 }
445 _ => false,
446 }
447 }
448
449 #[must_use]
451 pub fn with_source(mut self, source: AgentStreamSource) -> Self {
452 self.source = Some(source);
453 self
454 }
455
456 #[must_use]
458 pub const fn with_sequence(mut self, sequence: usize) -> Self {
459 self.sequence = sequence;
460 self
461 }
462
463 pub fn to_raw_json(&self) -> serde_json::Result<serde_json::Value> {
469 serde_json::to_value(self)
470 }
471}