From a141d5debce61b5e0e9ba0c28a43a62216d11d73 Mon Sep 17 00:00:00 2001
From: MB-X Bilibili Pipeline <mbx-bili-pipeline@localhost>
Date: Sat, 05 Sep 2026 22:31:14 +0800
Subject: [PATCH] chore(project-info): archive Bilibili dynamic acb922edb4e9

---
 dev/project-dev/bili_authenticated_extension/native_host.py | 1487 +++++++++++++++++++++++++++++++++++++++++++++++++++++++----
 1 files changed, 1,377 insertions(+), 110 deletions(-)

diff --git a/dev/project-dev/bili_authenticated_extension/native_host.py b/dev/project-dev/bili_authenticated_extension/native_host.py
index c5c2bb8..1dd791c 100644
--- a/dev/project-dev/bili_authenticated_extension/native_host.py
+++ b/dev/project-dev/bili_authenticated_extension/native_host.py
@@ -14,22 +14,30 @@
 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,
@@ -41,22 +49,32 @@
     )
     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,
@@ -68,13 +86,18 @@
     )
     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:
@@ -90,12 +113,517 @@
     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",
@@ -104,45 +632,94 @@
         "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(
@@ -171,16 +748,25 @@
     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
@@ -189,6 +775,7 @@
         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
@@ -196,10 +783,41 @@
         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)
@@ -243,6 +861,12 @@
         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:
@@ -254,19 +878,25 @@
             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 (
@@ -275,11 +905,12 @@
             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()
@@ -294,15 +925,60 @@
                 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"]
@@ -311,15 +987,18 @@
                 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()
@@ -327,6 +1006,9 @@
                 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"
@@ -344,7 +1026,10 @@
         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
@@ -375,43 +1060,88 @@
         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":
@@ -420,9 +1150,26 @@
             "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
@@ -436,7 +1183,10 @@
         "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
@@ -445,12 +1195,31 @@
         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
@@ -491,6 +1260,7 @@
     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:
@@ -502,6 +1272,7 @@
                 cleanup_run_directory,
                 fixed_stage_root,
                 prepare_run_directory,
+                recover_published_task,
                 run_authenticated_task,
                 sanitized_environment,
             )
@@ -513,6 +1284,7 @@
                 cleanup_run_directory,
                 fixed_stage_root,
                 prepare_run_directory,
+                recover_published_task,
                 run_authenticated_task,
                 sanitized_environment,
             )
@@ -526,7 +1298,66 @@
         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)
@@ -534,12 +1365,8 @@
             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
@@ -554,6 +1381,7 @@
             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(
@@ -571,22 +1399,31 @@
         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
@@ -657,11 +1494,15 @@
         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",
@@ -670,7 +1511,91 @@
         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:
@@ -687,6 +1612,203 @@
     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()
@@ -699,70 +1821,215 @@
     )
     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:

--
Gitblit v1.9.3