diff --git a/assets/styling/chat.css b/assets/styling/chat.css index cc5d352..007b548 100644 --- a/assets/styling/chat.css +++ b/assets/styling/chat.css @@ -146,6 +146,33 @@ color: var(--color-subtle); } +/* Error state */ +.chat-error { + text-align: center; + padding: 1.5rem 1rem; + color: var(--color-error); + font-size: 0.8rem; +} + +.chat-error p { + margin-bottom: 0.75rem; +} + +.chat-error button { + padding: 0.35rem 0.75rem; + background: transparent; + color: var(--color-error); + border: 1px solid var(--color-error); + cursor: pointer; + font-family: var(--font-ui); + font-size: 0.75rem; +} + +.chat-error button:hover { + background: var(--color-error); + color: var(--color-base); +} + /* Empty state */ .chat-empty { text-align: center; diff --git a/src/catalog.rs b/src/catalog.rs index cd594b6..ab268f9 100644 --- a/src/catalog.rs +++ b/src/catalog.rs @@ -244,7 +244,6 @@ use { jacquard_common::deps::fluent_uri::Uri, jacquard_common::types::ident::AtIdentifier, mini_moka_wasm::sync::Cache as MokaCache, - std::collections::HashMap, std::sync::Arc, std::sync::OnceLock, std::time::Duration, @@ -254,6 +253,7 @@ use { }; use std::collections::BTreeMap; +use std::collections::HashMap; #[cfg(feature = "server")] #[derive(Debug, thiserror::Error)] enum ThumbnailError { @@ -276,6 +276,8 @@ static THUMBNAIL_CACHE: OnceLock> = OnceLock::new #[cfg(feature = "server")] static IONOSPHERE_CACHE: OnceLock> = OnceLock::new(); #[cfg(feature = "server")] +static CHAT_COLOUR_CACHE: OnceLock>> = OnceLock::new(); +#[cfg(feature = "server")] static CLIENT: OnceLock> = OnceLock::new(); #[cfg(feature = "server")] static FFMPEG_INIT: std::sync::Once = std::sync::Once::new(); @@ -285,6 +287,8 @@ const CACHE_TTL: Duration = Duration::from_secs(3600); const THUMBNAIL_CACHE_MAX: u64 = 200; #[cfg(feature = "server")] const IONOSPHERE_CACHE_MAX: u64 = 500; +#[cfg(feature = "server")] +const CHAT_COLOUR_CACHE_MAX: u64 = 1000; #[cfg(feature = "server")] fn shared_client() -> Arc { @@ -779,6 +783,16 @@ fn ionosphere_cache() -> &'static MokaCache { }) } +#[cfg(feature = "server")] +fn chat_colour_cache() -> &'static MokaCache> { + CHAT_COLOUR_CACHE.get_or_init(|| { + MokaCache::builder() + .max_capacity(CHAT_COLOUR_CACHE_MAX) + .time_to_live(CACHE_TTL) + .build() + }) +} + #[cfg(feature = "server")] async fn resolve_ionosphere(client: &Arc, subject: &str) -> IonosphereData { let cache = ionosphere_cache(); @@ -1284,6 +1298,77 @@ pub async fn get_chat_messages( Ok(ChatMessageBatch { messages, has_more }) } +/// Fetch chat profile colours for a set of author DIDs. +/// +/// Each author stores a `place.stream.chat.profile` record in their own repo at +/// `at://author_did/place.stream.chat.profile/self`. Results are cached with +/// a 1-hour TTL. Missing or erroring profiles return `None` for that DID — +/// callers should fall back to the CSS default colour. +#[server] +pub async fn get_chat_colours( + author_dids: Vec, +) -> Result>, ServerFnError> { + use vodplace_api::place_stream::chat::profile::ProfileRecord; + + let client = shared_client(); + let cache = chat_colour_cache(); + + let fetch_futures = author_dids.into_iter().map(|did| { + let client = client.clone(); + let cache = cache; + async move { + let did_str = did.as_str().to_string(); + let cache_key = SmolStr::new(&did_str); + // Check cache first + if let Some(cached) = cache.get(&cache_key) { + return (did_str, cached); + } + + // Construct the profile URI: at://did/place.stream.chat.profile/self + let profile_uri_str = + format!("at://{}/place.stream.chat.profile/self", did.as_str()); + let profile_uri: AtUri = match profile_uri_str.parse() { + Ok(u) => u, + Err(e) => { + warn!(did = %did.as_str(), "failed to parse profile URI: {e}"); + cache.insert(cache_key, None); + return (did_str, None); + } + }; + + let colour = match client.get_record::(&profile_uri).await { + Ok(resp) => match resp.into_output() { + Ok(output) => output.value.color.map(|c| { + ( + c.red.clamp(0, 255) as u8, + c.green.clamp(0, 255) as u8, + c.blue.clamp(0, 255) as u8, + ) + }), + Err(e) => { + warn!(did = %did.as_str(), "failed to parse chat profile: {e}"); + None + } + }, + Err(_) => { + // Profile not found or PDS unreachable — silently return None. + None + } + }; + + cache.insert(cache_key, colour); + (did_str, colour) + } + }); + + let results: HashMap> = join_all(fetch_futures) + .await + .into_iter() + .collect(); + + Ok(results) +} + /// Fetch one page of messages for a given direction, store them in cache, and update the cursor. /// /// Tries slingshot first; falls back to direct Constellation fetch if slingshot fails. diff --git a/src/player/chat_sidebar.rs b/src/player/chat_sidebar.rs index b32ee98..bdc2b3d 100644 --- a/src/player/chat_sidebar.rs +++ b/src/player/chat_sidebar.rs @@ -1,10 +1,15 @@ +use std::collections::HashMap; + use dioxus::prelude::*; use jacquard::deps::smol_str::SmolStr; -use crate::catalog::{get_chat_messages, resolve_chat_context}; -use crate::chat::{AnchorSource, ChatMessageView}; +use crate::catalog::{get_chat_colours, get_chat_messages, resolve_chat_context}; +use crate::chat::{AnchorSource, ChatContext, ChatMessageView}; use crate::player::time_format::format_time; +/// Maximum number of manual retries offered before showing a persistent error. +const MAX_AUTO_RETRIES: u32 = 2; + /// Format a VOD-relative millisecond offset as MM:SS or H:MM:SS. /// /// Negative values (messages before stream start) are clamped to 0. @@ -24,21 +29,34 @@ pub fn ChatSidebar( video_uri: SmolStr, anchor_source: Signal>, /// Controls whether the sidebar is shown; passed from Watch component. - chat_visible: Signal, + mut chat_visible: Signal, on_seek: EventHandler, ) -> Element { let mut chat_messages: Signal> = use_signal(Vec::new); let mut loading: Signal = use_signal(|| true); let mut error: Signal> = use_signal(|| None); let mut user_scrolled: Signal = use_signal(|| false); + let mut author_colours: Signal>> = + use_signal(HashMap::new); + // Counts how many retries the user has triggered; capped at MAX_AUTO_RETRIES. + let mut retry_count: Signal = use_signal(|| 0); + + // Chat context saved after initial load — needed for refetching on seek/prefetch. + let mut chat_ctx: Signal> = use_signal(|| None); + + // Previous playback time — used for seek detection (>10s jump). + let mut prev_playback_time: Signal = use_signal(|| 0.0_f64); // Fetch chat data: resolve context then get messages. let video_uri_for_resource = video_uri.clone(); - let _chat_resource = use_resource(move || { + let mut chat_resource = use_resource(move || { let uri_str = video_uri_for_resource.clone(); async move { use jacquard::types::aturi::AtUri; + loading.set(true); + error.set(None); + let at_uri: AtUri = match uri_str.parse() { Ok(u) => u, Err(e) => { @@ -90,11 +108,30 @@ pub fn ChatSidebar( // Surface the anchor source to the parent for "timing approximate" indicator. anchor_source.set(Some(ctx.anchor.source)); + // Determine initial anchor position: prefer saved resume position from + // localStorage (same key as controls.rs), then fall back to stream start. + #[cfg(all(target_family = "wasm", target_os = "unknown"))] + let initial_anchor_ms: i64 = { + use jacquard::deps::smol_str::format_smolstr; + let storage_key = format_smolstr!("vodplace:playback:{uri_str}"); + match ::get::( + &storage_key, + ) { + Ok(secs) if secs > 0.0 && secs.is_finite() => { + // Saved value is VOD-relative seconds; convert to wall-clock ms. + ctx.anchor.timestamp_ms + (secs * 1000.0) as i64 + } + _ => ctx.anchor.timestamp_ms, + } + }; + #[cfg(not(all(target_family = "wasm", target_os = "unknown")))] + let initial_anchor_ms: i64 = ctx.anchor.timestamp_ms; + let batch = match get_chat_messages( - ctx.streamer_did, + ctx.streamer_did.clone(), ctx.window_start_ms, ctx.window_end_ms, - ctx.anchor.timestamp_ms, + initial_anchor_ms, ) .await { @@ -106,37 +143,200 @@ pub fn ChatSidebar( } }; + // Collect unique author DIDs for colour fetching. + let unique_dids: Vec = { + let mut seen = std::collections::HashSet::new(); + batch + .messages + .iter() + .filter(|m| seen.insert(m.author_did.as_str().to_string())) + .map(|m| m.author_did.clone()) + .collect() + }; + chat_messages.set(batch.messages); + chat_ctx.set(Some(ctx)); loading.set(false); + + // Fetch colours in the background; don't block message display. + if !unique_dids.is_empty() { + match get_chat_colours(unique_dids).await { + Ok(colours) => author_colours.set(colours), + Err(e) => { + // Colour fetch failure is non-fatal; messages are already shown. + tracing::warn!("failed to fetch chat colours: {e}"); + } + } + } } }); + // Sliding window manager: seek detection and prefetch on approach. + // + // Runs on every playback tick. Detects large time jumps (seeks > 10s) to + // positions outside the buffer, triggering a targeted refetch. Also prefetches + // when within 30s of the last buffered message. + #[cfg(all(target_family = "wasm", target_os = "unknown"))] + { + use_effect(move || { + let current_time = *current_playback_time.read(); + let prev_time = *prev_playback_time.read(); + + // Update prev time tracking. + prev_playback_time.set(current_time); + + // Skip if still loading or no context yet. + let ctx = match chat_ctx.read().clone() { + Some(c) => c, + None => return, + }; + if *loading.read() { + return; + } + + let current_ms = (current_time * 1000.0) as i64; + let messages = chat_messages.read(); + + // Seek detection: large jump (>10s) to a time outside the buffer. + let jump_secs = (current_time - prev_time).abs(); + if jump_secs > 10.0 && !messages.is_empty() { + let buffer_start = messages.first().map(|m| m.vod_relative_ms).unwrap_or(0); + let buffer_end = messages.last().map(|m| m.vod_relative_ms).unwrap_or(0); + let seek_outside_buffer = + current_ms < buffer_start - 30_000 || current_ms > buffer_end + 30_000; + + if seek_outside_buffer { + // Clear buffer and fetch a new window around the seek target. + loading.set(true); + drop(messages); + chat_messages.set(Vec::new()); + + let seek_wall_ms = ctx.anchor.timestamp_ms + current_ms; + spawn(async move { + match get_chat_messages( + ctx.streamer_did.clone(), + ctx.window_start_ms, + ctx.window_end_ms, + seek_wall_ms, + ) + .await + { + Ok(b) => { + chat_messages.set(b.messages); + } + Err(e) => { + error.set(Some(format!("failed to fetch chat messages: {e}"))); + } + } + loading.set(false); + }); + return; + } + } + + // Prefetch on approach: when within 30s of the last buffered message, + // extend the window and fetch more. Only if the server has more data. + let buffer_end_ms = messages.last().map(|m| m.vod_relative_ms).unwrap_or(0); + let secs_to_end = ((buffer_end_ms - current_ms) as f64) / 1000.0; + drop(messages); + + if secs_to_end < 30.0 { + let current_wall_ms = ctx.anchor.timestamp_ms + current_ms; + let extended_end_ms = ctx.window_end_ms + 5 * 60_000; + let mut updated_ctx = ctx.clone(); + updated_ctx.window_end_ms = extended_end_ms; + chat_ctx.set(Some(updated_ctx)); + + spawn(async move { + match get_chat_messages( + ctx.streamer_did.clone(), + ctx.window_start_ms, + extended_end_ms, + current_wall_ms, + ) + .await + { + Ok(b) => { + // Merge new messages into buffer, maintaining sort order. + let existing = chat_messages.read().clone(); + let existing_tids: std::collections::HashSet<_> = + existing.iter().map(|m| m.tid.clone()).collect(); + let mut merged = existing; + for new_msg in b.messages { + if !existing_tids.contains(&new_msg.tid) { + let pos = merged.partition_point(|m| { + m.vod_relative_ms < new_msg.vod_relative_ms + }); + merged.insert(pos, new_msg); + } + } + chat_messages.set(merged); + } + Err(_) => { + // Prefetch failure is non-fatal — existing buffer remains. + } + } + }); + } + }); + } + // Auto-scroll: when playback time changes and user hasn't manually scrolled, - // scroll the now-divider into view. + // scroll the now-divider into view. When user has scrolled, check if the + // now-divider re-enters the viewport and reattach auto-scroll. #[cfg(all(target_family = "wasm", target_os = "unknown"))] { let user_scrolled_for_effect = user_scrolled; use_effect(move || { let _t = current_playback_time.read(); - if *user_scrolled_for_effect.read() { - return; - } + let is_user_scrolled = *user_scrolled_for_effect.read(); + let doc = web_sys::window().and_then(|w| w.document()); - if let Some(doc) = doc { - if let Ok(Some(el)) = doc.query_selector(".chat-now-divider") { - let opts = web_sys::ScrollIntoViewOptions::new(); - opts.set_behavior(web_sys::ScrollBehavior::Smooth); - opts.set_block(web_sys::ScrollLogicalPosition::Nearest); - el.scroll_into_view_with_scroll_into_view_options(&opts); + let Some(doc) = doc else { + return; + }; + let Ok(Some(el)) = doc.query_selector(".chat-now-divider") else { + return; + }; + + if is_user_scrolled { + // Check if the now-divider is visible in its scroll container. + // The container is .chat-sidebar (grandparent: divider is inside + // .chat-message-list which is inside .chat-sidebar). + if let Some(container) = el + .parent_element() + .and_then(|p| p.parent_element()) + { + let container_rect = container.get_bounding_client_rect(); + let el_rect = el.get_bounding_client_rect(); + let visible = el_rect.top() >= container_rect.top() + && el_rect.bottom() <= container_rect.bottom(); + if visible { + user_scrolled.set(false); + } } + return; } + + // Auto-scroll now-divider into view with smooth/nearest behavior. + let opts = web_sys::ScrollIntoViewOptions::new(); + opts.set_behavior(web_sys::ScrollBehavior::Smooth); + opts.set_block(web_sys::ScrollLogicalPosition::Nearest); + el.scroll_into_view_with_scroll_into_view_options(&opts); }); } - // Hidden state: render nothing when sidebar is toggled off. - // Full collapsed/toggle UI is implemented in Task 6. + // Hidden state: when the sidebar is toggled off, render only the show-toggle + // tab so the user can reopen it. This stays inside the sidebar grid area. if !*chat_visible.read() { - return rsx! {}; + return rsx! { + div { + class: "chat-show-toggle", + title: "Show chat", + onclick: move |_| chat_visible.set(true), + "Chat" + } + }; } let current_ms = (*current_playback_time.read() * 1000.0) as i64; @@ -146,6 +346,11 @@ pub fn ChatSidebar( // is >= current_ms. The divider renders before that message. let divider_idx = messages.partition_point(|m| m.vod_relative_ms < current_ms); + // Build a set of tids currently in the buffer for reply orphan detection. + // Used to decide whether reply-indicator clicks can scroll to a parent. + let buffered_tids: std::collections::HashSet = + messages.iter().map(|m| m.tid.clone()).collect(); + rsx! { div { class: "chat-sidebar", @@ -164,6 +369,12 @@ pub fn ChatSidebar( } } } + button { + class: "chat-hide-toggle", + title: "Hide chat", + onclick: move |_| chat_visible.set(false), + "\u{00D7}" + } } if *loading.read() { @@ -171,6 +382,15 @@ pub fn ChatSidebar( } else if let Some(ref err_msg) = *error.read() { div { class: "chat-error", p { "{err_msg}" } + if *retry_count.read() < MAX_AUTO_RETRIES { + button { + onclick: move |_| { + retry_count += 1; + chat_resource.restart(); + }, + "Retry" + } + } } } else if messages.is_empty() { div { class: "chat-empty", "No chat for this video." } @@ -184,6 +404,14 @@ pub fn ChatSidebar( let timestamp_label = format_chat_timestamp(msg.vod_relative_ms); let msg_text = msg.text.clone(); let tid = msg.tid.clone(); + let author_did_key = msg.author_did.as_str().to_string(); + let colour = author_colours + .read() + .get(&author_did_key) + .and_then(|c| *c); + let author_style = colour + .map(|(r, g, b)| format!("color: rgb({r},{g},{b})")) + .unwrap_or_default(); rsx! { // "Now" divider before the first message that is in the future. @@ -205,7 +433,11 @@ pub fn ChatSidebar( span { class: "chat-timestamp", "{timestamp_label}" } span { class: "chat-message-content", - span { class: "chat-author", "{author_label}" } + span { + class: "chat-author", + style: "{author_style}", + "{author_label}" + } span { class: "chat-text", "{msg_text}" } } }