//! URL extraction and the per-URL media pipeline: bounded job channel + //! worker pool, link-cache fast path, fetch, task build and send dispatch. use super::{log_key, reply}; use crate::ctx::AppContext; use crate::link_cache::{CachedMedia, CachedMediaKind, CachedPost}; use crate::media_sender::MediaSender; use crate::send::{self, Delivery, MediaItemPayload, MediaRef, Task}; use crate::state::ChatData; use std::collections::{HashMap, HashSet}; use std::future::Future; use std::sync::LazyLock; use teloxide::RequestError; use teloxide::types::{ChatAction, ChatId, Message, MessageEntityKind, MessageId}; use x_media::media::Media; /// One fetch per post at a time, keyed by the normalized cache key. Two chats /// posting the same link at the same moment (or a batch forward and a queued /// retry) used to run two full fetches: two sets of source requests, and for an /// ugoira or a bsky video two ffmpeg encodes of the same post. The first caller /// runs it and the rest wait for its result. The entry is dropped the moment /// the fetch settles, so this dedupes what is *concurrent* and never answers /// from an old result: a repeat later fetches again, and a failure is not /// cached (the user may well retry it). static IN_FLIGHT_FETCHES: LazyLock>>> = LazyLock::new(Default::default); /// A fetched post (or the error that stopped it), shared as-is: the error side /// is not `Clone`, so callers read it through the `Arc` — the same shape the /// send paths use for `&Fetched`. type FetchOutcome = Result, x_media::site::FetchError>; /// The channel a sharer publishes its result on, and waiters subscribe to. type SharedFetch = tokio::sync::broadcast::Sender>; /// The sharing core, over the caller's own map so the sharing rules can /// be tested without a network fetch. /// /// A caller that finds a live entry subscribes to it and waits; the caller that /// created the entry runs `fetch` and publishes the result. Two things keep /// that from stranding a request: the entry is removed by a guard (so a /// cancelled fetch cannot leave waiters subscribed to a channel nothing will /// ever write to), and a waiter whose sharer vanished fetches for itself. async fn shared_fetch( map: &parking_lot::Mutex>>, key: &str, fetch: F, ) -> std::sync::Arc where T: Send + Sync + 'static, F: FnOnce() -> Fut, Fut: Future, { let (sender, leader) = { let mut map = map.lock(); match map.get(key) { Some(sender) => (sender.clone(), false), None => { let (sender, _) = tokio::sync::broadcast::channel(1); map.insert(key.to_string(), sender.clone()); (sender, true) } } }; if !leader { // `Err` means the entry is gone without a value: the sharer was // cancelled, or it finished just as this caller subscribed (the // message predates the subscription). Fetch for ourselves instead of // failing a link that is perfectly fetchable. let mut receiver = sender.subscribe(); // The sender clone taken from the map is dropped first: held, it would // keep the channel open past the sharer's exit (a broadcast channel // closes when *all* senders are gone), and `recv` would wait forever // instead of reporting that the sharer vanished. drop(sender); match receiver.recv().await { Ok(shared) => return shared, Err(_) => return std::sync::Arc::new(fetch().await), } } // Removes the entry on every exit path, cancellation included. let _guard = InFlightFetch { map, key }; let outcome = std::sync::Arc::new(fetch().await); // No receiver is the common case, not an error: a lone caller has nobody // to publish to. let _ = sender.send(std::sync::Arc::clone(&outcome)); outcome } /// Drops the in-flight entry it was created for, however the fetch ends. struct InFlightFetch<'a, T> { map: &'a parking_lot::Mutex>>, key: &'a str, } impl Drop for InFlightFetch<'_, T> { fn drop(&mut self) { self.map.lock().remove(self.key); } } /// Extracts URL and text-link entities (text + caption), deduped in order. /// /// The offset work (`parse_entities` turning entities into slices of the /// message text) is teloxide's; the two decisions that are ours are /// [`url_of`] and [`dedupe_urls`], which is why they are separate and tested. pub fn extract_urls(message: &Message) -> Vec { let entities = message .parse_entities() .into_iter() .flatten() .chain(message.parse_caption_entities().into_iter().flatten()); dedupe_urls( entities .filter_map(|entity| url_of(entity.kind(), entity.text())) .collect(), ) } /// The URL an entity carries: a bare `Url` entity is its own text, a /// `TextLink` is its target (its display text is often a different string). /// Every other entity kind (bold, code, hashtag, …) carries none. fn url_of(kind: &MessageEntityKind, text: &str) -> Option { match kind { MessageEntityKind::Url => Some(text.to_string()), MessageEntityKind::TextLink { url } => Some(url.to_string()), _ => None, } } /// Keeps the first occurrence of each link, in order. Dedup is by the /// normalized post id, so variant URLs of the same post (`/status/1` vs /// `/status/1/photo/1`, or a text link whose target equals a pasted URL) are /// sent once; URLs no site claims (and plain text that is no URL) fall back to /// exact-string dedup. fn dedupe_urls(urls: Vec) -> Vec { let mut seen = HashSet::new(); let mut out = Vec::with_capacity(urls.len()); for url in urls { let key = x_media::site::cache_key(&url).unwrap_or_else(|| url.clone()); if seen.insert(key) { out.push(url); } } out } /// For locally produced media (encoded ugoira MP4) the thumbnail URL is a /// hotlink-protected remote URL Telegram may not fetch; let Telegram generate /// its own thumbnail instead. fn thumbnail_for(media: &Media) -> Option { let url = media.url(); if url.starts_with("http://") || url.starts_with("https://") { // An empty thumbnail string (misskey video/gif files without a // thumbnailUrl) must not reach Telegram; let it generate its own. media .thumbnail_url() .map(str::to_string) .filter(|t| !t.is_empty()) } else { None } } /// The link-cache snapshot of a freshly fetched post: its caption and raw /// render fields, with no media yet — the send fills in the Telegram file ids /// and persists the entry. `None` for a post that carries no render data /// (nothing to rebuild a caption format from later). pub(super) fn cached_snapshot(fetched: &x_media::site::Fetched) -> Option { fetched .render_fields() .map(|(author, author_url, title, content, tags)| CachedPost { url: fetched.source_url.clone(), caption: fetched.caption.clone(), title: title.to_string(), content: content.to_string(), author: author.to_string(), author_url: author_url.to_string(), tags: tags.to_string(), sensitive: fetched.sensitive, media: vec![], }) } pub(super) fn media_to_payload(media: &Media, sensitive: bool) -> Option { let url = media.url(); // Only adapter-produced temp files may be non-HTTP. A bare path from an // upstream JSON field must never reach InputFile::file(). let is_local_temp = std::path::Path::new(url) .file_name() .and_then(|name| name.to_str()) .is_some_and(|name| name.starts_with(x_media::TEMP_FILE_PREFIX)); if !url.starts_with("http://") && !url.starts_with("https://") && !is_local_temp { log::warn!("dropping media with a non-URL upstream path"); return None; } let media_ref = MediaRef::Source(url.to_string()); let fallback_url = media .smaller_url() .filter(|url| url.starts_with("http://") || url.starts_with("https://")) .map(str::to_string); Some(match media { // A gif inside a group becomes a video item; a lone gif takes the // animation path (see url_media). Media::Illustration { .. } => MediaItemPayload::Photo { media: media_ref, has_spoiler: sensitive, fallback_url, }, Media::Video { .. } => MediaItemPayload::Video { media: media_ref, has_spoiler: sensitive, thumbnail: thumbnail_for(media), fallback_url, }, Media::Animated { .. } => MediaItemPayload::Video { media: media_ref, has_spoiler: sensitive, thumbnail: thumbnail_for(media), fallback_url, }, }) } /// Sends a task and handles the outcome: post-send actions on success, retry /// enqueue on retryable failure, reply + link-cache invalidation on /// permanent failure (a stale cached file id must not repeat forever). async fn dispatch_send( ctx: &AppContext<'_>, chat_id: i64, reply_to: MessageId, task: &Task, url: &str, started: std::time::Instant, ) { let result = match task { Task::SendAnimation { .. } => send::send_animation(ctx, task).await, Task::SendMediaSequence { .. } => send::send_media_sequence(ctx, task).await, Task::ForwardMessages { .. } => unreachable!(), }; // Fetch + cache lookup + upload: the whole wait the user sat through. let ms = started.elapsed().as_millis(); match result { Ok(message_ids) => { log::info!( "sent {} message(s) for [key={}] chat={chat_id} in {ms}ms", message_ids.len(), log_key(url) ); send::post_send_actions(ctx, task, message_ids).await; send::settle_task(ctx, task, send::Settled::Sent).await; } Err(send::SendError::Retryable { delay_seconds, task, }) => { log::info!( "send for [key={}] chat={chat_id} failed after {ms}ms, queued for retry in {delay_seconds:.1}s", log_key(url) ); // Name the post and the wait: "queued for retry" alone left the // user guessing which link it was and how long the wait is. The // promise is made only when the retry was really persisted — an // enqueue that failed (DB write) would leave the user waiting for // a retry nothing can deliver. let promised = if send::enqueue_retry(ctx.task_queue, &task, delay_seconds).await { format!( "Send failed for {} — retrying in {delay_seconds:.0}s.", log_key(url) ) } else { format!( "Send failed for {} and the retry could not be queued — please send the link again.", log_key(url) ) }; let _ = reply(ctx.sender, chat_id, reply_to, promised).await; } Err(send::SendError::Permanent { message: err_message, task, }) => { send::settle_task(ctx, &task, send::Settled::Failed).await; log::error!( "send for [key={}] chat={chat_id} failed permanently after {ms}ms: {err_message}", log_key(url) ); let _ = reply( ctx.sender, chat_id, reply_to, format!("Send failed: {err_message}"), ) .await; } } } /// A cached entry's item as a send payload: the Telegram file id when the entry /// still has one, otherwise the source URL. /// /// The URL case is a *degraded* entry — `send::post_send`'s `invalidate_cache` /// drops the file ids of an entry whose cached send failed permanently, keeping /// the URLs, because a stale file id says nothing about the media. Sending from /// the URL costs Telegram a fetch (or the upload fallback a download) and saves /// the whole source round trip, including a ugoira encode or an HLS remux. fn cached_media_payload(media: &CachedMedia, sensitive: bool) -> MediaItemPayload { // The entry's file id when it has one, else its source URL (a permanently // failed send degrades the entry and clears the id). let source = if media.file_id.is_empty() { MediaRef::Source(media.url.clone()) } else { MediaRef::FileId(media.file_id.clone()) }; // A degraded item carries no smaller variant: a fresh fetch's would, but // the item is what the source itself sent, so an oversize is handled by the // upload fallback rather than by a URL that was never recorded. let fallback_url = None; match media.kind { CachedMediaKind::Photo => MediaItemPayload::Photo { media: source, has_spoiler: sensitive, fallback_url, }, CachedMediaKind::Video => MediaItemPayload::Video { media: source, has_spoiler: sensitive, thumbnail: None, fallback_url, }, CachedMediaKind::Animation => MediaItemPayload::Animation { media: source, has_spoiler: sensitive, }, } } /// Whether a send also runs the chat's post-send actions. `/test` sends with /// them suppressed so a test can never forward to the channel or open the /// edit-before-forward prompt; a normal link uses whatever the chat is /// configured with. #[derive(Clone, Copy, PartialEq, Eq, Debug)] pub(crate) enum PostSend { /// Apply the chat's `forward_channel_id` / `edit_before_forward`. FromChat, /// Send only: no channel forward, no edit prompt. Suppressed, } /// Builds the send task from ready-made items, sharing the payload shape /// between the fresh-fetch, link-cache and `/test` paths. #[allow(clippy::too_many_arguments)] fn build_send_task( chat_data: &ChatData, chat_id: i64, reply_to_message_id: i64, source_url: String, caption: String, items: Vec, cache_data: Option, post_send: PostSend, ) -> Task { // Notification ids stay set in both modes: a queued retry that // dead-letters should still tell the chat. let (edit_before_forward, forward_channel_id) = match post_send { PostSend::FromChat => (chat_data.edit_before_forward, chat_data.forward_channel_id), PostSend::Suppressed => (false, None), }; Task::from_items( Delivery { chat_id, reply_to_message_id, edit_before_forward, forward_channel_id, notify_chat_id: Some(chat_id), notify_message_id: Some(reply_to_message_id), }, source_url, caption, items, cache_data, ) } /// The per-URL pipeline: link cache → fetch → build → send → post-send. /// /// `post_send` selects whether the chat's forward/edit settings apply: the URL /// workers pass [`PostSend::FromChat`], the `/test` command /// [`PostSend::Suppressed`]. Everything else (cache write, retry enqueue, /// dead-letter notification) is identical. /// /// Wraps [`url_media_inner`] with the chat-action keep-alive: Telegram expires /// an action indicator after ~5s, while a fetch (ugoira encode, HLS remux) plus /// a download-and-reupload fallback routinely takes longer — without the /// refresh the chat shows nothing and the bot reads as stalled. pub(crate) async fn url_media( ctx: &AppContext<'_>, chat_id: i64, reply_to_message_id: i64, url: &str, post_send: PostSend, ) { // Shared with the pipeline: once the media types are known the indicator // switches from "typing" to "sending photo/video". let hint = parking_lot::Mutex::new(ActionHint::Typing); run_with_chat_action( ctx.sender, chat_id, &hint, url_media_inner(ctx, chat_id, reply_to_message_id, url, post_send, &hint), ) .await; } /// Runs `pipeline` while keeping the chat's action indicator alive: Telegram /// expires an action after ~5s, while a fetch (ugoira encode, HLS remux) plus a /// download-and-reupload fallback routinely takes longer. The pipeline updates /// `hint` when it knows what it is sending. /// /// No action is ever awaited *ahead* of the pipeline: doing that held the loop /// — and with it the fetch the user is waiting for — for a Telegram round trip, /// once before the pipeline was polled at all and again every /// [`ACTION_REFRESH`]. The in-flight send is held and polled *beside* the /// pipeline instead: the opening indicator still goes out before the pipeline's /// own first call (that is what it is for), but a slow API can no longer delay /// anything but the next indicator. async fn run_with_chat_action>( sender: &dyn MediaSender, chat_id: i64, hint: &parking_lot::Mutex, pipeline: F, ) { let warn = |e: RequestError| { // Cosmetic indicator: a failure degrades the experience, it does not // break the send (a group where the bot cannot send actions). log::warn!("send_chat_action failed for chat {chat_id}: {e}"); }; // The guard is released before the await: a parking_lot guard held across // it makes the future !Send, and the URL workers spawn these. let mut action = Some(sender.send_chat_action(ChatId(chat_id), hint.lock().action())); tokio::pin!(pipeline); loop { tokio::select! { // `biased` fixes the order below: the indicator is polled ahead of // the pipeline, and a finished pipeline returns without arming the // refresh timer (no stray actions). biased; // `select!` evaluates every branch's future expression eagerly, so // the `None` case is an inert block: the guard is what keeps it // from being polled (and from unwrapping a `None`). result = async { action.as_mut().unwrap().await }, if action.is_some() => { action = None; if let Err(e) = result { warn(e); } } () = &mut pipeline => return, () = tokio::time::sleep(ACTION_REFRESH) => { // One action in flight at a time: re-arming while the previous // send is still unanswered would drop it mid-request. if action.is_none() { action = Some(sender.send_chat_action(ChatId(chat_id), hint.lock().action())); } } } } } /// How often the chat-action indicator is refreshed while a pipeline runs. /// Telegram's indicator lasts ~5s; refreshing slightly inside that keeps it /// on-screen continuously. const ACTION_REFRESH: std::time::Duration = std::time::Duration::from_secs(4); /// What the chat action should say. Unknown before the fetch, so the pipeline /// starts with `Typing` and switches as soon as the media types are known. #[derive(Clone, Copy)] enum ActionHint { Typing, Photo, Video, } impl ActionHint { /// Photos make Telegram label the send "sending photo"; video/animation /// only payloads get "sending video". A mixed post takes the photo label /// (the group's first item is always a photo, see `photos_first`). fn for_items(items: &[MediaItemPayload]) -> Self { if items .iter() .any(|item| matches!(item, MediaItemPayload::Photo { .. })) { Self::Photo } else { Self::Video } } fn action(self) -> ChatAction { match self { Self::Typing => ChatAction::Typing, Self::Photo => ChatAction::UploadPhoto, Self::Video => ChatAction::UploadVideo, } } } /// User-facing text for a failed fetch. The [`FetchError`] class is what tells /// the user whether the post is gone, withheld or the source is refusing /// requests; a single generic sentence threw that away. fn fetch_error_message(err: &x_media::site::FetchError) -> String { use x_media::site::FetchError; match err { FetchError::NotFound => "Post not found (deleted, private or unavailable).".to_string(), FetchError::Sensitive => concat!( "This post's media is withheld (age-restricted). ", "The bot owner must set TWITTER_AUTH_TOKEN to fetch it." ) .to_string(), FetchError::Blocked => { "The source site refused the request (risk control). Try again later.".to_string() } FetchError::Disabled { site } => { format!("{} support is disabled on this bot.", site_title(site)) } FetchError::RateLimited { .. } | FetchError::Transient(_) | FetchError::Http(_) => { "The source site is unavailable right now (tried 3 times). Try again later.".to_string() } FetchError::MediaPrep(_) => concat!( "Could not prepare this post's media (its download or encode failed). ", "Try again later." ) .to_string(), // Parse/shape surprises, pixiv auth details, oversized media: nothing // actionable for the user beyond "this did not work". _ => "Failed to fetch media from this link.".to_string(), } } /// Site ids are lowercase ASCII (`pixiv`); user-facing text capitalizes the /// first letter. fn site_title(site: &str) -> String { let mut chars = site.chars(); match chars.next() { Some(first) => first.to_uppercase().collect::() + chars.as_str(), None => String::new(), } } #[allow(clippy::too_many_arguments)] async fn url_media_inner( ctx: &AppContext<'_>, chat_id: i64, reply_to_message_id: i64, url: &str, post_send: PostSend, hint: &parking_lot::Mutex, ) { let reply_to = MessageId(reply_to_message_id as i32); // Whole-link timer for the result lines: fetch (ugoira encode, HLS remux // included) + cache lookup + upload — the wait the user actually had. let started = std::time::Instant::now(); // Link cache: a post sent before is re-sent from Telegram file ids — // no source-site request, no download, no upload. Keyed by the // normalized post id so x.com / fxtwitter / /photo/N variants collide. if let Some(key) = x_media::site::cache_key(url) && let Some(cached) = ctx.link_cache.get(&key, ctx.config.link_cache_ttl).await { log::debug!("link cache hit for {key}"); let chat_data = ctx.chat_store.get(chat_id).await; // Cache keys are prefixed with the site id ("twitter:…"), matching // the value a fresh fetch would read from Fetched::site_id. The key // came out of cache_key, so its prefix is a registered id by // construction — splitting it off is the whole lookup. let site = key.split(':').next().unwrap_or(""); let format = chat_data.format_for(site); // One call for both: `caption_from_fields` returns the truncated // built-in caption itself when the chat has no format for this site. let caption = x_media::site::caption_from_fields( &format, &cached.caption, &cached.url, &cached.author, &cached.author_url, &cached.title, &cached.content, &cached.tags, ); let items: Vec = cached .media .iter() .map(|m| cached_media_payload(m, cached.sensitive)) .collect(); // The indicator switches to "sending photo/video" once the kinds are // known; `items` is moved into the task below. *hint.lock() = ActionHint::for_items(&items); let task = build_send_task( &chat_data, chat_id, reply_to_message_id, cached.url.clone(), caption, items, Some(cached), post_send, ); dispatch_send(ctx, chat_id, reply_to, &task, url, started).await; return; } log::debug!("fetching [key={}]", log_key(url)); log::trace!("fetching {url}"); // One fetch per post at a time: a concurrent duplicate of this link waits // for *this* fetch instead of running its own. let outcome = match x_media::site::cache_key(url) { Some(key) => shared_fetch(&IN_FLIGHT_FETCHES, &key, || x_media::site::fetch(url)).await, // A URL no site claims (reached only through `/test`): nothing to key // the sharing on, and the dispatcher answers without a request. None => std::sync::Arc::new(x_media::site::fetch(url).await), }; match &*outcome { // Unsupported links are ignored silently (Python parity). Ok(None) => { // The URL itself is user data, so only `trace` names the link; // `debug` just records that the message was looked at. log::debug!("no site pattern matches the link; ignoring"); log::trace!("no site pattern matches {url}"); } // Retries exhausted: notify the user (Rust-only requirement 3). Err(e) => { log::error!("fetch [key={}]: {e}", log_key(url)); let _ = reply(ctx.sender, chat_id, reply_to, fetch_error_message(e)).await; } Ok(Some(fetched)) => { let chat_data = ctx.chat_store.get(chat_id).await; // Per-site caption format override (empty -> built-in caption). let format = chat_data.format_for(fetched.site_id); if fetched.media.is_empty() { // A post with no media is still a post: its text goes out as a // message through the same caption the media path would attach. let text = fetched .render_fields() .map(|(_, _, title, content, _)| x_media::site::compose_text(title, content)) .unwrap_or_default(); let caption = fetched.caption_with(&format); let caption = send::quote_long_caption(&caption, &text, ctx.config.caption_quote_text_chars); send::send_text_post(ctx, chat_id, reply_to.0 as i64, caption.into_owned()).await; return; } let caption = fetched.caption_with(&format); let cache_data = cached_snapshot(fetched); let items: Vec = fetched .media .iter() .filter_map(|media| media_to_payload(media, fetched.sensitive)) .collect(); if items.is_empty() { let text = fetched .render_fields() .map(|(_, _, title, content, _)| x_media::site::compose_text(title, content)) .unwrap_or_default(); let caption = fetched.caption_with(&format); let caption = send::quote_long_caption(&caption, &text, ctx.config.caption_quote_text_chars); send::send_text_post(ctx, chat_id, reply_to.0 as i64, caption.into_owned()).await; return; } *hint.lock() = ActionHint::for_items(&items); let task = build_send_task( &chat_data, chat_id, reply_to_message_id, fetched.source_url.clone(), caption, items, cache_data, post_send, ); // Hand the keep-alive temp dir (ugoira / bsky remux MP4) to the // retry registry: a queued retry runs after this function returns // and the fetch's own TempDir is dropped, so without this the // local file would be gone by the time the retry sends it. Shared // (Arc), so a second send of the same post holds its own reference. if let Some(dir) = fetched.keep_alive() { send::KEEP_ALIVE.lock().push(dir); } dispatch_send(ctx, chat_id, reply_to, &task, url, started).await; } } } #[cfg(test)] mod tests { use super::*; use crate::ctx::test_support::{TestStores, cached_photo, permanent_error}; use crate::media_sender::test_support::{MockSender, Outcome}; use std::time::Duration; #[tokio::test] async fn cache_hit_sends_file_ids_and_degrades_on_permanent_failure() { let stores = TestStores::new(); let sender = MockSender::scripted( vec![Outcome::GroupErr, Outcome::MessageErr], permanent_error, ); let ctx = stores.ctx(&sender); stores.link_cache().put("twitter:1", &cached_photo()).await; url_media(&ctx, 1, 2, "https://x.com/u/status/1", PostSend::FromChat).await; // The cached file id went out as a group send; the permanent failure // then triggered the fire-and-forget reply (its mock error is fine). assert_eq!( sender.calls(), vec!["send_chat_action", "send_media_group", "send_message"] ); // The entry is degraded, not dropped: the next request re-sends from // the source URL without a fetch. let entry = stores .link_cache() .get("twitter:1", Duration::from_secs(3600)) .await .expect("a stale file id must not cost the whole entry"); assert!(entry.media[0].file_id.is_empty()); assert_eq!(entry.media[0].url, "https://pbs.twimg.com/media/photo.jpg"); } /// A degraded entry (see the test above) sends the source URL: no fetch, /// no download of the post's data, and no upload of our own — Telegram /// fetches the media it is pointed at. The absence of a `send_message` /// (which the fetch-error path would emit) is what proves no fetch ran. #[tokio::test] async fn a_degraded_entry_sends_by_url_without_fetching() { let stores = TestStores::new(); let sender = MockSender::scripted(vec![Outcome::GroupOk], permanent_error); let ctx = stores.ctx(&sender); let mut cached = cached_photo(); cached.media[0].file_id.clear(); stores.link_cache().put("twitter:1", &cached).await; url_media(&ctx, 1, 2, "https://x.com/u/status/1", PostSend::FromChat).await; assert_eq!(sender.calls(), vec!["send_chat_action", "send_media_group"]); // Still cached, still degraded: a degraded entry keeps serving. let entry = stores .link_cache() .get("twitter:1", Duration::from_secs(3600)) .await .expect("the entry must stay"); assert!(entry.media[0].file_id.is_empty()); } /// A repeat request for a single-gif post re-sends the cached file id /// instead of failing: the animation path used to hand the id to /// `input_file_for`, which read it as a local path and answered "local /// media file missing" (permanent), so the second request of a gif link /// always failed and only the third — from the degraded entry — worked. #[tokio::test] async fn a_cached_animation_sends_by_file_id() { let stores = TestStores::new(); let sender = MockSender::scripted(vec![Outcome::AnimationOk], permanent_error); let ctx = stores.ctx(&sender); let mut entry = cached_photo(); entry.media = vec![CachedMedia { kind: CachedMediaKind::Animation, file_id: "AgAC-gif".into(), url: "https://p/1.gif".into(), }]; stores.link_cache().put("twitter:1", &entry).await; url_media(&ctx, 1, 2, "https://x.com/u/status/1", PostSend::FromChat).await; assert_eq!(sender.calls(), vec!["send_chat_action", "send_animation"]); assert_eq!( sender.animation_files(), vec!["AgAC-gif"], "a cached gif must go out as its file id, not as an upload" ); } /// The two payload shapes a cached item can take, asserted directly: the /// file id when there is one, the source URL when the entry was degraded. #[test] fn cached_media_payload_prefers_the_file_id_over_the_url() { let with_id = CachedMedia { kind: CachedMediaKind::Photo, file_id: "AgAC".into(), url: "https://p/1.jpg".into(), }; match cached_media_payload(&with_id, true) { MediaItemPayload::Photo { media: MediaRef::FileId(media), has_spoiler, .. } => { assert_eq!(media, "AgAC", "a cached send must go by file id"); assert!(has_spoiler); } other => panic!("expected a photo payload by file id, got {other:?}"), } let degraded = CachedMedia { kind: CachedMediaKind::Video, file_id: String::new(), url: "https://v/1.mp4".into(), }; match cached_media_payload(°raded, false) { MediaItemPayload::Video { media: MediaRef::Source(media), fallback_url, .. } => { assert_eq!(media, "https://v/1.mp4", "a degraded send goes by URL"); assert!(fallback_url.is_none()); } other => panic!("expected a video payload by URL, got {other:?}"), } } #[test] fn media_to_payload_rejects_an_upstream_local_path() { let media = Media::Illustration { url: "/etc/passwd".into(), thumbnail_url: Some("/etc/passwd".into()), fallback_url: Some("https://safe.example/fallback.jpg".into()), }; assert!(media_to_payload(&media, false).is_none()); } /// The caption-quote threshold matches the post's text inside the caption, /// so a long-text cache hit is quoted and a short-text one is not. #[tokio::test] async fn cache_hit_quotes_a_long_text_caption() { let mut stores = TestStores::new(); stores.config_mut().caption_quote_text_chars = 3; let prefix = "https://x.com/u/status/1\na: "; for (text, expected) in [ ( "abc", format!("{prefix}
abc
"), ), ("ab", format!("{prefix}ab")), ] { let sender = MockSender::scripted(vec![Outcome::GroupOk], permanent_error); let ctx = stores.ctx(&sender); let mut entry = cached_photo(); entry.caption = format!("{prefix}{text}"); entry.title = String::new(); entry.content = text.into(); stores.link_cache().put("twitter:1", &entry).await; url_media(&ctx, 1, 2, "https://x.com/u/status/1", PostSend::FromChat).await; assert_eq!(sender.captions(), vec![expected], "text {text:?}"); } } /// A post with no media of its own is delivered as its text instead of /// "No media found" (live: reaching the branch needs a real fetch, since /// `Fetched` cannot be built outside x-media). Run with /// `cargo test -p xmedia-bot -- --ignored live`. #[tokio::test] #[ignore = "live network: fetches the post from its site"] async fn live_a_text_only_link_is_sent_as_text() { let stores = TestStores::new(); let sender = MockSender::scripted(vec![Outcome::MessageOk], permanent_error); let ctx = stores.ctx(&sender); url_media( &ctx, 1, 2, "https://x.com/i/status/1992471125734142256", PostSend::FromChat, ) .await; assert_eq!( sender.calls(), vec!["send_chat_action", "send_message"], "a media-less post is one message, not a media send" ); let text = &sender.messages()[0]; assert!(text.contains("https://x.com/"), "{text}"); assert!(text.contains("1992471125734142256"), "{text}"); } #[tokio::test] async fn unsupported_url_is_ignored_silently() { let stores = TestStores::new(); let sender = MockSender::scripted(vec![], permanent_error); let ctx = stores.ctx(&sender); // No cache key → the fetch dispatcher returns Ok(None) without any // network; nothing is sent or replied. url_media( &ctx, 1, 2, "https://example.com/not-a-post", PostSend::FromChat, ) .await; assert_eq!(sender.calls(), vec!["send_chat_action"]); } // ── One fetch per post ─────────────────────────────────────────────── /// Two callers asking for the same post while its fetch is in flight run /// one fetch between them: the duplicate (a second chat, a batch forward /// and a retry) waits for that result instead of paying for its own. #[tokio::test] async fn concurrent_callers_share_one_fetch() { let map = parking_lot::Mutex::new(HashMap::new()); let calls = std::sync::atomic::AtomicUsize::new(0); let fetch = || async { calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); tokio::time::sleep(Duration::from_millis(50)).await; 7u32 }; let (first, second) = tokio::join!( shared_fetch(&map, "twitter:1", fetch), shared_fetch(&map, "twitter:1", fetch) ); assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 1); assert_eq!(*first, 7); assert!(std::sync::Arc::ptr_eq(&first, &second), "one shared result"); assert!( map.lock().is_empty(), "the entry must not outlive the fetch" ); } /// The dedup is *concurrent* only. A caller arriving after the fetch /// settled fetches again: the source may have changed, and a failure is /// deliberately not cached (the user is told to try again). #[tokio::test] async fn a_later_call_fetches_again() { let map = parking_lot::Mutex::new(HashMap::new()); let calls = std::sync::atomic::AtomicUsize::new(0); let counting = |value: u32| { let calls = &calls; async move { calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); value } }; let first = shared_fetch(&map, "twitter:1", || counting(1)).await; let second = shared_fetch(&map, "twitter:1", || counting(2)).await; assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 2); assert_eq!((*first, *second), (1, 2)); assert!(!std::sync::Arc::ptr_eq(&first, &second)); } /// A cancelled fetch must not strand the callers that joined it: a live /// entry whose sharer is gone holds a sender, and the waiters would wait /// for a value that can never come. They fetch for themselves. #[tokio::test] async fn a_cancelled_fetch_does_not_strand_waiters() { let map = parking_lot::Mutex::new(HashMap::new()); let calls = std::sync::atomic::AtomicUsize::new(0); // A fetch that never finishes, cancelled by the timeout below once it // has installed its entry. let slow = shared_fetch(&map, "twitter:1", || async { calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); std::future::pending::<()>().await; 0u32 }); assert!( tokio::time::timeout(Duration::from_millis(50), slow) .await .is_err(), "the sharer must still be waiting when it is cancelled" ); // Its entry is gone, and the next caller fetches its own value. let value = shared_fetch(&map, "twitter:1", || async { calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); 5u32 }) .await; assert_eq!(*value, 5); assert!(map.lock().is_empty()); } // ── Send modes: the URL flow vs `/test` ───────────────────────────── /// A chat that has both post-send actions configured. async fn seed_post_send_settings(ctx: &AppContext<'_>) { ctx.chat_store .update(1, |data| { data.forward_channel_id = Some(2); data.edit_before_forward = true; }) .await .unwrap(); } #[tokio::test] async fn chat_settings_apply_to_the_normal_link_flow() { let stores = TestStores::new(); let sender = MockSender::scripted(vec![Outcome::GroupOk, Outcome::MessageOk], permanent_error); let ctx = stores.ctx(&sender); stores.link_cache().put("twitter:1", &cached_photo()).await; seed_post_send_settings(&ctx).await; url_media(&ctx, 1, 2, "https://x.com/u/status/1", PostSend::FromChat).await; // Media group, then the edit prompt (edit-before-forward wins over the // channel forward, which only runs once the prompt is confirmed). assert_eq!( sender.calls(), vec!["send_chat_action", "send_media_group", "send_message"] ); } #[tokio::test] async fn test_mode_sends_the_media_without_forwarding_or_editing() { let stores = TestStores::new(); // Only the group send is scripted: any forward (copy_messages) or edit // prompt (send_message) would panic with "unexpected outcome". let sender = MockSender::scripted(vec![Outcome::GroupOk], permanent_error); let ctx = stores.ctx(&sender); stores.link_cache().put("twitter:1", &cached_photo()).await; seed_post_send_settings(&ctx).await; url_media(&ctx, 1, 2, "https://x.com/u/status/1", PostSend::Suppressed).await; assert_eq!(sender.calls(), vec!["send_chat_action", "send_media_group"]); // The send is otherwise ordinary: the post stays cached (this is also // the retention control for the eviction case above). assert!( stores .link_cache() .get("twitter:1", Duration::from_secs(3600)) .await .is_some() ); } #[test] fn send_mode_decides_whether_chat_actions_ride_along() { let chat = ChatData { forward_channel_id: Some(2), edit_before_forward: true, ..ChatData::default() }; let with_chat = build_send_task( &chat, 1, 2, "https://x.com/u/status/1".into(), "cap".into(), vec![], None, PostSend::FromChat, ); let Task::SendMediaSequence { edit_before_forward, forward_channel_id, .. } = with_chat else { panic!("expected a media sequence task"); }; assert!(edit_before_forward); assert_eq!(forward_channel_id, Some(2)); let suppressed = build_send_task( &chat, 1, 2, "https://x.com/u/status/1".into(), "cap".into(), vec![], None, PostSend::Suppressed, ); let Task::SendMediaSequence { edit_before_forward, forward_channel_id, notify_chat_id, .. } = suppressed else { panic!("expected a media sequence task"); }; assert!(!edit_before_forward, "`/test` must not open an edit prompt"); assert_eq!(forward_channel_id, None, "`/test` must not forward"); // Dead-letter notification still reaches the chat that asked. assert_eq!(notify_chat_id, Some(1)); } #[test] fn only_url_carrying_entities_yield_a_link() { use teloxide::types::MessageEntityKind; let post = "https://x.com/u/status/1"; assert_eq!( url_of(&MessageEntityKind::Url, post), Some(post.to_string()) ); // A text link keeps its target, not the words the user sees. assert_eq!( url_of( &MessageEntityKind::TextLink { url: url::Url::parse(post).unwrap() }, "the post" ), Some(post.to_string()) ); for kind in [ MessageEntityKind::Bold, MessageEntityKind::Code, MessageEntityKind::Hashtag, MessageEntityKind::Mention, ] { assert_eq!(url_of(&kind, "#nope"), None, "{kind:?}"); } // Plain text that is no URL still comes back when it is a `Url` entity // (the fetch layer ignores what no site claims, and the user gets the // one explanatory reply it produces). assert_eq!( url_of(&MessageEntityKind::Url, "https://example.com/x"), Some("https://example.com/x".to_string()) ); } #[test] fn dedupe_by_post_and_by_exact_text() { // Two variants of one post, and a text link to the same post: one entry, // the first one seen. assert_eq!( dedupe_urls(vec![ "https://x.com/u/status/1".into(), "https://x.com/u/status/1/photo/1".into(), "https://x.com/u/status/1".into(), ]), vec!["https://x.com/u/status/1"] ); // Different posts both survive, in order. assert_eq!( dedupe_urls(vec![ "https://x.com/a/status/1".into(), "https://x.com/b/status/2".into(), ]), vec!["https://x.com/a/status/1", "https://x.com/b/status/2"] ); // A URL no site claims: exact-string dedup only. assert_eq!( dedupe_urls(vec![ "https://example.com/a".into(), "https://example.com/a".into(), "https://example.com/b".into(), ]), vec!["https://example.com/a", "https://example.com/b"] ); assert!(dedupe_urls(vec![]).is_empty()); } #[test] fn fetch_errors_map_to_distinct_user_messages() { use x_media::site::FetchError; let disabled = fetch_error_message(&FetchError::Disabled { site: "pixiv" }); assert_eq!(disabled, "Pixiv support is disabled on this bot."); assert_eq!( fetch_error_message(&FetchError::NotFound), "Post not found (deleted, private or unavailable)." ); let sensitive = fetch_error_message(&FetchError::Sensitive); assert!(sensitive.contains("TWITTER_AUTH_TOKEN"), "{sensitive}"); let blocked = fetch_error_message(&FetchError::Blocked); assert!(blocked.contains("refused"), "{blocked}"); } #[test] fn action_hint_follows_the_media_kind() { use MediaItemPayload::{Animation, Photo, Video}; let photo = || Photo { media: MediaRef::Source("https://p/1.jpg".into()), has_spoiler: false, fallback_url: None, }; let video = || Video { media: MediaRef::Source("https://v/1.mp4".into()), has_spoiler: false, thumbnail: None, fallback_url: None, }; // Unknown before the fetch: the pipeline starts on "typing". assert!(matches!(ActionHint::Typing.action(), ChatAction::Typing)); assert!(matches!( ActionHint::for_items(&[photo()]).action(), ChatAction::UploadPhoto )); assert!(matches!( ActionHint::for_items(&[ video(), Animation { media: MediaRef::Source("https://v/2.mp4".into()), has_spoiler: false, } ]) .action(), ChatAction::UploadVideo )); // A mixed post takes the photo label: `photos_first` always leads with // a photo, which is what Telegram shows. assert!(matches!( ActionHint::for_items(&[video(), photo()]).action(), ChatAction::UploadPhoto )); } #[tokio::test(start_paused = true)] async fn a_long_pipeline_keeps_the_chat_action_alive() { let sender = MockSender::scripted(vec![], permanent_error); let hint = parking_lot::Mutex::new(ActionHint::Typing); // Three refresh windows of work: Telegram would have dropped the // indicator twice without the keep-alive. let pipeline = async { tokio::time::sleep(ACTION_REFRESH * 3).await }; run_with_chat_action(&sender, 1, &hint, pipeline).await; let actions = sender .calls() .iter() .filter(|call| **call == "send_chat_action") .count(); assert_eq!(actions, 3, "expected the initial action plus two refreshes"); } }