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