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