Skip to main content

relay_knowledge/api/
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)]
76mod tests {
77    use serde_json::json;
78
79    use super::*;
80    use crate::{
81        api::{InterfaceKind, RequestContext},
82        domain::GraphVersion,
83    };
84
85    #[test]
86    fn project_status_event_only_attaches_payload_to_item() {
87        let context = RequestContext::with_ids(InterfaceKind::Cli, "req", "trace");
88        let response = ProjectStatusResponse {
89            project_name: "relay-knowledge".to_owned(),
90            metadata: ApiMetadata::graph_only(&context, GraphVersion::ZERO),
91            runtime: RuntimeStatus {
92                config_dir: "/config".to_owned(),
93                data_dir: "/data".to_owned(),
94                state_dir: "/state".to_owned(),
95                cache_dir: "/cache".to_owned(),
96                log_dir: "/logs".to_owned(),
97                temp_dir: "/tmp".to_owned(),
98                runtime_dir: "/run".to_owned(),
99                service_dir: "/service".to_owned(),
100                storage_topology: "single_sqlite".to_owned(),
101                http_bind: "127.0.0.1:8791".to_owned(),
102                http_request_timeout_ms: 30000,
103                http_graceful_shutdown_timeout_ms: 10000,
104                http_max_request_body_bytes: 1024,
105                http_proxy_configured: false,
106                http_no_proxy_rules: 0,
107                http_ssl_verify: true,
108                qos_max_connections: 1,
109                qos_max_in_flight_requests: 1,
110                qos_max_queue_depth: 1,
111                qos_current_connections: 0,
112                qos_current_in_flight_requests: 0,
113                qos_current_queued_requests: 0,
114                qos_admitted_total: 0,
115                qos_queued_total: 0,
116                qos_rejected_total: 0,
117                qos_timed_out_total: 0,
118                qos_cancelled_total: 0,
119                qos_dropped_total: 0,
120                worker_embedding_endpoint_configured: false,
121                worker_ocr_endpoint_configured: false,
122                worker_vision_endpoint_configured: false,
123                worker_extractor_endpoint_configured: false,
124                worker_max_in_flight: 2,
125                code_index_max_in_flight: 2,
126                silent_updates_enabled: false,
127                file_index_enabled: false,
128                file_index_root_count: 0,
129                file_index_max_depth: 32,
130                file_index_max_file_bytes: 512 * 1024 * 1024,
131                file_index_scan_interval_ms: 900_000,
132                file_index_scan_timeout_ms: 300_000,
133                file_index_max_files_per_root: 50_000,
134                file_query_timeout_ms: 750,
135                semantic_backend_mode: "local".to_owned(),
136                vector_backend_mode: "local".to_owned(),
137                rerank_backend_mode: "local".to_owned(),
138                rerank_model: Some("relay-local-deterministic-rerank-v1".to_owned()),
139                rerank_candidate_multiplier: 4,
140                rerank_max_candidates: 64,
141                rerank_timeout_ms: 100,
142                embedding_provider: None,
143                embedding_base_url: None,
144                embedding_api_key_configured: false,
145                text_embedding_model: "relay-local-hash-ann-v1".to_owned(),
146                image_embedding_model: "relay-local-image-hash-v1".to_owned(),
147                embedding_dimension: 16,
148                embedding_batch_size: None,
149                embedding_timeout_ms: None,
150                embedding_max_concurrency: None,
151                model_profiles: crate::model_provider::ModelProfileRuntimeSummary {
152                    loaded: true,
153                    profile_count: 0,
154                    default_profile: None,
155                    error: None,
156                },
157                telemetry: crate::observability::ObservabilityRuntime::new(
158                    crate::observability::TelemetryConfig::from_environment(
159                        &crate::env::TelemetryEnvOverrides::default(),
160                    ),
161                )
162                .status(),
163            },
164        };
165
166        let started =
167            ApiStreamEvent::project_status(StreamEventKind::Started, &response, Some("starting"));
168        let item = ApiStreamEvent::project_status(StreamEventKind::Item, &response, None);
169
170        assert_eq!(started.project_name, None);
171        assert_eq!(started.message, Some("starting".to_owned()));
172        assert_eq!(item.project_name, Some("relay-knowledge".to_owned()));
173        assert!(item.runtime.is_some());
174    }
175
176    #[test]
177    fn operation_event_carries_generic_payload() {
178        let context = RequestContext::with_ids(InterfaceKind::Api, "req", "trace");
179        let metadata = ApiMetadata::graph_only(&context, GraphVersion::ZERO);
180
181        let event = ApiStreamEvent::operation(
182            StreamEventKind::Item,
183            "health",
184            metadata,
185            None,
186            Some(json!({"healthy": true})),
187        );
188
189        assert_eq!(event.operation, "health");
190        assert_eq!(event.payload, Some(json!({"healthy": true})));
191    }
192}