1#!/usr/bin/env python3
2import sys
3import json
4import os
5import queue
6import re
7import smtplib
8import subprocess
9import threading
10import time
11import urllib.request
12import xml.etree.ElementTree as ET
13from email.message import EmailMessage
14from email.utils import formatdate
15from html import unescape
16from xml.sax.saxutils import escape
17
18import yaml
19
20os.umask(0o007)
21
22YT_CONFIG_DIR = os.environ.get("STUDIO_YT_CONFIG", "/srv/clover/Documents/Config/Youtube Downloader")
23FEED_CONFIG = os.path.join(YT_CONFIG_DIR, "feed.yaml")
24STATE_DIR = os.environ.get("STUDIO_YT_STATE", "/srv/prod/ytdl/data")
25INTERVAL = int(os.environ.get("CHECK_INTERVAL", "1800"))
26SMTP_HOST = os.environ.get("SMTP_HOST", "")
27SMTP_PORT = int(os.environ.get("SMTP_PORT", "465"))
28SMTP_USER = os.environ.get("SMTP_USER", "")
29SMTP_PASS = os.environ.get("SMTP_PASS", "")
30MAIL_FROM = os.environ.get("MAIL_FROM", f"yt-feed@{os.environ.get('STUDIO_DOMAIN', 'studio.test')}")
31MAIL_TO = os.environ.get("MAIL_TO", os.environ.get("STUDIO_OWNER_EMAIL", "account@paperclover.net"))
32BASE_URL = os.environ.get("BASE_URL", f"https://snowglobe.{os.environ.get('STUDIO_DOMAIN', 'studio.test')}/youtube").rstrip("/")
33READ_ONLY = os.environ.get("STUDIO_MEDIA_READ_ONLY") == "true"
34
35MEDIA_ROOT = os.environ.get("STUDIO_YT_MEDIA", "/srv/clover/Media")
36INDIE_DIR = os.path.join(MEDIA_ROOT, "Indie Shows")
37INDEP_DIR = os.path.join(MEDIA_ROOT, "Videos/Independent")
38MUSIC_DIR = os.path.join(MEDIA_ROOT, "Intake - Music")
39VIDEO_EXTS = (".webm", ".mp4", ".mkv")
40
41ATOM = "{http://www.w3.org/2005/Atom}"
42YT = "{http://www.youtube.com/xml/schemas/2015}"
43MEDIA = "{http://search.yahoo.com/mrss/}"
44UA = {"User-Agent": "Mozilla/5.0 (The Snow Globe; +https://paperclover.net)"}
45SEEN_CAP = 300
46
47state_lock = threading.Lock()
48job_queue = queue.Queue()
49shows_cache = []
50
51# RSS keeps working during a YouTube bot wall, so probing can resume downloads later.
52WALL_RE = re.compile(r"confirm you.re not a bot", re.I)
53WALL_PROBE_INTERVAL = int(os.environ.get("WALL_PROBE_INTERVAL", "10800"))
54PROBE_VIDEO = "https://www.youtube.com/watch?v=jNQXAC9IVRw"
55
56
57class WallError(Exception):
58 pass
59
60
61def wall_active():
62 return load_json("wall.json", {}).get("walled", False)
63
64
65def set_wall(walled):
66 with state_lock:
67 save_json("wall.json", {"walled": walled, "since": int(time.time())})
68 log(f"bot wall {'detected — downloads paused' if walled else 'lifted — downloads resumed'}")
69
70
71def wall_prober():
72 while True:
73 time.sleep(WALL_PROBE_INTERVAL if wall_active() else 600)
74 if not wall_active():
75 continue
76 probe = subprocess.run(
77 ["yt-dlp", "--simulate", "--print", "%(id)s", PROBE_VIDEO],
78 capture_output=True, text=True, timeout=120)
79 if probe.returncode == 0:
80 set_wall(False)
81 resolve_stragglers()
82 else:
83 log("wall probe: still walled")
84
85
86def resolve_stragglers():
87 with state_lock:
88 pending = load_json("pending.json", {})
89 for v in pending.values():
90 if v.get("unresolved"):
91 threading.Thread(target=resolve_and_update, args=(v["id"], v["link"]),
92 daemon=True).start()
93
94
95def log(msg):
96 print(msg, file=sys.stderr, flush=True)
97
98
99def fetch(url):
100 req = urllib.request.Request(url, headers=UA)
101 with urllib.request.urlopen(req, timeout=30) as resp:
102 return resp.read().decode("utf-8", errors="replace")
103
104
105def state_path(name):
106 return os.path.join(STATE_DIR, name)
107
108
109def load_json(name, fallback):
110 try:
111 with open(state_path(name)) as f:
112 return json.load(f)
113 except (FileNotFoundError, json.JSONDecodeError):
114 return fallback
115
116
117def save_json(name, data):
118 tmp = state_path(name) + ".tmp"
119 with open(tmp, "w") as f:
120 json.dump(data, f, indent=1)
121 os.replace(tmp, state_path(name))
122
123
124def safe_name(name):
125 return re.sub(r'[/\\:*?"<>|]', "-", name).strip(" .") or "untitled"
126
127
128def stable_key(v):
129 return v.get("stub_id") or v["id"]
130
131
132def safe_listdir(p):
133 try:
134 return os.listdir(p)
135 except OSError:
136 return []
137
138
139def nfo_fields(file):
140 try:
141 with open(file) as stream:
142 text = stream.read()
143 except OSError:
144 return {name: "" for name in ("title", "plot")}
145 return {name: unescape(match.group(1)).strip() if (match := re.search(
146 rf"<{name}>(.*?)</{name}>", text, re.S)) else "" for name in ("title", "plot")}
147
148
149def library_entries():
150 entries = []
151 for show in sorted(safe_listdir(INDIE_DIR)):
152 root = os.path.join(INDIE_DIR, show)
153 if not os.path.isdir(root):
154 continue
155 for sub in sorted(safe_listdir(root)):
156 season = re.fullmatch(r"Season (\d+)", sub)
157 if not season:
158 continue
159 folder = os.path.join(root, sub)
160 for name in sorted(safe_listdir(folder)):
161 if not name.lower().endswith(VIDEO_EXTS):
162 continue
163 stem = os.path.splitext(name)[0]
164 episode = re.match(r"S\d+E(\d+) - (.*)", stem)
165 nfo = nfo_fields(os.path.join(folder, stem + ".nfo"))
166 entries.append({
167 "type": "indie", "path": os.path.relpath(os.path.join(folder, name), MEDIA_ROOT),
168 "context": f"{show} · S{int(season.group(1))}E{int(episode.group(1)) if episode else 1}",
169 "title": nfo["title"] or (episode.group(2) if episode else stem),
170 "link": nfo["plot"] if nfo["plot"].startswith(("http://", "https://")) else "",
171 "season": int(season.group(1)), "episode": int(episode.group(1)) if episode else 1,
172 })
173 for channel in sorted(safe_listdir(INDEP_DIR)):
174 folder = os.path.join(INDEP_DIR, channel)
175 if not os.path.isdir(folder):
176 continue
177 for name in sorted(safe_listdir(folder)):
178 if not name.lower().endswith(VIDEO_EXTS):
179 continue
180 stem = os.path.splitext(name)[0]
181 dated = re.match(r"(\d{4}-\d{2}-\d{2}) - (.*)", stem)
182 nfo = nfo_fields(os.path.join(folder, stem + ".nfo"))
183 entries.append({
184 "type": "independent", "path": os.path.relpath(os.path.join(folder, name), MEDIA_ROOT),
185 "context": f"{channel} · {dated.group(1) if dated else ''}",
186 "title": nfo["title"] or (dated.group(2) if dated else stem),
187 "link": nfo["plot"] if nfo["plot"].startswith(("http://", "https://")) else "",
188 "season": None, "episode": None,
189 })
190 return entries
191
192
193def rename_entry(args):
194 full = os.path.realpath(os.path.join(MEDIA_ROOT, args["path"]))
195 indie = os.path.commonpath((full, INDIE_DIR)) == INDIE_DIR
196 independent = os.path.commonpath((full, INDEP_DIR)) == INDEP_DIR
197 if not (indie or independent) or not os.path.isfile(full):
198 raise ValueError("Video is outside the managed YouTube libraries")
199 title = args["title"].strip()
200 if not title:
201 raise ValueError("Video title is empty")
202 folder, name = os.path.split(full)
203 stem, ext = os.path.splitext(name)
204 if ext.lower() not in VIDEO_EXTS:
205 raise ValueError("Video format is unsupported")
206 if indie:
207 season, episode = args["season"], args["episode"]
208 if not isinstance(season, int) or not isinstance(episode, int) or min(season, episode) < 1:
209 raise ValueError("Season and episode must be positive numbers")
210 new_stem = f"S{season:02d}E{episode:02d} - {safe_name(title)}"
211 updates = {"title": title, "season": season, "episode": episode}
212 else:
213 date = re.match(r"(\d{4}-\d{2}-\d{2}) - ", stem)
214 new_stem = (date.group(1) + " - " if date else "") + safe_name(title)
215 updates = {"title": title}
216 sidecars = [name for name in safe_listdir(folder) if name.startswith(stem)
217 and (name[len(stem):].startswith(".") or name[len(stem):].startswith("-thumb"))]
218 moves = [(os.path.join(folder, name), os.path.join(folder, new_stem + name[len(stem):])) for name in sidecars]
219 if any(src != dst and os.path.exists(dst) for src, dst in moves):
220 raise FileExistsError("A file already has that title")
221 for src, dst in moves:
222 if src != dst:
223 os.rename(src, dst)
224 nfo_file = os.path.join(folder, new_stem + ".nfo")
225 try:
226 with open(nfo_file) as stream:
227 text = stream.read()
228 except FileNotFoundError:
229 return None
230 for key, value in updates.items():
231 encoded = f"<{key}>{escape(str(value))}</{key}>"
232 text = re.sub(rf"<{key}>.*?</{key}>", lambda _: encoded, text, count=1, flags=re.S) if re.search(
233 rf"<{key}>.*?</{key}>", text, re.S) else re.sub(
234 r"(\n</[A-Za-z]+>\s*)\Z", lambda match: f"\n {encoded}{match.group(1)}", text)
235 with open(nfo_file, "w") as stream:
236 stream.write(text)
237 return None
238
239
240def channel_id_for(url, cache):
241 if url in cache:
242 return cache[url]
243 m = re.search(r"/channel/(UC[0-9A-Za-z_-]{22})", url)
244 if not m:
245 html = fetch(url)
246 m = re.search(r"channel_id=(UC[0-9A-Za-z_-]{22})", html) or re.search(
247 r'"channelId":"(UC[0-9A-Za-z_-]{22})"', html
248 )
249 if not m:
250 raise ValueError(f"could not resolve channel id for {url}")
251 cache[url] = m.group(1)
252 return cache[url]
253
254
255def feed_entries(channel_id):
256 text = fetch(f"https://www.youtube.com/feeds/videos.xml?channel_id={channel_id}")
257 entries = []
258 for e in ET.fromstring(text).findall(ATOM + "entry"):
259 vid = e.find(YT + "videoId")
260 title = e.find(ATOM + "title")
261 link = e.find(ATOM + "link")
262 published = e.find(ATOM + "published")
263 thumb = e.find(f"{MEDIA}group/{MEDIA}thumbnail")
264 if vid is None or vid.text is None or title is None:
265 continue
266 entries.append({
267 "id": vid.text,
268 "title": title.text or "(untitled)",
269 "link": link.get("href") if link is not None else f"https://youtu.be/{vid.text}",
270 "published": published.text[:10] if published is not None and published.text else "",
271 "thumb": thumb.get("url") if thumb is not None else "",
272 })
273 return entries
274
275
276def send_notification(channel, entry):
277 msg = EmailMessage()
278 slug = re.sub(r"[^a-z0-9]+", "-", channel.lower()).strip("-")
279 review = BASE_URL
280 msg["Subject"] = f"[yt] {channel}: {entry['title']}"
281 msg["From"] = MAIL_FROM
282 msg["To"] = MAIL_TO
283 msg["Date"] = formatdate(localtime=True)
284 msg["Message-ID"] = f"<yt-{entry['id']}@yt-feed.paperclover.net>"
285 msg["References"] = f"<yt-channel-{slug}@yt-feed.paperclover.net>"
286 msg["In-Reply-To"] = f"<yt-channel-{slug}@yt-feed.paperclover.net>"
287 msg.set_content(
288 f"{channel} uploaded: {entry['title']}\n\n"
289 f" {entry['link']}\n published {entry['published']}\n\n"
290 f"review and ingest: {review}\n"
291 )
292 h = escape(entry["title"])
293 msg.add_alternative(
294 f'<div style="font-family:sans-serif">'
295 f'<p><b>{escape(channel)}</b> uploaded:</p>'
296 f'<p><a href="{review or entry["link"]}">'
297 f'<img src="{entry["thumb"]}" alt="" width="320" style="display:block;border-radius:8px"></a></p>'
298 f'<p><a href="{review or entry["link"]}">{h}</a> · {entry["published"]}</p>'
299 f'<p><a href="{review}">review &amp; ingest →</a> · <a href="{entry["link"]}">watch</a></p>'
300 f"</div>",
301 subtype="html",
302 )
303 with smtplib.SMTP_SSL(SMTP_HOST, SMTP_PORT, timeout=30) as s:
304 s.login(SMTP_USER, SMTP_PASS)
305 s.send_message(msg)
306
307
308def poll_once():
309 with open(FEED_CONFIG) as f:
310 channels = (yaml.safe_load(f) or {}).get("channels") or {}
311 with state_lock:
312 ids = load_json("channel-ids.json", {})
313 seen = load_json("seen.json", {})
314 pending = load_json("pending.json", {})
315 for name, url in channels.items():
316 try:
317 cid = channel_id_for(url, ids)
318 entries = feed_entries(cid)
319 except Exception as e:
320 log(f"{name}: fetch failed: {e}")
321 continue
322 if cid not in seen:
323 seen[cid] = [e["id"] for e in entries]
324 log(f"{name}: now tracking ({len(entries)} existing videos skipped)")
325 continue
326 known = set(seen[cid])
327 for entry in reversed(entries):
328 if entry["id"] in known:
329 continue
330 entry["channel"] = name
331 pending[entry["id"]] = entry
332 seen[cid].append(entry["id"])
333 known.add(entry["id"])
334 log(f"{name}: queued {entry['title']!r}")
335 try:
336 send_notification(name, entry)
337 except Exception as e:
338 log(f"{name}: notification failed: {e}")
339 seen[cid] = seen[cid][-SEEN_CAP:]
340 with state_lock:
341 save_json("channel-ids.json", ids)
342 save_json("seen.json", seen)
343 save_json("pending.json", pending)
344
345
346def poller():
347 while True:
348 try:
349 poll_once()
350 except Exception as e:
351 log(f"poll failed: {e}")
352 time.sleep(INTERVAL)
353
354
355def update_job(job_id, **fields):
356 with state_lock:
357 jobs = load_json("jobs.json", [])
358 for j in jobs:
359 if j["id"] == job_id:
360 j.update(fields)
361 save_json("jobs.json", jobs)
362
363
364def resolve_url(url):
365 out = subprocess.run(
366 ["yt-dlp", "--no-playlist", "--print",
367 "%(id)s\t%(title)s\t%(channel)s\t%(upload_date>%Y-%m-%d)s\t%(thumbnail)s", url],
368 capture_output=True, text=True, timeout=90)
369 if out.returncode != 0:
370 if WALL_RE.search(out.stderr):
371 raise WallError(url)
372 raise ValueError(out.stderr[-300:])
373 vid, title, channel, published, thumb = out.stdout.strip().split("\t")
374 return {"id": vid, "title": title, "channel": channel, "published": published,
375 "thumb": thumb, "link": f"https://www.youtube.com/watch?v={vid}"}
376
377
378def resolve_and_update(stub_id, url):
379 try:
380 entry = resolve_url(url)
381 except WallError:
382 set_wall(True)
383 return
384 except Exception as e:
385 log(f"resolve failed for {url}: {e}")
386 return
387 with state_lock:
388 pending = load_json("pending.json", {})
389 # if the stub is gone the user already ingested or skipped it
390 if stub_id in pending:
391 del pending[stub_id]
392 entry["stub_id"] = stub_id
393 pending[entry["id"]] = entry
394 save_json("pending.json", pending)
395
396
397def write_nfo(filepath, root, tags):
398 base, _ = os.path.splitext(filepath)
399 body = "\n".join(f" <{k}>{escape(str(v))}</{k}>" for k, v in tags.items() if v != "")
400 with open(base + ".nfo", "w") as f:
401 f.write(f"<?xml version='1.0' encoding='utf-8'?>\n<{root}>\n{body}\n</{root}>\n")
402
403
404def run_job(job):
405 for delay in (0, 30, 90):
406 if delay:
407 update_job(job["id"], status="retrying", progress="")
408 time.sleep(delay)
409 try:
410 if run_job_once(job):
411 return
412 except WallError:
413 raise
414 except Exception as e:
415 log(f"job {job['id']} attempt failed: {e}")
416 update_job(job["id"], status="error", progress="")
417
418
419def run_job_once(job):
420 if job.get("unresolved"):
421 update_job(job["id"], status="resolving")
422 meta = resolve_url(job["url"])
423 job.update(title=meta["title"], channel=meta["channel"],
424 published=meta["published"], url=meta["link"], unresolved=False)
425 update_job(job["id"], title=meta["title"])
426 dest, url = job["dest"], job["url"]
427 if dest == "indie":
428 outdir = os.path.join(INDIE_DIR, safe_name(job["show"]), f"Season {job['season']}")
429 prefix = f"S{job['season']:02d}E{job['episode']:02d}"
430 ep_name = safe_name(job.get("ep_title") or job["title"])
431 out = f"{outdir}/{prefix} - {ep_name}.%(ext)s"
432 thumb_out = f"thumbnail:{outdir}/{prefix} - {ep_name}-thumb.%(ext)s"
433 elif dest == "independent":
434 outdir = os.path.join(INDEP_DIR, safe_name(job["channel"]))
435 out = f"{outdir}/%(upload_date>%Y-%m-%d)s - %(title)s.%(ext)s"
436 thumb_out = f"thumbnail:{outdir}/%(upload_date>%Y-%m-%d)s - %(title)s.%(ext)s"
437 else:
438 outdir = os.path.join(MUSIC_DIR, safe_name(job["channel"]))
439 out = f"{outdir}/%(title)s.%(ext)s"
440 thumb_out = None
441 os.makedirs(outdir, exist_ok=True)
442 cmd = ["yt-dlp", "--newline", "--no-playlist", "--embed-chapters",
443 "--sleep-requests", "0.75",
444 "--print", "after_move:filepath", "--no-simulate", "-o", out]
445 if dest == "music":
446 cmd += ["-x"]
447 else:
448 cmd += ["-f", "bv*+ba/b", "--write-thumbnail", "--convert-thumbnails", "jpg",
449 "-o", thumb_out]
450 cmd.append(url)
451 update_job(job["id"], status="downloading")
452 filepath = None
453 tail = []
454 proc = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True)
455 for line in proc.stdout:
456 line = line.rstrip()
457 tail = (tail + [line])[-30:]
458 m = re.search(r"\[download\]\s+([\d.]+%)", line)
459 if m:
460 update_job(job["id"], progress=m.group(1))
461 elif line.startswith("/"):
462 filepath = line
463 proc.wait()
464 if proc.returncode != 0 or (dest != "music" and not filepath):
465 if WALL_RE.search("\n".join(tail)):
466 raise WallError(url)
467 log(f"job {job['id']} attempt failed (exit {proc.returncode})")
468 return False
469 if dest == "indie":
470 write_nfo(filepath, "episodedetails", {
471 "title": job.get("ep_title") or job["title"],
472 "season": job["season"], "episode": job["episode"],
473 "aired": job.get("published", ""), "plot": url,
474 })
475 elif dest == "independent":
476 write_nfo(filepath, "movie", {
477 "title": job["title"], "premiered": job.get("published", ""), "plot": url,
478 })
479 update_job(job["id"], status="done", progress="")
480 log(f"job {job['id']} done: {filepath or job['title']}")
481 return True
482
483
484def worker():
485 while True:
486 job = job_queue.get()
487 if wall_active():
488 update_job(job["id"], status="waiting", progress="")
489 while wall_active():
490 time.sleep(30)
491 try:
492 run_job(job)
493 except WallError:
494 set_wall(True)
495 update_job(job["id"], status="queued", progress="")
496 job_queue.put(job)
497 except Exception as e:
498 update_job(job["id"], status="error")
499 log(f"job {job['id']} crashed: {e}")
500
501
502def writable():
503 if READ_ONLY:
504 raise PermissionError("Media is read-only on this host")
505
506
507def subscriptions():
508 try:
509 with open(os.path.join(YT_CONFIG_DIR, "subscriptions.yaml")) as f:
510 data = yaml.safe_load(f) or {}
511 except FileNotFoundError:
512 data = {}
513 preset = next((key for key in data if key != "__preset__"), "Jellyfin TV Show by Date | only-after | flat-videos")
514 sections = data.get(preset) or {}
515 section = next(iter(sections), "= Independent Creators")
516 channels = sections.get(section) or {}
517 result = []
518 fields = ("download_after", "title_include_keywords", "title_exclude_keywords",
519 "description_include_keywords", "description_exclude_keywords")
520 for name, value in channels.items():
521 if isinstance(value, str):
522 result.append({"name": name.lstrip("~"), "url": value, "rules": {}})
523 elif isinstance(value, dict):
524 result.append({"name": name.lstrip("~"), "url": value.get("url", ""),
525 "rules": {key: value[key] for key in fields if key in value}})
526 return data, preset, section, result
527
528
529def command(name, args):
530 if name == "snapshot":
531 with state_lock:
532 pending = list(load_json("pending.json", {}).values())
533 jobs = load_json("jobs.json", [])
534 wall = load_json("wall.json", {})
535 upscale = load_json("upscale-status.json", {})
536 enabled = load_json("upscale-enabled.json", {"enabled": True})["enabled"]
537 return {
538 "pending": [{"key": stable_key(v), "title": v["title"], "channel": v.get("channel", ""),
539 "published": v.get("published", ""), "duration": None,
540 "thumb": v.get("thumb"), "link": v["link"], "resolving": bool(v.get("unresolved"))}
541 for v in pending],
542 "shows": shows_cache,
543 "jobs": [{"id": j["id"], "title": j["title"], "channel": j.get("channel", ""),
544 "destination": j.get("dest_label", j.get("dest", "")),
545 "status": j["status"], "url": j["url"],
546 "folder": job_folder(j), "progress": progress(j.get("progress")),
547 "queuedAt": int(j["id"][1:]) / 1000} for j in jobs],
548 "wall": {"walled": bool(wall.get("walled")),
549 "nextProbe": wall.get("since", 0) + WALL_PROBE_INTERVAL if wall.get("walled") else None},
550 "archive": {"started": None, "finished": None, "walled": False, "next": 0},
551 "upscaler": {"enabled": enabled, "running": bool(upscale.get("running")),
552 "done": upscale.get("done", 0), "total": upscale.get("total", 0),
553 "current": upscale.get("current", ""), "errors": upscale.get("errors", 0)},
554 }
555 if name == "channels":
556 try:
557 with open(FEED_CONFIG) as f:
558 notify = (yaml.safe_load(f) or {}).get("channels") or {}
559 except FileNotFoundError:
560 notify = {}
561 return {"notify": [{"name": key, "url": value} for key, value in notify.items()],
562 "archive": subscriptions()[3]}
563 if name == "configs":
564 files = []
565 for filename in sorted(safe_listdir(YT_CONFIG_DIR)):
566 if filename.endswith((".yaml", ".yml")):
567 with open(os.path.join(YT_CONFIG_DIR, filename)) as f:
568 files.append({"name": filename, "body": f.read()})
569 return files
570 if name == "saveConfig":
571 filename = args["name"]
572 if filename not in safe_listdir(YT_CONFIG_DIR) or not filename.endswith((".yaml", ".yml")):
573 raise ValueError("Unknown YouTube config file")
574 target = os.path.join(YT_CONFIG_DIR, filename)
575 with open(target) as f:
576 if f.read() != args["original"]:
577 raise ValueError("Config changed since it was opened. Reload before saving")
578 yaml.safe_load(args["body"])
579 tmp = target + ".tmp"
580 with open(tmp, "w") as f:
581 f.write(args["body"])
582 os.replace(tmp, target)
583 return None
584 writable()
585 if name == "rename":
586 return rename_entry(args)
587 if name == "add":
588 with state_lock:
589 pending = load_json("pending.json", {})
590 for i, url in enumerate(args["urls"]):
591 key = f"u{int(time.time() * 1000)}{i}"
592 pending[key] = {"id": key, "title": url, "channel": "", "published": "",
593 "thumb": "", "link": url, "unresolved": True}
594 threading.Thread(target=resolve_and_update, args=(key, url), daemon=True).start()
595 save_json("pending.json", pending)
596 return None
597 if name in ("ingest", "skip"):
598 with state_lock:
599 pending = load_json("pending.json", {})
600 key = next((key for key, v in pending.items() if stable_key(v) == args["key"]), None)
601 if key is None:
602 raise KeyError("Video is no longer in the queue")
603 v = pending.pop(key)
604 save_json("pending.json", pending)
605 if name == "skip":
606 return None
607 choice = args["choice"]
608 job = {"id": f"j{int(time.time() * 1000)}", "url": v["link"], "title": v["title"],
609 "channel": v.get("channel", "unknown"), "published": v.get("published", ""),
610 "unresolved": v.get("unresolved", False), "status": "queued", "progress": "",
611 "dest": choice["dest"]}
612 if choice["dest"] == "indie":
613 job.update(show=choice["show"], season=choice["season"], episode=choice["episode"],
614 ep_title=choice["title"],
615 dest_label=f"{choice['show']} S{choice['season']:02d}E{choice['episode']:02d}")
616 else:
617 job["dest_label"] = "Independent" if choice["dest"] == "independent" else "Intake - Music"
618 jobs = load_json("jobs.json", [])
619 jobs.insert(0, job)
620 save_json("jobs.json", jobs[:50])
621 job_queue.put(job)
622 return None
623 if name == "retry":
624 with state_lock:
625 jobs = load_json("jobs.json", [])
626 job = next((j for j in jobs if j["id"] == args["id"]), None)
627 if not job:
628 raise KeyError("Download is no longer in history")
629 job.update(status="queued", progress="")
630 save_json("jobs.json", jobs)
631 job_queue.put(job)
632 return None
633 if name == "setUpscaler":
634 with state_lock:
635 save_json("upscale-enabled.json", {"enabled": args["enabled"]})
636 return None
637 if name == "setChannels":
638 try:
639 with open(FEED_CONFIG) as f:
640 feed = yaml.safe_load(f) or {}
641 except FileNotFoundError:
642 feed = {}
643 feed["channels"] = {ch["name"]: ch["url"] for ch in args["notify"]}
644 data, preset, section, _ = subscriptions()
645 fields = ("download_after", "title_include_keywords", "title_exclude_keywords",
646 "description_include_keywords", "description_exclude_keywords")
647 archive = {}
648 for ch in args["archive"]:
649 rules = {key: value for key in fields if (value := ch["rules"].get(key))}
650 archive[("~" if rules else "") + ch["name"]] = {"url": ch["url"], **rules} if rules else ch["url"]
651 sections = data.get(preset) or {}
652 sections[section] = archive
653 data[preset] = sections
654 for filename, body in (("feed.yaml", feed), ("subscriptions.yaml", data)):
655 target = os.path.join(YT_CONFIG_DIR, filename)
656 tmp = target + ".tmp"
657 with open(tmp, "w") as f:
658 yaml.safe_dump(body, f, sort_keys=False, allow_unicode=True)
659 os.replace(tmp, target)
660 return None
661 raise ValueError(f"Unknown YouTube operation: {name}")
662
663
664def job_folder(job):
665 if job["dest"] == "indie":
666 return os.path.join(INDIE_DIR, safe_name(job["show"]), f"Season {job['season']}")
667 if job["dest"] == "independent":
668 return os.path.join(INDEP_DIR, safe_name(job.get("channel", "unknown")))
669 return os.path.join(MUSIC_DIR, safe_name(job.get("channel", "unknown")))
670
671
672def progress(value):
673 try:
674 return float(str(value).rstrip("%")) / 100
675 except ValueError:
676 return None
677
678
679def indie_shows():
680 shows = {}
681 for name in sorted(safe_listdir(INDIE_DIR)):
682 path = os.path.join(INDIE_DIR, name)
683 if not os.path.isdir(path):
684 continue
685 seasons = {}
686 for sub in safe_listdir(path):
687 match = re.fullmatch(r"Season (\d+)", sub)
688 if match and os.path.isdir(os.path.join(path, sub)):
689 seasons[int(match.group(1))] = sum(f.lower().endswith(VIDEO_EXTS)
690 for f in safe_listdir(os.path.join(path, sub)))
691 shows[name] = seasons
692 return shows
693
694
695def refresh_shows():
696 global shows_cache
697 while True:
698 try:
699 shows = indie_shows()
700 shows_cache = [{"name": name, "seasons": [{"number": n, "episodes": count}
701 for n, count in seasons.items()]} for name, seasons in shows.items()]
702 except Exception as error:
703 log(f"show scan failed: {error}")
704 time.sleep(300)
705
706
707if __name__ == "__main__":
708 if sys.argv[1:] == ["--library"]:
709 print(json.dumps(library_entries()), flush=True)
710 sys.exit(0)
711 threading.Thread(target=refresh_shows, daemon=True).start()
712 if not READ_ONLY:
713 with state_lock:
714 for job in reversed(load_json("jobs.json", [])):
715 if job.get("status") in ("queued", "resolving", "downloading", "retrying"):
716 job_queue.put(job)
717 for loop in (poller, worker, wall_prober):
718 threading.Thread(target=loop, daemon=True).start()
719 for line in sys.stdin:
720 try:
721 request = json.loads(line)
722 result = command(request["method"], request.get("args", {}))
723 response = {"id": request["id"], "result": result}
724 except Exception as error:
725 response = {"id": request.get("id") if "request" in locals() else None,
726 "error": str(error)}
727 print(json.dumps(response), flush=True)