cut over-engineering found by a repo-wide audit
Removes ~120 lines net with no behaviour change beyond the two notes below. - embed the repositories in PresetService/SubscriptionService instead of hand-writing 15 forwarding methods; Create/Update collapse into Save - replace orderComments' explicit stack and three inline anonymous structs with plain recursion - add atomicWrite and use it for the marker, the info.json sidecar and the rollback, replacing three copies of temp-write-rename - fold getEnvIntMin into getEnvInt; the port is now bounded too - delete imageExts (the mediaExts check already skipped images), the test-only extractAttempts counter, and the execer interface - assorted shrinks: GetByRelPath reuses scannedItem, one pass over the heatmap, a loop over the startup mkdirs, errBadRequest -> fmt.Errorf, LiveDownload unexported behind ProgressCache.Start The item marker is now written 0644 rather than 0600, matching its info.json sibling. applyMetadata builds the sidecar in memory instead of staging a temp file before the marker commit, which widens the rollback window slightly; the rollback itself is unchanged and still tested.
Mcmd/vidarchive/main.go
@@ -37,17 +37,15 @@ func main() {
checkDependencies(cfg)
if err := os.MkdirAll(cfg.DataDir, 0o755); err != nil {
log.Fatalf("Failed to create data dir: %v", err)
}
if err := os.MkdirAll(cfg.LibraryDir, 0o755); err != nil {
log.Fatalf("Failed to create library dir: %v", err)
}
if err := os.MkdirAll(cfg.TempDir, 0o755); err != nil {
log.Fatalf("Failed to create temp dir: %v", err)
}
if err := os.MkdirAll(filepath.Join(cfg.DataDir, "archives"), 0o755); err != nil {
log.Fatalf("Failed to create archives dir: %v", err)
for _, dir := range []string{
cfg.DataDir,
cfg.LibraryDir,
cfg.TempDir,
filepath.Join(cfg.DataDir, "archives"),
} {
if err := os.MkdirAll(dir, 0o755); err != nil {
log.Fatalf("Failed to create %s: %v", dir, err)
}
}
db, err := database.New(cfg)
Minternal/config/config.go
@@ -26,7 +26,7 @@ func New() *Config {
dataDir := getEnv("VIDARCHIVE_DATA_DIR", "./data")
return &Config{
Port: getEnvInt("VIDARCHIVE_PORT", 8080),
Port: getEnvInt("VIDARCHIVE_PORT", 8080, 1),
DataDir: dataDir,
DBPath: getEnv("VIDARCHIVE_DB_PATH", filepath.Join(dataDir, "vidarchive.db")),
LibraryDir: getEnv("VIDARCHIVE_LIBRARY_DIR", filepath.Join(dataDir, "library")),
@@ -35,8 +35,8 @@ func New() *Config {
FFmpegPath: getEnv("VIDARCHIVE_FFMPEG_PATH", "ffmpeg"),
FFprobePath: getEnv("VIDARCHIVE_FFPROBE_PATH", "ffprobe"),
BaseURL: getEnv("VIDARCHIVE_BASE_URL", ""),
Workers: getEnvIntMin("VIDARCHIVE_WORKERS", 2, 1),
SchedulerInterval: getEnvIntMin("VIDARCHIVE_SCHEDULER_INTERVAL", 60, 1),
Workers: getEnvInt("VIDARCHIVE_WORKERS", 2, 1),
SchedulerInterval: getEnvInt("VIDARCHIVE_SCHEDULER_INTERVAL", 60, 1),
}
}
@@ -51,7 +51,10 @@ func getEnv(key, defaultVal string) string {
return defaultVal
}
func getEnvInt(key string, defaultVal int) int {
// 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.
func getEnvInt(key string, defaultVal, min int) int {
v := os.Getenv(key)
if v == "" {
return defaultVal
@@ -61,14 +64,6 @@ func getEnvInt(key string, defaultVal int) int {
fmt.Fprintf(os.Stderr, "invalid %s: %v, using default %d\n", key, err, defaultVal)
return defaultVal
}
return i
}
// getEnvIntMin is getEnvInt with a lower bound. A value below min is rejected in
// favour of the default: a zero or negative worker count or tick interval would
// otherwise stall downloads or spin the scheduler.
func getEnvIntMin(key string, defaultVal, min int) int {
i := getEnvInt(key, defaultVal)
if i < min {
fmt.Fprintf(os.Stderr, "invalid %s: %d is below the minimum %d, using default %d\n", key, i, min, defaultVal)
return defaultVal
Minternal/handler/settings.go
@@ -49,7 +49,7 @@ func (h *Handler) CreatePreset(w http.ResponseWriter, r *http.Request) {
return
}
if err := h.presetSvc.Create(preset); err != nil {
if err := h.presetSvc.Save(preset); err != nil {
redirectWithError(w, r, "/settings", "Couldn't create this preset.", err)
return
}
@@ -134,7 +134,7 @@ func (h *Handler) UpdatePreset(w http.ResponseWriter, r *http.Request) {
return
}
if err := h.presetSvc.Update(preset); err != nil {
if err := h.presetSvc.Save(preset); err != nil {
redirectWithError(w, r, "/settings", "Couldn't update this preset.", err)
return
}
Minternal/handler/subscription.go
@@ -2,6 +2,7 @@ package handler
import (
"database/sql"
"fmt"
"log"
"net/http"
"strconv"
@@ -43,15 +44,15 @@ func (h *Handler) subscriptionFromForm(r *http.Request) (*models.Subscription, e
name := strings.TrimSpace(r.FormValue("name"))
url := strings.TrimSpace(r.FormValue("url"))
if name == "" {
return nil, errBadRequest("a name is required")
return nil, fmt.Errorf("a name is required")
}
if url == "" {
return nil, errBadRequest("URL is required")
return nil, fmt.Errorf("URL is required")
}
outputDir := strings.TrimSpace(r.FormValue("output_dir"))
if outputDir == "" {
return nil, errBadRequest("an output directory is required — the subscription owns this folder")
return nil, fmt.Errorf("an output directory is required — the subscription owns this folder")
}
refreshMode := r.FormValue("refresh_mode")
@@ -64,7 +65,7 @@ func (h *Handler) subscriptionFromForm(r *http.Request) (*models.Subscription, e
scheduleKind := r.FormValue("schedule_kind")
cronExpr, err := h.subscriptionSvc.CronExprFor(scheduleKind, strings.TrimSpace(r.FormValue("cron_expr")))
if err != nil {
return nil, errBadRequest(err.Error())
return nil, err
}
sub := &models.Subscription{
@@ -217,8 +218,3 @@ func (h *Handler) DeleteSubscription(w http.ResponseWriter, r *http.Request) {
}
redirectWithSuccess(w, r, "/subscriptions", "Subscription deleted.")
}
// errBadRequest is a small sentinel-style error carrying a user-facing message.
type errBadRequest string
func (e errBadRequest) Error() string { return string(e) }
Minternal/repository/preset.go
@@ -7,12 +7,6 @@ import (
"vidarchive/internal/models"
)
// execer is satisfied by both *sql.DB and *sql.Tx, so the insert/update/clear
// helpers can run either directly or inside a transaction (see Save).
type execer interface {
Exec(query string, args ...interface{}) (sql.Result, error)
}
type PresetRepository struct {
db *sql.DB
}
@@ -21,8 +15,8 @@ func NewPresetRepository(db *sql.DB) *PresetRepository {
return &PresetRepository{db: db}
}
func insertPreset(e execer, p *models.Preset) error {
result, err := e.Exec(
func insertPreset(tx *sql.Tx, p *models.Preset) error {
result, err := tx.Exec(
`INSERT INTO presets (name, description, is_default, format_mode, format, quality, custom_format, extract_audio, audio_format, embed_subs, sub_langs, embed_thumbnail, embed_metadata, write_info_json, write_comments, comment_sort, max_comments, comment_extractor_args, custom_flags)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
p.Name, p.Description, boolToInt(p.IsDefault), p.FormatMode, p.Format, p.Quality, p.CustomFormat,
@@ -42,8 +36,8 @@ func insertPreset(e execer, p *models.Preset) error {
return nil
}
func updatePreset(e execer, p *models.Preset) error {
_, err := e.Exec(
func updatePreset(tx *sql.Tx, p *models.Preset) error {
_, err := tx.Exec(
`UPDATE presets SET name=?, description=?, is_default=?, format_mode=?, format=?, quality=?, custom_format=?, extract_audio=?, audio_format=?, embed_subs=?, sub_langs=?, embed_thumbnail=?, embed_metadata=?, write_info_json=?, write_comments=?, comment_sort=?, max_comments=?, comment_extractor_args=?, custom_flags=?
WHERE id=?`,
p.Name, p.Description, boolToInt(p.IsDefault), p.FormatMode, p.Format, p.Quality, p.CustomFormat,
@@ -55,8 +49,8 @@ func updatePreset(e execer, p *models.Preset) error {
return err
}
func clearDefault(e execer) error {
_, err := e.Exec(`UPDATE presets SET is_default = 0`)
func clearDefault(tx *sql.Tx) error {
_, err := tx.Exec(`UPDATE presets SET is_default = 0`)
return err
}
Minternal/service/download.go
@@ -239,7 +239,7 @@ func (s *DownloadService) ExecuteDownload(parent context.Context, d *models.Down
defer cancel()
defer s.registerActive(d.ID, cancel)()
s.cache.Set(d.ID, &LiveDownload{})
s.cache.Start(d.ID)
defer s.cache.Delete(d.ID)
var preset *models.Preset
Minternal/service/engagement.go
@@ -78,14 +78,11 @@ func (s *LibraryService) GetEngagement(relPath string) ([]models.Comment, []mode
}
comments = orderComments(comments)
maxValue := 0.0
maxValue, endTime := 0.0, 0.0
for _, h := range raw.Heatmap {
if h.Value > maxValue && !math.IsNaN(h.Value) && !math.IsInf(h.Value, 0) {
maxValue = h.Value
}
}
endTime := 0.0
for _, h := range raw.Heatmap {
if h.EndTime > endTime {
endTime = h.EndTime
}
@@ -119,6 +116,10 @@ func (s *LibraryService) GetEngagement(relPath string) ([]models.Comment, []mode
return comments, heatmap, nil
}
// maxCommentDepth caps the rendered indent level; style.css defines
// .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.
@@ -154,41 +155,22 @@ func orderComments(comments []models.Comment) []models.Comment {
ordered := make([]models.Comment, 0, len(comments))
visited := make([]bool, len(comments))
stack := make([]struct {
index int
depth int
}, 0, len(comments))
walk := func(index, depth int) {
stack = append(stack, struct {
index int
depth int
}{index, depth})
for len(stack) > 0 {
last := len(stack) - 1
entry := stack[last]
stack = stack[:last]
if visited[entry.index] {
continue
}
visited[entry.index] = true
comment := comments[entry.index]
if parentIndex, ok := byID[strings.TrimSpace(comment.Parent)]; ok && parentIndex != entry.index {
comment.ReplyTo = comments[parentIndex].Author
}
if entry.depth > 4 {
comment.Depth = 4
} else {
comment.Depth = entry.depth
}
ordered = append(ordered, comment)
childrenForComment := children[entry.index]
for i := len(childrenForComment) - 1; i >= 0; i-- {
stack = append(stack, struct {
index int
depth int
}{childrenForComment[i], entry.depth + 1})
}
// ponytail: plain recursion, so stack depth is the reply-chain length. Swap
// in an explicit stack if a sidecar ever carries a chain deep enough to care.
var walk func(index, depth int)
walk = func(index, depth int) {
if visited[index] {
return
}
visited[index] = true
comment := comments[index]
if parentIndex, ok := byID[strings.TrimSpace(comment.Parent)]; ok && parentIndex != index {
comment.ReplyTo = comments[parentIndex].Author
}
comment.Depth = min(depth, maxCommentDepth)
ordered = append(ordered, comment)
for _, child := range children[index] {
walk(child, depth+1)
}
}
@@ -196,11 +178,9 @@ func orderComments(comments []models.Comment) []models.Comment {
walk(root, 0)
}
// Cycles or references to invalid parents are rendered as top-level comments
// rather than being dropped.
// rather than being dropped. Already-walked indices are a no-op.
for i := range comments {
if !visited[i] {
walk(i, 0)
}
walk(i, 0)
}
return ordered
}
Minternal/service/execute_download_test.go
@@ -312,7 +312,7 @@ func TestExecuteDownloadRejectsReservedPresetFlag(t *testing.T) {
e.fakeYTDLP(t, `touch "$SCRATCH/ran"`)
preset := &models.Preset{Name: "Bad", CustomFlags: "-o /tmp/anywhere.mp4"}
if err := e.svc.presetSvc.Create(preset); err != nil {
if err := e.svc.presetSvc.Save(preset); err != nil {
t.Fatalf("create preset: %v", err)
}
Minternal/service/library.go
@@ -1,6 +1,7 @@
package service
import (
"bytes"
"encoding/json"
"fmt"
"log"
@@ -31,10 +32,6 @@ var audioExts = map[string]struct{}{
".mp3": {}, ".wav": {}, ".flac": {}, ".aac": {}, ".opus": {}, ".m4a": {},
}
var imageExts = map[string]struct{}{
".webp": {}, ".jpg": {}, ".jpeg": {}, ".png": {}, ".gif": {}, ".bmp": {},
}
type LibraryService struct {
libraryDir string
// ffmpegPath/ffprobePath are the binaries used for thumbnail extraction,
@@ -52,10 +49,6 @@ type LibraryService struct {
// run, so we trust ffmpeg's verdict and don't re-run it on every request.
thumbFailed sync.Map
// extractAttempts counts ffmpeg extraction runs, so the negative cache can
// be verified to actually prevent a re-run.
extractAttempts int32
// 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
@@ -246,16 +239,10 @@ func (s *LibraryService) GetByRelPath(relPath string) (*models.LibraryItem, erro
if err != nil {
return nil, fmt.Errorf("item not found")
}
markerPath := filepath.Join(itemDir, itemMarkerName)
if _, err := os.Stat(markerPath); err != nil {
if _, err := os.Stat(filepath.Join(itemDir, itemMarkerName)); err != nil {
return nil, fmt.Errorf("item not found")
}
item, err := s.scanItem(itemDir, relPath)
if err != nil {
return nil, err
}
s.putCachedScan(relPath, item)
return item, nil
return s.scannedItem(itemDir, relPath)
}
func (s *LibraryService) scanItem(itemDir, relPath string) (*models.LibraryItem, error) {
@@ -389,25 +376,42 @@ func (s *LibraryService) readMetadata(itemDir string) (models.ItemMetadata, erro
}
func (s *LibraryService) writeMetadata(itemDir string, metadata models.ItemMetadata) error {
markerPath := filepath.Join(itemDir, itemMarkerName)
// Write to a uniquely-named temp file and rename so a crash mid-encode can't
// leave a truncated marker, and two concurrent writers never collide on a
// shared temp path.
f, err := os.CreateTemp(itemDir, ".vidarchive-item-*.tmp")
var buf bytes.Buffer
if err := toml.NewEncoder(&buf).Encode(metadata); err != nil {
return err
}
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.
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.
func atomicWrite(path string, data []byte, mode os.FileMode) error {
f, err := os.CreateTemp(filepath.Dir(path), ".vidarchive-*.tmp")
if err != nil {
return err
}
tmpPath := f.Name()
if err := toml.NewEncoder(f).Encode(metadata); err != nil {
// A no-op once the rename below has succeeded.
defer os.Remove(tmpPath)
if err := f.Chmod(mode); err != nil {
f.Close()
os.Remove(tmpPath)
return err
}
if _, err := f.Write(data); err != nil {
f.Close()
return err
}
// Close explicitly: a deferred close would hide a flush error on the write.
if err := f.Close(); err != nil {
os.Remove(tmpPath)
return err
}
return os.Rename(tmpPath, markerPath)
return os.Rename(tmpPath, path)
}
func (s *LibraryService) listItemFiles(itemDir string) ([]models.MediaFile, string, error) {
@@ -434,9 +438,6 @@ func (s *LibraryService) listItemFiles(itemDir string) ([]models.MediaFile, stri
infoJSONFiles = append(infoJSONFiles, path)
continue
}
if _, ok := imageExts[ext]; ok {
continue
}
if _, ok := mediaExts[ext]; !ok {
continue
}
Minternal/service/preset.go
@@ -9,36 +9,16 @@ import (
"vidarchive/internal/repository"
)
// 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.
type PresetService struct {
repo *repository.PresetRepository
*repository.PresetRepository
}
func NewPresetService(repo *repository.PresetRepository) *PresetService {
return &PresetService{repo: repo}
}
func (s *PresetService) Create(p *models.Preset) error {
return s.repo.Save(p)
}
func (s *PresetService) GetByID(id int64) (*models.Preset, error) {
return s.repo.GetByID(id)
}
func (s *PresetService) GetAll() ([]*models.Preset, error) {
return s.repo.GetAll()
}
func (s *PresetService) GetDefault() (*models.Preset, error) {
return s.repo.GetDefault()
}
func (s *PresetService) Update(p *models.Preset) error {
return s.repo.Save(p)
}
func (s *PresetService) Delete(id int64) error {
return s.repo.Delete(id)
return &PresetService{PresetRepository: repo}
}
func (s *PresetService) BuildArgs(p *models.Preset, formatOverride, customFlags string) []string {
Minternal/service/progress_cache.go
@@ -5,34 +5,36 @@ import (
"sync"
)
type LiveDownload struct {
Logs strings.Builder
flushed int // length of Logs already persisted to the DB
type liveDownload struct {
logs strings.Builder
flushed int // length of logs already persisted to the DB
}
type ProgressCache struct {
mu sync.RWMutex
data map[int64]*LiveDownload
data map[int64]*liveDownload
}
func NewProgressCache() *ProgressCache {
return &ProgressCache{
data: make(map[int64]*LiveDownload),
data: make(map[int64]*liveDownload),
}
}
func (c *ProgressCache) Set(id int64, d *LiveDownload) {
// Start begins buffering live output for a download, discarding anything a
// previous run of the same id left behind.
func (c *ProgressCache) Start(id int64) {
c.mu.Lock()
defer c.mu.Unlock()
c.data[id] = d
c.data[id] = &liveDownload{}
}
func (c *ProgressCache) AppendLog(id int64, line string) {
c.mu.Lock()
defer c.mu.Unlock()
if d, ok := c.data[id]; ok {
d.Logs.WriteString(line)
d.Logs.WriteByte('\n')
d.logs.WriteString(line)
d.logs.WriteByte('\n')
}
}
@@ -43,7 +45,7 @@ func (c *ProgressCache) Snapshot(id int64) string {
c.mu.RLock()
defer c.mu.RUnlock()
if d, ok := c.data[id]; ok {
return d.Logs.String()
return d.logs.String()
}
return ""
}
@@ -64,7 +66,7 @@ func (c *ProgressCache) FlushLogs(id int64) string {
if !ok {
return ""
}
full := d.Logs.String()
full := d.logs.String()
if d.flushed >= len(full) {
return ""
}
Minternal/service/progress_cache_test.go
@@ -7,7 +7,7 @@ import (
func TestProgressCacheFlushAndSnapshot(t *testing.T) {
c := NewProgressCache()
c.Set(1, &LiveDownload{})
c.Start(1)
c.AppendLog(1, "line1")
c.AppendLog(1, "line2")
@@ -40,7 +40,7 @@ func TestProgressCacheFlushAndSnapshot(t *testing.T) {
// must be serialized by the cache lock. Run with -race to catch a regression.
func TestProgressCacheConcurrentSnapshot(t *testing.T) {
c := NewProgressCache()
c.Set(1, &LiveDownload{})
c.Start(1)
var wg sync.WaitGroup
wg.Add(2)
Minternal/service/subscription.go
@@ -14,13 +14,16 @@ import (
"vidarchive/internal/repository"
)
// 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.
type SubscriptionService struct {
repo *repository.SubscriptionRepository
cfg *config.Config
*repository.SubscriptionRepository
cfg *config.Config
}
func NewSubscriptionService(repo *repository.SubscriptionRepository, cfg *config.Config) *SubscriptionService {
return &SubscriptionService{repo: repo, cfg: cfg}
return &SubscriptionService{SubscriptionRepository: repo, cfg: cfg}
}
// cronForKind maps a schedule kind to a standard 5-field cron expression. Preset
@@ -78,51 +81,12 @@ func (s *SubscriptionService) ArchivePath(id int64) string {
return filepath.Join(s.cfg.DataDir, "archives", fmt.Sprintf("sub-%d.txt", id))
}
func (s *SubscriptionService) GetAll() ([]*models.Subscription, error) {
return s.repo.GetAll()
}
func (s *SubscriptionService) GetByID(id int64) (*models.Subscription, error) {
return s.repo.GetByID(id)
}
func (s *SubscriptionService) GetDue(now time.Time) ([]*models.Subscription, error) {
return s.repo.GetDue(now)
}
func (s *SubscriptionService) Create(sub *models.Subscription) error {
return s.repo.Create(sub)
}
func (s *SubscriptionService) Update(sub *models.Subscription) error {
return s.repo.Update(sub)
}
func (s *SubscriptionService) SetEnabled(id int64, enabled bool) error {
return s.repo.SetEnabled(id, enabled)
}
func (s *SubscriptionService) MarkRun(id int64, lastRunAt, nextRunAt time.Time, status string) error {
return s.repo.MarkRun(id, lastRunAt, nextRunAt, status)
}
// MarkManualRun records a run started from the subscriptions page, leaving the
// schedule's next run time untouched.
func (s *SubscriptionService) MarkManualRun(id int64, lastRunAt time.Time, status string) error {
return s.repo.MarkManualRun(id, lastRunAt, status)
}
// SetLastStatus updates the recorded outcome of the current/last run.
func (s *SubscriptionService) SetLastStatus(id int64, status string) error {
return s.repo.SetLastStatus(id, status)
}
// 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.
func (s *SubscriptionService) Delete(id int64) error {
if err := s.repo.Delete(id); err != nil {
if err := s.SubscriptionRepository.Delete(id); err != nil {
return err
}
if err := os.Remove(s.ArchivePath(id)); err != nil && !os.IsNotExist(err) {
Minternal/service/subscription_run.go
@@ -152,32 +152,6 @@ func isEmptyJSONArray(value json.RawMessage) bool {
return len(values) == 0
}
// restoreFileAtomically puts data back at path without exposing a partial file.
// It is used to roll back the marker if installing the staged info sidecar fails.
func restoreFileAtomically(path string, data []byte) error {
tmp, err := os.CreateTemp(filepath.Dir(path), ".vidarchive-restore-*.tmp")
if err != nil {
return err
}
tmpPath := tmp.Name()
defer os.Remove(tmpPath)
if err := tmp.Chmod(0o644); err != nil {
tmp.Close()
return err
}
if _, err := tmp.Write(data); err != nil {
tmp.Close()
return err
}
if err := tmp.Close(); err != nil {
return err
}
if err := os.Rename(tmpPath, path); err != nil {
return err
}
return nil
}
// applyMetadata refreshes an existing item's marker and info sidecar from a
// fresh info.json without touching its media.
func (s *DownloadService) applyMetadata(existing string, info infoJSON, sourceInfoJSON string) error {
@@ -202,15 +176,11 @@ func (s *DownloadService) applyMetadata(existing string, info infoJSON, sourceIn
return fmt.Errorf("read existing marker: %w", err)
}
// Stage the sidecar before changing the marker. The marker is committed first;
// if installing the sidecar then fails, restore the old marker so an ordinary
// 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.
stagedInfo := ""
defer func() {
if stagedInfo != "" {
_ = os.Remove(stagedInfo)
}
}()
var infoData []byte
if sourceInfoJSON != "" {
data, err := os.ReadFile(sourceInfoJSON)
if err != nil {
@@ -226,35 +196,19 @@ func (s *DownloadService) applyMetadata(existing string, info infoJSON, sourceIn
return fmt.Errorf("merge refreshed info JSON: %w", err)
}
}
tmp, err := os.CreateTemp(existing, ".info-json-*.tmp")
if err != nil {
return fmt.Errorf("create refreshed info JSON: %w", err)
}
stagedInfo = tmp.Name()
if err := tmp.Chmod(0o644); err != nil {
tmp.Close()
return fmt.Errorf("set refreshed info JSON permissions: %w", err)
}
if _, err := tmp.Write(data); err != nil {
tmp.Close()
return fmt.Errorf("write refreshed info JSON: %w", err)
}
if err := tmp.Close(); err != nil {
return fmt.Errorf("close refreshed info JSON: %w", err)
}
infoData = data
}
if err := s.librarySvc.writeMetadata(existing, meta); err != nil {
return err
}
if stagedInfo != "" {
if err := os.Rename(stagedInfo, filepath.Join(existing, "info.json")); err != nil {
if restoreErr := restoreFileAtomically(markerPath, oldMarker); restoreErr != nil {
if infoData != nil {
if err := atomicWrite(filepath.Join(existing, "info.json"), infoData, markerFileMode); err != nil {
if restoreErr := atomicWrite(markerPath, oldMarker, markerFileMode); restoreErr != nil {
return fmt.Errorf("install refreshed info JSON: %v; restore marker: %w", err, restoreErr)
}
return fmt.Errorf("install refreshed info JSON: %w", err)
}
stagedInfo = ""
}
if rel, err := filepath.Rel(s.cfg.LibraryDir, existing); err == nil {
s.librarySvc.evictCachedScan(filepath.ToSlash(rel))
Minternal/service/thumbnail.go
@@ -8,7 +8,6 @@ import (
"path/filepath"
"strings"
"sync"
"sync/atomic"
"time"
"vidarchive/internal/models"
@@ -71,7 +70,6 @@ func (s *LibraryService) ensureThumbnailForFile(mf models.MediaFile) (string, bo
s.thumbSem <- struct{}{}
defer func() { <-s.thumbSem }()
atomic.AddInt32(&s.extractAttempts, 1)
path, err := s.extractThumbnail(mf)
if err != nil {
log.Printf("thumbnail extraction failed for %s: %v", mf.Filepath, err)
Minternal/service/thumbnail_test.go
@@ -7,7 +7,6 @@ import (
"path/filepath"
"strings"
"sync"
"sync/atomic"
"testing"
)
@@ -207,16 +206,32 @@ func TestThumbnailUsesEmbeddedAttachment(t *testing.T) {
}
func TestThumbnailNegativeCacheSkipsReextraction(t *testing.T) {
svc, dir := newLibrary(t)
dir, scratch := t.TempDir(), t.TempDir()
// Stand-in ffmpeg that always fails and records every invocation, so the
// extraction attempts can be counted without a probe into the service.
counter := filepath.Join(scratch, "attempts")
fakeFFmpeg := writeScript(t, filepath.Join(scratch, "ffmpeg"), "echo x >> "+counter+"\nexit 1")
attempts := func() int {
data, err := os.ReadFile(counter)
if err != nil {
return 0
}
return strings.Count(string(data), "\n")
}
newSvc := func() *LibraryService {
return NewLibraryService(dir, fakeFFmpeg, filepath.Join(scratch, "no-ffprobe"))
}
// A file ffmpeg cannot extract a thumbnail from: every attempt fails.
writeItem(t, dir, "bad", "name = \"B\"\nduration = -1\n", map[string]string{
"broken.mp4": "not actually a video",
})
svc := newSvc()
if _, ok := svc.ThumbnailForFile("bad", "broken.mp4"); ok {
t.Fatal("expected extraction to fail for a non-video file")
}
attempts1 := atomic.LoadInt32(&svc.extractAttempts)
attempts1 := attempts()
if attempts1 == 0 {
t.Fatal("expected at least one extraction attempt")
}
@@ -226,17 +241,16 @@ func TestThumbnailNegativeCacheSkipsReextraction(t *testing.T) {
if _, ok := svc.ThumbnailForFile("bad", "broken.mp4"); ok {
t.Fatal("expected the cached failure to persist")
}
if attempts2 := atomic.LoadInt32(&svc.extractAttempts); attempts2 != attempts1 {
if attempts2 := attempts(); attempts2 != attempts1 {
t.Errorf("negative cache should prevent re-extraction; attempts %d -> %d", attempts1, attempts2)
}
// The cache is in-process only: a fresh service (≈ a restart) retries.
fresh := NewLibraryService(dir, "ffmpeg", "ffprobe")
if _, ok := fresh.ThumbnailForFile("bad", "broken.mp4"); ok {
if _, ok := newSvc().ThumbnailForFile("bad", "broken.mp4"); ok {
t.Fatal("fresh service still fails (file is unextractable)")
}
if atomic.LoadInt32(&fresh.extractAttempts) == 0 {
t.Error("a fresh service should retry extraction, not inherit the negative cache")
if attempts3 := attempts(); attempts3 <= attempts1 {
t.Errorf("a fresh service should retry extraction, not inherit the negative cache; attempts %d -> %d", attempts1, attempts3)
}
}
Minternal/service/ytdlp_live_test.go
@@ -99,7 +99,7 @@ func TestLiveExecuteDownload(t *testing.T) {
e, videoURL, videoID := liveEnv(t)
preset := &models.Preset{Name: "Live", WriteInfoJSON: true}
if err := e.svc.presetSvc.Create(preset); err != nil {
if err := e.svc.presetSvc.Save(preset); err != nil {
t.Fatalf("create preset: %v", err)
}