Skip to main content

pulse_pixelstream_types/
reports.rs

1use myko::prelude::*;
2use std::sync::Arc;
3
4use myko::hyphae::{Definite, JoinExt, MapExt, Materialize};
5use myko::report::{ReportContext, ReportHandler};
6use myko::{myko_report, myko_report_output};
7
8use crate::control_lock::{ControlLockQuery, GetControlLocksByQuery};
9use crate::stream::StreamId;
10use crate::viewer::{GetViewersByQuery, Viewer, ViewerQuery};
11
12/// One live subscription per stream: who's present + who holds control.
13#[myko_report_output]
14pub struct StreamPresence {
15    pub viewers: Vec<Viewer>,
16    pub control_holder: Option<String>, // viewer_id holding the wheel, or None
17}
18
19#[myko_report(StreamPresence)]
20pub struct WatchStream {
21    pub stream_id: StreamId,
22}
23
24impl ReportHandler for WatchStream {
25    type Output = StreamPresence;
26
27    fn compute(&self, ctx: ReportContext) -> impl Materialize<Arc<Self::Output>, Definite> {
28        let sid = self.stream_id.clone();
29        let viewers = ctx
30            .query_map(GetViewersByQuery(ViewerQuery {
31                stream_id: Some(IdFilter::Eq(sid.clone())),
32                ..Default::default()
33            }))
34            .items()
35            .materialize();
36        let locks = ctx
37            .query_map(GetControlLocksByQuery(ControlLockQuery {
38                stream_id: Some(IdFilter::Eq(sid)),
39                ..Default::default()
40            }))
41            .items()
42            .materialize();
43
44        // Reactive join: recomputes whenever a viewer or the lock changes.
45        viewers.join(locks).map(|(viewers, locks)| {
46            Arc::new(StreamPresence {
47                viewers: viewers.iter().map(|v| (**v).clone()).collect(),
48                control_holder: locks.first().map(|l| l.viewer_id.clone()),
49            })
50        })
51    }
52}