//! Notification subscriber: orchestrates the album-finishing pipeline. //! //! Background task that subscribes to tora's torrent notification stream, //! catches up on missed FINISHED torrents on startup, and on each FINISHED //! event looks up album metadata, matches downloaded files to tracks, and //! writes metadata back into musicfs. use std::path::Path; use std::time::Duration; use anyhow::Result; use sqlx::PgPool; use uuid::Uuid; use crate::db::{get_pending_monitored_albums, mark_album_completed}; use crate::generated::metadata::metadata_service_client::MetadataServiceClient; use crate::generated::metadata::{ GetAlbumRequest, GetAlbumTracksRequest, Provider, SearchAlbumsRequest, get_album_request::Identifier as AlbumIdentifier, }; use crate::generated::musicfs::client_control_client::ClientControlClient; use crate::generated::musicfs::{ListFilesRequest, UpdateMusicMetadataRequest}; use crate::generated::torrent::notification::Detail as NotificationDetail; use crate::generated::torrent::torrents_client::TorrentsClient; use crate::generated::torrent::{ ListRequest, NotificationKind, NotificationsRequest, State, StatusRequest, }; use crate::matching::{MatcherFactory, TrackMatcher, is_audio_file}; use crate::metadata::MetadataService; use crate::musicfs::MusicfsService; use crate::torrent::TorrentsService; // tonic re-exports tokio_stream — do NOT add tokio-stream to Cargo.toml. use tonic::codegen::tokio_stream::StreamExt; /// Spawn the notification subscriber as a fire-and-forget background task. #[allow(clippy::too_many_arguments)] pub fn spawn( pool: PgPool, torrents: TorrentsService, musicfs: MusicfsService, metadata: MetadataService, ) { tokio::spawn(async move { run(pool, torrents, musicfs, metadata).await; }); tracing::info!("notification subscriber started"); } async fn run( pool: PgPool, torrents: TorrentsService, musicfs: MusicfsService, metadata: MetadataService, ) { // Phase 0: reconcile with tora, discover untracked finished torrents. if let Err(e) = reconcile_with_tora(&pool, &torrents, &metadata).await { tracing::warn!(error = %e, "reconciliation pass failed"); } // Phase 1: catch up on albums that finished while we were down. if let Err(e) = catch_up_missed(&pool, &torrents, &musicfs, &metadata).await { tracing::warn!(error = %e, "catch-up pass failed"); } // Phase 2: live subscription with unconditional 5s reconnect. loop { if let Err(e) = subscribe_loop(&pool, &torrents, &musicfs, &metadata).await { tracing::warn!( error = %e, code = ?e.code(), "notification stream disconnected; reconnecting in 5s" ); tokio::time::sleep(Duration::from_secs(5)).await; } } } /// Parse a torrent directory name into (artist, album_title). /// Handles the common "Artist - Album - Year" format. /// Falls back to (empty, full_name) when the format doesn't match. fn parse_torrent_name(name: &str) -> (String, String) { let parts: Vec<&str> = name.split(" - ").collect(); match parts.as_slice() { [artist, album, ..] => (artist.trim().to_string(), album.trim().to_string()), _ => (String::new(), name.trim().to_string()), } } /// Discover torrents in tora that finished but were never tracked via /// MonitorAlbum (e.g. downloaded before the notification subscriber existed). /// For each, search the metadata service for a matching album and insert it /// into monitored_albums so the catch-up pass processes it. async fn reconcile_with_tora( pool: &PgPool, torrents: &TorrentsService, metadata: &MetadataService, ) -> Result<()> { tracing::info!("reconciliation: starting — querying tora for all torrents"); let list_resp = TorrentsClient::new(torrents.channel()) .list(ListRequest {}) .await? .into_inner(); let tracked_ids: Vec = sqlx::query_scalar("SELECT torrent_id FROM monitored_albums") .fetch_all(pool) .await?; let tracked: std::collections::HashSet = tracked_ids.into_iter().collect(); tracing::info!( total_torrents = list_resp.torrents.len(), already_tracked = tracked.len(), "reconciliation: processing torrents" ); let mut checked = 0u32; let mut inserted = 0u32; let mut skipped_not_finished = 0u32; let mut skipped_tracked = 0u32; let mut no_match = 0u32; let mut errors = 0u32; for torrent in &list_resp.torrents { tracing::debug!( id = %torrent.id, name = %torrent.name, state = torrent.state, downloaded_bytes = torrent.downloaded_bytes, total_bytes = torrent.total_bytes, "reconciliation: checking torrent" ); checked += 1; if torrent.state != State::Finished as i32 { skipped_not_finished += 1; continue; } let torrent_uuid: Uuid = match torrent.id.parse() { Ok(u) => u, Err(_) => { errors += 1; continue; } }; if tracked.contains(&torrent_uuid) { skipped_tracked += 1; tracing::debug!(id = %torrent.id, "reconciliation: already tracked, skipping"); continue; } let (artist, album_query) = parse_torrent_name(&torrent.name); tracing::debug!( artist = %artist, album = %album_query, "reconciliation: searching metadata service" ); let search_resp = match MetadataServiceClient::new(metadata.channel()) .search_albums(SearchAlbumsRequest { query: album_query, artist, limit: 1, offset: 0, provider: Provider::Unspecified as i32, album_types: vec![], }) .await { Ok(resp) => resp.into_inner(), Err(e) => { tracing::warn!( error = %e, name = %torrent.name, "reconcile: metadata search failed" ); errors += 1; continue; } }; tracing::debug!( results = search_resp.albums.len(), total_matches = search_resp.total, "reconciliation: search returned" ); let Some(album) = search_resp.albums.first() else { tracing::warn!( name = %torrent.name, "reconcile: no album found in metadata service" ); no_match += 1; continue; }; match crate::db::insert_monitored_album(pool, &album.id, torrent_uuid, "reconciled").await { Ok(_) => { tracing::info!( name = %torrent.name, album_id = %album.id, "reconcile: tracked previously untracked torrent" ); inserted += 1; } Err(e) => { tracing::warn!( error = %e, name = %torrent.name, "reconcile: failed to insert monitored_album" ); errors += 1; } } } tracing::info!( checked, inserted, skipped_not_finished, skipped_tracked, no_match, errors, "reconciliation: complete" ); Ok(()) } /// Replay any monitored album whose torrent reached FINISHED while the /// subscriber was offline. Per-album failures are logged and skipped. async fn catch_up_missed( pool: &PgPool, torrents: &TorrentsService, musicfs: &MusicfsService, metadata: &MetadataService, ) -> Result<()> { let pending = get_pending_monitored_albums(pool).await?; for (torrent_id, album_id) in pending { let result: Result<()> = async { let status = TorrentsClient::new(torrents.channel()) .status(StatusRequest { id: torrent_id.to_string(), }) .await? .into_inner(); if status.state == State::Finished as i32 { handle_finished(pool, torrents, musicfs, metadata, torrent_id, album_id).await?; } Ok(()) } .await; if let Err(e) = result { tracing::warn!(error = %e, %torrent_id, "catch-up: failed to process album"); } } Ok(()) } /// Subscribe to the live notification stream and dispatch FINISHED events. /// Returns when the stream ends or errors; the caller owns reconnect policy. async fn subscribe_loop( pool: &PgPool, torrents: &TorrentsService, musicfs: &MusicfsService, metadata: &MetadataService, ) -> Result<(), tonic::Status> { let mut client = TorrentsClient::new(torrents.channel()); let response = client .notifications(NotificationsRequest { ids: vec![], kinds: vec![NotificationKind::StateChanged as i32], }) .await?; let stream = response.into_inner(); tokio::pin!(stream); while let Some(item) = stream.next().await { let notification = item?; if notification.kind != NotificationKind::StateChanged as i32 { continue; } let Some(NotificationDetail::StateChanged(changed)) = notification.detail.as_ref() else { continue; }; if changed.current != State::Finished as i32 { continue; } let Some(torrent) = notification.torrent.as_ref() else { continue; }; let torrent_id_str = torrent.id.clone(); // Give tora's move-on-completion a moment to settle. tokio::time::sleep(Duration::from_secs(2)).await; let torrent_uuid: Uuid = match torrent_id_str.parse() { Ok(u) => u, Err(e) => { tracing::warn!( error = %e, id = %torrent_id_str, "notification: unparseable torrent id" ); continue; } }; let album_id: Option = sqlx::query_scalar("SELECT album_id FROM monitored_albums WHERE torrent_id = $1") .bind(torrent_uuid) .fetch_optional(pool) .await .map_err(|e| tonic::Status::internal(format!("db lookup failed: {e}")))?; let Some(album_id) = album_id else { continue; }; if let Err(e) = handle_finished(pool, torrents, musicfs, metadata, torrent_uuid, album_id).await { tracing::warn!(error = %e, %torrent_uuid, "live: failed to process album"); } } Ok(()) } /// Process a single FINISHED album end-to-end: idempotency check, torrent /// status lookup, metadata + tracks fetch, file filtering, matching, and /// metadata writeback. Idempotent — a second call for the same torrent is a /// no-op once `mark_album_completed` has run. #[allow(clippy::too_many_arguments)] async fn handle_finished( pool: &PgPool, torrents: &TorrentsService, musicfs: &MusicfsService, metadata: &MetadataService, torrent_id: Uuid, album_id: String, ) -> Result<()> { let already_completed: bool = sqlx::query_scalar( "SELECT EXISTS(SELECT 1 FROM monitored_albums \ WHERE torrent_id = $1 AND completed_at IS NOT NULL)", ) .bind(torrent_id) .fetch_one(pool) .await?; if already_completed { return Ok(()); } let status = TorrentsClient::new(torrents.channel()) .status(StatusRequest { id: torrent_id.to_string(), }) .await? .into_inner(); let album_resp = MetadataServiceClient::new(metadata.channel()) .get_album(GetAlbumRequest { provider: Provider::Unspecified as i32, identifier: Some(AlbumIdentifier::Id(album_id.clone())), }) .await? .into_inner(); let Some(album) = album_resp.album else { anyhow::bail!("metadata service returned no album for {album_id}"); }; let tracks_resp = MetadataServiceClient::new(metadata.channel()) .get_album_tracks(GetAlbumTracksRequest { album_id, provider: Provider::Unspecified as i32, }) .await? .into_inner(); let files_resp = ClientControlClient::new(musicfs.channel()) .list_files(ListFilesRequest {}) .await? .into_inner(); // output_path is tora's download dir; after move-on-completion the files // live under /completed//, so we filter by the // last path component, not by prefix. let dir_name = Path::new(&status.output_path) .file_name() .and_then(|n| n.to_str()) .unwrap_or(""); let mut filtered: Vec<_> = files_resp .files .iter() .filter(|f| is_audio_file(&f.name) && f.original_path.contains(dir_name)) .cloned() .collect(); filtered.sort_by(|a, b| a.name.cmp(&b.name)); let matcher = MatcherFactory::create(); let result = matcher.match_files_to_tracks(&filtered, &tracks_resp.tracks); let album_artist = album .artists .first() .and_then(|credit| credit.artist.as_ref()) .map(|artist| artist.name.clone()) .unwrap_or_else(|| "Unknown Artist".to_string()); for (file, track) in &result.matched { let needs_update = match file.metadata.as_ref() { None => true, Some(existing) => { existing.track_title != track.title || existing.track_number != track.track_number || existing.album != album.title || !existing.artist.iter().any(|a| a == &album_artist) } }; if !needs_update { continue; } ClientControlClient::new(musicfs.channel()) .update_music_metadata(UpdateMusicMetadataRequest { inode: file.inode, artist: vec![album_artist.clone()], album_artist: Some(album_artist.clone()), album: album.title.clone(), track_number: track.track_number, track_title: track.title.clone(), other_tags: vec![], }) .await?; tracing::debug!( inode = file.inode, name = %file.name, before_track_title = ?file.metadata.as_ref().map(|m| &m.track_title), before_track_number = ?file.metadata.as_ref().map(|m| m.track_number), before_album = ?file.metadata.as_ref().map(|m| &m.album), before_artist = ?file.metadata.as_ref().map(|m| &m.artist), after_track_title = %track.title, after_track_number = track.track_number, after_album = %album.title, after_artist = %album_artist, "metadata updated" ); } for file in &result.unmatched_files { tracing::warn!(inode = file.inode, name = %file.name, "no matching track for file"); } mark_album_completed(pool, torrent_id).await?; tracing::info!( %torrent_id, matched = result.matched.len(), unmatched_files = result.unmatched_files.len(), "album processing complete" ); Ok(()) }