| | |
| | | import os |
| | | import queue |
| | | import re |
| | | import stat |
| | | import subprocess |
| | | import sys |
| | | import threading |
| | | import time |
| | | import uuid |
| | | from pathlib import Path |
| | | from typing import Any, BinaryIO |
| | | |
| | | if __package__ in (None, ""): |
| | | sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) |
| | | from bili_authenticated_extension.constants import ( # type: ignore[import-not-found] |
| | | COMPLETION_CLOSURE_REQUIRED, |
| | | COOKIE_ACCESS_TERMINAL_CODES, |
| | | EXTENSION_BUILD, |
| | | EXPECTED_ORIGIN, |
| | | GRACEFUL_CANCEL_SECONDS, |
| | | JOB_WAIT_SECONDS, |
| | | MAX_INPUT_FRAME, |
| | | METADATA_TIMEOUT_SECONDS, |
| | | TARGET_BVID, |
| | | RELOAD_GENERATION, |
| | | THREAD_JOIN_SECONDS, |
| | | validate_postprocess_terminal, |
| | | validate_runtime_diagnostic, |
| | | ) |
| | | from bili_authenticated_extension.job import ( # type: ignore[import-not-found] |
| | | WindowsJob, |
| | |
| | | ) |
| | | from bili_authenticated_extension.protocol import ( # type: ignore[import-not-found] |
| | | ProtocolError, |
| | | maintenance_state, |
| | | read_frame, |
| | | safe_response, |
| | | strict_json_loads, |
| | | validate_message, |
| | | validate_origin_argv, |
| | | validate_media_complete_ack, |
| | | validate_media_complete_identity, |
| | | validate_worker_prepare, |
| | | write_frame, |
| | | ) |
| | | from bili_authenticated_extension.queue_state import QueueStore, ReloadStore # type: ignore[import-not-found] |
| | | else: |
| | | from .constants import ( |
| | | COMPLETION_CLOSURE_REQUIRED, |
| | | COOKIE_ACCESS_TERMINAL_CODES, |
| | | EXTENSION_BUILD, |
| | | EXPECTED_ORIGIN, |
| | | GRACEFUL_CANCEL_SECONDS, |
| | | JOB_WAIT_SECONDS, |
| | | MAX_INPUT_FRAME, |
| | | METADATA_TIMEOUT_SECONDS, |
| | | TARGET_BVID, |
| | | RELOAD_GENERATION, |
| | | THREAD_JOIN_SECONDS, |
| | | validate_postprocess_terminal, |
| | | validate_runtime_diagnostic, |
| | | ) |
| | | from .job import ( |
| | | WindowsJob, |
| | |
| | | ) |
| | | from .protocol import ( |
| | | ProtocolError, |
| | | maintenance_state, |
| | | read_frame, |
| | | safe_response, |
| | | strict_json_loads, |
| | | validate_message, |
| | | validate_origin_argv, |
| | | validate_media_complete_ack, |
| | | validate_media_complete_identity, |
| | | validate_worker_prepare, |
| | | write_frame, |
| | | ) |
| | | from .queue_state import QueueStore, ReloadStore |
| | | |
| | | |
| | | def _config_path() -> Path: |
| | |
| | | return digest.hexdigest().upper() |
| | | |
| | | |
| | | def preflight_configuration(path: Path) -> str | None: |
| | | """Perform stdlib-only collision/config checks before any Cookie is read.""" |
| | | _RECOVERY_MODE_NONE = "NONE" |
| | | _RECOVERY_MODE_EXACT_PAIR = "EXACT_PUBLISHED_PAIR" |
| | | _FOREGROUND_ERROR_CODES = frozenset({ |
| | | "E_FOREGROUND_LOCKED", |
| | | "E_FOREGROUND_WINDOW_ABSENT", |
| | | "E_FOREGROUND_WINDOW_AMBIGUOUS", |
| | | "E_FOREGROUND_LAUNCH_FAILED", |
| | | "E_FOREGROUND_PLATFORM_UNSUPPORTED", |
| | | }) |
| | | _RECOVERABLE_MAPPING_KEYS = frozenset({ |
| | | "schema_version", "bvid", "source", "published_at", "title", "local_file", |
| | | "bytes", "sha256", "duration_seconds", "remote_duration_seconds", |
| | | "local_duration_seconds", "duration_delta_seconds", "duration_tolerance_seconds", |
| | | "format_name", "video_codec", "audio_codec", "completed_at", "acquisition_mode", |
| | | "handoff_source_sha256", |
| | | }) |
| | | |
| | | |
| | | def _native_parent_window(arguments: list[str]) -> int: |
| | | for argument in arguments[1:]: |
| | | if argument.startswith("--parent-window="): |
| | | try: |
| | | value = int(argument.split("=", 1)[1], 10) |
| | | except ValueError: |
| | | return 0 |
| | | return value if value > 0 else 0 |
| | | return 0 |
| | | |
| | | |
| | | class _WindowsKnownFolderId(ctypes.Structure): |
| | | _fields_ = [ |
| | | ("Data1", ctypes.c_uint32), |
| | | ("Data2", ctypes.c_uint16), |
| | | ("Data3", ctypes.c_uint16), |
| | | ("Data4", ctypes.c_ubyte * 8), |
| | | ] |
| | | |
| | | @classmethod |
| | | def from_text(cls, value: str) -> "_WindowsKnownFolderId": |
| | | return cls.from_buffer_copy(uuid.UUID(value).bytes_le) |
| | | |
| | | |
| | | class _WindowsChromeForegroundAdapter: |
| | | """A title-free, profile-free adapter for one bounded Chrome foreground action.""" |
| | | |
| | | _SW_RESTORE = 9 |
| | | _GA_ROOT = 2 |
| | | _PROCESS_QUERY_LIMITED_INFORMATION = 0x1000 |
| | | _CHROME_RELATIVE_PATH = ("Google", "Chrome", "Application", "chrome.exe") |
| | | _KNOWN_FOLDER_IDS = ( |
| | | "905e63b6-c1bf-494e-b29c-65b732d3d21a", # ProgramFiles |
| | | "7c5a40ef-a0fb-4bfc-874a-c0f2e0b9fa8e", # ProgramFilesX86 |
| | | "f1b32785-6fba-4fcf-9d55-7b8e7f157091", # LocalAppData |
| | | ) |
| | | |
| | | def __init__(self) -> None: |
| | | if os.name != "nt": |
| | | raise OSError("unsupported platform") |
| | | from ctypes import wintypes |
| | | |
| | | self._wintypes = wintypes |
| | | self._user32 = ctypes.WinDLL("user32", use_last_error=True) |
| | | self._kernel32 = ctypes.WinDLL("kernel32", use_last_error=True) |
| | | self._shell32 = ctypes.WinDLL("shell32", use_last_error=True) |
| | | self._ole32 = ctypes.WinDLL("ole32", use_last_error=True) |
| | | self._user32.IsWindow.argtypes = [wintypes.HWND] |
| | | self._user32.IsWindow.restype = wintypes.BOOL |
| | | self._user32.IsWindowVisible.argtypes = [wintypes.HWND] |
| | | self._user32.IsWindowVisible.restype = wintypes.BOOL |
| | | self._user32.GetAncestor.argtypes = [wintypes.HWND, wintypes.UINT] |
| | | self._user32.GetAncestor.restype = wintypes.HWND |
| | | self._user32.GetClassNameW.argtypes = [wintypes.HWND, wintypes.LPWSTR, ctypes.c_int] |
| | | self._user32.GetClassNameW.restype = ctypes.c_int |
| | | self._user32.GetWindowRect.argtypes = [wintypes.HWND, ctypes.POINTER(wintypes.RECT)] |
| | | self._user32.GetWindowRect.restype = wintypes.BOOL |
| | | self._user32.GetWindowThreadProcessId.argtypes = [wintypes.HWND, ctypes.POINTER(wintypes.DWORD)] |
| | | self._user32.GetWindowThreadProcessId.restype = wintypes.DWORD |
| | | self._user32.ShowWindowAsync.argtypes = [wintypes.HWND, ctypes.c_int] |
| | | self._user32.ShowWindowAsync.restype = wintypes.BOOL |
| | | self._user32.BringWindowToTop.argtypes = [wintypes.HWND] |
| | | self._user32.BringWindowToTop.restype = wintypes.BOOL |
| | | self._user32.SetForegroundWindow.argtypes = [wintypes.HWND] |
| | | self._user32.SetForegroundWindow.restype = wintypes.BOOL |
| | | self._user32.GetForegroundWindow.argtypes = [] |
| | | self._user32.GetForegroundWindow.restype = wintypes.HWND |
| | | self._kernel32.OpenProcess.argtypes = [wintypes.DWORD, wintypes.BOOL, wintypes.DWORD] |
| | | self._kernel32.OpenProcess.restype = wintypes.HANDLE |
| | | self._kernel32.QueryFullProcessImageNameW.argtypes = [ |
| | | wintypes.HANDLE, wintypes.DWORD, wintypes.LPWSTR, ctypes.POINTER(wintypes.DWORD), |
| | | ] |
| | | self._kernel32.QueryFullProcessImageNameW.restype = wintypes.BOOL |
| | | self._kernel32.CloseHandle.argtypes = [wintypes.HANDLE] |
| | | self._kernel32.CloseHandle.restype = wintypes.BOOL |
| | | self._shell32.SHGetKnownFolderPath.argtypes = [ |
| | | ctypes.POINTER(_WindowsKnownFolderId), wintypes.DWORD, |
| | | wintypes.HANDLE, ctypes.POINTER(ctypes.c_void_p), |
| | | ] |
| | | self._shell32.SHGetKnownFolderPath.restype = ctypes.c_long |
| | | self._ole32.CoTaskMemFree.argtypes = [ctypes.c_void_p] |
| | | self._ole32.CoTaskMemFree.restype = None |
| | | |
| | | def _root(self, hwnd: int) -> int: |
| | | if hwnd <= 0 or not self._user32.IsWindow(hwnd): |
| | | return 0 |
| | | return int(self._user32.GetAncestor(hwnd, self._GA_ROOT) or hwnd) |
| | | |
| | | def _class_name(self, hwnd: int) -> str: |
| | | buffer = ctypes.create_unicode_buffer(256) |
| | | length = int(self._user32.GetClassNameW(hwnd, buffer, len(buffer))) |
| | | return buffer.value[:length] if length > 0 else "" |
| | | |
| | | def _image_is_chrome(self, hwnd: int) -> bool: |
| | | process_id = self._wintypes.DWORD(0) |
| | | self._user32.GetWindowThreadProcessId(hwnd, ctypes.byref(process_id)) |
| | | if process_id.value <= 0: |
| | | return False |
| | | handle = self._kernel32.OpenProcess( |
| | | self._PROCESS_QUERY_LIMITED_INFORMATION, False, process_id.value, |
| | | ) |
| | | if not handle: |
| | | return False |
| | | try: |
| | | size = self._wintypes.DWORD(32768) |
| | | buffer = ctypes.create_unicode_buffer(size.value) |
| | | if not self._kernel32.QueryFullProcessImageNameW(handle, 0, buffer, ctypes.byref(size)): |
| | | return False |
| | | return Path(buffer.value).name.casefold() == "chrome.exe" |
| | | finally: |
| | | self._kernel32.CloseHandle(handle) |
| | | |
| | | def is_chrome_window(self, hwnd: int) -> bool: |
| | | root = self._root(hwnd) |
| | | return bool( |
| | | root > 0 |
| | | and self._user32.IsWindowVisible(root) |
| | | and self._class_name(root) == "Chrome_WidgetWin_1" |
| | | and self._image_is_chrome(root) |
| | | ) |
| | | |
| | | def list_windows(self) -> list[int]: |
| | | windows: list[int] = [] |
| | | callback_type = ctypes.WINFUNCTYPE( |
| | | self._wintypes.BOOL, self._wintypes.HWND, self._wintypes.LPARAM, |
| | | ) |
| | | |
| | | @callback_type |
| | | def collect(hwnd: int, _parameter: int) -> bool: |
| | | value = int(hwnd) |
| | | if self.is_chrome_window(value): |
| | | windows.append(self._root(value)) |
| | | return True |
| | | |
| | | self._user32.EnumWindows(collect, 0) |
| | | return sorted(set(windows)) |
| | | |
| | | def bounds(self, hwnd: int) -> dict[str, int] | None: |
| | | rectangle = self._wintypes.RECT() |
| | | if not self._user32.GetWindowRect(hwnd, ctypes.byref(rectangle)): |
| | | return None |
| | | return { |
| | | "left": int(rectangle.left), "top": int(rectangle.top), |
| | | "width": int(rectangle.right - rectangle.left), |
| | | "height": int(rectangle.bottom - rectangle.top), |
| | | } |
| | | |
| | | def focus(self, hwnd: int) -> bool: |
| | | root = self._root(hwnd) |
| | | if not self.is_chrome_window(root): |
| | | return False |
| | | self._user32.ShowWindowAsync(root, self._SW_RESTORE) |
| | | self._user32.BringWindowToTop(root) |
| | | requested = bool(self._user32.SetForegroundWindow(root)) |
| | | deadline = time.monotonic() + 2.0 |
| | | while time.monotonic() < deadline: |
| | | foreground = self._root(int(self._user32.GetForegroundWindow() or 0)) |
| | | if foreground == root: |
| | | return True |
| | | time.sleep(0.05) |
| | | return requested and self._root(int(self._user32.GetForegroundWindow() or 0)) == root |
| | | |
| | | def _known_folder_roots(self) -> tuple[Path, ...]: |
| | | roots: list[Path] = [] |
| | | observed: set[str] = set() |
| | | for text in self._KNOWN_FOLDER_IDS: |
| | | folder_id = _WindowsKnownFolderId.from_text(text) |
| | | allocated = ctypes.c_void_p() |
| | | try: |
| | | result = int(self._shell32.SHGetKnownFolderPath( |
| | | ctypes.byref(folder_id), 0, None, ctypes.byref(allocated), |
| | | )) |
| | | if result != 0 or not allocated.value: |
| | | continue |
| | | root = Path(ctypes.wstring_at(allocated.value)) |
| | | key = os.path.normcase(os.path.abspath(os.fspath(root))) |
| | | if key not in observed: |
| | | observed.add(key) |
| | | roots.append(root) |
| | | finally: |
| | | if allocated.value: |
| | | self._ole32.CoTaskMemFree(allocated) |
| | | return tuple(roots) |
| | | |
| | | @staticmethod |
| | | def _local_absolute(path: Path) -> Path | None: |
| | | try: |
| | | value = Path(os.path.abspath(os.fspath(path))) |
| | | except (OSError, TypeError, ValueError): |
| | | return None |
| | | if not value.is_absolute() or value.anchor.startswith("\\\\"): |
| | | return None |
| | | if not re.fullmatch(r"[A-Za-z]:", value.drive): |
| | | return None |
| | | return value |
| | | |
| | | @classmethod |
| | | def _validated_chrome_candidate( |
| | | cls, root: Path, |
| | | ) -> tuple[Path, Path, tuple[tuple[str, int, int, int, int, int, int], ...]] | None: |
| | | lexical_root = cls._local_absolute(root) |
| | | if lexical_root is None: |
| | | return None |
| | | lexical_candidate = lexical_root.joinpath(*cls._CHROME_RELATIVE_PATH) |
| | | try: |
| | | canonical_root = lexical_root.resolve(strict=True) |
| | | canonical_candidate = lexical_candidate.resolve(strict=True) |
| | | except OSError: |
| | | return None |
| | | if ( |
| | | os.path.normcase(os.fspath(lexical_root)) |
| | | != os.path.normcase(os.fspath(canonical_root)) |
| | | or os.path.normcase(os.fspath(lexical_candidate)) |
| | | != os.path.normcase(os.fspath(canonical_candidate)) |
| | | ): |
| | | return None |
| | | try: |
| | | relative = canonical_candidate.relative_to(canonical_root) |
| | | except ValueError: |
| | | return None |
| | | if tuple(part.casefold() for part in relative.parts) != tuple( |
| | | part.casefold() for part in cls._CHROME_RELATIVE_PATH |
| | | ): |
| | | return None |
| | | |
| | | chain: list[Path] = [] |
| | | current = Path(canonical_root.anchor) |
| | | chain.append(current) |
| | | for part in canonical_root.parts[1:]: |
| | | current /= part |
| | | chain.append(current) |
| | | for part in cls._CHROME_RELATIVE_PATH: |
| | | current /= part |
| | | chain.append(current) |
| | | snapshots: list[tuple[str, int, int, int, int, int, int]] = [] |
| | | for index, item in enumerate(chain): |
| | | try: |
| | | identity = item.lstat() |
| | | if item.is_symlink() or _is_reparse(item): |
| | | return None |
| | | final = index == len(chain) - 1 |
| | | if final and not stat.S_ISREG(identity.st_mode): |
| | | return None |
| | | if not final and not stat.S_ISDIR(identity.st_mode): |
| | | return None |
| | | snapshots.append(( |
| | | os.path.normcase(os.fspath(item)), |
| | | int(identity.st_dev), int(identity.st_ino), int(identity.st_mode), |
| | | int(identity.st_size), int(identity.st_mtime_ns), int(identity.st_ctime_ns), |
| | | )) |
| | | except OSError: |
| | | return None |
| | | return canonical_candidate, canonical_root, tuple(snapshots) |
| | | |
| | | def _chrome_executable( |
| | | self, |
| | | ) -> tuple[Path, Path, tuple[tuple[str, int, int, int, int, int, int], ...]] | None: |
| | | for root in self._known_folder_roots(): |
| | | validated = self._validated_chrome_candidate(root) |
| | | if validated is not None: |
| | | return validated |
| | | return None |
| | | |
| | | def launch(self, canonical_target_url: str, prior_windows: set[int]) -> int | None: |
| | | discovered = self._chrome_executable() |
| | | if discovered is None: |
| | | return None |
| | | executable, trusted_root, identity = discovered |
| | | revalidated = self._validated_chrome_candidate(trusted_root) |
| | | if revalidated is None: |
| | | return None |
| | | current_executable, current_root, current_identity = revalidated |
| | | if ( |
| | | os.path.normcase(os.fspath(current_executable)) |
| | | != os.path.normcase(os.fspath(executable)) |
| | | or os.path.normcase(os.fspath(current_root)) |
| | | != os.path.normcase(os.fspath(trusted_root)) |
| | | or current_identity != identity |
| | | ): |
| | | return None |
| | | try: |
| | | subprocess.Popen( |
| | | [str(current_executable), "--new-window", canonical_target_url], |
| | | stdin=subprocess.DEVNULL, |
| | | stdout=subprocess.DEVNULL, |
| | | stderr=subprocess.DEVNULL, |
| | | close_fds=True, |
| | | ) |
| | | except OSError: |
| | | return None |
| | | deadline = time.monotonic() + 10.0 |
| | | while time.monotonic() < deadline: |
| | | created = [item for item in self.list_windows() if item not in prior_windows] |
| | | if len(created) == 1: |
| | | return created[0] |
| | | if len(created) > 1: |
| | | return -1 |
| | | time.sleep(0.1) |
| | | return None |
| | | |
| | | |
| | | def _window_bounds_match(observed: dict[str, int] | None, expected: dict[str, int]) -> bool: |
| | | if observed is None: |
| | | return False |
| | | return all(abs(observed[key] - expected[key]) <= 16 for key in ("left", "top", "width", "height")) |
| | | |
| | | |
| | | def _foreground_chrome_window( |
| | | canonical_target_url: str, |
| | | expected_bounds: dict[str, int], |
| | | parent_window: int, |
| | | *, |
| | | adapter: Any | None = None, |
| | | ) -> str | None: |
| | | """Restore/foreground one exact Chrome window or launch one visible target. |
| | | |
| | | Returns a fixed sanitized error code, or ``None`` on success. No title, |
| | | command line, URL query, profile, page, or secret data is inspected. |
| | | """ |
| | | |
| | | if adapter is None: |
| | | try: |
| | | adapter = _WindowsChromeForegroundAdapter() |
| | | except OSError: |
| | | return "E_FOREGROUND_PLATFORM_UNSUPPORTED" |
| | | windows = adapter.list_windows() |
| | | parent = adapter._root(parent_window) if parent_window > 0 else 0 |
| | | if parent > 0 and parent in windows and adapter.is_chrome_window(parent): |
| | | candidates = [parent] |
| | | else: |
| | | candidates = [ |
| | | hwnd for hwnd in windows if _window_bounds_match(adapter.bounds(hwnd), expected_bounds) |
| | | ] |
| | | if len(candidates) > 1: |
| | | return "E_FOREGROUND_WINDOW_AMBIGUOUS" |
| | | if len(candidates) == 1: |
| | | return None if adapter.focus(candidates[0]) else "E_FOREGROUND_LOCKED" |
| | | if windows: |
| | | return "E_FOREGROUND_WINDOW_ABSENT" |
| | | launched = adapter.launch(canonical_target_url, set(windows)) |
| | | if launched == -1: |
| | | return "E_FOREGROUND_WINDOW_AMBIGUOUS" |
| | | if launched is None: |
| | | return "E_FOREGROUND_LAUNCH_FAILED" |
| | | return None if adapter.focus(launched) else "E_FOREGROUND_LOCKED" |
| | | |
| | | |
| | | def _is_reparse(path: Path) -> bool: |
| | | try: |
| | | return bool(path.lstat().st_file_attributes & 0x400) |
| | | except AttributeError: |
| | | return path.is_symlink() |
| | | |
| | | |
| | | def _ordinary_child(path: Path, parent: Path) -> os.stat_result: |
| | | value = path.lstat() |
| | | if ( |
| | | not stat.S_ISREG(value.st_mode) |
| | | or path.is_symlink() |
| | | or _is_reparse(path) |
| | | or path.resolve(strict=True).parent != parent.resolve(strict=True) |
| | | ): |
| | | raise OSError("non-ordinary destination child") |
| | | return value |
| | | |
| | | |
| | | def _read_small_ordinary_json(path: Path, parent: Path) -> dict[str, Any]: |
| | | before = _ordinary_child(path, parent) |
| | | if before.st_size <= 0 or before.st_size > 64 * 1024: |
| | | raise OSError("mapping size") |
| | | descriptor: int | None = None |
| | | handle: int | None = None |
| | | try: |
| | | if os.name == "nt": |
| | | import msvcrt # noqa: PLC0415 |
| | | |
| | | kernel32 = ctypes.WinDLL("kernel32", use_last_error=True) |
| | | create_file = kernel32.CreateFileW |
| | | create_file.argtypes = ( |
| | | ctypes.c_wchar_p, ctypes.c_uint32, ctypes.c_uint32, ctypes.c_void_p, |
| | | ctypes.c_uint32, ctypes.c_uint32, ctypes.c_void_p, |
| | | ) |
| | | create_file.restype = ctypes.c_void_p |
| | | handle = create_file( |
| | | str(path), 0x80000000, 0x00000001, None, 3, |
| | | 0x00000080 | 0x00200000 | 0x08000000, None, |
| | | ) |
| | | if handle in (None, ctypes.c_void_p(-1).value): |
| | | raise OSError(ctypes.get_last_error(), "mapping open") |
| | | descriptor = msvcrt.open_osfhandle( |
| | | int(handle), os.O_RDONLY | getattr(os, "O_BINARY", 0) |
| | | ) |
| | | handle = None |
| | | else: |
| | | flags = os.O_RDONLY | getattr(os, "O_BINARY", 0) | getattr(os, "O_NOFOLLOW", 0) |
| | | descriptor = os.open(path, flags) |
| | | with os.fdopen(descriptor, "rb", closefd=True) as source: |
| | | descriptor = None |
| | | handle_stat = os.fstat(source.fileno()) |
| | | payload = source.read(64 * 1024 + 1) |
| | | after = _ordinary_child(path, parent) |
| | | finally: |
| | | if descriptor is not None: |
| | | os.close(descriptor) |
| | | if handle is not None: |
| | | ctypes.WinDLL("kernel32", use_last_error=True).CloseHandle(ctypes.c_void_p(handle)) |
| | | identity = lambda value: (value.st_dev, value.st_ino, value.st_mode, value.st_size) |
| | | if ( |
| | | len(payload) > 64 * 1024 or payload == b"" |
| | | or identity(before) != identity(handle_stat) or identity(before) != identity(after) |
| | | ): |
| | | raise OSError("mapping drift") |
| | | value = strict_json_loads(payload) |
| | | if not isinstance(value, dict): |
| | | raise OSError("mapping root") |
| | | return value |
| | | |
| | | |
| | | def _recovery_stage_is_clear(bvid: str) -> bool: |
| | | local = os.environ.get("LOCALAPPDATA") |
| | | if not local: |
| | | return False |
| | | local_root = Path(local) |
| | | if not local_root.is_absolute() or str(local_root).startswith("\\\\"): |
| | | return False |
| | | stage = local_root / "project-info" / "bili-auth-ingress" / bvid |
| | | if not stage.exists(): |
| | | return True |
| | | return stage.is_dir() and not stage.is_symlink() and not _is_reparse(stage) and next(stage.iterdir(), None) is None |
| | | |
| | | |
| | | def _classify_job_destination(config: dict[str, Any], job: dict[str, Any]) -> tuple[str | None, str]: |
| | | """Classify only an exact published pair as a recovery candidate. |
| | | |
| | | This gate is deliberately not a success decision. It only lets the |
| | | worker's locked, handle-backed verifier observe the pair before any Cookie |
| | | or network access. Any ambiguity remains a collision. |
| | | """ |
| | | try: |
| | | destination = Path(config["destination"]).resolve(strict=True) |
| | | prefix = f"{job['bvid']}.".casefold() |
| | | candidates = [child for child in destination.iterdir() if child.name.casefold().startswith(prefix)] |
| | | lineage = job.get("lineage") |
| | | closure_required = ( |
| | | isinstance(lineage, dict) |
| | | and lineage.get("predecessor_terminal_error_code") == COMPLETION_CLOSURE_REQUIRED |
| | | ) |
| | | if not candidates: |
| | | if closure_required: |
| | | return "E_EXISTS", _RECOVERY_MODE_NONE |
| | | return None, _RECOVERY_MODE_NONE |
| | | formal_name = f"{job['bvid']}.mkv" |
| | | mapping_name = f"{job['bvid']}.download.json" |
| | | if {child.name for child in candidates} != {formal_name, mapping_name} or len(candidates) != 2: |
| | | return "E_EXISTS", _RECOVERY_MODE_NONE |
| | | formal = destination / formal_name |
| | | mapping = destination / mapping_name |
| | | formal_stat = _ordinary_child(formal, destination) |
| | | persisted = _read_small_ordinary_json(mapping, destination) |
| | | if ( |
| | | set(persisted) != _RECOVERABLE_MAPPING_KEYS |
| | | or persisted.get("schema_version") != "1.0" |
| | | or persisted.get("bvid") != job["bvid"] |
| | | or persisted.get("source") != job["canonical_url"] |
| | | or persisted.get("published_at") != job["published_at"] |
| | | or persisted.get("local_file") != formal_name |
| | | or persisted.get("acquisition_mode") != "authorized_browser_file_handoff" |
| | | or isinstance(persisted.get("bytes"), bool) |
| | | or persisted.get("bytes") != formal_stat.st_size |
| | | or not isinstance(persisted.get("sha256"), str) |
| | | or re.fullmatch(r"[0-9a-f]{64}", persisted["sha256"]) is None |
| | | or persisted.get("handoff_source_sha256") != persisted["sha256"] |
| | | or job.get("creator_uid") not in config.get("creator_allowlist", ()) |
| | | or not _recovery_stage_is_clear(job["bvid"]) |
| | | or _ordinary_child(formal, destination).st_size != formal_stat.st_size |
| | | ): |
| | | return "E_EXISTS", _RECOVERY_MODE_NONE |
| | | return None, _RECOVERY_MODE_EXACT_PAIR |
| | | except (KeyError, OSError, ProtocolError, TypeError, ValueError): |
| | | return "E_EXISTS", _RECOVERY_MODE_NONE |
| | | |
| | | |
| | | def load_runtime_configuration(path: Path) -> dict[str, Any]: |
| | | """Load the generic local queue/tool configuration without importing the worker.""" |
| | | expected = { |
| | | "schema", |
| | | "target", |
| | | "canonical_url", |
| | | "creator_allowlist", |
| | | "queue_path", |
| | | "queue_state_path", |
| | | "queue_lock_path", |
| | | "reload_state_path", |
| | | "reload_generation", |
| | | "required_extension_build", |
| | | "ffmpeg", |
| | | "ffmpeg_sha256", |
| | | "ffprobe", |
| | |
| | | "bridge_python_sha256", |
| | | "bridge_script", |
| | | "bridge_script_sha256", |
| | | "batch_json", |
| | | "batch_json_sha256", |
| | | "yt_dlp_executable", |
| | | "yt_dlp_executable_sha256", |
| | | "destination", |
| | | "creator_name", |
| | | "formal_manifest_path", |
| | | "processing_handoff_path", |
| | | } |
| | | raw = strict_json_loads(path.read_bytes()) |
| | | if set(raw) != expected or raw["schema"] != 2: |
| | | raise ProtocolError("E_CONFIG") |
| | | creators = raw["creator_allowlist"] |
| | | if ( |
| | | not isinstance(creators, list) or not creators or len(creators) > 64 |
| | | or any(not isinstance(item, str) or not re.fullmatch(r"[1-9][0-9]{0,19}", item) for item in creators) |
| | | or creators != sorted(set(creators)) |
| | | ): |
| | | raise ProtocolError("E_CONFIG") |
| | | if raw["required_extension_build"] != EXTENSION_BUILD or raw["reload_generation"] != RELOAD_GENERATION: |
| | | raise ProtocolError("E_CONFIG") |
| | | for name in ( |
| | | "queue_path", "queue_state_path", "queue_lock_path", "reload_state_path", |
| | | "formal_manifest_path", "processing_handoff_path", |
| | | ): |
| | | value = raw[name] |
| | | if not isinstance(value, str): |
| | | raise ProtocolError("E_CONFIG") |
| | | candidate = Path(value) |
| | | if not candidate.is_absolute() or str(candidate).startswith("\\\\"): |
| | | raise ProtocolError("E_CONFIG") |
| | | parent = candidate.parent.resolve(strict=True) |
| | | if not parent.is_dir() or parent.is_symlink() or (candidate.exists() and (not candidate.is_file() or candidate.is_symlink())): |
| | | raise ProtocolError("E_CONFIG") |
| | | creator_name = raw["creator_name"] |
| | | if ( |
| | | not isinstance(creator_name, str) or not creator_name.strip() |
| | | or len(creator_name.encode("utf-8")) > 240 |
| | | or any(ord(ch) < 32 or ord(ch) == 127 for ch in creator_name) |
| | | or not Path(raw["formal_manifest_path"]).is_file() |
| | | or raw["formal_manifest_path"] == raw["processing_handoff_path"] |
| | | ): |
| | | raise ProtocolError("E_CONFIG") |
| | | raw[name] = str(candidate.resolve(strict=False)) |
| | | try: |
| | | raw = strict_json_loads(path.read_bytes()) |
| | | if set(raw) != expected or raw["schema"] != 1: |
| | | return "E_CONFIG" |
| | | if raw["target"] != TARGET_BVID or raw["canonical_url"] != "https://www.bilibili.com/video/BV1HA3o6oEJJ": |
| | | return "E_CONFIG" |
| | | for name in ("ffmpeg", "ffprobe", "bridge_python", "bridge_script", "batch_json", "yt_dlp_executable"): |
| | | for name in ("ffmpeg", "ffprobe", "bridge_python", "bridge_script", "yt_dlp_executable"): |
| | | value = raw[name] |
| | | expected_hash = raw[f"{name}_sha256"] |
| | | if not isinstance(value, str) or not isinstance(expected_hash, str): |
| | | return "E_CONFIG" |
| | | raise ProtocolError("E_CONFIG") |
| | | candidate = Path(value) |
| | | if not candidate.is_absolute() or str(candidate).startswith("\\\\"): |
| | | return "E_CONFIG" |
| | | raise ProtocolError("E_CONFIG") |
| | | candidate = candidate.resolve(strict=True) |
| | | if not candidate.is_file() or candidate.is_symlink() or _hash_file(candidate) != expected_hash.upper(): |
| | | return "E_CONFIG_HASH" |
| | | if raw["bridge_script_sha256"].upper() != "749FC486B0F42315BD463F11771FE2A7C71CAB53DD9AC2E411CE82E1175DFF13": |
| | | return "E_CONFIG_HASH" |
| | | raise ProtocolError("E_CONFIG_HASH") |
| | | raw[name] = str(candidate) |
| | | if raw["bridge_script_sha256"].upper() != "00F11DAF8387160DB863C89F0B33AB8480422233FF42199189222C989C7ED07E": |
| | | raise ProtocolError("E_CONFIG_HASH") |
| | | destination = Path(raw["destination"]) |
| | | if not destination.is_absolute() or str(destination).startswith("\\\\"): |
| | | return "E_CONFIG" |
| | | raise ProtocolError("E_CONFIG") |
| | | destination = destination.resolve(strict=True) |
| | | if not destination.is_dir() or destination.is_symlink(): |
| | | return "E_CONFIG" |
| | | if any( |
| | | child.is_file() and child.name.casefold().startswith(f"{TARGET_BVID}.".casefold()) |
| | | for child in destination.iterdir() |
| | | ): |
| | | return "E_EXISTS" |
| | | except (OSError, ProtocolError, KeyError, TypeError, ValueError): |
| | | raise ProtocolError("E_CONFIG") |
| | | raw["destination"] = str(destination) |
| | | except OSError as exc: |
| | | raise ProtocolError("E_CONFIG") from exc |
| | | return raw |
| | | |
| | | |
| | | def preflight_configuration(path: Path) -> str | None: |
| | | try: |
| | | load_runtime_configuration(path) |
| | | except ProtocolError as exc: |
| | | return exc.code |
| | | except (OSError, KeyError, TypeError, ValueError): |
| | | return "E_CONFIG" |
| | | return None |
| | | |
| | | |
| | | def preflight_job(config: dict[str, Any], job: dict[str, Any]) -> str | None: |
| | | return _classify_job_destination(config, job)[0] |
| | | |
| | | |
| | | def _queue_store(config: dict[str, Any]) -> QueueStore: |
| | | return QueueStore( |
| | | Path(config["queue_path"]), Path(config["queue_state_path"]), |
| | | Path(config["queue_lock_path"]), frozenset(config["creator_allowlist"]), |
| | | ) |
| | | |
| | | |
| | | def _worker_command( |
| | |
| | | return int(msvcrt.get_osfhandle(fd)) |
| | | |
| | | |
| | | def _media_complete_ack_test_seam(_stage: str) -> None: |
| | | """No-op production seam for Host-exit/ACK-loss counterexamples.""" |
| | | |
| | | return None |
| | | |
| | | |
| | | class WorkerTask: |
| | | """One worker, its task job, and sanitized control channel.""" |
| | | |
| | | def __init__(self, config_path: Path | None = None) -> None: |
| | | def __init__( |
| | | self, config_path: Path | None = None, queue_store: QueueStore | None = None, |
| | | ) -> None: |
| | | self.job: WindowsJob | None = None |
| | | self.process: Any = None |
| | | self.cancel_handle = 0 |
| | | self.commit_handle = 0 |
| | | self.config_path = _config_path() if config_path is None else config_path.resolve(strict=True) |
| | | self.input_writer: BinaryIO | None = None |
| | | self.input_lock = threading.Lock() |
| | | self.control_reader: BinaryIO | None = None |
| | | self.control_queue: queue.Queue[dict[str, Any] | None] = queue.Queue() |
| | | self.control_thread: threading.Thread | None = None |
| | |
| | | self.error_code: str | None = None |
| | | self.formal_filename: str | None = None |
| | | self.mapping_filename: str | None = None |
| | | self.diagnostic: dict[str, object] | None = None |
| | | self.task_nonce: str | None = None |
| | | self.terminal = False |
| | | self.task_started_at: float | None = None |
| | |
| | | self.prepared = False |
| | | self.secret_started = False |
| | | self.prepare_id: str | None = None |
| | | self.job_spec: dict[str, Any] | None = None |
| | | self.lease_id: str | None = None |
| | | self.queue_terminal_recorded = False |
| | | self.media_complete = False |
| | | self.media_identity: dict[str, Any] | None = None |
| | | self.postprocess_recovery_binding: dict[str, Any] | None = None |
| | | self.queue_store = queue_store |
| | | self.closing = False |
| | | self.recovery_mode = _RECOVERY_MODE_NONE |
| | | |
| | | def prepare(self, page_proof: dict[str, Any], prepare_id: str) -> None: |
| | | def prepare( |
| | | self, |
| | | job_spec: dict[str, Any], |
| | | lease_id: str, |
| | | page_proof: dict[str, Any] | None, |
| | | prepare_id: str, |
| | | recovery_mode: str = _RECOVERY_MODE_NONE, |
| | | ) -> None: |
| | | if self.process is not None: |
| | | raise ProtocolError("E_BUSY") |
| | | if recovery_mode not in {_RECOVERY_MODE_NONE, _RECOVERY_MODE_EXACT_PAIR}: |
| | | raise ProtocolError("E_PREPARE") |
| | | if recovery_mode == _RECOVERY_MODE_NONE: |
| | | if not isinstance(page_proof, dict) or not isinstance(page_proof.get("task_nonce"), str): |
| | | raise ProtocolError("E_PREPARE") |
| | | self.task_nonce = page_proof["task_nonce"] |
| | | else: |
| | | # Exact-pair recovery is deliberately independent of page, browser, |
| | | # Cookie, and network state. A proof may be present for compatibility |
| | | # with an already-prepared caller, but is never required or consumed. |
| | | self.task_nonce = None |
| | | self.prepare_id = prepare_id |
| | | self.job_spec = dict(job_spec) |
| | | self.lease_id = lease_id |
| | | self.recovery_mode = recovery_mode |
| | | input_read_fd, input_write_fd = os.pipe() |
| | | control_read_fd, control_write_fd = os.pipe() |
| | | input_read_handle = _fd_handle(input_read_fd) |
| | |
| | | self.control_reader = os.fdopen(control_read_fd, "rb", buffering=0) |
| | | self.control_thread = threading.Thread(target=self._read_control, name="bili-auth-control", daemon=True) |
| | | self.control_thread.start() |
| | | with self.input_lock: |
| | | write_frame( |
| | | self.input_writer, |
| | | {"schema": 3, "type": "worker_prepare", "job": job_spec, |
| | | "lease_id": lease_id, "recovery_mode": recovery_mode}, |
| | | ) |
| | | deadline = time.monotonic() + METADATA_TIMEOUT_SECONDS |
| | | while time.monotonic() < deadline: |
| | | try: |
| | |
| | | if message is None: |
| | | break |
| | | if message.get("type") == "ready" and message.get("code") == "READY_PLUGIN_DISABLED": |
| | | self.task_nonce = page_proof["task_nonce"] |
| | | self.prepare_id = prepare_id |
| | | self.phase = "READY" |
| | | self.prepared = True |
| | | return |
| | | if message.get("type") == "ready" and message.get("code") == "E_PLUGIN_BOUNDARY": |
| | | self.error_code = "E_PLUGIN_BOUNDARY" |
| | | self.phase = "FAILED" |
| | | self.terminal = True |
| | | break |
| | | if message.get("type") == "terminal": |
| | | self._apply_control(message) |
| | | if self.phase == "COMPLETE": |
| | | return |
| | | break |
| | | self.error_code = "E_PLUGIN_BOUNDARY" |
| | | error_code = self.error_code if isinstance(self.error_code, str) and re.fullmatch(r"E_[A-Z0-9_]{1,48}", self.error_code) else "E_PLUGIN_BOUNDARY" |
| | | self.error_code = error_code |
| | | self.phase = "FAILED" |
| | | self.terminal = True |
| | | self.terminate() |
| | | raise ProtocolError("E_PLUGIN_BOUNDARY") |
| | | raise ProtocolError(error_code) |
| | | |
| | | def start(self, start_message: dict[str, Any]) -> None: |
| | | if ( |
| | |
| | | or self.input_writer is None |
| | | or start_message["page_proof"]["task_nonce"] != self.task_nonce |
| | | or start_message["prepare_id"] != self.prepare_id |
| | | or start_message["job"] != self.job_spec |
| | | or start_message["lease_id"] != self.lease_id |
| | | ): |
| | | raise ProtocolError("E_PREPARE") |
| | | write_frame(self.input_writer, start_message) |
| | | self.input_writer.close() |
| | | self.input_writer = None |
| | | with self.input_lock: |
| | | write_frame(self.input_writer, start_message) |
| | | self.secret_started = True |
| | | self.phase = "CHECKING" |
| | | self.task_started_at = time.monotonic() |
| | |
| | | if payload is None: |
| | | break |
| | | value = strict_json_loads(payload) |
| | | if _valid_control_message(value): |
| | | self.control_queue.put(value) |
| | | if _valid_control_message(value, self.job_spec, self.lease_id): |
| | | if value.get("type") == "progress" and value.get("phase") == "MEDIA_COMPLETE": |
| | | self._accept_media_complete(value) |
| | | else: |
| | | self.control_queue.put(value) |
| | | else: |
| | | self.control_queue.put({"schema": 1, "type": "terminal", "phase": "FAILED", "error_code": "E_CONTROL"}) |
| | | break |
| | | except BaseException: |
| | | self.control_queue.put({"schema": 1, "type": "terminal", "phase": "FAILED", "error_code": "E_CONTROL"}) |
| | | if not self.closing: |
| | | self.control_queue.put({"schema": 1, "type": "terminal", "phase": "FAILED", "error_code": "E_CONTROL"}) |
| | | finally: |
| | | self.control_queue.put(None) |
| | | |
| | | def _accept_media_complete(self, value: dict[str, Any]) -> None: |
| | | """Durably persist and reread exact media identity before ACKing Worker.""" |
| | | |
| | | if ( |
| | | self.queue_store is None or self.job_spec is None or self.lease_id is None |
| | | or self.input_writer is None |
| | | or value["job"] != self.job_spec or value["lease_id"] != self.lease_id |
| | | ): |
| | | raise ProtocolError("E_MEDIA_COMPLETE") |
| | | media = validate_media_complete_identity(value["media"], self.job_spec) |
| | | persisted = self.queue_store.mark_media_complete( |
| | | self.job_spec, self.lease_id, int(time.time() * 1000), media, |
| | | ) |
| | | if persisted.get("media") != media: |
| | | raise ProtocolError("E_QUEUE_WRITE") |
| | | self.media_identity = dict(media) |
| | | self.media_complete = True |
| | | self.phase = "MEDIA_COMPLETE" |
| | | self.progress = 100 |
| | | _media_complete_ack_test_seam("BEFORE_ACK") |
| | | acknowledgement = { |
| | | "schema": 3, |
| | | "type": "media_complete_ack", |
| | | "job_id": self.job_spec["job_id"], |
| | | "lease_id": self.lease_id, |
| | | "media": media, |
| | | } |
| | | validate_media_complete_ack( |
| | | acknowledgement, self.job_spec, self.lease_id, media, |
| | | ) |
| | | with self.input_lock: |
| | | if self.input_writer is None: |
| | | raise ProtocolError("E_MEDIA_COMPLETE") |
| | | write_frame(self.input_writer, acknowledgement) |
| | | _media_complete_ack_test_seam("AFTER_ACK") |
| | | self.control_queue.put({ |
| | | "schema": 1, "type": "progress", "phase": "MEDIA_COMPLETE", |
| | | "progress": 100, "job": self.job_spec, "lease_id": self.lease_id, |
| | | "media": media, |
| | | }) |
| | | |
| | | def _apply_control(self, value: dict[str, Any]) -> None: |
| | | message_type = value["type"] |
| | |
| | | self.phase_started_at = time.monotonic() |
| | | self.phase = value["phase"] |
| | | self.progress = value["progress"] |
| | | if value["phase"] == "MEDIA_COMPLETE": |
| | | self.media_complete = True |
| | | elif message_type == "terminal": |
| | | self.phase = value["phase"] |
| | | self.progress = 100 if self.phase == "COMPLETE" else self.progress |
| | | self.error_code = value.get("error_code") |
| | | self.formal_filename = value.get("formal_filename") |
| | | self.mapping_filename = value.get("mapping_filename") |
| | | self.diagnostic = value.get("diagnostic") |
| | | self.terminal = True |
| | | |
| | | def poll(self) -> None: |
| | | def _drain_control_queue(self) -> None: |
| | | while True: |
| | | try: |
| | | value = self.control_queue.get_nowait() |
| | |
| | | break |
| | | if value is not None: |
| | | self._apply_control(value) |
| | | |
| | | def poll(self) -> None: |
| | | self._drain_control_queue() |
| | | if self.process is not None and self.process.wait(0) and not self.terminal: |
| | | self.phase = "FAILED" |
| | | self.error_code = "E_WORKER_EXIT" |
| | |
| | | if self.phase == "CHECKING" and self.phase_started_at is not None: |
| | | if now - self.phase_started_at >= METADATA_TIMEOUT_SECONDS: |
| | | return "E_METADATA_TIMEOUT" |
| | | if self.phase in {"DOWNLOADING", "MERGING", "VALIDATING", "PUBLISHING"}: |
| | | if self.phase in { |
| | | "DOWNLOADING", "MERGING", "VALIDATING", "PUBLISHING", |
| | | "MEDIA_COMPLETE", "POSTPROCESS_PENDING", |
| | | }: |
| | | if now - self.task_started_at >= 7_200: |
| | | return "E_DOWNLOAD_TIMEOUT" |
| | | return None |
| | |
| | | return True |
| | | |
| | | def terminate(self) -> None: |
| | | if self.closing: |
| | | return |
| | | self._drain_control_queue() |
| | | if self.terminal: |
| | | self.close() |
| | | return |
| | | if self.cancel_handle: |
| | | try: |
| | | set_event(self.cancel_handle) |
| | | except BaseException: |
| | | pass |
| | | if self.process is not None and not self.process.wait(GRACEFUL_CANCEL_SECONDS): |
| | | if self.input_writer is not None: |
| | | writer, self.input_writer = self.input_writer, None |
| | | try: |
| | | writer.close() |
| | | except (BrokenPipeError, OSError, ValueError): |
| | | pass |
| | | deadline = time.monotonic() + GRACEFUL_CANCEL_SECONDS |
| | | while self.process is not None and time.monotonic() < deadline: |
| | | self._drain_control_queue() |
| | | if self.terminal or self.process.wait(0): |
| | | break |
| | | time.sleep(0.02) |
| | | if self.process is not None and not self.process.wait(0): |
| | | if self.job is not None: |
| | | self.job.terminate() |
| | | if not self.process.wait(JOB_WAIT_SECONDS): |
| | | self.error_code = "E_PROCESS_OWNERSHIP" |
| | | if self.control_thread is not None and self.control_thread is not threading.current_thread(): |
| | | self.control_thread.join(THREAD_JOIN_SECONDS) |
| | | self._drain_control_queue() |
| | | self.close() |
| | | |
| | | def close(self) -> None: |
| | | if self.closing: |
| | | return |
| | | self.closing = True |
| | | if self.input_writer is not None: |
| | | self.input_writer.close() |
| | | self.input_writer = None |
| | | writer, self.input_writer = self.input_writer, None |
| | | try: |
| | | writer.close() |
| | | except (BrokenPipeError, OSError, ValueError): |
| | | pass |
| | | if self.control_reader is not None: |
| | | self.control_reader.close() |
| | | self.control_reader = None |
| | | reader, self.control_reader = self.control_reader, None |
| | | try: |
| | | reader.close() |
| | | except (BrokenPipeError, OSError, ValueError): |
| | | pass |
| | | if self.control_thread is not None and self.control_thread is not threading.current_thread(): |
| | | self.control_thread.join(THREAD_JOIN_SECONDS) |
| | | self.control_thread = None |
| | | control_thread, self.control_thread = self.control_thread, None |
| | | control_thread.join(THREAD_JOIN_SECONDS) |
| | | if self.process is not None: |
| | | self.process.close() |
| | | self.process = None |
| | | process, self.process = self.process, None |
| | | try: |
| | | process.close() |
| | | except BaseException: |
| | | pass |
| | | if self.job is not None: |
| | | self.job.close() |
| | | self.job = None |
| | | job, self.job = self.job, None |
| | | try: |
| | | job.close() |
| | | except BaseException: |
| | | pass |
| | | if self.cancel_handle: |
| | | close_handles(self.cancel_handle) |
| | | self.cancel_handle = 0 |
| | | cancel_handle, self.cancel_handle = self.cancel_handle, 0 |
| | | try: |
| | | close_handles(cancel_handle) |
| | | except BaseException: |
| | | pass |
| | | if self.commit_handle: |
| | | close_handles(self.commit_handle) |
| | | self.commit_handle = 0 |
| | | commit_handle, self.commit_handle = self.commit_handle, 0 |
| | | try: |
| | | close_handles(commit_handle) |
| | | except BaseException: |
| | | pass |
| | | |
| | | |
| | | def _valid_control_message(value: dict[str, Any]) -> bool: |
| | | def _valid_control_message( |
| | | value: dict[str, Any], expected_job: dict[str, Any] | str | None = None, |
| | | expected_lease: str | None = None, |
| | | ) -> bool: |
| | | if value.get("schema") != 1 or value.get("type") not in {"ready", "progress", "terminal"}: |
| | | return False |
| | | if value["type"] == "ready": |
| | |
| | | "E_PLUGIN_BOUNDARY", |
| | | } |
| | | if value["type"] == "progress": |
| | | if value.get("phase") == "MEDIA_COMPLETE": |
| | | if not isinstance(expected_job, dict): |
| | | return False |
| | | if set(value) != { |
| | | "schema", "type", "phase", "progress", "job", "lease_id", "media", |
| | | }: |
| | | return False |
| | | if value["job"] != expected_job or value["lease_id"] != expected_lease: |
| | | return False |
| | | try: |
| | | validate_media_complete_identity(value["media"], expected_job) |
| | | except ProtocolError: |
| | | return False |
| | | return value["progress"] == 100 and not isinstance(value["progress"], bool) |
| | | return ( |
| | | set(value) == {"schema", "type", "phase", "progress"} |
| | | and value["phase"] in {"CHECKING", "DOWNLOADING", "MERGING", "VALIDATING", "PUBLISHING"} |
| | | and value["phase"] in { |
| | | "CHECKING", "DOWNLOADING", "MERGING", "VALIDATING", "PUBLISHING", |
| | | "POSTPROCESS_PENDING", |
| | | } |
| | | and isinstance(value["progress"], int) |
| | | and not isinstance(value["progress"], bool) |
| | | and 0 <= value["progress"] <= 100 |
| | |
| | | "mapping_filename", |
| | | "cookie_stream_closed", |
| | | } |
| | | if set(value) != allowed or value["phase"] not in {"COMPLETE", "FAILED", "CANCELED"}: |
| | | keys = set(value) |
| | | if keys not in (allowed, allowed | {"diagnostic"}) or value["phase"] not in { |
| | | "COMPLETE", "FAILED", "POSTPROCESS_FAILED", "CANCELED" |
| | | }: |
| | | return False |
| | | if not isinstance(value["cookie_stream_closed"], bool): |
| | | return False |
| | |
| | | not isinstance(error_code, str) or not re.fullmatch(r"E_[A-Z0-9_]{1,48}", error_code) |
| | | ): |
| | | return False |
| | | if error_code == COMPLETION_CLOSURE_REQUIRED: |
| | | return False |
| | | if "diagnostic" in value: |
| | | try: |
| | | if value["phase"] == "POSTPROCESS_FAILED": |
| | | validate_postprocess_terminal(error_code, value["diagnostic"]) |
| | | else: |
| | | validate_runtime_diagnostic(value["diagnostic"]) |
| | | except ValueError: |
| | | return False |
| | | if value["phase"] not in {"FAILED", "POSTPROCESS_FAILED"}: |
| | | return False |
| | | elif value["phase"] == "POSTPROCESS_FAILED": |
| | | try: |
| | | validate_postprocess_terminal(error_code, None) |
| | | except ValueError: |
| | | return False |
| | | if value["phase"] == "COMPLETE": |
| | | return ( |
| | | expected_job is not None |
| | | and |
| | | error_code is None |
| | | and value["cookie_stream_closed"] is True |
| | | and value["formal_filename"] == f"{TARGET_BVID}.mkv" |
| | | and value["mapping_filename"] == f"{TARGET_BVID}.download.json" |
| | | and value["formal_filename"] == f"{expected_job['bvid'] if isinstance(expected_job, dict) else expected_job}.mkv" |
| | | and value["mapping_filename"] == f"{expected_job['bvid'] if isinstance(expected_job, dict) else expected_job}.download.json" |
| | | ) |
| | | return ( |
| | | value["formal_filename"] is None |
| | |
| | | control = _open_inherited_handle(control_handle, "wb") |
| | | worker_input = _open_inherited_handle(input_handle, "rb") |
| | | cookie_closed = False |
| | | media_complete = False |
| | | prepared_run: Path | None = None |
| | | stage_root: Path | None = None |
| | | try: |
| | |
| | | cleanup_run_directory, |
| | | fixed_stage_root, |
| | | prepare_run_directory, |
| | | recover_published_task, |
| | | run_authenticated_task, |
| | | sanitized_environment, |
| | | ) |
| | |
| | | cleanup_run_directory, |
| | | fixed_stage_root, |
| | | prepare_run_directory, |
| | | recover_published_task, |
| | | run_authenticated_task, |
| | | sanitized_environment, |
| | | ) |
| | |
| | | os.environ.clear() |
| | | os.environ.update(safe_environment) |
| | | config = HostConfig.load(config_path) |
| | | stage_root = fixed_stage_root() |
| | | prepare_payload = read_frame(worker_input, MAX_INPUT_FRAME) |
| | | if prepare_payload is None: |
| | | raise CancelRequested() |
| | | worker_prepare = validate_worker_prepare(strict_json_loads(prepare_payload)) |
| | | job_spec = worker_prepare["job"] |
| | | lease_id = worker_prepare["lease_id"] |
| | | recovery_mode = worker_prepare["recovery_mode"] |
| | | |
| | | def report( |
| | | phase: str, progress: int, media_identity: dict[str, Any] | None = None, |
| | | ) -> None: |
| | | nonlocal media_complete |
| | | if phase == "MEDIA_COMPLETE": |
| | | media = validate_media_complete_identity(media_identity, job_spec) |
| | | write_frame( |
| | | control, |
| | | { |
| | | "schema": 1, "type": "progress", "phase": phase, |
| | | "progress": progress, "job": job_spec, |
| | | "lease_id": lease_id, "media": media, |
| | | }, |
| | | ) |
| | | acknowledgement_payload = read_frame(worker_input, MAX_INPUT_FRAME) |
| | | if acknowledgement_payload is None: |
| | | raise ProtocolError("E_MEDIA_COMPLETE") |
| | | validate_media_complete_ack( |
| | | strict_json_loads(acknowledgement_payload), |
| | | job_spec, lease_id, media, |
| | | ) |
| | | media_complete = True |
| | | return |
| | | if media_identity is not None: |
| | | raise ProtocolError("E_CONTROL") |
| | | write_frame(control, { |
| | | "schema": 1, "type": "progress", "phase": phase, |
| | | "progress": progress, |
| | | }) |
| | | |
| | | if recovery_mode == _RECOVERY_MODE_EXACT_PAIR: |
| | | formal, mapping = recover_published_task( |
| | | config, |
| | | job_spec, |
| | | cancel_check=lambda: is_event_set(cancel_handle), |
| | | report=report, |
| | | commit_begin=lambda: set_event(commit_handle), |
| | | ) |
| | | write_frame( |
| | | control, |
| | | { |
| | | "schema": 1, |
| | | "type": "terminal", |
| | | "phase": "COMPLETE", |
| | | "error_code": None, |
| | | "formal_filename": formal, |
| | | "mapping_filename": mapping, |
| | | "cookie_stream_closed": True, |
| | | }, |
| | | ) |
| | | return 0 |
| | | stage_root = fixed_stage_root(job_spec["bvid"]) |
| | | prepared_run = prepare_run_directory(stage_root) |
| | | write_frame(control, {"schema": 1, "type": "ready", "code": "READY_PLUGIN_DISABLED"}) |
| | | payload = read_frame(worker_input, MAX_INPUT_FRAME) |
| | |
| | | raise CancelRequested() |
| | | start = strict_json_loads(payload) |
| | | validate_message(start) |
| | | |
| | | def report(phase: str, progress: int) -> None: |
| | | write_frame( |
| | | control, |
| | | {"schema": 1, "type": "progress", "phase": phase, "progress": progress}, |
| | | ) |
| | | if start["job"] != job_spec: |
| | | raise ProtocolError("E_JOB") |
| | | |
| | | def closure_report(closed: bool) -> None: |
| | | nonlocal cookie_closed |
| | |
| | | prepared_run_directory=prepared_run, |
| | | commit_begin=lambda: set_event(commit_handle), |
| | | closure_report=closure_report, |
| | | recovery_required=recovery_mode == _RECOVERY_MODE_EXACT_PAIR, |
| | | ) |
| | | prepared_run = None |
| | | write_frame( |
| | |
| | | return 0 |
| | | except BaseException as exc: |
| | | error_code = getattr(exc, "code", None) |
| | | phase = "CANCELED" if type(exc).__name__ == "CancelRequested" else "FAILED" |
| | | diagnostic = getattr(exc, "diagnostic", None) |
| | | phase = ( |
| | | "CANCELED" if type(exc).__name__ == "CancelRequested" |
| | | else "POSTPROCESS_FAILED" if media_complete |
| | | else "FAILED" |
| | | ) |
| | | if not isinstance(error_code, str) or not error_code.startswith("E_"): |
| | | error_code = None if phase == "CANCELED" else "E_WORKER" |
| | | try: |
| | | write_frame( |
| | | control, |
| | | { |
| | | "schema": 1, |
| | | "type": "terminal", |
| | | "phase": phase, |
| | | "error_code": error_code, |
| | | "formal_filename": None, |
| | | "mapping_filename": None, |
| | | "cookie_stream_closed": cookie_closed, |
| | | }, |
| | | ) |
| | | terminal = { |
| | | "schema": 1, |
| | | "type": "terminal", |
| | | "phase": phase, |
| | | "error_code": error_code, |
| | | "formal_filename": None, |
| | | "mapping_filename": None, |
| | | "cookie_stream_closed": cookie_closed, |
| | | } |
| | | if phase == "POSTPROCESS_FAILED": |
| | | validate_postprocess_terminal(error_code, diagnostic) |
| | | elif diagnostic is not None: |
| | | validate_runtime_diagnostic(diagnostic) |
| | | if diagnostic is not None: |
| | | terminal["diagnostic"] = diagnostic |
| | | write_frame(control, terminal) |
| | | except BaseException: |
| | | pass |
| | | return 32 |
| | |
| | | incoming.put(None) |
| | | |
| | | |
| | | def _task_response(task: WorkerTask | None, preflight_error: str | None = None) -> dict[str, Any]: |
| | | def _task_response( |
| | | task: WorkerTask | None, |
| | | preflight_error: str | None = None, |
| | | maintenance: dict[str, Any] | None = None, |
| | | ) -> dict[str, Any]: |
| | | if task is None: |
| | | if preflight_error: |
| | | return safe_response("status", "FAILED", error_code=preflight_error) |
| | | return safe_response("status", "READY") |
| | | return safe_response("status", "FAILED", error_code=preflight_error, maintenance=maintenance) |
| | | return safe_response("status", "READY", maintenance=maintenance) |
| | | task.poll() |
| | | return safe_response( |
| | | "status", |
| | |
| | | error_code=task.error_code, |
| | | formal_filename=task.formal_filename, |
| | | mapping_filename=task.mapping_filename, |
| | | job=task.job_spec, |
| | | lease_id=task.lease_id, |
| | | maintenance=maintenance, |
| | | ) |
| | | |
| | | |
| | | def _presecret_poll_recovery( |
| | | queue_store: QueueStore, |
| | | runtime_config: dict[str, Any], |
| | | claimed_job: dict[str, Any], |
| | | claimed_lease: str, |
| | | maintenance: dict[str, Any] | None = None, |
| | | ) -> tuple[WorkerTask | None, dict[str, Any] | None]: |
| | | """Close an exact published pair before browser, Cookie, or network work.""" |
| | | recovery_binding = queue_store.postprocess_recovery_claim( |
| | | claimed_job, claimed_lease, |
| | | ) |
| | | collision_error, recovery_mode = _classify_job_destination(runtime_config, claimed_job) |
| | | if collision_error is not None: |
| | | if recovery_binding is None: |
| | | queue_store.mark_terminal( |
| | | claimed_job, claimed_lease, int(time.time() * 1000), |
| | | complete=False, error_code=collision_error, |
| | | ) |
| | | phase = "FAILED" |
| | | else: |
| | | queue_store.mark_postprocess_failed( |
| | | claimed_job, claimed_lease, int(time.time() * 1000), |
| | | error_code=collision_error, |
| | | ) |
| | | phase = "POSTPROCESS_FAILED" |
| | | return None, safe_response( |
| | | "poll", phase, error_code=collision_error, job=claimed_job, |
| | | lease_id=claimed_lease, maintenance=maintenance, |
| | | ) |
| | | if recovery_mode != _RECOVERY_MODE_EXACT_PAIR: |
| | | if recovery_binding is not None: |
| | | queue_store.mark_postprocess_failed( |
| | | claimed_job, claimed_lease, int(time.time() * 1000), |
| | | error_code="E_EXISTS", |
| | | ) |
| | | return None, safe_response( |
| | | "poll", "POSTPROCESS_FAILED", error_code="E_EXISTS", |
| | | job=claimed_job, lease_id=claimed_lease, maintenance=maintenance, |
| | | ) |
| | | return None, None |
| | | |
| | | task = WorkerTask(queue_store=queue_store) |
| | | task.postprocess_recovery_binding = recovery_binding |
| | | try: |
| | | if recovery_binding is None: |
| | | queue_store.mark_started( |
| | | claimed_job, claimed_lease, int(time.time() * 1000), |
| | | ) |
| | | task.prepare( |
| | | claimed_job, claimed_lease, None, os.urandom(16).hex(), |
| | | _RECOVERY_MODE_EXACT_PAIR, |
| | | ) |
| | | if not task.terminal: |
| | | task.phase = "FAILED" |
| | | task.error_code = "E_PLUGIN_BOUNDARY" |
| | | task.terminal = True |
| | | task.terminate() |
| | | response = safe_response( |
| | | "poll", task.phase, progress=task.progress, |
| | | error_code=task.error_code, formal_filename=task.formal_filename, |
| | | mapping_filename=task.mapping_filename, job=claimed_job, |
| | | lease_id=claimed_lease, maintenance=maintenance, |
| | | ) |
| | | _record_task_terminal(queue_store, task) |
| | | return task, response |
| | | except BaseException as exc: |
| | | error_code = getattr(exc, "code", None) |
| | | if not isinstance(error_code, str) or not re.fullmatch(r"E_[A-Z0-9_]{1,48}", error_code): |
| | | error_code = "E_PLUGIN_BOUNDARY" |
| | | task.job_spec = dict(claimed_job) |
| | | task.lease_id = claimed_lease |
| | | task.phase = "FAILED" |
| | | task.error_code = error_code |
| | | task.terminal = True |
| | | _record_task_terminal(queue_store, task) |
| | | return task, safe_response( |
| | | "poll", "FAILED", error_code=error_code, job=claimed_job, |
| | | lease_id=claimed_lease, maintenance=maintenance, |
| | | ) |
| | | |
| | | |
| | | def _start_prepared_task(task: WorkerTask | None, message: dict[str, Any]) -> str | None: |
| | |
| | | return None |
| | | |
| | | |
| | | def _reject_claimed_job( |
| | | queue_store: QueueStore | None, |
| | | claimed_job: dict[str, Any] | None, |
| | | claimed_lease: str | None, |
| | | message: dict[str, Any], |
| | | maintenance: dict[str, Any] | None = None, |
| | | *, |
| | | now_ms: int | None = None, |
| | | ) -> dict[str, Any]: |
| | | if ( |
| | | queue_store is None or claimed_job is None or claimed_lease is None |
| | | or message["job_id"] != claimed_job["job_id"] or message["lease_id"] != claimed_lease |
| | | ): |
| | | return safe_response("reject", "FAILED", error_code="E_LEASE", maintenance=maintenance) |
| | | queue_store.reject_claim( |
| | | claimed_job, |
| | | claimed_lease, |
| | | int(time.time() * 1000) if now_ms is None else now_ms, |
| | | message["error_code"], |
| | | message.get("diagnostic"), |
| | | ) |
| | | return safe_response( |
| | | "reject", "FAILED", error_code=message["error_code"], |
| | | job=claimed_job, lease_id=claimed_lease, maintenance=maintenance, |
| | | ) |
| | | |
| | | |
| | | def _foreground_claimed_job( |
| | | queue_store: QueueStore | None, |
| | | claimed_job: dict[str, Any] | None, |
| | | claimed_lease: str | None, |
| | | message: dict[str, Any], |
| | | maintenance: dict[str, Any] | None, |
| | | *, |
| | | attempted: bool, |
| | | parent_window: int, |
| | | now_ms: int | None = None, |
| | | activator: Any = _foreground_chrome_window, |
| | | ) -> tuple[dict[str, Any], bool]: |
| | | """Execute at most one claim-bound, pre-secret foreground request.""" |
| | | |
| | | if ( |
| | | queue_store is None or claimed_job is None or claimed_lease is None |
| | | or message["job_id"] != claimed_job["job_id"] |
| | | or message["lease_id"] != claimed_lease |
| | | ): |
| | | error_code = "E_LEASE" |
| | | consumed = attempted |
| | | elif attempted: |
| | | error_code = "E_FOREGROUND_REPLAY" |
| | | consumed = True |
| | | else: |
| | | consumed = True |
| | | try: |
| | | queue_store.assert_claim( |
| | | claimed_job, claimed_lease, |
| | | int(time.time() * 1000) if now_ms is None else now_ms, |
| | | ) |
| | | error_code = activator( |
| | | claimed_job["canonical_url"], message["window_bounds"], parent_window, |
| | | ) |
| | | except ProtocolError: |
| | | error_code = "E_LEASE" |
| | | except BaseException: |
| | | error_code = "E_FOREGROUND_PLATFORM_UNSUPPORTED" |
| | | if error_code is not None and error_code not in _FOREGROUND_ERROR_CODES: |
| | | error_code = "E_FOREGROUND_PLATFORM_UNSUPPORTED" |
| | | return safe_response( |
| | | "foreground", "READY" if error_code is None else "FAILED", |
| | | error_code=error_code, job=claimed_job, lease_id=claimed_lease, |
| | | maintenance=maintenance, |
| | | ), consumed |
| | | |
| | | |
| | | def _abort_prepared_task( |
| | | queue_store: QueueStore | None, |
| | | task: WorkerTask | None, |
| | | claimed_job: dict[str, Any] | None, |
| | | claimed_lease: str | None, |
| | | message: dict[str, Any], |
| | | maintenance: dict[str, Any] | None = None, |
| | | *, |
| | | now_ms: int | None = None, |
| | | ) -> dict[str, Any]: |
| | | """Close one prepared worker and persist the extension's safe pre-start terminal.""" |
| | | if ( |
| | | queue_store is None or task is None or claimed_job is None or claimed_lease is None |
| | | or message["job_id"] != claimed_job["job_id"] |
| | | or message["lease_id"] != claimed_lease |
| | | or message["prepare_id"] != task.prepare_id |
| | | or task.job_spec != claimed_job or task.lease_id != claimed_lease |
| | | or not task.prepared or task.secret_started or task.terminal |
| | | ): |
| | | return safe_response( |
| | | "abort_prepare", "FAILED", error_code="E_PREPARE", |
| | | prepare_id=message.get("prepare_id"), maintenance=maintenance, |
| | | ) |
| | | # Terminate while the task is still nonterminal so the owned worker receives |
| | | # cancellation/EOF and cannot survive the durable queue terminal. |
| | | task.terminate() |
| | | earlier_error = ( |
| | | task.error_code |
| | | if task.phase == "FAILED" and isinstance(task.error_code, str) |
| | | and re.fullmatch(r"E_[A-Z0-9_]{1,48}", task.error_code) |
| | | else None |
| | | ) |
| | | task.phase = "FAILED" |
| | | task.error_code = earlier_error or COOKIE_ACCESS_TERMINAL_CODES[message["error_reason"]] |
| | | task.terminal = True |
| | | if not _record_task_terminal( |
| | | queue_store, |
| | | task, |
| | | now_ms=int(time.time() * 1000) if now_ms is None else now_ms, |
| | | ): |
| | | return safe_response( |
| | | "abort_prepare", "FAILED", error_code="E_PREPARE", |
| | | prepare_id=message["prepare_id"], maintenance=maintenance, |
| | | ) |
| | | return safe_response( |
| | | "abort_prepare", "FAILED", error_code=earlier_error or message["error_code"], |
| | | prepare_id=message["prepare_id"], job=claimed_job, |
| | | lease_id=claimed_lease, maintenance=maintenance, |
| | | ) |
| | | |
| | | |
| | | def _record_task_terminal( |
| | | queue_store: QueueStore | None, |
| | | task: WorkerTask | None, |
| | | *, |
| | | fallback_error: str | None = None, |
| | | now_ms: int | None = None, |
| | | ) -> bool: |
| | | """Persist one task terminal, preserving an earlier deterministic worker error. |
| | | |
| | | A Host/port disconnect is only a fallback. Control messages already queued |
| | | before teardown win; teardown-induced CANCELED/E_WORKER_EXIT states do not. |
| | | """ |
| | | if ( |
| | | queue_store is None or task is None or task.queue_terminal_recorded |
| | | or task.job_spec is None or task.lease_id is None |
| | | ): |
| | | return False |
| | | task.poll() |
| | | terminal_before_shutdown = task.terminal |
| | | if not task.terminal and fallback_error is not None: |
| | | task.terminate() |
| | | task.poll() |
| | | if ( |
| | | not task.terminal |
| | | or task.phase == "CANCELED" |
| | | or task.error_code in {None, "E_WORKER_EXIT"} |
| | | ): |
| | | task.phase = "FAILED" |
| | | task.error_code = fallback_error |
| | | task.terminal = True |
| | | if not task.terminal: |
| | | return False |
| | | task_media_complete = bool(getattr(task, "media_complete", False)) |
| | | postprocess_recovery = getattr(task, "postprocess_recovery_binding", None) is not None |
| | | if (task_media_complete or postprocess_recovery) and task.phase != "COMPLETE": |
| | | task.phase = "POSTPROCESS_FAILED" |
| | | complete = task.phase == "COMPLETE" |
| | | if complete: |
| | | error_code = None |
| | | elif task.phase == "CANCELED" and terminal_before_shutdown: |
| | | error_code = "E_CANCEL" |
| | | elif isinstance(task.error_code, str) and re.fullmatch(r"E_[A-Z0-9_]{1,48}", task.error_code): |
| | | error_code = task.error_code |
| | | else: |
| | | error_code = fallback_error or "E_WORKER" |
| | | terminal_now = int(time.time() * 1000) if now_ms is None else now_ms |
| | | if task_media_complete: |
| | | if task.media_identity is None: |
| | | raise ProtocolError("E_QUEUE_STATE") |
| | | queue_store.mark_media_complete( |
| | | task.job_spec, task.lease_id, terminal_now, task.media_identity, |
| | | ) |
| | | if task.phase == "POSTPROCESS_FAILED": |
| | | queue_store.mark_postprocess_failed( |
| | | task.job_spec, task.lease_id, terminal_now, |
| | | error_code=error_code, |
| | | diagnostic=getattr(task, "diagnostic", None), |
| | | ) |
| | | else: |
| | | queue_store.mark_terminal( |
| | | task.job_spec, |
| | | task.lease_id, |
| | | terminal_now, |
| | | complete=complete, |
| | | error_code=error_code, |
| | | diagnostic=getattr(task, "diagnostic", None), |
| | | ) |
| | | task.queue_terminal_recorded = True |
| | | task.close() |
| | | return True |
| | | |
| | | |
| | | def broker_main(arguments: list[str]) -> int: |
| | | validate_origin_argv(arguments, EXPECTED_ORIGIN) |
| | | protocol_output = _duplicate_protocol_output() |
| | |
| | | ) |
| | | reader.start() |
| | | task: WorkerTask | None = None |
| | | claimed_job: dict[str, Any] | None = None |
| | | claimed_lease: str | None = None |
| | | foreground_attempted = False |
| | | hello_complete = False |
| | | current_extension_build: str | None = None |
| | | preflight_error = preflight_configuration(_config_path()) |
| | | runtime_config: dict[str, Any] | None = None |
| | | queue_store: QueueStore | None = None |
| | | reload_store: ReloadStore | None = None |
| | | current_maintenance = maintenance_state() |
| | | if preflight_error is None: |
| | | runtime_config = load_runtime_configuration(_config_path()) |
| | | queue_store = _queue_store(runtime_config) |
| | | reload_store = ReloadStore(Path(runtime_config["reload_state_path"]), runtime_config["reload_generation"]) |
| | | try: |
| | | while True: |
| | | if task is not None: |
| | | task.poll() |
| | | _record_task_terminal(queue_store, task) |
| | | try: |
| | | message = incoming.get(timeout=0.1) |
| | | except queue.Empty: |
| | | continue |
| | | if message is None: |
| | | if task is not None and not task.terminal: |
| | | task.terminate() |
| | | _record_task_terminal(queue_store, task, fallback_error="E_HOST_DISCONNECT") |
| | | return 0 |
| | | message_type = message["type"] |
| | | if not hello_complete: |
| | | if message_type != "hello": |
| | | raise ProtocolError() |
| | | hello_complete = True |
| | | current_extension_build = message["extension_build"] |
| | | if reload_store is not None: |
| | | current_maintenance = reload_store.status(current_extension_build, int(time.time() * 1000)) |
| | | write_frame( |
| | | protocol_output, |
| | | safe_response( |
| | | "hello", |
| | | "FAILED" if preflight_error else "READY", |
| | | error_code=preflight_error, |
| | | maintenance=current_maintenance, |
| | | ), |
| | | ) |
| | | continue |
| | | if message_type == "hello": |
| | | raise ProtocolError() |
| | | if message_type == "status": |
| | | current_preflight = preflight_configuration(_config_path()) |
| | | write_frame(protocol_output, _task_response(task, current_preflight)) |
| | | elif message_type == "start": |
| | | start_error = _start_prepared_task(task, message) |
| | | if start_error is None: |
| | | write_frame(protocol_output, _task_response(task)) |
| | | else: |
| | | write_frame(protocol_output, safe_response("start", "FAILED", error_code=start_error)) |
| | | elif message_type == "prepare": |
| | | preflight_error = preflight_configuration(_config_path()) |
| | | if preflight_error: |
| | | write_frame(protocol_output, safe_response("prepare", "FAILED", error_code=preflight_error, prepare_id=message["prepare_id"])) |
| | | if message_type == "reload_begin": |
| | | if reload_store is None or current_extension_build is None: |
| | | write_frame(protocol_output, safe_response("reload_begin", "FAILED", error_code="E_RELOAD", maintenance=current_maintenance)) |
| | | continue |
| | | reload_store.begin(current_extension_build, message["reload_token"], int(time.time() * 1000)) |
| | | write_frame(protocol_output, safe_response("reload_begin", "RELOAD_REQUIRED", maintenance=current_maintenance)) |
| | | elif current_maintenance.get("reload_required"): |
| | | write_frame(protocol_output, safe_response(message_type, "FAILED", error_code="E_RELOAD_REQUIRED", maintenance=current_maintenance)) |
| | | elif message_type == "poll": |
| | | if queue_store is None: |
| | | write_frame(protocol_output, safe_response("poll", "FAILED", error_code=preflight_error or "E_CONFIG", maintenance=current_maintenance)) |
| | | continue |
| | | if task is not None and not task.terminal: |
| | | write_frame(protocol_output, safe_response("prepare", "FAILED", error_code="E_BUSY", prepare_id=message["prepare_id"])) |
| | | write_frame(protocol_output, _task_response(task, maintenance=current_maintenance)) |
| | | continue |
| | | claimed = queue_store.claim_next(int(time.time() * 1000)) |
| | | if claimed is None: |
| | | claimed_job = None |
| | | claimed_lease = None |
| | | foreground_attempted = False |
| | | write_frame(protocol_output, safe_response("poll", "IDLE", maintenance=current_maintenance)) |
| | | else: |
| | | claimed_job, claimed_lease = claimed |
| | | foreground_attempted = False |
| | | preflight_error = preflight_configuration(_config_path()) |
| | | if preflight_error is not None or runtime_config is None: |
| | | recovery_binding = queue_store.postprocess_recovery_claim( |
| | | claimed_job, claimed_lease, |
| | | ) |
| | | terminal_error = preflight_error or "E_CONFIG" |
| | | if recovery_binding is None: |
| | | queue_store.mark_terminal( |
| | | claimed_job, claimed_lease, int(time.time() * 1000), |
| | | complete=False, error_code=terminal_error, |
| | | ) |
| | | terminal_phase = "FAILED" |
| | | else: |
| | | queue_store.mark_postprocess_failed( |
| | | claimed_job, claimed_lease, int(time.time() * 1000), |
| | | error_code=terminal_error, |
| | | ) |
| | | terminal_phase = "POSTPROCESS_FAILED" |
| | | write_frame(protocol_output, safe_response( |
| | | "poll", terminal_phase, error_code=terminal_error, |
| | | job=claimed_job, lease_id=claimed_lease, |
| | | maintenance=current_maintenance, |
| | | )) |
| | | continue |
| | | task, recovery_response = _presecret_poll_recovery( |
| | | queue_store, runtime_config, claimed_job, claimed_lease, |
| | | current_maintenance, |
| | | ) |
| | | if recovery_response is not None: |
| | | write_frame(protocol_output, recovery_response) |
| | | else: |
| | | write_frame(protocol_output, safe_response( |
| | | "poll", "READY", job=claimed_job, lease_id=claimed_lease, |
| | | maintenance=current_maintenance, |
| | | )) |
| | | elif message_type == "foreground": |
| | | if task is not None and not task.terminal: |
| | | foreground_response = safe_response( |
| | | "foreground", "FAILED", error_code="E_LEASE", |
| | | job=claimed_job, lease_id=claimed_lease, |
| | | maintenance=current_maintenance, |
| | | ) |
| | | else: |
| | | foreground_response, foreground_attempted = _foreground_claimed_job( |
| | | queue_store, claimed_job, claimed_lease, message, current_maintenance, |
| | | attempted=foreground_attempted, |
| | | parent_window=_native_parent_window(arguments), |
| | | ) |
| | | write_frame(protocol_output, foreground_response) |
| | | elif message_type == "reject": |
| | | rejected = _reject_claimed_job( |
| | | queue_store, claimed_job, claimed_lease, message, current_maintenance |
| | | ) |
| | | write_frame(protocol_output, rejected) |
| | | if rejected["error_code"] != "E_LEASE": |
| | | claimed_job = None |
| | | claimed_lease = None |
| | | elif message_type == "abort_prepare": |
| | | aborted = _abort_prepared_task( |
| | | queue_store, task, claimed_job, claimed_lease, message, current_maintenance |
| | | ) |
| | | write_frame(protocol_output, aborted) |
| | | if aborted["error_code"] != "E_PREPARE": |
| | | claimed_job = None |
| | | claimed_lease = None |
| | | elif message_type == "status": |
| | | if task is None or message["job_id"] != (task.job_spec or {}).get("job_id") or message["lease_id"] != task.lease_id: |
| | | write_frame(protocol_output, safe_response("status", "FAILED", error_code="E_LEASE", maintenance=current_maintenance)) |
| | | continue |
| | | current_preflight = preflight_configuration(_config_path()) |
| | | write_frame(protocol_output, _task_response(task, current_preflight, current_maintenance)) |
| | | elif message_type == "start": |
| | | if queue_store is None or message["job"] != claimed_job or message["lease_id"] != claimed_lease: |
| | | write_frame(protocol_output, safe_response("start", "FAILED", error_code="E_LEASE", job=claimed_job, lease_id=claimed_lease, maintenance=current_maintenance)) |
| | | continue |
| | | queue_store.mark_started(message["job"], message["lease_id"], int(time.time() * 1000)) |
| | | start_error = _start_prepared_task(task, message) |
| | | if start_error is None: |
| | | write_frame(protocol_output, _task_response(task, maintenance=current_maintenance)) |
| | | else: |
| | | queue_store.mark_terminal(message["job"], message["lease_id"], int(time.time() * 1000), complete=False, error_code=start_error) |
| | | write_frame(protocol_output, safe_response("start", "FAILED", error_code=start_error, job=claimed_job, lease_id=claimed_lease, maintenance=current_maintenance)) |
| | | elif message_type == "prepare": |
| | | preflight_error = preflight_configuration(_config_path()) |
| | | if message["job"] != claimed_job or message["lease_id"] != claimed_lease or queue_store is None or runtime_config is None: |
| | | write_frame(protocol_output, safe_response("prepare", "FAILED", error_code="E_LEASE", prepare_id=message["prepare_id"], job=claimed_job, lease_id=claimed_lease, maintenance=current_maintenance)) |
| | | continue |
| | | queue_store.assert_claim(message["job"], message["lease_id"], int(time.time() * 1000)) |
| | | recovery_mode = _RECOVERY_MODE_NONE |
| | | if preflight_error is None: |
| | | preflight_error, recovery_mode = _classify_job_destination(runtime_config, message["job"]) |
| | | if preflight_error: |
| | | queue_store.mark_terminal(message["job"], message["lease_id"], int(time.time() * 1000), complete=False, error_code=preflight_error) |
| | | write_frame(protocol_output, safe_response("prepare", "FAILED", error_code=preflight_error, prepare_id=message["prepare_id"], job=claimed_job, lease_id=claimed_lease, maintenance=current_maintenance)) |
| | | continue |
| | | if task is not None and not task.terminal: |
| | | write_frame(protocol_output, safe_response("prepare", "FAILED", error_code="E_BUSY", prepare_id=message["prepare_id"], job=claimed_job, lease_id=claimed_lease, maintenance=current_maintenance)) |
| | | continue |
| | | if task is not None: |
| | | task.close() |
| | | task = WorkerTask() |
| | | task = WorkerTask(queue_store=queue_store) |
| | | try: |
| | | task.prepare(message["page_proof"], message["prepare_id"]) |
| | | write_frame(protocol_output, safe_response("prepare", "READY", prepare_id=message["prepare_id"])) |
| | | except BaseException: |
| | | write_frame(protocol_output, safe_response("prepare", "FAILED", error_code="E_PLUGIN_BOUNDARY", prepare_id=message["prepare_id"])) |
| | | if recovery_mode == _RECOVERY_MODE_EXACT_PAIR: |
| | | queue_store.mark_started( |
| | | message["job"], message["lease_id"], int(time.time() * 1000) |
| | | ) |
| | | task.prepare(message["job"], message["lease_id"], message["page_proof"], message["prepare_id"], recovery_mode) |
| | | if task.terminal: |
| | | response = safe_response( |
| | | "prepare", task.phase, progress=task.progress, |
| | | error_code=task.error_code, |
| | | formal_filename=task.formal_filename, |
| | | mapping_filename=task.mapping_filename, |
| | | prepare_id=message["prepare_id"], job=claimed_job, |
| | | lease_id=claimed_lease, maintenance=current_maintenance, |
| | | ) |
| | | _record_task_terminal(queue_store, task) |
| | | write_frame(protocol_output, response) |
| | | else: |
| | | write_frame(protocol_output, safe_response("prepare", "READY", prepare_id=message["prepare_id"], job=claimed_job, lease_id=claimed_lease, maintenance=current_maintenance)) |
| | | except BaseException as exc: |
| | | error_code = getattr(exc, "code", None) |
| | | if not isinstance(error_code, str) or not re.fullmatch(r"E_[A-Z0-9_]{1,48}", error_code): |
| | | error_code = "E_PLUGIN_BOUNDARY" |
| | | task.phase = "FAILED" |
| | | task.error_code = error_code |
| | | task.terminal = True |
| | | _record_task_terminal(queue_store, task) |
| | | write_frame(protocol_output, safe_response("prepare", "FAILED", error_code=error_code, prepare_id=message["prepare_id"], job=claimed_job, lease_id=claimed_lease, maintenance=current_maintenance)) |
| | | elif message_type == "cancel": |
| | | canceled = task is not None and task.cancel(message["task_nonce"]) |
| | | canceled = task is not None and message["job_id"] == (task.job_spec or {}).get("job_id") and message["lease_id"] == task.lease_id and task.cancel(message["task_nonce"]) |
| | | write_frame( |
| | | protocol_output, |
| | | safe_response("cancel", "CANCELED" if canceled else "FAILED", error_code=None if canceled else "E_CANCEL"), |
| | | safe_response("cancel", "CANCELED" if canceled else "FAILED", error_code=None if canceled else "E_CANCEL", job=None if task is None else task.job_spec, lease_id=None if task is None else task.lease_id, maintenance=current_maintenance), |
| | | ) |
| | | except BaseException: |
| | | if task is not None and not task.terminal: |
| | | task.terminate() |
| | | try: |
| | | _record_task_terminal(queue_store, task, fallback_error="E_HOST") |
| | | except BaseException: |
| | | pass |
| | | try: |
| | | write_frame(protocol_output, safe_response("error", "FAILED", error_code="E_PROTOCOL")) |
| | | except BaseException: |