various fixes

AuthorKonata <konata@posteo.jp>
Date
Commitc338db6fda8056abec2ddeb68d73063b04b61a10
Parentcfabced
13 files changed, 352 insertions(+), 171 deletions(-)
▾Mcmd/vidarchive/main.go
@@ -4,6 +4,7 @@ import (
"fmt"
"log"
"os"
"os/exec"
"path/filepath"
"time"
@@ -19,6 +20,8 @@ import (
func main() {
cfg := config.New()
checkDependencies(cfg)
if err := os.MkdirAll(cfg.DataDir, 0755); err != nil {
log.Fatalf("Failed to create data dir: %v", err)
}
@@ -44,7 +47,7 @@ func main() {
subscriptionRepo := repository.NewSubscriptionRepository(db)
presetSvc := service.NewPresetService(presetRepo)
librarySvc := service.NewLibraryService(cfg.LibraryDir)
librarySvc := service.NewLibraryService(cfg.LibraryDir, cfg.FFmpegPath, cfg.FFprobePath)
settingsSvc := service.NewSettingsService(settingsRepo)
subscriptionSvc := service.NewSubscriptionService(subscriptionRepo, cfg)
downloadSvc := service.NewDownloadService(downloadRepo, librarySvc, presetSvc, settingsSvc, subscriptionSvc, cfg)
@@ -73,3 +76,19 @@ func main() {
log.Fatalf("Server error: %v", err)
}
}
// checkDependencies warns (loudly, but without aborting) when an external binary
// VidArchive relies on isn't on PATH, so a missing yt-dlp/ffmpeg surfaces as a
// clear startup message rather than an opaque per-download failure later.
func checkDependencies(cfg *config.Config) {
deps := []struct{ label, path string }{
{"yt-dlp", cfg.YTDLPPath},
{"ffmpeg", cfg.FFmpegPath},
{"ffprobe", cfg.FFprobePath},
}
for _, dep := range deps {
if _, err := exec.LookPath(dep.path); err != nil {
log.Printf("WARNING: %s (%q) not found in PATH; related features will fail until it is installed", dep.label, dep.path)
}
}
}
▾Minternal/config/config.go
@@ -15,6 +15,8 @@ type Config struct {
LibraryDir string
TempDir string
YTDLPPath string
FFmpegPath string
FFprobePath string
BaseURL string
Workers int
RefreshInterval int
@@ -31,6 +33,8 @@ func New() *Config {
LibraryDir: getEnv("VIDARCHIVE_LIBRARY_DIR", filepath.Join(dataDir, "library")),
TempDir: getEnv("VIDARCHIVE_TEMP_DIR", filepath.Join(dataDir, "temp")),
YTDLPPath: getEnv("VIDARCHIVE_YTDLP_PATH", "yt-dlp"),
FFmpegPath: getEnv("VIDARCHIVE_FFMPEG_PATH", "ffmpeg"),
FFprobePath: getEnv("VIDARCHIVE_FFPROBE_PATH", "ffprobe"),
BaseURL: getEnv("VIDARCHIVE_BASE_URL", ""),
Workers: getEnvInt("VIDARCHIVE_WORKERS", 2),
RefreshInterval: getEnvInt("VIDARCHIVE_REFRESH_INTERVAL", 5),
▾Minternal/models/models.go
@@ -110,7 +110,6 @@ type ItemMetadata struct {
SourceURL string `toml:"source_url"`
Description string `toml:"description"`
VideoID string `toml:"video_id"`
Extractor string `toml:"extractor"`
YtdlpFlags string `toml:"ytdlp_flags"`
FileDurations map[string]int `toml:"file_durations"`
}
▾Minternal/repository/download.go
@@ -140,6 +140,21 @@ func (r *DownloadRepository) MarkError(id int64, errMsg string) error {
return err
}
// HasActiveForSubscription reports whether the given subscription already has a
// download that is queued or in progress, so the scheduler can avoid stacking a
// second run on top of one that hasn't finished.
func (r *DownloadRepository) HasActiveForSubscription(subID int64) (bool, error) {
var n int
err := r.db.QueryRow(
`SELECT COUNT(*) FROM downloads WHERE subscription_id = ? AND status IN ('queued', 'downloading')`,
subID,
).Scan(&n)
if err != nil {
return false, err
}
return n > 0, nil
}
func (r *DownloadRepository) Delete(id int64) error {
_, err := r.db.Exec(`DELETE FROM downloads WHERE id = ?`, id)
return err
▾Minternal/server/server_test.go
@@ -45,7 +45,7 @@ func setupTestServer(t *testing.T) (*Server, *config.Config, func()) {
subscriptionRepo := repository.NewSubscriptionRepository(db)
presetSvc := service.NewPresetService(presetRepo)
librarySvc := service.NewLibraryService(cfg.LibraryDir)
librarySvc := service.NewLibraryService(cfg.LibraryDir, cfg.FFmpegPath, cfg.FFprobePath)
settingsSvc := service.NewSettingsService(settingsRepo)
subscriptionSvc := service.NewSubscriptionService(subscriptionRepo, cfg)
downloadSvc := service.NewDownloadService(downloadRepo, librarySvc, presetSvc, settingsSvc, subscriptionSvc, cfg)
▾Minternal/service/download.go
@@ -2,6 +2,7 @@ package service
import (
"bufio"
"bytes"
"database/sql"
"encoding/json"
"fmt"
@@ -9,7 +10,6 @@ import (
"os"
"os/exec"
"path/filepath"
"regexp"
"sort"
"strconv"
"strings"
@@ -118,6 +118,12 @@ func (s *DownloadService) GetQueued(limit int) ([]*models.Download, error) {
return s.repo.GetQueued(limit)
}
// HasActiveForSubscription reports whether the subscription already has a queued
// or in-progress download, so the scheduler can skip stacking another run.
func (s *DownloadService) HasActiveForSubscription(subID int64) (bool, error) {
return s.repo.HasActiveForSubscription(subID)
}
func (s *DownloadService) Delete(id int64) error {
s.killProcess(id)
s.cache.Delete(id)
@@ -146,13 +152,18 @@ func (s *DownloadService) ResetStalledDownloads() error {
}
func (s *DownloadService) ListFormats(url string) ([]*models.FormatInfo, error) {
cmd := exec.Command(s.cfg.YTDLPPath, "-F", "--no-warnings", url)
output, err := cmd.CombinedOutput()
// Use machine-readable JSON (-J) rather than scraping the human "-F" table,
// whose columns/separators shift between yt-dlp versions. stderr is captured
// separately so warnings can't corrupt the JSON on stdout.
cmd := exec.Command(s.cfg.YTDLPPath, "-J", "--no-warnings", url)
var stderr bytes.Buffer
cmd.Stderr = &stderr
output, err := cmd.Output()
if err != nil {
return nil, fmt.Errorf("yt-dlp -F failed: %w\nOutput: %s", err, string(output))
return nil, fmt.Errorf("yt-dlp -J failed: %w\n%s", err, stderr.String())
}
return parseFormatList(string(output)), nil
return parseFormatJSON(output)
}
// ExecuteDownload runs the download for d. The bool reports whether this call
@@ -223,7 +234,7 @@ func (s *DownloadService) ExecuteDownload(d *models.Download) (bool, error) {
if sub != nil {
// Always write info.json so the import step can read the stable identity
// (id + extractor) used to match/replace existing items.
// (yt-dlp's video id) used to match/replace existing items.
args = append(args, "--write-info-json")
switch sub.RefreshMode {
case "skip":
@@ -259,10 +270,24 @@ func (s *DownloadService) ExecuteDownload(d *models.Download) (bool, error) {
// The main pass ran with --skip-download, so the temp dir holds only
// info.json files: refresh existing items in place and fetch new ones.
if err := s.refreshAndAddNew(d, preset, tempDownloadDir, ytdlpFlags); err != nil {
log.Printf("Download %d metadata refresh failed: %v", d.ID, err)
s.finalizeError(d.ID, err)
return false, err
}
} else {
imported, err := s.importDownloadedItems(d, tempDownloadDir, mode, ytdlpFlags)
if err != nil {
s.finalizeError(d.ID, err)
return false, err
}
// A plain (non-subscription) download that yields nothing is a failure, not
// a silent "completed". Subscription modes legitimately import zero (skip
// mode, or a metadata refresh with no new entries), so only enforce this for
// plain runs.
if sub == nil && imported == 0 {
err := fmt.Errorf("yt-dlp finished but no media files were downloaded")
s.finalizeError(d.ID, err)
return false, err
}
} else if err := s.importDownloadedItems(d, tempDownloadDir, mode, ytdlpFlags); err != nil {
log.Printf("Download %d import failed: %v", d.ID, err)
}
if sub != nil && sub.PruneRemoved {
@@ -280,7 +305,12 @@ func (s *DownloadService) ExecuteDownload(d *models.Download) (bool, error) {
// progress cache and periodically flushing it to the download's persisted log.
// It registers the process so Delete can kill it, and returns the exit error.
func (s *DownloadService) runYTDLP(d *models.Download, args []string) error {
cmd := exec.Command(s.cfg.YTDLPPath, args...)
// --newline forces yt-dlp to emit each progress update on its own line. Without
// it, progress is rewritten in place with carriage returns, so a long download
// becomes one ever-growing line that overflows the reader's buffer and stalls
// the pipe — hanging the download. See the hardened scanner below.
fullArgs := append([]string{"--newline"}, args...)
cmd := exec.Command(s.cfg.YTDLPPath, fullArgs...)
cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true}
stdout, err := cmd.StdoutPipe()
@@ -319,9 +349,15 @@ func (s *DownloadService) runYTDLP(d *models.Download, args []string) error {
}()
scanner := bufio.NewScanner(stdout)
// Allow long lines (a single yt-dlp message can exceed the 64 KiB default)
// rather than letting the scanner abort and leave the pipe unread.
scanner.Buffer(make([]byte, 0, 64*1024), 1024*1024)
for scanner.Scan() {
s.cache.AppendLog(d.ID, scanner.Text())
}
if err := scanner.Err(); err != nil {
log.Printf("download %d: error reading yt-dlp output: %v", d.ID, err)
}
close(done)
if logs := s.cache.FlushLogs(d.ID); logs != "" {
@@ -344,8 +380,18 @@ func (s *DownloadService) appendCookies(args []string) ([]string, func()) {
if err != nil {
return args, cleanup
}
tmpFile.WriteString(cookies)
tmpFile.Close()
// A short write would hand yt-dlp a truncated cookies file; on any write/close
// failure, drop the temp file and proceed without cookies rather than silently
// using a broken one.
if _, err := tmpFile.WriteString(cookies); err != nil {
tmpFile.Close()
os.Remove(tmpFile.Name())
return args, cleanup
}
if err := tmpFile.Close(); err != nil {
os.Remove(tmpFile.Name())
return args, cleanup
}
return append(args, "--cookies", tmpFile.Name()), func() { os.Remove(tmpFile.Name()) }
}
@@ -378,18 +424,22 @@ func (s *DownloadService) resolveBaseLibraryDir(d *models.Download) (string, err
return baseLibraryDir, nil
}
func (s *DownloadService) importDownloadedItems(d *models.Download, tempDownloadDir, mode, ytdlpFlags string) error {
// importDownloadedItems moves each downloaded item from the temp dir into the
// library and returns the number of items successfully imported. Per-item
// failures are logged and skipped (a playlist with a few bad entries still
// imports the rest); a non-nil error means the import couldn't even start.
func (s *DownloadService) importDownloadedItems(d *models.Download, tempDownloadDir, mode, ytdlpFlags string) (int, error) {
entries, err := os.ReadDir(tempDownloadDir)
if err != nil {
return err
return 0, err
}
baseLibraryDir, err := s.resolveBaseLibraryDir(d)
if err != nil {
return err
return 0, err
}
if err := os.MkdirAll(baseLibraryDir, 0755); err != nil {
return err
return 0, err
}
var itemDirs []string
@@ -404,14 +454,20 @@ func (s *DownloadService) importDownloadedItems(d *models.Download, tempDownload
}
sort.Strings(itemDirs)
imported := 0
for _, itemDir := range itemDirs {
if err := s.importItemDir(d.URL, itemDir, baseLibraryDir, mode, ytdlpFlags); err != nil {
log.Printf("warning: failed to import item %s: %v", itemDir, err)
continue
}
imported++
}
os.Remove(tempDownloadDir)
return nil
// RemoveAll (not Remove): leftover item dirs from failed imports, plus any
// orphaned thumbnail/image files yt-dlp left behind, would otherwise keep the
// temp dir non-empty and leak it forever.
os.RemoveAll(tempDownloadDir)
return imported, nil
}
func (s *DownloadService) importItemDir(url, itemDir, baseLibraryDir, mode, ytdlpFlags string) error {
@@ -452,13 +508,13 @@ func (s *DownloadService) importItemDir(url, itemDir, baseLibraryDir, mode, ytdl
}
name := s.deriveItemName(itemDir, infoJSONPath, mediaFiles)
videoID, extractor := readInfoIdentity(infoJSONPath)
videoID := readInfoID(infoJSONPath)
// Overwrite mode: replace the existing copy of this video in place rather than
// creating a duplicate folder. Removing the old dir lets uniqueDir reuse its
// name (or land on the new title if it changed upstream).
if mode == "overwrite" && videoID != "" {
if existing, ok := s.librarySvc.FindByVideoID(baseLibraryDir, extractor, videoID); ok {
if existing, ok := s.librarySvc.FindByVideoID(baseLibraryDir, videoID); ok {
if rel, err := filepath.Rel(s.cfg.LibraryDir, existing); err == nil {
s.librarySvc.evictCachedScan(filepath.ToSlash(rel))
}
@@ -488,7 +544,7 @@ func (s *DownloadService) importItemDir(url, itemDir, baseLibraryDir, mode, ytdl
// serving pages. Files we can't probe simply get no duration.
fileDurations := make(map[string]int)
for _, entry := range mediaFiles {
if d, ok := probeDuration(filepath.Join(targetDir, entry.Name())); ok {
if d, ok := probeDuration(s.cfg.FFprobePath, filepath.Join(targetDir, entry.Name())); ok {
fileDurations[entry.Name()] = d
}
}
@@ -509,7 +565,6 @@ func (s *DownloadService) importItemDir(url, itemDir, baseLibraryDir, mode, ytdl
Name: name,
SourceURL: url,
VideoID: videoID,
Extractor: extractor,
YtdlpFlags: ytdlpFlags,
FileDurations: fileDurations,
}
@@ -527,24 +582,25 @@ func (s *DownloadService) importItemDir(url, itemDir, baseLibraryDir, mode, ytdl
return nil
}
// readInfoIdentity extracts the stable identity (yt-dlp id + extractor) from an
// info.json. Returns empty strings when the file is absent or unreadable.
func readInfoIdentity(infoJSONPath string) (id, extractor string) {
// readInfoID extracts the stable item identity (yt-dlp's video id) from an
// info.json. The id alone is sufficient to match items within a single
// subscription's owned directory, so the extractor is not used. Returns an
// empty string when the file is absent or unreadable.
func readInfoID(infoJSONPath string) string {
if infoJSONPath == "" {
return "", ""
return ""
}
data, err := os.ReadFile(infoJSONPath)
if err != nil {
return "", ""
return ""
}
var info struct {
ID string `json:"id"`
Extractor string `json:"extractor"`
ID string `json:"id"`
}
if err := json.Unmarshal(data, &info); err != nil {
return "", ""
return ""
}
return info.ID, info.Extractor
return info.ID
}
// refreshAndAddNew handles a metadata-mode run. The main pass used
@@ -575,12 +631,12 @@ func (s *DownloadService) refreshAndAddNew(d *models.Download, preset *models.Pr
if infoJSONPath == "" {
continue
}
videoID, extractor := readInfoIdentity(infoJSONPath)
videoID := readInfoID(infoJSONPath)
if videoID == "" {
continue
}
if existing, ok := s.librarySvc.FindByVideoID(baseLibraryDir, extractor, videoID); ok {
if err := s.applyMetadata(existing, infoJSONPath, videoID, extractor); err != nil {
if existing, ok := s.librarySvc.FindByVideoID(baseLibraryDir, videoID); ok {
if err := s.applyMetadata(existing, infoJSONPath, videoID); err != nil {
log.Printf("warning: failed to refresh metadata for %s: %v", itemDir, err)
}
continue
@@ -617,7 +673,7 @@ func (s *DownloadService) downloadFresh(d *models.Download, preset *models.Prese
runErr := s.runYTDLP(d, args)
// Import whatever succeeded even if some entries errored.
if err := s.importDownloadedItems(d, tempDir, "", ytdlpFlags); err != nil {
if _, err := s.importDownloadedItems(d, tempDir, "", ytdlpFlags); err != nil {
log.Printf("warning: failed to import new metadata-mode items: %v", err)
}
return runErr
@@ -625,7 +681,7 @@ func (s *DownloadService) downloadFresh(d *models.Download, preset *models.Prese
// applyMetadata rewrites an existing item's marker (name/description/identity)
// from a fresh info.json without touching its media.
func (s *DownloadService) applyMetadata(existing, infoJSONPath, videoID, extractor string) error {
func (s *DownloadService) applyMetadata(existing, infoJSONPath, videoID string) error {
data, err := os.ReadFile(infoJSONPath)
if err != nil {
return err
@@ -648,7 +704,6 @@ func (s *DownloadService) applyMetadata(existing, infoJSONPath, videoID, extract
meta.SourceURL = info.WebpageURL
}
meta.VideoID = videoID
meta.Extractor = extractor
if err := s.librarySvc.writeMetadata(existing, meta); err != nil {
return err
@@ -723,11 +778,12 @@ func (s *DownloadService) pruneSubscription(d *models.Download, sub *models.Subs
}
}
// enumeratePlaylistIDs lists the current "extractor:id" set for a URL without
// enumeratePlaylistIDs lists the current video-id set for a URL without
// downloading, using yt-dlp --flat-playlist. Cookies are applied so private
// playlists enumerate correctly.
// playlists enumerate correctly. Ids alone are sufficient to match items within
// a subscription's own directory (see FindByVideoID / PruneToIDSet).
func (s *DownloadService) enumeratePlaylistIDs(url string) (map[string]bool, error) {
args := []string{"--flat-playlist", "--no-warnings", "--print", "%(extractor)s %(id)s"}
args := []string{"--flat-playlist", "--no-warnings", "--print", "%(id)s"}
args, cleanup := s.appendCookies(args)
defer cleanup()
@@ -740,13 +796,12 @@ func (s *DownloadService) enumeratePlaylistIDs(url string) (map[string]bool, err
keep := make(map[string]bool)
for _, line := range strings.Split(string(out), "\n") {
fields := strings.Fields(line)
if len(fields) != 2 {
id := strings.TrimSpace(line)
// yt-dlp prints "NA" for a missing field; never treat that as a real id.
if id == "" || id == "NA" {
continue
}
if key := videoKey(fields[0], fields[1]); key != "" {
keep[key] = true
}
keep[id] = true
}
return keep, nil
}
@@ -811,13 +866,10 @@ func sanitizeDirName(name string) string {
return name
}
// resolutionRe matches a yt-dlp resolution column like "1920x1080".
var resolutionRe = regexp.MustCompile(`^\d+x\d+$`)
// probeDuration returns the duration of a media file in whole seconds. The bool
// is false when ffprobe is unavailable or the file has no usable duration.
func probeDuration(path string) (int, bool) {
out, err := exec.Command("ffprobe", "-v", "error",
func probeDuration(ffprobePath, path string) (int, bool) {
out, err := exec.Command(ffprobePath, "-v", "error",
"-show_entries", "format=duration",
"-of", "default=nw=1:nk=1", path).Output()
if err != nil {
@@ -830,54 +882,100 @@ func probeDuration(path string) (int, bool) {
return int(f + 0.5), true
}
func parseFormatList(output string) []*models.FormatInfo {
lines := strings.Split(output, "\n")
var formats []*models.FormatInfo
// ytFormat mirrors the subset of yt-dlp's per-format JSON (-J) we surface.
// Numeric fields are pointers so an absent value (null/omitted) is distinct
// from a real zero.
type ytFormat struct {
FormatID string `json:"format_id"`
Ext string `json:"ext"`
Resolution string `json:"resolution"`
Width *int `json:"width"`
Height *int `json:"height"`
FPS *float64 `json:"fps"`
VCodec string `json:"vcodec"`
ACodec string `json:"acodec"`
AudioChannels *int `json:"audio_channels"`
Filesize *int64 `json:"filesize"`
FilesizeApprox *int64 `json:"filesize_approx"`
FormatNote string `json:"format_note"`
}
inFormats := false
for _, line := range lines {
line = strings.TrimSpace(line)
if line == "" {
continue
}
// parseFormatJSON reads yt-dlp's single-JSON dump (-J) and returns the available
// formats. For a single video the formats live at the top level; for a playlist
// URL we fall back to the first entry's formats so the picker still shows
// something useful.
func parseFormatJSON(data []byte) ([]*models.FormatInfo, error) {
var top struct {
Formats []ytFormat `json:"formats"`
Entries []struct {
Formats []ytFormat `json:"formats"`
} `json:"entries"`
}
if err := json.Unmarshal(data, &top); err != nil {
return nil, fmt.Errorf("parse yt-dlp JSON: %w", err)
}
if strings.Contains(line, "ID") && strings.Contains(line, "EXT") {
inFormats = true
continue
}
raw := top.Formats
if len(raw) == 0 && len(top.Entries) > 0 {
raw = top.Entries[0].Formats
}
if !inFormats {
continue
}
formats := make([]*models.FormatInfo, 0, len(raw))
for _, f := range raw {
formats = append(formats, f.toFormatInfo())
}
return formats, nil
}
parts := strings.Fields(line)
if len(parts) >= 4 {
format := &models.FormatInfo{
ID: parts[0],
Ext: parts[1],
}
func (f ytFormat) toFormatInfo() *models.FormatInfo {
fi := &models.FormatInfo{
ID: f.FormatID,
Ext: f.Ext,
Note: f.FormatNote,
}
for i, part := range parts {
// A resolution token is strictly <digits>x<digits> (e.g. 1920x1080);
// matching on a literal "x" anywhere misclassified notes/codecs.
if resolutionRe.MatchString(part) {
format.Resolution = part
if i+1 < len(parts) {
format.FPS = parts[i+1]
}
break
}
}
switch {
case f.Resolution != "":
fi.Resolution = f.Resolution
case f.Width != nil && f.Height != nil && *f.Width > 0 && *f.Height > 0:
fi.Resolution = fmt.Sprintf("%dx%d", *f.Width, *f.Height)
}
if len(parts) > 3 {
format.Note = strings.Join(parts[3:], " ")
}
if f.FPS != nil && *f.FPS > 0 {
fi.FPS = strconv.FormatFloat(*f.FPS, 'f', -1, 64)
}
if f.AudioChannels != nil && *f.AudioChannels > 0 {
fi.Channels = strconv.Itoa(*f.AudioChannels)
}
formats = append(formats, format)
}
// Prefer the video codec; fall back to the audio codec for audio-only formats.
if f.VCodec != "" && f.VCodec != "none" {
fi.Codec = f.VCodec
} else if f.ACodec != "" && f.ACodec != "none" {
fi.Codec = f.ACodec
}
return formats
if f.Filesize != nil && *f.Filesize > 0 {
fi.FileSize = humanizeBytes(*f.Filesize)
} else if f.FilesizeApprox != nil && *f.FilesizeApprox > 0 {
fi.FileSize = "~" + humanizeBytes(*f.FilesizeApprox)
}
return fi
}
// humanizeBytes renders a byte count as a compact human-readable size.
func humanizeBytes(n int64) string {
const unit = 1024
if n < unit {
return fmt.Sprintf("%dB", n)
}
div, exp := int64(unit), 0
for m := n / unit; m >= unit; m /= unit {
div *= unit
exp++
}
return fmt.Sprintf("%.1f%ciB", float64(n)/float64(div), "KMGTPE"[exp])
}
func sqlNullInt64(v int64) sql.NullInt64 {
@@ -894,6 +992,7 @@ var reservedFlags = map[string]string{
"--paths": "the download path",
"--cookies": "cookies (set these in Settings instead)",
"--no-cookies": "cookies (set these in Settings instead)",
"--newline": "progress output formatting (VidArchive sets this to stream logs)",
}
// reservedSubscriptionFlags are additionally reserved for subscription runs,
▾Minternal/service/download_test.go
@@ -94,27 +94,53 @@ func dirEntry(t *testing.T, dir, name string) os.DirEntry {
return nil
}
func TestParseFormatList(t *testing.T) {
output := `[info] Available formats:
ID EXT RESOLUTION FPS
18 mp4 640x360 30
137 mp4 1920x1080 60
233 m4a audio_only_xtra 128k
`
formats := parseFormatList(output)
func TestParseFormatJSON(t *testing.T) {
data := []byte(`{
"id": "vid",
"formats": [
{"format_id": "18", "ext": "mp4", "resolution": "640x360", "fps": 30, "vcodec": "avc1", "acodec": "mp4a", "format_note": "360p"},
{"format_id": "137", "ext": "mp4", "width": 1920, "height": 1080, "fps": 60, "vcodec": "avc1", "acodec": "none", "filesize": 1048576, "format_note": "1080p"},
{"format_id": "233", "ext": "m4a", "resolution": "audio only", "vcodec": "none", "acodec": "mp4a", "audio_channels": 2, "format_note": "audio"}
]
}`)
formats, err := parseFormatJSON(data)
if err != nil {
t.Fatalf("parseFormatJSON: %v", err)
}
if len(formats) != 3 {
t.Fatalf("expected 3 formats, got %d: %+v", len(formats), formats)
}
if formats[0].ID != "18" || formats[0].Ext != "mp4" {
if formats[0].ID != "18" || formats[0].Ext != "mp4" || formats[0].Resolution != "640x360" {
t.Errorf("format[0] = %+v", formats[0])
}
if formats[0].FPS != "30" {
t.Errorf("format[0] fps = %q, want 30", formats[0].FPS)
}
// Resolution is derived from width/height when no resolution string is present.
if formats[1].Resolution != "1920x1080" {
t.Errorf("format[1] resolution = %q, want 1920x1080", formats[1].Resolution)
}
// A non-resolution token containing the letter 'x' must NOT be parsed as a
// resolution (the old "contains x" heuristic misclassified it).
if formats[2].Resolution != "" {
t.Errorf("format[2] resolution = %q, want empty (token has 'x' but isn't WxH)", formats[2].Resolution)
if formats[1].FileSize != "1.0MiB" {
t.Errorf("format[1] filesize = %q, want 1.0MiB", formats[1].FileSize)
}
// Audio-only format: codec falls back to acodec and channels are populated.
if formats[2].Codec != "mp4a" {
t.Errorf("format[2] codec = %q, want mp4a", formats[2].Codec)
}
if formats[2].Channels != "2" {
t.Errorf("format[2] channels = %q, want 2", formats[2].Channels)
}
}
func TestParseFormatJSONPlaylistFallback(t *testing.T) {
// A playlist dump exposes formats under the first entry, not at the top level.
data := []byte(`{"_type":"playlist","entries":[{"id":"a","formats":[{"format_id":"18","ext":"mp4"}]}]}`)
formats, err := parseFormatJSON(data)
if err != nil {
t.Fatalf("parseFormatJSON: %v", err)
}
if len(formats) != 1 || formats[0].ID != "18" {
t.Fatalf("expected 1 format from entry fallback, got %+v", formats)
}
}
@@ -125,7 +151,7 @@ func TestImportItemDir(t *testing.T) {
requireFFmpeg(t)
libDir := t.TempDir()
svc := &DownloadService{cfg: &config.Config{LibraryDir: libDir}}
svc := &DownloadService{cfg: &config.Config{LibraryDir: libDir, FFprobePath: "ffprobe"}}
src := t.TempDir()
makeTestVideo(t, filepath.Join(src, "raw.mp4"))
@@ -220,7 +246,7 @@ func TestImportDownloadedItemsRejectsOutputTraversal(t *testing.T) {
URL: "u",
OutputDir: sql.NullString{String: "../escape", Valid: true},
}
if err := svc.importDownloadedItems(d, tempDir, "", ""); err == nil {
if _, err := svc.importDownloadedItems(d, tempDir, "", ""); err == nil {
t.Error("expected path-traversal output dir to be rejected")
}
}
▾Minternal/service/library.go
@@ -41,6 +41,11 @@ const maxConcurrentThumbnails = 3
type LibraryService struct {
libraryDir string
// ffmpegPath/ffprobePath are the binaries used for thumbnail extraction,
// subtitle conversion, and media probing. Configurable so non-PATH installs
// (e.g. a pinned build) can be pointed at directly.
ffmpegPath string
ffprobePath string
// thumbLocks maps a media filepath -> *sync.Mutex to serialize extraction
// per file (so two callers never write the same temp file at once). Entries
// are bounded by the number of distinct media files ever requested, not by
@@ -80,12 +85,14 @@ type scanCacheEntry struct {
// scanCacheTTL is how long a scanned item is reused before being re-read.
const scanCacheTTL = 10 * time.Second
func NewLibraryService(libraryDir string) *LibraryService {
func NewLibraryService(libraryDir, ffmpegPath, ffprobePath string) *LibraryService {
return &LibraryService{
libraryDir: libraryDir,
thumbSem: make(chan struct{}, maxConcurrentThumbnails),
scanCache: make(map[string]scanCacheEntry),
scanTTL: scanCacheTTL,
libraryDir: libraryDir,
ffmpegPath: ffmpegPath,
ffprobePath: ffprobePath,
thumbSem: make(chan struct{}, maxConcurrentThumbnails),
scanCache: make(map[string]scanCacheEntry),
scanTTL: scanCacheTTL,
}
}
@@ -306,21 +313,15 @@ func (s *LibraryService) scanItem(itemDir, relPath string) (*models.LibraryItem,
}
}
// Backfill the stable identity (yt-dlp id + extractor) from info.json so
// Backfill the stable identity (yt-dlp's video id) from info.json so
// pre-existing items gain an identity on their next scan. Subscriptions match
// and prune items by this key (see videoKey / FindByVideoID).
// and prune items by this id (see FindByVideoID / PruneToIDSet).
if metadata.VideoID == "" {
if id, ok := infoString(info, "id"); ok && id != "" {
metadata.VideoID = id
dirty = true
}
}
if metadata.Extractor == "" {
if ex, ok := infoString(info, "extractor"); ok && ex != "" {
metadata.Extractor = ex
dirty = true
}
}
if metadata.FileDurations == nil {
metadata.FileDurations = make(map[string]int)
@@ -632,8 +633,8 @@ func (s *LibraryService) findExistingThumbnail(path string) (string, bool) {
// yt-dlp embeds into MKV with --embed-thumbnail — or -1 if there is none. Such
// covers are attachment streams, not attached_pic video streams, so they must be
// dumped with -dump_attachment rather than mapped like a normal stream.
func findImageAttachment(path string) int {
cmd := exec.Command("ffprobe", "-v", "error", "-show_streams", "-of", "json", path)
func (s *LibraryService) findImageAttachment(path string) int {
cmd := exec.Command(s.ffprobePath, "-v", "error", "-show_streams", "-of", "json", path)
output, err := cmd.Output()
if err != nil {
return -1
@@ -699,7 +700,7 @@ func (s *LibraryService) extractThumbnail(mf models.MediaFile) (string, error) {
outExt := filepath.Ext(outputPath)
tmpPath := strings.TrimSuffix(outputPath, outExt) + ".tmp" + outExt
os.Remove(tmpPath)
cmd := exec.Command("ffmpeg", append(args, tmpPath)...)
cmd := exec.Command(s.ffmpegPath, append(args, tmpPath)...)
if output, err := cmd.CombinedOutput(); err != nil {
os.Remove(tmpPath)
return "", fmt.Errorf("ffmpeg failed: %v\n%s", err, string(output))
@@ -718,7 +719,7 @@ func (s *LibraryService) extractThumbnail(mf models.MediaFile) (string, error) {
// Prefer an embedded image attachment (e.g. yt-dlp's cover.webp/cover.jpg in
// MKV). These are attachment streams, not mappable video streams, so dump the
// raw bytes with -dump_attachment then transcode to a canonical WebP.
if idx := findImageAttachment(mf.Filepath); idx >= 0 {
if idx := s.findImageAttachment(mf.Filepath); idx >= 0 {
if raw, err := os.CreateTemp("", "vidarchive-attachment-*"); err == nil {
rawPath := raw.Name()
raw.Close()
@@ -727,7 +728,7 @@ func (s *LibraryService) extractThumbnail(mf models.MediaFile) (string, error) {
fmt.Sprintf("-dump_attachment:t:%d", idx), rawPath,
"-i", mf.Filepath, "-y", "-t", "0", "-f", "null", "-",
}
if out, err := exec.Command("ffmpeg", dumpArgs...).CombinedOutput(); err != nil {
if out, err := exec.Command(s.ffmpegPath, dumpArgs...).CombinedOutput(); err != nil {
record("attachment-dump", fmt.Errorf("%v\n%s", err, out))
} else if path, err := tryWrite(webpPath, []string{"-i", rawPath, "-c:v", "libwebp"}); err == nil {
return path, nil
@@ -833,23 +834,14 @@ func (s *LibraryService) Delete(relPath string) error {
return os.RemoveAll(itemDir)
}
// videoKey is the stable per-item identity used to match items across runs:
// yt-dlp's extractor + id, matching the form used in --download-archive and the
// "%(extractor)s %(id)s" enumeration. Empty when either component is missing.
func videoKey(extractor, id string) string {
if extractor == "" || id == "" {
return ""
}
return extractor + ":" + id
}
// FindByVideoID returns the absolute directory of the item under baseDir whose
// marker matches the given extractor+id, scanning only that directory (the
// subscription's owned folder). ok is false when no match is found or the key
// is incomplete.
func (s *LibraryService) FindByVideoID(baseDir, extractor, id string) (string, bool) {
want := videoKey(extractor, id)
if want == "" {
// marker matches the given yt-dlp video id, scanning only that directory (the
// subscription's owned folder). Matching on the id alone is safe here because
// each subscription owns a single source, so ids don't collide across
// extractors within the folder. ok is false when no match is found or id is
// empty.
func (s *LibraryService) FindByVideoID(baseDir, id string) (string, bool) {
if id == "" {
return "", false
}
entries, err := os.ReadDir(baseDir)
@@ -865,7 +857,7 @@ func (s *LibraryService) FindByVideoID(baseDir, extractor, id string) (string, b
continue
}
meta, _ := s.readOrCreateMetadata(itemDir)
if videoKey(meta.Extractor, meta.VideoID) == want {
if meta.VideoID == id {
return itemDir, true
}
}
@@ -892,8 +884,7 @@ func (s *LibraryService) PruneToIDSet(baseDir string, keep map[string]bool) (int
continue
}
meta, _ := s.readOrCreateMetadata(itemDir)
key := videoKey(meta.Extractor, meta.VideoID)
if key == "" || keep[key] {
if meta.VideoID == "" || keep[meta.VideoID] {
continue
}
if rel, err := filepath.Rel(s.libraryDir, itemDir); err == nil {
@@ -948,7 +939,7 @@ func (s *LibraryService) GetSubtitles(ctx context.Context, relPath string) ([]mo
if mf.IsAudio {
continue
}
streams, err := extractSubtitleInfo(mf.Filepath)
streams, err := s.extractSubtitleInfo(mf.Filepath)
if err != nil || len(streams) == 0 {
continue
}
@@ -962,7 +953,7 @@ func (s *LibraryService) GetSubtitles(ctx context.Context, relPath string) ([]mo
lang = fmt.Sprintf("track%d", stream.Index)
}
outPath := filepath.Join(cacheDir, lang+".vtt")
if err := extractSubtitleToVTT(mf.Filepath, outPath, stream.Index); err != nil {
if err := s.extractSubtitleToVTT(mf.Filepath, outPath, stream.Index); err != nil {
continue
}
tracks = append(tracks, models.SubtitleTrack{
@@ -983,8 +974,8 @@ type subtitleStream struct {
Label string
}
func extractSubtitleInfo(path string) ([]subtitleStream, error) {
cmd := exec.Command("ffprobe",
func (s *LibraryService) extractSubtitleInfo(path string) ([]subtitleStream, error) {
cmd := exec.Command(s.ffprobePath,
"-v", "error",
"-show_streams",
"-select_streams", "s",
@@ -1032,8 +1023,8 @@ func extractSubtitleInfo(path string) ([]subtitleStream, error) {
return streams, nil
}
func extractSubtitleToVTT(inputPath, outputPath string, streamIndex int) error {
cmd := exec.Command("ffmpeg",
func (s *LibraryService) extractSubtitleToVTT(inputPath, outputPath string, streamIndex int) error {
cmd := exec.Command(s.ffmpegPath,
"-i", inputPath,
"-map", fmt.Sprintf("0:s:%d", streamIndex),
"-f", "webvtt",
@@ -1072,10 +1063,10 @@ func (s *LibraryService) GetMetadata(ctx context.Context, relPath, filename stri
return nil, fmt.Errorf("no media file")
}
return probeMedia(target.Filepath)
return s.probeMedia(target.Filepath)
}
func probeMedia(path string) (*MediaMetadata, error) {
func (s *LibraryService) probeMedia(path string) (*MediaMetadata, error) {
info, err := os.Stat(path)
if err != nil {
return nil, err
@@ -1085,7 +1076,7 @@ func probeMedia(path string) (*MediaMetadata, error) {
FileSize: info.Size(),
}
cmd := exec.Command("ffprobe",
cmd := exec.Command(s.ffprobePath,
"-v", "error",
"-show_format",
"-show_streams",
▾Minternal/service/library_test.go
@@ -19,7 +19,7 @@ import (
func newLibrary(t *testing.T) (*LibraryService, string) {
t.Helper()
dir := t.TempDir()
return NewLibraryService(dir), dir
return NewLibraryService(dir, "ffmpeg", "ffprobe"), dir
}
// writeItem creates an item directory with a marker and the given files
@@ -756,7 +756,7 @@ func TestThumbnailUsesEmbeddedAttachment(t *testing.T) {
// 100x100 cover so it's distinguishable from a 64x64 video frame.
makeVideoWithCoverAttachment(t, mkv, "100x100")
if idx := findImageAttachment(mkv); idx < 0 {
if idx := svc.findImageAttachment(mkv); idx < 0 {
t.Fatal("findImageAttachment did not find the embedded cover")
}
@@ -797,7 +797,7 @@ func TestThumbnailNegativeCacheSkipsReextraction(t *testing.T) {
}
// The cache is in-process only: a fresh service (≈ a restart) retries.
fresh := NewLibraryService(dir)
fresh := NewLibraryService(dir, "ffmpeg", "ffprobe")
if _, ok := fresh.ThumbnailForFile(ctx, "bad", "broken.mp4"); ok {
t.Fatal("fresh service still fails (file is unextractable)")
}
▾Minternal/service/preset.go
@@ -59,7 +59,7 @@ func (s *PresetService) BuildArgs(p *models.Preset, formatOverride, customFlags
args = append(args, "-f", p.CustomFormat)
} else if p.FormatMode == "preset" && p.Format != "" {
if p.Quality != "" {
args = append(args, "-f", fmt.Sprintf("%s[height<=%s]/%s", p.Format, p.Quality, p.Format))
args = append(args, "-f", applyHeightCap(p.Format, p.Quality))
} else {
args = append(args, "-f", p.Format)
}
@@ -109,6 +109,17 @@ func (s *PresetService) BuildArgs(p *models.Preset, formatOverride, customFlags
return args
}
// applyHeightCap adds a max-height filter to a format selector. For a combined
// "video+audio" selector it binds the filter to the video stream only (e.g.
// "bestvideo[height<=720]+bestaudio"), since applying it to the whole expression
// would wrongly constrain the audio selector. It always appends an unfiltered
// fallback so a missing capped variant still resolves.
func applyHeightCap(format, maxHeight string) string {
parts := strings.Split(format, "+")
parts[0] = fmt.Sprintf("%s[height<=%s]", parts[0], maxHeight)
return strings.Join(parts, "+") + "/" + format
}
func (s *PresetService) EffectiveFlags(p *models.Preset, formatOverride, customFlags string) string {
args := s.BuildArgs(p, formatOverride, customFlags)
return strings.Join(args, " ")
▾Minternal/service/subscription_library_test.go
@@ -9,34 +9,32 @@ import (
func TestFindByVideoID(t *testing.T) {
lib, dir := newLibrary(t)
writeItem(t, dir, "alpha", "name = \"Alpha\"\nvideo_id = \"aaa\"\nextractor = \"youtube\"\n", nil)
writeItem(t, dir, "beta", "name = \"Beta\"\nvideo_id = \"bbb\"\nextractor = \"youtube\"\n", nil)
// Same id but different extractor must not match.
writeItem(t, dir, "gamma", "name = \"Gamma\"\nvideo_id = \"aaa\"\nextractor = \"vimeo\"\n", nil)
writeItem(t, dir, "alpha", "name = \"Alpha\"\nvideo_id = \"aaa\"\n", nil)
writeItem(t, dir, "beta", "name = \"Beta\"\nvideo_id = \"bbb\"\n", nil)
got, ok := lib.FindByVideoID(dir, "youtube", "aaa")
got, ok := lib.FindByVideoID(dir, "aaa")
if !ok || filepath.Base(got) != "alpha" {
t.Fatalf("expected to find alpha, got %q ok=%v", got, ok)
}
if _, ok := lib.FindByVideoID(dir, "youtube", "zzz"); ok {
if _, ok := lib.FindByVideoID(dir, "zzz"); ok {
t.Error("expected no match for unknown id")
}
// Incomplete key never matches.
if _, ok := lib.FindByVideoID(dir, "", "aaa"); ok {
t.Error("expected no match for empty extractor")
// An empty id never matches.
if _, ok := lib.FindByVideoID(dir, ""); ok {
t.Error("expected no match for empty id")
}
}
func TestPruneToIDSet(t *testing.T) {
lib, dir := newLibrary(t)
writeItem(t, dir, "keep", "name = \"Keep\"\nvideo_id = \"k1\"\nextractor = \"youtube\"\n", nil)
writeItem(t, dir, "drop", "name = \"Drop\"\nvideo_id = \"d1\"\nextractor = \"youtube\"\n", nil)
writeItem(t, dir, "keep", "name = \"Keep\"\nvideo_id = \"k1\"\n", nil)
writeItem(t, dir, "drop", "name = \"Drop\"\nvideo_id = \"d1\"\n", nil)
// An item without identity must never be pruned (we can't positively id it).
writeItem(t, dir, "noid", "name = \"NoId\"\n", nil)
keep := map[string]bool{"youtube:k1": true}
keep := map[string]bool{"k1": true}
removed, err := lib.PruneToIDSet(dir, keep)
if err != nil {
t.Fatalf("prune: %v", err)
▾Minternal/worker/pool.go
@@ -95,8 +95,10 @@ func (p *Pool) queueChecker() {
}
func (p *Pool) checkQueue() {
// Check how many items are in queue vs active
downloads, err := p.downloadSvc.GetAll("queued", "date")
// Pull at most a channel's worth of the oldest queued downloads rather than
// loading every queued row each tick. Anything beyond the buffer is picked up
// on a later tick; duplicates are harmless (MarkStarted claims atomically).
downloads, err := p.downloadSvc.GetQueued(cap(p.queue))
if err != nil {
log.Printf("Queue check error: %v", err)
return
▾Minternal/worker/scheduler.go
@@ -110,6 +110,23 @@ func (s *Scheduler) run(sub *models.Subscription, now time.Time) {
next = now.Add(24 * time.Hour)
}
// Skip this run if the previous one is still queued or downloading — a run
// longer than the interval would otherwise stack duplicate downloads. Advance
// next_run_at so we don't re-evaluate it every tick, preserving the last run.
if active, aErr := s.downloadSvc.HasActiveForSubscription(sub.ID); aErr != nil {
log.Printf("scheduler: subscription %d active-run check failed: %v", sub.ID, aErr)
} else if active {
lastRun := now
if sub.LastRunAt.Valid {
lastRun = sub.LastRunAt.Time
}
log.Printf("scheduler: subscription %d still has an active run, skipping until %s", sub.ID, next.Format(time.RFC3339))
if mErr := s.subscriptionSvc.MarkRun(sub.ID, lastRun, next, "skipped"); mErr != nil {
log.Printf("scheduler: failed to mark subscription %d run: %v", sub.ID, mErr)
}
return
}
d, err := s.downloadSvc.CreateForSubscription(sub)
if err != nil {
log.Printf("scheduler: failed to queue subscription %d: %v", sub.ID, err)