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