fix: split long message forwards

This commit is contained in:
2026-09-25 01:49:29 +08:00
parent d8b9a06453
commit 7b816d4153
4 changed files with 89 additions and 31 deletions
@@ -107,6 +107,7 @@ async fn handle_callback(
from_chat_id: edit.chat_id, from_chat_id: edit.chat_id,
to_chat_id: channel_id, to_chat_id: channel_id,
message_ids: edit.forward_message_ids.clone(), message_ids: edit.forward_message_ids.clone(),
forward_offset: 0,
notify_chat_id: Some(chat_id), notify_chat_id: Some(chat_id),
notify_message_id: Some(prompt_message_id), notify_message_id: Some(prompt_message_id),
}; };
+1
View File
@@ -253,6 +253,7 @@ mod tests {
from_chat_id: 1, from_chat_id: 1,
to_chat_id: 2, to_chat_id: 2,
message_ids: vec![3], message_ids: vec![3],
forward_offset: 0,
notify_chat_id: None, notify_chat_id: None,
notify_message_id: None, notify_message_id: None,
})); }));
+77 -22
View File
@@ -150,6 +150,8 @@ pub enum Task {
from_chat_id: i64, from_chat_id: i64,
to_chat_id: i64, to_chat_id: i64,
message_ids: Vec<i64>, message_ids: Vec<i64>,
#[serde(default)]
forward_offset: usize,
notify_chat_id: Option<i64>, notify_chat_id: Option<i64>,
notify_message_id: Option<i64>, notify_message_id: Option<i64>,
}, },
@@ -698,32 +700,38 @@ pub(crate) async fn send_text_post(
} }
} }
/// Copies already-sent messages to the forward channel. No download fallback:
/// the files are already on Telegram's servers.
pub async fn forward_messages(ctx: &AppContext<'_>, task: &Task) -> Result<(), SendError> { pub async fn forward_messages(ctx: &AppContext<'_>, task: &Task) -> Result<(), SendError> {
let Task::ForwardMessages { let Task::ForwardMessages {
from_chat_id, from_chat_id,
to_chat_id, to_chat_id,
message_ids, message_ids,
forward_offset,
.. ..
} = task } = task
else { else {
unreachable!("forward_messages requires a ForwardMessages task") unreachable!("forward_messages requires a ForwardMessages task")
}; };
let message_ids = message_ids let mut offset = (*forward_offset).min(message_ids.len());
while offset < message_ids.len() {
let end = (offset + 100).min(message_ids.len());
let ids = message_ids[offset..end]
.iter() .iter()
.map(|id| MessageId(*id as i32)) .copied()
.collect::<Vec<_>>(); .map(|id| MessageId(id as i32))
match ctx .collect();
if let Err(e) = ctx
.sender .sender
.copy_messages( .copy_messages(ChatId(*to_chat_id), ChatId(*from_chat_id), ids)
ChatId(*to_chat_id),
ChatId(*from_chat_id),
message_ids.clone(),
)
.await .await
{ {
Ok(_) => { let mut retry_task = task.clone();
if let Task::ForwardMessages { forward_offset, .. } = &mut retry_task {
*forward_offset = offset;
}
return Err(classify_to_send_error(&e, retry_task, "media fetch failed"));
}
offset = end;
}
log::info!( log::info!(
"copied {} message(s) from {} to {}", "copied {} message(s) from {} to {}",
message_ids.len(), message_ids.len(),
@@ -731,13 +739,6 @@ pub async fn forward_messages(ctx: &AppContext<'_>, task: &Task) -> Result<(), S
to_chat_id to_chat_id
); );
Ok(()) Ok(())
}
Err(e) => Err(classify_to_send_error(
&e,
task.clone(),
"media fetch failed",
)),
}
} }
#[cfg(test)] #[cfg(test)]
@@ -950,6 +951,7 @@ mod tests {
from_chat_id: 1, from_chat_id: 1,
to_chat_id: 2, to_chat_id: 2,
message_ids: vec![1], message_ids: vec![1],
forward_offset: 0,
notify_chat_id: None, notify_chat_id: None,
notify_message_id: None, notify_message_id: None,
}; };
@@ -1447,6 +1449,7 @@ mod tests {
from_chat_id: 1, from_chat_id: 1,
to_chat_id: 2, to_chat_id: 2,
message_ids: vec![3], message_ids: vec![3],
forward_offset: 0,
notify_chat_id: None, notify_chat_id: None,
notify_message_id: None, notify_message_id: None,
}; };
@@ -1475,11 +1478,63 @@ mod tests {
)); ));
} }
#[tokio::test]
async fn long_forward_is_split_into_telegram_batches() {
let sender = MockSender::scripted(vec![Outcome::CopyOk, Outcome::CopyOk], || {
api_error("unused")
});
let stores = TestStores::new();
let ctx = stores.ctx(&sender);
let task = Task::ForwardMessages {
from_chat_id: 1,
to_chat_id: 2,
message_ids: (1..=201).collect(),
forward_offset: 0,
notify_chat_id: None,
notify_message_id: None,
};
forward_messages(&ctx, &task).await.unwrap();
assert_eq!(
sender.calls(),
vec!["copy_messages", "copy_messages", "copy_messages"]
);
}
#[tokio::test]
async fn failed_forward_resumes_after_completed_batches() {
use teloxide::types::Seconds;
let sender = MockSender::scripted(vec![Outcome::CopyOk, Outcome::CopyErr], || {
RequestError::RetryAfter(Seconds::from_seconds(7))
});
let stores = TestStores::new();
let ctx = stores.ctx(&sender);
let task = Task::ForwardMessages {
from_chat_id: 1,
to_chat_id: 2,
message_ids: (1..=201).collect(),
forward_offset: 0,
notify_chat_id: None,
notify_message_id: None,
};
match forward_messages(&ctx, &task).await {
Err(SendError::Retryable { task, .. }) => {
let Task::ForwardMessages { forward_offset, .. } = *task else {
panic!("retry task is not a forward");
};
assert_eq!(forward_offset, 100);
}
other => panic!("expected retryable forward, got {other:?}"),
}
assert_eq!(
sender.calls(),
vec!["copy_messages", "copy_messages"],
"the first completed batch must not be replayed"
);
}
#[tokio::test] #[tokio::test]
async fn dead_letter_releases_keep_alive_temp_media() { async fn dead_letter_releases_keep_alive_temp_media() {
// A task that exhausts its retries is dead-lettered by the queue
// without the handler running again: the keep-alive temp dir the
// fetch pipeline handed over must not outlive the task.
let dir = tempfile::tempdir().unwrap(); let dir = tempfile::tempdir().unwrap();
let file = dir.path().join("ugoira.mp4"); let file = dir.path().join("ugoira.mp4");
std::fs::write(&file, b"not-a-real-mp4").unwrap(); std::fs::write(&file, b"not-a-real-mp4").unwrap();
+1
View File
@@ -363,6 +363,7 @@ pub(crate) async fn post_send_actions(ctx: &AppContext<'_>, task: &Task, message
from_chat_id: chat_id, from_chat_id: chat_id,
to_chat_id: channel_id, to_chat_id: channel_id,
message_ids, message_ids,
forward_offset: 0,
notify_chat_id, notify_chat_id,
notify_message_id, notify_message_id,
}; };