#!/usr/bin/env python3
|
"""Trusted, injectable controller for one supported Chrome refresh run.
|
|
Production orchestration supplies a reviewed adapter for the supported Chrome
|
calls. Tests use a deterministic fake implementing the same four methods. This
|
module owns dispatch state, monotonic timing, extraction adaptation, attestation,
|
and the no-overwrite evidence commit; callers never select runtime outcomes.
|
"""
|
from __future__ import annotations
|
|
import hashlib
|
import json
|
import time
|
from datetime import datetime, timedelta, timezone
|
from pathlib import Path
|
from typing import Any, Callable, Mapping, Protocol
|
|
import bili_dynamic_collector as core
|
import bili_dynamic_refresh as refresh
|
|
|
class SupportedChromeAdapter(Protocol):
|
def open_tabs(self) -> list[Mapping[str, Any]]: ...
|
def reload(self, tab: Mapping[str, Any], timeout_seconds: int) -> None: ...
|
def goto(self, url: str, timeout_seconds: int) -> Mapping[str, Any]: ...
|
def evaluate(self, tab: Mapping[str, Any], source: str, timeout_seconds: int) -> Mapping[str, Any]: ...
|
|
|
def _millis(clock: Callable[[], float]) -> int:
|
return int(clock() * 1000)
|
|
|
def _sha_material(value: str) -> str:
|
return hashlib.sha256(value.encode("utf-8")).hexdigest()
|
|
|
def _adapt_extract(config: core.CollectorConfig, pending: Mapping[str, Any], raw: Mapping[str, Any], observed_at: str) -> tuple[dict[str, Any], list[dict[str, Any]], dict[str, Any] | None]:
|
cards: list[dict[str, Any]] = []
|
items: list[dict[str, Any]] = []
|
positions: list[tuple[int, str, Mapping[str, Any]]] = []
|
for entry in raw.get("cards", []):
|
positions.append((int(entry["position"]), "card", entry))
|
for entry in raw.get("unparsed_nodes", []):
|
positions.append((int(entry["position"]), "unparsed", entry))
|
positions.sort()
|
remap = {old: index for index, (old, _, _) in enumerate(positions)}
|
for old_position, kind, entry in positions:
|
if kind != "card":
|
continue
|
normalized = core.normalize_item(entry, config, len(cards))
|
cards.append({
|
"position": remap[old_position],
|
"identifiers": {"dynamic_id": normalized["dynamic_id"], "opus_id": normalized["opus_id"], "bvid": normalized["bvid"]},
|
"stable_keys": normalized["dedupe_keys"], "published_at": normalized["published_at"],
|
"content_type": normalized["content_type"], "source_url": normalized["source_url"],
|
})
|
items.append({
|
**{key: normalized[key] for key in ("dynamic_id", "opus_id", "bvid", "content_type", "published_at", "title", "source_url")},
|
"body_text": entry.get("body_text") or "", "body_complete": bool(entry.get("body_complete")),
|
"duration_seconds": entry.get("duration_seconds"), "artifacts": [],
|
})
|
unparsed = [{"position": remap[old], "node_fingerprint_sha256": _sha_material(str(entry["fingerprint_material"])), "reason_code": entry["reason_code"]} for old, kind, entry in positions if kind == "unparsed"]
|
observation = {"ordinal": 0, "observed_at": observed_at, "cursor_before": 0, "cursor_after": len(positions), "visible_node_count": len(positions), "complete_card_count": len(cards), "unparsed_node_count": len(unparsed), "cards": cards, "unparsed_nodes": unparsed, "limit_hit": raw.get("limit_hit", "NONE")}
|
marker = None
|
if raw.get("terminal_marker_text"):
|
text = str(raw["terminal_marker_text"])
|
marker = {"observation_ordinal": 0, "kind": "EXACT_END_OF_FEED", "selector_id": "DYNAMIC_FEED_END_TEXT", "normalized_text": text, "marker_sha256": hashlib.sha256(f"DYNAMIC_FEED_END_TEXT\n{text}".encode()).hexdigest()}
|
return observation, items, marker
|
|
|
def _commit_create_new(path: Path, evidence: Mapping[str, Any]) -> None:
|
payload = core.canonical_json_bytes(evidence, newline=False)
|
core.ensure_directory(path.parent, create=True)
|
partial = path.parent / f".{path.name}.{hashlib.sha256(payload).hexdigest()[:16]}.partial"
|
core.lexical_lstat_chain(partial, allow_missing_leaf=True)
|
try:
|
with partial.open("xb") as stream:
|
stream.write(payload); stream.flush(); __import__("os").fsync(stream.fileno())
|
refresh._rename_no_overwrite(partial, path)
|
refresh._fsync_directory(path.parent)
|
if path.read_bytes() != payload:
|
raise core.CollectorError("E_EVIDENCE_COMMIT", "Controller evidence durable readback mismatch.", safety=True)
|
finally:
|
try: partial.unlink()
|
except FileNotFoundError: pass
|
|
|
def _discard_exact_created(path: Path, evidence: Mapping[str, Any]) -> None:
|
"""Remove only the exact evidence created by this controller invocation."""
|
payload = core.canonical_json_bytes(evidence, newline=False)
|
core.lexical_lstat_chain(path, allow_missing_leaf=False)
|
if not path.is_file() or path.read_bytes() != payload:
|
raise core.CollectorError(
|
"E_EVIDENCE_COMMIT",
|
"Late controller evidence no longer matches the owned payload.",
|
safety=True,
|
)
|
path.unlink()
|
refresh._fsync_directory(path.parent)
|
if path.exists():
|
raise core.CollectorError(
|
"E_EVIDENCE_COMMIT",
|
"Late controller evidence cleanup did not persist.",
|
safety=True,
|
)
|
|
|
def run_once(config: core.CollectorConfig, begin: Mapping[str, Any], adapter: SupportedChromeAdapter, *, monotonic: Callable[[], float] = time.monotonic, wall_now: Callable[[], datetime] = lambda: datetime.now(timezone.utc)) -> Path:
|
pending = refresh._load_pending(config)
|
if not pending or pending["run_id"] != begin.get("run_id") or pending["schema_version"] != refresh.PENDING_SCHEMA:
|
raise core.CollectorError("E_CONTROLLER_PENDING", "Controller requires the exact active runtime-v2 pending.", safety=True)
|
run_ms = _millis(monotonic); deadline_ms = run_ms + config.refresh.overall_deadline_seconds * 1000
|
action_start = action_finish = observation_start = observation_finish = None
|
action_dispatched = False; action = None; action_outcome = "PRE_DISPATCH_ERROR"; observation_outcome = "NOT_ATTEMPTED"; raw = None
|
try:
|
tabs = [tab for tab in adapter.open_tabs() if tab.get("url") == config.creator_dynamic_url]
|
except Exception:
|
tabs = None
|
if tabs is None or len(tabs) > 1:
|
tab = None
|
else:
|
try:
|
action_start = _millis(monotonic); action_dispatched = True
|
if tabs: action = "reload"; tab = tabs[0]; adapter.reload(tab, config.refresh.refresh_action_timeout_seconds)
|
else: action = "navigate"; tab = adapter.goto(config.creator_dynamic_url, config.refresh.refresh_action_timeout_seconds)
|
action_outcome = "CONFIRMED"
|
except TimeoutError:
|
action_outcome = "TIMEOUT"; tab = tabs[0] if tabs else {"url": config.creator_dynamic_url}
|
except Exception:
|
action_outcome = "POST_DISPATCH_ERROR"; tab = tabs[0] if tabs else {"url": config.creator_dynamic_url}
|
finally:
|
if action_dispatched: action_finish = _millis(monotonic)
|
if action_dispatched and action_finish is not None and action_finish < deadline_ms:
|
try:
|
observation_start = _millis(monotonic)
|
extractor_call = (
|
refresh.EXTRACTOR_SOURCE.read_text(encoding="utf-8")
|
+ "\nprojectInfoCollectVisibleDynamicNodes({page_internal_settle_timeout_ms:15000})"
|
)
|
raw = adapter.evaluate(tab, extractor_call, min(config.refresh.observation_timeout_seconds, max(0, (deadline_ms - observation_start) // 1000)))
|
observation_outcome = "READABLE"
|
except TimeoutError: observation_outcome = "TIMEOUT"
|
except PermissionError: observation_outcome = "ACCESS_BLOCKED"
|
except Exception: observation_outcome = "ERROR"
|
finally: observation_finish = _millis(monotonic)
|
elif action_dispatched:
|
observation_outcome = "DEADLINE_EXHAUSTED"
|
write_ms = _millis(monotonic)
|
if write_ms > deadline_ms:
|
raise core.CollectorError("E_OVERALL_DEADLINE", "Runtime controller exceeded its total deadline.", safety=True)
|
wall_started = core.parse_datetime(pending["started_at"], "pending.started_at")
|
action_wall_ms = 0 if action_start is None else action_finish - action_start
|
observation_wall_ms = 0 if observation_start is None else observation_finish - observation_start
|
refresh_finished_at = wall_started + timedelta(milliseconds=action_wall_ms)
|
read_finished_at = refresh_finished_at + timedelta(milliseconds=observation_wall_ms)
|
observed_at = core.canonical_datetime(read_finished_at)
|
observation = items = marker = None
|
if observation_outcome == "READABLE" and raw is not None:
|
observation, items, marker = _adapt_extract(config, pending, raw, observed_at)
|
page_outcome = {"READABLE":"READABLE", "TIMEOUT":"UNREADABLE_TIMEOUT", "ERROR":"UNREADABLE_ERROR", "ACCESS_BLOCKED":"ACCESS_BLOCKED", "NOT_ATTEMPTED":"UNREADABLE_ERROR", "DEADLINE_EXHAUSTED":"UNREADABLE_TIMEOUT"}[observation_outcome]
|
diagnostics = {"PRE_DISPATCH_ERROR":"ACTION_PRE_DISPATCH", "POST_DISPATCH_ERROR":"ACTION_POST_DISPATCH", "TIMEOUT":"ACTION_TIMEOUT"}.get(action_outcome) or {"READABLE":"NONE", "TIMEOUT":"OBSERVATION_TIMEOUT", "ERROR":"OBSERVATION_ERROR", "ACCESS_BLOCKED":"ACCESS_INTERSTITIAL", "DEADLINE_EXHAUSTED":"DEADLINE_EXHAUSTED"}.get(observation_outcome, "ACTION_PRE_DISPATCH")
|
contract_raw = refresh.OBSERVATION_CONTRACT.read_bytes(); runtime_raw = refresh.RUNTIME_CONTRACT.read_bytes()
|
creator = raw.get("creator", {}) if raw else {"uid": None, "name": "", "profile_url": None}
|
evidence = {"schema_version": 3, "run_id": pending["run_id"], "transport": "codex_chrome_visible_page", "requested_url": config.creator_dynamic_url, "final_url": raw.get("final_url", config.creator_dynamic_url) if raw else config.creator_dynamic_url, "refresh_action": action, "refresh_count": 1 if action_dispatched else 0, "refresh_started_at": pending["started_at"], "refresh_finished_at": core.canonical_datetime(refresh_finished_at), "read_finished_at": observed_at, "page_outcome": page_outcome, "page_title": raw.get("page_title", "") if raw else "", "creator": creator, "extractor": {"contract_id": json.loads(contract_raw)["contract_id"], "contract_sha256": hashlib.sha256(contract_raw).hexdigest(), "parser_version": json.loads(contract_raw)["parser_version"], "parser_sha256": hashlib.sha256(refresh.EXTRACTOR_SOURCE.read_bytes()).hexdigest()}, "page_observation": None if observation is None else {"schema_version": 1, "limits": json.loads(contract_raw)["limits"], "observations": [observation], "terminal_marker": marker}, "items": items or [], "discovery_summary": {"status": "NOT_USED", "item_count": 0}, "safe_diagnostics": {"code": diagnostics, "overall_deadline_seconds": 120, "refresh_action_timeout_seconds": 35, "observation_timeout_seconds": 45}, "runtime_contract": {"contract_id": pending["runtime_contract"]["contract_id"], "contract_bytes": len(runtime_raw), "contract_sha256": hashlib.sha256(runtime_raw).hexdigest()}, "runtime_observation": {"refresh_action_outcome": action_outcome, "refresh_action_elapsed_ms": action_wall_ms, "refresh_count": 1 if action_dispatched else 0, "observation_outcome": observation_outcome, "observation_elapsed_ms": observation_wall_ms, "observation_count": 0 if observation_start is None else 1}}
|
refresh._attest_controller_evidence(pending, evidence, action_dispatched=action_dispatched, monotonic_run_started_ms=run_ms, monotonic_action_started_ms=action_start, monotonic_action_finished_ms=action_finish, monotonic_observation_started_ms=observation_start, monotonic_observation_finished_ms=observation_finish, monotonic_evidence_write_started_ms=write_ms)
|
path = Path(pending["evidence_path"])
|
_commit_create_new(path, evidence)
|
committed_ms = _millis(monotonic)
|
transitioned_at = wall_now()
|
wall_deadline = core.parse_datetime(pending["deadline_at"], "pending.deadline_at")
|
if committed_ms > deadline_ms or transitioned_at > wall_deadline:
|
_discard_exact_created(path, evidence)
|
raise core.CollectorError(
|
"E_OVERALL_DEADLINE",
|
"Controller evidence did not durably complete within the total deadline.",
|
safety=True,
|
)
|
refresh.bind_controller_evidence(config, pending, path, transitioned_at=transitioned_at)
|
return path
|