comment cleanup
Massets.go
@@ -1,18 +1,14 @@
// Package vidarchive embeds the web assets (templates and static files) into the
// binary so it runs correctly regardless of the working directory or where the
// binary is installed — no reliance on locating the source tree at runtime.
// Package vidarchive embeds the web assets into the binary, so nothing has to
// locate the source tree at runtime.
package vidarchive
import "embed"
// TemplatesFS holds the HTML templates under web/templates.
//
//go:embed web/templates/*.html
var TemplatesFS embed.FS
// StaticFS holds the static assets (CSS, icons) under web/static. The all: prefix
// ensures files that would otherwise be skipped (none currently, but future
// dotfiles) are still embedded.
// The all: prefix also embeds files embed would otherwise skip, such as a
// dotfile added later.
//
//go:embed all:web/static
var StaticFS embed.FS
Mcmd/vidarchive/main.go
@@ -21,15 +21,14 @@ import (
)
func main() {
// All configuration is environment variables (see README). An argument is
// a mistake, so fail instead of silently starting the server.
// Configuration is environment variables only, so an argument is a mistake.
if len(os.Args) > 1 {
fmt.Fprintf(os.Stderr, "vidarchive takes no arguments; configure it with VIDARCHIVE_* environment variables (got %q)\n", os.Args[1:])
os.Exit(2)
}
// Cancelled on SIGINT/SIGTERM so the server drains, workers stop, running
// yt-dlp children are killed and the database is checkpointed and closed.
// SIGINT/SIGTERM drains the server, stops the workers and their yt-dlp
// children, then checkpoints and closes the database.
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
@@ -85,9 +84,8 @@ func main() {
srvErr := srv.Start(ctx)
// Shut down in dependency order: stop scheduling new runs, then stop the
// workers (killing any running yt-dlp), then checkpoint and close the DB.
// This runs on both a signal and a server error, so nothing is left orphaned.
// Dependency order: no new runs, then no workers, then the database. Reached
// on both a signal and a server error, so nothing is left orphaned.
slog.Info("shutting down")
scheduler.Stop()
workerPool.Stop()
@@ -104,8 +102,8 @@ func main() {
}
}
// fatal logs a startup or shutdown failure and exits. slog has no Fatal, and
// log.Fatalf would report it at info level once slog owns the log package.
// fatal exists because slog has no Fatal, and log.Fatalf would report at info
// level once slog owns the log package.
func fatal(msg string, args ...any) {
slog.Error(msg, args...)
os.Exit(1)
@@ -115,9 +113,8 @@ func fatal(msg string, args ...any) {
// with Apache's tools (apache2-utils on Debian, httpd-tools on Fedora).
const hashHint = `htpasswd -bnBC 12 "" 'your-password' | tr -d ':\n'`
// checkCredentials refuses to start without a usable login. A hash that only
// fails at the login form would look like a forgotten password rather than a
// misconfiguration.
// checkCredentials refuses to start without a usable login: a hash that only
// fails at the login form looks like a forgotten password, not a misconfiguration.
func checkCredentials(cfg *config.Config) {
if cfg.Username == "" || cfg.PasswordHash == "" {
fatal("VIDARCHIVE_USERNAME and VIDARCHIVE_PASSWORD_HASH are required", "hint", hashHint)
@@ -127,9 +124,8 @@ func checkCredentials(cfg *config.Config) {
}
}
// 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.
// checkDependencies warns without aborting, so a missing yt-dlp or ffmpeg shows
// up at startup rather than as an opaque per-download failure later.
func checkDependencies(cfg *config.Config) {
deps := []struct{ label, path string }{
{"yt-dlp", cfg.YTDLPPath},
Minternal/config/config.go
@@ -49,9 +49,8 @@ func New() *Config {
}
}
// SetupLogging installs the process-wide slog handler. slog.SetDefault also
// redirects the stdlib log package through it, so output from dependencies is
// formatted and levelled the same way.
// SetupLogging also redirects the stdlib log package, so output from
// dependencies is formatted and levelled the same way.
func (c *Config) SetupLogging() {
var level slog.Level
if err := level.UnmarshalText([]byte(c.LogLevel)); err != nil {
@@ -78,9 +77,9 @@ func getEnv(key, defaultVal string) string {
return defaultVal
}
// getEnvInt reads an int with a lower bound. A value below min is rejected in
// favour of the default: a zero or negative worker count, tick interval or port
// would otherwise stall downloads, spin the scheduler or fail to bind.
// getEnvInt falls back to the default below min: a zero or negative worker
// count, tick interval or port stalls downloads, spins the scheduler or fails to
// bind.
func getEnvInt(key string, defaultVal, min int) int {
v := os.Getenv(key)
if v == "" {
Minternal/database/database.go
@@ -45,8 +45,8 @@ func New(cfg *config.Config) (*sql.DB, error) {
return db, nil
}
// Checkpoint flushes the write-ahead log into the main database file. Call it
// before closing so a killed process doesn't leave a large uncheckpointed WAL.
// Checkpoint flushes the WAL into the main database file. Call it before closing
// so a killed process does not leave a large uncheckpointed WAL.
func Checkpoint(db *sql.DB) error {
if _, err := db.Exec("PRAGMA wal_checkpoint(TRUNCATE)"); err != nil {
return fmt.Errorf("wal checkpoint: %w", err)
@@ -132,10 +132,9 @@ func migrate(db *sql.DB) error {
FOREIGN KEY (preset_id) REFERENCES presets(id)
)`},
{14, `ALTER TABLE downloads ADD COLUMN subscription_id INTEGER`},
// Rebuild downloads so its foreign keys survive enforcement: deleting a
// preset or subscription now nulls the reference instead of failing, which
// matches how ExecuteDownload already degrades to the default preset.
// subscription_id gains the foreign key it never had.
// Rebuilt so the foreign keys survive enforcement: deleting a preset or
// subscription nulls the reference instead of failing, matching how
// ExecuteDownload degrades to the default preset.
{15, `CREATE TABLE downloads_new (
id INTEGER PRIMARY KEY AUTOINCREMENT,
url TEXT NOT NULL,
@@ -186,9 +185,8 @@ func migrate(db *sql.DB) error {
created_at FROM subscriptions;
DROP TABLE subscriptions;
ALTER TABLE subscriptions_new RENAME TO subscriptions;`},
// Indexes for the hot predicates: the queue poll and claim filter on
// status, the scheduler's dedup check filters on subscription_id, and
// GetDue filters on enabled + next_run_at.
// The hot predicates: queue poll and claim on status, the scheduler's dedup
// check on subscription_id, GetDue on enabled + next_run_at.
{17, `CREATE INDEX IF NOT EXISTS idx_downloads_status ON downloads(status);
CREATE INDEX IF NOT EXISTS idx_downloads_subscription_id ON downloads(subscription_id);
CREATE INDEX IF NOT EXISTS idx_downloads_created_at ON downloads(created_at);
Minternal/handler/auth.go
@@ -18,10 +18,9 @@ import (
"vidarchive/internal/service"
)
// Authentication is mandatory: nothing here has a disabled path, and the server
// refuses to start without credentials (see checkCredentials in main). The
// session is a signed cookie rather than server-side state; sessionKey explains
// what signs it.
// Authentication is mandatory: there is no disabled path, and the server refuses
// to start without credentials. The session is a signed cookie, not server-side
// state.
const sessionCookie = "session"
const sessionTTL = 30 * 24 * time.Hour
@@ -31,10 +30,9 @@ const sessionTTL = 30 * 24 * time.Hour
// roughly two seconds, the most a login can take and still fit in loginWait.
const maxLoginCost = 15
// maxConcurrentLogins caps how many password checks run at once. /login is
// public and bcrypt is deliberately slow, so without this an unauthenticated
// flood pins every core. loginWait bounds how long a request queues for a slot
// rather than parking a goroutine indefinitely.
// /login is public and bcrypt is deliberately slow, so without a cap an
// unauthenticated flood pins every core. loginWait bounds the queueing rather
// than parking a goroutine indefinitely.
const (
maxConcurrentLogins = 2
loginWait = 3 * time.Second
@@ -44,8 +42,8 @@ const (
// being read into memory before the credentials are even looked at.
const maxLoginBody = 4 << 10
// sessionKey derives the cookie signing key. The stored secret is mixed in so
// Logout can rotate it and make every cookie issued so far stop verifying.
// sessionKey mixes in the stored secret so Logout can rotate it and make every
// cookie issued so far stop verifying.
func (h *Handler) sessionKey() []byte {
h.sessionMu.RLock()
secret := h.sessionSecret
@@ -55,8 +53,8 @@ func (h *Handler) sessionKey() []byte {
return sum[:]
}
// signExpiry returns the MAC binding a session to its expiry time. Signing the
// expiry is what stops a client from extending its own session.
// signExpiry binds a session to its expiry, which stops a client from extending
// its own session.
func (h *Handler) signExpiry(exp int64) string {
mac := hmac.New(sha256.New, h.sessionKey())
fmt.Fprintf(mac, "%d", exp)
@@ -107,15 +105,12 @@ func (h *Handler) hasSession(r *http.Request) bool {
return time.Now().Unix() < exp
}
// publicPath reports the routes reachable without a session: the login form
// itself, the assets it needs to render, and the health check a monitor scrapes.
func publicPath(path string) bool {
return path == "/login" || path == "/healthz" || strings.HasPrefix(path, "/static/")
}
// RequireAuth guards every route. It is a single middleware rather than
// per-route wrapping, so a new route is protected by default. Forgetting to
// guard one is the failure mode that matters here.
// RequireAuth is a single middleware rather than per-route wrapping, so a new
// route is protected by default.
func (h *Handler) RequireAuth(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if publicPath(r.URL.Path) || h.hasSession(r) {
@@ -142,7 +137,7 @@ func (h *Handler) Login(w http.ResponseWriter, r *http.Request) {
}
// Take a slot before hashing, so the work is bounded no matter how many
// requests arrive. Giving up after loginWait keeps a flood from queueing.
// requests arrive.
select {
case h.loginSem <- struct{}{}:
defer func() { <-h.loginSem }()
@@ -166,14 +161,14 @@ func (h *Handler) Login(w http.ResponseWriter, r *http.Request) {
http.Redirect(w, r, "/", http.StatusSeeOther)
}
// Logout clears the browser's cookie and rotates the signing secret, so a copy
// of that cookie taken beforehand stops working too. There is one account, so
// invalidating every session is exactly the intent.
// Logout rotates the signing secret, so a copy of the cookie taken beforehand
// stops working too. There is one account, so invalidating every session is the
// intent.
func (h *Handler) Logout(w http.ResponseWriter, r *http.Request) {
h.clearSession(w)
if err := h.rotateSessionSecret(); err != nil {
// This browser is signed out either way, since its cookie is gone. Only a
// copy taken elsewhere survives, which is worth a loud log.
// This browser is signed out either way; only a copy taken elsewhere
// survives.
slog.Error("failed to rotate the session secret; cookies issued earlier stay valid", "err", err)
}
http.Redirect(w, r, "/login", http.StatusSeeOther)
@@ -201,8 +196,8 @@ func newSessionSecret() (string, error) {
return hex.EncodeToString(b[:]), nil
}
// loadSessionSecret returns the stored signing secret, creating one on first
// run. Persisting it is what lets sessions survive a restart.
// loadSessionSecret creates a secret on first run. Persisting it lets sessions
// survive a restart.
func loadSessionSecret(settingsSvc *service.SettingsService) (string, error) {
secret, err := settingsSvc.GetSessionSecret()
if err != nil {
@@ -217,10 +212,8 @@ func loadSessionSecret(settingsSvc *service.SettingsService) (string, error) {
return secret, settingsSvc.SetSessionSecret(secret)
}
// ValidatePasswordHash reports whether VIDARCHIVE_PASSWORD_HASH is a bcrypt
// hash this server can actually use. Startup calls it, so a malformed or
// absurdly expensive hash fails there instead of turning into a login that
// mysteriously fails or never returns.
// ValidatePasswordHash runs at startup too, so a malformed or absurdly expensive
// hash fails there instead of turning into a login that never returns.
func ValidatePasswordHash(encoded string) error {
cost, err := bcrypt.Cost([]byte(encoded))
if err != nil {
@@ -232,12 +225,9 @@ func ValidatePasswordHash(encoded string) error {
return nil
}
// VerifyPassword reports whether password matches the stored bcrypt hash. The
// comparison is constant time.
//
// The hash is re-validated first. Startup already did that, but this check is
// what bounds the work: deriving against a cost-31 hash would hold a login slot
// for hours, and reading the cost header is free by comparison.
// VerifyPassword compares in constant time. Re-validating the hash is what
// bounds the work: deriving against a cost-31 hash would hold a login slot for
// hours, and reading the cost header is free by comparison.
func VerifyPassword(encoded, password string) bool {
if ValidatePasswordHash(encoded) != nil {
return false
Minternal/handler/handler.go
@@ -31,10 +31,9 @@ type Handler struct {
subscriptionSvc *service.SubscriptionService
workerPool *worker.Pool
// loginSem caps concurrent password checks; see maxConcurrentLogins.
loginSem chan struct{}
// sessionSecret is the stored cookie signing secret, cached here because
// every guarded request reads it. Logout rotates it.
// sessionSecret is cached here because every guarded request reads it; Logout
// rotates it.
sessionMu sync.RWMutex
sessionSecret string
tools toolCache
@@ -65,7 +64,6 @@ func New(cfg *config.Config, presetSvc *service.PresetService, downloadSvc *serv
}, nil
}
// loadTemplates builds the template set.
func loadTemplates(presetSvc *service.PresetService) (*template.Template, error) {
tmpl := template.New("").Funcs(template.FuncMap{
"formatDuration": formatDuration,
@@ -79,16 +77,14 @@ func loadTemplates(presetSvc *service.PresetService) (*template.Template, error)
"presetFlags": func(p *models.Preset) string { return presetSvc.EffectiveFlags(p, "", "") },
"urlEncodePath": util.URLEncodePath,
"sub": func(a, b int) int { return a - b },
// emptyPreset / newSubscription supply a zero value so the shared create and
// edit form partials can be rendered from the create page too. newSubscription
// carries the create-time defaults (overwrite mode, daily schedule).
// A zero value lets the create page render the form partials it shares with
// the edit page.
"emptyPreset": func() *models.Preset { return nil },
"newSubscription": func() *models.Subscription {
return &models.Subscription{RefreshMode: "overwrite", ScheduleKind: "daily"}
},
// dict builds a map from alternating key/value args, so a template can pass
// more than one value into a sub-template (e.g. the subscription form needs
// both the subscription and the preset list).
// dict takes alternating key/value args, so a template can pass more than
// one value into a sub-template.
"dict": func(values ...interface{}) (map[string]interface{}, error) {
if len(values)%2 != 0 {
return nil, fmt.Errorf("dict expects an even number of arguments")
@@ -103,9 +99,8 @@ func loadTemplates(presetSvc *service.PresetService) (*template.Template, error)
}
return m, nil
},
// isLongText reports whether text spans more than ~2 lines, so the detail
// view can make long descriptions collapsible. Uses rune count (not bytes)
// so CJK text isn't flagged early.
// isLongText drives the collapsible description. Rune count, not bytes, so
// CJK text is not flagged early.
"isLongText": func(s string) bool {
return strings.Count(s, "\n") >= 2 || len([]rune(s)) > 180
},
@@ -140,8 +135,8 @@ type PageData struct {
Flash *Flash
}
// Flash is a one-shot message shown to the user after a redirect (the
// Post/Redirect/Get pattern). Kind is "success" or "error".
// Flash is a one-shot message shown after a redirect. Kind is "success" or
// "error".
type Flash struct {
Kind string
Message string
@@ -151,9 +146,7 @@ const cookieMaxAge = 365 * 24 * 60 * 60
const flashCookie = "flash"
// setFlash stashes a one-shot message in a short-lived cookie. The next rendered
// page reads and clears it (see consumeFlash), so the message appears once after
// the redirect and never again.
// setFlash stashes the message in a short-lived cookie; consumeFlash clears it.
func setFlash(w http.ResponseWriter, kind, message string) {
http.SetCookie(w, &http.Cookie{
Name: flashCookie,
@@ -168,8 +161,6 @@ func setFlash(w http.ResponseWriter, kind, message string) {
func flashSuccess(w http.ResponseWriter, message string) { setFlash(w, "success", message) }
func flashError(w http.ResponseWriter, message string) { setFlash(w, "error", message) }
// consumeFlash reads the flash cookie (if any) and immediately expires it, so a
// message is shown exactly once.
func consumeFlash(w http.ResponseWriter, r *http.Request) *Flash {
c, err := r.Cookie(flashCookie)
if err != nil || c.Value == "" {
@@ -191,8 +182,7 @@ func consumeFlash(w http.ResponseWriter, r *http.Request) *Flash {
if !ok {
return nil
}
// The cookie is client-editable, so don't let an arbitrary kind flow into the
// banner's class name — clamp it to the two we render.
// The cookie is client-editable, and kind becomes the banner's class name.
if kind != "success" {
kind = "error"
}
@@ -247,8 +237,8 @@ func (h *Handler) renderWithRequest(w http.ResponseWriter, r *http.Request, cont
}
}
// settingsOrDefault loads settings for page rendering, falling back to sane
// defaults (rather than a nil deref) if the store can't be read.
// settingsOrDefault falls back to defaults rather than a nil deref when the
// store can't be read.
func (h *Handler) settingsOrDefault() *models.Settings {
settings, err := h.settingsSvc.GetAll()
if err != nil {
Minternal/handler/health.go
@@ -13,24 +13,20 @@ import (
"vidarchive/internal/util"
)
// healthProbeTimeout caps the external version lookups so the endpoint always
// answers promptly, even when a tool is wedged.
// healthProbeTimeout keeps the endpoint answering promptly with a wedged tool.
const healthProbeTimeout = 3 * time.Second
// toolCacheTTL bounds how stale the reported versions may be. /healthz is
// public and each lookup forks a process, so without a cache anyone who can
// reach the port costs the host three processes per request.
// toolCacheTTL: /healthz is public and each lookup forks a process, so without a
// cache anyone who reaches the port costs the host three processes per request.
const toolCacheTTL = time.Minute
// toolRetryTTL is the shorter TTL used when a probe timed out. A timeout says
// less about the tool than a real answer, so it is worth re-probing sooner. It
// is still cached, though: not caching it hands an unauthenticated caller three
// forks per request for as long as the tool stays wedged.
// toolRetryTTL applies when a probe timed out, which says less about the tool
// than a real answer. Still cached, or a wedged tool hands an unauthenticated
// caller three forks per request.
const toolRetryTTL = 10 * time.Second
// toolCache holds the last tool version probe. The mutex is held across the
// probe itself, so a burst of requests collapses into one set of forks. ttl
// varies per result; see toolRetryTTL.
// toolCache holds the last version probe. The mutex is held across the probe
// itself, so a burst of requests collapses into one set of forks.
type toolCache struct {
mu sync.Mutex
at time.Time
@@ -43,22 +39,20 @@ type healthResponse struct {
Dependencies map[string]string `json:"dependencies"`
}
// Health reports whether the app can reach its database and external tools. It
// returns 503 when the database is unreachable, so a container healthcheck or
// orchestrator can act on it. A missing yt-dlp/ffmpeg is reported but does not
// fail the check: the UI still works, only downloads are affected.
// Health returns 503 only when the database is unreachable, so a container
// healthcheck can act on it. A missing yt-dlp or ffmpeg is reported but passes:
// the UI still works, only downloads are affected.
func (h *Handler) Health(w http.ResponseWriter, r *http.Request) {
// Bound the whole probe: a health check that can hang is worse than useless,
// since an orchestrator reads a stuck request as "still starting".
// An orchestrator reads a stuck request as "still starting", so the whole
// probe is bounded.
ctx, cancel := context.WithTimeout(r.Context(), healthProbeTimeout)
defer cancel()
resp := healthResponse{Status: "ok", Dependencies: make(map[string]string, 4)}
// The database is probed first because it alone decides the status code, and
// it must get the budget before a wedged tool can spend it — otherwise a slow
// yt-dlp reports the database as unreachable and the container is restarted
// for the wrong reason.
// The database goes first because it alone decides the status code, and it
// must get the timeout budget before a wedged tool spends it. Otherwise a slow
// yt-dlp gets the container restarted for a database that was fine.
status := http.StatusOK
if err := h.downloadSvc.Ping(ctx); err != nil {
resp.Status = "unhealthy"
@@ -77,8 +71,6 @@ func (h *Handler) Health(w http.ResponseWriter, r *http.Request) {
_ = json.NewEncoder(w).Encode(resp)
}
// toolVersions returns the external tools' versions, re-probing them at most
// once per toolCacheTTL.
func (h *Handler) toolVersions(ctx context.Context) map[string]string {
h.tools.mu.Lock()
defer h.tools.mu.Unlock()
@@ -100,9 +92,8 @@ func (h *Handler) toolVersions(ctx context.Context) map[string]string {
return results
}
// toolVersion returns the first line of the tool's version output. It reports
// "not installed" when the binary is missing and "timed out" when it doesn't
// answer within ctx.
// toolVersion reports "not installed" for a missing binary and "timed out" when
// the tool does not answer within ctx.
func toolVersion(ctx context.Context, path, versionArg string) string {
out, err := util.KillableCommand(ctx, path, versionArg).Output()
if err != nil {
Minternal/handler/helpers.go
@@ -10,8 +10,7 @@ import (
"github.com/go-chi/chi/v5"
)
// parseID reads the {id} route parameter. It writes a 400 and reports false when
// the value isn't a valid id, so callers can simply return.
// parseID writes the 400 itself and reports false, so callers can just return.
func parseID(w http.ResponseWriter, r *http.Request) (int64, bool) {
id, err := strconv.ParseInt(chi.URLParam(r, "id"), 10, 64)
if err != nil {
@@ -21,9 +20,8 @@ func parseID(w http.ResponseWriter, r *http.Request) (int64, bool) {
return id, true
}
// sortFromRequest returns the sort order for a listing page. An explicit ?sort=
// is remembered in a cookie; without one the last choice is restored, so the
// order survives the auto-refresh that reloads these pages.
// sortFromRequest remembers an explicit ?sort= in a cookie and restores it
// otherwise, so the order survives the auto-refresh that reloads these pages.
func sortFromRequest(w http.ResponseWriter, r *http.Request, cookie, defaultSort string) string {
sortBy := r.URL.Query().Get("sort")
if sortBy == "" {
@@ -33,17 +31,14 @@ func sortFromRequest(w http.ResponseWriter, r *http.Request, cookie, defaultSort
return sortBy
}
// pageSize is how many items a listing page renders. The library page fires one
// thumbnail request per media file, so an unpaginated page of a large archive is
// both a slow render and a request storm.
// pageSize: the library page fires one thumbnail request per media file, so an
// unpaginated large archive is both a slow render and a request storm.
const pageSize = 50
// Pagination carries what a listing template needs for its prev/next links.
// Query holds the request's other parameters (filter, sort, path) with a
// trailing '&', so a link keeps them. It is template.URL because
// url.Values.Encode already percent-encoded it: as a plain string the template
// escaper would treat the whole query as one value and re-encode its '&' and
// '='.
// trailing '&', so a link keeps them. It is template.URL because Encode already
// percent-encoded it: as a plain string the template escaper would treat the
// whole query as one value and re-encode its '&' and '='.
type Pagination struct {
Page int
Prev int
@@ -52,9 +47,9 @@ type Pagination struct {
Query template.URL
}
// pageFromRequest reads the 1-based ?page=. Anything unparseable, below 1, or
// large enough that page*pageSize would overflow is page 1: a hand-edited URL
// must not produce a negative offset, which slices a listing out of range.
// pageFromRequest reads the 1-based ?page=. Unparseable, below 1, or large
// enough to overflow page*pageSize all mean page 1: a hand-edited URL must not
// produce a negative offset.
func pageFromRequest(r *http.Request) int {
p, err := strconv.Atoi(r.URL.Query().Get("page"))
if err != nil || p < 1 || p > math.MaxInt/pageSize {
@@ -63,8 +58,7 @@ func pageFromRequest(r *http.Request) int {
return p
}
// newPagination describes the page links for a request. hasNext comes from the
// caller because each listing determines it differently.
// hasNext comes from the caller: each listing determines it differently.
func newPagination(r *http.Request, page int, hasNext bool) Pagination {
q := r.URL.Query()
q.Del("page")
@@ -81,17 +75,15 @@ func newPagination(r *http.Request, page int, hasNext bool) Pagination {
}
}
// serverError logs the underlying failure and shows the user a generic message.
// Internal error strings can carry filesystem paths and SQL text, so they are
// kept out of the response.
// serverError keeps the error out of the response: internal error strings can
// carry filesystem paths and SQL text.
func (h *Handler) serverError(w http.ResponseWriter, r *http.Request, context string, err error) {
slog.Error(context, "method", r.Method, "path", r.URL.Path, "err", err)
http.Error(w, "Something went wrong. Please try again.", http.StatusInternalServerError)
}
// redirectWithError flashes a message and redirects, the standard
// Post/Redirect/Get failure path for form submissions. The underlying error is
// logged rather than shown.
// redirectWithError is the Post/Redirect/Get failure path for forms. The
// underlying error is logged rather than shown.
func redirectWithError(w http.ResponseWriter, r *http.Request, path, message string, err error) {
if err != nil {
slog.Error(message, "method", r.Method, "path", r.URL.Path, "err", err)
@@ -100,15 +92,13 @@ func redirectWithError(w http.ResponseWriter, r *http.Request, path, message str
http.Redirect(w, r, path, http.StatusSeeOther)
}
// redirectWithSuccess flashes a confirmation and redirects.
func redirectWithSuccess(w http.ResponseWriter, r *http.Request, path, message string) {
flashSuccess(w, message)
http.Redirect(w, r, path, http.StatusSeeOther)
}
// parseForm parses a form body and reports whether it succeeded. On failure it
// writes a 400 with a generic message: ParseForm's own error names the offending
// bytes, which is noise to the user and detail we don't need to hand out.
// parseForm writes a generic 400 itself: ParseForm's own error names the
// offending bytes, which is noise to the user and detail not worth handing out.
func parseForm(w http.ResponseWriter, r *http.Request) bool {
if err := r.ParseForm(); err != nil {
http.Error(w, "Invalid form", http.StatusBadRequest)
Minternal/handler/library.go
@@ -25,8 +25,7 @@ func (h *Handler) Library(w http.ResponseWriter, r *http.Request) {
}
total := len(items)
// Only the items are paginated. Folders are few and act as navigation, so
// they stay on every page.
// Folders are navigation, not content, so they stay on every page.
page := pageFromRequest(r)
start := min((page-1)*pageSize, len(items))
end := min(start+pageSize, len(items))
@@ -58,9 +57,8 @@ func (h *Handler) Library(w http.ResponseWriter, r *http.Request) {
}
func normalizeRelPath(r *http.Request) string {
// chi gives the raw, still-encoded wildcard. Decode it as a URL path, where
// '+' is a literal plus (only query strings treat '+' as space) — so an item
// directory named "a+b" round-trips correctly.
// chi hands over the still-encoded wildcard. Decoding it as a path, not a
// query, keeps '+' literal, so an item directory named "a+b" round-trips.
relPath := chi.URLParam(r, "*")
relPath = strings.Trim(relPath, "/")
if decoded, err := url.PathUnescape(relPath); err == nil {
@@ -168,8 +166,7 @@ func (h *Handler) ServeMediaItem(w http.ResponseWriter, r *http.Request) {
h.serveThumbnail(strings.TrimSuffix(relPath, "/thumbnail"), w, r)
return
}
// Subtitle tracks are addressed as <item>/subtitles/<lang>, with the language
// as a trailing path segment (see LibraryService.GetSubtitles).
// Subtitle tracks are addressed as <item>/subtitles/<lang>.
if i := strings.LastIndex(relPath, "/subtitles/"); i >= 0 {
item := relPath[:i]
lang := relPath[i+len("/subtitles/"):]
@@ -203,8 +200,7 @@ func (h *Handler) serveThumbnail(relPath string, w http.ResponseWriter, r *http.
return
}
// Fall back to an icon, matched to the requested file's type (or the item's
// primary file when no specific file was requested).
// Fall back to an icon matching the requested file's type.
icon := "video-icon.svg"
if item, err := h.librarySvc.GetByRelPath(relPath); err == nil {
if isAudioFile(item, filename) {
@@ -220,8 +216,7 @@ func (h *Handler) serveThumbnail(relPath string, w http.ResponseWriter, r *http.
w.Write(data)
}
// isAudioFile reports whether the named file (or, if unnamed, the first media
// file) of an item is audio.
// isAudioFile checks the named file, or the first media file when unnamed.
func isAudioFile(item *models.LibraryItem, filename string) bool {
if filename != "" {
for _, mf := range item.MediaFiles {
Minternal/handler/queue.go
@@ -13,8 +13,7 @@ func (h *Handler) Downloads(w http.ResponseWriter, r *http.Request) {
status := r.URL.Query().Get("status")
sortBy := sortFromRequest(w, r, "queue_sort", "date")
// One row beyond the page tells us whether a next page exists without a
// second COUNT query.
// One row beyond the page answers "is there a next one" without a COUNT query.
page := pageFromRequest(r)
downloads, err := h.downloadSvc.GetAll(status, sortBy, pageSize+1, (page-1)*pageSize)
if err != nil {
@@ -47,9 +46,8 @@ func (h *Handler) Downloads(w http.ResponseWriter, r *http.Request) {
})
}
// backToQueue returns the queue page the action was triggered from, so a
// redirect keeps the filter, sort and page the user was looking at instead of
// dropping them on an unfiltered page 1.
// backToQueue keeps the filter, sort and page the action was triggered from,
// instead of dropping the user on an unfiltered page 1.
func backToQueue(r *http.Request) string {
ref := localReferer(r)
if strings.HasPrefix(ref, "/queue") {
@@ -154,16 +152,15 @@ func (h *Handler) CancelDownload(w http.ResponseWriter, r *http.Request) {
redirectWithSuccess(w, r, backToQueue(r), "Download cancelled.")
}
// clearable maps the queue page's bulk-clear buttons to the status each one
// removes. Only these are accepted, so a hand-made form can't delete rows the
// page never offers to clear (in particular 'downloading').
// clearable is the accepted set, so a hand-made form cannot delete rows the page
// never offers to clear, 'downloading' in particular.
var clearable = map[string]string{
"completed": "Completed downloads cleared.",
"error": "Failed downloads cleared.",
}
// ClearDownloads removes finished downloads: all of them, or just one status
// when the form names one.
// ClearDownloads removes every finished download, or one status when the form
// names one.
func (h *Handler) ClearDownloads(w http.ResponseWriter, r *http.Request) {
status := r.FormValue("status")
if status == "" {
@@ -194,8 +191,7 @@ func (h *Handler) DownloadForm(w http.ResponseWriter, r *http.Request) {
return
}
// A missing default preset is normal (the user may not have set one); the
// template handles a nil DefaultPreset, so this isn't surfaced as an error.
// A missing default preset is normal, and the template handles a nil one.
defaultPreset, _ := h.presetSvc.GetDefault()
url := r.URL.Query().Get("url")
Minternal/handler/settings.go
@@ -56,16 +56,16 @@ func (h *Handler) CreatePreset(w http.ResponseWriter, r *http.Request) {
redirectWithSuccess(w, r, "/settings", "Preset created.")
}
// applyPresetForm copies the preset form fields onto p and validates them. It is
// shared by create and update so the two can't drift apart as fields are added.
// applyPresetForm is shared by create and update, so the two cannot drift apart
// as fields are added.
func applyPresetForm(p *models.Preset, r *http.Request) error {
name := strings.TrimSpace(r.FormValue("name"))
if name == "" {
return fmt.Errorf("A preset needs a name.")
}
// Mirrors the radio options on the settings form; empty means "unspecified"
// and BuildArgs applies its own default.
// Mirrors the settings form's radio options. Empty means unspecified, and
// BuildArgs applies its own default.
formatMode := r.FormValue("format_mode")
switch formatMode {
case "", "default", "preset", "custom":
@@ -99,8 +99,8 @@ func applyPresetForm(p *models.Preset, r *http.Request) error {
}
p.MaxComments = maxComments
}
// Comments are stored in the info JSON sidecar. Keep the dependent options
// consistent even when a client submits the form without JavaScript.
// Comments live in the info JSON sidecar, so the dependent options must stay
// consistent even for a client that submits the form without JavaScript.
if !p.WriteInfoJSON || !p.WriteComments {
p.WriteComments = false
p.CommentSort = ""
@@ -197,17 +197,15 @@ func (h *Handler) Theme(w http.ResponseWriter, r *http.Request) {
http.Redirect(w, r, localReferer(r), http.StatusSeeOther)
}
// localReferer returns where to send the user back to. Only the path is kept:
// Referer is attacker-controlled, so honouring its host would make this route
// an open redirect.
// localReferer keeps only the path: Referer is attacker-controlled, so honouring
// its host would make this route an open redirect.
func localReferer(r *http.Request) string {
ref, err := url.Parse(r.Header.Get("Referer"))
if err != nil || !strings.HasPrefix(ref.Path, "/") {
return "/"
}
// "//host" and "/\host" are protocol-relative URLs to browsers, so a path
// starting with them would redirect off-site. Only a single leading slash
// followed by a normal path segment stays local.
// Browsers read "//host" and "/\host" as protocol-relative URLs, so such a
// path still redirects off-site.
if strings.HasPrefix(ref.Path, "//") || strings.HasPrefix(ref.Path, `/\`) {
return "/"
}
Minternal/handler/subscription.go
@@ -37,9 +37,9 @@ func (h *Handler) Subscriptions(w http.ResponseWriter, r *http.Request) {
})
}
// subscriptionFromForm builds and validates a Subscription from form values,
// shared by create and update. It resolves/validates the schedule into a canonical
// cron expression and requires an output directory (the subscription owns it).
// subscriptionFromForm is shared by create and update. It canonicalizes the
// schedule into a cron expression and requires an output directory, which the
// subscription then owns.
func (h *Handler) subscriptionFromForm(r *http.Request) (*models.Subscription, error) {
name := strings.TrimSpace(r.FormValue("name"))
url := strings.TrimSpace(r.FormValue("url"))
@@ -132,7 +132,7 @@ func (h *Handler) UpdateSubscription(w http.ResponseWriter, r *http.Request) {
}
sub.ID = id
sub.Enabled = existing.Enabled
// Recompute the next run from the (possibly changed) schedule.
// The schedule may have changed.
if next, err := h.subscriptionSvc.ComputeNextRun(sub, time.Now()); err == nil {
sub.NextRunAt = sql.NullTime{Time: next, Valid: true}
}
@@ -176,9 +176,8 @@ func (h *Handler) RunSubscription(w http.ResponseWriter, r *http.Request) {
return
}
// Refuse a second run while one is still in flight, the same guard the
// scheduler applies. Two runs of one subscription share an output directory,
// so in overwrite mode they race: one deletes the item the other just wrote.
// The same guard the scheduler applies: two runs share an output directory, so
// in overwrite mode one deletes the item the other just wrote.
active, err := h.downloadSvc.HasActiveForSubscription(id)
if err != nil {
redirectWithError(w, r, "/subscriptions", "Couldn't start this subscription run.", err)
@@ -196,8 +195,8 @@ func (h *Handler) RunSubscription(w http.ResponseWriter, r *http.Request) {
}
h.workerPool.Submit(download)
// Record the manual run so the subscriptions page shows it. The schedule's
// next run time is deliberately left alone.
// The manual run shows on the subscriptions page, but must not move the
// schedule's next run time.
if err := h.subscriptionSvc.MarkManualRun(id, time.Now(), "queued"); err != nil {
slog.Error("failed to record manual run", "subscription_id", id, "err", err)
}
Minternal/models/models.go
@@ -63,9 +63,8 @@ type Download struct {
CreatedAt time.Time
}
// Subscription is a saved URL that the scheduler re-downloads on a recurring
// schedule. Each subscription owns its OutputDir: refresh/dedup and pruning are
// scoped to that directory (see service.DownloadService refresh modes).
// Subscription is a saved URL the scheduler re-downloads. Each one owns its
// OutputDir: dedup and pruning are scoped to that directory.
type Subscription struct {
ID int64
Name string
Minternal/repository/download.go
@@ -43,8 +43,8 @@ func (r *DownloadRepository) GetByID(id int64) (*models.Download, error) {
return scanDownload(row)
}
// GetAll lists one page of downloads. created_at is indexed, so the ORDER BY
// doesn't sort the whole table to serve a page.
// GetAll lists one page. created_at is indexed, so the ORDER BY does not sort
// the whole table to serve it.
func (r *DownloadRepository) GetAll(status, sortBy string, limit, offset int) ([]*models.Download, error) {
query := `SELECT ` + downloadColumns + ` FROM downloads WHERE 1=1`
var args []interface{}
@@ -79,7 +79,6 @@ func (r *DownloadRepository) Ping(ctx context.Context) error {
return r.db.QueryRowContext(ctx, `SELECT 1`).Scan(&one)
}
// IDsByStatus returns the ids of downloads in the given status.
func (r *DownloadRepository) IDsByStatus(status string) ([]int64, error) {
rows, err := r.db.Query(`SELECT id FROM downloads WHERE status = ?`, status)
if err != nil {
@@ -124,9 +123,8 @@ func (r *DownloadRepository) AppendLogs(id int64, logs string) error {
return err
}
// MarkStarted atomically transitions a download from 'queued' to 'downloading'.
// It reports whether this call actually claimed it: false means another worker
// already started it, so the caller must not process it again.
// MarkStarted claims a download atomically. False means another worker got
// there first, so the caller must not process it again.
func (r *DownloadRepository) MarkStarted(id int64) (bool, error) {
return r.affected(
`UPDATE downloads SET status = 'downloading', started_at = CURRENT_TIMESTAMP WHERE id = ? AND status = 'queued'`,
@@ -134,10 +132,8 @@ func (r *DownloadRepository) MarkStarted(id int64) (bool, error) {
)
}
// Requeue moves a finished-but-unsuccessful download back to 'queued' so the
// queue checker picks it up again. Only 'error' and 'cancelled' rows qualify:
// re-queuing a running or completed one would duplicate work. The bool reports
// whether a row actually changed.
// Requeue accepts only 'error' and 'cancelled' rows: re-queuing a running or
// completed one would duplicate work. The bool reports whether a row changed.
func (r *DownloadRepository) Requeue(id int64) (bool, error) {
return r.affected(
`UPDATE downloads SET status = 'queued', error_message = NULL, logs = NULL,
@@ -147,9 +143,8 @@ func (r *DownloadRepository) Requeue(id int64) (bool, error) {
)
}
// CancelQueued marks a not-yet-started download as cancelled. A running one is
// stopped by cancelling its context instead (see DownloadService.Cancel), which
// is why the status guard matters: it must not overwrite a row a worker owns.
// CancelQueued only touches a not-yet-started download. A running one is stopped
// by cancelling its context, so this must not overwrite a row a worker owns.
func (r *DownloadRepository) CancelQueued(id int64) (bool, error) {
return r.affected(
`UPDATE downloads SET status = 'cancelled', completed_at = CURRENT_TIMESTAMP
@@ -158,7 +153,6 @@ func (r *DownloadRepository) CancelQueued(id int64) (bool, error) {
)
}
// affected runs an UPDATE and reports whether it matched a row.
func (r *DownloadRepository) affected(query string, args ...interface{}) (bool, error) {
res, err := r.db.Exec(query, args...)
if err != nil {
@@ -187,9 +181,8 @@ 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.
// HasActiveForSubscription lets the scheduler avoid stacking a second run on top
// of one that has not finished.
func (r *DownloadRepository) HasActiveForSubscription(subID int64) (bool, error) {
var n int
err := r.db.QueryRow(
@@ -212,8 +205,6 @@ func (r *DownloadRepository) DeleteAll() error {
return err
}
// DeleteByStatus removes every download in one status, backing the queue page's
// "Clear completed" / "Clear failed" actions.
func (r *DownloadRepository) DeleteByStatus(status string) error {
_, err := r.db.Exec(`DELETE FROM downloads WHERE status = ?`, status)
return err
Minternal/repository/preset.go
@@ -54,11 +54,9 @@ func clearDefault(tx *sql.Tx) error {
return err
}
// Save persists p, inserting it when it has no id yet and updating it otherwise.
// When p is the new default it first clears the previous default, doing all of it
// within a single transaction so a mid-sequence failure can't leave the library
// with no default set. Create and Update both route through here, so the
// exclusive-default invariant can't be bypassed.
// Save clears the previous default in the same transaction, so a mid-sequence
// failure cannot leave the library with no default set. Create and Update both
// route through here, so the exclusive-default invariant can't be bypassed.
func (r *PresetRepository) Save(p *models.Preset) error {
tx, err := r.db.Begin()
if err != nil {
Minternal/repository/subscription.go
@@ -55,8 +55,7 @@ func (r *SubscriptionRepository) GetAll() ([]*models.Subscription, error) {
return scanSubscriptions(rows)
}
// GetDue returns enabled subscriptions whose next run is at or before now (or
// has never run). The scheduler uses this to decide what to enqueue.
// GetDue returns enabled subscriptions due at or before now, or never run.
func (r *SubscriptionRepository) GetDue(now time.Time) ([]*models.Subscription, error) {
rows, err := r.db.Query(
`SELECT `+subscriptionColumns+`
@@ -90,8 +89,6 @@ func (r *SubscriptionRepository) SetEnabled(id int64, enabled bool) error {
return err
}
// MarkRun records that a scheduled run was started, together with the next run
// time the scheduler computed for it.
func (r *SubscriptionRepository) MarkRun(id int64, lastRunAt, nextRunAt time.Time, status string) error {
_, err := r.db.Exec(
`UPDATE subscriptions SET last_run_at = ?, next_run_at = ?, last_status = ? WHERE id = ?`,
@@ -110,10 +107,8 @@ func (r *SubscriptionRepository) MarkManualRun(id int64, lastRunAt time.Time, st
return err
}
// SetLastStatus records the outcome of a run that has already started, without
// touching the run times. The download worker calls this as the run progresses,
// so the subscription reflects what actually happened rather than staying on the
// status it had when it was queued.
// SetLastStatus leaves the run times alone. The worker calls it as the run
// progresses, so the subscription does not stay on the status it was queued with.
func (r *SubscriptionRepository) SetLastStatus(id int64, status string) error {
_, err := r.db.Exec(`UPDATE subscriptions SET last_status = ? WHERE id = ?`, status, id)
return err
Minternal/server/handler_test.go
@@ -22,7 +22,6 @@ func getWith(router http.Handler, path string, cookies []*http.Cookie) *httptest
return w
}
// flash returns the (kind, message) of a flash cookie set on the response.
func flash(w *httptest.ResponseRecorder) (kind, message string, ok bool) {
for _, c := range w.Result().Cookies() {
if c.Name == "flash" && c.Value != "" {
@@ -105,7 +104,6 @@ func TestPresetUpdateAndDelete(t *testing.T) {
}
id := newestPresetID(t, router)
// Update renames it.
if w := postForm(router, "/settings/presets/"+id, validPresetForm("Renamed")); w.Code != http.StatusSeeOther {
t.Fatalf("update: got %d", w.Code)
}
@@ -113,7 +111,6 @@ func TestPresetUpdateAndDelete(t *testing.T) {
t.Error("updated preset name not reflected")
}
// Delete removes it.
w := postForm(router, "/settings/presets/"+id+"/delete", nil)
if w.Code != http.StatusSeeOther {
t.Fatalf("delete: got %d", w.Code)
@@ -232,7 +229,6 @@ func TestSubscriptionEditFormRenders(t *testing.T) {
}
}
// findCookie returns the value of a cookie set on the response.
func findCookie(w *httptest.ResponseRecorder, name string) (string, bool) {
for _, c := range w.Result().Cookies() {
if c.Name == name {
Minternal/server/queue_test.go
@@ -238,8 +238,8 @@ func TestPaginationRejectsOverflowingPage(t *testing.T) {
createItem(t, cfg.LibraryDir, "item", "Item", map[string]string{"video.mp4": "dummy"})
insertDownload(t, db, "completed", "")
// A page number large enough to overflow page*pageSize used to slice the
// library listing with a negative index, which panicked into a 500.
// A page number large enough to overflow page*pageSize would slice the library
// listing with a negative index and panic into a 500.
for _, page := range []string{"9223372036854775807", "184467440737095518", "99999999999999999999", "-1", "abc"} {
for _, path := range []string{"/library?page=", "/queue?page="} {
if got := getWith(router, path+page, nil); got.Code != http.StatusOK {
Minternal/server/server.go
@@ -43,8 +43,7 @@ func (s *Server) setupRoutes() {
s.router.Use(s.securityHeaders)
s.router.Use(s.handler.RequireAuth)
// Serve embedded static assets (fs.Sub strips the web/static prefix). The
// error is only possible for an invalid constant path, so it can't occur here.
// The error is only possible for an invalid constant path.
staticFS, _ := fs.Sub(vidarchive.StaticFS, "web/static")
s.router.Handle("/static/*", http.StripPrefix("/static/", http.FileServer(http.FS(staticFS))))
@@ -91,8 +90,8 @@ func (s *Server) setupRoutes() {
}
// requestLogger replaces chi's middleware.Logger, which writes its own
// preformatted (and ANSI-coloured) line. This one emits the same facts as slog
// attributes, so a request can be filtered and correlated like anything else.
// preformatted line. These are slog attributes, so a request can be filtered and
// correlated like anything else.
func (s *Server) requestLogger(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
ww := middleware.NewWrapResponseWriter(w, r.ProtoMajor)
@@ -112,8 +111,7 @@ func (s *Server) securityHeaders(next http.Handler) http.Handler {
w.Header().Set("Referrer-Policy", "strict-origin-when-cross-origin")
// The app ships no JavaScript, so the strict policy costs nothing and is
// sent regardless of scheme — plain-HTTP deployments were previously left
// with no CSP at all.
// sent regardless of scheme.
w.Header().Set("Content-Security-Policy", "default-src 'self'; script-src 'none'; style-src 'self' 'unsafe-inline'; media-src 'self' blob:;")
if s.cfg.IsHTTPS() {
@@ -128,9 +126,9 @@ func (s *Server) Router() http.Handler {
return s.router
}
// Start serves until ctx is cancelled, then drains in-flight requests within
// shutdownTimeout. Media streaming rules out a WriteTimeout, but a header
// deadline still bounds a client that connects and never completes a request.
// Start drains in-flight requests within shutdownTimeout once ctx is cancelled.
// Media streaming rules out a WriteTimeout, but the header deadline still bounds
// a client that connects and never completes a request.
func (s *Server) Start(ctx context.Context) error {
srv := &http.Server{
Addr: fmt.Sprintf(":%d", s.cfg.Port),
Minternal/server/server_test.go
@@ -153,8 +153,8 @@ func TestLibraryEngagementViews(t *testing.T) {
t.Fatalf("comments view status/body unexpected: %d %s", all.Code, all.Body.String())
}
// A real item whose path ends in "comments" must remain reachable through
// the item route now that the full comments view has its own namespace.
// The full comments view has its own namespace, so an item whose path ends in
// "comments" must still be reachable through the item route.
createItem(t, cfg.LibraryDir, "playlist/comments", "Nested Comments Item", map[string]string{
"video.mp4": "dummy video",
})
@@ -273,10 +273,9 @@ func TestMediaFileQueryDecoding(t *testing.T) {
}
}
// TestSubtitleServedByPathSegment guards the regression where subtitle tracks
// are linked as <item>/subtitles/<lang> (language as a trailing path segment),
// but the handler only recognized the "/subtitles" suffix with a ?lang= query.
// The path form fell through to media serving and returned 400 "Missing file".
// Subtitle tracks are linked as <item>/subtitles/<lang>, with the language as a
// trailing path segment. Recognizing only the ?lang= query form makes that fall
// through to media serving and return 400 "Missing file".
func TestSubtitleServedByPathSegment(t *testing.T) {
srv, cfg, cleanup := setupTestServer(t)
defer cleanup()
@@ -393,9 +392,9 @@ func TestMultiFileCardThumbnailURLsDecodeToFilenames(t *testing.T) {
srv, cfg, cleanup := setupTestServer(t)
defer cleanup()
// Filenames with spaces are the case that broke: the template must emit a
// query value that the handler decodes back to the exact filename (the bug
// was double-escaping spaces to %2b, which decodes to '+').
// Filenames with spaces are the hard case: the template must emit a query
// value the handler decodes back to the exact filename, not one that
// double-escapes a space to %2b and decodes to '+'.
files := map[string]string{
"01 - Color Bars.mp4": "v",
"02 - Test Pattern.mp4": "v",
@@ -424,14 +423,10 @@ func TestMultiFileCardThumbnailURLsDecodeToFilenames(t *testing.T) {
}
}
// TestNestedFolderLinkRoundTrip guards the double-encoding regression: a folder
// whose name contains a space was linked with urlEncodePath *inside* a ?path=
// query, which html/template then re-escaped (%20 -> %2520). Clicking the link
// landed on a path the server decoded to "playlist%20test%202" — a directory
// that doesn't exist — so the folder rendered empty and the breadcrumb showed
// the literal "%20". The link must round-trip: its decoded ?path must be the
// real directory, the item inside must render, and the breadcrumb must show the
// human-readable name.
// A folder name containing a space is the double-encoding case: urlEncodePath
// inside a ?path= query gets re-escaped by html/template (%20 -> %2520), and the
// link then lands on a directory that does not exist. The decoded ?path must be
// the real directory, and the breadcrumb must show the readable name.
func TestNestedFolderLinkRoundTrip(t *testing.T) {
srv, cfg, cleanup := setupTestServer(t)
defer cleanup()
Minternal/service/download.go
@@ -65,9 +65,8 @@ func (s *DownloadService) Create(url string, presetID *int64, formatOverride, cu
return d, nil
}
// CreateForSubscription queues a download for a subscription run, copying its
// download options and tagging it with the subscription id so ExecuteDownload
// applies the right refresh mode and pruning.
// CreateForSubscription copies the subscription's download options. The
// subscription id is what makes ExecuteDownload apply refresh mode and pruning.
func (s *DownloadService) CreateForSubscription(sub *models.Subscription) (*models.Download, error) {
d := &models.Download{
URL: sub.URL,
@@ -114,8 +113,6 @@ 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)
}
@@ -126,10 +123,9 @@ func (s *DownloadService) Delete(id int64) error {
return s.repo.Delete(id)
}
// registerActive records the cancel func for a claimed download and returns a
// release func. Registration happens at claim time rather than after the process
// spawns, so a delete arriving during setup, between yt-dlp and the import, or
// mid-import still stops the work instead of silently letting it finish.
// registerActive returns a release func. Registration happens at claim time
// rather than after the process spawns, so a delete arriving during setup or
// mid-import still stops the work.
func (s *DownloadService) registerActive(id int64, cancel context.CancelFunc) func() {
s.activeMu.Lock()
s.active[id] = cancel
@@ -156,8 +152,8 @@ func (s *DownloadService) cancelDownload(id int64) bool {
return ok
}
// CancelAll stops every download currently in flight. Used on shutdown and when
// clearing the queue, so no yt-dlp child outlives the rows that described it.
// CancelAll stops every in-flight download, so no yt-dlp child outlives the row
// that described it.
func (s *DownloadService) CancelAll() {
s.activeMu.Lock()
cancels := make([]context.CancelFunc, 0, len(s.active))
@@ -172,8 +168,8 @@ func (s *DownloadService) CancelAll() {
}
}
// Retry puts a failed or cancelled download back in the queue. The worker pool's
// queue checker picks it up on its next tick, so nothing is submitted here.
// Retry re-queues a download. The pool's queue checker picks it up on its next
// tick, so nothing is submitted here.
func (s *DownloadService) Retry(id int64) error {
ok, err := s.repo.Requeue(id)
if err != nil {
@@ -187,33 +183,29 @@ func (s *DownloadService) Retry(id int64) error {
}
// Cancel stops a download without deleting its row, so it stays visible and
// retryable. The queued case is marked here; a running one is stopped by
// cancelling its context, and ExecuteDownload records the 'cancelled' status.
//
// The order matters: marking first means a worker racing to claim the row finds
// it no longer 'queued' and skips it, instead of starting work we just cancelled.
// retryable. Marking the row first means a worker racing to claim it finds it no
// longer 'queued' and skips it.
func (s *DownloadService) Cancel(id int64) error {
wasQueued, err := s.repo.CancelQueued(id)
if err != nil {
return err
}
// Neither queued nor running: the row is already finished, or it is a stale
// 'downloading' row no worker owns. Say so instead of reporting success.
// Neither queued nor running: already finished, or a stale 'downloading' row
// no worker owns.
if !s.cancelDownload(id) && !wasQueued {
return fmt.Errorf("download %d is not running", id)
}
return nil
}
// DeleteByStatus clears one status' worth of rows. Only finished statuses are
// offered in the UI, so nothing it removes needs cancelling first.
// DeleteByStatus needs no cancelling: only finished statuses are offered in the UI.
func (s *DownloadService) DeleteByStatus(status string) error {
return s.repo.DeleteByStatus(status)
}
func (s *DownloadService) DeleteAll() error {
// Clearing the queue must also stop what is running; otherwise yt-dlp keeps
// going and imports into the library after its row is gone.
// Clearing the queue must stop what is running; otherwise yt-dlp imports into
// the library after its row is gone.
s.CancelAll()
return s.repo.DeleteAll()
}
@@ -222,9 +214,9 @@ func (s *DownloadService) Ping(ctx context.Context) error {
return s.repo.Ping(ctx)
}
// ResetStalledDownloads re-queues downloads left mid-flight by a previous run and
// discards their temp directories. Without the cleanup the re-run imports into a
// fresh uniqueDir and the library ends up with a duplicate of the same item.
// ResetStalledDownloads re-queues downloads left mid-flight by a previous run.
// The temp dirs go too: otherwise the re-run imports into a fresh uniqueDir and
// the library ends up with a duplicate.
func (s *DownloadService) ResetStalledDownloads() error {
ids, err := s.repo.IDsByStatus("downloading")
if err != nil {
@@ -242,30 +234,29 @@ func (s *DownloadService) ResetStalledDownloads() error {
return s.repo.UpdateStatusWhere("downloading", "queued")
}
// tempDirFor returns the scratch directory a download writes into.
func (s *DownloadService) tempDirFor(id int64) string {
return filepath.Join(s.cfg.TempDir, strconv.FormatInt(id, 10))
}
// tempNewDirFor returns the second-pass scratch directory used by metadata mode.
// tempNewDirFor is the second-pass scratch directory, used by metadata mode.
func (s *DownloadService) tempNewDirFor(id int64) string {
return s.tempDirFor(id) + "-new"
}
// tempDirsFor returns every scratch directory a download owns. ResetStalledDownloads
// clears these, so the two builders above must stay the only places that name them.
// tempDirsFor must list every scratch directory a download owns:
// ResetStalledDownloads relies on it, so keep the two builders above the only
// places that name one.
func (s *DownloadService) tempDirsFor(id int64) []string {
return []string{s.tempDirFor(id), s.tempNewDirFor(id)}
}
// ExecuteDownload runs the download for d. The bool reports whether this call
// actually processed it: false means another worker already claimed it (Submit
// and the queue checker can both enqueue the same row within the 2s poll window),
// so the caller should not log it as completed. A cancelled run returns
// ErrCancelled.
// processed it: false means another worker already claimed it (Submit and the
// queue checker can both enqueue the same row within the 2s poll window). A
// cancelled run returns ErrCancelled.
//
// parent belongs to the worker pool: deriving from it means a shutdown cancels
// the download even if it lands before this call registers its own cancel func.
// parent belongs to the worker pool, so a shutdown cancels the download even if
// it lands before this call registers its own cancel func.
func (s *DownloadService) ExecuteDownload(parent context.Context, d *models.Download) (bool, error) {
claimed, err := s.repo.MarkStarted(d.ID)
if err != nil {
@@ -275,8 +266,6 @@ func (s *DownloadService) ExecuteDownload(parent context.Context, d *models.Down
return false, nil
}
// Registered before any work starts so Delete/CancelAll can interrupt every
// phase, not just the window where yt-dlp happens to be running.
ctx, cancel := context.WithCancel(parent)
defer cancel()
defer s.registerActive(d.ID, cancel)()
@@ -308,8 +297,8 @@ func (s *DownloadService) ExecuteDownload(parent context.Context, d *models.Down
}
}
// Reject custom flags that clash with options VidArchive sets itself, before
// spending any work — the download fails with a message naming the offender.
// Reject flags that clash with options VidArchive sets itself before spending
// any work.
isSubscription := d.SubscriptionID.Valid
for _, flags := range []string{d.CustomFlags, preset.CustomFlags} {
if err := checkReservedFlags(flags, isSubscription); err != nil {
@@ -318,28 +307,25 @@ func (s *DownloadService) ExecuteDownload(parent context.Context, d *models.Down
}
}
// From here the run is really under way, so a subscription shows "downloading"
// instead of the "queued" the scheduler recorded.
// The run is really under way now; the scheduler recorded "queued".
s.recordSubscriptionStatus(d, "downloading")
tempDownloadDir := s.tempDirFor(d.ID)
if err := os.MkdirAll(tempDownloadDir, 0o755); err != nil {
// MarkStarted already moved the row to "downloading"; returning without
// finalizing would strand it there until the next restart.
// finalizing strands it there until the next restart.
err = fmt.Errorf("create temp download dir: %w", err)
s.finalizeError(d, err)
return false, err
}
// Own the temp dir's lifetime here, where it's created, so it's removed on
// every exit path — including a failed yt-dlp run or an early return that
// crashes mid-import. The import helpers below no longer clean it up.
// The temp dir's lifetime is owned here, not by the import helpers below, so
// every exit path removes it.
defer os.RemoveAll(tempDownloadDir)
args := s.presetSvc.BuildArgs(preset, d.FormatOverride, d.CustomFlags)
// Record the meaningful flags (format/audio/subs/custom) that shaped this
// download, before the internal plumbing (cookies, -P/-o, URL) is appended,
// so each imported item can show how it was fetched.
// Snapshot the flags that shaped this download before the plumbing (cookies,
// -P/-o, URL) is appended, so each imported item can show how it was fetched.
ytdlpFlags := strings.Join(args, " ")
var cookieCleanup func()
@@ -347,20 +333,18 @@ func (s *DownloadService) ExecuteDownload(parent context.Context, d *models.Down
defer cookieCleanup()
if sub != nil {
// Always write info.json so the import step can read the stable identity
// (yt-dlp's video id) used to match/replace existing items.
// info.json carries yt-dlp's video id, which the import uses to match and
// replace existing items.
if !slices.Contains(args, "--write-info-json") {
args = append(args, "--write-info-json")
}
switch sub.RefreshMode {
case "skip":
// Let yt-dlp skip entries already recorded — no re-download.
archive := s.subscriptionSvc.ArchivePath(sub.ID)
if err := os.MkdirAll(filepath.Dir(archive), 0o755); err == nil {
args = append(args, "--download-archive", archive)
}
case "metadata":
// Refresh metadata only; don't fetch media.
args = append(args, "--skip-download")
}
}
@@ -371,9 +355,9 @@ func (s *DownloadService) ExecuteDownload(parent context.Context, d *models.Down
runErr := s.runYTDLP(ctx, d, args)
// Cancellation wins over both the run error and the import: a cancel that
// lands just after yt-dlp exited 0 leaves runErr nil, and the item must not
// reach the library after the user removed it.
// Cancellation wins over runErr and the import: a cancel landing just after
// yt-dlp exits 0 leaves runErr nil, and the item must not reach the library
// after the user removed it.
if ctx.Err() != nil {
return false, s.finalizeCancelled(parent, d)
}
@@ -383,19 +367,15 @@ func (s *DownloadService) ExecuteDownload(parent context.Context, d *models.Down
mode = sub.RefreshMode
}
// Post-process before marking completed, so the download stays "downloading"
// until everything is really done — including metadata mode's second pass,
// which downloads any genuinely new entries as full items.
//
// This runs even when runErr is set: yt-dlp exits non-zero if a single
// playlist entry fails, while the other entries downloaded fine. Skipping the
// import would throw those away with the temp dir. The outcome is decided
// below, once the imported count is known.
// Post-processing runs before the row is marked completed, and runs even when
// runErr is set: yt-dlp exits non-zero if a single playlist entry fails, and
// skipping the import would throw the successful entries away with the temp
// dir. The outcome is decided below, once the imported count is known.
var postErr error
partial := ""
if mode == "metadata" {
// 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.
// info.json files.
postErr = s.refreshAndAddNew(ctx, d, preset, tempDownloadDir, ytdlpFlags)
if postErr == nil && runErr != nil {
partial = fmt.Sprintf("VidArchive: yt-dlp exited with an error (%v); the metadata refresh finished anyway. Check the log above for failed entries.", runErr)
@@ -407,10 +387,8 @@ func (s *DownloadService) ExecuteDownload(parent context.Context, d *models.Down
postErr = err
case imported == 0 && runErr != nil:
postErr = runErr
// 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.
// Subscription modes legitimately import zero (skip mode, or a metadata
// refresh with no new entries); a plain download that yields nothing failed.
case imported == 0 && sub == nil:
postErr = fmt.Errorf("yt-dlp finished but no media files were downloaded")
case runErr != nil:
@@ -422,8 +400,7 @@ func (s *DownloadService) ExecuteDownload(parent context.Context, d *models.Down
s.pruneSubscription(ctx, d, sub)
}
// Same cancellation check as above: a stop during post-processing must not
// be recorded as "completed".
// A stop during post-processing must not be recorded as "completed" either.
if ctx.Err() != nil {
return false, s.finalizeCancelled(parent, d)
}
@@ -432,8 +409,8 @@ func (s *DownloadService) ExecuteDownload(parent context.Context, d *models.Down
return false, postErr
}
// Record the partial failure in the download's own log. The run counts as
// completed, so nothing else would tell the user some entries failed.
// The run counts as completed, so the log is the only place that can tell the
// user some entries failed.
if partial != "" {
s.cache.AppendLog(d.ID, partial)
s.flushLogs(d.ID)
@@ -447,21 +424,16 @@ func (s *DownloadService) ExecuteDownload(parent context.Context, d *models.Down
return true, nil
}
// runYTDLP executes yt-dlp with args, streaming combined output into the live
// progress cache and periodically flushing it to the download's persisted log.
// Cancelling ctx kills the whole process group and makes this return.
// runYTDLP streams yt-dlp's combined output into the progress cache, flushing it
// to the persisted log as it goes. Cancelling ctx kills the whole process group.
func (s *DownloadService) runYTDLP(ctx context.Context, d *models.Download, args []string) error {
// --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.
//
// --no-write-playlist-metafiles suppresses the playlist-level info.json that
// --write-info-json would also produce. It lands in the first item dir and
// would be imported as if it were an item.
//
// --socket-timeout bounds a stalled connection. Without it a dead socket pins
// a worker forever. There is no inactivity killer beyond this.
// --newline: 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.
// --no-write-playlist-metafiles: the playlist-level info.json lands in the
// first item dir and would be imported as if it were an item.
// --socket-timeout: the only bound on a stalled connection. Without it a dead
// socket pins a worker forever.
fullArgs := append([]string{"--newline", "--no-write-playlist-metafiles", "--socket-timeout", "30"}, args...)
cmd := util.KillableCommand(ctx, s.cfg.YTDLPPath, fullArgs...)
@@ -490,8 +462,8 @@ func (s *DownloadService) runYTDLP(ctx context.Context, d *models.Download, args
}()
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.
// A single yt-dlp message can exceed the 64 KiB default, and an aborted
// scanner leaves the pipe unread.
scanner.Buffer(make([]byte, 0, 64*1024), 1024*1024)
for scanner.Scan() {
s.cache.AppendLog(d.ID, scanner.Text())
@@ -506,10 +478,8 @@ func (s *DownloadService) runYTDLP(ctx context.Context, d *models.Download, args
return cmd.Wait()
}
// appendCookies writes the saved cookies (if any) to a temp file and appends a
// --cookies flag. The returned cleanup saves the cookies yt-dlp left behind and
// removes the temp file. It is always safe to call, even when no cookies were
// configured.
// appendCookies writes the saved cookies, if any, to a temp file. The returned
// cleanup saves back what yt-dlp left behind and is always safe to call.
func (s *DownloadService) appendCookies(args []string) ([]string, func()) {
cookies, err := s.settingsSvc.GetCookies()
if err != nil || strings.TrimSpace(cookies) == "" {
@@ -525,16 +495,14 @@ func (s *DownloadService) appendCookies(args []string) ([]string, func()) {
}
}
// saveRefreshedCookies stores back what yt-dlp wrote to the cookie file.
//
// yt-dlp rewrites the jar on exit. YouTube rotates session cookies on use, so
// the snapshot we sent is stale once the run ends. Keeping the old snapshot and
// replaying it later gets the session invalidated, and the user has to export
// cookies again. sent is what we wrote, so an untouched file saves nothing.
// saveRefreshedCookies stores back what yt-dlp wrote to the cookie file. YouTube
// rotates session cookies on use, so replaying the snapshot we sent invalidates
// the session and the user has to export cookies again. sent is that snapshot,
// so an untouched file saves nothing.
func (s *DownloadService) saveRefreshedCookies(path, sent string) {
data, err := os.ReadFile(path)
// A missing file means yt-dlp never got that far. Empty content would wipe
// working cookies, so treat it as nothing to do.
// A missing file means yt-dlp never got that far; empty content would wipe
// working cookies.
if err != nil || strings.TrimSpace(string(data)) == "" || string(data) == sent {
return
}
@@ -543,9 +511,8 @@ func (s *DownloadService) saveRefreshedCookies(path, sent string) {
}
}
// writeCookiesFile writes cookies to a temp file in the app's own temp dir (the
// same volume the rest of the run uses). On any failure the partial file is
// removed — a truncated cookies file must not be handed to yt-dlp.
// writeCookiesFile removes the partial file on any failure: a truncated cookies
// file must not be handed to yt-dlp.
func (s *DownloadService) writeCookiesFile(cookies string) (string, error) {
if err := os.MkdirAll(s.cfg.TempDir, 0o755); err != nil {
return "", err
@@ -574,10 +541,8 @@ func (s *DownloadService) finalizeError(d *models.Download, err error) {
s.recordSubscriptionStatus(d, "error")
}
// recordSubscriptionStatus mirrors a subscription download's state onto the
// subscription row. Without it the row keeps the status it had when the
// scheduler queued it, so the subscriptions page reports "queued" long after the
// run finished — or failed.
// recordSubscriptionStatus mirrors the download's state onto the subscription
// row, which otherwise keeps the "queued" the scheduler wrote.
func (s *DownloadService) recordSubscriptionStatus(d *models.Download, status string) {
if !d.SubscriptionID.Valid || s.subscriptionSvc == nil {
return
@@ -588,23 +553,19 @@ func (s *DownloadService) recordSubscriptionStatus(d *models.Download, status st
}
}
// ErrCancelled reports that a download was deliberately stopped (deleted, queue
// cleared, or shutdown) rather than having failed. Callers distinguish it so a
// cancellation isn't logged as an error.
// ErrCancelled reports a deliberate stop (deleted, queue cleared, or shutdown)
// rather than a failure, so callers don't log it as an error.
var ErrCancelled = errors.New("download cancelled")
// finalizeCancelled records a stopped download.
//
// A shutdown (parent already cancelled) deliberately leaves the row
// "downloading": ResetStalledDownloads re-queues it on the next start, so
// stopping the server resumes the download instead of losing it. Only a
// user-initiated cancel is terminal. The row may already be deleted in that
// case — cancellation usually arrives via Delete — so a missing row is fine.
// finalizeCancelled leaves the row "downloading" on shutdown (parent already
// cancelled), so ResetStalledDownloads resumes it on the next start. Only a
// user-initiated cancel is terminal, and its row is often already deleted, so a
// missing row is fine.
func (s *DownloadService) finalizeCancelled(parent context.Context, d *models.Download) error {
s.flushLogs(d.ID)
// A shutdown leaves the subscription status alone too: the run resumes on the
// next start, so it is still in progress rather than cancelled.
// The subscription status is left alone too: the run resumes, so it is still
// in progress rather than cancelled.
if parent.Err() != nil {
return ErrCancelled
}
@@ -616,7 +577,6 @@ func (s *DownloadService) finalizeCancelled(parent context.Context, d *models.Do
return ErrCancelled
}
// flushLogs persists whatever output has accumulated for a download.
func (s *DownloadService) flushLogs(id int64) {
logs := s.cache.FlushLogs(id)
if logs == "" {
Minternal/service/engagement.go
@@ -12,10 +12,9 @@ import (
"vidarchive/internal/models"
)
// GetEngagement reads optional comments and a playback heatmap from the item's
// yt-dlp info sidecar. These fields are intentionally not copied into the
// marker: the sidecar remains the source of truth and old items simply return
// empty data when the fields are absent.
// GetEngagement reads comments and a playback heatmap from the yt-dlp sidecar,
// which stays their source of truth. Nothing is copied into the marker, so an
// old item just returns empty data.
func (s *LibraryService) GetEngagement(relPath string) ([]models.Comment, []models.HeatmapSegment, error) {
item, err := s.GetByRelPath(relPath)
if err != nil {
@@ -120,9 +119,9 @@ func (s *LibraryService) GetEngagement(relPath string) ([]models.Comment, []mode
// .comment-depth-1 through .comment-depth-4.
const maxCommentDepth = 4
// orderComments groups replies beneath their parent while preserving the
// source order among siblings. Depth is capped for presentation so malformed
// or unusually deep reply chains cannot make the UI progressively narrower.
// orderComments groups replies under their parent, keeping the source order
// among siblings. Depth is capped so a deep or malformed reply chain cannot make
// the UI progressively narrower.
func orderComments(comments []models.Comment) []models.Comment {
if len(comments) < 2 {
return comments
Minternal/service/execute_download_test.go
@@ -404,9 +404,8 @@ func TestExecuteDownloadShutdownLeavesRowResumable(t *testing.T) {
}
}
// 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.
// A cancel that lands after yt-dlp exited 0 must not be recorded as completed.
// This one arrives during the import, so the yt-dlp error is nil.
func TestExecuteDownloadCancelAfterSuccessIsNotCompleted(t *testing.T) {
e := newExecEnv(t)
e.fakeYTDLP(t, e.writeItem(t, "item-00001", "My Clip", "abc123"))
Minternal/service/formats.go
@@ -12,14 +12,12 @@ import (
)
// ListFormats returns the formats offered for url. The second result names the
// playlist entry the formats came from, and is empty for a plain video URL: a
// playlist's entries can differ, so the caller must say which one this is.
// playlist entry they came from, empty for a plain video URL: entries can
// differ, so the caller has to say which one this is.
func (s *DownloadService) ListFormats(ctx context.Context, url string) ([]*models.FormatInfo, string, error) {
// 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.
// -I 1 limits a playlist URL to its first entry, which is all the parsing
// below reads anyway.
// -J rather than scraping the human "-F" table, whose columns shift between
// yt-dlp versions. stderr is captured separately so warnings can't corrupt the
// JSON on stdout. -I 1 is all the parsing below reads anyway.
args := []string{"-J", "--no-warnings", "-I", "1"}
args, cleanup := s.appendCookies(args)
defer cleanup()
@@ -36,9 +34,8 @@ func (s *DownloadService) ListFormats(ctx context.Context, url string) ([]*model
return parseFormatJSON(output)
}
// 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.
// ytFormat mirrors the subset of yt-dlp's per-format JSON VidArchive surfaces.
// Numeric fields are pointers so an absent value stays distinct from a real zero.
type ytFormat struct {
FormatID string `json:"format_id"`
Ext string `json:"ext"`
@@ -54,11 +51,9 @@ type ytFormat struct {
FormatNote string `json:"format_note"`
}
// 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 the dump is restricted to one entry (-I 1), so the formats come from that
// entry and the second result names it, so the UI can say the list describes
// the first item rather than the whole playlist.
// parseFormatJSON takes the top-level formats for a single video. A playlist
// dump carries them under its one entry (-I 1) instead, and that entry's title
// becomes the second result.
func parseFormatJSON(data []byte) ([]*models.FormatInfo, string, error) {
var top struct {
Formats []ytFormat `json:"formats"`
@@ -108,7 +103,6 @@ func (f ytFormat) toFormatInfo() *models.FormatInfo {
fi.Channels = strconv.Itoa(*f.AudioChannels)
}
// 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" {
Minternal/service/import.go
@@ -21,17 +21,13 @@ import (
"vidarchive/internal/util"
)
// resolveBaseLibraryDir returns the absolute library directory a download writes
// into, applying the optional per-download OutputDir while rejecting any path
// that escapes the library root.
// resolveBaseLibraryDir applies the optional per-download OutputDir. Both
// branches go through the library guard: it resolves symlinks, which a plain
// prefix check does not, and the cache keys are derived from the resolved root.
func (s *DownloadService) resolveBaseLibraryDir(d *models.Download) (string, error) {
// Both branches go through the library service's guard, so every caller gets
// a symlink-resolved path. Cache keys are derived from that resolved root.
if !d.OutputDir.Valid || d.OutputDir.String == "" {
return s.librarySvc.ResolveWithinLibrary("")
}
// Reuse the library service's guard so both entry points enforce the boundary
// the same way — it resolves symlinks, which a plain prefix check does not.
dir, err := s.librarySvc.ResolveWithinLibrary(d.OutputDir.String)
if err != nil {
return "", fmt.Errorf("invalid output directory: %w", err)
@@ -39,14 +35,10 @@ func (s *DownloadService) resolveBaseLibraryDir(d *models.Download) (string, err
return dir, nil
}
// 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.
//
// A cancel stops the import between items and returns ctx.Err() with the count
// imported so far. Importing a long playlist takes real time (a move plus an
// ffprobe per file), so a deleted download must not keep filling the library.
// importDownloadedItems returns the number of items imported. Per-item failures
// are logged and skipped, so a playlist with a few bad entries still imports the
// rest; a non-nil error means the import couldn't even start. A cancel stops
// between items and returns ctx.Err() with the count imported so far.
func (s *DownloadService) importDownloadedItems(ctx context.Context, d *models.Download, tempDownloadDir, mode, ytdlpFlags string) (int, error) {
entries, err := os.ReadDir(tempDownloadDir)
if err != nil {
@@ -79,8 +71,8 @@ func (s *DownloadService) importDownloadedItems(ctx context.Context, d *models.D
return imported, err
}
if err := s.importItemDir(ctx, d.URL, itemDir, baseLibraryDir, mode, ytdlpFlags); err != nil {
// A cancelled item isn't a bad item: stop instead of logging a warning
// for it and every one that follows.
// A cancelled item is not a bad item, and every one after it would warn
// too.
if ctx.Err() != nil {
return imported, ctx.Err()
}
@@ -90,8 +82,7 @@ func (s *DownloadService) importDownloadedItems(ctx context.Context, d *models.D
imported++
}
// The temp dir (and any leftovers from failed imports) is removed by the
// caller's deferred cleanup, so partial state never leaks even on a crash.
// Leftovers from failed imports go with the caller's deferred temp-dir cleanup.
return imported, nil
}
@@ -136,17 +127,15 @@ func (s *DownloadService) importItemDir(ctx context.Context, url, itemDir, baseL
name := s.deriveItemName(itemDir, info, mediaFiles)
videoID := info.ID
// Last point at which nothing has been written to the library yet: give up
// here on a cancel rather than part-way through, which would leave a folder
// with some of its files and no marker — or, in overwrite mode, delete the
// existing item and not replace it.
// Last point at which nothing has been written to the library yet. Cancelling
// later leaves a folder with no marker, or in overwrite mode deletes the
// existing item without replacing it.
if err := ctx.Err(); err != nil {
return err
}
// 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).
// Removing the old dir avoids a duplicate folder and lets uniqueDir reuse its
// name, or take the new title if it changed upstream.
if mode == "overwrite" && videoID != "" {
if existing, ok := s.librarySvc.FindByVideoID(baseLibraryDir, videoID); ok {
s.librarySvc.evictCachedDir(existing)
@@ -174,9 +163,8 @@ func (s *DownloadService) importItemDir(ctx context.Context, url, itemDir, baseL
}
}
// Probe each media file's duration once, here in the worker (off the request
// path), and cache it in the marker so the library never has to probe while
// serving pages. Files we can't probe simply get no duration.
// Probing once here, in the worker, is what keeps ffprobe off the request
// path. An unprobeable file simply gets no duration.
fileDurations := make(map[string]int)
for _, entry := range mediaFiles {
if d, ok := probeDuration(ctx, s.cfg.FFprobePath, filepath.Join(targetDir, entry.Name())); ok {
@@ -207,9 +195,9 @@ func (s *DownloadService) importItemDir(ctx context.Context, url, itemDir, baseL
return s.librarySvc.writeMetadata(targetDir, metadata)
}
// infoJSON is the subset of yt-dlp's info.json VidArchive reads. ID is the
// stable item identity; within a single subscription's own directory it is
// enough to match items, so the extractor is not needed.
// infoJSON is the subset of yt-dlp's info.json VidArchive reads. ID alone
// identifies an item within one subscription's directory, so the extractor is
// not needed.
type infoJSON struct {
ID string `json:"id"`
// Type is yt-dlp's "_type": "playlist" marks a playlist-level sidecar rather
@@ -220,9 +208,8 @@ type infoJSON struct {
WebpageURL string `json:"webpage_url"`
}
// readInfoJSON parses an info.json. A missing, unreadable or malformed file
// yields a zero-value struct: every caller treats absent fields as "unknown"
// and falls back, so there is nothing to distinguish.
// readInfoJSON yields a zero-value struct for a missing, unreadable or malformed
// file: every caller already falls back on an absent field.
func readInfoJSON(infoJSONPath string) infoJSON {
var info infoJSON
if infoJSONPath == "" {
@@ -239,14 +226,12 @@ func readInfoJSON(infoJSONPath string) infoJSON {
return info
}
// isInfoJSON reports whether a file name is yt-dlp's metadata sidecar. yt-dlp
// writes "<title>.info.json" next to the media, but a bare "info.json" is what
// an already-imported item holds.
// isInfoJSON accepts both forms: yt-dlp writes "<title>.info.json", while an
// already-imported item holds a bare "info.json".
func isInfoJSON(name string) bool {
return name == "info.json" || strings.HasSuffix(name, ".info.json")
}
// findInfoJSON returns the path to an info.json directly inside itemDir, or "".
func findInfoJSON(itemDir string) string {
entries, err := os.ReadDir(itemDir)
if err != nil {
@@ -264,8 +249,8 @@ func findInfoJSON(itemDir string) string {
return ""
}
// deriveItemName names the imported item after its title, falling back to the
// largest media file's base name when there is no usable info.json.
// deriveItemName falls back to the largest media file's base name when there is
// no usable info.json title.
func (s *DownloadService) deriveItemName(itemDir string, info infoJSON, mediaFiles []os.DirEntry) string {
if info.Title != "" {
return sanitizeDirName(info.Title)
@@ -286,10 +271,9 @@ func (s *DownloadService) deriveItemName(itemDir string, info infoJSON, mediaFil
return sanitizeDirName(base)
}
// uniqueDir returns a directory under base that does not exist yet, appending
// "-1", "-2", ... until it finds one. A stat error other than "does not exist"
// is returned rather than treated as "taken": every candidate would fail the
// same way, so the loop would never end.
// uniqueDir appends "-1", "-2", ... until the name is free. A stat error other
// than "does not exist" is returned rather than treated as taken: every
// candidate would fail the same way and the loop would never end.
func (s *DownloadService) uniqueDir(base, name string) (string, error) {
dir := filepath.Join(base, name)
for i := 0; ; i++ {
@@ -307,10 +291,8 @@ func (s *DownloadService) uniqueDir(base, name string) (string, error) {
}
}
// moveFile moves src to dst, falling back to copy-and-delete when the two are on
// different filesystems. The temp and library directories are independently
// configurable, so they can legitimately live on separate mounts — where a plain
// rename fails with EXDEV.
// moveFile falls back to copy-and-delete on EXDEV: the temp and library
// directories are configured separately and can live on different mounts.
func moveFile(src, dst string) error {
if err := os.Rename(src, dst); err == nil {
return nil
@@ -363,9 +345,8 @@ func sanitizeDirName(name string) string {
)
name = replacer.Replace(name)
name = strings.TrimSpace(name)
// A name made only of dots resolves to the parent ("..") or to the target
// directory itself ("."), so joining it would place the item outside the
// library. Titles come from remote metadata, so refuse them here.
// Titles come from remote metadata, and a name of only dots resolves to the
// parent or to the target dir itself, placing the item outside the library.
if strings.Trim(name, ".") == "" {
name = "untitled"
}
@@ -376,8 +357,8 @@ func sanitizeDirName(name string) string {
// uniqueDir's "-N" suffix.
const maxDirNameBytes = 240
// truncateDirName shortens name to maxDirNameBytes without splitting a rune.
// Over the limit, every filesystem call on the name fails with ENAMETOOLONG.
// truncateDirName cuts on a rune boundary. Over the limit, every filesystem
// call on the name fails with ENAMETOOLONG.
func truncateDirName(name string) string {
if len(name) <= maxDirNameBytes {
return name
Minternal/service/import_test.go
@@ -207,7 +207,6 @@ func TestImportDownloadedItemsRejectsOutputTraversal(t *testing.T) {
librarySvc: NewLibraryService(libDir, "ffmpeg", "ffprobe"),
}
// A temp download dir with one item subdir.
tempDir := t.TempDir()
itemDir := filepath.Join(tempDir, "item-00001")
if err := os.MkdirAll(itemDir, 0o755); err != nil {
@@ -339,7 +338,6 @@ func TestImportItemDirOverwriteReplacesExisting(t *testing.T) {
librarySvc: NewLibraryService(libDir, "ffmpeg", "ffprobe"),
}
// First import establishes the item.
first := t.TempDir()
makeTestVideo(t, filepath.Join(first, "raw.mp4"))
if err := os.WriteFile(filepath.Join(first, "clip.info.json"), []byte(`{"id":"vid1","title":"Old Title"}`), 0o644); err != nil {
Minternal/service/library.go
@@ -33,30 +33,25 @@ var audioExts = map[string]struct{}{
}
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.
libraryDir string
ffmpegPath string
ffprobePath string
// thumbLocks holds a per-media-file mutex serializing extraction so two
// callers never write the same temp file at once. Entries are reference
// counted and removed once nobody holds them. Evicting a lock that is still
// held would let a second caller take a fresh one and reintroduce the race.
// thumbLocks serializes extraction per media file, so two callers never write
// the same temp file at once. Entries are reference counted: evicting a lock
// someone still holds would let the next caller take a fresh one and
// reintroduce the race.
thumbMu sync.Mutex
thumbLocks map[string]*refLock
thumbSem chan struct{}
// thumbFailed records media filepaths whose extraction already failed this
// run, so we trust ffmpeg's verdict and don't re-run it on every request.
// Bounded: an evicted path just means one more ffmpeg attempt.
// thumbFailed remembers paths whose extraction already failed this run, so
// ffmpeg is not re-run on every request. An evicted path costs one retry.
thumbFailed *lru[bool]
// scanCache memoizes scanned items for a short TTL so the listing page (which
// fans out one thumbnail request per media file) and quick auto-refreshes
// don't re-parse each item's marker + info.json on every request. Only item
// scans are cached — the directory listing itself is always read fresh, so
// newly added/removed items and subfolders appear immediately.
// scanCache memoizes item scans for a short TTL so the listing page does not
// re-parse every marker + info.json per request. Only scans are cached: the
// directory listing is always read fresh, so added and removed items appear
// immediately.
scanCache *lru[scanCacheEntry]
scanTTL time.Duration
}
@@ -66,12 +61,10 @@ type scanCacheEntry struct {
at time.Time
}
// scanCacheTTL is how long a scanned item is reused before being re-read.
const scanCacheTTL = 10 * time.Second
// maxCachedPaths bounds each per-path cache. Without it both grow with the
// number of distinct files touched over the process lifetime, which is fine for
// a personal archive and not for a large one.
// maxCachedPaths bounds each per-path cache. Unbounded, both grow with every
// distinct file touched over the process lifetime.
const maxCachedPaths = 1024
func NewLibraryService(libraryDir, ffmpegPath, ffprobePath string) *LibraryService {
@@ -109,12 +102,9 @@ func (s *LibraryService) evictCachedScan(relPath string) {
s.scanCache.Delete(relPath)
}
// evictCachedDir drops the cached scan for an absolute item directory.
//
// The key must be relative to the symlink-resolved root, because that is what
// resolveItemDir stored it under. Computing it against the raw libraryDir misses
// whenever that is a symlink, leaving a deleted or replaced item visible until
// the TTL expires.
// evictCachedDir keys off the symlink-resolved root, because that is what
// resolveItemDir stored the entry under. Against the raw libraryDir it misses
// whenever that is a symlink.
func (s *LibraryService) evictCachedDir(itemDir string) {
rel, err := filepath.Rel(s.libraryRoot(), itemDir)
if err != nil {
@@ -129,7 +119,6 @@ type refLock struct {
refs int
}
// lockThumbFile locks the mutex guarding path and returns its release func.
func (s *LibraryService) lockThumbFile(path string) func() {
s.thumbMu.Lock()
l, ok := s.thumbLocks[path]
@@ -153,7 +142,6 @@ func (s *LibraryService) lockThumbFile(path string) func() {
}
}
// scannedItem returns a cached scan if fresh, otherwise scans and caches it.
func (s *LibraryService) scannedItem(itemDir, relPath string) (*models.LibraryItem, error) {
if item, ok := s.getCachedScan(relPath); ok {
return item, nil
@@ -166,14 +154,12 @@ func (s *LibraryService) scannedItem(itemDir, relPath string) (*models.LibraryIt
return item, nil
}
// ResolveWithinLibrary resolves a caller-supplied relative directory against the
// library root and rejects anything that escapes it. The directory need not
// exist yet, so it is safe to use when choosing a download's output location.
// ResolveWithinLibrary rejects anything that escapes the library root. The
// directory need not exist yet, so it also works for a download's output dir.
func (s *LibraryService) ResolveWithinLibrary(relPath string) (string, error) {
return s.resolveItemDir(filepath.Clean(relPath))
}
// libraryRoot returns the symlink-resolved library root.
func (s *LibraryService) libraryRoot() string {
if base, err := filepath.EvalSymlinks(s.libraryDir); err == nil {
return base
@@ -195,8 +181,8 @@ func (s *LibraryService) resolveItemDir(relPath string) (string, error) {
if !strings.HasPrefix(cleanDir, base+string(filepath.Separator)) && cleanDir != base {
return "", fmt.Errorf("invalid path")
}
// Return the cleaned/symlink-resolved path we just validated, so callers do
// I/O on exactly the path that passed the boundary check.
// Return the path that passed the boundary check, so callers do I/O on
// exactly that one.
return cleanDir, nil
}
@@ -273,8 +259,6 @@ func (s *LibraryService) GetAll(path, sortBy, filter string) ([]*models.LibraryI
return items, folders, nil
}
// GetByRelPath returns the item at relPath, reusing a recent cached scan when
// available (see scanCache).
func (s *LibraryService) GetByRelPath(relPath string) (*models.LibraryItem, error) {
relPath = strings.Trim(relPath, "/")
@@ -313,8 +297,8 @@ func (s *LibraryService) scanItem(itemDir, relPath string) (*models.LibraryItem,
}
}
// dirty tracks whether we derived any new metadata worth persisting, so a
// plain listing or detail view doesn't rewrite the marker file on every read.
// Only derived metadata makes this dirty, so a plain listing does not rewrite
// the marker on every read.
dirty := false
if metadata.Name == "" {
@@ -336,9 +320,8 @@ func (s *LibraryService) scanItem(itemDir, relPath string) (*models.LibraryItem,
dirty = backfillString(&metadata.Description, info, "description") || dirty
}
// 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 id (see FindByVideoID / PruneToIDSet).
// Backfilling the video id gives pre-existing items an identity on their next
// scan; subscriptions match and prune by it.
if metadata.VideoID == "" {
dirty = backfillString(&metadata.VideoID, info, "id") || dirty
}
@@ -347,11 +330,9 @@ func (s *LibraryService) scanItem(itemDir, relPath string) (*models.LibraryItem,
metadata.FileDurations = make(map[string]int)
}
// Per-file durations come from the marker's file_durations map (populated at
// import time). For a single-file item we also seed it from info.json's
// duration, which covers the common case without a probe. We deliberately do
// not run ffprobe here — keeping it off the scan/listing path is the point of
// the marker cache. Files without a known duration simply show no badge.
// Durations come from the marker (written at import), plus info.json for a
// single-file item. No ffprobe here: keeping it off the scan path is the point
// of the marker cache. A file without a known duration shows no badge.
for i := range mediaFiles {
mf := &mediaFiles[i]
if d, ok := metadata.FileDurations[mf.Filename]; ok {
@@ -367,8 +348,7 @@ func (s *LibraryService) scanItem(itemDir, relPath string) (*models.LibraryItem,
}
}
// The item-level duration is the sum of known per-file durations, used only
// for the "duration" sort — there is no single "overall" duration shown.
// The item-level duration only feeds the "duration" sort; nothing displays it.
total := 0
for _, mf := range mediaFiles {
if mf.Duration > 0 {
@@ -396,10 +376,9 @@ func (s *LibraryService) scanItem(itemDir, relPath string) (*models.LibraryItem,
return item, nil
}
// backfillString sets *field from the first non-empty string value among
// info's keys, reporting whether it changed anything. An empty value in
// info.json must not count as a change: the marker would be rewritten on every
// scan forever without ever gaining a value.
// backfillString takes the first non-empty value among keys and reports whether
// it changed anything. An empty value must not count as a change: the marker
// would be rewritten on every scan forever without ever gaining a value.
func backfillString(field *string, info map[string]interface{}, keys ...string) bool {
for _, k := range keys {
if v, ok := infoString(info, k); ok && v != "" {
@@ -430,13 +409,11 @@ func (s *LibraryService) writeMetadata(itemDir string, metadata models.ItemMetad
return atomicWrite(filepath.Join(itemDir, itemMarkerName), buf.Bytes(), markerFileMode)
}
// markerFileMode matches the info.json sidecar: both sit in the library next to
// the media and are meant to be readable (and hand-editable) by the operator.
// markerFileMode matches the info.json sidecar: hand-editable by the operator.
const markerFileMode = 0o644
// atomicWrite writes data to path via a uniquely-named temp file in the same
// directory and a rename, so a crash mid-write can't leave a truncated file and
// two concurrent writers never collide on a shared temp path.
// atomicWrite writes via a uniquely-named temp file in the same directory, so a
// crash can't leave a truncated file and two writers can't collide.
func atomicWrite(path string, data []byte, mode os.FileMode) error {
f, err := os.CreateTemp(filepath.Dir(path), ".vidarchive-*.tmp")
if err != nil {
@@ -519,9 +496,8 @@ func infoString(info map[string]interface{}, key string) (string, bool) {
if info == nil {
return "", false
}
// Only accept genuine strings: title/url/description are always strings in
// yt-dlp output, and stringifying an arbitrary JSON value (map, slice) would
// store junk like "map[...]" into the field.
// Only genuine strings: stringifying a map or slice would store junk like
// "map[...]" into the field.
if s, ok := info[key].(string); ok {
return s, true
}
@@ -551,10 +527,9 @@ func mediaFileStem(mf models.MediaFile) string {
return strings.TrimSuffix(filepath.Base(mf.Filename), filepath.Ext(mf.Filename))
}
// primaryMediaFile picks the representative file for an item: the largest video
// file, or — if there are none — the largest file overall. Returns nil for an
// item with no media files. Used for the item-level thumbnail and as the default
// target for metadata, keeping those two consistent.
// primaryMediaFile picks the largest video file, or the largest file overall if
// the item has no video. It backs both the item thumbnail and the default
// metadata target, which keeps those two consistent.
func primaryMediaFile(item *models.LibraryItem) *models.MediaFile {
size := func(mf *models.MediaFile) int64 {
if info, err := os.Stat(mf.Filepath); err == nil {
@@ -569,7 +544,6 @@ func primaryMediaFile(item *models.LibraryItem) *models.MediaFile {
case best == nil:
best = mf
case best.IsAudio && !mf.IsAudio:
// Prefer any video over audio.
best = mf
case best.IsAudio == mf.IsAudio && size(mf) > size(best):
best = mf
@@ -596,25 +570,20 @@ func (s *LibraryService) Delete(relPath string) error {
if err != nil {
return err
}
// resolveItemDir maps ""/"." to the library root and resolves symlinks, so a
// result equal to the root (reachable via a URL-encoded slash, "sub/..", or a
// symlink pointing back at the root) must be refused — deleting it would
// wipe the entire library.
// resolveItemDir maps ""/"." to the root and resolves symlinks, so the root is
// reachable via a URL-encoded slash, "sub/..", or a symlink pointing back at
// it. Deleting that would wipe the entire library.
if itemDir == s.libraryRoot() {
return fmt.Errorf("refusing to delete library root")
}
// Evict the cached scan so the deletion is reflected immediately rather than
// lingering until the TTL expires.
// Evict so the deletion shows immediately instead of at TTL expiry.
s.evictCachedScan(strings.Trim(relPath, "/"))
return os.RemoveAll(itemDir)
}
// FindByVideoID returns the absolute directory of the item under baseDir whose
// 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.
// FindByVideoID scans only baseDir, the folder one subscription owns. Matching
// on the yt-dlp video id alone is safe there: a single source can't collide with
// another extractor's ids.
func (s *LibraryService) FindByVideoID(baseDir, id string) (string, bool) {
if id == "" {
return "", false
@@ -633,10 +602,9 @@ func (s *LibraryService) FindByVideoID(baseDir, id string) (string, bool) {
return match, match != ""
}
// eachItemDir walks the marked library items directly under baseDir, reading
// each one's metadata, and calls fn until it returns false. Entries that aren't
// items, or whose metadata can't be read, are skipped — a single bad item must
// not abort a scan. logLabel names the caller in those skip messages.
// eachItemDir calls fn for each marked item directly under baseDir until fn
// returns false. Unreadable items are skipped: one bad item must not abort a
// scan. logLabel names the caller in those skip messages.
func (s *LibraryService) eachItemDir(baseDir, logLabel string, fn func(itemDir string, meta models.ItemMetadata) bool) error {
entries, err := os.ReadDir(baseDir)
if err != nil {
@@ -664,11 +632,9 @@ func (s *LibraryService) eachItemDir(baseDir, logLabel string, fn func(itemDir s
return nil
}
// PruneToIDSet deletes items directly under baseDir whose identity key is not in
// keep. It is used to mirror a subscription's source: entries removed upstream
// are removed locally. Items without a known identity key are left untouched (we
// never delete something we can't positively identify). Returns the number
// removed.
// PruneToIDSet deletes items under baseDir whose identity key is not in keep,
// mirroring a subscription's source. An item without an identity key is left
// alone: never delete what can't be positively identified.
func (s *LibraryService) PruneToIDSet(baseDir string, keep map[string]bool) (int, error) {
removed := 0
err := s.eachItemDir(baseDir, "PruneToIDSet", func(itemDir string, meta models.ItemMetadata) bool {
Minternal/service/library_test.go
@@ -53,7 +53,6 @@ func makeTestVideo(t *testing.T, path string) {
makeTestVideoSize(t, path, "64x64")
}
// makeTestVideoSize is makeTestVideo with an explicit WxH size.
func makeTestVideoSize(t *testing.T, path, size string) {
t.Helper()
// 3s so the default 1s thumbnail seek lands on a real frame.
Minternal/service/lru.go
@@ -5,11 +5,8 @@ import (
"sync"
)
// lru is a fixed-size map keyed by path, dropping the least recently used entry
// when it is full.
//
// Only use it for values that are safe to lose: an eviction means the work is
// redone, never that correctness changes.
// lru is a fixed-size map keyed by path. Only for values that are safe to lose:
// an eviction means the work is redone, never that correctness changes.
type lru[V any] struct {
mu sync.Mutex
maxSize int
@@ -23,8 +20,7 @@ type lruEntry[V any] struct {
}
func newLRU[V any](maxSize int) *lru[V] {
// A size below 1 would evict each entry as it is stored, turning the cache
// into a silent miss on every lookup.
// Below 1, each entry is evicted as it is stored: a silent miss every time.
return &lru[V]{
maxSize: max(maxSize, 1),
order: list.New(),
Minternal/service/media_probe.go
@@ -14,9 +14,8 @@ import (
"vidarchive/internal/util"
)
// GetMetadata probes media details for the named file within an item. An empty
// filename (or one that doesn't match) falls back to the item's primary media
// file, so the detail view shows metadata for whichever file is selected.
// GetMetadata falls back to the item's primary media file when filename is empty
// or does not match.
func (s *LibraryService) GetMetadata(relPath, filename string) (*MediaMetadata, error) {
item, err := s.GetByRelPath(relPath)
if err != nil {
@@ -44,14 +43,12 @@ func (s *LibraryService) GetMetadata(relPath, filename string) (*MediaMetadata,
const probeTimeout = 30 * time.Second
// probeStreams runs ffprobe once and returns the container format and every
// stream in it. It is the single ffprobe entry point: callers that only care
// about one stream kind filter the result themselves, rather than each
// re-declaring the same command and JSON shape.
// probeStreams is the single ffprobe entry point. A caller that wants one stream
// kind filters the result itself, rather than re-declaring the same command and
// JSON shape.
func (s *LibraryService) probeStreams(path string) (ffprobeOutput, error) {
// Reading a local file's headers is quick. A longer run means ffprobe is
// stuck on a truncated or unreadable file, so cut it off rather than block
// the request that asked for it.
// Reading local headers is quick, so a longer run means ffprobe is stuck on a
// truncated or unreadable file.
ctx, cancel := context.WithTimeout(context.Background(), probeTimeout)
defer cancel()
@@ -82,8 +79,7 @@ func (s *LibraryService) probeMedia(path string) (*MediaMetadata, error) {
FileSize: info.Size(),
}
// A file ffprobe can't read still has a size worth showing, so a probe
// failure degrades to the size-only metadata rather than erroring.
// A file ffprobe can't read still has a size worth showing.
probe, err := s.probeStreams(path)
if err != nil {
return meta, nil
Minternal/service/preset.go
@@ -10,9 +10,8 @@ import (
)
// PresetService adds yt-dlp argument building on top of preset storage. The
// repository is embedded rather than wrapped in forwarding methods: every one of
// its methods (Save, GetByID, GetAll, GetDefault, Delete) is part of this
// service's contract anyway.
// repository is embedded rather than wrapped: every one of its methods is part
// of this service's contract anyway.
type PresetService struct {
*repository.PresetRepository
}
@@ -24,7 +23,6 @@ func NewPresetService(repo *repository.PresetRepository) *PresetService {
func (s *PresetService) BuildArgs(p *models.Preset, formatOverride, customFlags string) []string {
var args []string
// Format selection
if formatOverride != "" {
args = append(args, "-f", formatOverride)
} else if p.FormatMode == "custom" && p.CustomFormat != "" {
@@ -36,9 +34,8 @@ func (s *PresetService) BuildArgs(p *models.Preset, formatOverride, customFlags
args = append(args, "-f", p.Format)
}
}
// If FormatMode is "default" or empty, don't pass -f (let yt-dlp choose)
// A "default" or empty FormatMode passes no -f at all: yt-dlp chooses.
// Audio extraction
if p.ExtractAudio {
args = append(args, "--extract-audio")
if p.AudioFormat != "" {
@@ -46,7 +43,6 @@ func (s *PresetService) BuildArgs(p *models.Preset, formatOverride, customFlags
}
}
// Subtitles
if p.EmbedSubs {
args = append(args, "--embed-subs")
if p.SubLangs != "" {
@@ -54,18 +50,15 @@ func (s *PresetService) BuildArgs(p *models.Preset, formatOverride, customFlags
}
}
// Thumbnail
if p.EmbedThumbnail {
args = append(args, "--embed-thumbnail")
}
// Metadata
if p.EmbedMetadata {
args = append(args, "--embed-metadata")
}
// Info JSON and comments. Comments are stored in the info JSON sidecar, so
// enabling comment collection implicitly enables sidecar writing.
// Comments live in the info JSON sidecar, so collecting them implies writing it.
if p.WriteInfoJSON || p.WriteComments {
args = append(args, "--write-info-json")
}
@@ -80,10 +73,9 @@ func (s *PresetService) BuildArgs(p *models.Preset, formatOverride, customFlags
if p.MaxComments > 0 {
commentArgs = append(commentArgs, "max_comments="+strconv.Itoa(p.MaxComments))
}
// yt-dlp's syntax is IE_KEY:ARG1=VAL1,VAL2;ARG2=VAL: arguments are separated
// by ";", only the values of one argument by ",". Passing --extractor-args
// twice for the same key makes the second replace the first, so all youtube
// arguments must go into a single flag.
// yt-dlp's syntax is IE_KEY:ARG1=VAL1,VAL2;ARG2=VAL. A second --extractor-args
// for the same key replaces the first, so every youtube argument has to go
// into one flag.
extra := strings.TrimSpace(p.CommentExtractorArgs)
if rest, ok := cutYoutubePrefix(extra); ok {
if rest != "" {
@@ -98,15 +90,14 @@ func (s *PresetService) BuildArgs(p *models.Preset, formatOverride, customFlags
args = append(args, "--extractor-args", extra)
}
// Custom flags
args = append(args, splitCustomFlags(p.CustomFlags)...)
args = append(args, splitCustomFlags(customFlags)...)
return args
}
// cutYoutubePrefix strips a leading "youtube:" (any case) from extractor args.
// The bool reports whether the value targeted the youtube extractor at all.
// cutYoutubePrefix strips a leading "youtube:", any case. The bool reports
// whether the value targeted the youtube extractor at all.
func cutYoutubePrefix(s string) (string, bool) {
if len(s) < len("youtube:") || !strings.EqualFold(s[:len("youtube:")], "youtube:") {
return s, false
@@ -114,10 +105,9 @@ func cutYoutubePrefix(s string) (string, bool) {
return strings.TrimSpace(s[len("youtube:"):]), true
}
// splitCustomFlags splits a custom-flags string into arguments. BuildArgs cannot
// report an error, so malformed quoting falls back to whitespace splitting. The
// download path already rejected such input in checkReservedFlags, so only the
// EffectiveFlags display can reach the fallback.
// splitCustomFlags falls back to whitespace splitting on malformed quoting,
// because BuildArgs cannot report an error. Only the EffectiveFlags display
// reaches that fallback; the download path rejects such input beforehand.
func splitCustomFlags(s string) []string {
if strings.TrimSpace(s) == "" {
return nil
@@ -129,11 +119,9 @@ func splitCustomFlags(s string) []string {
return fields
}
// 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.
// applyHeightCap binds the filter to the video stream only, since a combined
// "video+audio" selector would otherwise constrain the audio too. The unfiltered
// fallback keeps a missing capped variant resolvable.
func applyHeightCap(format, maxHeight string) string {
parts := strings.Split(format, "+")
parts[0] = fmt.Sprintf("%s[height<=%s]", parts[0], maxHeight)
Minternal/service/progress_cache.go
@@ -21,8 +21,7 @@ func NewProgressCache() *ProgressCache {
}
}
// Start begins buffering live output for a download, discarding anything a
// previous run of the same id left behind.
// Start discards anything a previous run of the same id left behind.
func (c *ProgressCache) Start(id int64) {
c.mu.Lock()
defer c.mu.Unlock()
@@ -38,9 +37,8 @@ func (c *ProgressCache) AppendLog(id int64, line string) {
}
}
// Snapshot returns the full in-memory log buffer for a live download, read under
// the lock so it never races the writer in AppendLog (strings.Builder is not
// safe for concurrent read/write). Empty string if the download isn't live.
// Snapshot reads under the lock so it never races AppendLog: strings.Builder is
// not safe for concurrent read/write. Empty string if the download isn't live.
func (c *ProgressCache) Snapshot(id int64) string {
c.mu.RLock()
defer c.mu.RUnlock()
@@ -56,9 +54,8 @@ func (c *ProgressCache) Delete(id int64) {
delete(c.data, id)
}
// FlushLogs returns only the log content appended since the last flush and
// advances the flushed mark, so the caller can append (rather than rewrite the
// whole buffer) to the DB. Empty string when there's nothing new.
// FlushLogs returns only what was appended since the last call and advances the
// mark, so the caller appends to the DB instead of rewriting the whole buffer.
func (c *ProgressCache) FlushLogs(id int64) string {
c.mu.Lock()
defer c.mu.Unlock()
Minternal/service/settings.go
@@ -45,9 +45,8 @@ func (s *SettingsService) GetCookies() (string, error) {
return s.repo.Get("cookies")
}
// GetSessionSecret returns the secret that signs login cookies, or "" when none
// has been generated yet. It lives in the database so sessions survive a
// restart, and is rotated on sign-out so cookies issued earlier stop verifying.
// GetSessionSecret returns the secret that signs login cookies, "" when none has
// been generated. It lives in the database so sessions survive a restart.
func (s *SettingsService) GetSessionSecret() (string, error) {
return s.repo.Get("session_secret")
}
Minternal/service/subscription.go
@@ -15,8 +15,8 @@ import (
)
// SubscriptionService adds schedule resolution and archive-file ownership on top
// of subscription storage. The repository is embedded rather than wrapped in
// forwarding methods; Delete below deliberately shadows the repository's.
// of subscription storage. The repository is embedded rather than wrapped, and
// Delete below deliberately shadows the repository's.
type SubscriptionService struct {
*repository.SubscriptionRepository
cfg *config.Config
@@ -26,9 +26,8 @@ func NewSubscriptionService(repo *repository.SubscriptionRepository, cfg *config
return &SubscriptionService{SubscriptionRepository: repo, cfg: cfg}
}
// cronForKind maps a schedule kind to a standard 5-field cron expression. Preset
// kinds use canonical expressions; "cron" passes the user's custom expression
// through. Preset times default to 03:00 local to avoid the top-of-hour rush.
// cronForKind maps a schedule kind to a 5-field cron expression. The preset
// kinds run at 03:00 local, away from the top-of-hour rush.
func cronForKind(kind, custom string) (string, error) {
switch kind {
case "hourly":
@@ -49,8 +48,6 @@ func cronForKind(kind, custom string) (string, error) {
}
}
// CronExprFor resolves and validates the effective cron expression for a kind +
// custom expression, returning an error for unknown kinds or unparseable cron.
func (s *SubscriptionService) CronExprFor(kind, custom string) (string, error) {
expr, err := cronForKind(kind, custom)
if err != nil {
@@ -75,16 +72,14 @@ func (s *SubscriptionService) ComputeNextRun(sub *models.Subscription, from time
return sched.Next(from), nil
}
// ArchivePath is the per-subscription yt-dlp download-archive file used by
// "skip" refresh mode to avoid re-downloading already-fetched entries.
// ArchivePath is the yt-dlp download-archive file "skip" refresh mode reads.
func (s *SubscriptionService) ArchivePath(id int64) string {
return filepath.Join(s.cfg.DataDir, "archives", fmt.Sprintf("sub-%d.txt", id))
}
// Delete removes the subscription and the download-archive file "skip" mode
// keeps for it. Leaving the archive behind would make a later subscription that
// reuses the id silently skip entries it never downloaded. The library
// directory is left alone: the archived media is the point of the tool.
// Delete takes the archive file with it: left behind, a later subscription
// reusing the id would silently skip entries it never downloaded. The library
// directory stays: the archived media is the point of the tool.
func (s *SubscriptionService) Delete(id int64) error {
if err := s.SubscriptionRepository.Delete(id); err != nil {
return err
Minternal/service/subscription_library_test.go
@@ -53,9 +53,9 @@ func TestPruneToIDSet(t *testing.T) {
}
}
// A symlinked library root used to break cache eviction: the cache is keyed off
// the resolved root, so computing the key from the raw configured path missed,
// and a pruned item stayed visible until the scan TTL expired.
// The scan cache is keyed off the resolved root, so under a symlinked library
// root a key computed from the raw configured path misses and a pruned item
// stays visible until the TTL expires.
func TestPruneEvictsCacheUnderSymlinkedRoot(t *testing.T) {
parent := t.TempDir()
real := filepath.Join(parent, "real")
Minternal/service/subscription_run.go
@@ -14,14 +14,12 @@ import (
"vidarchive/internal/util"
)
// refreshAndAddNew handles a metadata-mode run. The main pass used
// --skip-download, so tempDownloadDir holds only info.json files. Existing
// library items have their markers refreshed in place; entries with no existing
// match are genuinely new and are downloaded as full items in a second pass.
// refreshAndAddNew handles a metadata-mode run, where tempDownloadDir holds only
// info.json files. Matched items have their markers refreshed in place; the rest
// are genuinely new and are downloaded as full items in a second pass.
func (s *DownloadService) refreshAndAddNew(ctx context.Context, d *models.Download, preset *models.Preset, tempDownloadDir, ytdlpFlags string) error {
// The library dir is not created here: a refresh that matches everything
// writes nothing, and importDownloadedItems creates it when a second pass
// actually has an item to add.
// The library dir is left uncreated: a refresh that matches everything writes
// nothing, and importDownloadedItems creates it when a second pass adds one.
baseLibraryDir, err := s.resolveBaseLibraryDir(d)
if err != nil {
return err
@@ -49,10 +47,9 @@ func (s *DownloadService) refreshAndAddNew(ctx context.Context, d *models.Downlo
if info.ID == "" {
continue
}
// A playlist-level info.json describes the source, not an item. Its id
// never matches a library item, so it would look new and trigger a
// download of the whole playlist. --no-write-playlist-metafiles already
// prevents it; this also covers sidecars written by older runs.
// A playlist-level info.json describes the source, not an item: its id
// never matches, so it would look new and re-download the whole playlist.
// --no-write-playlist-metafiles prevents it now; older runs left some.
if info.Type == "playlist" {
continue
}
@@ -73,9 +70,8 @@ func (s *DownloadService) refreshAndAddNew(ctx context.Context, d *models.Downlo
return s.downloadFresh(ctx, d, preset, newURLs, ytdlpFlags)
}
// downloadFresh fetches the given item URLs as full downloads (media + info.json)
// and imports them into the download's library directory. Metadata mode uses this
// to add entries that don't exist in the library yet.
// downloadFresh fetches the given URLs as full downloads, media included. It is
// how metadata mode adds entries the library does not have yet.
func (s *DownloadService) downloadFresh(ctx context.Context, d *models.Download, preset *models.Preset, urls []string, ytdlpFlags string) error {
tempDir := s.tempNewDirFor(d.ID)
if err := os.MkdirAll(tempDir, 0o755); err != nil {
@@ -98,9 +94,8 @@ func (s *DownloadService) downloadFresh(ctx context.Context, d *models.Download,
if ctx.Err() != nil {
return runErr
}
// Import whatever succeeded even if some entries errored. yt-dlp exits
// non-zero when one entry fails, so its error only counts as a failure when
// nothing at all was imported.
// yt-dlp exits non-zero when a single entry fails, so its error only counts
// as a failure when nothing at all was imported.
imported, err := s.importDownloadedItems(ctx, d, tempDir, "", ytdlpFlags)
if err != nil {
slog.Warn("failed to import new metadata-mode items", "err", err)
@@ -111,10 +106,9 @@ func (s *DownloadService) downloadFresh(ctx context.Context, d *models.Download,
return runErr
}
// mergeInfoJSON keeps fields from the existing sidecar that are absent from a
// metadata-only refresh. In particular, comments and heatmap data are expensive
// to reacquire and must not disappear just because the refresh preset does not
// request them.
// mergeInfoJSON keeps fields the refresh did not fetch. Comments and heatmap
// data are expensive to reacquire and must not disappear just because the
// refresh preset does not request them.
func mergeInfoJSON(oldData, newData []byte) ([]byte, error) {
var oldObject, newObject map[string]json.RawMessage
if err := json.Unmarshal(newData, &newObject); err != nil {
@@ -152,8 +146,7 @@ func isEmptyJSONArray(value json.RawMessage) bool {
return len(values) == 0
}
// applyMetadata refreshes an existing item's marker and info sidecar from a
// fresh info.json without touching its media.
// applyMetadata refreshes an item's marker and sidecar without touching media.
func (s *DownloadService) applyMetadata(existing string, info infoJSON, sourceInfoJSON string) error {
meta, err := s.librarySvc.readMetadata(existing)
if err != nil {
@@ -176,10 +169,9 @@ func (s *DownloadService) applyMetadata(existing string, info infoJSON, sourceIn
return fmt.Errorf("read existing marker: %w", err)
}
// Read and merge the sidecar before touching the marker, so a bad source file
// fails before anything is committed. The marker is written first; if the
// sidecar then fails to install, the old marker is restored so an ordinary
// I/O error cannot leave the two metadata files out of sync.
// Merge the sidecar before touching the marker, so a bad source file fails
// before anything is committed. If installing the sidecar still fails, the old
// marker is restored: the two files must not end up out of sync.
var infoData []byte
if sourceInfoJSON != "" {
data, err := os.ReadFile(sourceInfoJSON)
@@ -214,10 +206,9 @@ func (s *DownloadService) applyMetadata(existing string, info infoJSON, sourceIn
return nil
}
// pruneSubscription mirrors the source by deleting items in the subscription's
// directory that are no longer present upstream. It enumerates the current id
// set with a cheap flat-playlist listing; it never prunes when that enumeration
// fails or returns nothing, so a dead URL or network error can't wipe the dir.
// pruneSubscription deletes items no longer present upstream. It never prunes
// when the enumeration fails or comes back empty, so a dead URL or network error
// can't wipe the directory.
func (s *DownloadService) pruneSubscription(ctx context.Context, d *models.Download, sub *models.Subscription) {
baseLibraryDir, err := s.resolveBaseLibraryDir(d)
if err != nil {
@@ -245,10 +236,8 @@ func (s *DownloadService) pruneSubscription(ctx context.Context, d *models.Downl
}
}
// 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. Ids alone are sufficient to match items within
// a subscription's own directory (see FindByVideoID / PruneToIDSet).
// enumeratePlaylistIDs lists the current video-id set without downloading.
// Cookies are applied so private playlists enumerate correctly.
func (s *DownloadService) enumeratePlaylistIDs(ctx context.Context, url string) (map[string]bool, error) {
args := []string{"--flat-playlist", "--no-warnings", "--print", "%(id)s"}
Minternal/service/subtitles.go
@@ -20,9 +20,9 @@ func (s *LibraryService) SubtitleDir(relPath string) string {
return filepath.Join(itemDir, subtitlesDirName)
}
// GetSubtitlePath returns the .vtt path for a language. It errors when the item
// can't be resolved: joining onto an empty dir would yield a bare relative name
// that the caller would then serve relative to the process working directory.
// GetSubtitlePath errors when the item can't be resolved: joining onto an empty
// dir yields a bare relative name, which the caller would serve relative to the
// process working directory.
func (s *LibraryService) GetSubtitlePath(relPath, lang string) (string, error) {
dir := s.SubtitleDir(relPath)
if dir == "" {
@@ -102,9 +102,9 @@ func (s *LibraryService) extractSubtitleInfo(path string) ([]subtitleStream, err
func selectSubtitleStreams(all []ffprobeStream) []subtitleStream {
var streams []subtitleStream
// Index is an ffmpeg "0:s:N" selector, so it must count every subtitle
// stream. Counting only the convertible ones would shift the selector past a
// bitmap track (PGS, DVD subs) and extract the wrong stream.
// Index is an ffmpeg "0:s:N" selector, so every subtitle stream counts.
// Skipping the inconvertible ones shifts it past a bitmap track (PGS, DVD
// subs) and extracts the wrong stream.
subIndex := 0
for _, stream := range all {
if stream.CodecType != "subtitle" {
@@ -136,8 +136,7 @@ func selectSubtitleStreams(all []ffprobeStream) []subtitleStream {
return streams
}
// A subtitle track is small, but ffmpeg still walks the whole container to
// find it, so allow more than a probe and less than a download.
// A subtitle track is small, but ffmpeg walks the whole container to find it.
const subtitleExtractTimeout = 5 * time.Minute
func (s *LibraryService) extractSubtitleToVTT(inputPath, outputPath string, streamIndex int) error {
Minternal/service/thumbnail.go
@@ -13,16 +13,13 @@ import (
"vidarchive/internal/util"
)
// maxConcurrentThumbnails caps how many ffmpeg extraction processes may run at
// once, so a freshly loaded library page (which fires one thumbnail request per
// visible item) cannot spawn an unbounded ffmpeg storm.
// maxConcurrentThumbnails bounds the ffmpeg storm a freshly loaded library page
// would otherwise spawn, one process per visible item.
const maxConcurrentThumbnails = 3
// ThumbnailForFile returns the thumbnail for a specific media file within an
// item, extracting it on demand if needed. The bool is false when no thumbnail
// is available (file not found, audio-only, or extraction failed) so the caller
// can serve an icon. An empty/unmatched filename yields false — every thumbnail
// is keyed to a specific media file.
// ThumbnailForFile extracts on demand if needed. The bool is false when no
// thumbnail is available at all (file not found, audio-only, extraction failed),
// so the caller can serve an icon instead.
func (s *LibraryService) ThumbnailForFile(relPath, filename string) (string, bool) {
if filename == "" {
return "", false
@@ -44,11 +41,9 @@ func (s *LibraryService) ThumbnailForFile(relPath, filename string) (string, boo
return "", false
}
// ensureThumbnailForFile returns an existing thumbnail for mf or extracts one,
// serializing concurrent extraction of the same file via a per-path mutex. It
// re-checks the disk under the lock so that whichever caller wins the race does
// the work and the rest reuse the result. A prior in-process failure short-
// circuits to avoid re-running ffmpeg on every request (see thumbFailed).
// ensureThumbnailForFile serializes extraction of one file via a per-path mutex,
// re-checking the disk under the lock so the loser of the race reuses the
// winner's result instead of extracting again.
func (s *LibraryService) ensureThumbnailForFile(mf models.MediaFile) (string, bool) {
release := s.lockThumbFile(mf.Filepath)
defer release()
@@ -57,9 +52,8 @@ func (s *LibraryService) ensureThumbnailForFile(mf models.MediaFile) (string, bo
return path, true
}
// Trust a prior failure for this run rather than re-running ffmpeg every
// request; a restart clears thumbFailed and retries. Checked after the disk
// so a thumbnail that appears later (e.g. added manually) still wins.
// A prior failure is trusted for this run; a restart clears thumbFailed and
// retries. Checked after the disk so a manually added thumbnail still wins.
if _, failed := s.thumbFailed.Get(mf.Filepath); failed {
return "", false
}
@@ -87,11 +81,10 @@ func (s *LibraryService) findExistingThumbnail(path string) (string, bool) {
return "", false
}
// findImageAttachment returns the ordinal (0-based among attachment streams) of
// the best image attachment in a container — e.g. the cover.jpg/cover.webp that
// 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.
// findImageAttachment returns the ordinal among attachment streams of the best
// image attachment, or -1. The covers yt-dlp embeds into MKV are attachment
// streams, not attached_pic video streams, so they need -dump_attachment rather
// than a normal stream mapping.
func (s *LibraryService) findImageAttachment(path string) int {
probe, err := s.probeStreams(path)
if err != nil {
@@ -140,12 +133,12 @@ func (s *LibraryService) extractThumbnail(mf models.MediaFile) (string, error) {
jpgPath := base + ".jpg"
var attempts []thumbAttempt
// Collect every attempt's error so a genuine failure surfaces all of them
// rather than only the last fallback's stderr.
// Every attempt's error is collected, so a genuine failure surfaces all of
// them rather than only the last fallback's stderr.
var attemptErrs []string
// Prefer an embedded image attachment (e.g. yt-dlp's cover.webp/cover.jpg in
// MKV): dump its raw bytes, then transcode to a canonical WebP.
// An embedded image attachment is the best source: dump the raw bytes, then
// transcode to a canonical WebP.
if idx := s.findImageAttachment(mf.Filepath); idx >= 0 {
if rawPath, err := s.dumpAttachment(mf.Filepath, idx); err != nil {
attemptErrs = append(attemptErrs, fmt.Sprintf("attachment-dump: %v", err))
@@ -158,8 +151,8 @@ func (s *LibraryService) extractThumbnail(mf models.MediaFile) (string, error) {
}
}
// Next try an embedded cover art video stream (attached_pic), e.g. mp3/mp4,
// falling back to jpeg if libwebp or webp encoding fails.
// Next best is cover art carried as an attached_pic video stream, as mp3 and
// mp4 do. The jpg attempt covers a build without libwebp.
embedded := []string{"-i", mf.Filepath, "-map", "0:v", "-map", "-0:V", "-vframes", "1"}
attempts = append(attempts,
thumbAttempt{"embedded-webp", webpPath, append(embedded, "-c:v", "libwebp")},
@@ -192,13 +185,12 @@ func (s *LibraryService) extractThumbnail(mf models.MediaFile) (string, error) {
return "", fmt.Errorf("all thumbnail extraction attempts failed:\n%s", strings.Join(attemptErrs, "\n"))
}
// Seeking one frame out of a file is fast. A longer run means ffmpeg is stuck,
// and a stuck thumbnail must not hold a worker or a request forever.
// Seeking one frame is fast, so a longer run means ffmpeg is stuck, and a stuck
// thumbnail must not hold a request forever.
const thumbnailTimeout = 2 * time.Minute
// dumpAttachment extracts the raw bytes of attachment stream idx (e.g. the
// cover.jpg/cover.webp yt-dlp embeds into MKV with --embed-thumbnail) into a
// temp file and returns its path. The caller removes the file.
// dumpAttachment writes attachment stream idx to a temp file. The caller removes
// that file.
func (s *LibraryService) dumpAttachment(path string, idx int) (string, error) {
raw, err := os.CreateTemp("", "vidarchive-attachment-*")
if err != nil {
@@ -219,9 +211,8 @@ func (s *LibraryService) dumpAttachment(path string, idx int) (string, error) {
return rawPath, nil
}
// tryWriteThumbnail runs ffmpeg with args to produce outputPath. The temp file
// keeps the final extension so ffmpeg can infer the output muxer (it cannot for
// a bare ".tmp" suffix), then is atomically renamed into place.
// tryWriteThumbnail keeps the final extension on its temp file, because ffmpeg
// infers the output muxer from it and cannot for a bare ".tmp".
func (s *LibraryService) tryWriteThumbnail(outputPath string, args []string) (string, error) {
outExt := filepath.Ext(outputPath)
tmpPath := strings.TrimSuffix(outputPath, outExt) + ".tmp" + outExt
Minternal/service/thumbnail_test.go
@@ -33,7 +33,6 @@ func makeVideoWithCoverAttachment(t *testing.T, outPath, coverSize string) {
}
}
// probeImageSize returns the pixel dimensions of an image/video file.
func probeImageSize(t *testing.T, path string) (int, int) {
t.Helper()
out, err := exec.Command("ffprobe", "-v", "error", "-select_streams", "v:0",
@@ -198,8 +197,8 @@ func TestThumbnailUsesEmbeddedAttachment(t *testing.T) {
if !ok {
t.Fatal("thumbnail extraction failed")
}
// The embedded cover (100x100) must be used in preference to a video frame
// (which would be 64x64) — this is the regression the refactor introduced.
// The embedded cover (100x100) must win over a video frame, which would be
// 64x64.
if w, h := probeImageSize(t, path); w != 100 || h != 100 {
t.Errorf("thumbnail is %dx%d, expected 100x100 from the embedded cover (got a video frame instead)", w, h)
}
Minternal/service/ytdlp_flags.go
@@ -13,9 +13,9 @@ func splitFlags(s string) ([]string, error) {
return shellwords.Parse(s)
}
// reservedFlags are yt-dlp options VidArchive always sets itself; user custom
// flags must not pass them (or a conflicting inverse). The value describes what
// the option controls, for the failure message.
// reservedFlags are options VidArchive sets itself, so a custom flag must not
// pass them or a conflicting inverse. The value describes the option for the
// failure message.
var reservedFlags = map[string]string{
"-o": "the output template",
"--output": "the output template",
@@ -28,11 +28,10 @@ var reservedFlags = map[string]string{
"--write-playlist-metafiles": "playlist metadata files (VidArchive imports per-item metadata only)",
"--no-write-playlist-metafiles": "playlist metadata files (VidArchive imports per-item metadata only)",
// These options are blocked as a guard against accidental misuse. The list is
// not a security boundary: yt-dlp accepts unambiguous option prefixes (e.g.
// --exec-b) and --alias can define new options, and options such as
// --ffmpeg-location or --plugin-dirs are not listed. Custom flags are
// operator-controlled by design.
// These guard against accidental misuse only. This is not a security boundary:
// yt-dlp accepts unambiguous prefixes like --exec-b, --alias can define new
// options, and --ffmpeg-location and --plugin-dirs are not listed at all.
// Custom flags are operator-controlled by design.
"--exec": "running external commands (not permitted)",
"--exec-before-download": "running external commands (not permitted)",
"--postprocessor-args": "post-processor arguments (not permitted)",
@@ -62,7 +61,6 @@ func checkReservedFlags(customFlags string, isSubscription bool) error {
return fmt.Errorf("invalid custom flags: %w", err)
}
for _, tok := range tokens {
// Both "--flag value" and "--flag=value" name the same option.
name, _, _ := strings.Cut(tok, "=")
desc, ok := reservedFlags[name]
Minternal/util/exec.go
@@ -9,14 +9,11 @@ import (
const killWaitDelay = time.Second
// KillableCommand builds a command whose cancellation actually stops it, and is
// the single place that plumbing lives — every context-aware external command
// goes through here.
//
// yt-dlp, ffmpeg and ffprobe spawn helpers that inherit the output pipe, and
// both Wait and Output block reading it until every holder exits. So killing
// only the direct child does not end the call; the group signal does, with
// WaitDelay as the backstop for whatever survives it.
// KillableCommand is the single entry point for context-aware external commands.
// yt-dlp, ffmpeg and ffprobe spawn helpers that inherit the output pipe, and Wait
// and Output both block on it until every holder exits. Killing the direct child
// alone does not end the call; the group signal does, with WaitDelay as the
// backstop for whatever survives it.
func KillableCommand(ctx context.Context, path string, args ...string) *exec.Cmd {
cmd := exec.CommandContext(ctx, path, args...)
cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true}
Minternal/util/util.go
@@ -24,8 +24,7 @@ func FormatClock(seconds int) string {
return fmt.Sprintf("%d:%02d", m, s)
}
// FormatBytes renders a byte count as a compact human-readable size using
// 1024-based (binary) units with the matching iB suffix.
// FormatBytes uses 1024-based units with the matching iB suffix.
func FormatBytes(n int64) string {
const unit = 1024
if n < unit {
Minternal/worker/pool.go
@@ -32,8 +32,7 @@ func New(downloadSvc *service.DownloadService, workers int) *Pool {
}
func (p *Pool) Start() {
// Reset any downloads that were in progress during a previous run, and clear
// the temp dirs they left behind so a re-run doesn't import a duplicate.
// A previous run's temp dirs go too, or the re-run imports a duplicate.
if err := p.downloadSvc.ResetStalledDownloads(); err != nil {
slog.Error("failed to reset stalled downloads", "err", err)
}
@@ -47,17 +46,16 @@ func (p *Pool) Start() {
go p.queueChecker()
}
// Stop cancels running downloads and waits for the workers to return, so the
// process doesn't exit while a yt-dlp child is still writing to the temp dir.
// Stop waits for the workers, so the process does not exit while a yt-dlp child
// is still writing to the temp dir.
func (p *Pool) Stop() {
p.cancel()
p.downloadSvc.CancelAll()
p.wg.Wait()
}
// Submit enqueues a download without blocking. When the buffer is full the row
// stays "queued" in the database and the queue checker picks it up on a later
// tick, so a burst of submissions can't stall the HTTP handler that made them.
// Submit never blocks: on a full buffer the row stays "queued" for the queue
// checker, so a burst can't stall the HTTP handler that made it.
func (p *Pool) Submit(d *models.Download) {
if p.ctx.Err() != nil {
return
@@ -112,9 +110,8 @@ func (p *Pool) queueChecker() {
}
func (p *Pool) checkQueue() {
// 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).
// At most a channel's worth per tick; the rest waits for a later one.
// Duplicates are harmless because MarkStarted claims atomically.
downloads, err := p.downloadSvc.GetQueued(cap(p.queue))
if err != nil {
slog.Error("queue check failed", "err", err)
Minternal/worker/scheduler.go
@@ -11,9 +11,9 @@ import (
"vidarchive/internal/service"
)
// Scheduler periodically enqueues downloads for subscriptions whose next run is
// due. It reuses the worker Pool and the normal download pipeline; refresh-mode
// and pruning behaviour live in DownloadService.ExecuteDownload.
// Scheduler enqueues downloads for subscriptions whose next run is due. It
// reuses the normal download pipeline, so refresh-mode and pruning behaviour
// live in DownloadService.ExecuteDownload.
type Scheduler struct {
subscriptionSvc *service.SubscriptionService
downloadSvc *service.DownloadService
@@ -45,9 +45,8 @@ func (s *Scheduler) Start() {
go s.loop()
}
// Stop cancels the loop and waits for it to return. Waiting matters: main
// closes the database next, and a checkDue still in flight would query a closed
// handle or queue a run nobody will service.
// Stop waits for the loop because main closes the database next: a checkDue
// still in flight would query a closed handle or queue a run nobody services.
func (s *Scheduler) Stop() {
s.cancel()
s.wg.Wait()
@@ -59,8 +58,7 @@ func (s *Scheduler) loop() {
ticker := time.NewTicker(s.interval)
defer ticker.Stop()
// Check once promptly on startup so a subscription that came due during
// downtime doesn't wait a full interval.
// A subscription that came due during downtime should not wait a full interval.
s.checkDue()
for {
@@ -73,9 +71,8 @@ func (s *Scheduler) loop() {
}
}
// backfillNextRuns gives any enabled subscription without a next run time one,
// without running it — so a fresh restart doesn't fire every subscription that
// happens to have a null next_run_at.
// backfillNextRuns fills in a missing next run time without running anything, so
// a restart does not fire every subscription with a null next_run_at.
func (s *Scheduler) backfillNextRuns() {
subs, err := s.subscriptionSvc.GetAll()
if err != nil {
@@ -119,9 +116,8 @@ 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.
// A run longer than the interval would stack duplicate downloads. next_run_at
// still advances, so this is not re-evaluated every tick.
if active, aErr := s.downloadSvc.HasActiveForSubscription(sub.ID); aErr != nil {
slog.Error("scheduler: active-run check failed", "subscription_id", sub.ID, "err", aErr)
} else if active {
Minternal/worker/scheduler_test.go
@@ -201,7 +201,6 @@ func TestRunSkipsWhileAnEarlierRunIsActive(t *testing.T) {
if err := e.subRepo.MarkRun(sub.ID, lastRun, time.Now().Add(-time.Hour), "queued"); err != nil {
t.Fatal(err)
}
// An in-flight download for this subscription.
active := &models.Download{URL: sub.URL, Status: "downloading"}
active.SubscriptionID.Int64, active.SubscriptionID.Valid = sub.ID, true
if err := e.repo.Create(active); err != nil {