Skip to main content

uptrakit_openapi_client/
events_stream.rs

1//! Typed SSE streaming method for admin events.
2//!
3//! Provides [`UptrakitClient::stream_events`] which connects to the
4//! `GET /api/v1/events/stream` endpoint and returns a typed stream of
5//! admin events for the authenticated user's tenant.
6
7use crate::sse::{self, RawSseEvent, SseError};
8use crate::{Result, UptrakitClient};
9use rootcause::prelude::*;
10use uuid::Uuid;
11
12/// A typed SSE event from the admin events stream.
13///
14/// Mirrors [`AdminEvent`](uptrakit_web_api_types::events::AdminEvent) variants
15/// with an additional [`Unknown`](Self::Unknown) catch-all for forward
16/// compatibility when the server emits new event types.
17#[derive(Debug, Clone)]
18pub enum AdminSseEvent {
19    /// A host's metadata was updated.
20    HostUpdated { id: Uuid },
21    /// A new host was created.
22    HostCreated { id: Uuid },
23    /// A host was deactivated / deleted.
24    HostDeleted { id: Uuid },
25    /// A service's status changed.
26    ServiceStatusChanged { id: Uuid, status: String },
27    /// A software item was updated.
28    SoftwareItemUpdated { id: Uuid },
29    /// A new software item was created.
30    SoftwareItemCreated { id: Uuid },
31    /// A version check completed for a host + software item pair.
32    VersionCheckCompleted {
33        host_id: Uuid,
34        software_item_id: Uuid,
35    },
36    /// A software update was created and dispatched to the agent.
37    UpdateTriggered {
38        update_history_id: Uuid,
39        host_id: Uuid,
40        software_item_id: Uuid,
41    },
42    /// A software update started executing.
43    UpdateStarted {
44        update_history_id: Uuid,
45        host_id: Uuid,
46        software_item_id: Uuid,
47        /// Whether the update was dispatched in interactive mode (PTY allocated).
48        interactive: bool,
49    },
50    /// A software update completed.
51    UpdateCompleted {
52        update_history_id: Uuid,
53        host_id: Uuid,
54        software_item_id: Uuid,
55        status: String,
56    },
57    /// Autodiscovery completed for a host.
58    DiscoveryCompleted { host_id: Uuid },
59    /// Host packages changed.
60    HostPackagesChanged { host_id: Uuid },
61    /// A batch host package update completed.
62    BatchHostPackageUpdateCompleted { host_id: Uuid },
63    /// A system service's status changed.
64    SystemServiceStatusChanged { id: Uuid, status: String },
65    /// A scheduled task completed execution.
66    SchedulerTaskCompleted { task_id: Uuid },
67    /// All tenant data was reset.
68    DataReset,
69    /// An unrecognised event type from a newer server version.
70    Unknown { event_type: String, data: String },
71}
72
73/// Errors specific to admin event streaming.
74#[derive(Debug, thiserror::Error)]
75pub enum StreamError {
76    #[error("SSE transport error: {0}")]
77    Sse(#[from] SseError),
78
79    #[error("failed to parse SSE event data: {0}")]
80    Parse(#[from] serde_json::Error),
81}
82
83impl UptrakitClient {
84    /// Connect to the admin events SSE stream and return a stream of typed events.
85    ///
86    /// The returned stream yields [`AdminSseEvent`] values for the authenticated
87    /// user's tenant. The stream stays open indefinitely (server pushes events as
88    /// state changes occur) and should be cancelled by the caller when no longer
89    /// needed.
90    ///
91    /// Uses an 86400s (24h) timeout, matching other long-lived SSE connections.
92    #[cfg_attr(feature = "tracing", tracing::instrument(skip_all))]
93    pub async fn stream_events(
94        &self,
95    ) -> Result<impl futures_util::Stream<Item = std::result::Result<AdminSseEvent, StreamError>>>
96    {
97        let url = format!("{}{}", self.base_url, crate::paths::events::STREAM);
98
99        let req = self
100            .http
101            .get(&url)
102            .bearer_auth(self.token_or_err()?)
103            .header("Accept", "text/event-stream")
104            .timeout(std::time::Duration::from_secs(86400));
105
106        let resp = req.send().await.context_to()?;
107
108        let status = resp.status();
109        if status == reqwest::StatusCode::UNAUTHORIZED {
110            bail!(crate::ClientError::NotAuthenticated);
111        }
112        if status.is_client_error() || status.is_server_error() {
113            let text = resp.text().await.context_to()?;
114            let message = crate::extract_error_message(&text);
115            bail!(crate::ClientError::Api { status, message });
116        }
117
118        let raw_stream = sse::parse_sse_stream(resp);
119
120        let typed_stream = futures_util::StreamExt::filter_map(raw_stream, |result| async move {
121            match result {
122                Ok(event) => Some(parse_typed_event(event)),
123                Err(e) => Some(Err(StreamError::Sse(e))),
124            }
125        });
126
127        Ok(typed_stream)
128    }
129}
130
131/// Helper for parsing a JSON `data` field with a single `id` key.
132fn parse_id(data: &str) -> std::result::Result<Uuid, serde_json::Error> {
133    #[derive(serde::Deserialize)]
134    struct Id {
135        id: Uuid,
136    }
137    serde_json::from_str::<Id>(data).map(|v| v.id)
138}
139
140/// Helper for parsing a JSON `data` field with `id` and `status` keys.
141fn parse_id_status(data: &str) -> std::result::Result<(Uuid, String), serde_json::Error> {
142    #[derive(serde::Deserialize)]
143    struct IdStatus {
144        id: Uuid,
145        status: String,
146    }
147    serde_json::from_str::<IdStatus>(data).map(|v| (v.id, v.status))
148}
149
150/// Helper for parsing a JSON `data` field with a single `host_id` key.
151fn parse_host_id(data: &str) -> std::result::Result<Uuid, serde_json::Error> {
152    #[derive(serde::Deserialize)]
153    struct HostId {
154        host_id: Uuid,
155    }
156    serde_json::from_str::<HostId>(data).map(|v| v.host_id)
157}
158
159/// Parse a raw SSE event into a typed [`AdminSseEvent`].
160///
161/// Unknown event types are returned as [`AdminSseEvent::Unknown`] for forward
162/// compatibility — the client never drops events from a newer server.
163fn parse_typed_event(event: RawSseEvent) -> std::result::Result<AdminSseEvent, StreamError> {
164    match event.event_type.as_str() {
165        "host_updated" => Ok(AdminSseEvent::HostUpdated {
166            id: parse_id(&event.data)?,
167        }),
168        "host_created" => Ok(AdminSseEvent::HostCreated {
169            id: parse_id(&event.data)?,
170        }),
171        "host_deleted" => Ok(AdminSseEvent::HostDeleted {
172            id: parse_id(&event.data)?,
173        }),
174        "service_status_changed" => {
175            let (id, status) = parse_id_status(&event.data)?;
176            Ok(AdminSseEvent::ServiceStatusChanged { id, status })
177        }
178        "software_item_updated" => Ok(AdminSseEvent::SoftwareItemUpdated {
179            id: parse_id(&event.data)?,
180        }),
181        "software_item_created" => Ok(AdminSseEvent::SoftwareItemCreated {
182            id: parse_id(&event.data)?,
183        }),
184        "version_check_completed" => {
185            #[derive(serde::Deserialize)]
186            struct Payload {
187                host_id: Uuid,
188                software_item_id: Uuid,
189            }
190            let p: Payload = serde_json::from_str(&event.data)?;
191            Ok(AdminSseEvent::VersionCheckCompleted {
192                host_id: p.host_id,
193                software_item_id: p.software_item_id,
194            })
195        }
196        "update_triggered" => {
197            #[derive(serde::Deserialize)]
198            struct Payload {
199                update_history_id: Uuid,
200                host_id: Uuid,
201                software_item_id: Uuid,
202            }
203            let p: Payload = serde_json::from_str(&event.data)?;
204            Ok(AdminSseEvent::UpdateTriggered {
205                update_history_id: p.update_history_id,
206                host_id: p.host_id,
207                software_item_id: p.software_item_id,
208            })
209        }
210        "update_started" => {
211            #[derive(serde::Deserialize)]
212            struct Payload {
213                update_history_id: Uuid,
214                host_id: Uuid,
215                software_item_id: Uuid,
216                #[serde(default)]
217                interactive: bool,
218            }
219            let p: Payload = serde_json::from_str(&event.data)?;
220            Ok(AdminSseEvent::UpdateStarted {
221                update_history_id: p.update_history_id,
222                host_id: p.host_id,
223                software_item_id: p.software_item_id,
224                interactive: p.interactive,
225            })
226        }
227        "update_completed" => {
228            #[derive(serde::Deserialize)]
229            struct Payload {
230                update_history_id: Uuid,
231                host_id: Uuid,
232                software_item_id: Uuid,
233                status: String,
234            }
235            let p: Payload = serde_json::from_str(&event.data)?;
236            Ok(AdminSseEvent::UpdateCompleted {
237                update_history_id: p.update_history_id,
238                host_id: p.host_id,
239                software_item_id: p.software_item_id,
240                status: p.status,
241            })
242        }
243        "discovery_completed" => Ok(AdminSseEvent::DiscoveryCompleted {
244            host_id: parse_host_id(&event.data)?,
245        }),
246        "host_packages_changed" => Ok(AdminSseEvent::HostPackagesChanged {
247            host_id: parse_host_id(&event.data)?,
248        }),
249        "batch_host_package_update_completed" => {
250            Ok(AdminSseEvent::BatchHostPackageUpdateCompleted {
251                host_id: parse_host_id(&event.data)?,
252            })
253        }
254        "system_service_status_changed" => {
255            let (id, status) = parse_id_status(&event.data)?;
256            Ok(AdminSseEvent::SystemServiceStatusChanged { id, status })
257        }
258        "scheduler_task_completed" => {
259            #[derive(serde::Deserialize)]
260            struct Payload {
261                task_id: Uuid,
262            }
263            let p: Payload = serde_json::from_str(&event.data)?;
264            Ok(AdminSseEvent::SchedulerTaskCompleted { task_id: p.task_id })
265        }
266        "data_reset" => Ok(AdminSseEvent::DataReset),
267        _ => Ok(AdminSseEvent::Unknown {
268            event_type: event.event_type,
269            data: event.data,
270        }),
271    }
272}
273
274#[cfg(test)]
275mod tests {
276    use super::*;
277    use crate::sse::RawSseEvent;
278
279    fn make_event(event_type: &str, data: &str) -> RawSseEvent {
280        RawSseEvent {
281            event_type: event_type.to_string(),
282            data: data.to_string(),
283            id: None,
284        }
285    }
286
287    #[test]
288    fn parse_host_updated() {
289        let event = make_event(
290            "host_updated",
291            r#"{"id":"550e8400-e29b-41d4-a716-446655440000"}"#,
292        );
293        let result = parse_typed_event(event).unwrap();
294        assert!(matches!(result, AdminSseEvent::HostUpdated { id } if !id.is_nil()));
295    }
296
297    #[test]
298    fn parse_host_created() {
299        let event = make_event(
300            "host_created",
301            r#"{"id":"550e8400-e29b-41d4-a716-446655440000"}"#,
302        );
303        let result = parse_typed_event(event).unwrap();
304        assert!(matches!(result, AdminSseEvent::HostCreated { .. }));
305    }
306
307    #[test]
308    fn parse_host_deleted() {
309        let event = make_event(
310            "host_deleted",
311            r#"{"id":"550e8400-e29b-41d4-a716-446655440000"}"#,
312        );
313        let result = parse_typed_event(event).unwrap();
314        assert!(matches!(result, AdminSseEvent::HostDeleted { .. }));
315    }
316
317    #[test]
318    fn parse_service_status_changed() {
319        let event = make_event(
320            "service_status_changed",
321            r#"{"id":"550e8400-e29b-41d4-a716-446655440000","status":"approved"}"#,
322        );
323        let result = parse_typed_event(event).unwrap();
324        assert!(
325            matches!(result, AdminSseEvent::ServiceStatusChanged { status, .. } if status == "approved")
326        );
327    }
328
329    #[test]
330    fn parse_software_item_updated() {
331        let event = make_event(
332            "software_item_updated",
333            r#"{"id":"550e8400-e29b-41d4-a716-446655440000"}"#,
334        );
335        let result = parse_typed_event(event).unwrap();
336        assert!(matches!(result, AdminSseEvent::SoftwareItemUpdated { .. }));
337    }
338
339    #[test]
340    fn parse_software_item_created() {
341        let event = make_event(
342            "software_item_created",
343            r#"{"id":"550e8400-e29b-41d4-a716-446655440000"}"#,
344        );
345        let result = parse_typed_event(event).unwrap();
346        assert!(matches!(result, AdminSseEvent::SoftwareItemCreated { .. }));
347    }
348
349    #[test]
350    fn parse_version_check_completed() {
351        let event = make_event(
352            "version_check_completed",
353            r#"{"host_id":"550e8400-e29b-41d4-a716-446655440001","software_item_id":"550e8400-e29b-41d4-a716-446655440002"}"#,
354        );
355        let result = parse_typed_event(event).unwrap();
356        assert!(matches!(
357            result,
358            AdminSseEvent::VersionCheckCompleted { .. }
359        ));
360    }
361
362    #[test]
363    fn parse_update_triggered() {
364        let event = make_event(
365            "update_triggered",
366            r#"{"update_history_id":"550e8400-e29b-41d4-a716-446655440001","host_id":"550e8400-e29b-41d4-a716-446655440002","software_item_id":"550e8400-e29b-41d4-a716-446655440003"}"#,
367        );
368        let result = parse_typed_event(event).unwrap();
369        assert!(matches!(result, AdminSseEvent::UpdateTriggered { .. }));
370    }
371
372    #[test]
373    fn parse_update_started() {
374        let event = make_event(
375            "update_started",
376            r#"{"update_history_id":"550e8400-e29b-41d4-a716-446655440001","host_id":"550e8400-e29b-41d4-a716-446655440002","software_item_id":"550e8400-e29b-41d4-a716-446655440003","interactive":true}"#,
377        );
378        let result = parse_typed_event(event).unwrap();
379        assert!(matches!(
380            result,
381            AdminSseEvent::UpdateStarted {
382                interactive: true,
383                ..
384            }
385        ));
386    }
387
388    #[test]
389    fn parse_update_started_without_interactive_defaults_false() {
390        // Older server versions may not send the `interactive` field.
391        let event = make_event(
392            "update_started",
393            r#"{"update_history_id":"550e8400-e29b-41d4-a716-446655440001","host_id":"550e8400-e29b-41d4-a716-446655440002","software_item_id":"550e8400-e29b-41d4-a716-446655440003"}"#,
394        );
395        let result = parse_typed_event(event).unwrap();
396        assert!(matches!(
397            result,
398            AdminSseEvent::UpdateStarted {
399                interactive: false,
400                ..
401            }
402        ));
403    }
404
405    #[test]
406    fn parse_update_completed() {
407        let event = make_event(
408            "update_completed",
409            r#"{"update_history_id":"550e8400-e29b-41d4-a716-446655440001","host_id":"550e8400-e29b-41d4-a716-446655440002","software_item_id":"550e8400-e29b-41d4-a716-446655440003","status":"completed"}"#,
410        );
411        let result = parse_typed_event(event).unwrap();
412        assert!(
413            matches!(result, AdminSseEvent::UpdateCompleted { status, .. } if status == "completed")
414        );
415    }
416
417    #[test]
418    fn parse_discovery_completed() {
419        let event = make_event(
420            "discovery_completed",
421            r#"{"host_id":"550e8400-e29b-41d4-a716-446655440000"}"#,
422        );
423        let result = parse_typed_event(event).unwrap();
424        assert!(matches!(result, AdminSseEvent::DiscoveryCompleted { .. }));
425    }
426
427    #[test]
428    fn parse_host_packages_changed() {
429        let event = make_event(
430            "host_packages_changed",
431            r#"{"host_id":"550e8400-e29b-41d4-a716-446655440000"}"#,
432        );
433        let result = parse_typed_event(event).unwrap();
434        assert!(matches!(result, AdminSseEvent::HostPackagesChanged { .. }));
435    }
436
437    #[test]
438    fn parse_batch_host_package_update_completed() {
439        let event = make_event(
440            "batch_host_package_update_completed",
441            r#"{"host_id":"550e8400-e29b-41d4-a716-446655440000"}"#,
442        );
443        let result = parse_typed_event(event).unwrap();
444        assert!(matches!(
445            result,
446            AdminSseEvent::BatchHostPackageUpdateCompleted { .. }
447        ));
448    }
449
450    #[test]
451    fn parse_system_service_status_changed() {
452        let event = make_event(
453            "system_service_status_changed",
454            r#"{"id":"550e8400-e29b-41d4-a716-446655440000","status":"rejected"}"#,
455        );
456        let result = parse_typed_event(event).unwrap();
457        assert!(
458            matches!(result, AdminSseEvent::SystemServiceStatusChanged { status, .. } if status == "rejected")
459        );
460    }
461
462    #[test]
463    fn parse_scheduler_task_completed() {
464        let event = make_event(
465            "scheduler_task_completed",
466            r#"{"task_id":"550e8400-e29b-41d4-a716-446655440000"}"#,
467        );
468        let result = parse_typed_event(event).unwrap();
469        assert!(matches!(
470            result,
471            AdminSseEvent::SchedulerTaskCompleted { .. }
472        ));
473    }
474
475    #[test]
476    fn parse_data_reset() {
477        let event = make_event("data_reset", "{}");
478        let result = parse_typed_event(event).unwrap();
479        assert!(matches!(result, AdminSseEvent::DataReset));
480    }
481
482    #[test]
483    fn parse_unknown_event_returns_unknown() {
484        let event = make_event("future_event", r#"{"foo":"bar"}"#);
485        let result = parse_typed_event(event).unwrap();
486        assert!(
487            matches!(result, AdminSseEvent::Unknown { event_type, .. } if event_type == "future_event")
488        );
489    }
490
491    #[test]
492    fn parse_malformed_data_returns_error() {
493        let event = make_event("host_updated", "not json");
494        assert!(parse_typed_event(event).is_err());
495    }
496}