use std::time::Duration;
use crate::emulator::Button;
#[derive(serde::Deserialize, Debug)]
struct RawTimestamp {
label: String,
ts: f64,
}
#[derive(Debug)]
pub struct TimestampEntry {
pub label: String,
pub ts: Duration,
}
#[derive(serde::Deserialize, Debug)]
#[serde(tag = "type")]
enum RawCommand {
#[serde(rename = "buttons")]
Buttons {
#[serde(default)]
buttons: Vec<Button>,
#[serde(default)]
timestamps: Vec<RawTimestamp>,
},
#[serde(rename = "reset")]
Reset {},
}
#[derive(Debug)]
pub enum Command {
Buttons {
buttons: Vec<Button>,
viewer_id: String,
timestamps: Vec<TimestampEntry>,
},
Reset,
ViewerLeft {
viewer_id: String,
},
}
pub async fn handle_viewers(
viewer_origin: &mut moq_net::announce::Consumer,
cmd_tx: &tokio::sync::mpsc::Sender<Command>,
) -> anyhow::Result<()> {
loop {
let Some(moq_net::announce::Update { path, broadcast }) = viewer_origin.next().await else {
break;
};
let viewer_id = path.to_string();
if let Some(broadcast) = broadcast {
tracing::info!(%viewer_id, "viewer connected");
let cmd_tx = cmd_tx.clone();
let vid = viewer_id.clone();
tokio::spawn(async move {
if let Err(e) = handle_viewer_commands(&vid, broadcast, &cmd_tx).await {
tracing::warn!(viewer_id = %vid, error = %e, "viewer command error");
}
tracing::info!(viewer_id = %vid, "viewer disconnected");
let _ = cmd_tx.send(Command::ViewerLeft { viewer_id: vid }).await;
});
} else {
tracing::info!(%viewer_id, "viewer went offline");
let _ = cmd_tx
.send(Command::ViewerLeft {
viewer_id: viewer_id.clone(),
})
.await;
}
}
Ok(())
}
async fn handle_viewer_commands(
viewer_id: &str,
broadcast: moq_net::broadcast::Consumer,
cmd_tx: &tokio::sync::mpsc::Sender<Command>,
) -> anyhow::Result<()> {
let track = broadcast.track("command")?.subscribe(None).await?;
let mut commands =
moq_json::snapshot::Consumer::<RawCommand>::new(track, moq_json::snapshot::ConsumerConfig::default());
loop {
let command = match commands.next().await {
Ok(Some(command)) => command,
Ok(None) => break,
Err(moq_json::Error::Json(err)) => {
tracing::warn!(%viewer_id, %err, "invalid command");
continue;
}
Err(err) => return Err(err.into()),
};
match command {
RawCommand::Buttons { buttons, timestamps } => {
let timestamps: Vec<_> = timestamps
.into_iter()
.filter_map(|t| {
let ts = Duration::try_from_secs_f64(t.ts / 1000.0).ok()?;
Some(TimestampEntry { label: t.label, ts })
})
.collect();
let _ = cmd_tx
.send(Command::Buttons {
buttons,
viewer_id: viewer_id.to_string(),
timestamps,
})
.await;
}
RawCommand::Reset {} => {
tracing::info!(%viewer_id, "reset");
let _ = cmd_tx.send(Command::Reset).await;
}
}
}
Ok(())
}