subscription_run.go
⎇
Raw
1package service
2
3import (
4 "context"
5 "encoding/json"
6 "fmt"
7 "log"
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 log.Printf("warning: failed to refresh metadata for %s: %v", itemDir, 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 log.Printf("warning: failed to import new metadata-mode items: %v", 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// restoreFileAtomically puts data back at path without exposing a partial file.
156// It is used to roll back the marker if installing the staged info sidecar fails.
157func restoreFileAtomically(path string, data []byte) error {
158 tmp, err := os.CreateTemp(filepath.Dir(path), ".vidarchive-restore-*.tmp")
159 if err != nil {
160 return err
161 }
162 tmpPath := tmp.Name()
163 defer os.Remove(tmpPath)
164 if err := tmp.Chmod(0o644); err != nil {
165 tmp.Close()
166 return err
167 }
168 if _, err := tmp.Write(data); err != nil {
169 tmp.Close()
170 return err
171 }
172 if err := tmp.Close(); err != nil {
173 return err
174 }
175 if err := os.Rename(tmpPath, path); err != nil {
176 return err
177 }
178 return nil
179}
180
181// applyMetadata refreshes an existing item's marker and info sidecar from a
182// fresh info.json without touching its media.
183func (s *DownloadService) applyMetadata(existing string, info infoJSON, sourceInfoJSON string) error {
184 meta, err := s.librarySvc.readMetadata(existing)
185 if err != nil {
186 return fmt.Errorf("read metadata for %s: %w", existing, err)
187 }
188 if info.Title != "" {
189 meta.Name = info.Title
190 }
191 if info.Description != "" {
192 meta.Description = info.Description
193 }
194 if meta.SourceURL == "" && info.WebpageURL != "" {
195 meta.SourceURL = info.WebpageURL
196 }
197 meta.VideoID = info.ID
198
199 markerPath := filepath.Join(existing, itemMarkerName)
200 oldMarker, err := os.ReadFile(markerPath)
201 if err != nil {
202 return fmt.Errorf("read existing marker: %w", err)
203 }
204
205 // Stage the sidecar before changing the marker. The marker is committed first;
206 // if installing the sidecar then fails, restore the old marker so an ordinary
207 // I/O error cannot leave the two metadata files out of sync.
208 stagedInfo := ""
209 defer func() {
210 if stagedInfo != "" {
211 _ = os.Remove(stagedInfo)
212 }
213 }()
214 if sourceInfoJSON != "" {
215 data, err := os.ReadFile(sourceInfoJSON)
216 if err != nil {
217 return fmt.Errorf("read refreshed info JSON: %w", err)
218 }
219 if oldInfoJSON := findInfoJSON(existing); oldInfoJSON != "" {
220 oldData, err := os.ReadFile(oldInfoJSON)
221 if err != nil {
222 return fmt.Errorf("read existing info JSON: %w", err)
223 }
224 data, err = mergeInfoJSON(oldData, data)
225 if err != nil {
226 return fmt.Errorf("merge refreshed info JSON: %w", err)
227 }
228 }
229 tmp, err := os.CreateTemp(existing, ".info-json-*.tmp")
230 if err != nil {
231 return fmt.Errorf("create refreshed info JSON: %w", err)
232 }
233 stagedInfo = tmp.Name()
234 if err := tmp.Chmod(0o644); err != nil {
235 tmp.Close()
236 return fmt.Errorf("set refreshed info JSON permissions: %w", err)
237 }
238 if _, err := tmp.Write(data); err != nil {
239 tmp.Close()
240 return fmt.Errorf("write refreshed info JSON: %w", err)
241 }
242 if err := tmp.Close(); err != nil {
243 return fmt.Errorf("close refreshed info JSON: %w", err)
244 }
245 }
246
247 if err := s.librarySvc.writeMetadata(existing, meta); err != nil {
248 return err
249 }
250 if stagedInfo != "" {
251 if err := os.Rename(stagedInfo, filepath.Join(existing, "info.json")); err != nil {
252 if restoreErr := restoreFileAtomically(markerPath, oldMarker); restoreErr != nil {
253 return fmt.Errorf("install refreshed info JSON: %v; restore marker: %w", err, restoreErr)
254 }
255 return fmt.Errorf("install refreshed info JSON: %w", err)
256 }
257 stagedInfo = ""
258 }
259 if rel, err := filepath.Rel(s.cfg.LibraryDir, existing); err == nil {
260 s.librarySvc.evictCachedScan(filepath.ToSlash(rel))
261 }
262 return nil
263}
264
265// pruneSubscription mirrors the source by deleting items in the subscription's
266// directory that are no longer present upstream. It enumerates the current id
267// set with a cheap flat-playlist listing; it never prunes when that enumeration
268// fails or returns nothing, so a dead URL or network error can't wipe the dir.
269func (s *DownloadService) pruneSubscription(ctx context.Context, d *models.Download, sub *models.Subscription) {
270 baseLibraryDir, err := s.resolveBaseLibraryDir(d)
271 if err != nil {
272 log.Printf("subscription %d prune skipped: %v", sub.ID, err)
273 return
274 }
275
276 keep, err := s.enumeratePlaylistIDs(ctx, sub.URL)
277 if err != nil {
278 log.Printf("subscription %d prune skipped: enumeration failed: %v", sub.ID, err)
279 return
280 }
281 if len(keep) == 0 {
282 log.Printf("subscription %d prune skipped: source returned no entries", sub.ID)
283 return
284 }
285
286 removed, err := s.librarySvc.PruneToIDSet(baseLibraryDir, keep)
287 if err != nil {
288 log.Printf("subscription %d prune error: %v", sub.ID, err)
289 return
290 }
291 if removed > 0 {
292 log.Printf("subscription %d pruned %d item(s) removed upstream", sub.ID, removed)
293 }
294}
295
296// enumeratePlaylistIDs lists the current video-id set for a URL without
297// downloading, using yt-dlp --flat-playlist. Cookies are applied so private
298// playlists enumerate correctly. Ids alone are sufficient to match items within
299// a subscription's own directory (see FindByVideoID / PruneToIDSet).
300func (s *DownloadService) enumeratePlaylistIDs(ctx context.Context, url string) (map[string]bool, error) {
301 args := []string{"--flat-playlist", "--no-warnings", "--print", "%(id)s"}
302
303 args, cleanup := s.appendCookies(args)
304 defer cleanup()
305 args = append(args, url)
306
307 out, err := util.KillableCommand(ctx, s.cfg.YTDLPPath, args...).Output()
308 if err != nil {
309 return nil, err
310 }
311
312 keep := make(map[string]bool)
313 for _, line := range strings.Split(string(out), "\n") {
314 id := strings.TrimSpace(line)
315 // yt-dlp prints "NA" for a missing field; never treat that as a real id.
316 if id == "" || id == "NA" {
317 continue
318 }
319 keep[id] = true
320 }
321 return keep, nil
322}
323