diff --git a/internal/monitor_workflow.go b/internal/monitor_workflow.go index 2fdb0bd..47a6b89 100644 --- a/internal/monitor_workflow.go +++ b/internal/monitor_workflow.go @@ -417,7 +417,7 @@ func (w *monitorWorkflow) run(ctx context.Context) error { } if dbAlbum != nil { w.publisher.SetAlbumID(dbAlbum.ID) - w.service.saveTorrentAndDownload(ctx, dbAlbum.ID, best) + w.service.saveTorrentAndDownload(ctx, dbAlbum.ID, w.workflowRunID, best) } else { log.Warn().Str("album_id", w.req.AlbumId).Msg("album not in DB after persist attempt, skipping torrent/download persistence") } diff --git a/internal/service.go b/internal/service.go index 6871627..055046c 100644 --- a/internal/service.go +++ b/internal/service.go @@ -589,7 +589,7 @@ func (service *MusicAgregatorService) MonitorAlbum(ctx context.Context, req *pb. dbAlbum, _ = service.metadata.GetAlbumByExternalID(ctx, album.GetId()) } if dbAlbum != nil { - service.saveTorrentAndDownload(ctx, dbAlbum.ID, best) + service.saveTorrentAndDownload(ctx, dbAlbum.ID, "", best) } else { log.Warn().Str("album_id", req.GetAlbumId()).Msg("album not in DB after persist attempt, skipping torrent/download persistence") } @@ -840,7 +840,7 @@ func (service *MusicAgregatorService) addToTorrentClient(best parsedItem) error return nil } -func (service *MusicAgregatorService) saveTorrentAndDownload(ctx context.Context, dbAlbumID string, best parsedItem) { +func (service *MusicAgregatorService) saveTorrentAndDownload(ctx context.Context, dbAlbumID string, workflowRunID string, best parsedItem) { quality := normalizeQuality(pb.QualityType_QUALITY_UNSPECIFIED, best.rel.BitDepth, best.rel.SampleRate) dbTorrent := &database.Torrent{ @@ -898,6 +898,9 @@ func (service *MusicAgregatorService) saveTorrentAndDownload(ctx context.Context DownloadID: download.ID, TorrentHash: best.rel.InfoHash, CheckInterval: 30 * time.Second, + AlbumID: dbAlbumID, + Quality: quality, + WorkflowRunID: workflowRunID, }, &river.InsertOpts{ ScheduledAt: time.Now().Add(30 * time.Second), }) diff --git a/internal/workers/poll_download.go b/internal/workers/poll_download.go index d574421..efc5d45 100644 --- a/internal/workers/poll_download.go +++ b/internal/workers/poll_download.go @@ -2,6 +2,7 @@ package workers import ( "context" + "fmt" "time" "github.com/jackc/pgx/v5" @@ -10,6 +11,7 @@ import ( "homelab.lan/music-agregator/internal/analysis" "homelab.lan/music-agregator/internal/database" + "homelab.lan/music-agregator/internal/eventbus" "homelab.lan/music-agregator/internal/metadata" "homelab.lan/music-agregator/internal/musicfs" "homelab.lan/music-agregator/internal/torrent" @@ -19,24 +21,29 @@ type PollDownloadArgs struct { DownloadID string `json:"download_id"` TorrentHash string `json:"torrent_hash"` CheckInterval time.Duration `json:"check_interval"` + AlbumID string `json:"album_id"` + Quality string `json:"quality"` + WorkflowRunID string `json:"workflow_run_id"` } func (PollDownloadArgs) Kind() string { return "poll_download" } type PollDownloadWorker struct { river.WorkerDefaults[PollDownloadArgs] - TorrentClient torrent.TorrentClient - Downloads *database.DownloadRepository - DownloadFiles *database.DownloadFileRepository - AlbumReleases *database.AlbumReleaseRepository - TrackReleases *database.TrackReleaseRepository - RiverClient *river.Client[pgx.Tx] - PathMapper *torrent.PathMapper - Analyzer *analysis.ReleaseAnalyzer - MusicFSClient *musicfs.Client - MusicFSOriginID string + TorrentClient torrent.TorrentClient + Downloads *database.DownloadRepository + DownloadFiles *database.DownloadFileRepository + AlbumReleases *database.AlbumReleaseRepository + TrackReleases *database.TrackReleaseRepository + RiverClient *river.Client[pgx.Tx] + PathMapper *torrent.PathMapper + Analyzer *analysis.ReleaseAnalyzer + MusicFSClient *musicfs.Client + MusicFSOriginID string MusicFSOriginRoot string - MetadataService *metadata.MetadataService + MetadataService *metadata.MetadataService + EventBus *eventbus.EventBus + AlbumEvents *database.AlbumEventRepository } func (w *PollDownloadWorker) Work(ctx context.Context, job *river.Job[PollDownloadArgs]) error { @@ -53,6 +60,7 @@ func (w *PollDownloadWorker) Work(ctx context.Context, job *river.Job[PollDownlo if len(results) == 0 { log.Warn().Str("hash", args.TorrentHash).Msg("torrent not found in client, marking failed") w.Downloads.SetFailed(ctx, args.DownloadID, "torrent not found in client") + w.publishEvent(ctx, args, "error", "", "torrent not found in client") return nil } @@ -65,6 +73,7 @@ func (w *PollDownloadWorker) Work(ctx context.Context, job *river.Job[PollDownlo case t.State == "error": log.Warn().Str("hash", args.TorrentHash).Str("state", t.State).Msg("torrent in error state") w.Downloads.SetFailed(ctx, args.DownloadID, "torrent error state") + w.publishEvent(ctx, args, "error", "", "torrent error state") return nil default: @@ -74,6 +83,8 @@ func (w *PollDownloadWorker) Work(ctx context.Context, job *river.Job[PollDownlo Float64("progress", t.Progress*100). Int64("dlspeed", t.DlSpeed). Msg("download in progress") + w.publishEvent(ctx, args, "status", "", + fmt.Sprintf("Downloading: %.1f%% ↓%s", t.Progress*100, formatSpeed(t.DlSpeed))) return w.reschedule(ctx, args) } } @@ -127,6 +138,8 @@ func (w *PollDownloadWorker) onCompleted(ctx context.Context, args PollDownloadA w.syncAndEnrichMusicFS(ctx, args, contentPath) } + w.publishEvent(ctx, args, "result", "MONITOR_STEP_COMPLETE", "Download completed") + return nil } @@ -207,6 +220,8 @@ func (w *PollDownloadWorker) RecoverOrphanedDownloads(ctx context.Context) { DownloadID: d.ID, TorrentHash: d.QbitHash, CheckInterval: 30 * time.Second, + AlbumID: d.AlbumID, + Quality: d.Quality, }, &river.InsertOpts{ ScheduledAt: time.Now().Add(5 * time.Second), UniqueOpts: river.UniqueOpts{ @@ -220,3 +235,48 @@ func (w *PollDownloadWorker) RecoverOrphanedDownloads(ctx context.Context) { } } } + +func (w *PollDownloadWorker) publishEvent(ctx context.Context, args PollDownloadArgs, eventType, step, message string) { + if w.EventBus == nil { + return + } + + topic := args.AlbumID + ":" + args.Quality + + var seq int64 + if w.AlbumEvents != nil && args.AlbumID != "" { + event := &database.AlbumEvent{ + WorkflowRunID: args.WorkflowRunID, + AlbumID: args.AlbumID, + EventType: eventType, + Step: step, + Message: message, + } + if err := w.AlbumEvents.Create(ctx, event); err != nil { + log.Error().Err(err).Msg("failed to persist download event") + } else { + seq = event.Seq + } + } + + w.EventBus.Publish(topic, &eventbus.Event{ + Seq: seq, + WorkflowRunID: args.WorkflowRunID, + AlbumID: args.AlbumID, + Quality: args.Quality, + EventType: eventType, + Step: step, + Message: message, + }) +} + +func formatSpeed(bytesPerSec int64) string { + switch { + case bytesPerSec >= 1<<20: + return fmt.Sprintf("%.1f MB/s", float64(bytesPerSec)/float64(1<<20)) + case bytesPerSec >= 1<<10: + return fmt.Sprintf("%.0f KB/s", float64(bytesPerSec)/float64(1<<10)) + default: + return fmt.Sprintf("%d B/s", bytesPerSec) + } +}