mirror of
https://github.com/TheFunny/TelegramTwitterMediaBot.git
synced 2026-09-29 00:12:12 +00:00
perf(x-media): give the slot-holding fallback download its own budget
download_media_limited had one hard-coded total (600s) for every caller, and its heaviest caller — the bot's upload fallback — holds a PREP slot (and its memory reservation) for the whole transfer: six slow-but-alive downloads (a byte every 29s satisfies the idle window) could stall the fallback chain for ten minutes, queue retries included. The budget is a parameter now: the fallback passes 300s of its own (50 MiB in 300s ≈ 1.4 Mbit/s; a slower link is better served by retrying toward the item's smaller URL than by pinning a slot), while bsky's in-fetch HLS segments keep the generous 600s DOWNLOAD_TOTAL_TIMEOUT, now pub(crate) and re-exported for them. download_too_slow reports whichever budget it got.
This commit is contained in:
@@ -131,10 +131,10 @@ fn concat_list(files: &mut [(usize, std::path::PathBuf)]) -> String {
|
|||||||
/// the end of a 500-segment video meant downloading the entire thing twice
|
/// the end of a 500-segment video meant downloading the entire thing twice
|
||||||
/// more, so the second attempt belongs on the request that actually failed.
|
/// more, so the second attempt belongs on the request that actually failed.
|
||||||
async fn fetch_hls(url: &str, cap: u64) -> Result<bytes::Bytes, String> {
|
async fn fetch_hls(url: &str, cap: u64) -> Result<bytes::Bytes, String> {
|
||||||
match crate::site::download_media_limited(url, cap).await {
|
match crate::site::download_media_limited(url, cap, crate::site::DOWNLOAD_TOTAL_TIMEOUT).await {
|
||||||
Err(FetchError::Http(_) | FetchError::Transient(_)) => {
|
Err(FetchError::Http(_) | FetchError::Transient(_)) => {
|
||||||
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
|
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
|
||||||
crate::site::download_media_limited(url, cap)
|
crate::site::download_media_limited(url, cap, crate::site::DOWNLOAD_TOTAL_TIMEOUT)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| e.to_string())
|
.map_err(|e| e.to_string())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -16,15 +16,18 @@ use std::time::Duration;
|
|||||||
/// [`DOWNLOAD_TOTAL_TIMEOUT`].
|
/// [`DOWNLOAD_TOTAL_TIMEOUT`].
|
||||||
const DOWNLOAD_IDLE_TIMEOUT: Duration = Duration::from_secs(30);
|
const DOWNLOAD_IDLE_TIMEOUT: Duration = Duration::from_secs(30);
|
||||||
|
|
||||||
/// Absolute ceiling for one media download, on top of the idle window. A server
|
/// Absolute ceiling for one media download, on top of the idle window: a
|
||||||
/// that drips a byte every 29 s keeps [`next_chunk`] satisfied indefinitely, and
|
/// server that drips a byte every 29 s keeps [`next_chunk`] satisfied
|
||||||
/// on the bot's side each such download holds one of the process-wide upload-prep
|
/// indefinitely, and a transfer that trickles forever holds whatever the
|
||||||
/// slots (`send::upload`'s `PREP_SLOTS`) for as long as it lasts. Generous on
|
/// caller pinned to it — a fetch permit for an in-flight post, a prep slot
|
||||||
/// purpose: the legitimate cases are big — an ugoira frame zip runs to hundreds
|
/// for the bot's upload fallback. Generous on purpose: the legitimate cases
|
||||||
/// of MB and an HLS remux pulls a whole video — and a slow link is not an error.
|
/// are big — an ugoira frame zip runs to hundreds of MB and an HLS remux
|
||||||
/// Checked between chunks, so a transfer that completes just over the budget is
|
/// pulls a whole video — so this is the budget for downloads *inside a
|
||||||
/// kept rather than thrown away.
|
/// fetch*, while the slot-holding fallback passes its own shorter one (see
|
||||||
const DOWNLOAD_TOTAL_TIMEOUT: Duration = Duration::from_secs(600);
|
/// [`download_media_limited`]'s `total`). Checked between chunks, so a
|
||||||
|
/// transfer that completes just over the budget is kept rather than thrown
|
||||||
|
/// away.
|
||||||
|
pub(crate) const DOWNLOAD_TOTAL_TIMEOUT: Duration = Duration::from_secs(600);
|
||||||
|
|
||||||
/// The error a download reports when it spends its whole budget without
|
/// The error a download reports when it spends its whole budget without
|
||||||
/// finishing. Retryable: the transfer may simply have been unlucky, and a retry
|
/// finishing. Retryable: the transfer may simply have been unlucky, and a retry
|
||||||
@@ -89,8 +92,8 @@ pub(crate) static CLIENT: LazyLock<reqwest::Client> =
|
|||||||
/// impossible to deliver at all (the size cap said 512 MiB, the clock said 30s).
|
/// impossible to deliver at all (the size cap said 512 MiB, the clock said 30s).
|
||||||
/// What a stalled connection cannot do is hang a worker: the head and every
|
/// What a stalled connection cannot do is hang a worker: the head and every
|
||||||
/// chunk are bounded by [`DOWNLOAD_IDLE_TIMEOUT`] (see [`next_chunk`]), and a
|
/// chunk are bounded by [`DOWNLOAD_IDLE_TIMEOUT`] (see [`next_chunk`]), and a
|
||||||
/// transfer that keeps trickling but never finishes is bounded by
|
/// transfer that keeps trickling but never finishes is bounded by the
|
||||||
/// [`DOWNLOAD_TOTAL_TIMEOUT`].
|
/// caller's total budget (see [`download_media_limited`]).
|
||||||
static MEDIA_CLIENT: LazyLock<reqwest::Client> = LazyLock::new(|| build_client(None));
|
static MEDIA_CLIENT: LazyLock<reqwest::Client> = LazyLock::new(|| build_client(None));
|
||||||
|
|
||||||
/// The error a download reports when it stops making progress.
|
/// The error a download reports when it stops making progress.
|
||||||
@@ -232,7 +235,16 @@ fn apply_media_headers(mut request: reqwest::RequestBuilder, url: &str) -> reqwe
|
|||||||
/// cannot fetch a media URL itself (hotlink protection), the bot downloads
|
/// cannot fetch a media URL itself (hotlink protection), the bot downloads
|
||||||
/// the file and uploads it via multipart. Site-appropriate headers come from
|
/// the file and uploads it via multipart. Site-appropriate headers come from
|
||||||
/// each site's `media_headers` (pixiv image hosts need `Referer`).
|
/// each site's `media_headers` (pixiv image hosts need `Referer`).
|
||||||
pub async fn download_media_limited(url: &str, max_bytes: u64) -> Result<bytes::Bytes, FetchError> {
|
///
|
||||||
|
/// `total` is this caller's whole-transfer budget. The bot's upload fallback
|
||||||
|
/// holds a prep slot (and its memory reservation) while this runs, so it
|
||||||
|
/// passes a shorter one of its own; bsky's in-fetch segments take the
|
||||||
|
/// generous [`super::DOWNLOAD_TOTAL_TIMEOUT`].
|
||||||
|
pub async fn download_media_limited(
|
||||||
|
url: &str,
|
||||||
|
max_bytes: u64,
|
||||||
|
total: Duration,
|
||||||
|
) -> Result<bytes::Bytes, FetchError> {
|
||||||
let response = send_download(media_request(url)?).await?;
|
let response = send_download(media_request(url)?).await?;
|
||||||
if let Some(len) = response.content_length()
|
if let Some(len) = response.content_length()
|
||||||
&& len > max_bytes
|
&& len > max_bytes
|
||||||
@@ -243,8 +255,8 @@ pub async fn download_media_limited(url: &str, max_bytes: u64) -> Result<bytes::
|
|||||||
let mut buf = Vec::new();
|
let mut buf = Vec::new();
|
||||||
let started = std::time::Instant::now();
|
let started = std::time::Instant::now();
|
||||||
while let Some(chunk) = next_chunk(&mut response).await? {
|
while let Some(chunk) = next_chunk(&mut response).await? {
|
||||||
if started.elapsed() > DOWNLOAD_TOTAL_TIMEOUT {
|
if started.elapsed() > total {
|
||||||
return Err(download_too_slow(DOWNLOAD_TOTAL_TIMEOUT));
|
return Err(download_too_slow(total));
|
||||||
}
|
}
|
||||||
buf.extend_from_slice(&chunk);
|
buf.extend_from_slice(&chunk);
|
||||||
if buf.len() as u64 > max_bytes {
|
if buf.len() as u64 > max_bytes {
|
||||||
@@ -358,7 +370,10 @@ mod tests {
|
|||||||
#[ignore = "live network: requires outbound HTTPS to httpbin.org"]
|
#[ignore = "live network: requires outbound HTTPS to httpbin.org"]
|
||||||
async fn live_redirect_into_the_hosts_network_is_refused() {
|
async fn live_redirect_into_the_hosts_network_is_refused() {
|
||||||
let url = "https://httpbin.org/redirect-to?url=http://169.254.169.254/latest/meta-data/";
|
let url = "https://httpbin.org/redirect-to?url=http://169.254.169.254/latest/meta-data/";
|
||||||
match download_media_limited(url, u64::MAX).await.unwrap_err() {
|
match download_media_limited(url, u64::MAX, DOWNLOAD_TOTAL_TIMEOUT)
|
||||||
|
.await
|
||||||
|
.unwrap_err()
|
||||||
|
{
|
||||||
// A policy refusal reaches the caller wrapped by reqwest.
|
// A policy refusal reaches the caller wrapped by reqwest.
|
||||||
FetchError::Http(e) => assert!(e.is_redirect(), "got {e}"),
|
FetchError::Http(e) => assert!(e.is_redirect(), "got {e}"),
|
||||||
FetchError::Blocked => {}
|
FetchError::Blocked => {}
|
||||||
@@ -375,13 +390,15 @@ mod tests {
|
|||||||
"http://169.254.169.254/latest/meta-data/",
|
"http://169.254.169.254/latest/meta-data/",
|
||||||
"http://127.0.0.1:9/secret",
|
"http://127.0.0.1:9/secret",
|
||||||
] {
|
] {
|
||||||
let err = download_media_limited(url, u64::MAX).await.unwrap_err();
|
let err = download_media_limited(url, u64::MAX, DOWNLOAD_TOTAL_TIMEOUT)
|
||||||
|
.await
|
||||||
|
.unwrap_err();
|
||||||
assert!(matches!(err, FetchError::Blocked), "{url}: got {err:?}");
|
assert!(matches!(err, FetchError::Blocked), "{url}: got {err:?}");
|
||||||
}
|
}
|
||||||
// A malformed URL is refused the same way instead of becoming a
|
// A malformed URL is refused the same way instead of becoming a
|
||||||
// retryable transport error.
|
// retryable transport error.
|
||||||
assert!(matches!(
|
assert!(matches!(
|
||||||
download_media_limited("not a url", u64::MAX)
|
download_media_limited("not a url", u64::MAX, DOWNLOAD_TOTAL_TIMEOUT)
|
||||||
.await
|
.await
|
||||||
.unwrap_err(),
|
.unwrap_err(),
|
||||||
FetchError::Blocked
|
FetchError::Blocked
|
||||||
@@ -409,7 +426,9 @@ mod tests {
|
|||||||
other => panic!("expected illustration media, got {other:?}"),
|
other => panic!("expected illustration media, got {other:?}"),
|
||||||
};
|
};
|
||||||
assert!(url.contains("i.pximg.net"));
|
assert!(url.contains("i.pximg.net"));
|
||||||
let bytes = download_media_limited(&url, u64::MAX).await.unwrap();
|
let bytes = download_media_limited(&url, u64::MAX, DOWNLOAD_TOTAL_TIMEOUT)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
assert!(!bytes.is_empty());
|
assert!(!bytes.is_empty());
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -22,7 +22,7 @@ pub mod twitter;
|
|||||||
|
|
||||||
pub use pixiv::PixivError;
|
pub use pixiv::PixivError;
|
||||||
|
|
||||||
pub(crate) use download::CLIENT;
|
pub(crate) use download::{CLIENT, DOWNLOAD_TOTAL_TIMEOUT};
|
||||||
pub use download::{download_media_limited, download_media_to_file};
|
pub use download::{download_media_limited, download_media_to_file};
|
||||||
|
|
||||||
/// The result of fetching a post: canonical URL, HTML caption, the post's
|
/// The result of fetching a post: canonical URL, HTML caption, the post's
|
||||||
|
|||||||
@@ -34,6 +34,17 @@ static PREP_SLOTS: LazyLock<tokio::sync::Semaphore> =
|
|||||||
/// post was lost.
|
/// post was lost.
|
||||||
pub(super) const MAX_MEDIA_UPLOAD_BYTES: u64 = 50 * 1024 * 1024;
|
pub(super) const MAX_MEDIA_UPLOAD_BYTES: u64 = 50 * 1024 * 1024;
|
||||||
|
|
||||||
|
/// Whole-transfer budget for one fallback download. The prep slot (and the
|
||||||
|
/// non-photo memory reservation) is held while this runs, and the idle window
|
||||||
|
/// alone lets a server drip one byte every 29 s forever — so this path caps
|
||||||
|
/// its own transfers well below the in-fetch default: 50 MiB in 300 s needs
|
||||||
|
/// about 1.4 Mbit/s, and a much slower link is better served by the retry
|
||||||
|
/// path toward the item's smaller fallback URL than by pinning a slot for
|
||||||
|
/// ten minutes.
|
||||||
|
/// ponytail: if slow-link reports show up, move the download out of the prep
|
||||||
|
/// slot (slot = decode/upload only) instead of raising this again.
|
||||||
|
const FALLBACK_DOWNLOAD_TOTAL: std::time::Duration = std::time::Duration::from_secs(300);
|
||||||
|
|
||||||
/// Infers a file extension from magic bytes so Telegram detects the mime type
|
/// Infers a file extension from magic bytes so Telegram detects the mime type
|
||||||
/// on multipart uploads.
|
/// on multipart uploads.
|
||||||
pub(super) fn sniff_ext(bytes: &[u8]) -> &'static str {
|
pub(super) fn sniff_ext(bytes: &[u8]) -> &'static str {
|
||||||
@@ -111,7 +122,13 @@ async fn download_to_temp(
|
|||||||
} else {
|
} else {
|
||||||
Some(photo::reserve_memory(MAX_MEDIA_UPLOAD_BYTES).await)
|
Some(photo::reserve_memory(MAX_MEDIA_UPLOAD_BYTES).await)
|
||||||
};
|
};
|
||||||
let bytes = match x_media::site::download_media_limited(media_url, limit).await {
|
let bytes = match x_media::site::download_media_limited(
|
||||||
|
media_url,
|
||||||
|
limit,
|
||||||
|
FALLBACK_DOWNLOAD_TOTAL,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
{
|
||||||
Ok(bytes) => bytes,
|
Ok(bytes) => bytes,
|
||||||
Err(e) => return Err(classify_download_error(e)),
|
Err(e) => return Err(classify_download_error(e)),
|
||||||
};
|
};
|
||||||
|
|||||||
Reference in New Issue
Block a user