pulse-pixelstream-types 0.14.1

Shared Myko entity and command types for the Pulse Pixelstream recording cell.
Documentation
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};

/// One live subscription per stream: who's present + who holds control.
#[myko_report_output]
pub struct StreamPresence {
    pub viewers: Vec<Viewer>,
    pub control_holder: Option<String>, // viewer_id holding the wheel, or None
}

#[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();

        // Reactive join: recomputes whenever a viewer or the lock changes.
        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()),
            })
        })
    }
}