Skip to main content

starweaver_runtime/
stream.rs

1//! Runtime stream-result type and compatibility exports for raw stream protocol records.
2
3use serde::{Deserialize, Serialize};
4
5use crate::AgentResult;
6
7pub use starweaver_stream::{
8    AgentSidebandEvent, AgentSidebandEventCategory, AgentStreamEvent, AgentStreamRecord,
9    AgentStreamSink, AgentStreamSource, AgentStreamSourceKind,
10};
11
12/// Result returned by collection-based stream runs.
13#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
14pub struct AgentStreamResult {
15    /// Final agent result.
16    pub result: AgentResult,
17    /// Events captured while the run progressed.
18    pub events: Vec<AgentStreamRecord>,
19}
20
21impl AgentStreamResult {
22    /// Return captured stream records.
23    #[must_use]
24    pub fn events(&self) -> &[AgentStreamRecord] {
25        &self.events
26    }
27
28    /// Return the final result.
29    #[must_use]
30    pub const fn result(&self) -> &AgentResult {
31        &self.result
32    }
33
34    /// Project captured raw runtime stream records into JSON values.
35    ///
36    /// # Errors
37    ///
38    /// Returns a serialization error if a nested event payload cannot be encoded.
39    pub fn raw_json_records(&self) -> serde_json::Result<Vec<serde_json::Value>> {
40        self.events
41            .iter()
42            .map(AgentStreamRecord::to_raw_json)
43            .collect()
44    }
45}
46
47pub(crate) fn push_stream_event(
48    events: &mut Option<&mut Vec<AgentStreamRecord>>,
49    event: AgentStreamEvent,
50) {
51    if let Some(events) = events.as_deref_mut() {
52        events.push(AgentStreamRecord::new(events.len(), event));
53    }
54}
55
56pub(crate) fn push_stream_record(
57    events: &mut Option<&mut Vec<AgentStreamRecord>>,
58    record: AgentStreamRecord,
59) {
60    if let Some(events) = events.as_deref_mut() {
61        events.push(record.with_sequence(events.len()));
62    }
63}