Skip to main content

relay_knowledge/api/contracts/
stream.rs

1use serde::{Deserialize, Serialize};
2
3use super::{ApiMetadata, ErrorKind, ProjectStatusResponse, RuntimeStatus};
4
5/// Stream event categories for newline-delimited JSON output.
6#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
7#[serde(rename_all = "kebab-case")]
8pub enum StreamEventKind {
9    Started,
10    Progress,
11    Item,
12    Completed,
13    Failed,
14}
15
16/// A single streaming API event. Each serialized event is one NDJSON line.
17#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
18pub struct ApiStreamEvent {
19    pub event: StreamEventKind,
20    pub operation: String,
21    #[serde(skip_serializing_if = "Option::is_none")]
22    pub message: Option<String>,
23    #[serde(skip_serializing_if = "Option::is_none")]
24    pub project_name: Option<String>,
25    #[serde(skip_serializing_if = "Option::is_none")]
26    pub runtime: Option<RuntimeStatus>,
27    #[serde(skip_serializing_if = "Option::is_none")]
28    pub payload: Option<serde_json::Value>,
29    #[serde(skip_serializing_if = "Option::is_none")]
30    pub error_kind: Option<ErrorKind>,
31    #[serde(skip_serializing_if = "Option::is_none")]
32    pub metadata: Option<ApiMetadata>,
33}
34
35impl ApiStreamEvent {
36    /// Creates a stream event for the project status operation.
37    pub fn project_status(
38        event: StreamEventKind,
39        response: &ProjectStatusResponse,
40        message: Option<&str>,
41    ) -> Self {
42        Self {
43            event,
44            operation: "project.status".to_owned(),
45            message: message.map(str::to_owned),
46            project_name: (event == StreamEventKind::Item).then(|| response.project_name.clone()),
47            runtime: (event == StreamEventKind::Item).then(|| response.runtime.clone()),
48            payload: None,
49            error_kind: None,
50            metadata: Some(response.metadata.clone()),
51        }
52    }
53
54    /// Creates a generic streaming event for non-status operations.
55    pub fn operation(
56        event: StreamEventKind,
57        operation: impl Into<String>,
58        metadata: ApiMetadata,
59        message: Option<&str>,
60        payload: Option<serde_json::Value>,
61    ) -> Self {
62        Self {
63            event,
64            operation: operation.into(),
65            message: message.map(str::to_owned),
66            project_name: None,
67            runtime: None,
68            payload,
69            error_kind: None,
70            metadata: Some(metadata),
71        }
72    }
73}
74
75#[cfg(test)]
76#[path = "stream_tests.rs"]
77mod tests;