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 },
42 UpdateStarted {
44 update_history_id: Uuid,
45 host_id: Uuid,
46 software_item_id: Uuid,
47 interactive: bool,
49 },
50 UpdateCompleted {
52 update_history_id: Uuid,
53 host_id: Uuid,
54 software_item_id: Uuid,
55 status: String,
56 },
57 DiscoveryCompleted { host_id: Uuid },
59 HostPackagesChanged { host_id: Uuid },
61 BatchHostPackageUpdateCompleted { host_id: Uuid },
63 SystemServiceStatusChanged { id: Uuid, status: String },
65 SchedulerTaskCompleted { task_id: Uuid },
67 DataReset,
69 Unknown { event_type: String, data: String },
71}
72
73#[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 #[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
131fn 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
140fn 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
150fn 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
159fn 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 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}