refactor(send): introduce MediaSender seam; add mock-based fallback tests

This commit is contained in:
2026-08-14 22:04:16 +08:00
parent c9e72fda70
commit 50206a9056
3 changed files with 464 additions and 52 deletions
+1
View File
@@ -11,6 +11,7 @@ mod config;
mod db;
mod handlers;
mod link_cache;
mod media_sender;
mod photo;
mod queue;
mod send;
+247
View File
@@ -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<Box<dyn Future<Output = T> + 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<InputMedia>,
) -> BoxFuture<'_, Result<Vec<Message>, 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<Message, RequestError>>;
/// Copies messages between chats (forward to channel).
fn copy_messages(
&self,
to: ChatId,
from: ChatId,
ids: Vec<MessageId>,
) -> BoxFuture<'_, Result<Vec<MessageId>, 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<MessageId>,
reply_markup: Option<InlineKeyboardMarkup>,
) -> BoxFuture<'_, Result<Message, RequestError>>;
}
impl MediaSender for Bot {
fn send_media_group(
&self,
chat_id: ChatId,
reply_to: MessageId,
items: Vec<InputMedia>,
) -> BoxFuture<'_, Result<Vec<Message>, RequestError>> {
Box::pin(async move {
// `<Bot as Requester>::` disambiguates from this trait's same-named
// method (teloxide's API lives in the `Requester` trait).
<Bot as Requester>::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<Message, RequestError>> {
Box::pin(async move {
let mut request = <Bot as Requester>::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<MessageId>,
) -> BoxFuture<'_, Result<Vec<MessageId>, RequestError>> {
Box::pin(async move { <Bot as Requester>::copy_messages(self, to, from, ids).await })
}
fn send_message(
&self,
chat_id: ChatId,
text: String,
reply_to: Option<MessageId>,
reply_markup: Option<InlineKeyboardMarkup>,
) -> BoxFuture<'_, Result<Message, RequestError>> {
Box::pin(async move {
let mut request = <Bot as Requester>::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<Vec<Outcome>>,
cursor: Mutex<usize>,
calls: Mutex<Vec<&'static str>>,
/// Builds the error every `*Err` outcome returns (RequestError is not
/// cloneable, so the factory recreates it per call).
error: Box<dyn Fn() -> RequestError + Send + Sync>,
}
impl MockSender {
pub(crate) fn scripted(
script: Vec<Outcome>,
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<InputMedia>,
) -> BoxFuture<'_, Result<Vec<Message>, 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<Message, RequestError>> {
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<MessageId>,
) -> BoxFuture<'_, Result<Vec<MessageId>, 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<MessageId>,
_reply_markup: Option<InlineKeyboardMarkup>,
) -> BoxFuture<'_, Result<Message, RequestError>> {
Box::pin(async move { panic!("send_message should not be called in these tests") })
}
}
}
+216 -52
View File
@@ -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<Vec<i64>, SendError> {
pub async fn send_media_sequence(
sender: &dyn MediaSender,
task: &Task,
) -> Result<Vec<i64>, 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<Vec<i64>, 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<Vec<i64>, 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<Vec<i64>, Sen
}
async fn send_animation_inner(
bot: &Bot,
sender: &dyn MediaSender,
chat_id: i64,
reply_to: i64,
caption: &str,
spoiler: bool,
file: InputFile,
) -> Result<Message, RequestError> {
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<Vec<i64>, SendError> {
pub async fn send_animation(sender: &dyn MediaSender, task: &Task) -> Result<Vec<i64>, SendError> {
let Task::SendAnimation {
chat_id,
reply_to_message_id,
@@ -1062,7 +1059,7 @@ pub async fn send_animation(bot: &Bot, task: &Task) -> Result<Vec<i64>, 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<Vec<i64>, 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<Vec<i64>, 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::<Vec<_>>();
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<String, String>) -> 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<i64>,
message_id: Option<i64>,
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<i64>) {
pub async fn post_send_actions(sender: &dyn MediaSender, task: &Task, message_ids: Vec<i64>) {
let (
chat_id,
reply_to,
@@ -1231,14 +1226,15 @@ pub async fn post_send_actions(bot: &Bot, task: &Task, message_ids: Vec<i64>) {
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<i64>) {
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<i64>) {
}
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<Vec<i64>, SendError> {
async fn send_media_or_animation(
sender: &dyn MediaSender,
task: &Task,
) -> Result<Vec<i64>, 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 { .. })
));
}
}