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