package service import ( "context" "os" "path/filepath" "strings" "testing" "time" "vidarchive/internal/config" "vidarchive/internal/models" "vidarchive/internal/repository" ) // execEnv is a DownloadService wired to a real database and temp directories, // with yt-dlp and ffprobe replaced by scripts the test controls. type execEnv struct { svc *DownloadService repo *repository.DownloadRepository subRepo *repository.SubscriptionRepository settings *SettingsService cfg *config.Config // scratch holds the marker/signal files the fake tools read and write. scratch string } func newExecEnv(t *testing.T) *execEnv { t.Helper() db := setupTestDB(t) t.Cleanup(func() { db.Close() }) root := t.TempDir() cfg := &config.Config{ LibraryDir: filepath.Join(root, "library"), TempDir: filepath.Join(root, "temp"), YTDLPPath: "/bin/false", FFmpegPath: "ffmpeg", FFprobePath: "ffprobe", } for _, dir := range []string{cfg.LibraryDir, cfg.TempDir} { if err := os.MkdirAll(dir, 0755); err != nil { t.Fatal(err) } } downloadRepo := repository.NewDownloadRepository(db) subRepo := repository.NewSubscriptionRepository(db) settingsSvc := NewSettingsService(repository.NewSettingsRepository(db)) svc := NewDownloadService( downloadRepo, NewLibraryService(cfg.LibraryDir, cfg.FFmpegPath, cfg.FFprobePath), NewPresetService(repository.NewPresetRepository(db)), settingsSvc, NewSubscriptionService(subRepo, cfg), cfg, ) return &execEnv{svc: svc, repo: downloadRepo, subRepo: subRepo, settings: settingsSvc, cfg: cfg, scratch: root} } // fakeYTDLP installs a stand-in for yt-dlp. body is shell run with $DEST set to // the directory yt-dlp was told to write into (its -P argument), $SCRATCH set to // the test's scratch dir, and $RUN set to the invocation count, so a body can // behave differently on the metadata second pass. func (e *execEnv) fakeYTDLP(t *testing.T, body string) { t.Helper() e.cfg.YTDLPPath = writeScript(t, filepath.Join(e.scratch, "yt-dlp"), ` DEST="" prev="" for a in "$@"; do if [ "$prev" = "-P" ]; then DEST="$a"; fi prev="$a" done export DEST export SCRATCH="`+e.scratch+`" RUN=$(( $(cat "$SCRATCH/runs" 2>/dev/null || echo 0) + 1 )) echo "$RUN" > "$SCRATCH/runs" export RUN `+body) } // fakeFFprobe installs a stand-in for ffprobe. It reports a 3 second duration, // which is what the real one would say about the test videos. func (e *execEnv) fakeFFprobe(t *testing.T, body string) { t.Helper() e.cfg.FFprobePath = writeScript(t, filepath.Join(e.scratch, "ffprobe"), ` export SCRATCH="`+e.scratch+`" `+body+` echo 3.0`) } func writeScript(t *testing.T, path, body string) string { t.Helper() if err := os.WriteFile(path, []byte("#!/bin/sh\n"+body+"\n"), 0755); err != nil { t.Fatal(err) } return path } // writeItem is shell that creates one importable item under $DEST. The media // file is a copy of a real video, so the mimetype classification in // importItemDir sees actual video content. func (e *execEnv) writeItem(t *testing.T, dirName, title, videoID string) string { t.Helper() requireFFmpeg(t) src := filepath.Join(e.scratch, "source-"+dirName+".mp4") makeTestVideo(t, src) return ` mkdir -p "$DEST/` + dirName + `" cp "` + src + `" "$DEST/` + dirName + `/clip.mp4" printf '{"id":"` + videoID + `","title":"` + title + `","webpage_url":"https://example.com/` + videoID + `"}' \ > "$DEST/` + dirName + `/clip.info.json" ` } // queue inserts a queued download, the state ExecuteDownload expects. func (e *execEnv) queue(t *testing.T, d *models.Download) *models.Download { t.Helper() d.Status = "queued" if d.URL == "" { d.URL = "https://example.com/watch" } if err := e.repo.Create(d); err != nil { t.Fatalf("create download: %v", err) } return d } func (e *execEnv) status(t *testing.T, id int64) *models.Download { t.Helper() got, err := e.repo.GetByID(id) if err != nil { t.Fatalf("reload download %d: %v", id, err) } return got } // waitForFile blocks until path exists. The fake tools touch a file to say they // have started, which lets a test act at a known point instead of sleeping. func waitForFile(t *testing.T, path string) { t.Helper() deadline := time.Now().Add(10 * time.Second) for time.Now().Before(deadline) { if _, err := os.Stat(path); err == nil { return } time.Sleep(5 * time.Millisecond) } t.Fatalf("timed out waiting for %s", path) } func libraryEntries(t *testing.T, dir string) []string { t.Helper() entries, err := os.ReadDir(dir) if err != nil { t.Fatal(err) } var names []string for _, entry := range entries { names = append(names, entry.Name()) } return names } func TestExecuteDownloadCompletes(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, e.writeItem(t, "item-00001", "My Clip", "abc123")) d := e.queue(t, &models.Download{}) processed, err := e.svc.ExecuteDownload(context.Background(), d) if err != nil { t.Fatalf("ExecuteDownload: %v", err) } if !processed { t.Error("processed = false, want true") } got := e.status(t, d.ID) if got.Status != "completed" { t.Errorf("status = %q, want completed", got.Status) } if !got.CompletedAt.Valid { t.Error("completed_at not set") } if names := libraryEntries(t, e.cfg.LibraryDir); len(names) != 1 || names[0] != "My Clip" { t.Errorf("library = %v, want [My Clip]", names) } if _, err := os.Stat(e.svc.tempDirFor(d.ID)); !os.IsNotExist(err) { t.Error("temp download dir was not removed") } } // A download already claimed by another worker must be left alone: Submit and // the 2s queue checker can both enqueue the same row. func TestExecuteDownloadSkipsAlreadyClaimed(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, `touch "$SCRATCH/ran"`) d := e.queue(t, &models.Download{}) if _, err := e.repo.MarkStarted(d.ID); err != nil { t.Fatal(err) } processed, err := e.svc.ExecuteDownload(context.Background(), d) if err != nil { t.Errorf("error = %v, want nil", err) } if processed { t.Error("processed = true, want false for an already-claimed download") } if _, err := os.Stat(filepath.Join(e.scratch, "ran")); err == nil { t.Error("yt-dlp ran for a download this worker did not claim") } } func TestExecuteDownloadRecordsYTDLPFailure(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, `echo "ERROR: video unavailable" >&2; exit 3`) d := e.queue(t, &models.Download{}) processed, err := e.svc.ExecuteDownload(context.Background(), d) if err == nil { t.Fatal("expected an error from a failing yt-dlp") } if processed { t.Error("processed = true, want false") } got := e.status(t, d.ID) if got.Status != "error" { t.Errorf("status = %q, want error", got.Status) } if !got.ErrorMessage.Valid || got.ErrorMessage.String == "" { t.Error("error_message not recorded") } if !strings.Contains(got.Logs.String, "video unavailable") { t.Errorf("yt-dlp output not persisted to logs: %q", got.Logs.String) } } // yt-dlp can exit 0 having downloaded nothing (every entry filtered out). For a // plain download that is a failure, not a silent success. func TestExecuteDownloadEmptyResultIsFailure(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, `mkdir -p "$DEST/item-00001"; exit 0`) d := e.queue(t, &models.Download{}) if _, err := e.svc.ExecuteDownload(context.Background(), d); err == nil { t.Fatal("expected an error when no media was downloaded") } got := e.status(t, d.ID) if got.Status != "error" { t.Errorf("status = %q, want error", got.Status) } if !strings.Contains(got.ErrorMessage.String, "no media files") { t.Errorf("error_message = %q, want it to mention no media files", got.ErrorMessage.String) } } // yt-dlp exits non-zero when a single playlist entry fails. The entries that did // download must still be imported, and the run must read as completed with a // note in the log. func TestExecuteDownloadPartialPlaylistFailureCompletes(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, e.writeItem(t, "item-00001", "My Clip", "abc123")+` echo "ERROR: [youtube] bad2: Video unavailable" >&2 exit 1`) d := e.queue(t, &models.Download{}) processed, err := e.svc.ExecuteDownload(context.Background(), d) if err != nil { t.Fatalf("ExecuteDownload: %v", err) } if !processed { t.Error("processed = false, want true") } got := e.status(t, d.ID) if got.Status != "completed" { t.Errorf("status = %q, want completed", got.Status) } if names := libraryEntries(t, e.cfg.LibraryDir); len(names) != 1 || names[0] != "My Clip" { t.Errorf("library = %v, want [My Clip], the entry that did download", names) } if !strings.Contains(got.Logs.String, "1 item(s) were imported anyway") { t.Errorf("logs do not report the partial failure: %q", got.Logs.String) } } // The same failure with nothing downloaded is a plain error. func TestExecuteDownloadTotalFailureIsError(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, `echo "ERROR: [youtube] bad: Video unavailable" >&2; exit 1`) d := e.queue(t, &models.Download{}) if _, err := e.svc.ExecuteDownload(context.Background(), d); err == nil { t.Fatal("expected an error when nothing was downloaded") } if got := e.status(t, d.ID); got.Status != "error" { t.Errorf("status = %q, want error", got.Status) } } // A reserved flag must fail before yt-dlp is ever started, and a preset's flags // are checked as well as the download's own. func TestExecuteDownloadRejectsReservedPresetFlag(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, `touch "$SCRATCH/ran"`) preset := &models.Preset{Name: "Bad", CustomFlags: "-o /tmp/anywhere.mp4"} if err := e.svc.presetSvc.Create(preset); err != nil { t.Fatalf("create preset: %v", err) } d := e.queue(t, &models.Download{PresetID: sqlNullInt64(preset.ID)}) if _, err := e.svc.ExecuteDownload(context.Background(), d); err == nil { t.Fatal("expected the reserved flag to be rejected") } if got := e.status(t, d.ID); got.Status != "error" { t.Errorf("status = %q, want error", got.Status) } if _, err := os.Stat(filepath.Join(e.scratch, "ran")); err == nil { t.Error("yt-dlp ran despite a reserved flag") } } // A user-initiated cancel while yt-dlp runs is terminal: the row is recorded // "cancelled" and nothing reaches the library. func TestExecuteDownloadCancelDuringRun(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, `touch "$SCRATCH/started"; sleep 60`) d := e.queue(t, &models.Download{}) done := make(chan error, 1) go func() { _, err := e.svc.ExecuteDownload(context.Background(), d) done <- err }() waitForFile(t, filepath.Join(e.scratch, "started")) e.svc.cancelDownload(d.ID) select { case err := <-done: if err != ErrCancelled { t.Errorf("error = %v, want ErrCancelled", err) } case <-time.After(20 * time.Second): t.Fatal("ExecuteDownload did not return after cancel") } if got := e.status(t, d.ID); got.Status != "cancelled" { t.Errorf("status = %q, want cancelled", got.Status) } if names := libraryEntries(t, e.cfg.LibraryDir); len(names) != 0 { t.Errorf("library = %v, want empty", names) } } // A shutdown is not a user cancel: the row must stay "downloading" so that // ResetStalledDownloads re-queues it on the next start instead of losing it. func TestExecuteDownloadShutdownLeavesRowResumable(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, `touch "$SCRATCH/started"; sleep 60`) d := e.queue(t, &models.Download{}) // The parent context stands in for the worker pool's, which Stop cancels. parent, stop := context.WithCancel(context.Background()) defer stop() done := make(chan error, 1) go func() { _, err := e.svc.ExecuteDownload(parent, d) done <- err }() waitForFile(t, filepath.Join(e.scratch, "started")) stop() select { case err := <-done: if err != ErrCancelled { t.Errorf("error = %v, want ErrCancelled", err) } case <-time.After(20 * time.Second): t.Fatal("ExecuteDownload did not return after shutdown") } if got := e.status(t, d.ID); got.Status != "downloading" { t.Fatalf("status = %q, want downloading so the restart can resume it", got.Status) } if err := e.svc.ResetStalledDownloads(); err != nil { t.Fatalf("ResetStalledDownloads: %v", err) } if got := e.status(t, d.ID); got.Status != "queued" { t.Errorf("status after restart = %q, want queued", got.Status) } } // A cancel that lands after yt-dlp exited 0 must not be recorded as a completed // download. The cancel here arrives during the import, so the yt-dlp error is // nil — the case that used to fall through to MarkCompleted. func TestExecuteDownloadCancelAfterSuccessIsNotCompleted(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, e.writeItem(t, "item-00001", "My Clip", "abc123")) // ffprobe runs during the import, after yt-dlp has already succeeded. Block // there so the cancel lands at exactly that point. e.fakeFFprobe(t, ` touch "$SCRATCH/probing" while [ ! -f "$SCRATCH/proceed" ]; do sleep 0.05; done`) d := e.queue(t, &models.Download{}) done := make(chan error, 1) go func() { _, err := e.svc.ExecuteDownload(context.Background(), d) done <- err }() waitForFile(t, filepath.Join(e.scratch, "probing")) e.svc.cancelDownload(d.ID) if err := os.WriteFile(filepath.Join(e.scratch, "proceed"), nil, 0644); err != nil { t.Fatal(err) } select { case err := <-done: if err != ErrCancelled { t.Errorf("error = %v, want ErrCancelled", err) } case <-time.After(20 * time.Second): t.Fatal("ExecuteDownload did not return after cancel") } if got := e.status(t, d.ID); got.Status != "cancelled" { t.Errorf("status = %q, want cancelled (a cancelled run must not read as completed)", got.Status) } } // The same rule for a metadata-mode subscription refresh, whose second pass // returns the yt-dlp error rather than the import error — so a cancel there also // leaves a nil error behind. func TestExecuteDownloadMetadataCancelIsNotCompleted(t *testing.T) { e := newExecEnv(t) // Pass 1 (--skip-download) writes metadata for an entry the library does not // have; pass 2 downloads it as a full item. e.fakeYTDLP(t, ` if [ "$RUN" = "1" ]; then mkdir -p "$DEST/item-00001" printf '{"id":"new1","title":"Fresh","webpage_url":"https://example.com/new1"}' > "$DEST/item-00001/clip.info.json" else `+e.writeItem(t, "item-00001", "Fresh", "new1")+` fi`) e.fakeFFprobe(t, ` touch "$SCRATCH/probing" while [ ! -f "$SCRATCH/proceed" ]; do sleep 0.05; done`) sub := &models.Subscription{ Name: "Channel", URL: "https://example.com/channel", Enabled: true, RefreshMode: "metadata", ScheduleKind: "daily", } if err := e.subRepo.Create(sub); err != nil { t.Fatalf("create subscription: %v", err) } d := e.queue(t, &models.Download{SubscriptionID: sqlNullInt64(sub.ID)}) done := make(chan error, 1) go func() { _, err := e.svc.ExecuteDownload(context.Background(), d) done <- err }() waitForFile(t, filepath.Join(e.scratch, "probing")) e.svc.cancelDownload(d.ID) if err := os.WriteFile(filepath.Join(e.scratch, "proceed"), nil, 0644); err != nil { t.Fatal(err) } select { case err := <-done: if err != ErrCancelled { t.Errorf("error = %v, want ErrCancelled", err) } case <-time.After(20 * time.Second): t.Fatal("ExecuteDownload did not return after cancel") } if got := e.status(t, d.ID); got.Status != "cancelled" { t.Errorf("status = %q, want cancelled", got.Status) } } // A restart must clear the temp dirs of downloads it re-queues: leaving them // behind makes the re-run import into a fresh uniqueDir and the library ends up // holding the same item twice. func TestResetStalledDownloadsClearsTempDirs(t *testing.T) { e := newExecEnv(t) d := e.queue(t, &models.Download{}) if _, err := e.repo.MarkStarted(d.ID); err != nil { t.Fatal(err) } dirs := e.svc.tempDirsFor(d.ID) if len(dirs) != 2 { t.Fatalf("tempDirsFor returned %d dirs, want the main and second-pass dirs", len(dirs)) } for _, dir := range dirs { if err := os.MkdirAll(dir, 0755); err != nil { t.Fatal(err) } if err := os.WriteFile(filepath.Join(dir, "partial.mp4.part"), []byte("x"), 0644); err != nil { t.Fatal(err) } } // A download that is not stalled must keep whatever it owns. other := e.queue(t, &models.Download{}) keep := e.svc.tempDirFor(other.ID) if err := os.MkdirAll(keep, 0755); err != nil { t.Fatal(err) } if err := e.svc.ResetStalledDownloads(); err != nil { t.Fatalf("ResetStalledDownloads: %v", err) } if got := e.status(t, d.ID); got.Status != "queued" { t.Errorf("status = %q, want queued", got.Status) } for _, dir := range dirs { if _, err := os.Stat(dir); !os.IsNotExist(err) { t.Errorf("stale temp dir %s was not removed", dir) } } if _, err := os.Stat(keep); err != nil { t.Errorf("temp dir of a queued download was removed: %v", err) } } // Deleting a download from the queue stops the work it started, so yt-dlp does // not keep running for a row that no longer exists. func TestDeleteCancelsRunningDownload(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, `touch "$SCRATCH/started"; sleep 60`) d := e.queue(t, &models.Download{}) done := make(chan error, 1) go func() { _, err := e.svc.ExecuteDownload(context.Background(), d) done <- err }() waitForFile(t, filepath.Join(e.scratch, "started")) if err := e.svc.Delete(d.ID); err != nil { t.Fatalf("Delete: %v", err) } select { case err := <-done: if err != ErrCancelled { t.Errorf("error = %v, want ErrCancelled", err) } case <-time.After(20 * time.Second): t.Fatal("deleting the download did not stop it") } if _, err := e.repo.GetByID(d.ID); err == nil { t.Error("download row still exists after Delete") } if names := libraryEntries(t, e.cfg.LibraryDir); len(names) != 0 { t.Errorf("library = %v, want empty", names) } } // Clearing the queue must also stop what is running, for the same reason. func TestDeleteAllCancelsRunningDownload(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, `touch "$SCRATCH/started"; sleep 60`) d := e.queue(t, &models.Download{}) done := make(chan error, 1) go func() { _, err := e.svc.ExecuteDownload(context.Background(), d) done <- err }() waitForFile(t, filepath.Join(e.scratch, "started")) if err := e.svc.DeleteAll(); err != nil { t.Fatalf("DeleteAll: %v", err) } select { case err := <-done: if err != ErrCancelled { t.Errorf("error = %v, want ErrCancelled", err) } case <-time.After(20 * time.Second): t.Fatal("clearing the queue did not stop the running download") } } // yt-dlp rewrites the cookie file on exit with rotated session cookies. Those // must be saved back, or the stored snapshot goes stale and the session dies. func TestExecuteDownloadSavesRefreshedCookies(t *testing.T) { e := newExecEnv(t) if err := e.settings.SetCookies("# Netscape HTTP Cookie File\nold-session\n"); err != nil { t.Fatalf("SetCookies: %v", err) } e.fakeYTDLP(t, ` cookies="" prev="" for a in "$@"; do if [ "$prev" = "--cookies" ]; then cookies="$a"; fi prev="$a" done [ -n "$cookies" ] || { echo "no --cookies argument" >&2; exit 1; } printf '# Netscape HTTP Cookie File\nnew-session\n' > "$cookies" `+e.writeItem(t, "item-00001", "My Clip", "abc123")) d := e.queue(t, &models.Download{}) if _, err := e.svc.ExecuteDownload(context.Background(), d); err != nil { t.Fatalf("ExecuteDownload: %v", err) } got, err := e.settings.GetCookies() if err != nil { t.Fatalf("GetCookies: %v", err) } if want := "# Netscape HTTP Cookie File\nnew-session\n"; got != want { t.Errorf("cookies = %q, want %q", got, want) } } // An untouched cookie file must leave the stored cookies exactly as they were. func TestExecuteDownloadKeepsUntouchedCookies(t *testing.T) { e := newExecEnv(t) const stored = "# Netscape HTTP Cookie File\nold-session\n" if err := e.settings.SetCookies(stored); err != nil { t.Fatalf("SetCookies: %v", err) } e.fakeYTDLP(t, e.writeItem(t, "item-00001", "My Clip", "abc123")) d := e.queue(t, &models.Download{}) if _, err := e.svc.ExecuteDownload(context.Background(), d); err != nil { t.Fatalf("ExecuteDownload: %v", err) } got, err := e.settings.GetCookies() if err != nil { t.Fatalf("GetCookies: %v", err) } if got != stored { t.Errorf("cookies = %q, want them unchanged (%q)", got, stored) } } // A download whose scratch directory can't be created must be finalized as an // error, not left in "downloading" — which also keeps its subscription busy. func TestExecuteDownloadTempDirFailureIsError(t *testing.T) { e := newExecEnv(t) e.fakeYTDLP(t, `exit 0`) // A regular file where the temp tree needs a directory makes MkdirAll fail // with ENOTDIR. blocker := filepath.Join(e.scratch, "blocker") if err := os.WriteFile(blocker, nil, 0644); err != nil { t.Fatal(err) } e.cfg.TempDir = filepath.Join(blocker, "temp") sub := &models.Subscription{Name: "s", URL: "https://example.com", OutputDir: "feed"} if err := e.subRepo.Create(sub); err != nil { t.Fatal(err) } d := e.queue(t, &models.Download{SubscriptionID: sqlNullInt64(sub.ID)}) if _, err := e.svc.ExecuteDownload(context.Background(), d); err == nil { t.Fatal("ExecuteDownload succeeded, want an error") } got := e.status(t, d.ID) if got.Status != "error" { t.Errorf("status = %q, want error", got.Status) } if !got.ErrorMessage.Valid || !strings.Contains(got.ErrorMessage.String, "temp download dir") { t.Errorf("error_message = %q, want it to name the temp dir", got.ErrorMessage.String) } active, err := e.svc.HasActiveForSubscription(sub.ID) if err != nil { t.Fatal(err) } if active { t.Error("subscription still reported as having an active run") } }