1use crate::sse::{self, RawSseEvent, SseError};
8use crate::{Result, UptrakitClient};
9use rootcause::prelude::*;
10use uuid::Uuid;
11
12#[derive(Debug, Clone)]
18pub enum AdminSseEvent {
19 HostUpdated { id: Uuid },
21 HostCreated { id: Uuid },
23 HostDeleted { id: Uuid },
25 ServiceStatusChanged { id: Uuid, status: String },
27 SoftwareItemUpdated { id: Uuid },
29 SoftwareItemCreated { id: Uuid },
31 VersionCheckCompleted {
33 host_id: Uuid,
34 software_item_id: Uuid,
35 },
36 UpdateTriggered {
38 update_history_id: Uuid,
39 host_id: Uuid,
40 software_item_id: Uuid,
41 status: String,
43 },
44 UpdateStarted {
46 update_history_id: Uuid,
47 host_id: Uuid,
48 software_item_id: Uuid,
49 interactive: bool,
51 },
52 UpdateCompleted {
54 update_history_id: Uuid,
55 host_id: Uuid,
56 software_item_id: Uuid,
57 status: String,
58 },
59 DiscoveryCompleted { host_id: Uuid },
61 HostPackagesChanged { host_id: Uuid },
63 BatchHostPackageUpdateCompleted { host_id: Uuid },
65 SystemServiceStatusChanged { id: Uuid, status: String },
67 SchedulerTaskCompleted { task_id: Uuid },
69 DataReset,
71 Unknown { event_type: String, data: String },
73}
74
75#[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 #[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
133fn 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
142fn 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
152fn 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
161fn 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 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}