use std::task::{Poll, ready};
use moq_net::{PathOwned, announce, broadcast, origin};
use crate::path::{Kind, parse};
#[derive(Clone)]
pub struct Event {
pub identity: PathOwned,
pub kind: Kind,
pub path: PathOwned,
pub broadcast: Option<broadcast::Consumer>,
}
impl std::fmt::Debug for Event {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Event")
.field("identity", &self.identity)
.field("kind", &self.kind)
.field("path", &self.path)
.field("online", &self.broadcast.is_some())
.finish()
}
}
pub struct Room {
announced: announce::Consumer,
local: Option<PathOwned>,
}
impl std::fmt::Debug for Room {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Room")
.field("local", &self.local)
.finish_non_exhaustive()
}
}
impl Room {
pub fn new(origin: &origin::Consumer, local: Option<PathOwned>) -> Self {
Self {
announced: origin.announced(),
local,
}
}
pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<Event>> {
loop {
let Some(update) = ready!(self.announced.poll_next(waiter)) else {
return Poll::Ready(None);
};
let Some(parsed) = parse(&update.path) else {
continue;
};
if self.local.as_ref().is_some_and(|id| *id == parsed.identity) {
continue;
}
return Poll::Ready(Some(Event {
identity: parsed.identity,
kind: parsed.kind,
path: update.path,
broadcast: update.broadcast,
}));
}
}
pub async fn next(&mut self) -> Option<Event> {
kio::wait(|waiter| self.poll_next(waiter)).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use moq_net::{Origin, Path, broadcast::Route};
#[tokio::test]
async fn yields_camera_and_skips_local_and_unknown() {
let origin = Origin::random().produce();
let local = Path::new("alice").to_owned();
let mut room = Room::new(&origin.consume(), Some(local));
let _alice = origin
.create_broadcast("alice/camera", Route::announced())
.expect("alice camera");
let _bob = origin
.create_broadcast("bob/camera", Route::announced())
.expect("bob camera");
let _noise = origin
.create_broadcast("bob/chat", Route::announced())
.expect("unknown kind");
let event = room.next().await.expect("bob camera");
assert_eq!(event.identity.as_str(), "bob");
assert_eq!(event.kind, Kind::Camera);
assert!(event.broadcast.is_some());
drop(_bob);
let gone = room.next().await.expect("bob unannounce");
assert_eq!(gone.identity.as_str(), "bob");
assert_eq!(gone.kind, Kind::Camera);
assert!(gone.broadcast.is_none());
}
#[tokio::test]
async fn yields_screen() {
let origin = Origin::random().produce();
let mut room = Room::new(&origin.consume(), None);
let _screen = origin
.create_broadcast("bob/screen", Route::announced())
.expect("bob screen");
let event = room.next().await.expect("bob screen");
assert_eq!(event.identity.as_str(), "bob");
assert_eq!(event.kind, Kind::Screen);
assert!(event.broadcast.is_some());
}
}