diff --git a/crates/xmedia-bot/src/main.rs b/crates/xmedia-bot/src/main.rs index c8fefe4..4d451cc 100644 --- a/crates/xmedia-bot/src/main.rs +++ b/crates/xmedia-bot/src/main.rs @@ -11,6 +11,7 @@ mod config; mod db; mod handlers; mod link_cache; +mod media_sender; mod photo; mod queue; mod send; diff --git a/crates/xmedia-bot/src/media_sender.rs b/crates/xmedia-bot/src/media_sender.rs new file mode 100644 index 0000000..c95d3e2 --- /dev/null +++ b/crates/xmedia-bot/src/media_sender.rs @@ -0,0 +1,247 @@ +//! Send abstraction: the message-sending surface [`send`](crate::send) +//! needs, so the send pipeline can be tested with a scripted mock instead of +//! a live teloxide `Bot`. + +use std::future::Future; +use std::pin::Pin; +use teloxide::RequestError; +use teloxide::prelude::Requester; +use teloxide::prelude::*; +use teloxide::types::{ + ChatId, InlineKeyboardMarkup, InputFile, InputMedia, Message, MessageId, ParseMode, + ReplyParameters, +}; + +/// Boxed, `Send` future returned by a [`MediaSender`] method (`async fn` in +/// traits is not dyn-compatible). +type BoxFuture<'a, T> = Pin + Send + 'a>>; + +/// The message-sending surface the send pipeline uses. The production +/// implementation is teloxide's [`Bot`]; tests inject a scripted mock to +/// cover the fallback and classification logic without touching the +/// Telegram API. +pub trait MediaSender: Send + Sync { + /// Sends a media group, replying to `reply_to`. + fn send_media_group( + &self, + chat_id: ChatId, + reply_to: MessageId, + items: Vec, + ) -> BoxFuture<'_, Result, RequestError>>; + + /// Sends a lone animation, replying to `reply_to`. + fn send_animation<'a>( + &'a self, + chat_id: ChatId, + reply_to: MessageId, + caption: &'a str, + spoiler: bool, + file: InputFile, + ) -> BoxFuture<'a, Result>; + + /// Copies messages between chats (forward to channel). + fn copy_messages( + &self, + to: ChatId, + from: ChatId, + ids: Vec, + ) -> BoxFuture<'_, Result, RequestError>>; + + /// Sends a plain text message, optionally replying to `reply_to` and + /// attaching `reply_markup`. + fn send_message( + &self, + chat_id: ChatId, + text: String, + reply_to: Option, + reply_markup: Option, + ) -> BoxFuture<'_, Result>; +} + +impl MediaSender for Bot { + fn send_media_group( + &self, + chat_id: ChatId, + reply_to: MessageId, + items: Vec, + ) -> BoxFuture<'_, Result, RequestError>> { + Box::pin(async move { + // `::` disambiguates from this trait's same-named + // method (teloxide's API lives in the `Requester` trait). + ::send_media_group(self, chat_id, items) + .reply_parameters(ReplyParameters::new(reply_to).allow_sending_without_reply()) + .await + }) + } + + fn send_animation<'a>( + &'a self, + chat_id: ChatId, + reply_to: MessageId, + caption: &'a str, + spoiler: bool, + file: InputFile, + ) -> BoxFuture<'a, Result> { + Box::pin(async move { + let mut request = ::send_animation(self, chat_id, file) + .caption(caption) + .parse_mode(ParseMode::Html) + .reply_parameters(ReplyParameters::new(reply_to).allow_sending_without_reply()); + if spoiler { + request = request.has_spoiler(true); + } + request.await + }) + } + + fn copy_messages( + &self, + to: ChatId, + from: ChatId, + ids: Vec, + ) -> BoxFuture<'_, Result, RequestError>> { + Box::pin(async move { ::copy_messages(self, to, from, ids).await }) + } + + fn send_message( + &self, + chat_id: ChatId, + text: String, + reply_to: Option, + reply_markup: Option, + ) -> BoxFuture<'_, Result> { + Box::pin(async move { + let mut request = ::send_message(self, chat_id, text); + if let Some(reply_to) = reply_to { + request = request + .reply_parameters(ReplyParameters::new(reply_to).allow_sending_without_reply()); + } + if let Some(markup) = reply_markup { + request = request.reply_markup(markup); + } + request.await + }) + } +} + +/// Test support: a scripted [`MediaSender`] mock (no Telegram API involved). +#[cfg(test)] +pub(crate) mod test_support { + use super::*; + use std::sync::Mutex; + + /// One scripted outcome, consumed front-to-back; the last entry repeats + /// for further calls of the same method kind. + #[derive(Clone, Copy, Debug, PartialEq, Eq)] + pub(crate) enum Outcome { + GroupOk, + GroupErr, + AnimationErr, + CopyOk, + CopyErr, + } + + /// Replays a script and records the method names that were called. + pub(crate) struct MockSender { + script: Mutex>, + cursor: Mutex, + calls: Mutex>, + /// Builds the error every `*Err` outcome returns (RequestError is not + /// cloneable, so the factory recreates it per call). + error: Box RequestError + Send + Sync>, + } + + impl MockSender { + pub(crate) fn scripted( + script: Vec, + error: impl Fn() -> RequestError + Send + Sync + 'static, + ) -> Self { + MockSender { + script: Mutex::new(script), + cursor: Mutex::new(0), + calls: Mutex::new(Vec::new()), + error: Box::new(error), + } + } + + /// Method names in call order (e.g. `["send_media_group", + /// "send_media_group"]` proves the fallback re-sent). + pub(crate) fn calls(&self) -> Vec<&'static str> { + self.calls.lock().unwrap().clone() + } + + fn next(&self, kind: &'static str) -> Outcome { + self.calls.lock().unwrap().push(kind); + let script = self.script.lock().unwrap(); + let mut cursor = self.cursor.lock().unwrap(); + if script.is_empty() { + panic!("mock script exhausted: {kind}"); + } + let idx = (*cursor).min(script.len() - 1); + *cursor = idx + 1; + script[idx] + } + + fn error(&self) -> RequestError { + (self.error)() + } + } + + impl MediaSender for MockSender { + fn send_media_group( + &self, + _chat_id: ChatId, + _reply_to: MessageId, + _items: Vec, + ) -> BoxFuture<'_, Result, RequestError>> { + Box::pin(async move { + match self.next("send_media_group") { + Outcome::GroupOk => Ok(Vec::new()), + Outcome::GroupErr => Err(self.error()), + other => panic!("unexpected outcome {other:?} for send_media_group"), + } + }) + } + + fn send_animation<'a>( + &'a self, + _chat_id: ChatId, + _reply_to: MessageId, + _caption: &'a str, + _spoiler: bool, + _file: InputFile, + ) -> BoxFuture<'a, Result> { + Box::pin(async move { + match self.next("send_animation") { + Outcome::AnimationErr => Err(self.error()), + other => panic!("unexpected outcome {other:?} for send_animation"), + } + }) + } + + fn copy_messages( + &self, + _to: ChatId, + _from: ChatId, + _ids: Vec, + ) -> BoxFuture<'_, Result, RequestError>> { + Box::pin(async move { + match self.next("copy_messages") { + Outcome::CopyOk => Ok(vec![MessageId(1)]), + Outcome::CopyErr => Err(self.error()), + other => panic!("unexpected outcome {other:?} for copy_messages"), + } + }) + } + + fn send_message( + &self, + _chat_id: ChatId, + _text: String, + _reply_to: Option, + _reply_markup: Option, + ) -> BoxFuture<'_, Result> { + Box::pin(async move { panic!("send_message should not be called in these tests") }) + } + } +} diff --git a/crates/xmedia-bot/src/send.rs b/crates/xmedia-bot/src/send.rs index 31a2c3d..6cc3d47 100644 --- a/crates/xmedia-bot/src/send.rs +++ b/crates/xmedia-bot/src/send.rs @@ -5,6 +5,7 @@ use crate::handlers::{CHAT_STORE, LINK_CACHE, TASK_QUEUE, log_key}; use crate::link_cache::{CachedMedia, CachedMediaKind, CachedPost}; +use crate::media_sender::MediaSender; use crate::photo::{self, MAX_UPLOAD_BYTES, PhotoPrep}; use crate::queue::QueueError; use crate::state::{EditMessage, unix_now}; @@ -15,7 +16,7 @@ use std::sync::LazyLock; use teloxide::prelude::*; use teloxide::types::{ ChatId, InlineKeyboardButton, InlineKeyboardMarkup, InputFile, InputMedia, InputMediaAnimation, - InputMediaPhoto, InputMediaVideo, Message, MessageId, ParseMode, ReplyParameters, + InputMediaPhoto, InputMediaVideo, Message, MessageId, ParseMode, }; use teloxide::{ApiError, RequestError}; use tempfile::NamedTempFile; @@ -380,6 +381,7 @@ pub fn classify_request_error(e: &RequestError) -> Classification { } } +#[derive(Debug)] pub enum SendError { Retryable { delay_seconds: f64, task: Task }, Permanent { message: String, task: Task }, @@ -807,7 +809,7 @@ async fn prepare_upload_item( /// original order. Returns the fallback-error without the task attached; /// callers wrap it with the updated task state. async fn send_batch_via_upload( - bot: &Bot, + sender: &dyn MediaSender, chat_id: i64, reply_to: i64, batch: &[MediaItemPayload], @@ -859,11 +861,8 @@ async fn send_batch_via_upload( .map(|m| m.expect("every upload item was prepared")) .collect(); // `keep_alive` holds the temp files until the group request completes. - let result = bot - .send_media_group(ChatId(chat_id), items) - .reply_parameters( - ReplyParameters::new(MessageId(reply_to as i32)).allow_sending_without_reply(), - ) + let result = sender + .send_media_group(ChatId(chat_id), MessageId(reply_to as i32), items) .await; drop(keep_alive); match result { @@ -918,7 +917,10 @@ fn updated_sequence_task(task: &Task, batch_index: usize, sent_message_ids: Vec< /// Sends the media batches starting at `task.batch_index`, extending /// `sent_message_ids`. Returns all sent message ids on full success; on /// failure returns a [`SendError`] whose task carries the resumed state. -pub async fn send_media_sequence(bot: &Bot, task: &Task) -> Result, SendError> { +pub async fn send_media_sequence( + sender: &dyn MediaSender, + task: &Task, +) -> Result, SendError> { let Task::SendMediaSequence { chat_id, reply_to_message_id, @@ -954,11 +956,8 @@ pub async fn send_media_sequence(bot: &Bot, task: &Task) -> Result, Sen }); } }; - match bot - .send_media_group(ChatId(chat_id), items) - .reply_parameters( - ReplyParameters::new(MessageId(reply_to as i32)).allow_sending_without_reply(), - ) + match sender + .send_media_group(ChatId(chat_id), MessageId(reply_to as i32), items) .await { Ok(messages) => { @@ -980,7 +979,7 @@ pub async fn send_media_sequence(bot: &Bot, task: &Task) -> Result, Sen .unwrap_or_else(|| "?".into()) ); match send_batch_via_upload( - bot, + sender, chat_id, reply_to, batch, @@ -1011,28 +1010,26 @@ pub async fn send_media_sequence(bot: &Bot, task: &Task) -> Result, Sen } async fn send_animation_inner( - bot: &Bot, + sender: &dyn MediaSender, chat_id: i64, reply_to: i64, caption: &str, spoiler: bool, file: InputFile, ) -> Result { - let mut request = bot - .send_animation(ChatId(chat_id), file) - .caption(caption) - .parse_mode(ParseMode::Html) - .reply_parameters( - ReplyParameters::new(MessageId(reply_to as i32)).allow_sending_without_reply(), - ); - if spoiler { - request = request.has_spoiler(true); - } - request.await + sender + .send_animation( + ChatId(chat_id), + MessageId(reply_to as i32), + caption, + spoiler, + file, + ) + .await } /// Sends a lone animation (gif), URL first with the download fallback. -pub async fn send_animation(bot: &Bot, task: &Task) -> Result, SendError> { +pub async fn send_animation(sender: &dyn MediaSender, task: &Task) -> Result, SendError> { let Task::SendAnimation { chat_id, reply_to_message_id, @@ -1062,7 +1059,7 @@ pub async fn send_animation(bot: &Bot, task: &Task) -> Result, SendErro }); } }; - match send_animation_inner(bot, chat_id, reply_to, caption, has_spoiler, url_file).await { + match send_animation_inner(sender, chat_id, reply_to, caption, has_spoiler, url_file).await { Ok(message) => { let id = message.id.0 as i64; cache_animation_send(task, &message).await; @@ -1089,7 +1086,7 @@ pub async fn send_animation(bot: &Bot, task: &Task) -> Result, SendErro // Hold the temp file until the request completes. let _keep_alive = keep_alive; match send_animation_inner( - bot, + sender, chat_id, reply_to, caption, @@ -1115,7 +1112,7 @@ pub async fn send_animation(bot: &Bot, task: &Task) -> Result, SendErro /// Copies already-sent messages to the forward channel. No download fallback: /// the files are already on Telegram's servers. -pub async fn forward_messages(bot: &Bot, task: &Task) -> Result<(), SendError> { +pub async fn forward_messages(sender: &dyn MediaSender, task: &Task) -> Result<(), SendError> { let Task::ForwardMessages { from_chat_id, to_chat_id, @@ -1129,7 +1126,7 @@ pub async fn forward_messages(bot: &Bot, task: &Task) -> Result<(), SendError> { .iter() .map(|id| MessageId(*id as i32)) .collect::>(); - match bot + match sender .copy_messages( ChatId(*to_chat_id), ChatId(*from_chat_id), @@ -1169,26 +1166,24 @@ pub fn build_edit_markup(templates: &HashMap) -> InlineKeyboardM /// Notifies a chat about a dead-lettered task (skips when `notify_chat_id` is /// absent). pub async fn notify_failure( - bot: &Bot, + sender: &dyn MediaSender, chat_id: Option, message_id: Option, message: &str, ) { let Some(chat_id) = chat_id else { return }; - let mut request = bot.send_message(ChatId(chat_id), message); - if let Some(message_id) = message_id { - request = request.reply_parameters( - ReplyParameters::new(MessageId(message_id as i32)).allow_sending_without_reply(), - ); - } - if let Err(e) = request.await { + let reply_to = message_id.map(|id| MessageId(id as i32)); + if let Err(e) = sender + .send_message(ChatId(chat_id), message.to_string(), reply_to, None) + .await + { log::error!("failed to notify about failed task: {e}"); } } /// After a successful send: either open the edit-before-forward prompt or /// forward to the configured channel (with retry/queue handling). -pub async fn post_send_actions(bot: &Bot, task: &Task, message_ids: Vec) { +pub async fn post_send_actions(sender: &dyn MediaSender, task: &Task, message_ids: Vec) { let ( chat_id, reply_to, @@ -1231,14 +1226,15 @@ pub async fn post_send_actions(bot: &Bot, task: &Task, message_ids: Vec) { if edit_before_forward { let keyboard = build_edit_markup(&CHAT_STORE.get(chat_id).await.template); - match bot - .send_message(ChatId(chat_id), "Reply to edit message.") - .reply_markup(keyboard) - .reply_parameters( - ReplyParameters::new(MessageId(reply_to as i32)).allow_sending_without_reply(), + let prompt = sender + .send_message( + ChatId(chat_id), + "Reply to edit message.".to_string(), + Some(MessageId(reply_to as i32)), + Some(keyboard), ) - .await - { + .await; + match prompt { Ok(prompt) => { log::info!( "edit-before-forward prompt {} opened for {} message(s)", @@ -1279,7 +1275,7 @@ pub async fn post_send_actions(bot: &Bot, task: &Task, message_ids: Vec) { notify_chat_id, notify_message_id, }; - match forward_messages(bot, &forward_task).await { + match forward_messages(sender, &forward_task).await { Ok(()) => {} Err(SendError::Retryable { delay_seconds, @@ -1297,7 +1293,7 @@ pub async fn post_send_actions(bot: &Bot, task: &Task, message_ids: Vec) { } Err(SendError::Permanent { message, .. }) => { notify_failure( - bot, + sender, notify_chat_id, notify_message_id, &format!("Task failed after retries: {message}"), @@ -1381,10 +1377,13 @@ pub async fn handle_task(payload: serde_json::Value) -> Result<(), QueueError> { } } -async fn send_media_or_animation(bot: &Bot, task: &Task) -> Result, SendError> { +async fn send_media_or_animation( + sender: &dyn MediaSender, + task: &Task, +) -> Result, SendError> { match task { - Task::SendMediaSequence { .. } => send_media_sequence(bot, task).await, - Task::SendAnimation { .. } => send_animation(bot, task).await, + Task::SendMediaSequence { .. } => send_media_sequence(sender, task).await, + Task::SendAnimation { .. } => send_animation(sender, task).await, Task::ForwardMessages { .. } => unreachable!(), } } @@ -1656,4 +1655,169 @@ mod tests { assert_eq!(sniff_ext(b"\x00\x00\x00\x18ftypisom"), "mp4"); assert_eq!(sniff_ext(b"something else"), "bin"); } + + // ── MediaSender-mock tests: fallback trigger + error classification ── + + use crate::media_sender::test_support::{MockSender, Outcome}; + + /// Telegram's "I could not fetch this URL" error, which triggers the + /// download-and-reupload fallback. + fn media_fetch_error() -> RequestError { + RequestError::Api(ApiError::Unknown("Bad Request: WEBPAGE_MEDIA_EMPTY".into())) + } + + fn sequence_task(media: &str) -> Task { + Task::SendMediaSequence { + chat_id: 1, + reply_to_message_id: 2, + caption: "cap".into(), + media_batches: vec![vec![MediaItemPayload::Photo { + media: media.to_string(), + has_spoiler: false, + fallback_url: None, + file_id: false, + }]], + batch_index: 0, + sent_message_ids: vec![], + source_url: "https://x.com/u/status/1".into(), + edit_before_forward: false, + forward_channel_id: None, + notify_chat_id: Some(1), + notify_message_id: Some(2), + cache_data: None, + } + } + + #[tokio::test] + async fn media_group_fetch_failure_falls_back_then_permanent() { + // A local file avoids any network in the fallback (the prep pipeline + // uploads local paths directly). The first group send fails with a + // media-fetch error → the download-reupload fallback runs → the + // reupload also fails → Permanent. + let dir = tempfile::tempdir().unwrap(); + let file = dir.path().join("media.jpg"); + std::fs::write(&file, b"not-a-real-jpeg").unwrap(); + let sender = MockSender::scripted( + vec![Outcome::GroupErr, Outcome::GroupErr], + media_fetch_error, + ); + let task = sequence_task(file.to_str().unwrap()); + let result = send_media_sequence(&sender, &task).await; + assert!( + matches!(result, Err(SendError::Permanent { .. })), + "got {result:?}" + ); + // Two group sends: the original + the fallback reupload. + assert_eq!(sender.calls(), vec!["send_media_group", "send_media_group"]); + } + + #[tokio::test] + async fn media_group_retry_after_classifies_retryable_without_fallback() { + use teloxide::types::Seconds; + // RetryAfter is not a media-fetch failure: no fallback, straight to a + // retryable error carrying the Telegram delay. + let dir = tempfile::tempdir().unwrap(); + let file = dir.path().join("media.jpg"); + std::fs::write(&file, b"not-a-real-jpeg").unwrap(); + let sender = MockSender::scripted(vec![Outcome::GroupErr], || { + RequestError::RetryAfter(Seconds::from_seconds(7)) + }); + let task = sequence_task(file.to_str().unwrap()); + let result = send_media_sequence(&sender, &task).await; + match result { + Err(SendError::Retryable { delay_seconds, .. }) => { + assert_eq!(delay_seconds, 7.0) + } + other => panic!("expected Retryable, got {other:?}"), + } + assert_eq!(sender.calls(), vec!["send_media_group"]); + } + + #[tokio::test] + async fn animation_fetch_failure_falls_back_then_permanent() { + let dir = tempfile::tempdir().unwrap(); + let file = dir.path().join("gif.mp4"); + std::fs::write(&file, b"not-a-real-mp4").unwrap(); + let sender = MockSender::scripted( + vec![Outcome::AnimationErr, Outcome::AnimationErr], + media_fetch_error, + ); + let task = Task::SendAnimation { + chat_id: 1, + reply_to_message_id: 2, + caption: "cap".into(), + animation: MediaItemPayload::Animation { + media: file.to_string_lossy().into_owned(), + has_spoiler: false, + file_id: false, + }, + source_url: "https://x.com/u/status/1".into(), + edit_before_forward: false, + forward_channel_id: None, + notify_chat_id: Some(1), + notify_message_id: Some(2), + cache_data: None, + }; + let result = send_animation(&sender, &task).await; + assert!( + matches!(result, Err(SendError::Permanent { .. })), + "got {result:?}" + ); + assert_eq!(sender.calls(), vec!["send_animation", "send_animation"]); + } + + #[tokio::test] + async fn media_group_success_and_forward_ok() { + // GroupOk: the group send succeeds (empty message list → no file ids + // collected, the batch counts as sent). CopyOk: the forward succeeds. + let dir = tempfile::tempdir().unwrap(); + let file = dir.path().join("media.jpg"); + std::fs::write(&file, b"not-a-real-jpeg").unwrap(); + let sender = MockSender::scripted(vec![Outcome::GroupOk], media_fetch_error); + let task = sequence_task(file.to_str().unwrap()); + let result = send_media_sequence(&sender, &task).await; + assert!(result.is_ok(), "got {result:?}"); + + let sender = MockSender::scripted(vec![Outcome::CopyOk], media_fetch_error); + let task = Task::ForwardMessages { + from_chat_id: 1, + to_chat_id: 2, + message_ids: vec![3], + notify_chat_id: None, + notify_message_id: None, + }; + assert!(forward_messages(&sender, &task).await.is_ok()); + } + + #[tokio::test] + async fn forward_classifies_retry_after_and_permanent() { + use teloxide::types::Seconds; + let task = Task::ForwardMessages { + from_chat_id: 1, + to_chat_id: 2, + message_ids: vec![3], + notify_chat_id: None, + notify_message_id: None, + }; + // RetryAfter → Retryable with the Telegram delay. + let sender = MockSender::scripted(vec![Outcome::CopyErr], || { + RequestError::RetryAfter(Seconds::from_seconds(7)) + }); + match forward_messages(&sender, &task).await { + Err(SendError::Retryable { delay_seconds, .. }) => { + assert_eq!(delay_seconds, 7.0) + } + other => panic!("expected Retryable, got {other:?}"), + } + // A generic API error → Permanent. + let sender = MockSender::scripted(vec![Outcome::CopyErr], || { + RequestError::Api(ApiError::Unknown( + "Bad Request: message is not modified".into(), + )) + }); + assert!(matches!( + forward_messages(&sender, &task).await, + Err(SendError::Permanent { .. }) + )); + } }