mirror of
https://github.com/TheFunny/TelegramTwitterMediaBot.git
synced 2026-10-07 01:32:13 +00:00
perf(x-media): stream the frame zip through a tokio file handle
download_media_to_file took &mut std::fs::File and wrote every network chunk with a sync write_all on the executor thread — for a ugoira frame zip (up to 512 MiB) that is the whole download stalling a runtime worker, while bsky's remux had already been moved to tokio::fs for exactly this reason (its comment: a multi-megabyte std::fs::write blocks the executor thread). The signature takes &mut tokio::fs::File now; the single caller (pixiv's ugoira path) clones the NamedTempFile's handle — the clone shares the file offset, so the ZipArchive extraction in spawn_blocking reads what the download wrote — and drops it after the download hands the bytes to the OS.
This commit is contained in:
@@ -29,11 +29,8 @@ 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
|
||||||
/// of the post restarts the download.
|
/// of the post restarts the download.
|
||||||
fn download_too_slow() -> FetchError {
|
fn download_too_slow(total: Duration) -> FetchError {
|
||||||
FetchError::Transient(format!(
|
FetchError::Transient(format!("download exceeded {}s", total.as_secs()))
|
||||||
"download exceeded {}s",
|
|
||||||
DOWNLOAD_TOTAL_TIMEOUT.as_secs()
|
|
||||||
))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Builds a client with the shared configuration (browser User-Agent, the
|
/// Builds a client with the shared configuration (browser User-Agent, the
|
||||||
@@ -247,7 +244,7 @@ pub async fn download_media_limited(url: &str, max_bytes: u64) -> Result<bytes::
|
|||||||
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() > DOWNLOAD_TOTAL_TIMEOUT {
|
||||||
return Err(download_too_slow());
|
return Err(download_too_slow(DOWNLOAD_TOTAL_TIMEOUT));
|
||||||
}
|
}
|
||||||
buf.extend_from_slice(&chunk);
|
buf.extend_from_slice(&chunk);
|
||||||
if buf.len() as u64 > max_bytes {
|
if buf.len() as u64 > max_bytes {
|
||||||
@@ -261,14 +258,15 @@ pub async fn download_media_limited(url: &str, max_bytes: u64) -> Result<bytes::
|
|||||||
/// moment the body crosses `max_bytes` (or when a declared Content-Length
|
/// moment the body crosses `max_bytes` (or when a declared Content-Length
|
||||||
/// already exceeds it). Unlike [`download_media_limited`] the body is never
|
/// already exceeds it). Unlike [`download_media_limited`] the body is never
|
||||||
/// buffered in memory — used for large files (e.g. the pixiv ugoira frame
|
/// buffered in memory — used for large files (e.g. the pixiv ugoira frame
|
||||||
/// zip, which can be hundreds of MB) that would otherwise spike RAM.
|
/// zip, which can be hundreds of MB) that would otherwise spike RAM. Writes
|
||||||
/// Returns the number of bytes written.
|
/// go through the tokio handle so a sync write never stalls an executor
|
||||||
|
/// thread for the length of the download. Returns the number of bytes written.
|
||||||
pub async fn download_media_to_file(
|
pub async fn download_media_to_file(
|
||||||
url: &str,
|
url: &str,
|
||||||
max_bytes: u64,
|
max_bytes: u64,
|
||||||
out: &mut std::fs::File,
|
out: &mut tokio::fs::File,
|
||||||
) -> Result<u64, FetchError> {
|
) -> Result<u64, FetchError> {
|
||||||
use std::io::Write;
|
use tokio::io::AsyncWriteExt;
|
||||||
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
|
||||||
@@ -280,13 +278,13 @@ pub async fn download_media_to_file(
|
|||||||
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() > DOWNLOAD_TOTAL_TIMEOUT {
|
||||||
return Err(download_too_slow());
|
return Err(download_too_slow(DOWNLOAD_TOTAL_TIMEOUT));
|
||||||
}
|
}
|
||||||
total += chunk.len() as u64;
|
total += chunk.len() as u64;
|
||||||
if total > max_bytes {
|
if total > max_bytes {
|
||||||
return Err(FetchError::TooLarge);
|
return Err(FetchError::TooLarge);
|
||||||
}
|
}
|
||||||
out.write_all(&chunk).map_err(FetchError::Io)?;
|
out.write_all(&chunk).await.map_err(FetchError::Io)?;
|
||||||
}
|
}
|
||||||
Ok(total)
|
Ok(total)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -224,12 +224,23 @@ impl PixivAPI {
|
|||||||
// Stream the frame zip to a temp file instead of buffering it in
|
// Stream the frame zip to a temp file instead of buffering it in
|
||||||
// memory: ugoira zips can be hundreds of MB, and the old
|
// memory: ugoira zips can be hundreds of MB, and the old
|
||||||
// download_media_limited path spiked RAM up to the size cap.
|
// download_media_limited path spiked RAM up to the size cap.
|
||||||
let mut zip_file = tempfile::Builder::new()
|
let zip_file = tempfile::Builder::new()
|
||||||
.prefix(crate::TEMP_FILE_PREFIX)
|
.prefix(crate::TEMP_FILE_PREFIX)
|
||||||
.suffix(".zip")
|
.suffix(".zip")
|
||||||
.tempfile()
|
.tempfile()
|
||||||
.map_err(|e| PixivError::Api(format!("temp zip failed: {e}")))?;
|
.map_err(|e| PixivError::Api(format!("temp zip failed: {e}")))?;
|
||||||
crate::site::download_media_to_file(&zip_url, 512 * 1024 * 1024, zip_file.as_file_mut())
|
// Stream through a tokio handle: a sync write per chunk would stall
|
||||||
|
// an executor thread for the whole (up to 512 MiB) download. The
|
||||||
|
// clone shares the file offset with `zip_file`, so the extraction
|
||||||
|
// below reads what was written, and dropping it after the download
|
||||||
|
// hands every byte to the OS.
|
||||||
|
let mut zip_out = tokio::fs::File::from_std(
|
||||||
|
zip_file
|
||||||
|
.as_file()
|
||||||
|
.try_clone()
|
||||||
|
.map_err(|e| PixivError::Api(format!("temp zip clone failed: {e}")))?,
|
||||||
|
);
|
||||||
|
crate::site::download_media_to_file(&zip_url, 512 * 1024 * 1024, &mut zip_out)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| match e {
|
.map_err(|e| match e {
|
||||||
FetchError::Http(e) => PixivError::Http(e),
|
FetchError::Http(e) => PixivError::Http(e),
|
||||||
@@ -243,6 +254,7 @@ impl PixivAPI {
|
|||||||
}
|
}
|
||||||
other => PixivError::Api(format!("frame zip download failed: {other}")),
|
other => PixivError::Api(format!("frame zip download failed: {other}")),
|
||||||
})?;
|
})?;
|
||||||
|
drop(zip_out);
|
||||||
let frame_delays = metadata.frames.iter().map(|f| f.delay).collect::<Vec<_>>();
|
let frame_delays = metadata.frames.iter().map(|f| f.delay).collect::<Vec<_>>();
|
||||||
let result =
|
let result =
|
||||||
tokio::task::spawn_blocking(move || -> Result<(String, tempfile::TempDir), String> {
|
tokio::task::spawn_blocking(move || -> Result<(String, tempfile::TempDir), String> {
|
||||||
|
|||||||
Reference in New Issue
Block a user