diff --git a/src/main.rs b/src/main.rs index 7bae696..60d548a 100644 --- a/src/main.rs +++ b/src/main.rs @@ -9,6 +9,7 @@ mod playback; mod playlist; mod room; mod source; +mod state; mod transport; mod types; mod web; diff --git a/src/playback.rs b/src/playback.rs index de9199b..206fa08 100644 --- a/src/playback.rs +++ b/src/playback.rs @@ -1,11 +1,17 @@ //! Playback lifecycle management. //! -//! The [`Player`] struct owns the lifecycle of the currently-playing track: -//! starting a new track (resolving its stream URL and transcoding to fMP4), -//! aborting playback, and detecting when a track has ended. +//! The [`Player`] struct owns the abort handle for the currently-running +//! transcoding pipeline. It does **not** track what's playing — that's the +//! [`PlaybackState`](crate::state::PlaybackState) machine's job. //! -//! This separates playback concerns from room orchestration — the room actor -//! tells the player *what* to play, not *how* to play it. +//! # Lifecycle +//! +//! 1. **Start**: [`Player::start`] receives a fully-formed [`ActiveTrackInfo`], +//! resolves the URL via the source registry, and spawns ffmpeg transcoding. +//! 2. **End**: The spawned task sends [`RoomCommand::TrackEnded`] when ffmpeg +//! finishes or errors. The room actor calls [`Player::on_track_ended`] to +//! acknowledge. +//! 3. **Skip**: [`Player::abort`] signals the transcoding task to stop. use std::path::PathBuf; use std::sync::Arc; @@ -14,30 +20,22 @@ use std::time::Duration; use tokio::sync::{mpsc, watch}; use crate::media; -use crate::playlist::QueueItem; use crate::source::SourceRegistry; use crate::transport::TrackPublishers; use crate::types::{ActiveTrackInfo, RoomCommand}; -/// Manages the lifecycle of the currently-playing track. -/// -/// # Lifecycle -/// -/// 1. **Start**: [`Player::start`] resolves the URL via the source registry, -/// spawns ffmpeg transcoding, and stores an abort handle. -/// 2. **End**: The spawned task sends [`RoomCommand::TrackEnded`] when ffmpeg -/// finishes or errors. The room actor calls [`Player::on_track_ended`] to -/// acknowledge. -/// 3. **Skip**: [`Player::abort`] signals the transcoding task to stop. +/// Manages the abort handle for the currently-running transcoding pipeline. /// -/// The player is **not** responsible for queue management — that belongs to -/// the room actor. +/// The player is **not** responsible for queue or state management — that +/// belongs to the [`PlaybackState`](crate::state::PlaybackState) machine in the +/// room actor. pub(crate) struct Player { source_registry: Arc, cache_dir: PathBuf, - /// Channel to send [`RoomCommand::TrackEnded`] back to the room actor. pipeline_tx: mpsc::Sender, - active_track: Option, + /// The id of the currently-spawned pipeline, if any. + track_id: Option, + /// Signal handle to abort the current pipeline. abort_tx: Option>, } @@ -51,32 +49,22 @@ impl Player { source_registry, cache_dir, pipeline_tx, - active_track: None, + track_id: None, abort_tx: None, } } - /// Start playing a queue item. + /// Start transcoding a track. /// - /// Resolves its URL to a playable stream via the source registry, spawns - /// ffmpeg transcoding to fMP4, and returns the active track info for state - /// publishing. + /// `info` must be a fully-formed [`ActiveTrackInfo`] with the `started_at` + /// timestamp already set by the caller. /// - /// If a track is already playing, it is **not** aborted — call [`abort`] - /// first. - /// - /// [`abort`]: Player::abort - pub(crate) fn start( - &mut self, - item: &QueueItem, - publishers: TrackPublishers, - ) -> ActiveTrackInfo { + /// If a pipeline is already running, call [`abort`](Self::abort) first. + pub(crate) fn start(&mut self, info: &ActiveTrackInfo, publishers: TrackPublishers) { let (abort_tx, abort_rx) = watch::channel(false); - let info = ActiveTrackInfo::from_item(item); - let track_id = item.id; - + let track_id = info.id; let source_registry = self.source_registry.clone(); - let url = item.url.clone(); + let url = info.url.clone(); let cache_dir = self.cache_dir.clone(); let pipeline_tx = self.pipeline_tx.clone(); @@ -93,8 +81,8 @@ impl Player { // 2. Transcode to fMP4 via ffmpeg. let group_id = track_id as u64; - let result = media::transcode_to_fmp4(&stream_url, publishers, abort_rx, group_id) - .await; + let result = + media::transcode_to_fmp4(&stream_url, publishers, abort_rx, group_id).await; if let Err(e) = &result { tracing::warn!(%url, error = %e, "player: transcoding finished with error"); @@ -105,54 +93,33 @@ impl Player { send_track_ended(&pipeline_tx, track_id).await; }); - self.active_track = Some(info.clone()); + self.track_id = Some(track_id); self.abort_tx = Some(abort_tx); - - info } - /// Abort the currently playing track. - /// - /// Returns the DB id of the aborted track, or `None` if nothing was playing. - pub(crate) fn abort(&mut self) -> Option { - let id = self.active_track.take().map(|t| t.id); + /// Abort the current pipeline. Idempotent. + pub(crate) fn abort(&mut self) { + self.track_id = None; if let Some(abort) = self.abort_tx.take() { let _ = abort.send(true); } - id } /// Acknowledge that a track has ended. /// - /// Returns the track's DB id if `item_id` matches the currently active - /// track, or `None` if it doesn't match (e.g. it was already skipped). - pub(crate) fn on_track_ended(&mut self, item_id: i64) -> Option { - let is_current = self - .active_track - .as_ref() - .map(|t| t.id == item_id) - .unwrap_or(false); - if is_current { - self.active_track = None; + /// Returns `true` if `item_id` matched the currently-tracked pipeline, + /// `false` if it was a stale event from an older pipeline. + pub(crate) fn on_track_ended(&mut self, item_id: i64) -> bool { + if self.track_id == Some(item_id) { + self.track_id = None; self.abort_tx = None; - Some(item_id) + true } else { - None + false } } - - /// Is a track currently playing? - pub(crate) fn is_playing(&self) -> bool { - self.active_track.is_some() - } - - /// Get a snapshot of the currently playing track, if any. - pub(crate) fn active_track(&self) -> Option { - self.active_track.clone() - } } -/// Send [`RoomCommand::TrackEnded`] with best-effort logging on failure. async fn send_track_ended(tx: &mpsc::Sender, item_id: i64) { if let Err(e) = tx.send(RoomCommand::TrackEnded { item_id }).await { tracing::warn!(item_id, error = %e, "player: failed to send TrackEnded"); @@ -174,7 +141,6 @@ mod tests { use std::path::Path; use std::sync::Arc; - /// A mock source that resolves all URLs to a fake stream. struct MockPlayerSource; #[async_trait::async_trait] @@ -192,7 +158,6 @@ mod tests { url: &str, _cache_dir: &Path, ) -> Result { - // Fake resolve: return the URL with /stream appended. Ok(format!("{url}/stream")) } } @@ -216,7 +181,7 @@ mod tests { pool } - async fn make_queue_item(pool: &Pool) -> QueueItem { + async fn make_queue_item(pool: &Pool) -> i64 { let pl = Playlist::new(pool.clone(), "test-room"); pl.push(&TrackMeta { title: "Test".into(), @@ -227,73 +192,78 @@ mod tests { }) .await .unwrap() + .id + } + + fn test_info(id: i64) -> ActiveTrackInfo { + ActiveTrackInfo { + id, + title: "Test".into(), + url: "https://example.com/track".into(), + duration: "3:45".into(), + thumbnail: None, + started_at_wall: 1_700_000_000_000, + } } #[tokio::test] - async fn test_player_start_sets_active_track() { + async fn test_player_start_then_abort() { let pool = test_pool().await; - let item = make_queue_item(&pool).await; + let item_id = make_queue_item(&pool).await; let (tx, _rx) = mpsc::channel(256); let mut player = Player::new(test_registry(), PathBuf::from("/tmp"), tx); - assert!(!player.is_playing()); + player.start(&test_info(item_id), TrackPublishers::new()); + + player.abort(); - let info = player.start(&item, TrackPublishers::new()); - assert!(player.is_playing()); - assert_eq!(info.title, "Test"); - assert_eq!(info.id, item.id); + // on_track_ended with a stale id should return false. + assert!(!player.on_track_ended(99999)); } #[tokio::test] - async fn test_player_abort_clears_active_track() { + async fn test_player_abort_idempotent() { let pool = test_pool().await; - let item = make_queue_item(&pool).await; + let item_id = make_queue_item(&pool).await; let (tx, _rx) = mpsc::channel(256); let mut player = Player::new(test_registry(), PathBuf::from("/tmp"), tx); - player.start(&item, TrackPublishers::new()); - assert!(player.is_playing()); + player.start(&test_info(item_id), TrackPublishers::new()); - let aborted_id = player.abort(); - assert_eq!(aborted_id, Some(item.id)); - assert!(!player.is_playing()); - assert!(player.active_track().is_none()); + player.abort(); + player.abort(); // second abort is a no-op + // Did not panic = pass } #[tokio::test] - async fn test_player_abort_when_idle_returns_none() { + async fn test_player_abort_when_idle() { let (tx, _rx) = mpsc::channel(256); let mut player = Player::new(test_registry(), PathBuf::from("/tmp"), tx); - assert!(player.abort().is_none()); + player.abort(); // no-op } #[tokio::test] async fn test_player_on_track_ended_matches_current() { let pool = test_pool().await; - let item = make_queue_item(&pool).await; + let item_id = make_queue_item(&pool).await; let (tx, _rx) = mpsc::channel(256); let mut player = Player::new(test_registry(), PathBuf::from("/tmp"), tx); - player.start(&item, TrackPublishers::new()); + player.start(&test_info(item_id), TrackPublishers::new()); - let result = player.on_track_ended(item.id); - assert_eq!(result, Some(item.id)); - assert!(!player.is_playing()); + assert!(player.on_track_ended(item_id)); } #[tokio::test] async fn test_player_on_track_ended_ignores_old_id() { let pool = test_pool().await; - let item = make_queue_item(&pool).await; + let item_id = make_queue_item(&pool).await; let (tx, _rx) = mpsc::channel(256); let mut player = Player::new(test_registry(), PathBuf::from("/tmp"), tx); - player.start(&item, TrackPublishers::new()); - assert!(player.is_playing()); + player.start(&test_info(item_id), TrackPublishers::new()); - // An old/stale item_id should be ignored. - let result = player.on_track_ended(99999); - assert!(result.is_none()); - assert!(player.is_playing()); + // Stale id should be ignored. + assert!(!player.on_track_ended(99999)); } } diff --git a/src/playlist.rs b/src/playlist.rs index 6827845..b1662e0 100644 --- a/src/playlist.rs +++ b/src/playlist.rs @@ -84,6 +84,7 @@ impl Playlist { /// Returns the item with `started_playing_at` set, or `None` when the /// queue is empty. Unlike the old `pop_next`, this does NOT set /// `played = 1` — that happens later via `mark_finished`. + #[allow(dead_code)] pub(crate) async fn pop_next(&self) -> Result, sqlx::Error> { let mut tx = self.pool.begin().await?; @@ -142,6 +143,16 @@ impl Playlist { Ok(items) } + /// Mark an item as started playing (sets started_playing_at timestamp). + pub(crate) async fn mark_started(&self, id: i64, started_at: &str) -> Result<(), sqlx::Error> { + sqlx::query("UPDATE queue_items SET started_playing_at = ? WHERE id = ?") + .bind(started_at) + .bind(id) + .execute(&self.pool) + .await?; + Ok(()) + } + /// Mark a currently-playing item as finished (moves it from "playing" to "played"). pub(crate) async fn mark_finished(&self, id: i64) -> Result<(), sqlx::Error> { let now = now_iso(); diff --git a/src/room.rs b/src/room.rs index d1818c1..c1cdb53 100644 --- a/src/room.rs +++ b/src/room.rs @@ -1,11 +1,15 @@ //! Room state management. //! //! Each room is identified by a snowflake ID, persists metadata in -//! SQLite, and holds ephemeral playback state in memory. +//! SQLite, and holds ephemeral playback state in memory via a pure +//! [`PlaybackState`] machine. //! //! All mutations flow through a per-room mpsc event-loop actor, -//! eliminating races, torn lock windows, and the generation counter. -//! State snapshots are published automatically after every command. +//! following the impure-pure-impure sandwich: +//! +//! 1. **Gather** — fetch data (DB, wall clock) +//! 2. **Transition** — `PlaybackState::transition()` (pure) +//! 3. **Commit** — execute returned [`Effect`]s use std::collections::HashMap; use std::path::PathBuf; @@ -24,6 +28,7 @@ use crate::db; use crate::playback::Player; use crate::playlist::Playlist; use crate::source::SourceRegistry; +use crate::state::{self, Effect, Event, PlaybackState}; use crate::transport::TrackPublishers; use crate::types::{ ChatMessage, HistoryEntry, QueueSummary, RoomCommand, RoomHandle, RoomId, TrackMeta, @@ -117,6 +122,7 @@ struct RoomActor { pool: Pool, // Mutable state (owned by this actor, no locks needed) + state: PlaybackState, player: Player, client_count: u64, next_user_id: u64, @@ -174,6 +180,7 @@ impl Registry { room_id: id.clone(), publishers: publishers.clone(), pool: self.pool.clone(), + state: PlaybackState::default(), player, client_count: 0, next_user_id: 1, @@ -250,9 +257,6 @@ impl Registry { } /// Send a chat message: validate, send to actor, return best-effort ChatMessage. - /// - /// Validation happens synchronously; the actual DB insert and broadcast - /// occur in the actor. The returned `id` is 0 (not the DB-assigned id). #[cfg(test)] pub(crate) async fn send_chat( &self, @@ -313,9 +317,7 @@ impl Registry { .collect()) } - /// Register a client in the room, returning a monotonically increasing user id. - /// - /// Returns `None` if the room is not active in memory. + /// Register a client in the room. #[cfg(test)] pub(crate) async fn register_client(&self, room_id: &RoomId) -> Option { let handle = self.handle(room_id).await?; @@ -328,7 +330,7 @@ impl Registry { resp_rx.await.ok() } - /// Unregister a client, decrementing the client count for the room. + /// Unregister a client. #[cfg(test)] pub(crate) async fn unregister_client(&self, room_id: &RoomId) { if let Some(handle) = self.handle(room_id).await { @@ -344,14 +346,10 @@ impl Registry { } /// Sweep idle rooms that have been inactive longer than `timeout`. - /// - /// Sends Shutdown to the actor and removes the room from both - /// in-memory state and SQLite. Returns the list of swept room IDs. pub(crate) async fn sweep_idle(&self, timeout: Duration) -> Vec { let now = Utc::now(); let mut swept = Vec::new(); - // Identify idle rooms while holding only a read lock. let idle_rooms: Vec = { let inner = self.inner.read().await; inner @@ -388,9 +386,6 @@ impl Registry { } /// An owned handle to a room, returned by [`Registry::create`]. -/// -/// Dropping this handle does not destroy the room — rooms are torn -/// down by an idle timeout sweeper. pub(crate) struct RoomRef { pub id: RoomId, _registry: Registry, @@ -410,38 +405,72 @@ impl RoomActor { tracing::info!(room = %self.room_id, "room actor shut down"); } + /// Process a single command. + /// + /// Returns `true` if the actor should shut down. async fn process(&mut self, cmd: RoomCommand) -> bool { match cmd { RoomCommand::QueueTrack(meta) => { tracing::debug!(room = %self.room_id, title = %meta.title, "track queued"); + + // ── Gather ────────────────────────────────────────── let playlist = Playlist::new(self.pool.clone(), &self.room_id.0); - if let Err(e) = playlist.push(&meta).await { - tracing::warn!("failed to persist queue item: {e}"); - } - if !self.player.is_playing() { - self.start_next().await; - } - self.publish_state_snapshot().await; + let item = match playlist.push(&meta).await { + Ok(item) => item, + Err(e) => { + tracing::warn!("failed to persist queue item: {e}"); + return false; + } + }; + let now = Utc::now().timestamp_millis(); + + // ── Transition (pure) ─────────────────────────────── + let effects = self.state.transition( + &Event::TrackQueued(state::QueuedTrack { + id: item.id, + title: item.title, + url: item.url, + duration: item.duration, + thumbnail: item.thumbnail, + }), + now, + ); + + // ── Commit ────────────────────────────────────────── + self.execute_effects(effects).await; } + RoomCommand::Skip => { tracing::info!(room = %self.room_id, "skip requested"); - let finished_id = self.player.abort(); - if let Some(id) = finished_id { - let playlist = Playlist::new(self.pool.clone(), &self.room_id.0); - let _ = playlist.mark_finished(id).await; - } - self.start_next().await; - self.publish_state_snapshot().await; + + // ── Gather ────────────────────────────────────────── + let now = Utc::now().timestamp_millis(); + + // ── Transition (pure) ─────────────────────────────── + let effects = self.state.transition(&Event::Skip, now); + + // ── Commit ────────────────────────────────────────── + self.execute_effects(effects).await; } + RoomCommand::TrackEnded { item_id } => { tracing::debug!(room = %self.room_id, item_id, "track ended"); - if self.player.on_track_ended(item_id).is_some() { - let playlist = Playlist::new(self.pool.clone(), &self.room_id.0); - let _ = playlist.mark_finished(item_id).await; - self.start_next().await; - self.publish_state_snapshot().await; + + // Gate on player first — stale events are silently dropped. + if !self.player.on_track_ended(item_id) { + return false; } + + // ── Gather ────────────────────────────────────────── + let now = Utc::now().timestamp_millis(); + + // ── Transition (pure) ─────────────────────────────── + let effects = self.state.transition(&Event::TrackEnded { item_id }, now); + + // ── Commit ────────────────────────────────────────── + self.execute_effects(effects).await; } + RoomCommand::Chat { user_name, content, @@ -468,8 +497,8 @@ impl RoomActor { let payload = serde_json::to_vec(&json).unwrap_or_default(); self.publishers.publish_chat(0, Bytes::from(payload)); } - // No state snapshot needed for chat. } + RoomCommand::RegisterClient { resp } => { let id = self.next_user_id; self.next_user_id = self.next_user_id.wrapping_add(1); @@ -477,10 +506,12 @@ impl RoomActor { let _ = resp.send(id); self.publish_state_snapshot().await; } + RoomCommand::UnregisterClient => { self.client_count = self.client_count.saturating_sub(1); self.publish_state_snapshot().await; } + RoomCommand::Shutdown => { tracing::info!(room = %self.room_id, "actor shutdown requested"); self.player.abort(); @@ -490,63 +521,96 @@ impl RoomActor { false } - async fn start_next(&mut self) { - let playlist = Playlist::new(self.pool.clone(), &self.room_id.0); - let item = match playlist.pop_next().await { - Ok(Some(item)) => item, - Ok(None) => return, - Err(e) => { - tracing::warn!("failed to pop queue: {e}"); - return; + /// Execute a batch of effects from a state transition. + async fn execute_effects(&mut self, effects: Vec) { + for effect in effects { + match effect { + Effect::AbortPipeline => { + self.player.abort(); + } + + Effect::PersistStarted(id) => { + let now = Utc::now().format("%Y-%m-%dT%H:%M:%S%.3fZ").to_string(); + let playlist = Playlist::new(self.pool.clone(), &self.room_id.0); + let _ = playlist.mark_started(id, &now).await; + } + + Effect::StartPipeline(track) => { + let now = Utc::now().timestamp_millis(); + self.state.resolve_started_at(now); + let info = crate::types::ActiveTrackInfo { + id: track.id, + title: track.title, + url: track.url, + duration: track.duration, + thumbnail: track.thumbnail, + started_at_wall: now, + }; + self.player.start(&info, self.publishers.clone()); + } + + Effect::PersistFinished(id) => { + let playlist = Playlist::new(self.pool.clone(), &self.room_id.0); + let _ = playlist.mark_finished(id).await; + } + + Effect::PublishSnapshot => { + self.publish_state_snapshot().await; + } } - }; - tracing::info!(room = %self.room_id, title = %item.title, item_id = item.id, "starting track"); - self.player.start(&item, self.publishers.clone()); + } } + /// Publish the current state snapshot to all connected clients. + /// + /// Reads from in-memory [`PlaybackState`] — no SQLite queries. async fn publish_state_snapshot(&self) { - let current = self.player.active_track().map(|t| TrackState { + let current = self.state.active.as_ref().map(|t| TrackState { id: t.id, title: t.title.clone(), url: t.url.clone(), duration: t.duration.clone(), started_at: t.started_at_wall, }); - let client_count = self.client_count; - - let playlist = Playlist::new(self.pool.clone(), &self.room_id.0); - let upcoming = playlist.upcoming().await.unwrap_or_default(); - let history = playlist.history(50).await.unwrap_or_default(); - let queue: Vec = upcoming - .into_iter() - .map(|i| QueueSummary { - id: i.id, - title: i.title, - url: i.url, - duration: i.duration, - thumbnail: i.thumbnail, + let queue: Vec = self + .state + .queue + .iter() + .map(|q| QueueSummary { + id: q.id, + title: q.title.clone(), + url: q.url.clone(), + duration: q.duration.clone(), + thumbnail: q.thumbnail.clone(), }) .collect(); - let history_entries: Vec = history - .into_iter() - .map(|i| HistoryEntry { - id: i.id, - title: i.title, - url: i.url, - duration: i.duration, - thumbnail: i.thumbnail, - played_at: i.played_at.unwrap_or_default(), + let history: Vec = self + .state + .history + .iter() + .map(|h| HistoryEntry { + id: h.id, + title: h.title.clone(), + url: h.url.clone(), + duration: h.duration.clone(), + thumbnail: h.thumbnail.clone(), + played_at: h.played_at.clone(), }) .collect(); - tracing::trace!(room = %self.room_id, queue_len = %queue.len(), history_len = %history_entries.len(), clients = %client_count, "state snapshot published"); + let client_count = self.client_count; + + tracing::trace!( + room = %self.room_id, queue_len = %queue.len(), history_len = %history.len(), + clients = %client_count, "state snapshot published" + ); let snapshot = serde_json::json!({ "current_track": current, "queue": queue, - "history": history_entries, + "history": history, "clients": client_count, }); let payload = Bytes::from(serde_json::to_vec(&snapshot).unwrap_or_default()); @@ -569,7 +633,7 @@ mod tests { pool } - /// A mock source that handles all URLs (for tests that need playback). + /// A mock source that handles all URLs. struct PanaceaSource; #[async_trait::async_trait] @@ -710,12 +774,10 @@ mod tests { .await .unwrap(); - // Give the actor time to process both chat commands. tokio::time::sleep(Duration::from_millis(100)).await; let recent = reg.recent_chat(&room.id, 10).await.unwrap(); assert_eq!(recent.len(), 2); - // recent_chat returns most recent first. assert_eq!(recent[0].user_name, "bob"); assert_eq!(recent[0].content, "msg2"); assert_eq!(recent[1].user_name, "alice"); @@ -781,12 +843,11 @@ mod tests { let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Register and unregister — subsequent register should still work. reg.register_client(&room.id).await; reg.unregister_client(&room.id).await; let id = reg.register_client(&room.id).await; - assert_eq!(id, Some(2)); // next_user_id is monotonically increasing + assert_eq!(id, Some(2)); } #[tokio::test] @@ -794,9 +855,7 @@ mod tests { let pool = test_pool().await; let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Skip on an empty, idle room should not panic or deadlock. reg.skip(&room.id).await; - // If we got here without panic, the test passes. } #[tokio::test] @@ -815,7 +874,6 @@ mod tests { reg.push(&room.id, meta).await; tokio::time::sleep(Duration::from_millis(100)).await; - // Two rapid skips — must not deadlock or corrupt state. let r1 = reg.clone(); let rid1 = room.id.clone(); let h1 = tokio::spawn(async move { r1.skip(&rid1).await }); @@ -834,7 +892,6 @@ mod tests { .send_chat(&room.id, "alice", "hello", "message") .await .unwrap(); - // With the actor architecture, send_chat returns id: 0 (fire-and-forget). assert_eq!(msg.id, 0, "send_chat should return id 0 in actor mode"); } @@ -844,10 +901,8 @@ mod tests { let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Tiny sleep so the room has a non-zero idle duration. tokio::time::sleep(Duration::from_millis(1)).await; - // Immediately sweep with a 0-second timeout. let swept = reg.sweep_idle(Duration::from_secs(0)).await; assert!(swept.contains(&room.id), "room should be swept immediately"); assert!( @@ -862,9 +917,7 @@ mod tests { let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Register a client (triggers state snapshot, updates last_active). reg.register_client(&room.id).await; - // Wait a tiny bit, then sweep with a timeout slightly longer than the wait. tokio::time::sleep(Duration::from_millis(50)).await; let swept = reg.sweep_idle(Duration::from_millis(100)).await; assert!(!swept.contains(&room.id), "active room should not be swept"); @@ -885,10 +938,8 @@ mod tests { }; reg.push(&room.id, meta).await; - // Give actor time to process. tokio::time::sleep(Duration::from_millis(100)).await; - // The track should have been popped from the queue (marked as playing). let queue = reg.queue(&room.id).await; assert!( queue.is_empty(), @@ -902,7 +953,6 @@ mod tests { let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // First track starts playing immediately. reg.push( &room.id, TrackMeta { @@ -916,7 +966,6 @@ mod tests { .await; tokio::time::sleep(Duration::from_millis(100)).await; - // Second track should stay in the queue. reg.push( &room.id, TrackMeta { @@ -949,7 +998,6 @@ mod tests { let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Push a track, let it start playing, then skip. reg.push( &room.id, TrackMeta { @@ -966,7 +1014,6 @@ mod tests { reg.skip(&room.id).await; tokio::time::sleep(Duration::from_millis(100)).await; - // After skip with empty queue, nothing should be playing and no panic. let queue = reg.queue(&room.id).await; assert!( queue.is_empty(), @@ -980,7 +1027,6 @@ mod tests { let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Push two tracks and play the first. reg.push( &room.id, TrackMeta { @@ -994,7 +1040,6 @@ mod tests { .await; tokio::time::sleep(Duration::from_millis(100)).await; - // Push second track. reg.push( &room.id, TrackMeta { @@ -1008,20 +1053,16 @@ mod tests { .await; tokio::time::sleep(Duration::from_millis(50)).await; - // Two rapid skips. reg.skip(&room.id).await; reg.skip(&room.id).await; tokio::time::sleep(Duration::from_millis(100)).await; - // After double skip, B should be in history and next should be idle. let playlist = crate::playlist::Playlist::new(reg.pool.clone(), &room.id.0); let history = playlist.history(10).await.unwrap(); - // At least one track should be in history (A or B or both). assert!( !history.is_empty(), "at least one track should be in history after double skip" ); - // The later track in history should be the last skipped one. assert_eq!(history[0].title, "B", "last skipped track should be B"); } @@ -1031,14 +1072,12 @@ mod tests { let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Register many clients — overflow should not panic. for _ in 0..100 { let id = reg.register_client(&room.id).await; assert!(id.is_some(), "register_client should return Some(id)"); } tokio::time::sleep(Duration::from_millis(50)).await; - // Room should still be usable after many registrations. let exists = reg.exists(&room.id).await; assert!(exists, "room should still exist"); } @@ -1062,7 +1101,6 @@ mod tests { .await; tokio::time::sleep(Duration::from_millis(100)).await; - // Remove room while track is playing. reg.remove(&room.id).await; tokio::time::sleep(Duration::from_millis(100)).await; @@ -1078,11 +1116,8 @@ mod tests { let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Register a client (updates last_active via actor while we await). reg.register_client(&room.id).await; - // Immediately sweep with a short timeout — room should NOT be swept because - // register_client just updated last_active (elapsed time ≪ timeout). let swept = reg.sweep_idle(Duration::from_millis(10_000)).await; assert!( !swept.contains(&room.id), @@ -1096,17 +1131,14 @@ mod tests { let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Register 3 clients. let id1 = reg.register_client(&room.id).await; let id2 = reg.register_client(&room.id).await; let id3 = reg.register_client(&room.id).await; assert!(id1.is_some() && id2.is_some() && id3.is_some()); - // Unregister one. reg.unregister_client(&room.id).await; tokio::time::sleep(Duration::from_millis(50)).await; - // Register another — should get a new id. let id4 = reg.register_client(&room.id).await; assert!(id4.is_some(), "new client should get an id"); assert_ne!(id4, id1, "ids should be unique"); @@ -1149,7 +1181,6 @@ mod tests { let queue_a = reg.queue(&room_a.id).await; let queue_b = reg.queue(&room_b.id).await; - // Each room's first track was popped (played) so queue should be empty. assert!( queue_a.is_empty(), "room A should be empty after first track starts" @@ -1165,10 +1196,8 @@ mod tests { let pool = test_pool().await; let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Skip with no tracks ever queued. reg.skip(&room.id).await; tokio::time::sleep(Duration::from_millis(50)).await; - // No crash = pass. } #[tokio::test] @@ -1177,7 +1206,6 @@ mod tests { let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Push and skip. reg.push( &room.id, TrackMeta { @@ -1193,7 +1221,6 @@ mod tests { reg.skip(&room.id).await; tokio::time::sleep(Duration::from_millis(100)).await; - // Push another track — should play. reg.push( &room.id, TrackMeta { @@ -1207,11 +1234,9 @@ mod tests { .await; tokio::time::sleep(Duration::from_millis(100)).await; - // B should be popped from the queue (playing). let queue = reg.queue(&room.id).await; assert!(queue.is_empty(), "new track should be playing"); - // History should have A (skipped). let playlist = crate::playlist::Playlist::new(reg.pool.clone(), &room.id.0); let history = playlist.history(10).await.unwrap(); assert!(!history.is_empty(), "skipped track should be in history"); @@ -1223,11 +1248,9 @@ mod tests { let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Unregister with no clients registered — should not underflow. reg.unregister_client(&room.id).await; tokio::time::sleep(Duration::from_millis(50)).await; - // Room should still work. assert!(reg.exists(&room.id).await); } @@ -1237,7 +1260,6 @@ mod tests { let reg = Registry::new(pool, PathBuf::from("/tmp/moqbox"), test_registry()); let room = reg.create().await; - // Push 3 tracks, skip through all. for title in ["A", "B", "C"] { reg.push( &room.id, @@ -1253,26 +1275,21 @@ mod tests { } tokio::time::sleep(Duration::from_millis(100)).await; - // Queue should have B, C (A is playing). let mut queue = reg.queue(&room.id).await; assert_eq!(queue.len(), 2, "B and C should be queued"); - // Skip A → B plays. reg.skip(&room.id).await; tokio::time::sleep(Duration::from_millis(100)).await; queue = reg.queue(&room.id).await; assert_eq!(queue.len(), 1, "only C should remain after skipping A"); assert_eq!(queue[0].title, "C"); - // Skip B → C plays. reg.skip(&room.id).await; tokio::time::sleep(Duration::from_millis(100)).await; queue = reg.queue(&room.id).await; assert!(queue.is_empty(), "nothing should remain after skipping B"); - // Skip C → idle. reg.skip(&room.id).await; tokio::time::sleep(Duration::from_millis(100)).await; - // No crash = pass. } } diff --git a/src/state.rs b/src/state.rs new file mode 100644 index 0000000..fc7a845 --- /dev/null +++ b/src/state.rs @@ -0,0 +1,403 @@ +//! Pure playback state machine — deterministic, testable, IO-free. +//! +//! Owns the in-memory queue, active track, and history. All transitions +//! are pure functions returning [`Effect`]s that the caller executes. +//! +//! # Impure-pure-impure sandwich +//! +//! 1. **Gather** — caller fetches data (DB, wall clock) +//! 2. **Transition** — `PlaybackState::transition(event, now)` — pure +//! 3. **Commit** — caller executes returned effects (DB writes, pipeline start) +//! +//! # Testing +//! +//! ```ignore +//! let mut state = PlaybackState::default(); +//! let effects = state.transition(&Event::TrackQueued(my_track), 1_700_000_000_000); +//! assert!(effects.contains(&Effect::StartPipeline(...))); +//! ``` + +use crate::types::ActiveTrackInfo; + +/// A track waiting in the queue (has a DB id). +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct QueuedTrack { + pub id: i64, + pub title: String, + pub url: String, + pub duration: String, + pub thumbnail: Option, +} + +/// A finished track in history. +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct FinishedTrack { + pub id: i64, + pub title: String, + pub url: String, + pub duration: String, + pub thumbnail: Option, + pub played_at: String, +} + +/// Events that drive the playback state machine. +#[derive(Debug, Clone)] +pub(crate) enum Event { + /// A new track was added to the queue (with its DB-assigned id). + TrackQueued(QueuedTrack), + /// User requested to skip the current track. + Skip, + /// The transcoding pipeline finished (naturally or via abort). + TrackEnded { item_id: i64 }, +} + +/// Side effects to execute after a transition. +/// +/// The caller (RoomActor) must execute these in order. +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) enum Effect { + /// Signal the current transcoding pipeline to abort. + AbortPipeline, + /// Start transcoding a new track. RoomActor fills in the timestamp. + StartPipeline(QueuedTrack), + /// Mark a track as started in the DB (set started_playing_at). + PersistStarted(i64), + /// Mark a track as finished in the DB (set played=1, played_at). + PersistFinished(i64), + /// Broadcast the current state snapshot to all clients. + PublishSnapshot, +} + +/// Pure playback state — no IO handles, no DB connections. +#[derive(Debug, Clone)] +pub(crate) struct PlaybackState { + /// Currently playing track, if any. + pub active: Option, + /// Upcoming tracks (not yet started). + pub queue: Vec, + /// Recently finished tracks (newest first). + pub history: Vec, +} + +impl Default for PlaybackState { + fn default() -> Self { + Self { + active: None, + queue: Vec::new(), + history: Vec::new(), + } + } +} + +impl PlaybackState { + /// Pure transition function. + /// + /// - `event` — what happened + /// - `now` — current wall-clock epoch ms (gathered by caller for determinism) + /// + /// Returns effects for the caller to execute. Mutates `self` in place. + /// Does NO I/O. Deterministic given the same inputs. + pub fn transition(&mut self, event: &Event, now: i64) -> Vec { + match event { + Event::TrackQueued(track) => self.handle_queued(track), + Event::Skip => self.handle_skip(now), + Event::TrackEnded { item_id } => self.handle_track_ended(*item_id, now), + } + } + + /// A track was added to the queue. + fn handle_queued(&mut self, track: &QueuedTrack) -> Vec { + self.queue.push(track.clone()); + + // If nothing is playing, start immediately. + if self.active.is_none() { + self.advance(track.clone()) + } else { + vec![Effect::PublishSnapshot] + } + } + + /// User requested skip. + fn handle_skip(&mut self, _now: i64) -> Vec { + let mut effects = Vec::new(); + + // Abort any active track and move it to history. + if let Some(active) = self.active.take() { + effects.push(Effect::AbortPipeline); + effects.push(Effect::PersistFinished(active.id)); + self.history.insert( + 0, + FinishedTrack { + id: active.id, + title: active.title, + url: active.url, + duration: active.duration, + thumbnail: active.thumbnail, + played_at: self.format_played_at(), + }, + ); + } + + if effects.is_empty() { + return effects; // nothing was playing, nothing to do + } + + // Advance to next track. + if let Some(next) = self.queue.first().cloned() { + effects.extend(self.advance(next)); + } else if !effects.is_empty() { + effects.push(Effect::PublishSnapshot); + } + + effects + } + + /// A transcoding pipeline finished (naturally or after abort signal). + fn handle_track_ended(&mut self, item_id: i64, _now: i64) -> Vec { + // Only react if this matches the currently active track. + let is_current = self + .active + .as_ref() + .map(|t| t.id == item_id) + .unwrap_or(false); + + if !is_current { + return Vec::new(); // stale event, ignore + } + + let mut effects = Vec::new(); + + // Move to history. + if let Some(active) = self.active.take() { + effects.push(Effect::PersistFinished(active.id)); + self.history.insert( + 0, + FinishedTrack { + id: active.id, + title: active.title, + url: active.url, + duration: active.duration, + thumbnail: active.thumbnail, + played_at: self.format_played_at(), + }, + ); + } + + // Advance to next track. + if let Some(next) = self.queue.first().cloned() { + effects.extend(self.advance(next)); + } else { + effects.push(Effect::PublishSnapshot); + } + + effects + } + + /// Pop the first track from the queue and start playing it. + /// + /// Returns effects for starting the pipeline and publishing state. + fn advance(&mut self, next: QueuedTrack) -> Vec { + self.queue.retain(|t| t.id != next.id); + + // Store active track with placeholder timestamp — RoomActor calls + // `resolve_started_at` with the real wall-clock time. + let info = ActiveTrackInfo { + id: next.id, + title: next.title.clone(), + url: next.url.clone(), + duration: next.duration.clone(), + thumbnail: next.thumbnail.clone(), + started_at_wall: 0, + }; + self.active = Some(info); + + vec![ + Effect::PersistStarted(next.id), + Effect::StartPipeline(next), + Effect::PublishSnapshot, + ] + } + + /// Replace the placeholder `started_at_wall` on the active track with + /// the real timestamp. Called by RoomActor when it executes StartPipeline. + pub fn resolve_started_at(&mut self, real_started_at: i64) { + if let Some(ref mut active) = self.active { + active.started_at_wall = real_started_at; + } + } + + fn format_played_at(&self) -> String { + // Simple ISO-like timestamp for frontend display. + // This is a mock since we don't have Utc::now() in pure code. + // RoomActor will override via PersistFinished handler. + String::new() + } +} + +/// # Pure unit tests (no async, no DB, no IO) +#[cfg(test)] +mod tests { + use super::*; + + fn queued(title: &str, id: i64) -> QueuedTrack { + QueuedTrack { + id, + title: title.into(), + url: format!("https://example.com/{title}"), + duration: "3:45".into(), + thumbnail: None, + } + } + + #[test] + fn idle_room_starts_playing_on_first_track() { + let mut s = PlaybackState::default(); + let effects = s.transition(&Event::TrackQueued(queued("A", 1)), 0); + + assert!(s.active.is_some()); + assert_eq!(s.active.as_ref().unwrap().title, "A"); + assert!(s.queue.is_empty()); + assert!(effects.contains(&Effect::PublishSnapshot)); + assert!(effects + .iter() + .any(|e| matches!(e, Effect::StartPipeline(_)))); + } + + #[test] + fn second_track_stays_in_queue() { + let mut s = PlaybackState::default(); + s.transition(&Event::TrackQueued(queued("A", 1)), 0); + let effects = s.transition(&Event::TrackQueued(queued("B", 2)), 0); + + assert_eq!(s.queue.len(), 1); + assert_eq!(s.queue[0].title, "B"); + assert_eq!(s.active.as_ref().unwrap().title, "A"); + assert_eq!(effects, vec![Effect::PublishSnapshot]); + } + + #[test] + fn skip_advances_to_next_track() { + let mut s = PlaybackState::default(); + s.transition(&Event::TrackQueued(queued("A", 1)), 0); + s.transition(&Event::TrackQueued(queued("B", 2)), 0); + let effects = s.transition(&Event::Skip, 1_700_000_000_001); + + // Should have advanced to B. + assert_eq!(s.active.as_ref().unwrap().title, "B"); + assert!(s.queue.is_empty()); + + // A should be in history. + assert_eq!(s.history.len(), 1); + assert_eq!(s.history[0].title, "A"); + + // Effects: abort prev, persist finished, start B, publish snapshot. + assert!(effects.contains(&Effect::AbortPipeline)); + assert!(effects.contains(&Effect::PersistFinished(1))); + assert!(effects + .iter() + .any(|e| matches!(e, Effect::StartPipeline(t) if t.title == "B"))); + assert!(effects.contains(&Effect::PublishSnapshot)); + } + + #[test] + fn skip_on_idle_does_nothing() { + let mut s = PlaybackState::default(); + let effects = s.transition(&Event::Skip, 0); + + assert!(s.active.is_none()); + assert!(effects.is_empty()); + } + + #[test] + fn skip_last_track_goes_idle() { + let mut s = PlaybackState::default(); + s.transition(&Event::TrackQueued(queued("A", 1)), 0); + let effects = s.transition(&Event::Skip, 0); + + assert!(s.active.is_none()); + assert_eq!(s.history.len(), 1); + assert_eq!(s.history[0].title, "A"); + assert!(effects.contains(&Effect::AbortPipeline)); + assert!(effects.contains(&Effect::PersistFinished(1))); + assert!(effects.contains(&Effect::PublishSnapshot)); + // No StartPipeline since queue is empty. + assert!(!effects + .iter() + .any(|e| matches!(e, Effect::StartPipeline(_)))); + } + + #[test] + fn track_ended_advances_to_next() { + let mut s = PlaybackState::default(); + s.transition(&Event::TrackQueued(queued("A", 1)), 0); + s.transition(&Event::TrackQueued(queued("B", 2)), 0); + + let effects = s.transition(&Event::TrackEnded { item_id: 1 }, 0); + + assert_eq!(s.active.as_ref().unwrap().title, "B"); + assert_eq!(s.history.len(), 1); + assert_eq!(s.history[0].title, "A"); + assert!(effects.contains(&Effect::PersistFinished(1))); + assert!(effects.contains(&Effect::PublishSnapshot)); + } + + #[test] + fn track_ended_ignores_stale_id() { + let mut s = PlaybackState::default(); + s.transition(&Event::TrackQueued(queued("A", 1)), 0); + + // Stale TrackEnded for a different id. + let effects = s.transition(&Event::TrackEnded { item_id: 999 }, 0); + + assert!(s.active.is_some()); + assert_eq!(s.active.as_ref().unwrap().title, "A"); + assert!(effects.is_empty()); + } + + #[test] + fn track_ended_last_track_goes_idle() { + let mut s = PlaybackState::default(); + s.transition(&Event::TrackQueued(queued("A", 1)), 0); + + let effects = s.transition(&Event::TrackEnded { item_id: 1 }, 0); + + assert!(s.active.is_none()); + assert_eq!(s.history.len(), 1); + assert!(effects.contains(&Effect::PublishSnapshot)); + assert!(!effects + .iter() + .any(|e| matches!(e, Effect::StartPipeline(_)))); + } + + #[test] + fn resolve_started_at_updates_active_track() { + let mut s = PlaybackState::default(); + s.transition(&Event::TrackQueued(queued("A", 1)), 0); + + s.resolve_started_at(42); + assert_eq!(s.active.as_ref().unwrap().started_at_wall, 42); + } + + #[test] + fn skip_then_track_ended_stale_is_ignored() { + let mut s = PlaybackState::default(); + s.transition(&Event::TrackQueued(queued("A", 1)), 0); + s.transition(&Event::TrackQueued(queued("B", 2)), 0); + + // Skip A → B starts playing. + s.transition(&Event::Skip, 0); + assert_eq!(s.active.as_ref().unwrap().title, "B"); + + // Stale TrackEnded for A arrives — should be ignored. + let effects = s.transition(&Event::TrackEnded { item_id: 1 }, 0); + assert!(effects.is_empty()); + assert_eq!(s.active.as_ref().unwrap().title, "B"); + } + + #[test] + fn resolve_started_at_noop_when_idle() { + let mut s = PlaybackState::default(); + s.resolve_started_at(42); + assert!(s.active.is_none()); + } +} diff --git a/src/types.rs b/src/types.rs index ce01602..0ffeae3 100644 --- a/src/types.rs +++ b/src/types.rs @@ -160,28 +160,17 @@ pub(crate) enum RoomCommand { } /// Ephemeral info about the active track, stored in-memory only. -#[derive(Debug, Clone)] +#[derive(Debug, Clone, PartialEq, Eq)] pub(crate) struct ActiveTrackInfo { pub(crate) id: i64, pub(crate) title: String, pub(crate) url: String, pub(crate) duration: String, + pub(crate) thumbnail: Option, /// Epoch ms when the track started — for client-side elapsed computation. pub(crate) started_at_wall: i64, } -impl ActiveTrackInfo { - pub(crate) fn from_item(item: &crate::playlist::QueueItem) -> Self { - Self { - id: item.id, - title: item.title.clone(), - url: item.url.clone(), - duration: item.duration.clone(), - started_at_wall: chrono::Utc::now().timestamp_millis(), - } - } -} - /// Public handle to a room — allows sending commands and reading publishers. #[derive(Clone)] pub(crate) struct RoomHandle { diff --git a/templates/room.html b/templates/room.html index 6d162ae..6dcd80f 100644 --- a/templates/room.html +++ b/templates/room.html @@ -110,6 +110,10 @@ } /* Video */ + #video-container { + position: relative; + margin-bottom: 1.5rem; + } #video-player { width: 100%; max-height: 85vh; @@ -117,7 +121,40 @@ background: var(--crust); border-radius: var(--radius-lg); display: block; - margin-bottom: 1.5rem; + } + .volume-overlay { + position: absolute; + bottom: 0.75rem; + right: 0.75rem; + display: flex; + align-items: center; + gap: 0.5rem; + background: rgba(0, 0, 0, 0.6); + border-radius: var(--radius-md); + padding: 0.375rem 0.625rem; + opacity: 0; + transition: opacity 150ms ease-out; + pointer-events: none; + z-index: 10; + } + #video-container:hover .volume-overlay { + opacity: 1; + pointer-events: auto; + } + .volume-overlay .icon { + font-size: 1.125rem; + line-height: 1; + color: var(--text); + cursor: pointer; + user-select: none; + } + .volume-overlay input[type="range"] { + width: 80px; + height: 4px; + accent-color: var(--blue); + cursor: pointer; + background: transparent; + } box-shadow: 0 8px 32px rgba(0, 0, 0, 0.35), 0 2px 8px rgba(0, 0, 0, 0.2); @@ -515,7 +552,12 @@
-
+
+
+ 🔊 + +
+
@@ -707,10 +749,26 @@ document.querySelector(".main-col"); const video = document.createElement("video"); video.id = "video-player"; - video.controls = true; + video.controls = false; video.autoplay = true; video.muted = true; container.insertBefore(video, container.firstChild); + + // Wire up volume controls (element is in HTML, not created in JS) + const volSlider = document.getElementById("vol-slider"); + const volIcon = document.getElementById("vol-icon"); + volSlider.addEventListener("input", () => { + const v = parseFloat(volSlider.value); + video.volume = v; + volIcon.textContent = v === 0 ? "\u{1F507}" : v < 0.5 ? "\u{1F509}" : "\u{1F50A}"; + }); + volIcon.addEventListener("click", () => { + video.muted = !video.muted; + volIcon.textContent = video.muted ? "\u{1F507}" : "\u{1F50A}"; + volSlider.value = video.muted ? "0" : "1"; + }); + volSlider.addEventListener("click", (e) => e.stopPropagation()); + mediaSource = new MediaSource(); mediaSource.addEventListener("sourceopen", () => { debug("MediaSource sourceopen"); @@ -1004,13 +1062,14 @@ } }, 1000); - // Always-latest sync: seek video if difference > 3 seconds. + // Forward-only sync: if video is more than 3 seconds behind + // the live edge, seek forward to catch up. Never seek backward. const video = document.getElementById("video-player"); if (video && video.buffered.length > 0 && startedAt) { const target = (Date.now() - startedAt) / 1000; if (target > 0) { - const diff = Math.abs(video.currentTime - target); - if (diff > 3) { + const behind = target - video.currentTime; + if (behind > 3) { video.currentTime = target; } }