use anyhow::{Context, Result};
use bytes::Bytes;
use clap::Parser;
use std::collections::{BTreeMap, HashMap};
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::time::{Duration, Instant};
use url::Url;
#[cfg(feature = "jemalloc")]
#[global_allocator]
static ALLOC: moq_native::jemalloc::tikv_jemallocator::Jemalloc = moq_native::jemalloc::tikv_jemallocator::Jemalloc;
mod audio;
mod emulator;
mod input;
mod stats;
mod status;
mod video;
#[derive(Parser, Clone)]
pub struct Config {
#[arg(long)]
pub url: Url,
#[arg(long)]
pub rom: PathBuf,
#[arg(long)]
pub name: Option<String>,
#[arg(long, default_value = "boy")]
pub prefix: String,
#[arg(long)]
pub prefix_game: Option<String>,
#[arg(long)]
pub prefix_viewer: Option<String>,
#[arg(long)]
pub location: Option<String>,
#[command(flatten)]
pub client: moq_native::ClientConfig,
#[command(flatten)]
pub log: moq_native::Log,
}
struct Session {
video_encoder: video::VideoEncoder,
video_track: moq_net::TrackProducer,
audio_track: moq_net::TrackProducer,
video_active: AtomicBool,
audio_active: AtomicBool,
paused: AtomicBool,
resume: (Mutex<()>, Condvar),
location: Option<String>,
}
impl Session {
async fn run_track_monitor(&self, name: &str, track: &moq_net::TrackProducer, flag: &AtomicBool) {
loop {
if track.used().await.is_err() {
break;
}
tracing::info!("resuming {name}: viewer subscribed");
flag.store(true, Ordering::Release);
self.paused.store(false, Ordering::Release);
self.resume.1.notify_all();
if track.unused().await.is_err() {
break;
}
tracing::info!("pausing {name}: no viewers");
flag.store(false, Ordering::Release);
}
}
async fn run_pause_monitor(&self) {
loop {
let (v, a) = tokio::join!(self.video_track.unused(), self.audio_track.unused());
if v.is_err() || a.is_err() {
break;
}
tracing::info!("pausing emulation: no viewers");
self.paused.store(true, Ordering::Release);
tokio::select! {
Err(_) = self.video_track.used() => break,
Err(_) = self.audio_track.used() => break,
else => {},
}
tracing::info!("resuming emulation: viewer connected");
self.paused.store(false, Ordering::Release);
self.resume.1.notify_all();
}
self.paused.store(false, Ordering::Release);
self.resume.1.notify_all();
}
fn wait_for_resume(&self) {
tracing::info!("pausing encoding");
let (lock, cvar) = &self.resume;
let mut guard = lock.lock().unwrap();
while self.paused.load(Ordering::Acquire) {
guard = cvar.wait(guard).unwrap();
}
tracing::info!("resuming encoding");
}
fn publish_status(
&self,
emu: &emulator::Emulator,
viewer_latency: &HashMap<String, Vec<status::LatencyEntry>>,
game_stats: &stats::Stats,
publisher: &mut status::StatusPublisher,
) {
let held: Vec<_> = emu.pressed_buttons().iter().copied().collect();
let latency_map: BTreeMap<String, Vec<status::LatencyEntry>> =
viewer_latency.iter().map(|(k, v)| (k.clone(), v.clone())).collect();
let status = status::Status {
buttons: held,
latency: latency_map,
stats: game_stats.report(),
location: self.location.clone(),
};
publisher.publish(&status);
}
}
async fn run(config: &Config) -> Result<()> {
let rom_path = config.rom.clone();
let name = config.name.clone().unwrap_or_else(|| {
rom_path
.file_stem()
.and_then(|s| s.to_str())
.unwrap_or("unknown")
.to_string()
});
tracing::info!(rom = %rom_path.display(), %name, "starting Game Boy emulator");
let (cmd_tx, cmd_rx) = tokio::sync::mpsc::channel::<input::Command>(64);
let client = config.client.clone().init()?;
let mut broadcast = moq_net::Broadcast::new().produce();
let publish_origin = moq_net::Origin::random().produce();
let default_game_prefix = format!("{}/game", config.prefix);
let default_viewer_prefix = format!("{}/viewer", config.prefix);
let game_prefix = config.prefix_game.as_deref().unwrap_or(&default_game_prefix);
let viewer_prefix = config.prefix_viewer.as_deref().unwrap_or(&default_viewer_prefix);
let broadcast_path = format!("{game_prefix}/{name}");
publish_origin.publish_broadcast(&broadcast_path, broadcast.consume());
let viewer_path = format!("{viewer_prefix}/{name}");
let consume_origin = moq_net::Origin::random().produce();
let mut viewer_consumer = consume_origin
.with_root(&viewer_path)
.expect("viewer prefix should be valid")
.consume();
tracing::info!(url = %config.url, %name, broadcast = %broadcast_path, "connecting to relay");
let reconnect = client
.with_publish(publish_origin.consume())
.with_consume(consume_origin)
.reconnect(config.url.clone());
let catalog = moq_mux::catalog::Producer::new(&mut broadcast)?;
let video_encoder = video::VideoEncoder::spawn(broadcast.clone(), catalog.clone());
let audio_encoder = audio::AudioEncoder::new(broadcast.clone(), catalog.clone(), 44100)?;
let video_track = video_encoder.track.clone();
let audio_track = audio_encoder.track().clone();
let status_publisher = status::StatusPublisher::new(&mut broadcast)?;
let session = Arc::new(Session {
video_encoder,
video_track,
audio_track,
video_active: AtomicBool::new(false),
audio_active: AtomicBool::new(false),
paused: AtomicBool::new(true), resume: (Mutex::new(()), Condvar::new()),
location: config.location.clone(),
});
let s = session.clone();
tokio::spawn(async move { s.run_track_monitor("video", &s.video_track, &s.video_active).await });
let s = session.clone();
tokio::spawn(async move { s.run_track_monitor("audio", &s.audio_track, &s.audio_active).await });
let s = session.clone();
tokio::spawn(async move { s.run_pause_monitor().await });
let emulator_handle = tokio::task::spawn_blocking({
let session = session.clone();
move || run_emulator(session, &rom_path, audio_encoder, status_publisher, cmd_rx)
});
tokio::select! {
res = emulator_handle => res?.context("emulator error"),
res = reconnect.closed() => Ok(res?),
res = input::handle_viewers(&mut viewer_consumer, &cmd_tx) => res,
}
}
fn run_emulator(
session: Arc<Session>,
rom_path: &std::path::Path,
mut audio_encoder: audio::AudioEncoder,
mut status_publisher: status::StatusPublisher,
mut cmd_rx: tokio::sync::mpsc::Receiver<input::Command>,
) -> Result<()> {
let mut emu = emulator::Emulator::new(rom_path)?;
let start = Instant::now();
emu.tick();
let elapsed = start.elapsed();
let rgba = Bytes::from(emu.framebuffer());
let ts = hang::container::Timestamp::from_micros(elapsed.as_micros() as u64).context("timestamp overflow")?;
session.video_encoder.try_frame(rgba, ts);
let samples = emu.audio_samples();
if !samples.is_empty() {
audio_encoder.push_samples(&samples, elapsed)?;
}
let frame_duration = Duration::from_micros(16_742);
let mut next_frame = Instant::now();
let mut viewer_latency: HashMap<String, Vec<status::LatencyEntry>> = HashMap::new();
let mut game_stats = stats::Stats::new();
let mut was_audio_active = false;
loop {
if session.paused.load(Ordering::Acquire) {
session.wait_for_resume();
next_frame = Instant::now();
game_stats.reset_tick();
session.video_encoder.force_keyframe();
audio_encoder.reset_epoch();
}
{
let elapsed = start.elapsed();
let encode_ms = u32::try_from(session.video_encoder.encode_duration().as_millis()).unwrap_or(u32::MAX);
while let Ok(cmd) = cmd_rx.try_recv() {
match cmd {
input::Command::Buttons {
buttons,
viewer_id,
timestamps,
} => {
emu.set_buttons(&viewer_id, buttons.into_iter().collect());
let mut breakdown = Vec::new();
let entry = |label: &str, ms: u32| status::LatencyEntry {
label: label.to_string(),
ms,
};
breakdown.push(entry("encode", encode_ms));
let ms_saturating = |d: Duration| u32::try_from(d.as_millis()).unwrap_or(u32::MAX);
for t in ×tamps {
let latency = elapsed.saturating_sub(t.ts);
breakdown.push(entry(&t.label, ms_saturating(latency)));
}
if let Some(min_ts) = timestamps.iter().map(|t| t.ts).min() {
let latency = elapsed.saturating_sub(min_ts);
breakdown.push(entry("input", ms_saturating(latency)));
}
viewer_latency.insert(viewer_id, breakdown);
}
input::Command::ViewerLeft { viewer_id } => {
emu.viewer_left(&viewer_id);
viewer_latency.remove(&viewer_id);
}
input::Command::Reset => {
tracing::info!("resetting emulator (viewer request)");
emu.reset()?;
game_stats = stats::Stats::new();
}
}
}
}
let now = Instant::now();
if now < next_frame {
std::thread::sleep(next_frame - now);
}
next_frame += frame_duration;
let elapsed = start.elapsed();
let is_video = session.video_active.load(Ordering::Relaxed);
let is_audio = session.audio_active.load(Ordering::Relaxed);
game_stats.tick(is_video, is_audio);
emu.tick();
session.publish_status(&emu, &viewer_latency, &game_stats, &mut status_publisher);
if is_video {
let rgba = Bytes::from(emu.framebuffer());
let ts =
hang::container::Timestamp::from_micros(elapsed.as_micros() as u64).context("timestamp overflow")?;
session.video_encoder.try_frame(rgba, ts);
}
if is_audio {
if !was_audio_active {
audio_encoder.reset_epoch();
}
let samples = emu.audio_samples();
if !samples.is_empty() {
if let Err(e) = audio_encoder.push_samples(&samples, elapsed) {
tracing::warn!(error = %e, "audio encode error");
}
}
} else {
emu.audio_samples();
}
was_audio_active = is_audio;
}
}
#[tokio::main]
async fn main() -> Result<()> {
let config = Config::parse();
config.log.init()?;
#[cfg(feature = "jemalloc")]
let jemalloc = moq_native::jemalloc::run();
#[cfg(not(feature = "jemalloc"))]
let jemalloc = std::future::pending::<anyhow::Result<()>>();
let result = tokio::select! {
res = run(&config) => res,
Err(err) = jemalloc => Err(err).context("jemalloc profiler failed"),
res = tokio::signal::ctrl_c() => res.context("failed to listen for ctrl-c"),
};
match result {
Ok(()) => std::process::exit(0),
Err(err) => {
tracing::error!(err = %format!("{err:#}"), "exiting");
std::process::exit(1);
}
}
}