Files
DepressedCat 5f2aad00f4
Check extension and helper / check (push) Canceled after 0s
Publish YouTubeDL 1.4.5 and website on Gitea
2026-10-10 23:53:17 +03:00

400 lines
21 KiB
Python

"""Chrome/Edge native messaging host; no listening network port."""
from __future__ import annotations
import concurrent.futures
import json
import os
import shutil
import subprocess
import sys
import threading
import time
import uuid
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from host.common import EXTENSION_ID, VERSION, validate_download, youtube_url
from host.media import NO_WINDOW, ROOT, capabilities
from host.protocol import read_message, write_message
from host.session import validate_cookies
from host.bundle import history_path, installation_root
from host.processes import ProcessTree
class NativeHost:
def __init__(self, output):
self.output = output
self.output_lock = threading.Lock()
self.state_lock = threading.RLock()
self.capability_lock = threading.Lock()
self.executor = concurrent.futures.ThreadPoolExecutor(max_workers=4)
self.history_path = history_path() if installation_root() else ROOT / ".runtime" / "completed-jobs.json"
self.jobs = self.read_history()
self.processes = {}
self.closed = False
self._capabilities = None
self.picker_lock = threading.Lock()
self.auto_timer = None
self.auto_candidates = None
def read_history(self):
try:
history = json.loads(self.history_path.read_text(encoding="utf-8"))
if not isinstance(history, list):
return {}
return {job["id"]: job for job in history[:30] if isinstance(job, dict)
and isinstance(job.get("id"), str) and job.get("status") == "complete"
and isinstance(job.get("result"), dict) and isinstance(job["result"].get("path"), str)
and Path(job["result"]["path"]).is_absolute()}
except (OSError, ValueError):
return {}
def save_history(self):
# Keep completed file references so Reveal still works after a browser restart.
history = sorted((dict(job) for job in self.jobs.values() if job["status"] == "complete" and job.get("result", {}).get("path")),
key=lambda job: job.get("createdAt", 0), reverse=True)[:30]
temporary = self.history_path.with_suffix(".tmp")
try:
self.history_path.parent.mkdir(parents=True, exist_ok=True)
temporary.write_text(json.dumps(history, ensure_ascii=False), encoding="utf-8")
os.replace(temporary, self.history_path)
except OSError:
pass # Saving a history entry must not fail an otherwise completed download.
def send(self, message):
with self.output_lock:
if not self.closed:
try:
write_message(self.output, message)
except (BrokenPipeError, OSError):
self.closed = True
def reply(self, request_id, result=None, error=None):
self.send({"id": request_id, "ok": error is None, "result": result, "error": error})
def hardware(self):
with self.capability_lock:
if self._capabilities is None:
self._capabilities = capabilities()
return self._capabilities
def event(self, job_id, **changes):
with self.state_lock:
state = self.jobs[job_id]
if state.get("status") in {"cancelled", "complete", "error"} and changes.get("status") != state.get("status"):
return
if state.get("cancelling") and changes.get("status") == "complete":
return
if state.get("cancelling") and changes.get("status") is None and changes.get("stage") != "cancelling":
return
state.update(changes)
state["updatedAt"] = time.time()
if changes.get("status") == "complete":
self.save_history()
snapshot = dict(state)
self.send({"event": "job", "job": snapshot})
def kill(self, process):
tree = getattr(process, "ytdl_process_tree", None)
if isinstance(tree, ProcessTree):
tree.terminate()
return
if process.poll() is not None:
return
if os.name == "nt":
subprocess.run(["taskkill", "/PID", str(process.pid), "/T", "/F"], stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL, creationflags=NO_WINDOW, timeout=15)
else:
import signal
os.killpg(process.pid, signal.SIGTERM)
def worker(self, action, payload, job_id=None):
hardware = self.hardware() if action == "download" else {}
encoders = {key: hardware.get(key, []) for key in ("nvenc", "amf", "amf10bit")}
# CREATE_SUSPENDED + job assignment happen before even the Python
# launcher can create children. stdin is a second gate for the task.
process = ProcessTree.launch([sys.executable, "-X", "utf8", "-B", "-u", str(ROOT / "host" / "worker.py")],
stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE,
text=True, encoding="utf-8", errors="replace", creationflags=NO_WINDOW,
start_new_session=os.name != "nt")
token = job_id or str(uuid.uuid4())
with self.state_lock:
self.processes[token] = process
cancelled = self.closed or (job_id and (self.jobs[job_id]["status"] == "cancelled" or self.jobs[job_id].get("cancelling")))
if cancelled:
self.kill(process)
process.ytdl_process_tree.close()
with self.state_lock:
self.processes.pop(token, None)
return None
try:
process.stdin.write(json.dumps({"action": action, "payload": payload, "encoders": encoders}) + "\n")
process.stdin.close()
except BaseException:
self.kill(process)
process.ytdl_process_tree.close()
with self.state_lock:
self.processes.pop(token, None)
raise
errors = []
def stderr_reader():
for line in process.stderr:
errors.append(line)
if len(errors) > 40:
del errors[:20]
reader = threading.Thread(target=stderr_reader, daemon=True)
reader.start()
result, error, workspaces = None, None, []
workspace_parent = Path(payload.get("workFolder") or payload["folder"]).resolve() if action == "download" else (ROOT / ".runtime").resolve()
if action == "helper_update": workspace_parent = (installation_root() / "releases").resolve()
workspace_parents = {workspace_parent}
if action == "download": workspace_parents.add(Path(payload["folder"]).resolve())
try:
for line in process.stdout:
message = json.loads(line)
if message["type"] == "progress" and job_id:
self.event(job_id, **{k: v for k, v in message.items() if k != "type"})
elif message["type"] == "result":
result = message["result"]
elif message["type"] == "error":
error = message["message"]
elif message["type"] == "workspace" and action in {"download", "dependency_install", "helper_update"}:
candidate = Path(message["path"])
if candidate.parent.resolve() in workspace_parents and candidate.name.startswith(".ytdl-") and not candidate.is_symlink():
if candidate not in workspaces: workspaces.append(candidate)
process.wait()
reader.join(timeout=2)
if job_id and (self.jobs[job_id]["status"] == "cancelled" or self.jobs[job_id].get("cancelling")):
return None
if error or process.returncode or result is None:
raise ValueError(error or "Локальный помощник завершился с ошибкой. " + "".join(errors)[-1000:])
return result
finally:
if process.poll() is None:
self.kill(process)
process.ytdl_process_tree.close()
for stream in (process.stdin, process.stdout, process.stderr):
stream.close()
with self.state_lock:
self.processes.pop(token, None)
for workspace in workspaces:
if not workspace.exists() or workspace.is_symlink() or workspace.resolve().parent not in workspace_parents:
continue
retained = False
if action == "helper_update":
marker = json.loads((installation_root() / "installation.json").read_text(encoding="utf-8-sig"))
retained = Path(marker["current"]).resolve() == workspace.resolve()
if not retained: shutil.rmtree(workspace)
def run_job(self, job_id, payload, action="download"):
try:
try:
result = self.worker(action, payload, job_id)
finally:
if action in {"dependency_install", "helper_update"}:
with self.capability_lock:
self._capabilities = None
if result is not None:
self.event(job_id, status="complete", stage="complete", percent=100, result=result)
except Exception as error:
with self.state_lock:
cancelled = self.jobs[job_id]["status"] == "cancelled" or self.jobs[job_id].get("cancelling")
if not cancelled:
self.event(job_id, status="error", stage="error", error=str(error)[:2000])
def start_update(self, candidates=None):
if not installation_root(): raise ValueError("Для обновления нужен автономный помощник.")
with self.state_lock:
if self.processes or any(job["status"] == "running" for job in self.jobs.values()):
raise ValueError("Дождитесь завершения текущих операций перед обновлением.")
job_id = str(uuid.uuid4())
self.jobs[job_id] = {"id": job_id, "status": "running", "stage": "updating", "kind": "update", "title": "Обновление компонентов",
"percent": None, "createdAt": time.time(), "canCancel": False}
self.event(job_id)
self.executor.submit(self.run_job, job_id, {"candidates": candidates}, "helper_update")
return {"jobId": job_id}
def schedule_updates(self, delay=30):
if self.closed or not installation_root() or self.auto_timer is not None: return
self.auto_timer = threading.Timer(delay, self.auto_check)
self.auto_timer.daemon = True; self.auto_timer.start()
def auto_check(self):
from host.updates import check_updates, due, settings
try:
with self.state_lock:
if self.closed or self.processes or self._capabilities and self._capabilities.get("ready") is False or any(job["status"] == "running" for job in self.jobs.values()): return
if not settings()["enabled"]: self.auto_candidates = None; return
if self.auto_candidates is None:
if not due(): return
self.auto_candidates = check_updates()
settings(lastCheck=time.time())
with self.state_lock:
if self.closed or self.processes or any(job["status"] == "running" for job in self.jobs.values()): return
if self.auto_candidates: self.start_update(self.auto_candidates)
self.auto_candidates = None
except Exception:
# Quietly retry a failed automatic metadata check in an hour.
try: settings(lastCheck=time.time() - 86400 + 3600)
except OSError: pass
finally:
self.auto_timer = None
self.schedule_updates(60)
def pick_folder(self, current):
from host.folders import pick_folder
if not self.picker_lock.acquire(blocking=False):
raise ValueError("Диалог выбора папки уже открыт.")
try:
initial = current if isinstance(current, str) and Path(current).is_dir() else str(Path.home())
folder = pick_folder(initial)
return {"folder": folder or None}
finally:
self.picker_lock.release()
def handle(self, request):
request_id = request.get("id")
if not isinstance(request_id, str) or len(request_id) > 100:
return
action, payload = request.get("action"), request.get("payload", {})
try:
if not isinstance(payload, dict):
raise ValueError("Некорректное сообщение.")
if action == "hello":
from host.updates import settings
result = {**self.hardware(), "supportsBrowserCookies": True, "supportsComponentUpdates": bool(installation_root()), "updateSettings": settings()}
self.schedule_updates()
elif action == "ping":
result = {"alive": True, "pid": os.getpid(), "version": VERSION}
elif action == "update_settings":
from host.updates import settings
if not isinstance(payload.get("enabled"), bool): raise ValueError("Некорректная настройка обновлений.")
result = settings(payload["enabled"])
elif action == "helper_update":
result = self.start_update()
elif action == "helper_uninstall":
with self.state_lock:
if any(job["status"] == "running" for job in self.jobs.values()):
raise ValueError("Дождитесь завершения загрузок перед удалением помощника.")
from host.installation import unregister
result = unregister(ROOT)
installed = installation_root()
if installed and (installed / "YouTubeDL-Manager.exe").is_file():
# The uninstaller outlives the native port. Its output must
# never inherit the browser's native-messaging pipes.
log = installed / "data" / "uninstall.log"
log.parent.mkdir(parents=True, exist_ok=True)
with log.open("wb") as output:
subprocess.Popen([str(installed / "YouTubeDL-Manager.exe"), "--uninstall", "--silent", "--wait-pid", str(os.getpid())],
stdin=subprocess.DEVNULL, stdout=output, stderr=subprocess.STDOUT,
close_fds=True, creationflags=NO_WINDOW)
result["uninstallStarted"] = True
elif action == "inspect":
result = self.worker("inspect", {"url": youtube_url(payload.get("url")), "browserCookies": validate_cookies(payload.get("browserCookies"))})
elif action == "pick_folder":
result = self.pick_folder(payload.get("current"))
elif action == "download":
payload = {**validate_download(payload), "browserCookies": validate_cookies(payload.get("browserCookies"))}
with self.state_lock:
if any(j["status"] == "running" and j.get("kind") in {"dependency", "update"} for j in self.jobs.values()):
raise ValueError("Дождитесь завершения установки компонентов.")
if sum(j["status"] == "running" for j in self.jobs.values()) >= 2:
raise ValueError("Уже выполняются две загрузки. Дождитесь завершения одной.")
job_id = str(uuid.uuid4())
self.jobs[job_id] = {"id": job_id, "url": payload["url"], "status": "running",
"stage": "preparing", "percent": None, "createdAt": time.time()}
self.reply(request_id, {"jobId": job_id})
self.executor.submit(self.run_job, job_id, payload)
return
elif action == "dependency_install":
from host.dependencies import LABELS
key = payload.get("dependency")
if key not in {"node", "ffmpeg", "yt_dlp", "ejs"}:
raise ValueError("Неизвестный компонент.")
with self.state_lock:
if any(job["status"] == "running" for job in self.jobs.values()):
raise ValueError("Дождитесь завершения текущих операций перед установкой компонентов.")
job_id = str(uuid.uuid4())
self.jobs[job_id] = {"id": job_id, "status": "running", "stage": "installing", "kind": "dependency",
"title": LABELS[key], "percent": None, "createdAt": time.time()}
self.reply(request_id, {"jobId": job_id})
self.executor.submit(self.run_job, job_id, {"dependency": key}, action)
return
elif action == "jobs":
with self.state_lock:
result = [dict(job) for job in self.jobs.values()]
elif action == "cancel":
job_id = payload.get("jobId")
with self.state_lock:
job = self.jobs.get(job_id)
if not job or job["status"] != "running":
raise ValueError("Эта загрузка уже завершена.")
if job.get("canCancel") is False:
raise ValueError("Дождитесь завершения обновления компонентов.")
self.event(job_id, stage="cancelling", percent=None, cancelling=True, canCancel=False,
speed="", eta=None, message="Останавливаю процессы загрузки и обработки")
process = self.processes.get(job_id)
if process:
try:
self.kill(process)
except Exception as error:
self.event(job_id, status="error", stage="error", cancelling=False,
error="Не удалось остановить процессы: " + str(error))
raise
self.event(job_id, status="cancelled", stage="cancelled", percent=None, cancelling=False)
result = {"cancelled": True}
elif action == "reveal":
with self.state_lock:
job = self.jobs.get(payload.get("jobId"))
if not job or job["status"] != "complete":
raise ValueError("Файл ещё не готов.")
path = Path(job["result"]["path"])
if not path.is_file():
raise ValueError("Файл был перемещён или удалён.")
if os.name == "nt":
subprocess.Popen(["explorer.exe", "/select,", str(path)], creationflags=NO_WINDOW)
result = {"shown": True}
else:
raise ValueError("Неизвестная операция.")
self.reply(request_id, result)
except Exception as error:
self.reply(request_id, error=str(error)[:2000])
def close(self):
self.closed = True
if self.auto_timer: self.auto_timer.cancel()
with self.state_lock:
processes = list(self.processes.values())
for process in processes:
try:
self.kill(process)
except (OSError, subprocess.TimeoutExpired):
pass
self.executor.shutdown(wait=True, cancel_futures=True)
def main():
if len(sys.argv) > 1 and sys.argv[1] != f"chrome-extension://{EXTENSION_ID}/":
print("Недопустимый источник native messaging.", file=sys.stderr)
return 1
if os.name == "nt":
import msvcrt
msvcrt.setmode(sys.stdin.fileno(), os.O_BINARY)
msvcrt.setmode(sys.stdout.fileno(), os.O_BINARY)
host = NativeHost(sys.stdout.buffer)
try:
while (request := read_message(sys.stdin.buffer)) is not None:
if request.get("action") in {"cancel", "ping"}:
host.handle(request)
else:
host.executor.submit(host.handle, request)
except (EOFError, ValueError, OSError) as error:
print(str(error), file=sys.stderr)
finally:
host.close()
return 0
if __name__ == "__main__":
raise SystemExit(main())