Cai
2026-08-24 156ea25b402479f0abc54c558bbf87f9eaaa0422
dev/project-dev/bili_dynamic_collector.py
@@ -22,6 +22,12 @@
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
# Keep lazy refresh imports on the same module identity when this file is run as
# a script, so CollectorError and dataclass contracts are not duplicated.
if __name__ == "__main__":
    sys.modules.setdefault("bili_dynamic_collector", sys.modules[__name__])
SCHEMA_VERSION = 1
MANIFEST_SCHEMA_VERSION = 1
SECRET_KEY_PATTERN = re.compile(
@@ -57,6 +63,7 @@
    "PROCESSING_HANDOFF_CONFIRMED",
    "PROCESSING",
    "COMPLETE",
    "CONTENT_SAVED",
}
RETRYABLE_STATUSES = {
    "QUEUE_FAILED",
@@ -121,8 +128,26 @@
@dataclass(frozen=True)
class RefreshConfig:
    archive_dir: Path
    formal_manifest: Path
    intake_dir: Path
    overall_deadline_seconds: int
    refresh_action_timeout_seconds: int
    observation_timeout_seconds: int
    page_internal_settle_timeout_seconds: int
    max_refresh_count: int
    run_history_slots: int
    max_items: int
    max_images_per_item: int
    max_image_bytes: int
    max_text_bytes: int
@dataclass(frozen=True)
class CollectorConfig:
    creator_name: str
    creator_uid: str | None
    creator_dynamic_url: str
    timezone_name: str
    window_hours: int
@@ -134,6 +159,7 @@
    allowed_source_hosts: frozenset[str]
    allowed_video_extensions: frozenset[str]
    native_handoff: NativeHandoffRoute
    refresh: RefreshConfig | None = None
    @property
    def manifest_path(self) -> Path:
@@ -157,40 +183,165 @@
class StateLock:
    """Small fail-fast lock for overlapping local/scheduled invocations."""
    """Kernel-backed fail-fast lock; a dead process cannot strand ownership."""
    def __init__(self, path: Path) -> None:
        self.path = path
        self.fd: int | None = None
        self.owner_path = path.parent / f"{path.name}.owner.json"
    def _try_kernel_lock(self) -> None:
        assert self.fd is not None
        if os.name == "nt":
            import msvcrt
            os.lseek(self.fd, 0, os.SEEK_SET)
            try:
                msvcrt.locking(self.fd, msvcrt.LK_NBLCK, 1)
            except OSError as exc:
                raise CollectorError("E_STATE_LOCKED", "Another collector invocation is active.", safety=True) from exc
        else:  # pragma: no cover - exercised by non-Windows CI only
            import fcntl
            try:
                fcntl.flock(self.fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
            except OSError as exc:
                raise CollectorError("E_STATE_LOCKED", "Another collector invocation is active.", safety=True) from exc
    def _unlock(self) -> None:
        assert self.fd is not None
        if os.name == "nt":
            import msvcrt
            os.lseek(self.fd, 0, os.SEEK_SET)
            msvcrt.locking(self.fd, msvcrt.LK_UNLCK, 1)
        else:  # pragma: no cover
            import fcntl
            fcntl.flock(self.fd, fcntl.LOCK_UN)
    def __enter__(self) -> "StateLock":
        ensure_directory(self.path.parent, create=True)
        lexical_lstat_chain(self.path, allow_missing_leaf=True)
        if self.path.exists():
            info = os.lstat(self.path)
            if not self.path.is_file() or is_reparse(info):
                raise CollectorError("E_STATE_LOCK", "State lock path is not a regular file.", safety=True)
            try:
                lock_bytes = self.path.read_bytes()
            except PermissionError as exc:
                raise CollectorError("E_STATE_LOCKED", "Another collector invocation is active.", safety=True) from exc
            if lock_bytes not in (b"", b"\0"):
                raise CollectorError("E_STATE_LOCK", "Legacy or damaged state lock requires review.", safety=True)
        try:
            self.fd = os.open(self.path, os.O_CREAT | os.O_EXCL | os.O_WRONLY)
        except FileExistsError as exc:
            raise CollectorError(
                "E_STATE_LOCKED",
                "Another collector invocation is active, or a stale lock needs manual review.",
                details={"lock_path": str(self.path)},
                safety=True,
            ) from exc
        payload = json.dumps(
            {"pid": os.getpid(), "created_at": utc_now().isoformat()},
            ensure_ascii=False,
            sort_keys=True,
        ).encode("utf-8")
        os.write(self.fd, payload)
        os.fsync(self.fd)
            self.fd = os.open(self.path, os.O_CREAT | os.O_RDWR, 0o600)
        except PermissionError as exc:
            raise CollectorError("E_STATE_LOCKED", "Another collector invocation is active.", safety=True) from exc
        try:
            if os.fstat(self.fd).st_size == 0:
                os.write(self.fd, b"\0")
                os.fsync(self.fd)
            self._try_kernel_lock()
            if self.owner_path.exists():
                lexical_lstat_chain(self.owner_path, allow_missing_leaf=False)
                try:
                    prior = json.loads(self.owner_path.read_text(encoding="utf-8"))
                except (UnicodeDecodeError, json.JSONDecodeError) as exc:
                    raise CollectorError("E_STATE_LOCK", "State lock owner metadata is damaged.", safety=True) from exc
                if not isinstance(prior, dict) or set(prior) != {"pid", "process_created_at", "run_id", "acquired_at"}:
                    raise CollectorError("E_STATE_LOCK", "State lock owner metadata has an invalid schema.", safety=True)
                prior_pid = prior.get("pid")
                prior_created = prior.get("process_created_at")
                if not isinstance(prior_pid, int) or not isinstance(prior_created, str):
                    raise CollectorError("E_STATE_LOCK", "State lock owner identity is invalid.", safety=True)
                try:
                    actual_created = process_created_at(prior_pid)
                except ProcessLookupError:
                    actual_created = None
                except OSError as exc:
                    raise CollectorError("E_STATE_LOCK", "State lock owner identity is unprovable.", safety=True) from exc
                if actual_created is not None:
                    code = "E_STATE_LOCKED" if actual_created == prior_created else "E_STATE_LOCK_PID_REUSE"
                    raise CollectorError(code, "State lock owner metadata refers to a live process.", safety=True)
                quarantine = self.path.parent / f".{self.owner_path.name}.{hashlib.sha256(self.owner_path.read_bytes()).hexdigest()}.stale"
                if quarantine.exists():
                    raise CollectorError("E_STATE_LOCK", "State lock owner quarantine already exists.", safety=True)
                os.replace(self.owner_path, quarantine)
            else:
                quarantine = None
            payload = canonical_json_bytes(
                {
                    "pid": os.getpid(),
                    "process_created_at": process_created_at(os.getpid()),
                    "run_id": None,
                    "acquired_at": canonical_datetime(utc_now()),
                }
            )
            atomic_replace_bytes(self.owner_path, payload)
            if quarantine is not None:
                quarantine.unlink()
        except BaseException:
            try:
                self._unlock()
            except BaseException:
                pass
            os.close(self.fd)
            self.fd = None
            raise
        return self
    def __exit__(self, exc_type: object, exc: object, tb: object) -> None:
        if self.fd is not None:
            os.close(self.fd)
            try:
                if self.owner_path.exists():
                    self.owner_path.unlink()
                self._unlock()
            finally:
                os.close(self.fd)
            self.fd = None
def process_created_at(pid: int) -> str:
    """Return a stable process creation identity without exposing command lines."""
    if os.name == "nt":
        import ctypes
        from ctypes import wintypes
        PROCESS_QUERY_LIMITED_INFORMATION = 0x1000
        kernel32 = ctypes.WinDLL("kernel32", use_last_error=True)
        kernel32.OpenProcess.argtypes = (wintypes.DWORD, wintypes.BOOL, wintypes.DWORD)
        kernel32.OpenProcess.restype = wintypes.HANDLE
        kernel32.GetProcessTimes.argtypes = (
            wintypes.HANDLE,
            ctypes.POINTER(wintypes.FILETIME),
            ctypes.POINTER(wintypes.FILETIME),
            ctypes.POINTER(wintypes.FILETIME),
            ctypes.POINTER(wintypes.FILETIME),
        )
        kernel32.GetProcessTimes.restype = wintypes.BOOL
        kernel32.CloseHandle.argtypes = (wintypes.HANDLE,)
        kernel32.CloseHandle.restype = wintypes.BOOL
        handle = kernel32.OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, False, pid)
        if not handle:
            error = ctypes.get_last_error()
            if error in (87, 1168):
                raise ProcessLookupError(pid)
            raise OSError(error, "Could not query process identity")
        try:
            self.path.unlink()
        except FileNotFoundError:
            pass
            created = wintypes.FILETIME()
            exited = wintypes.FILETIME()
            kernel = wintypes.FILETIME()
            user = wintypes.FILETIME()
            if not kernel32.GetProcessTimes(handle, ctypes.byref(created), ctypes.byref(exited), ctypes.byref(kernel), ctypes.byref(user)):
                raise OSError(ctypes.get_last_error(), "Could not read process creation time")
            ticks = (created.dwHighDateTime << 32) | created.dwLowDateTime
            return str(ticks)
        finally:
            kernel32.CloseHandle(handle)
    stat = Path(f"/proc/{pid}/stat")  # pragma: no cover
    if not stat.exists():
        raise ProcessLookupError(pid)
    return stat.read_text(encoding="ascii").split()[21]
def utc_now() -> datetime:
@@ -414,6 +565,13 @@
    if any(not host or "/" in host or ":" in host for host in allowed_hosts):
        raise CollectorError("E_CONFIG", "allowed_source_hosts contains an invalid hostname.")
    dynamic_url = validate_url(creator.get("dynamic_url"), "creator.dynamic_url", allowed_hosts)
    creator_uid_raw = creator.get("uid")
    if creator_uid_raw is None:
        creator_uid = None
    elif isinstance(creator_uid_raw, (str, int)) and re.fullmatch(r"[1-9][0-9]{0,19}", str(creator_uid_raw)):
        creator_uid = str(creator_uid_raw)
    else:
        raise CollectorError("E_CONFIG", "creator.uid must be a positive decimal identifier.")
    window_hours = value.get("window_hours", 72)
    minimum_age = value.get("minimum_complete_age_seconds", 30)
    title_max = value.get("title_max_length", 48)
@@ -430,8 +588,54 @@
    if any(not re.fullmatch(r"\.[a-z0-9]{1,8}", item) for item in extensions):
        raise CollectorError("E_CONFIG", "allowed_video_extensions contains an invalid suffix.")
    base = path.parent
    refresh_raw = value.get("refresh")
    refresh: RefreshConfig | None = None
    if refresh_raw is not None:
        if not isinstance(refresh_raw, Mapping):
            raise CollectorError("E_CONFIG", "config.refresh must be an object.")
        if (
            creator_name.strip() != "青枫浦上Q"
            or creator_uid != "1420210197"
            or dynamic_url != "https://space.bilibili.com/1420210197/dynamic"
        ):
            raise CollectorError("E_CONFIG", "refresh supports only the registered creator identity.", safety=True)
        def bounded_int(field: str, default: int, lower: int, upper: int) -> int:
            raw = refresh_raw.get(field, default)
            if not isinstance(raw, int) or isinstance(raw, bool) or not lower <= raw <= upper:
                raise CollectorError("E_CONFIG", f"refresh.{field} must be {lower}..{upper}.")
            return raw
        legacy_timeout_fields = {
            "timeout_seconds", "page_ready_timeout_seconds", "dom_read_timeout_seconds"
        }
        if legacy_timeout_fields.intersection(refresh_raw):
            raise CollectorError("E_CONFIG", "Legacy refresh timeout fields are not valid for runtime-v2.")
        def frozen_int(field: str, expected: int) -> int:
            raw = refresh_raw.get(field)
            if not isinstance(raw, int) or isinstance(raw, bool) or raw != expected:
                raise CollectorError("E_CONFIG", f"refresh.{field} must equal {expected}.")
            return raw
        refresh = RefreshConfig(
            archive_dir=config_path(base, refresh_raw.get("archive_dir"), "refresh.archive_dir"),
            formal_manifest=config_path(base, refresh_raw.get("formal_manifest"), "refresh.formal_manifest"),
            intake_dir=config_path(base, refresh_raw.get("intake_dir"), "refresh.intake_dir"),
            overall_deadline_seconds=frozen_int("overall_deadline_seconds", 120),
            refresh_action_timeout_seconds=frozen_int("refresh_action_timeout_seconds", 35),
            observation_timeout_seconds=frozen_int("observation_timeout_seconds", 45),
            page_internal_settle_timeout_seconds=frozen_int("page_internal_settle_timeout_seconds", 15),
            max_refresh_count=bounded_int("max_refresh_count", 1, 1, 1),
            run_history_slots=bounded_int("run_history_slots", 168, 168, 168),
            max_items=bounded_int("max_items", 200, 1, 200),
            max_images_per_item=bounded_int("max_images_per_item", 20, 1, 20),
            max_image_bytes=bounded_int("max_image_bytes", 20 * 1024 * 1024, 1, 20 * 1024 * 1024),
            max_text_bytes=bounded_int("max_text_bytes", 2 * 1024 * 1024, 1, 2 * 1024 * 1024),
        )
    return CollectorConfig(
        creator_name=creator_name.strip(),
        creator_uid=creator_uid,
        creator_dynamic_url=dynamic_url,
        timezone_name=timezone_name,
        window_hours=window_hours,
@@ -443,6 +647,7 @@
        allowed_source_hosts=allowed_hosts,
        allowed_video_extensions=extensions,
        native_handoff=NativeHandoffRoute(**route_values),
        refresh=refresh,
    )
@@ -1412,6 +1617,11 @@
    handoff = subparsers.add_parser("handoff", help="Generate an unsent canonical Codex-native video handoff")
    handoff.add_argument("--output", type=Path, help="Optional no-overwrite Markdown output path")
    handoff.add_argument("--now", help="Optional offset-aware ISO-8601 clock for deterministic runs")
    begin = subparsers.add_parser("refresh-begin", help="Create one durable hourly browser-refresh run")
    begin.add_argument("--now", help="Optional offset-aware ISO-8601 clock for deterministic runs")
    commit = subparsers.add_parser("refresh-commit", help="Validate and commit one local browser evidence file")
    commit.add_argument("--input", required=True, type=Path, help="Final browser evidence JSON from refresh-begin")
    commit.add_argument("--now", help="Optional offset-aware ISO-8601 clock for deterministic runs")
    return parser
@@ -1431,10 +1641,19 @@
            if args.output:
                output = args.output if args.output.is_absolute() else Path.cwd() / args.output
            result = generate_handoff(config, output, now)
        elif args.command in {"refresh-begin", "refresh-commit"}:
            from bili_dynamic_refresh import refresh_begin, refresh_commit
            if config.refresh is None:
                raise CollectorError("E_CONFIG", "config.refresh is required for refresh commands.")
            if args.command == "refresh-begin":
                result = refresh_begin(config, config_path_value, now)
            else:
                result = refresh_commit(config, config_path_value, absolute_lexical(args.input), now)
        else:  # pragma: no cover - argparse owns this contract
            raise CollectorError("E_COMMAND", f"Unknown command: {args.command}")
    result = {"schema_version": SCHEMA_VERSION, "ok": True, **result}
    code = 4 if result.get("status") == "COMPLETE_WITH_RETAINED_SOURCE" else 0
    code = int(result.pop("exit_code", 4 if result.get("status") == "COMPLETE_WITH_RETAINED_SOURCE" else 0))
    return code, result