subscription_run.go
| 1 | package service |
| 2 | |
| 3 | import ( |
| 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. 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 | slog.Warn("metadata refresh failed", "dir", itemDir, "err", 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 | slog.Warn("failed to import new metadata-mode items", "err", 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 | // applyMetadata refreshes an existing item's marker and info sidecar from a |
| 156 | // fresh info.json without touching its media. |
| 157 | func (s *DownloadService) applyMetadata(existing string, info infoJSON, sourceInfoJSON string) error { |
| 158 | meta, err := s.librarySvc.readMetadata(existing) |
| 159 | if err != nil { |
| 160 | return fmt.Errorf("read metadata for %s: %w", existing, err) |
| 161 | } |
| 162 | if info.Title != "" { |
| 163 | meta.Name = info.Title |
| 164 | } |
| 165 | if info.Description != "" { |
| 166 | meta.Description = info.Description |
| 167 | } |
| 168 | if meta.SourceURL == "" && info.WebpageURL != "" { |
| 169 | meta.SourceURL = info.WebpageURL |
| 170 | } |
| 171 | meta.VideoID = info.ID |
| 172 | |
| 173 | markerPath := filepath.Join(existing, itemMarkerName) |
| 174 | oldMarker, err := os.ReadFile(markerPath) |
| 175 | if err != nil { |
| 176 | return fmt.Errorf("read existing marker: %w", err) |
| 177 | } |
| 178 | |
| 179 | // Read and merge the sidecar before touching the marker, so a bad source file |
| 180 | // fails before anything is committed. The marker is written first; if the |
| 181 | // sidecar then fails to install, the old marker is restored so an ordinary |
| 182 | // I/O error cannot leave the two metadata files out of sync. |
| 183 | var infoData []byte |
| 184 | if sourceInfoJSON != "" { |
| 185 | data, err := os.ReadFile(sourceInfoJSON) |
| 186 | if err != nil { |
| 187 | return fmt.Errorf("read refreshed info JSON: %w", err) |
| 188 | } |
| 189 | if oldInfoJSON := findInfoJSON(existing); oldInfoJSON != "" { |
| 190 | oldData, err := os.ReadFile(oldInfoJSON) |
| 191 | if err != nil { |
| 192 | return fmt.Errorf("read existing info JSON: %w", err) |
| 193 | } |
| 194 | data, err = mergeInfoJSON(oldData, data) |
| 195 | if err != nil { |
| 196 | return fmt.Errorf("merge refreshed info JSON: %w", err) |
| 197 | } |
| 198 | } |
| 199 | infoData = data |
| 200 | } |
| 201 | |
| 202 | if err := s.librarySvc.writeMetadata(existing, meta); err != nil { |
| 203 | return err |
| 204 | } |
| 205 | if infoData != nil { |
| 206 | if err := atomicWrite(filepath.Join(existing, "info.json"), infoData, markerFileMode); err != nil { |
| 207 | if restoreErr := atomicWrite(markerPath, oldMarker, markerFileMode); restoreErr != nil { |
| 208 | return fmt.Errorf("install refreshed info JSON: %v; restore marker: %w", err, restoreErr) |
| 209 | } |
| 210 | return fmt.Errorf("install refreshed info JSON: %w", err) |
| 211 | } |
| 212 | } |
| 213 | s.librarySvc.evictCachedDir(existing) |
| 214 | return nil |
| 215 | } |
| 216 | |
| 217 | // pruneSubscription mirrors the source by deleting items in the subscription's |
| 218 | // directory that are no longer present upstream. It enumerates the current id |
| 219 | // set with a cheap flat-playlist listing; it never prunes when that enumeration |
| 220 | // fails or returns nothing, so a dead URL or network error can't wipe the dir. |
| 221 | func (s *DownloadService) pruneSubscription(ctx context.Context, d *models.Download, sub *models.Subscription) { |
| 222 | baseLibraryDir, err := s.resolveBaseLibraryDir(d) |
| 223 | if err != nil { |
| 224 | slog.Warn("prune skipped", "subscription_id", sub.ID, "err", err) |
| 225 | return |
| 226 | } |
| 227 | |
| 228 | keep, err := s.enumeratePlaylistIDs(ctx, sub.URL) |
| 229 | if err != nil { |
| 230 | slog.Warn("prune skipped, enumeration failed", "subscription_id", sub.ID, "err", err) |
| 231 | return |
| 232 | } |
| 233 | if len(keep) == 0 { |
| 234 | slog.Warn("prune skipped, source returned no entries", "subscription_id", sub.ID) |
| 235 | return |
| 236 | } |
| 237 | |
| 238 | removed, err := s.librarySvc.PruneToIDSet(baseLibraryDir, keep) |
| 239 | if err != nil { |
| 240 | slog.Error("prune failed", "subscription_id", sub.ID, "err", err) |
| 241 | return |
| 242 | } |
| 243 | if removed > 0 { |
| 244 | slog.Info("pruned items removed upstream", "subscription_id", sub.ID, "removed", removed) |
| 245 | } |
| 246 | } |
| 247 | |
| 248 | // enumeratePlaylistIDs lists the current video-id set for a URL without |
| 249 | // downloading, using yt-dlp --flat-playlist. Cookies are applied so private |
| 250 | // playlists enumerate correctly. Ids alone are sufficient to match items within |
| 251 | // a subscription's own directory (see FindByVideoID / PruneToIDSet). |
| 252 | func (s *DownloadService) enumeratePlaylistIDs(ctx context.Context, url string) (map[string]bool, error) { |
| 253 | args := []string{"--flat-playlist", "--no-warnings", "--print", "%(id)s"} |
| 254 | |
| 255 | args, cleanup := s.appendCookies(args) |
| 256 | defer cleanup() |
| 257 | args = append(args, url) |
| 258 | |
| 259 | out, err := util.KillableCommand(ctx, s.cfg.YTDLPPath, args...).Output() |
| 260 | if err != nil { |
| 261 | return nil, err |
| 262 | } |
| 263 | |
| 264 | keep := make(map[string]bool) |
| 265 | for _, line := range strings.Split(string(out), "\n") { |
| 266 | id := strings.TrimSpace(line) |
| 267 | // yt-dlp prints "NA" for a missing field; never treat that as a real id. |
| 268 | if id == "" || id == "NA" { |
| 269 | continue |
| 270 | } |
| 271 | keep[id] = true |
| 272 | } |
| 273 | return keep, nil |
| 274 | } |
| 275 |