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, 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. |
| 20 | func (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. |
| 75 | func (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. |
| 112 | func 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 | |
| 141 | func 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. |
| 150 | func (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. |
| 212 | func (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. |
| 241 | func (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 |