perf: move audio metadata I/O to spawn_blocking + stream reextract_all

- Move all sync I/O (id3::Tag, mp3_duration) into spawn_blocking via
  extract_metadata_blocking() to avoid stalling Tokio worker threads
- Replace fetch_all with streaming .fetch() in reextract_all_audio_metadata
  for O(1) memory usage regardless of audio file count
- Consolidate get_duration_secs into extract_metadata_blocking, eliminating
  redundant file open (tag was read twice before)
This commit is contained in:
Diocrafts
2026-04-11 19:36:34 +02:00
parent 63bcd0ffe7
commit 78fcf5f08f
@@ -1,3 +1,4 @@
use futures::StreamExt;
use id3::{Tag, TagLike}; use id3::{Tag, TagLike};
use sqlx::{FromRow, PgPool}; use sqlx::{FromRow, PgPool};
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
@@ -55,53 +56,30 @@ impl AudioMetadataService {
self.blob_root.join(prefix).join(format!("{}.blob", hash)) self.blob_root.join(prefix).join(format!("{}.blob", hash))
} }
fn get_duration_secs(file_path: &Path) -> i32 { /// Extract ID3 tag and MP3 duration from a file.
match mp3_duration::from_path(file_path) { ///
Ok(dur) => dur.as_secs_f64().round() as i32, /// All I/O is synchronous (id3 + mp3_duration crates), so this MUST
Err(_) => { /// only be called inside `spawn_blocking`.
if let Ok(tag) = Tag::read_from_path(file_path) { fn extract_metadata_blocking(
tag.duration().unwrap_or(0) as i32
} else {
0
}
}
}
}
pub async fn extract_and_save(
&self,
file_id: &Uuid,
file_path: &Path, file_path: &Path,
) -> Result<(), DomainError> { ) -> Option<AudioMetadataFields> {
info!(
"AudioMetadataService: blob_root={:?}, file_id={}, file_path={:?}, exists={}",
self.blob_root,
file_id,
file_path,
file_path.exists()
);
if !file_path.exists() { if !file_path.exists() {
warn!("File does not exist: {:?}", file_path); warn!("File does not exist: {:?}", file_path);
return Ok(()); return None;
} }
let tag = match Tag::read_from_path(file_path) { let tag = match Tag::read_from_path(file_path) {
Ok(t) => t, Ok(t) => t,
Err(e) => { Err(e) => {
warn!("Failed to read ID3 tag from {:?}: {}", file_path, e); warn!("Failed to read ID3 tag from {:?}: {}", file_path, e);
return Ok(()); return None;
} }
}; };
let title = tag.title().map(|s| s.to_string()); let duration_secs = match mp3_duration::from_path(file_path) {
let artist = tag.artist().map(|s| s.to_string()); Ok(dur) => dur.as_secs_f64().round() as i32,
let album = tag.album().map(|s| s.to_string()); Err(_) => tag.duration().unwrap_or(0) as i32,
let genre = tag.genre().map(|s| s.to_string()); };
let track_number: Option<i32> = tag.track().map(|n| n as i32);
let disc_number: Option<i32> = tag.disc().map(|n| n as i32);
let year: Option<i32> = tag.year();
let duration_secs = Self::get_duration_secs(file_path);
let album_artist = let album_artist =
tag.frames() tag.frames()
@@ -111,9 +89,46 @@ impl AudioMetadataService {
_ => None, _ => None,
}); });
Some(AudioMetadataFields {
title: tag.title().map(|s| s.to_string()),
artist: tag.artist().map(|s| s.to_string()),
album: tag.album().map(|s| s.to_string()),
album_artist,
genre: tag.genre().map(|s| s.to_string()),
track_number: tag.track().map(|n| n as i32),
disc_number: tag.disc().map(|n| n as i32),
year: tag.year(),
duration_secs,
})
}
pub async fn extract_and_save(
&self,
file_id: &Uuid,
file_path: &Path,
) -> Result<(), DomainError> {
info!(
"AudioMetadataService: blob_root={:?}, file_id={}, file_path={:?}",
self.blob_root, file_id, file_path,
);
// ── Sync I/O on the blocking thread pool (never stalls Tokio workers) ──
let path = file_path.to_path_buf();
let metadata = tokio::task::spawn_blocking(move || {
Self::extract_metadata_blocking(&path)
})
.await
.map_err(|e| {
DomainError::internal_error("AudioMetadataService", format!("spawn_blocking join error: {e}"))
})?;
let Some(m) = metadata else {
return Ok(());
};
info!( info!(
"Extracted audio metadata for file {}: title={:?}, artist={:?}, album={:?}, duration={}s", "Extracted audio metadata for file {}: title={:?}, artist={:?}, album={:?}, duration={}s",
file_id, title, artist, album, duration_secs file_id, m.title, m.artist, m.album, m.duration_secs
); );
info!("Saving metadata to database for file_id={}", file_id); info!("Saving metadata to database for file_id={}", file_id);
@@ -139,15 +154,15 @@ impl AudioMetadataService {
"#, "#,
) )
.bind(file_id) .bind(file_id)
.bind(&title) .bind(&m.title)
.bind(&artist) .bind(&m.artist)
.bind(&album) .bind(&m.album)
.bind(&album_artist) .bind(&m.album_artist)
.bind(&genre) .bind(&m.genre)
.bind(track_number) .bind(m.track_number)
.bind(disc_number) .bind(m.disc_number)
.bind(year) .bind(m.year)
.bind(duration_secs) .bind(m.duration_secs)
.bind("MPEG") .bind("MPEG")
.execute(&*self.pool) .execute(&*self.pool)
.await .await
@@ -172,24 +187,27 @@ impl AudioMetadataService {
pub async fn reextract_all_audio_metadata( pub async fn reextract_all_audio_metadata(
&self, &self,
) -> Result<MetadataExtractionResult, DomainError> { ) -> Result<MetadataExtractionResult, DomainError> {
let audio_files = sqlx::query_as::<_, AudioFileRow>( // Stream rows one-by-one instead of fetch_all to keep memory O(1).
let mut stream = sqlx::query_as::<_, AudioFileRow>(
r#" r#"
SELECT id as file_id, blob_hash SELECT id as file_id, blob_hash
FROM storage.files FROM storage.files
WHERE mime_type LIKE 'audio/%' WHERE mime_type LIKE 'audio/%'
"#, "#,
) )
.fetch_all(&*self.pool) .fetch(&*self.pool);
.await
.map_err(|e| DomainError::database_error(format!("Failed to fetch audio files: {}", e)))?;
let total = audio_files.len(); let mut total: usize = 0;
let mut processed = 0; let mut processed: usize = 0;
let mut failed = 0; let mut failed: usize = 0;
info!("Starting metadata extraction for {} audio files", total); info!("Starting streaming metadata extraction for audio files");
for audio_file in audio_files { while let Some(row) = stream.next().await {
total += 1;
let audio_file = row.map_err(|e| {
DomainError::database_error(format!("Failed to fetch audio file row: {}", e))
})?;
let file_path = self.blob_path(&audio_file.blob_hash); let file_path = self.blob_path(&audio_file.blob_hash);
match self.extract_and_save(&audio_file.file_id, &file_path).await { match self.extract_and_save(&audio_file.file_id, &file_path).await {
Ok(()) => processed += 1, Ok(()) => processed += 1,
@@ -216,6 +234,19 @@ impl AudioMetadataService {
} }
} }
/// Extracted audio metadata fields transferred from the blocking thread.
struct AudioMetadataFields {
title: Option<String>,
artist: Option<String>,
album: Option<String>,
album_artist: Option<String>,
genre: Option<String>,
track_number: Option<i32>,
disc_number: Option<i32>,
year: Option<i32>,
duration_secs: i32,
}
#[derive(Debug, serde::Serialize)] #[derive(Debug, serde::Serialize)]
pub struct MetadataExtractionResult { pub struct MetadataExtractionResult {
pub total: usize, pub total: usize,