subscription_run.go
⎇
Raw
1package service
2
3import (
4 "context"
5 "encoding/json"
6 "fmt"
7 "log/slog"
8 "os"
9 "path/filepath"
10 "slices"
11 "strings"
12
13 "vidarchive/internal/models"
14 "vidarchive/internal/util"
15)
16
17// refreshAndAddNew handles a metadata-mode run. The main pass used
18// --skip-download, so tempDownloadDir holds only info.json files. Existing
19// library items have their markers refreshed in place; entries with no existing
20// match are genuinely new and are downloaded as full items in a second pass.
21func (s *DownloadService) refreshAndAddNew(ctx context.Context, d *models.Download, preset *models.Preset, tempDownloadDir, ytdlpFlags string) error {
22 // The library dir is not created here: a refresh that matches everything
23 // writes nothing, and importDownloadedItems creates it when a second pass
24 // actually has an item to add.
25 baseLibraryDir, err := s.resolveBaseLibraryDir(d)
26 if err != nil {
27 return err
28 }
29
30 entries, err := os.ReadDir(tempDownloadDir)
31 if err != nil {
32 return err
33 }
34
35 var newURLs []string
36 for _, entry := range entries {
37 if err := ctx.Err(); err != nil {
38 return err
39 }
40 if !entry.IsDir() || !strings.HasPrefix(entry.Name(), "item-") {
41 continue
42 }
43 itemDir := filepath.Join(tempDownloadDir, entry.Name())
44 infoJSONPath := findInfoJSON(itemDir)
45 if infoJSONPath == "" {
46 continue
47 }
48 info := readInfoJSON(infoJSONPath)
49 if info.ID == "" {
50 continue
51 }
52 // A playlist-level info.json describes the source, not an item. Its id
53 // never matches a library item, so it would look new and trigger a
54 // download of the whole playlist. --no-write-playlist-metafiles already
55 // prevents it; this also covers sidecars written by older runs.
56 if info.Type == "playlist" {
57 continue
58 }
59 if existing, ok := s.librarySvc.FindByVideoID(baseLibraryDir, info.ID); ok {
60 if err := s.applyMetadata(existing, info, infoJSONPath); err != nil {
61 slog.Warn("metadata refresh failed", "dir", itemDir, "err", err)
62 }
63 continue
64 }
65 if info.WebpageURL != "" {
66 newURLs = append(newURLs, info.WebpageURL)
67 }
68 }
69
70 if len(newURLs) == 0 {
71 return nil
72 }
73 return s.downloadFresh(ctx, d, preset, newURLs, ytdlpFlags)
74}
75
76// downloadFresh fetches the given item URLs as full downloads (media + info.json)
77// and imports them into the download's library directory. Metadata mode uses this
78// to add entries that don't exist in the library yet.
79func (s *DownloadService) downloadFresh(ctx context.Context, d *models.Download, preset *models.Preset, urls []string, ytdlpFlags string) error {
80 tempDir := s.tempNewDirFor(d.ID)
81 if err := os.MkdirAll(tempDir, 0o755); err != nil {
82 return err
83 }
84 defer os.RemoveAll(tempDir)
85
86 args := s.presetSvc.BuildArgs(preset, d.FormatOverride, d.CustomFlags)
87 args, cleanup := s.appendCookies(args)
88 defer cleanup()
89 if !slices.Contains(args, "--write-info-json") {
90 args = append(args, "--write-info-json")
91 }
92 args = append(args, "-P", tempDir)
93 args = append(args, "-o", "item-%(autonumber)05d/%(title)s.%(ext)s")
94 args = append(args, urls...)
95
96 runErr := s.runYTDLP(ctx, d, args)
97 // A cancelled second pass has nothing worth importing.
98 if ctx.Err() != nil {
99 return runErr
100 }
101 // Import whatever succeeded even if some entries errored. yt-dlp exits
102 // non-zero when one entry fails, so its error only counts as a failure when
103 // nothing at all was imported.
104 imported, err := s.importDownloadedItems(ctx, d, tempDir, "", ytdlpFlags)
105 if err != nil {
106 slog.Warn("failed to import new metadata-mode items", "err", err)
107 }
108 if imported > 0 {
109 return nil
110 }
111 return runErr
112}
113
114// mergeInfoJSON keeps fields from the existing sidecar that are absent from a
115// metadata-only refresh. In particular, comments and heatmap data are expensive
116// to reacquire and must not disappear just because the refresh preset does not
117// request them.
118func mergeInfoJSON(oldData, newData []byte) ([]byte, error) {
119 var oldObject, newObject map[string]json.RawMessage
120 if err := json.Unmarshal(newData, &newObject); err != nil {
121 return nil, err
122 }
123 if err := json.Unmarshal(oldData, &oldObject); err != nil {
124 return newData, nil
125 }
126 if newObject == nil {
127 return newData, nil
128 }
129
130 merged := make(map[string]json.RawMessage, len(oldObject)+len(newObject))
131 for key, value := range oldObject {
132 merged[key] = value
133 }
134 for key, value := range newObject {
135 merged[key] = value
136 }
137 for _, key := range []string{"comments", "heatmap"} {
138 oldValue, hadOldValue := oldObject[key]
139 newValue, hasNewValue := newObject[key]
140 if hadOldValue && (!hasNewValue || isEmptyJSONArray(newValue)) {
141 merged[key] = oldValue
142 }
143 }
144 return json.Marshal(merged)
145}
146
147func isEmptyJSONArray(value json.RawMessage) bool {
148 var values []json.RawMessage
149 if err := json.Unmarshal(value, &values); err != nil {
150 return false
151 }
152 return len(values) == 0
153}
154
155// applyMetadata refreshes an existing item's marker and info sidecar from a
156// fresh info.json without touching its media.
157func (s *DownloadService) applyMetadata(existing string, info infoJSON, sourceInfoJSON string) error {
158 meta, err := s.librarySvc.readMetadata(existing)
159 if err != nil {
160 return fmt.Errorf("read metadata for %s: %w", existing, err)
161 }
162 if info.Title != "" {
163 meta.Name = info.Title
164 }
165 if info.Description != "" {
166 meta.Description = info.Description
167 }
168 if meta.SourceURL == "" && info.WebpageURL != "" {
169 meta.SourceURL = info.WebpageURL
170 }
171 meta.VideoID = info.ID
172
173 markerPath := filepath.Join(existing, itemMarkerName)
174 oldMarker, err := os.ReadFile(markerPath)
175 if err != nil {
176 return fmt.Errorf("read existing marker: %w", err)
177 }
178
179 // Read and merge the sidecar before touching the marker, so a bad source file
180 // fails before anything is committed. The marker is written first; if the
181 // sidecar then fails to install, the old marker is restored so an ordinary
182 // I/O error cannot leave the two metadata files out of sync.
183 var infoData []byte
184 if sourceInfoJSON != "" {
185 data, err := os.ReadFile(sourceInfoJSON)
186 if err != nil {
187 return fmt.Errorf("read refreshed info JSON: %w", err)
188 }
189 if oldInfoJSON := findInfoJSON(existing); oldInfoJSON != "" {
190 oldData, err := os.ReadFile(oldInfoJSON)
191 if err != nil {
192 return fmt.Errorf("read existing info JSON: %w", err)
193 }
194 data, err = mergeInfoJSON(oldData, data)
195 if err != nil {
196 return fmt.Errorf("merge refreshed info JSON: %w", err)
197 }
198 }
199 infoData = data
200 }
201
202 if err := s.librarySvc.writeMetadata(existing, meta); err != nil {
203 return err
204 }
205 if infoData != nil {
206 if err := atomicWrite(filepath.Join(existing, "info.json"), infoData, markerFileMode); err != nil {
207 if restoreErr := atomicWrite(markerPath, oldMarker, markerFileMode); restoreErr != nil {
208 return fmt.Errorf("install refreshed info JSON: %v; restore marker: %w", err, restoreErr)
209 }
210 return fmt.Errorf("install refreshed info JSON: %w", err)
211 }
212 }
213 if rel, err := filepath.Rel(s.cfg.LibraryDir, existing); err == nil {
214 s.librarySvc.evictCachedScan(filepath.ToSlash(rel))
215 }
216 return nil
217}
218
219// pruneSubscription mirrors the source by deleting items in the subscription's
220// directory that are no longer present upstream. It enumerates the current id
221// set with a cheap flat-playlist listing; it never prunes when that enumeration
222// fails or returns nothing, so a dead URL or network error can't wipe the dir.
223func (s *DownloadService) pruneSubscription(ctx context.Context, d *models.Download, sub *models.Subscription) {
224 baseLibraryDir, err := s.resolveBaseLibraryDir(d)
225 if err != nil {
226 slog.Warn("prune skipped", "subscription_id", sub.ID, "err", err)
227 return
228 }
229
230 keep, err := s.enumeratePlaylistIDs(ctx, sub.URL)
231 if err != nil {
232 slog.Warn("prune skipped, enumeration failed", "subscription_id", sub.ID, "err", err)
233 return
234 }
235 if len(keep) == 0 {
236 slog.Warn("prune skipped, source returned no entries", "subscription_id", sub.ID)
237 return
238 }
239
240 removed, err := s.librarySvc.PruneToIDSet(baseLibraryDir, keep)
241 if err != nil {
242 slog.Error("prune failed", "subscription_id", sub.ID, "err", err)
243 return
244 }
245 if removed > 0 {
246 slog.Info("pruned items removed upstream", "subscription_id", sub.ID, "removed", removed)
247 }
248}
249
250// enumeratePlaylistIDs lists the current video-id set for a URL without
251// downloading, using yt-dlp --flat-playlist. Cookies are applied so private
252// playlists enumerate correctly. Ids alone are sufficient to match items within
253// a subscription's own directory (see FindByVideoID / PruneToIDSet).
254func (s *DownloadService) enumeratePlaylistIDs(ctx context.Context, url string) (map[string]bool, error) {
255 args := []string{"--flat-playlist", "--no-warnings", "--print", "%(id)s"}
256
257 args, cleanup := s.appendCookies(args)
258 defer cleanup()
259 args = append(args, url)
260
261 out, err := util.KillableCommand(ctx, s.cfg.YTDLPPath, args...).Output()
262 if err != nil {
263 return nil, err
264 }
265
266 keep := make(map[string]bool)
267 for _, line := range strings.Split(string(out), "\n") {
268 id := strings.TrimSpace(line)
269 // yt-dlp prints "NA" for a missing field; never treat that as a real id.
270 if id == "" || id == "NA" {
271 continue
272 }
273 keep[id] = true
274 }
275 return keep, nil
276}
277