400 lines
21 KiB
Python
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())
|