execute_download_test.go
⎇
Raw
1package service
2
3import (
4 "context"
5 "os"
6 "path/filepath"
7 "strings"
8 "testing"
9 "time"
10
11 "vidarchive/internal/config"
12 "vidarchive/internal/models"
13 "vidarchive/internal/repository"
14)
15
16// execEnv is a DownloadService wired to a real database and temp directories,
17// with yt-dlp and ffprobe replaced by scripts the test controls.
18type execEnv struct {
19 svc *DownloadService
20 repo *repository.DownloadRepository
21 subRepo *repository.SubscriptionRepository
22 cfg *config.Config
23 // scratch holds the marker/signal files the fake tools read and write.
24 scratch string
25}
26
27func newExecEnv(t *testing.T) *execEnv {
28 t.Helper()
29
30 db := setupTestDB(t)
31 t.Cleanup(func() { db.Close() })
32
33 root := t.TempDir()
34 cfg := &config.Config{
35 LibraryDir: filepath.Join(root, "library"),
36 TempDir: filepath.Join(root, "temp"),
37 YTDLPPath: "/bin/false",
38 FFmpegPath: "ffmpeg",
39 FFprobePath: "ffprobe",
40 }
41 for _, dir := range []string{cfg.LibraryDir, cfg.TempDir} {
42 if err := os.MkdirAll(dir, 0755); err != nil {
43 t.Fatal(err)
44 }
45 }
46
47 downloadRepo := repository.NewDownloadRepository(db)
48 subRepo := repository.NewSubscriptionRepository(db)
49 svc := NewDownloadService(
50 downloadRepo,
51 NewLibraryService(cfg.LibraryDir, cfg.FFmpegPath, cfg.FFprobePath),
52 NewPresetService(repository.NewPresetRepository(db)),
53 NewSettingsService(repository.NewSettingsRepository(db)),
54 NewSubscriptionService(subRepo, cfg),
55 cfg,
56 )
57
58 return &execEnv{svc: svc, repo: downloadRepo, subRepo: subRepo, cfg: cfg, scratch: root}
59}
60
61// fakeYTDLP installs a stand-in for yt-dlp. body is shell run with $DEST set to
62// the directory yt-dlp was told to write into (its -P argument), $SCRATCH set to
63// the test's scratch dir, and $RUN set to the invocation count, so a body can
64// behave differently on the metadata second pass.
65func (e *execEnv) fakeYTDLP(t *testing.T, body string) {
66 t.Helper()
67 e.cfg.YTDLPPath = writeScript(t, filepath.Join(e.scratch, "yt-dlp"), `
68DEST=""
69prev=""
70for a in "$@"; do
71 if [ "$prev" = "-P" ]; then DEST="$a"; fi
72 prev="$a"
73done
74export DEST
75export SCRATCH="`+e.scratch+`"
76RUN=$(( $(cat "$SCRATCH/runs" 2>/dev/null || echo 0) + 1 ))
77echo "$RUN" > "$SCRATCH/runs"
78export RUN
79`+body)
80}
81
82// fakeFFprobe installs a stand-in for ffprobe. It reports a 3 second duration,
83// which is what the real one would say about the test videos.
84func (e *execEnv) fakeFFprobe(t *testing.T, body string) {
85 t.Helper()
86 e.cfg.FFprobePath = writeScript(t, filepath.Join(e.scratch, "ffprobe"), `
87export SCRATCH="`+e.scratch+`"
88`+body+`
89echo 3.0`)
90}
91
92func writeScript(t *testing.T, path, body string) string {
93 t.Helper()
94 if err := os.WriteFile(path, []byte("#!/bin/sh\n"+body+"\n"), 0755); err != nil {
95 t.Fatal(err)
96 }
97 return path
98}
99
100// writeItem is shell that creates one importable item under $DEST. The media
101// file is a copy of a real video, so the mimetype classification in
102// importItemDir sees actual video content.
103func (e *execEnv) writeItem(t *testing.T, dirName, title, videoID string) string {
104 t.Helper()
105 requireFFmpeg(t)
106 src := filepath.Join(e.scratch, "source-"+dirName+".mp4")
107 makeTestVideo(t, src)
108 return `
109mkdir -p "$DEST/` + dirName + `"
110cp "` + src + `" "$DEST/` + dirName + `/clip.mp4"
111printf '{"id":"` + videoID + `","title":"` + title + `","webpage_url":"https://example.com/` + videoID + `"}' \
112 > "$DEST/` + dirName + `/clip.info.json"
113`
114}
115
116// queue inserts a queued download, the state ExecuteDownload expects.
117func (e *execEnv) queue(t *testing.T, d *models.Download) *models.Download {
118 t.Helper()
119 d.Status = "queued"
120 if d.URL == "" {
121 d.URL = "https://example.com/watch"
122 }
123 if err := e.repo.Create(d); err != nil {
124 t.Fatalf("create download: %v", err)
125 }
126 return d
127}
128
129func (e *execEnv) status(t *testing.T, id int64) *models.Download {
130 t.Helper()
131 got, err := e.repo.GetByID(id)
132 if err != nil {
133 t.Fatalf("reload download %d: %v", id, err)
134 }
135 return got
136}
137
138// waitForFile blocks until path exists. The fake tools touch a file to say they
139// have started, which lets a test act at a known point instead of sleeping.
140func waitForFile(t *testing.T, path string) {
141 t.Helper()
142 deadline := time.Now().Add(10 * time.Second)
143 for time.Now().Before(deadline) {
144 if _, err := os.Stat(path); err == nil {
145 return
146 }
147 time.Sleep(5 * time.Millisecond)
148 }
149 t.Fatalf("timed out waiting for %s", path)
150}
151
152func libraryEntries(t *testing.T, dir string) []string {
153 t.Helper()
154 entries, err := os.ReadDir(dir)
155 if err != nil {
156 t.Fatal(err)
157 }
158 var names []string
159 for _, entry := range entries {
160 names = append(names, entry.Name())
161 }
162 return names
163}
164
165func TestExecuteDownloadCompletes(t *testing.T) {
166 e := newExecEnv(t)
167 e.fakeYTDLP(t, e.writeItem(t, "item-00001", "My Clip", "abc123"))
168
169 d := e.queue(t, &models.Download{})
170 processed, err := e.svc.ExecuteDownload(context.Background(), d)
171 if err != nil {
172 t.Fatalf("ExecuteDownload: %v", err)
173 }
174 if !processed {
175 t.Error("processed = false, want true")
176 }
177
178 got := e.status(t, d.ID)
179 if got.Status != "completed" {
180 t.Errorf("status = %q, want completed", got.Status)
181 }
182 if !got.CompletedAt.Valid {
183 t.Error("completed_at not set")
184 }
185 if names := libraryEntries(t, e.cfg.LibraryDir); len(names) != 1 || names[0] != "My Clip" {
186 t.Errorf("library = %v, want [My Clip]", names)
187 }
188 if _, err := os.Stat(e.svc.tempDirFor(d.ID)); !os.IsNotExist(err) {
189 t.Error("temp download dir was not removed")
190 }
191}
192
193// A download already claimed by another worker must be left alone: Submit and
194// the 2s queue checker can both enqueue the same row.
195func TestExecuteDownloadSkipsAlreadyClaimed(t *testing.T) {
196 e := newExecEnv(t)
197 e.fakeYTDLP(t, `touch "$SCRATCH/ran"`)
198
199 d := e.queue(t, &models.Download{})
200 if _, err := e.repo.MarkStarted(d.ID); err != nil {
201 t.Fatal(err)
202 }
203
204 processed, err := e.svc.ExecuteDownload(context.Background(), d)
205 if err != nil {
206 t.Errorf("error = %v, want nil", err)
207 }
208 if processed {
209 t.Error("processed = true, want false for an already-claimed download")
210 }
211 if _, err := os.Stat(filepath.Join(e.scratch, "ran")); err == nil {
212 t.Error("yt-dlp ran for a download this worker did not claim")
213 }
214}
215
216func TestExecuteDownloadRecordsYTDLPFailure(t *testing.T) {
217 e := newExecEnv(t)
218 e.fakeYTDLP(t, `echo "ERROR: video unavailable" >&2; exit 3`)
219
220 d := e.queue(t, &models.Download{})
221 processed, err := e.svc.ExecuteDownload(context.Background(), d)
222 if err == nil {
223 t.Fatal("expected an error from a failing yt-dlp")
224 }
225 if processed {
226 t.Error("processed = true, want false")
227 }
228
229 got := e.status(t, d.ID)
230 if got.Status != "error" {
231 t.Errorf("status = %q, want error", got.Status)
232 }
233 if !got.ErrorMessage.Valid || got.ErrorMessage.String == "" {
234 t.Error("error_message not recorded")
235 }
236 if !strings.Contains(got.Logs.String, "video unavailable") {
237 t.Errorf("yt-dlp output not persisted to logs: %q", got.Logs.String)
238 }
239}
240
241// yt-dlp can exit 0 having downloaded nothing (every entry filtered out). For a
242// plain download that is a failure, not a silent success.
243func TestExecuteDownloadEmptyResultIsFailure(t *testing.T) {
244 e := newExecEnv(t)
245 e.fakeYTDLP(t, `mkdir -p "$DEST/item-00001"; exit 0`)
246
247 d := e.queue(t, &models.Download{})
248 if _, err := e.svc.ExecuteDownload(context.Background(), d); err == nil {
249 t.Fatal("expected an error when no media was downloaded")
250 }
251
252 got := e.status(t, d.ID)
253 if got.Status != "error" {
254 t.Errorf("status = %q, want error", got.Status)
255 }
256 if !strings.Contains(got.ErrorMessage.String, "no media files") {
257 t.Errorf("error_message = %q, want it to mention no media files", got.ErrorMessage.String)
258 }
259}
260
261// A reserved flag must fail before yt-dlp is ever started, and a preset's flags
262// are checked as well as the download's own.
263func TestExecuteDownloadRejectsReservedPresetFlag(t *testing.T) {
264 e := newExecEnv(t)
265 e.fakeYTDLP(t, `touch "$SCRATCH/ran"`)
266
267 preset := &models.Preset{Name: "Bad", CustomFlags: "-o /tmp/anywhere.mp4"}
268 if err := e.svc.presetSvc.Create(preset); err != nil {
269 t.Fatalf("create preset: %v", err)
270 }
271
272 d := e.queue(t, &models.Download{PresetID: sqlNullInt64(preset.ID)})
273 if _, err := e.svc.ExecuteDownload(context.Background(), d); err == nil {
274 t.Fatal("expected the reserved flag to be rejected")
275 }
276
277 if got := e.status(t, d.ID); got.Status != "error" {
278 t.Errorf("status = %q, want error", got.Status)
279 }
280 if _, err := os.Stat(filepath.Join(e.scratch, "ran")); err == nil {
281 t.Error("yt-dlp ran despite a reserved flag")
282 }
283}
284
285// A user-initiated cancel while yt-dlp runs is terminal: the row is recorded
286// "cancelled" and nothing reaches the library.
287func TestExecuteDownloadCancelDuringRun(t *testing.T) {
288 e := newExecEnv(t)
289 e.fakeYTDLP(t, `touch "$SCRATCH/started"; sleep 60`)
290
291 d := e.queue(t, &models.Download{})
292
293 done := make(chan error, 1)
294 go func() {
295 _, err := e.svc.ExecuteDownload(context.Background(), d)
296 done <- err
297 }()
298
299 waitForFile(t, filepath.Join(e.scratch, "started"))
300 e.svc.cancelDownload(d.ID)
301
302 select {
303 case err := <-done:
304 if err != ErrCancelled {
305 t.Errorf("error = %v, want ErrCancelled", err)
306 }
307 case <-time.After(20 * time.Second):
308 t.Fatal("ExecuteDownload did not return after cancel")
309 }
310
311 if got := e.status(t, d.ID); got.Status != "cancelled" {
312 t.Errorf("status = %q, want cancelled", got.Status)
313 }
314 if names := libraryEntries(t, e.cfg.LibraryDir); len(names) != 0 {
315 t.Errorf("library = %v, want empty", names)
316 }
317}
318
319// A shutdown is not a user cancel: the row must stay "downloading" so that
320// ResetStalledDownloads re-queues it on the next start instead of losing it.
321func TestExecuteDownloadShutdownLeavesRowResumable(t *testing.T) {
322 e := newExecEnv(t)
323 e.fakeYTDLP(t, `touch "$SCRATCH/started"; sleep 60`)
324
325 d := e.queue(t, &models.Download{})
326
327 // The parent context stands in for the worker pool's, which Stop cancels.
328 parent, stop := context.WithCancel(context.Background())
329 defer stop()
330
331 done := make(chan error, 1)
332 go func() {
333 _, err := e.svc.ExecuteDownload(parent, d)
334 done <- err
335 }()
336
337 waitForFile(t, filepath.Join(e.scratch, "started"))
338 stop()
339
340 select {
341 case err := <-done:
342 if err != ErrCancelled {
343 t.Errorf("error = %v, want ErrCancelled", err)
344 }
345 case <-time.After(20 * time.Second):
346 t.Fatal("ExecuteDownload did not return after shutdown")
347 }
348
349 if got := e.status(t, d.ID); got.Status != "downloading" {
350 t.Fatalf("status = %q, want downloading so the restart can resume it", got.Status)
351 }
352 if err := e.svc.ResetStalledDownloads(); err != nil {
353 t.Fatalf("ResetStalledDownloads: %v", err)
354 }
355 if got := e.status(t, d.ID); got.Status != "queued" {
356 t.Errorf("status after restart = %q, want queued", got.Status)
357 }
358}
359
360// A cancel that lands after yt-dlp exited 0 must not be recorded as a completed
361// download. The cancel here arrives during the import, so the yt-dlp error is
362// nil — the case that used to fall through to MarkCompleted.
363func TestExecuteDownloadCancelAfterSuccessIsNotCompleted(t *testing.T) {
364 e := newExecEnv(t)
365 e.fakeYTDLP(t, e.writeItem(t, "item-00001", "My Clip", "abc123"))
366 // ffprobe runs during the import, after yt-dlp has already succeeded. Block
367 // there so the cancel lands at exactly that point.
368 e.fakeFFprobe(t, `
369touch "$SCRATCH/probing"
370while [ ! -f "$SCRATCH/proceed" ]; do sleep 0.05; done`)
371
372 d := e.queue(t, &models.Download{})
373
374 done := make(chan error, 1)
375 go func() {
376 _, err := e.svc.ExecuteDownload(context.Background(), d)
377 done <- err
378 }()
379
380 waitForFile(t, filepath.Join(e.scratch, "probing"))
381 e.svc.cancelDownload(d.ID)
382 if err := os.WriteFile(filepath.Join(e.scratch, "proceed"), nil, 0644); err != nil {
383 t.Fatal(err)
384 }
385
386 select {
387 case err := <-done:
388 if err != ErrCancelled {
389 t.Errorf("error = %v, want ErrCancelled", err)
390 }
391 case <-time.After(20 * time.Second):
392 t.Fatal("ExecuteDownload did not return after cancel")
393 }
394
395 if got := e.status(t, d.ID); got.Status != "cancelled" {
396 t.Errorf("status = %q, want cancelled (a cancelled run must not read as completed)", got.Status)
397 }
398}
399
400// The same rule for a metadata-mode subscription refresh, whose second pass
401// returns the yt-dlp error rather than the import error — so a cancel there also
402// leaves a nil error behind.
403func TestExecuteDownloadMetadataCancelIsNotCompleted(t *testing.T) {
404 e := newExecEnv(t)
405 // Pass 1 (--skip-download) writes metadata for an entry the library does not
406 // have; pass 2 downloads it as a full item.
407 e.fakeYTDLP(t, `
408if [ "$RUN" = "1" ]; then
409 mkdir -p "$DEST/item-00001"
410 printf '{"id":"new1","title":"Fresh","webpage_url":"https://example.com/new1"}' > "$DEST/item-00001/clip.info.json"
411else
412`+e.writeItem(t, "item-00001", "Fresh", "new1")+`
413fi`)
414 e.fakeFFprobe(t, `
415touch "$SCRATCH/probing"
416while [ ! -f "$SCRATCH/proceed" ]; do sleep 0.05; done`)
417
418 sub := &models.Subscription{
419 Name: "Channel",
420 URL: "https://example.com/channel",
421 Enabled: true,
422 RefreshMode: "metadata",
423 ScheduleKind: "daily",
424 }
425 if err := e.subRepo.Create(sub); err != nil {
426 t.Fatalf("create subscription: %v", err)
427 }
428
429 d := e.queue(t, &models.Download{SubscriptionID: sqlNullInt64(sub.ID)})
430
431 done := make(chan error, 1)
432 go func() {
433 _, err := e.svc.ExecuteDownload(context.Background(), d)
434 done <- err
435 }()
436
437 waitForFile(t, filepath.Join(e.scratch, "probing"))
438 e.svc.cancelDownload(d.ID)
439 if err := os.WriteFile(filepath.Join(e.scratch, "proceed"), nil, 0644); err != nil {
440 t.Fatal(err)
441 }
442
443 select {
444 case err := <-done:
445 if err != ErrCancelled {
446 t.Errorf("error = %v, want ErrCancelled", err)
447 }
448 case <-time.After(20 * time.Second):
449 t.Fatal("ExecuteDownload did not return after cancel")
450 }
451
452 if got := e.status(t, d.ID); got.Status != "cancelled" {
453 t.Errorf("status = %q, want cancelled", got.Status)
454 }
455}
456
457// A restart must clear the temp dirs of downloads it re-queues: leaving them
458// behind makes the re-run import into a fresh uniqueDir and the library ends up
459// holding the same item twice.
460func TestResetStalledDownloadsClearsTempDirs(t *testing.T) {
461 e := newExecEnv(t)
462
463 d := e.queue(t, &models.Download{})
464 if _, err := e.repo.MarkStarted(d.ID); err != nil {
465 t.Fatal(err)
466 }
467
468 dirs := e.svc.tempDirsFor(d.ID)
469 if len(dirs) != 2 {
470 t.Fatalf("tempDirsFor returned %d dirs, want the main and second-pass dirs", len(dirs))
471 }
472 for _, dir := range dirs {
473 if err := os.MkdirAll(dir, 0755); err != nil {
474 t.Fatal(err)
475 }
476 if err := os.WriteFile(filepath.Join(dir, "partial.mp4.part"), []byte("x"), 0644); err != nil {
477 t.Fatal(err)
478 }
479 }
480
481 // A download that is not stalled must keep whatever it owns.
482 other := e.queue(t, &models.Download{})
483 keep := e.svc.tempDirFor(other.ID)
484 if err := os.MkdirAll(keep, 0755); err != nil {
485 t.Fatal(err)
486 }
487
488 if err := e.svc.ResetStalledDownloads(); err != nil {
489 t.Fatalf("ResetStalledDownloads: %v", err)
490 }
491
492 if got := e.status(t, d.ID); got.Status != "queued" {
493 t.Errorf("status = %q, want queued", got.Status)
494 }
495 for _, dir := range dirs {
496 if _, err := os.Stat(dir); !os.IsNotExist(err) {
497 t.Errorf("stale temp dir %s was not removed", dir)
498 }
499 }
500 if _, err := os.Stat(keep); err != nil {
501 t.Errorf("temp dir of a queued download was removed: %v", err)
502 }
503}
504
505// Deleting a download from the queue stops the work it started, so yt-dlp does
506// not keep running for a row that no longer exists.
507func TestDeleteCancelsRunningDownload(t *testing.T) {
508 e := newExecEnv(t)
509 e.fakeYTDLP(t, `touch "$SCRATCH/started"; sleep 60`)
510
511 d := e.queue(t, &models.Download{})
512
513 done := make(chan error, 1)
514 go func() {
515 _, err := e.svc.ExecuteDownload(context.Background(), d)
516 done <- err
517 }()
518
519 waitForFile(t, filepath.Join(e.scratch, "started"))
520 if err := e.svc.Delete(d.ID); err != nil {
521 t.Fatalf("Delete: %v", err)
522 }
523
524 select {
525 case err := <-done:
526 if err != ErrCancelled {
527 t.Errorf("error = %v, want ErrCancelled", err)
528 }
529 case <-time.After(20 * time.Second):
530 t.Fatal("deleting the download did not stop it")
531 }
532
533 if _, err := e.repo.GetByID(d.ID); err == nil {
534 t.Error("download row still exists after Delete")
535 }
536 if names := libraryEntries(t, e.cfg.LibraryDir); len(names) != 0 {
537 t.Errorf("library = %v, want empty", names)
538 }
539}
540
541// Clearing the queue must also stop what is running, for the same reason.
542func TestDeleteAllCancelsRunningDownload(t *testing.T) {
543 e := newExecEnv(t)
544 e.fakeYTDLP(t, `touch "$SCRATCH/started"; sleep 60`)
545
546 d := e.queue(t, &models.Download{})
547
548 done := make(chan error, 1)
549 go func() {
550 _, err := e.svc.ExecuteDownload(context.Background(), d)
551 done <- err
552 }()
553
554 waitForFile(t, filepath.Join(e.scratch, "started"))
555 if err := e.svc.DeleteAll(); err != nil {
556 t.Fatalf("DeleteAll: %v", err)
557 }
558
559 select {
560 case err := <-done:
561 if err != ErrCancelled {
562 t.Errorf("error = %v, want ErrCancelled", err)
563 }
564 case <-time.After(20 * time.Second):
565 t.Fatal("clearing the queue did not stop the running download")
566 }
567}
568