use anyhow::{Context, Result};
use bytes::Bytes;
use clap::Parser;
use std::collections::BTreeMap;
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::time::{Duration, Instant};
use url::Url;
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, default_value_t = 300)]
pub timeout: u64,
#[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_lite::TrackProducer,
audio_track: moq_lite::TrackProducer,
video_active: AtomicBool,
audio_active: AtomicBool,
paused: AtomicBool,
resume: (Mutex<()>, Condvar),
timeout: Duration,
location: Option<String>,
}
impl Session {
async fn run_track_monitor(&self, name: &str, track: &moq_lite::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, emu: &mut emulator::Emulator, game_stats: &mut stats::Stats) -> Result<()> {
tracing::info!("pausing encoding");
let (lock, cvar) = &self.resume;
let mut guard = lock.lock().unwrap();
let pause_start = Instant::now();
let mut reset_done = false;
while self.paused.load(Ordering::Acquire) {
if !reset_done && pause_start.elapsed() >= self.timeout {
tracing::info!("resetting emulator (paused too long)");
emu.reset()?;
*game_stats = stats::Stats::new();
reset_done = true;
}
if reset_done {
guard = cvar.wait(guard).unwrap();
} else {
let remaining = self.timeout.saturating_sub(pause_start.elapsed());
let (g, _) = cvar.wait_timeout(guard, remaining).unwrap();
guard = g;
}
}
tracing::info!("resuming encoding");
Ok(())
}
fn publish_status(
&self,
emu: &emulator::Emulator,
viewer_latency: &std::collections::HashMap<String, f64>,
game_stats: &stats::Stats,
publisher: &mut status::StatusPublisher,
) {
let held: Vec<_> = emu.pressed_buttons().iter().copied().collect();
let latency_map: BTreeMap<String, u32> = viewer_latency.iter().map(|(k, ms)| (k.clone(), *ms as u32)).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_lite::BroadcastProducer::default();
let publish_origin = moq_lite::Origin::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_lite::Origin::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 session_handle = client
.with_publish(publish_origin.consume())
.with_consume(consume_origin)
.connect(config.url.clone())
.await?;
let catalog = moq_mux::CatalogProducer::new(&mut broadcast)?;
let video_encoder = video::VideoEncoder::spawn(broadcast.clone(), catalog.clone());
ffmpeg_next::init().context("failed to init ffmpeg")?;
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()),
timeout: Duration::from_secs(config.timeout),
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 = session_handle.closed() => res.map_err(Into::into),
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();
session.video_encoder.bootstrap(&mut emu, start)?;
let frame_duration = Duration::from_micros(16_742);
let mut next_frame = Instant::now();
let mut viewer_latency: std::collections::HashMap<String, f64> = std::collections::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(&mut emu, &mut game_stats)?;
next_frame = Instant::now();
game_stats.reset_tick();
session.video_encoder.force_keyframe();
audio_encoder.reset_epoch();
}
let now = Instant::now();
if now < next_frame {
std::thread::sleep(next_frame - now);
}
next_frame += frame_duration;
let elapsed = start.elapsed();
let current_ts_ms = elapsed.as_secs_f64() * 1000.0;
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);
while let Ok(cmd) = cmd_rx.try_recv() {
match cmd {
input::Command::Buttons {
buttons,
viewer_id,
ts_ms,
} => {
emu.set_buttons(&viewer_id, buttons.into_iter().collect());
let latency = current_ts_ms - ts_ms;
if latency >= 0.0 {
viewer_latency.insert(viewer_id, latency);
}
}
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();
}
}
}
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();
tokio::select! {
res = run(&config) => res,
_ = tokio::signal::ctrl_c() => std::process::exit(0),
}
}