#!/usr/bin/env python3 """Configuration-driven Bilibili article/image collector and verifier. The browser side may supply only sanitized visible-page capture bundles. This module never reads browser storage, credentials, headers, or profiles and never performs network requests. It validates, publishes with CreateNew semantics, or verifies an existing append-only corpus. """ from __future__ import annotations import argparse import hashlib import json import os import re import stat import sys import uuid from dataclasses import dataclass from datetime import date, datetime, timedelta, timezone from pathlib import Path, PurePosixPath from typing import Any, Iterable, Mapping, Sequence from urllib.parse import urlsplit, urlunsplit from zoneinfo import ZoneInfo, ZoneInfoNotFoundError CONFIG_SCHEMA = 1 CAPTURE_SCHEMA = 1 MANIFEST_SCHEMA = 1 ALLOWED_TYPES = frozenset({"article", "text", "image"}) PENDING_STATES = frozenset({"VIDEO_ABSENT", "METADATA_NOT_READY", "OWNER_PENDING", "DIMENSIONS_PENDING"}) READY_STATE = "READY" ACCESS_STATES = frozenset({"ACCESS_BLOCKED", "LOGIN_REQUIRED", "CAPTCHA", "HTTP_412", "PAYWALL"}) SECRET_KEY = re.compile( r"(?:password|passwd|cookie|token|secret|authorization|captcha|session|localstorage|signed[_-]?url|" r"口令|密码|令牌|验证码|会话)", re.IGNORECASE, ) CONTROL = re.compile(r"[\x00-\x1f\x7f]") SHA256 = re.compile(r"[0-9A-F]{64}") UID = re.compile(r"[1-9][0-9]{0,19}") ITEM_ID = re.compile(r"[A-Za-z0-9_-]{1,128}") LEGACY_IMAGE_ID = re.compile(r"([A-Za-z0-9_-]{1,96}):image:([1-9][0-9]{0,3})") WINDOWS_BAD = re.compile(r'[<>:"/\\|?*\x00-\x1f]') WINDOWS_RESERVED = {"CON", "PRN", "AUX", "NUL", *(f"COM{i}" for i in range(1, 10)), *(f"LPT{i}" for i in range(1, 10))} FILE_ATTRIBUTE_REPARSE_POINT = 0x400 OWNED_PENDING = re.compile(r"\.bili-article-image\.pending\.[0-9a-f]{32}\.json") PROJECT_ROOT = Path(__file__).resolve().parents[2] class CollectorError(RuntimeError): def __init__(self, code: str, message: str, *, safety: bool = False) -> None: super().__init__(message) self.code = code self.message = message self.safety = safety @dataclass(frozen=True) class CollectorConfig: path: Path raw_bytes: bytes sha256: str creator_uid: str creator_name: str dynamic_url: str profile_url: str output_root: Path manifest_path: Path intake_root: Path summary_path: Path | None timezone_name: str date_start: datetime date_end: datetime include_types: frozenset[str] deadline_seconds: int observation_interval_ms: int stable_observations: int rerun_policy: str max_items: int max_body_bytes: int max_images_per_item: int max_image_bytes: int @property def tz(self) -> ZoneInfo: return ZoneInfo(self.timezone_name) def _exact_keys(value: Any, expected: Iterable[str], field: str) -> Mapping[str, Any]: expected_set = set(expected) if not isinstance(value, Mapping) or set(value) != expected_set: raise CollectorError("E_CONFIG_SCHEMA", f"{field} keys differ from the strict schema.", safety=True) return value def _reject_secrets(value: Any, path: str = "$") -> None: if isinstance(value, Mapping): for key, child in value.items(): if not isinstance(key, str) or SECRET_KEY.search(key): raise CollectorError("E_SECRET_FIELD", f"Secret-bearing field is forbidden at {path}.", safety=True) _reject_secrets(child, f"{path}.{key}") elif isinstance(value, list): for index, child in enumerate(value): _reject_secrets(child, f"{path}[{index}]") def _strict_json(path: Path, description: str) -> tuple[Any, bytes]: try: raw = path.read_bytes() except OSError as exc: raise CollectorError("E_INPUT", f"{description} is unreadable.", safety=True) from exc if raw.startswith(b"\xef\xbb\xbf") or b"\r" in raw or not raw.endswith(b"\n") or raw.endswith(b"\n\n"): raise CollectorError("E_INPUT_ENCODING", f"{description} must be strict UTF-8 LF with one final LF.", safety=True) try: text = raw[:-1].decode("utf-8", errors="strict") value = json.loads(text) except (UnicodeDecodeError, json.JSONDecodeError) as exc: raise CollectorError("E_INPUT_SCHEMA", f"{description} is not strict JSON.", safety=True) from exc _reject_secrets(value) return value, raw def _canonical_bytes(value: Any) -> bytes: return (json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":")) + "\n").encode("utf-8") def _canonical_payload(value: Any) -> bytes: return json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":")).encode("utf-8") def _sha256(path: Path) -> str: digest = hashlib.sha256() with path.open("rb") as handle: for chunk in iter(lambda: handle.read(1024 * 1024), b""): digest.update(chunk) return digest.hexdigest().upper() def _is_reparse(path: Path) -> bool: try: return bool(path.lstat().st_file_attributes & FILE_ATTRIBUTE_REPARSE_POINT) except AttributeError: return path.is_symlink() def _ordinary_file(path: Path) -> None: if not path.exists() or not path.is_file() or _is_reparse(path): raise CollectorError("E_PATH", "Required file is missing, non-ordinary, or reparse-backed.", safety=True) def _safe_existing_chain(path: Path, *, allow_missing_leaf: bool = False) -> None: candidate = path.resolve(strict=False) current = Path(candidate.anchor) parts = candidate.parts[1:] for index, part in enumerate(parts): current = current / part if not current.exists(): if allow_missing_leaf and index == len(parts) - 1: return continue if _is_reparse(current): raise CollectorError("E_PATH_REPARSE", "Path chain contains a reparse point.", safety=True) if index < len(parts) - 1 and not current.is_dir(): raise CollectorError("E_PATH", "Path chain contains a non-directory component.", safety=True) def _within(child: Path, parent: Path) -> bool: try: child.resolve(strict=False).relative_to(parent.resolve(strict=False)) return True except ValueError: return False def _safe_component(value: str, *, max_length: int = 80) -> str: cleaned = WINDOWS_BAD.sub("_", value).strip(" .") cleaned = re.sub(r"\s+", "", cleaned) if not cleaned or cleaned.upper() in WINDOWS_RESERVED: raise CollectorError("E_CONFIG", "creator.name cannot form a safe output component.", safety=True) return cleaned[:max_length].rstrip(" .") def _parse_datetime(value: Any, field: str) -> datetime: if not isinstance(value, str) or not value.strip(): raise CollectorError("E_CONFIG", f"{field} must be an offset-aware ISO-8601 string.") try: parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) except ValueError as exc: raise CollectorError("E_CONFIG", f"{field} is invalid.") from exc if parsed.tzinfo is None: raise CollectorError("E_CONFIG", f"{field} must include an offset.") return parsed.astimezone(timezone.utc) def _validate_page_url(value: Any, uid: str, *, dynamic: bool) -> str: if not isinstance(value, str) or CONTROL.search(value): raise CollectorError("E_CONFIG", "Configured page URL is invalid.", safety=True) parsed = urlsplit(value) expected_path = f"/{uid}/dynamic" if dynamic else f"/{uid}" if parsed.scheme != "https" or parsed.hostname != "space.bilibili.com" or parsed.query or parsed.fragment: raise CollectorError("E_CONFIG", "Configured page URL must be a query-free Bilibili space URL.", safety=True) if parsed.path.rstrip("/") != expected_path: raise CollectorError("E_CONFIG", "Configured page URL does not bind the creator UID.", safety=True) return urlunsplit(("https", "space.bilibili.com", expected_path, "", "")) def _bounded_int(value: Any, field: str, lower: int, upper: int) -> int: if not isinstance(value, int) or isinstance(value, bool) or not lower <= value <= upper: raise CollectorError("E_CONFIG", f"{field} must be {lower}..{upper}.") return value def load_config(path: Path, *, now: datetime | None = None) -> CollectorConfig: path = path.resolve(strict=True) if not _within(path, PROJECT_ROOT): raise CollectorError("E_PATH_ESCAPE", "config path must remain inside the project root.", safety=True) _ordinary_file(path) value, raw = _strict_json(path, "config") root = _exact_keys(value, {"schema_version", "creator", "page", "output", "selection", "readiness", "rerun", "verification", "limits"}, "config") if root["schema_version"] != CONFIG_SCHEMA: raise CollectorError("E_CONFIG_SCHEMA", "config.schema_version differs.") creator = _exact_keys(root["creator"], {"uid", "name"}, "creator") page = _exact_keys(root["page"], {"dynamic_url", "profile_url"}, "page") output = _exact_keys(root["output"], {"root", "intake_root", "manifest_name"}, "output") selection = _exact_keys(root["selection"], {"date_start", "date_end", "window_days", "timezone", "include_types"}, "selection") readiness = _exact_keys(root["readiness"], {"deadline_seconds", "observation_interval_ms", "stable_observations"}, "readiness") rerun = _exact_keys(root["rerun"], {"policy"}, "rerun") verification = _exact_keys(root["verification"], {"summary_path"}, "verification") limits = _exact_keys(root["limits"], {"max_items", "max_body_bytes", "max_images_per_item", "max_image_bytes"}, "limits") uid = creator["uid"] name = creator["name"] if not isinstance(uid, str) or UID.fullmatch(uid) is None: raise CollectorError("E_CONFIG", "creator.uid must be a positive decimal string.") if not isinstance(name, str) or not name.strip() or len(name.strip()) > 80 or CONTROL.search(name): raise CollectorError("E_CONFIG", "creator.name is invalid.") name = name.strip() dynamic_url = _validate_page_url(page["dynamic_url"], uid, dynamic=True) profile_url = _validate_page_url(page["profile_url"], uid, dynamic=False) project_root = PROJECT_ROOT root_value = output["root"] if root_value is None: output_root = project_root / "ana-data" / f"news-{_safe_component(name)}" elif isinstance(root_value, str) and root_value.strip(): raw_root = Path(root_value) output_root = raw_root if raw_root.is_absolute() else path.parent / raw_root else: raise CollectorError("E_CONFIG", "output.root must be null or a non-empty path.") output_root = output_root.resolve(strict=False) if not _within(output_root, project_root): raise CollectorError("E_PATH_ESCAPE", "output.root must remain inside the project root.", safety=True) _safe_existing_chain(output_root) intake_value = output["intake_root"] if not isinstance(intake_value, str) or not intake_value.strip(): raise CollectorError("E_CONFIG", "output.intake_root must be a non-empty path.") intake_root = Path(intake_value) intake_root = (intake_root if intake_root.is_absolute() else path.parent / intake_root).resolve(strict=False) if not _within(intake_root, project_root): raise CollectorError("E_PATH_ESCAPE", "output.intake_root must remain inside the project root.", safety=True) _safe_existing_chain(intake_root) manifest_name = output["manifest_name"] if manifest_name != "manifest.jsonl": raise CollectorError("E_CONFIG", "output.manifest_name must equal manifest.jsonl.") summary_value = verification["summary_path"] if summary_value is None: summary_path = None elif isinstance(summary_value, str) and summary_value.strip(): summary_relative = PurePosixPath(summary_value) if summary_relative.is_absolute() or ".." in summary_relative.parts or not summary_relative.parts: raise CollectorError("E_CONFIG", "verification.summary_path must be output-root relative.", safety=True) summary_path = (output_root / Path(*summary_relative.parts)).resolve(strict=False) if not _within(summary_path, output_root): raise CollectorError("E_PATH_ESCAPE", "verification.summary_path escapes output.root.", safety=True) _ordinary_file(summary_path) else: raise CollectorError("E_CONFIG", "verification.summary_path must be null or a relative path.") timezone_name = selection["timezone"] if not isinstance(timezone_name, str): raise CollectorError("E_CONFIG", "selection.timezone must be a zoneinfo name.") try: tz = ZoneInfo(timezone_name) except ZoneInfoNotFoundError as exc: raise CollectorError("E_CONFIG", "selection.timezone is unknown.") from exc explicit = selection["date_start"] is not None or selection["date_end"] is not None window_days = selection["window_days"] if explicit: if selection["date_start"] is None or selection["date_end"] is None or window_days is not None: raise CollectorError("E_CONFIG", "Use either date_start/date_end or window_days.") start = _parse_datetime(selection["date_start"], "selection.date_start") end = _parse_datetime(selection["date_end"], "selection.date_end") else: days = _bounded_int(window_days, "selection.window_days", 1, 366) current = (now or datetime.now(timezone.utc)).astimezone(tz) end_local = current start_local = current - timedelta(days=days) start, end = start_local.astimezone(timezone.utc), end_local.astimezone(timezone.utc) if start > end or end - start > timedelta(days=366): raise CollectorError("E_CONFIG", "Configured selection interval is invalid.") include_raw = selection["include_types"] if not isinstance(include_raw, list) or not include_raw or len(include_raw) != len(set(include_raw)): raise CollectorError("E_CONFIG", "selection.include_types must be a non-empty unique list.") include_types = frozenset(include_raw) if not include_types.issubset(ALLOWED_TYPES): raise CollectorError("E_CONFIG", "selection.include_types contains an unsupported type.") policy = rerun["policy"] if policy not in {"verify_only", "verify_or_append"}: raise CollectorError("E_CONFIG", "rerun.policy is unsupported.") return CollectorConfig( path=path, raw_bytes=raw, sha256=hashlib.sha256(raw).hexdigest().upper(), creator_uid=uid, creator_name=name, dynamic_url=dynamic_url, profile_url=profile_url, output_root=output_root, manifest_path=output_root / "manifest.jsonl", intake_root=intake_root, summary_path=summary_path, timezone_name=timezone_name, date_start=start, date_end=end, include_types=include_types, deadline_seconds=_bounded_int(readiness["deadline_seconds"], "readiness.deadline_seconds", 5, 600), observation_interval_ms=_bounded_int(readiness["observation_interval_ms"], "readiness.observation_interval_ms", 100, 10000), stable_observations=_bounded_int(readiness["stable_observations"], "readiness.stable_observations", 2, 5), rerun_policy=policy, max_items=_bounded_int(limits["max_items"], "limits.max_items", 1, 1000), max_body_bytes=_bounded_int(limits["max_body_bytes"], "limits.max_body_bytes", 1, 8 * 1024 * 1024), max_images_per_item=_bounded_int(limits["max_images_per_item"], "limits.max_images_per_item", 0, 50), max_image_bytes=_bounded_int(limits["max_image_bytes"], "limits.max_image_bytes", 1, 50 * 1024 * 1024), ) def _capture_page_identity(value: Mapping[str, Any], config: CollectorConfig) -> None: if value.get("creator_uid") != config.creator_uid or value.get("creator_name") != config.creator_name: raise CollectorError("E_CREATOR_IDENTITY", "Capture creator differs from config.", safety=True) if value.get("dynamic_url") != config.dynamic_url or value.get("profile_url") != config.profile_url: raise CollectorError("E_PAGE_IDENTITY", "Capture page proof differs from config.", safety=True) def _accepted_snapshot_sha256(items: Sequence[Mapping[str, Any]]) -> str: canonical_items = [] for item in items: published = item["published_at"] if not isinstance(published, datetime): raise CollectorError("E_READINESS_DIGEST", "Accepted snapshot timestamp is not normalized.", safety=True) if published.microsecond % 1000: raise CollectorError("E_READINESS_DIGEST", "Accepted snapshot timestamp exceeds browser millisecond precision.", safety=True) canonical_items.append({ "body_complete": True, "body_text": item["body"].decode("utf-8", errors="strict").rstrip("\n"), "image_count": len(item["images"]), "item_type": item["item_type"], "published_at_epoch_ms": int(published.timestamp() * 1000), "source_url": item["source_url"], "stable_id": item["stable_id"], "title": item["title"], }) canonical_items.sort(key=lambda item: item["stable_id"]) snapshot = {"items": canonical_items, "schema_version": 1} return hashlib.sha256(_canonical_payload(snapshot)).hexdigest().upper() def _stable_readiness(observations: Any, config: CollectorConfig, expected_fingerprint: str) -> tuple[int, int, str]: if not isinstance(observations, list) or not observations: raise CollectorError("E_READINESS", "Capture requires readiness observations.") streak = 0 prior_fingerprint: str | None = None previous_elapsed = -1 for index, raw in enumerate(observations): value = _exact_keys(raw, {"elapsed_ms", "state", "reason", "snapshot_sha256"}, f"observations[{index}]") elapsed = value["elapsed_ms"] state = value["state"] reason = value["reason"] fingerprint = value["snapshot_sha256"] if not isinstance(elapsed, int) or isinstance(elapsed, bool) or elapsed <= previous_elapsed or elapsed > config.deadline_seconds * 1000: raise CollectorError("E_READINESS", "Observation clock is invalid.", safety=True) previous_elapsed = elapsed if state in ACCESS_STATES: raise CollectorError("E_ACCESS_CONTROL", "Visible page reported an access-control stop.", safety=True) if state in PENDING_STATES: if reason != state or fingerprint is not None: raise CollectorError("E_READINESS", "Pending observation shape is invalid.", safety=True) streak, prior_fingerprint = 0, None continue if state != READY_STATE or reason != "READY" or not isinstance(fingerprint, str) or SHA256.fullmatch(fingerprint) is None: raise CollectorError("E_READINESS", "Observation state is unsupported.", safety=True) if fingerprint != expected_fingerprint: raise CollectorError("E_READINESS_DIGEST", "READY fingerprint does not bind the accepted items snapshot.", safety=True) if fingerprint == prior_fingerprint: streak += 1 else: streak, prior_fingerprint = 1, fingerprint if streak < config.stable_observations: raise CollectorError("E_READINESS_TIMEOUT", "Stable READY evidence is not the terminal observation suffix.", safety=True) return len(observations), previous_elapsed, expected_fingerprint def validate_capture(config: CollectorConfig, capture_path: Path) -> dict[str, Any]: capture_path = capture_path.resolve(strict=True) if not _within(capture_path, PROJECT_ROOT): raise CollectorError("E_PATH_ESCAPE", "capture path must remain inside the project root.", safety=True) _ordinary_file(capture_path) root, raw = _strict_json(capture_path, "capture") value = _exact_keys(root, {"schema_version", "creator_uid", "creator_name", "dynamic_url", "profile_url", "observations", "items"}, "capture") if value["schema_version"] != CAPTURE_SCHEMA: raise CollectorError("E_CAPTURE_SCHEMA", "capture.schema_version differs.") _capture_page_identity(value, config) items = value["items"] if not isinstance(items, list) or len(items) > config.max_items: raise CollectorError("E_CAPTURE_SCHEMA", "capture.items exceeds the configured bound.") normalized: list[dict[str, Any]] = [] seen: set[str] = set() for index, raw_item in enumerate(items): item = _exact_keys(raw_item, {"stable_id", "item_type", "title", "source_url", "published_at", "body_text", "body_complete", "images"}, f"items[{index}]") stable_id = item["stable_id"] item_type = item["item_type"] if not isinstance(stable_id, str) or ITEM_ID.fullmatch(stable_id) is None or stable_id in seen: raise CollectorError("E_ITEM_IDENTITY", "Item stable identity is invalid or duplicated.", safety=True) seen.add(stable_id) if item_type not in config.include_types: raise CollectorError("E_ITEM_TYPE", "Capture item type is outside configured include_types.", safety=True) source = urlsplit(str(item["source_url"])) if source.scheme != "https" or source.hostname != "www.bilibili.com" or source.query or source.fragment or source.path.rstrip("/") != f"/opus/{stable_id}": raise CollectorError("E_ITEM_IDENTITY", "Item source URL does not bind the stable ID.", safety=True) published = _parse_datetime(item["published_at"], f"items[{index}].published_at") if published < config.date_start or published > config.date_end: raise CollectorError("E_ITEM_WINDOW", "Capture item is outside the configured interval.", safety=True) title = item["title"] body = item["body_text"] if not isinstance(title, str) or not title.strip() or CONTROL.search(title) or not isinstance(body, str) or not body.strip() or not item["body_complete"]: raise CollectorError("E_CONTENT_INCOMPLETE", "Item title/body is incomplete.", safety=True) body_bytes = body.replace("\r\n", "\n").replace("\r", "\n").rstrip("\n").encode("utf-8") + b"\n" if len(body_bytes) > config.max_body_bytes: raise CollectorError("E_CONTENT_LIMIT", "Item body exceeds the configured limit.") images_raw = item["images"] if not isinstance(images_raw, list) or len(images_raw) > config.max_images_per_item: raise CollectorError("E_CONTENT_LIMIT", "Item images exceed the configured limit.") if item_type == "image" and not images_raw: raise CollectorError("E_CONTENT_INCOMPLETE", "Image item requires at least one original image.") images: list[dict[str, Any]] = [] for sequence, raw_image in enumerate(images_raw, 1): image = _exact_keys(raw_image, {"path", "bytes", "sha256", "extension"}, f"items[{index}].images[{sequence - 1}]") relative = PurePosixPath(str(image["path"])) if relative.is_absolute() or ".." in relative.parts or not relative.parts: raise CollectorError("E_ARTIFACT_PATH", "Image intake path is unsafe.", safety=True) source_path = (config.intake_root / Path(*relative.parts)).resolve(strict=False) if not _within(source_path, config.intake_root): raise CollectorError("E_ARTIFACT_PATH", "Image intake path escapes its root.", safety=True) _ordinary_file(source_path) extension = image["extension"] if extension not in {".jpg", ".jpeg", ".png", ".webp"} or source_path.suffix.lower() != extension: raise CollectorError("E_ARTIFACT", "Image extension is unsupported.") size = image["bytes"] digest = image["sha256"] if not isinstance(size, int) or size <= 0 or size > config.max_image_bytes or not isinstance(digest, str) or SHA256.fullmatch(digest) is None: raise CollectorError("E_ARTIFACT", "Image identity is invalid.") if source_path.stat().st_size != size or _sha256(source_path) != digest: raise CollectorError("E_ARTIFACT_HASH", "Image identity differs from intake bytes.", safety=True) head = source_path.read_bytes()[:12] if extension in {".jpg", ".jpeg"} and not head.startswith(b"\xff\xd8\xff"): raise CollectorError("E_ARTIFACT", "JPEG magic differs.") if extension == ".png" and not head.startswith(b"\x89PNG\r\n\x1a\n"): raise CollectorError("E_ARTIFACT", "PNG magic differs.") if extension == ".webp" and not (head.startswith(b"RIFF") and head[8:12] == b"WEBP"): raise CollectorError("E_ARTIFACT", "WebP magic differs.") images.append({"sequence": sequence, "source": source_path, "bytes": size, "sha256": digest, "extension": extension}) normalized.append({ "stable_id": stable_id, "item_type": item_type, "title": title.strip(), "source_url": urlunsplit(("https", "www.bilibili.com", f"/opus/{stable_id}", "", "")), "published_at": published, "body": body_bytes, "images": images, }) fingerprint = _accepted_snapshot_sha256(normalized) attempts, elapsed, fingerprint = _stable_readiness(value["observations"], config, fingerprint) return { "capture_bytes": len(raw), "capture_sha256": hashlib.sha256(raw).hexdigest().upper(), "readiness_attempts": attempts, "readiness_elapsed_ms": elapsed, "snapshot_sha256": fingerprint, "items": normalized, } def _same_file_identity(left: os.stat_result, right: os.stat_result) -> bool: left_inode = (getattr(left, "st_dev", 0), getattr(left, "st_ino", 0)) right_inode = (getattr(right, "st_dev", 0), getattr(right, "st_ino", 0)) return left_inode == right_inode and stat.S_ISREG(left.st_mode) and stat.S_ISREG(right.st_mode) def _freeze_image_payload(image: Mapping[str, Any], max_bytes: int) -> bytes: path = image["source"] _ordinary_file(path) before = path.lstat() flags = os.O_RDONLY | getattr(os, "O_BINARY", 0) | getattr(os, "O_NOFOLLOW", 0) try: descriptor = os.open(path, flags) except OSError as exc: raise CollectorError("E_ARTIFACT_DRIFT", "Image intake cannot be opened without following indirection.", safety=True) from exc try: opened = os.fstat(descriptor) if not _same_file_identity(before, opened) or _is_reparse(path): raise CollectorError("E_ARTIFACT_DRIFT", "Image intake identity changed before precommit.", safety=True) chunks: list[bytes] = [] total = 0 while True: chunk = os.read(descriptor, min(1024 * 1024, max_bytes + 1 - total)) if not chunk: break chunks.append(chunk) total += len(chunk) if total > max_bytes: raise CollectorError("E_ARTIFACT_DRIFT", "Image intake exceeds its configured bound at precommit.", safety=True) payload = b"".join(chunks) finally: os.close(descriptor) after = path.lstat() if ( not _same_file_identity(opened, after) or _is_reparse(path) or opened.st_size != after.st_size or getattr(opened, "st_mtime_ns", None) != getattr(after, "st_mtime_ns", None) ): raise CollectorError("E_ARTIFACT_DRIFT", "Image intake changed during the final precommit read.", safety=True) digest = hashlib.sha256(payload).hexdigest().upper() if len(payload) != image["bytes"] or digest != image["sha256"]: raise CollectorError("E_ARTIFACT_DRIFT", "Image intake bytes differ from the validated identity at precommit.", safety=True) extension = image["extension"] head = payload[:12] if extension in {".jpg", ".jpeg"} and not head.startswith(b"\xff\xd8\xff"): raise CollectorError("E_ARTIFACT_DRIFT", "JPEG magic differs at precommit.", safety=True) if extension == ".png" and not head.startswith(b"\x89PNG\r\n\x1a\n"): raise CollectorError("E_ARTIFACT_DRIFT", "PNG magic differs at precommit.", safety=True) if extension == ".webp" and not (head.startswith(b"RIFF") and head[8:12] == b"WEBP"): raise CollectorError("E_ARTIFACT_DRIFT", "WebP magic differs at precommit.", safety=True) return payload def _read_manifest(path: Path) -> list[dict[str, Any]]: if not path.exists(): return [] _ordinary_file(path) raw = path.read_bytes() if raw.startswith(b"\xef\xbb\xbf") or b"\r" in raw or (raw and not raw.endswith(b"\n")): raise CollectorError("E_MANIFEST", "Manifest encoding is invalid.", safety=True) events: list[dict[str, Any]] = [] for line_number, line in enumerate(raw.splitlines(), 1): if not line: raise CollectorError("E_MANIFEST", "Manifest contains a blank line.", safety=True) try: value = json.loads(line.decode("utf-8", errors="strict")) except (UnicodeDecodeError, json.JSONDecodeError) as exc: raise CollectorError("E_MANIFEST", f"Manifest line {line_number} is invalid.", safety=True) from exc if not isinstance(value, dict): raise CollectorError("E_MANIFEST", "Manifest row is not an object.", safety=True) events.append(value) return events def _manifest_snapshot(config: CollectorConfig) -> tuple[int, str]: if not config.manifest_path.exists(): return 0, hashlib.sha256(b"").hexdigest().upper() raw = config.manifest_path.read_bytes() return len(raw), hashlib.sha256(raw).hexdigest().upper() def _create_new(path: Path, payload: bytes) -> None: path.parent.mkdir(parents=True, exist_ok=True) _safe_existing_chain(path.parent) try: with path.open("xb") as handle: handle.write(payload) handle.flush() os.fsync(handle.fileno()) except FileExistsError as exc: raise CollectorError("E_TARGET_EXISTS", "CreateNew target already exists.", safety=True) from exc def _terminal_write(path: Path | None, value: Mapping[str, Any]) -> str: payload = _canonical_bytes(value) if path is None: return "STDOUT_ONLY" target = path.resolve(strict=False) if not _within(target, PROJECT_ROOT): raise CollectorError("E_PATH_ESCAPE", "terminal path must remain inside the project root.", safety=True) _safe_existing_chain(target, allow_missing_leaf=True) if target.exists(): _ordinary_file(target) if target.read_bytes() != payload: raise CollectorError("E_TERMINAL_CONFLICT", "Terminal path exists with different bytes.", safety=True) return "REUSED" _create_new(target, payload) return "CREATED" def _owned_recovery_evidence(config: CollectorConfig) -> list[Path]: if not config.output_root.exists(): return [] if not config.output_root.is_dir() or _is_reparse(config.output_root): raise CollectorError("E_PATH", "Output root is not an ordinary directory.", safety=True) evidence: list[Path] = [] for candidate in config.output_root.iterdir(): if OWNED_PENDING.fullmatch(candidate.name): _ordinary_file(candidate) evidence.append(candidate) return sorted(evidence, key=lambda item: item.name) def collect(config: CollectorConfig, capture_path: Path, terminal_path: Path | None) -> dict[str, Any]: if config.rerun_policy != "verify_or_append": raise CollectorError("E_RERUN_POLICY", "collect requires rerun.policy=verify_or_append.", safety=True) if _owned_recovery_evidence(config): raise CollectorError("E_RECOVERY_REQUIRED", "Owned pending evidence requires separate recovery.", safety=True) capture = validate_capture(config, capture_path) events = _read_manifest(config.manifest_path) latest = {str(row.get("stable_id")): row for row in events if isinstance(row.get("stable_id"), str)} plan: list[dict[str, Any]] = [] for item in capture["items"]: prior = latest.get(item["stable_id"]) if prior and prior.get("status") == "SAVED": continue local = item["published_at"].astimezone(config.tz) stem = f"{local:%Y%m%d-%H%M%S}_{item['item_type']}_{_safe_component(item['title'], max_length=48)}_{item['stable_id']}" text_name = f"{stem}.txt" images = [f"{stem}_{index:02d}{image['extension']}" for index, image in enumerate(item["images"], 1)] targets = [config.output_root / text_name, *(config.output_root / name for name in images)] if any(target.exists() for target in targets): raise CollectorError("E_TARGET_EXISTS", "A planned output target already exists.", safety=True) plan.append({"item": item, "text_name": text_name, "image_names": images, "targets": targets}) if not plan: terminal = { "schema_version": 1, "status": "NO_NEW_ITEMS", "creator_uid": config.creator_uid, "creator_name": config.creator_name, "config_sha256": config.sha256, "capture_sha256": capture["capture_sha256"], "input_items": len(capture["items"]), "new_items": 0, "mutation_count": 0, } terminal["terminal_disposition"] = _terminal_write(terminal_path, terminal) return terminal manifest_before = _manifest_snapshot(config) for planned in plan: planned["frozen_images"] = [ {"identity": image, "payload": _freeze_image_payload(image, config.max_image_bytes)} for image in planned["item"]["images"] ] if _manifest_snapshot(config) != manifest_before or any(target.exists() for planned in plan for target in planned["targets"]): raise CollectorError("E_PRECOMMIT_DRIFT", "Manifest or target state changed during image intake freeze.", safety=True) config.output_root.mkdir(parents=True, exist_ok=True) _safe_existing_chain(config.output_root) created: list[Path] = [] rows: list[bytes] = [] pending_path = config.output_root / f".bili-article-image.pending.{uuid.uuid4().hex}.json" pending = { "schema_version": 1, "status": "PUBLISH_PENDING", "config_sha256": config.sha256, "capture_sha256": capture["capture_sha256"], "manifest_bytes": manifest_before[0], "manifest_sha256": manifest_before[1], "targets": [target.name for planned in plan for target in planned["targets"]], } _create_new(pending_path, _canonical_bytes(pending)) try: for planned in plan: item = planned["item"] _create_new(planned["targets"][0], item["body"]) created.append(planned["targets"][0]) image_refs: list[dict[str, Any]] = [] for frozen, name, target in zip(planned["frozen_images"], planned["image_names"], planned["targets"][1:]): image = frozen["identity"] _create_new(target, frozen["payload"]) created.append(target) image_refs.append({"path": name, "bytes": image["bytes"], "sha256": image["sha256"]}) row = { "schema_version": MANIFEST_SCHEMA, "creator": config.creator_name, "creator_uid": config.creator_uid, "item_type": item["item_type"], "stable_id": item["stable_id"], "title": item["title"], "source_url": item["source_url"], "published_at": item["published_at"].isoformat(), "collected_at": datetime.now(timezone.utc).isoformat(), "status": "SAVED", "path": planned["text_name"], "bytes": len(item["body"]), "sha256": hashlib.sha256(item["body"]).hexdigest().upper(), "images": image_refs, "capture_method": "authenticated_visible_dom_config_bound", "readiness_observation_count": capture["readiness_attempts"], "config_sha256": config.sha256, } rows.append(_canonical_bytes(row)) if _manifest_snapshot(config) != manifest_before: raise CollectorError("E_PRECOMMIT_DRIFT", "Manifest changed before append.", safety=True) with config.manifest_path.open("ab") as handle: for row in rows: handle.write(row) handle.flush() os.fsync(handle.fileno()) appended = b"".join(rows) manifest_after = config.manifest_path.read_bytes() if not manifest_after.endswith(appended) or len(manifest_after) != manifest_before[0] + len(appended): raise CollectorError("E_RECOVERY_REQUIRED", "Manifest append readback is ambiguous.", safety=True) pending_path.unlink() except Exception as failure: try: manifest_unchanged = _manifest_snapshot(config) == manifest_before except Exception: manifest_unchanged = False cleanup_ok = manifest_unchanged if manifest_unchanged: for path in reversed(created): try: path.unlink() except OSError: cleanup_ok = False if cleanup_ok: try: pending_path.unlink() except OSError: cleanup_ok = False if not cleanup_ok: raise CollectorError("E_RECOVERY_REQUIRED", "Publish state is ambiguous; owned evidence was preserved.", safety=True) from failure raise terminal = { "schema_version": 1, "status": "CONTENT_SAVED", "creator_uid": config.creator_uid, "creator_name": config.creator_name, "config_sha256": config.sha256, "capture_sha256": capture["capture_sha256"], "input_items": len(capture["items"]), "new_items": len(plan), "artifact_count": len(created), "mutation_count": len(created) + 1, "manifest_bytes": config.manifest_path.stat().st_size, "manifest_sha256": _sha256(config.manifest_path), } terminal["terminal_disposition"] = _terminal_write(terminal_path, terminal) return terminal def _validate_artifact(root: Path, relative: Any, expected_bytes: Any, expected_sha: Any) -> Path: if not isinstance(relative, str): raise CollectorError("E_MANIFEST", "Artifact path is missing.", safety=True) parsed = PurePosixPath(relative) if parsed.is_absolute() or ".." in parsed.parts or not parsed.parts: raise CollectorError("E_MANIFEST", "Artifact path is unsafe.", safety=True) path = (root / Path(*parsed.parts)).resolve(strict=False) if not _within(path, root): raise CollectorError("E_MANIFEST", "Artifact path escapes output root.", safety=True) _ordinary_file(path) if not isinstance(expected_bytes, int) or expected_bytes < 1 or not isinstance(expected_sha, str) or SHA256.fullmatch(expected_sha) is None: raise CollectorError("E_MANIFEST", "Artifact identity is invalid.", safety=True) if path.stat().st_size != expected_bytes or _sha256(path) != expected_sha: raise CollectorError("E_MANIFEST_DRIFT", "Artifact bytes or SHA-256 drifted.", safety=True) return path def _verify_summary(config: CollectorConfig, manifest_latest: Mapping[str, Mapping[str, Any]]) -> tuple[int, int, int, int]: if config.summary_path is None: raise CollectorError("E_SUMMARY", "Summary path is absent.", safety=True) root, _ = _strict_json(config.summary_path, "verification summary") value = _exact_keys( root, { "schema_version", "creator", "creator_uid", "date_window", "generated_at", "item_count", "article_count", "text_count", "image_count", "body_bytes_total", "manifest", "rows", }, "verification summary", ) if value["schema_version"] != 1 or value["creator"] != config.creator_name or str(value["creator_uid"]) != config.creator_uid: raise CollectorError("E_SUMMARY_IDENTITY", "Verification summary creator differs from config.", safety=True) window = _exact_keys(value["date_window"], {"start", "end"}, "verification summary date_window") if _parse_datetime(window["start"], "summary.date_window.start") != config.date_start or _parse_datetime(window["end"], "summary.date_window.end") != config.date_end: raise CollectorError("E_SUMMARY_IDENTITY", "Verification summary interval differs from config.", safety=True) manifest = _exact_keys( value["manifest"], {"path", "prefix_lines", "prefix_sha256", "appended_rows", "appended_block_sha256", "final_lines", "final_bytes", "final_sha256"}, "verification summary manifest", ) manifest_raw = config.manifest_path.read_bytes() manifest_lines = manifest_raw.splitlines(keepends=True) final_lines = manifest.get("final_lines") if not isinstance(final_lines, int) or isinstance(final_lines, bool) or final_lines < 0 or final_lines > len(manifest_lines): raise CollectorError("E_SUMMARY_MANIFEST_DRIFT", "Verification summary manifest prefix length is invalid.", safety=True) frozen_prefix = b"".join(manifest_lines[:final_lines]) if ( manifest["path"] != "manifest.jsonl" or manifest["final_bytes"] != len(frozen_prefix) or manifest["final_sha256"] != hashlib.sha256(frozen_prefix).hexdigest().upper() ): raise CollectorError("E_SUMMARY_MANIFEST_DRIFT", "Verification summary does not bind the immutable manifest prefix.", safety=True) rows = value["rows"] if not isinstance(rows, list) or len(rows) != value["item_count"] or len(rows) > config.max_items: raise CollectorError("E_SUMMARY", "Verification summary item count is invalid.", safety=True) seen: set[str] = set() article_count = 0 text_count = 0 image_paths: set[str] = set() for index, raw_row in enumerate(rows): row = _exact_keys( raw_row, {"opus_id", "title", "content_type", "published_at", "source_url", "body_bytes", "body_sha256", "text", "images"}, f"verification summary rows[{index}]", ) stable_id = row["opus_id"] if not isinstance(stable_id, str) or ITEM_ID.fullmatch(stable_id) is None or stable_id in seen: raise CollectorError("E_SUMMARY", "Verification summary has an invalid or duplicate opus ID.", safety=True) seen.add(stable_id) item_type = row["content_type"] if item_type not in {"article", "text"} or item_type not in config.include_types: raise CollectorError("E_SUMMARY", "Verification summary item type is outside config.", safety=True) published = _parse_datetime(row["published_at"], f"summary.rows[{index}].published_at") source = urlsplit(str(row["source_url"])) if ( published < config.date_start or published > config.date_end or source.scheme != "https" or source.hostname != "www.bilibili.com" or source.query or source.fragment or source.path.rstrip("/") != f"/opus/{stable_id}" ): raise CollectorError("E_SUMMARY_IDENTITY", "Verification summary item proof is invalid.", safety=True) if not isinstance(row["body_bytes"], int) or row["body_bytes"] < 1 or not isinstance(row["body_sha256"], str) or SHA256.fullmatch(row["body_sha256"]) is None: raise CollectorError("E_SUMMARY", "Verification summary body identity is invalid.", safety=True) text = _exact_keys(row["text"], {"path", "bytes", "sha256"}, f"verification summary rows[{index}].text") _validate_artifact(config.output_root, text["path"], text["bytes"], text["sha256"]) manifest_row = manifest_latest.get(stable_id) if ( manifest_row is None or manifest_row.get("status") != "SAVED" or manifest_row.get("path") != text["path"] or manifest_row.get("bytes") != text["bytes"] or manifest_row.get("sha256") != text["sha256"] ): raise CollectorError("E_SUMMARY_MANIFEST_DRIFT", "Verification summary item differs from append-only manifest.", safety=True) if item_type == "article": article_count += 1 else: text_count += 1 images = row["images"] if not isinstance(images, list) or len(images) > config.max_images_per_item: raise CollectorError("E_SUMMARY", "Verification summary images are invalid.", safety=True) for image_index, raw_image in enumerate(images): image = _exact_keys(raw_image, {"path", "bytes", "sha256", "source_url"}, f"verification summary rows[{index}].images[{image_index}]") source_image = urlsplit(str(image["source_url"])) if source_image.scheme != "https" or source_image.hostname not in {"i0.hdslb.com", "i1.hdslb.com", "i2.hdslb.com"} or source_image.query or source_image.fragment or not source_image.path.startswith("/bfs/"): raise CollectorError("E_SUMMARY_IDENTITY", "Verification summary image source proof is invalid.", safety=True) path = _validate_artifact(config.output_root, image["path"], image["bytes"], image["sha256"]) image_paths.add(path.name) if ( value["article_count"] != article_count or value["text_count"] != text_count or value["image_count"] != len(image_paths) or value["item_count"] != article_count + text_count ): raise CollectorError("E_SUMMARY", "Verification summary aggregate counts differ from rows.", safety=True) return len(rows), article_count, text_count, len(image_paths) def verify(config: CollectorConfig, terminal_path: Path | None) -> dict[str, Any]: if _owned_recovery_evidence(config): raise CollectorError("E_RECOVERY_REQUIRED", "Owned pending evidence requires separate recovery.", safety=True) events = _read_manifest(config.manifest_path) latest: dict[str, dict[str, Any]] = {} for row in events: if row.get("creator") != config.creator_name: continue row_uid = row.get("creator_uid") if row_uid is not None and (isinstance(row_uid, bool) or str(row_uid) != config.creator_uid): raise CollectorError("E_CREATOR_IDENTITY", "Manifest creator UID conflicts with config.", safety=True) item_type = row.get("item_type") if item_type not in config.include_types: continue _reject_secrets(row, "$selected_manifest") published = _parse_datetime(row.get("published_at"), "manifest.published_at") if published < config.date_start or published > config.date_end: continue stable_id = row.get("stable_id") legacy_image = LEGACY_IMAGE_ID.fullmatch(stable_id) if isinstance(stable_id, str) else None if not isinstance(stable_id, str) or (ITEM_ID.fullmatch(stable_id) is None and legacy_image is None): raise CollectorError("E_MANIFEST", "Manifest stable ID is invalid.", safety=True) source = urlsplit(str(row.get("source_url", ""))) if legacy_image is not None: parent_id = row.get("source_parent_stable_id") if ( row.get("item_type") != "image" or parent_id != legacy_image.group(1) or source.scheme != "https" or source.hostname not in {"i0.hdslb.com", "i1.hdslb.com", "i2.hdslb.com"} or source.query or source.fragment or not source.path.startswith("/bfs/") ): raise CollectorError("E_MANIFEST", "Legacy image source proof is invalid.", safety=True) elif source.scheme != "https" or source.hostname != "www.bilibili.com" or source.query or source.fragment or source.path.rstrip("/") != f"/opus/{stable_id}": raise CollectorError("E_MANIFEST", "Manifest source URL is invalid.", safety=True) latest[stable_id] = row if config.summary_path is not None: item_count, article_count, non_article_count, image_count = _verify_summary(config, latest) else: article_count = 0 non_article_count = 0 image_paths: set[str] = set() item_count = 0 for stable_id, row in ([] if config.summary_path is not None else latest.items()): if row.get("status") != "SAVED": raise CollectorError("E_CORPUS_INCOMPLETE", "Latest in-window item is not SAVED.", safety=True) _validate_artifact(config.output_root, row.get("path"), row.get("bytes"), row.get("sha256")) if row.get("item_type") == "article": article_count += 1 else: non_article_count += 1 images = row.get("images") if images is not None: if not isinstance(images, list): raise CollectorError("E_MANIFEST", "Manifest images is not a list.", safety=True) for image in images: if not isinstance(image, Mapping): raise CollectorError("E_MANIFEST", "Manifest image is not an object.", safety=True) path = _validate_artifact(config.output_root, image.get("path"), image.get("bytes"), image.get("sha256")) image_paths.add(path.name) elif row.get("image_path") is not None: path = _validate_artifact(config.output_root, row.get("image_path"), row.get("image_bytes"), row.get("image_sha256")) image_paths.add(path.name) elif row.get("item_type") == "image": image_paths.add(Path(str(row["path"])).name) if config.summary_path is None: item_count = len(latest) image_count = len(image_paths) manifest_bytes, manifest_sha = _manifest_snapshot(config) terminal = { "schema_version": 1, "status": "CORPUS_VERIFIED", "creator_uid": config.creator_uid, "creator_name": config.creator_name, "config_sha256": config.sha256, "date_start": config.date_start.isoformat(), "date_end": config.date_end.isoformat(), "include_types": sorted(config.include_types), "item_count": item_count, "article_count": article_count, "text_image_dynamic_count": non_article_count, "original_image_count": image_count, "manifest_bytes": manifest_bytes, "manifest_sha256": manifest_sha, "mutation_count": 0, } terminal["terminal_disposition"] = _terminal_write(terminal_path, terminal) return terminal def build_parser() -> argparse.ArgumentParser: parser = argparse.ArgumentParser(description="Generic logged-session Bilibili article/image collector") parser.add_argument("--config", required=True, type=Path, help="Strict UTF-8 JSON config") subparsers = parser.add_subparsers(dest="command", required=True) validate = subparsers.add_parser("validate-capture", help="Validate one sanitized browser capture without publishing") validate.add_argument("--capture", required=True, type=Path) collect_parser = subparsers.add_parser("collect", help="CreateNew-publish one validated capture") collect_parser.add_argument("--capture", required=True, type=Path) collect_parser.add_argument("--terminal", type=Path) verify_parser = subparsers.add_parser("verify", help="Read-only verify the configured corpus") verify_parser.add_argument("--terminal", type=Path) return parser def run(argv: Sequence[str] | None = None) -> tuple[int, dict[str, Any]]: args = build_parser().parse_args(argv) try: config = load_config(args.config) if args.command == "validate-capture": capture = validate_capture(config, args.capture) result = { "schema_version": 1, "status": "CAPTURE_VALID", "config_sha256": config.sha256, "capture_sha256": capture["capture_sha256"], "item_count": len(capture["items"]), "readiness_attempts": capture["readiness_attempts"], "mutation_count": 0, } elif args.command == "collect": result = collect(config, args.capture, args.terminal) else: result = verify(config, args.terminal) return 0, result except CollectorError as exc: return (3 if exc.safety else 2), { "schema_version": 1, "status": "SAFETY_STOP" if exc.safety else "INPUT_ERROR", "error_code": exc.code, "message": exc.message, "mutation_count": 0, } except KeyboardInterrupt: return 130, {"schema_version": 1, "status": "INTERRUPTED", "mutation_count": 0} except Exception: return 1, {"schema_version": 1, "status": "INTERNAL_ERROR", "error_code": "E_INTERNAL", "mutation_count": 0} def main(argv: Sequence[str] | None = None) -> int: code, result = run(argv) print(json.dumps(result, ensure_ascii=False, sort_keys=True, separators=(",", ":"))) return code if __name__ == "__main__": raise SystemExit(main())