use myko::prelude::*;
use std::sync::Arc;
use myko::hyphae::{Definite, JoinExt, MapExt, Materialize};
use myko::report::{ReportContext, ReportHandler};
use myko::{myko_report, myko_report_output};
use crate::control_lock::{ControlLockQuery, GetControlLocksByQuery};
use crate::stream::StreamId;
use crate::viewer::{GetViewersByQuery, Viewer, ViewerQuery};
#[myko_report_output]
pub struct StreamPresence {
pub viewers: Vec<Viewer>,
pub control_holder: Option<String>, }
#[myko_report(StreamPresence)]
pub struct WatchStream {
pub stream_id: StreamId,
}
impl ReportHandler for WatchStream {
type Output = StreamPresence;
fn compute(&self, ctx: ReportContext) -> impl Materialize<Arc<Self::Output>, Definite> {
let sid = self.stream_id.clone();
let viewers = ctx
.query_map(GetViewersByQuery(ViewerQuery {
stream_id: Some(IdFilter::Eq(sid.clone())),
..Default::default()
}))
.items()
.materialize();
let locks = ctx
.query_map(GetControlLocksByQuery(ControlLockQuery {
stream_id: Some(IdFilter::Eq(sid)),
..Default::default()
}))
.items()
.materialize();
viewers.join(locks).map(|(viewers, locks)| {
Arc::new(StreamPresence {
viewers: viewers.iter().map(|v| (**v).clone()).collect(),
control_holder: locks.first().map(|l| l.viewer_id.clone()),
})
})
}
}