subscription_run.go
| 1 | package service |
| 2 | |
| 3 | import ( |
| 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. |
| 21 | func (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. |
| 79 | func (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. |
| 118 | func 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 | |
| 147 | func 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. |
| 157 | func 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. |
| 183 | func (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. |
| 269 | func (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). |
| 300 | func (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 |