from __future__ import annotations
|
|
import hashlib
|
import importlib.util
|
import json
|
import gc
|
import os
|
import shutil
|
import subprocess
|
import sys
|
import tempfile
|
import time
|
import unittest
|
from pathlib import Path
|
|
|
PROJECT_DEV = Path(__file__).resolve().parents[1]
|
if str(PROJECT_DEV) not in sys.path:
|
sys.path.insert(0, str(PROJECT_DEV))
|
COLLECTOR_PATH = PROJECT_DEV / "bili_dynamic_collector.py"
|
SPEC = importlib.util.spec_from_file_location("bili_dynamic_collector", COLLECTOR_PATH)
|
assert SPEC and SPEC.loader
|
collector = importlib.util.module_from_spec(SPEC)
|
sys.modules[SPEC.name] = collector
|
SPEC.loader.exec_module(collector)
|
import bili_dynamic_refresh as refresh
|
import bili_dynamic_refresh_controller as controller
|
|
|
NOW = "2026-08-13T10:00:00+08:00"
|
|
|
class BiliDynamicRefreshTests(unittest.TestCase):
|
def setUp(self) -> None:
|
self.temp = tempfile.TemporaryDirectory()
|
self.root = Path(self.temp.name)
|
self.state = self.root / "state"
|
self.archive = self.root / "archive"
|
self.intake = self.root / "intake"
|
self.downloads = self.root / "downloads"
|
self.videos = self.root / "videos"
|
for path in (self.archive, self.intake, self.downloads, self.videos):
|
path.mkdir()
|
self.config_path = self.root / "config.json"
|
self.write_json(
|
self.config_path,
|
{
|
"schema_version": 1,
|
"creator": {
|
"name": "青枫浦上Q",
|
"uid": "1420210197",
|
"dynamic_url": "https://space.bilibili.com/1420210197/dynamic",
|
},
|
"timezone": "Asia/Shanghai",
|
"window_hours": 72,
|
"minimum_complete_age_seconds": 0,
|
"title_max_length": 48,
|
"allowed_source_hosts": [
|
"space.bilibili.com", "www.bilibili.com", "t.bilibili.com", "b23.tv"
|
],
|
"allowed_video_extensions": [".mp4", ".mkv", ".mov", ".webm"],
|
"native_handoff": {
|
"project_id": "project-info",
|
"source_ai_id": "video-downloader",
|
"source_thread_id": "019fcc5d-798f-7ea1-8325-3a4d1f2dc5a5",
|
"source_role_instance_id": "case_analysis.video_downloader",
|
"target_ai_id": "media-processor",
|
"target_thread_id": "019fb7a4-bdfd-79f2-bd6b-e67e2b7d8efd",
|
"target_role_instance_id": "case_analysis.media_processor",
|
"reply_thread_id": "019fcc5d-798f-7ea1-8325-3a4d1f2dc5a5",
|
},
|
"paths": {"state_dir": "state", "download_dir": "downloads", "video_dir": "videos"},
|
"refresh": {
|
"archive_dir": "archive",
|
"formal_manifest": "archive/manifest.jsonl",
|
"intake_dir": "intake",
|
"overall_deadline_seconds": 120,
|
"refresh_action_timeout_seconds": 35,
|
"observation_timeout_seconds": 45,
|
"page_internal_settle_timeout_seconds": 15,
|
"max_refresh_count": 1,
|
"run_history_slots": 168,
|
"max_items": 200,
|
"max_images_per_item": 20,
|
"max_image_bytes": 20971520,
|
"max_text_bytes": 2097152,
|
},
|
},
|
)
|
|
def tearDown(self) -> None:
|
self.temp.cleanup()
|
|
def write_json(self, path: Path, value: object) -> None:
|
if isinstance(value, dict) and value.get("schema_version") == 3 and isinstance(value.get("controller_attestation"), dict):
|
config = collector.load_config(self.config_path)
|
pending = refresh._load_pending(config)
|
if pending is not None:
|
attestation = value["controller_attestation"]
|
refresh._attest_controller_evidence(
|
pending, value,
|
action_dispatched=bool(attestation["action_dispatched"]),
|
monotonic_run_started_ms=int(attestation["monotonic_run_started_ms"]),
|
monotonic_action_started_ms=attestation["monotonic_action_started_ms"],
|
monotonic_action_finished_ms=attestation["monotonic_action_finished_ms"],
|
monotonic_observation_started_ms=attestation["monotonic_observation_started_ms"],
|
monotonic_observation_finished_ms=attestation["monotonic_observation_finished_ms"],
|
monotonic_evidence_write_started_ms=int(attestation["monotonic_evidence_write_started_ms"]),
|
)
|
path.parent.mkdir(parents=True, exist_ok=True)
|
path.write_text(json.dumps(value, ensure_ascii=False), encoding="utf-8")
|
|
def reset_runtime_fixture(self) -> None:
|
"""Reset only disposable per-test runtime/formal trees between matrix rows."""
|
for path in (self.state, self.archive, self.intake):
|
if path.exists():
|
shutil.rmtree(path)
|
path.mkdir()
|
|
def begin(self) -> dict[str, object]:
|
code, result = collector.run(["--config", str(self.config_path), "refresh-begin", "--now", NOW])
|
self.assertEqual(0, code)
|
self.assertEqual("BROWSER_REFRESH_REQUIRED", result["status"])
|
return result
|
|
def formal_pending(self, config: collector.CollectorConfig) -> dict[str, object]:
|
return {
|
"schema_version": 2,
|
"run_id": "a" * 32,
|
"owner_nonce": "b" * 32,
|
"phase": "TRANSACTION_INTENT",
|
"task_id": refresh.TASK_ID,
|
"creator_uid": config.creator_uid,
|
"creator_dynamic_url": config.creator_dynamic_url,
|
"started_at": "2026-08-13T02:00:00Z",
|
"deadline_at": "2026-08-13T02:00:30Z",
|
"window_start": "2026-08-10T02:00:00Z",
|
"window_end": "2026-08-13T02:00:00Z",
|
"config_sha256": hashlib.sha256(self.config_path.read_bytes()).hexdigest(),
|
"state_manifest_preimage": refresh._identity(config.manifest_path).as_dict(),
|
"formal_manifest_preimage": refresh._identity(config.refresh.formal_manifest).as_dict(),
|
"evidence_path": str(config.state_dir / "refresh" / "incoming" / ("a" * 32 + ".json")),
|
"intake_root": str(config.refresh.intake_dir / ("a" * 32)),
|
"evidence_identity": {"bytes": 2, "sha256": hashlib.sha256(b"{}").hexdigest()},
|
"planned_terminal": None,
|
"transaction_identity": {
|
"transaction_id": "c" * 32,
|
"state_preimage": refresh._identity(config.manifest_path).as_dict(),
|
"formal_preimage": refresh._identity(config.refresh.formal_manifest).as_dict(),
|
"state_candidate": refresh._identity(config.manifest_path).as_dict(),
|
"formal_candidate": refresh._identity(config.refresh.formal_manifest).as_dict(),
|
"created_artifacts": [],
|
"input_item_count": 0,
|
"new_item_count": 0,
|
"formal_lock_claim": None,
|
},
|
"last_transition_at": "2026-08-13T02:00:00Z",
|
}
|
|
def extractor(self) -> dict[str, object]:
|
contract_raw = refresh.OBSERVATION_CONTRACT.read_bytes()
|
contract = json.loads(contract_raw)
|
return {
|
"contract_id": contract["contract_id"],
|
"contract_sha256": hashlib.sha256(contract_raw).hexdigest(),
|
"parser_version": contract["parser_version"],
|
"parser_sha256": hashlib.sha256(refresh.EXTRACTOR_SOURCE.read_bytes()).hexdigest(),
|
}
|
|
def evidence(
|
self,
|
begin: dict[str, object],
|
*,
|
outcome: str,
|
observations: list[dict[str, object]] | None,
|
marker: dict[str, object] | None,
|
items: list[dict[str, object]] | None = None,
|
discovery: str = "NOT_USED",
|
action_outcome: str = "CONFIRMED",
|
observation_outcome: str | None = None,
|
refresh_count: int = 1,
|
observation_count: int | None = None,
|
) -> dict[str, object]:
|
outcome_map = {
|
"READABLE": "READABLE",
|
"UNREADABLE_TIMEOUT": "TIMEOUT",
|
"UNREADABLE_ERROR": "ERROR",
|
"ACCESS_BLOCKED": "ACCESS_BLOCKED",
|
}
|
observation_outcome = observation_outcome or outcome_map[outcome]
|
if action_outcome == "PRE_DISPATCH_ERROR":
|
refresh_started = refresh_finished = read_finished = "2026-08-13T02:00:01Z"
|
action_elapsed = observation_elapsed = 0
|
observation_count = 0
|
refresh_action = None
|
else:
|
refresh_started = "2026-08-13T02:00:01Z"
|
refresh_finished = "2026-08-13T02:00:16Z"
|
read_finished = "2026-08-13T02:00:24Z"
|
action_elapsed = 15000
|
observation_elapsed = 8000
|
observation_count = 1 if observation_count is None else observation_count
|
refresh_action = "reload"
|
evidence = {
|
"schema_version": 3,
|
"run_id": begin["run_id"],
|
"transport": "codex_chrome_visible_page",
|
"requested_url": "https://space.bilibili.com/1420210197/dynamic",
|
"final_url": "https://space.bilibili.com/1420210197/dynamic",
|
"refresh_action": refresh_action,
|
"refresh_count": refresh_count,
|
"refresh_started_at": refresh_started,
|
"refresh_finished_at": refresh_finished,
|
"read_finished_at": read_finished,
|
"page_outcome": outcome,
|
"page_title": "青枫浦上Q个人动态-青枫浦上Q动态记录-哔哩哔哩视频",
|
"creator": {"uid": "1420210197", "name": "青枫浦上Q", "profile_url": "https://space.bilibili.com/1420210197"},
|
"extractor": self.extractor(),
|
"page_observation": None if observations is None else {
|
"schema_version": 1,
|
"limits": json.loads(refresh.OBSERVATION_CONTRACT.read_text(encoding="utf-8"))["limits"],
|
"observations": observations,
|
"terminal_marker": marker,
|
},
|
"items": items or [],
|
"discovery_summary": {"status": discovery, "item_count": 0},
|
"safe_diagnostics": {
|
"code": "OBSERVATION_TIMEOUT" if outcome == "UNREADABLE_TIMEOUT" else "OBSERVATION_ERROR" if outcome == "UNREADABLE_ERROR" else "ACCESS_INTERSTITIAL" if outcome == "ACCESS_BLOCKED" else "NONE",
|
"overall_deadline_seconds": 120,
|
"refresh_action_timeout_seconds": 35,
|
"observation_timeout_seconds": 45,
|
},
|
"runtime_contract": {
|
"contract_id": begin["runtime_contract_id"],
|
"contract_bytes": Path(refresh.RUNTIME_CONTRACT).stat().st_size,
|
"contract_sha256": begin["runtime_contract_sha256"],
|
},
|
"runtime_observation": {
|
"refresh_action_outcome": action_outcome,
|
"refresh_action_elapsed_ms": action_elapsed,
|
"refresh_count": refresh_count,
|
"observation_outcome": observation_outcome,
|
"observation_elapsed_ms": observation_elapsed,
|
"observation_count": observation_count,
|
},
|
}
|
config = collector.load_config(self.config_path)
|
pending = refresh._load_pending(config)
|
assert pending is not None
|
if action_outcome == "PRE_DISPATCH_ERROR":
|
action_started = action_finished = observation_started = observation_finished = None
|
else:
|
action_started = 1000
|
action_finished = 1000 + action_elapsed
|
if observation_count:
|
observation_started = action_finished
|
observation_finished = observation_started + observation_elapsed
|
else:
|
observation_started = observation_finished = None
|
write_started = max(
|
1000,
|
action_finished or 0,
|
observation_finished or 0,
|
)
|
refresh._attest_controller_evidence(
|
pending,
|
evidence,
|
action_dispatched=action_outcome != "PRE_DISPATCH_ERROR",
|
monotonic_run_started_ms=0,
|
monotonic_action_started_ms=action_started,
|
monotonic_action_finished_ms=action_finished,
|
monotonic_observation_started_ms=observation_started,
|
monotonic_observation_finished_ms=observation_finished,
|
monotonic_evidence_write_started_ms=write_started,
|
)
|
return evidence
|
|
def set_runtime_diagnostic(self, evidence: dict[str, object], code: str) -> None:
|
diagnostics = evidence["safe_diagnostics"]
|
assert isinstance(diagnostics, dict)
|
diagnostics["code"] = code
|
|
@staticmethod
|
def observation(cards: list[dict[str, object]], unparsed: list[dict[str, object]] | None = None) -> dict[str, object]:
|
unparsed = unparsed or []
|
return {
|
"ordinal": 0,
|
"observed_at": "2026-08-13T02:00:24Z",
|
"cursor_before": 0,
|
"cursor_after": len(cards) + len(unparsed),
|
"visible_node_count": len(cards) + len(unparsed),
|
"complete_card_count": len(cards),
|
"unparsed_node_count": len(unparsed),
|
"cards": cards,
|
"unparsed_nodes": unparsed,
|
"limit_hit": "NONE",
|
}
|
|
def commit(self, begin: dict[str, object], evidence: dict[str, object]) -> tuple[int, dict[str, object]]:
|
if "controller_attestation" in evidence:
|
config = collector.load_config(self.config_path)
|
pending = refresh._load_pending(config)
|
assert pending is not None
|
attestation = evidence["controller_attestation"]
|
assert isinstance(attestation, dict)
|
refresh._attest_controller_evidence(
|
pending,
|
evidence,
|
action_dispatched=bool(attestation["action_dispatched"]),
|
monotonic_run_started_ms=int(attestation["monotonic_run_started_ms"]),
|
monotonic_action_started_ms=attestation["monotonic_action_started_ms"],
|
monotonic_action_finished_ms=attestation["monotonic_action_finished_ms"],
|
monotonic_observation_started_ms=attestation["monotonic_observation_started_ms"],
|
monotonic_observation_finished_ms=attestation["monotonic_observation_finished_ms"],
|
monotonic_evidence_write_started_ms=int(attestation["monotonic_evidence_write_started_ms"]),
|
)
|
path = Path(str(begin["evidence_path"]))
|
self.write_json(path, evidence)
|
config = collector.load_config(self.config_path)
|
pending = refresh._load_pending(config)
|
assert pending is not None
|
refresh.bind_controller_evidence(
|
config, pending, path,
|
transitioned_at=collector.parse_datetime("2026-08-13T10:00:24+08:00", "test transitioned_at"),
|
)
|
return collector.run(["--config", str(self.config_path), "refresh-commit", "--input", str(path), "--now", "2026-08-13T10:00:25+08:00"])
|
|
@staticmethod
|
def bound_card_item(identifier: str = "456") -> tuple[dict[str, object], dict[str, object]]:
|
source = f"https://www.bilibili.com/opus/{identifier}"
|
card = {
|
"position": 0,
|
"identifiers": {"dynamic_id": identifier, "opus_id": identifier, "bvid": None},
|
"stable_keys": [f"dynamic:{identifier}", f"opus:{identifier}", f"url:{source}"],
|
"published_at": "2026-08-13T01:55:00Z",
|
"content_type": "text",
|
"source_url": source,
|
}
|
item = {
|
"dynamic_id": identifier, "opus_id": identifier, "bvid": None, "content_type": "text",
|
"published_at": "2026-08-13T01:55:00Z", "title": "offline body",
|
"source_url": source, "body_text": "offline body", "body_complete": True,
|
"duration_seconds": None, "artifacts": [],
|
}
|
return card, item
|
|
@staticmethod
|
def end_marker() -> dict[str, object]:
|
selector = "DYNAMIC_FEED_END_TEXT"
|
text = "没有更多动态了"
|
return {
|
"observation_ordinal": 0, "kind": "EXACT_END_OF_FEED", "selector_id": selector,
|
"normalized_text": text, "marker_sha256": hashlib.sha256(f"{selector}\n{text}".encode()).hexdigest(),
|
}
|
|
def test_real_timeout_and_412_replay_never_claim_no_new(self) -> None:
|
begin = self.begin()
|
before = (self.archive / "manifest.jsonl").read_bytes() if (self.archive / "manifest.jsonl").exists() else b""
|
evidence = self.evidence(begin, outcome="UNREADABLE_TIMEOUT", observations=None, marker=None, discovery="BLOCKED_412")
|
code, result = self.commit(begin, evidence)
|
self.assertEqual(4, code)
|
self.assertEqual("REFRESH_FAILED_PAGE_UNREADABLE", result["status"])
|
self.assertFalse(result["no_new_confirmed"])
|
self.assertFalse(result["formal_manifest_changed"])
|
after = (self.archive / "manifest.jsonl").read_bytes() if (self.archive / "manifest.jsonl").exists() else b""
|
self.assertEqual(before, after)
|
self.assertFalse((self.state / "refresh" / "pending.json").exists())
|
slot = json.loads(Path(str(result["run_evidence_path"])).read_text(encoding="utf-8"))
|
self.assertEqual("TERMINAL", slot["run_state"])
|
self.assertEqual("BLOCKED_412", evidence["discovery_summary"]["status"])
|
|
def test_empty_page_requires_exact_marker_for_no_new(self) -> None:
|
begin = self.begin()
|
config = collector.load_config(self.config_path)
|
pending_before = refresh._load_pending(config)
|
self.assertIsNotNone(pending_before)
|
observation = self.observation([])
|
marker = {
|
"observation_ordinal": 0,
|
"kind": "EXACT_END_OF_FEED",
|
"selector_id": "DYNAMIC_FEED_END_TEXT",
|
"normalized_text": "没有更多动态了",
|
"marker_sha256": hashlib.sha256("DYNAMIC_FEED_END_TEXT\n没有更多动态了".encode()).hexdigest(),
|
}
|
code, result = self.commit(begin, self.evidence(begin, outcome="READABLE", observations=[observation], marker=marker))
|
self.assertEqual(0, code)
|
self.assertEqual("REFRESH_CONFIRMED_NO_NEW", result["status"])
|
self.assertTrue(result["coverage_complete"])
|
slot_path = Path(str(result["run_evidence_path"]))
|
slot_payload = slot_path.read_bytes()
|
slot = json.loads(slot_payload)
|
receipt = slot["transaction_receipt"]
|
self.assertEqual("NO_FORMAL_CHANGE", receipt["business_commit_kind"])
|
self.assertEqual(receipt["state_preimage"], receipt["state_candidate"])
|
self.assertEqual(receipt["formal_preimage"], receipt["formal_candidate"])
|
latest_path = self.state / "refresh" / "runs" / "latest.json"
|
latest = json.loads(latest_path.read_text(encoding="utf-8"))
|
self.assertEqual(
|
{"schema_version", "slot_index", "hour_epoch", "run_id", "slot_bytes", "slot_sha256", "status", "terminal_at"},
|
set(latest),
|
)
|
self.assertEqual(len(slot_payload), latest["slot_bytes"])
|
self.assertEqual(hashlib.sha256(slot_payload).hexdigest(), latest["slot_sha256"])
|
|
latest_path.unlink()
|
replay_pending = dict(pending_before)
|
replay_pending["phase"] = "TERMINAL_RECORDED"
|
replay_pending["planned_terminal"] = {key: slot[key] for key in (
|
"terminal_at", "status", "error_code", "exit_code", "coverage_proof", "state_manifest",
|
"formal_manifest", "artifact_tree_sha256", "refresh_action", "refresh_count", "page_authoritative",
|
"coverage_complete", "evidence_sha256", "input_item_count", "new_item_count", "saved_artifact_count",
|
"transaction_receipt",
|
)}
|
replay_pending["planned_terminal"]["formal_manifest_changed"] = False
|
refresh._write_pending(config, replay_pending, create=True)
|
replay = refresh._recover_or_replay(
|
config, self.config_path, collector.parse_now("2026-08-13T10:00:26+08:00")
|
)
|
self.assertEqual("REFRESH_CONFIRMED_NO_NEW", replay["status"])
|
self.assertEqual(latest, json.loads(latest_path.read_text(encoding="utf-8")))
|
|
def test_unparsed_node_is_mutually_exclusive_and_forces_partial(self) -> None:
|
begin = self.begin()
|
node = {"position": 0, "node_fingerprint_sha256": "a" * 64, "reason_code": "PARSER_REJECTED"}
|
observation = self.observation([], [node])
|
marker = {
|
"observation_ordinal": 0,
|
"kind": "EXACT_END_OF_FEED",
|
"selector_id": "DYNAMIC_FEED_END_TEXT",
|
"normalized_text": "已经到底了",
|
"marker_sha256": hashlib.sha256("DYNAMIC_FEED_END_TEXT\n已经到底了".encode()).hexdigest(),
|
}
|
code, result = self.commit(begin, self.evidence(begin, outcome="READABLE", observations=[observation], marker=marker))
|
self.assertEqual(4, code)
|
self.assertEqual("PARTIAL_DISCOVERY_UNCONFIRMED", result["status"])
|
self.assertFalse(result["coverage_complete"])
|
|
def test_node_position_overlap_fails_before_formal_write(self) -> None:
|
begin = self.begin()
|
card = {
|
"position": 0,
|
"identifiers": {"dynamic_id": "123", "opus_id": None, "bvid": None},
|
"stable_keys": ["dynamic:123", "url:https://www.bilibili.com/opus/123"],
|
"published_at": "2026-08-13T01:55:00Z",
|
"content_type": "text",
|
"source_url": "https://www.bilibili.com/opus/123",
|
}
|
node = {"position": 0, "node_fingerprint_sha256": "b" * 64, "reason_code": "PARSER_REJECTED"}
|
observation = self.observation([card], [node])
|
path = Path(str(begin["evidence_path"]))
|
self.write_json(path, self.evidence(begin, outcome="READABLE", observations=[observation], marker=None))
|
with self.assertRaisesRegex(collector.CollectorError, "positions overlap"):
|
collector.run(["--config", str(self.config_path), "refresh-commit", "--input", str(path), "--now", "2026-08-13T10:00:25+08:00"])
|
self.assertFalse((self.archive / "manifest.jsonl").exists())
|
|
def test_complete_observed_card_with_missing_item_cannot_claim_no_new(self) -> None:
|
begin = self.begin()
|
card, _ = self.bound_card_item()
|
evidence = self.evidence(
|
begin, outcome="READABLE", observations=[self.observation([card])],
|
marker=self.end_marker(), items=[],
|
)
|
path = Path(str(begin["evidence_path"]))
|
self.write_json(path, evidence)
|
with self.assertRaises(collector.CollectorError) as failure:
|
collector.run(["--config", str(self.config_path), "refresh-commit", "--input", str(path), "--now", "2026-08-13T10:00:25+08:00"])
|
self.assertEqual("E_EVIDENCE_ITEM_BINDING", failure.exception.code)
|
self.assertFalse((self.archive / "manifest.jsonl").exists())
|
|
def test_extra_item_without_observed_card_fails_closed(self) -> None:
|
begin = self.begin()
|
_, item = self.bound_card_item()
|
evidence = self.evidence(
|
begin, outcome="READABLE", observations=[self.observation([])],
|
marker=self.end_marker(), items=[item],
|
)
|
path = Path(str(begin["evidence_path"]))
|
self.write_json(path, evidence)
|
with self.assertRaises(collector.CollectorError) as failure:
|
collector.run(["--config", str(self.config_path), "refresh-commit", "--input", str(path), "--now", "2026-08-13T10:00:25+08:00"])
|
self.assertEqual("E_EVIDENCE_ITEM_BINDING", failure.exception.code)
|
self.assertFalse((self.archive / "manifest.jsonl").exists())
|
|
def test_observation_item_content_binding_mismatch_fails_closed(self) -> None:
|
begin = self.begin()
|
card, item = self.bound_card_item()
|
item["content_type"] = "article"
|
evidence = self.evidence(
|
begin, outcome="READABLE", observations=[self.observation([card])],
|
marker=self.end_marker(), items=[item],
|
)
|
path = Path(str(begin["evidence_path"]))
|
self.write_json(path, evidence)
|
with self.assertRaises(collector.CollectorError) as failure:
|
collector.run(["--config", str(self.config_path), "refresh-commit", "--input", str(path), "--now", "2026-08-13T10:00:25+08:00"])
|
self.assertEqual("E_EVIDENCE_ITEM_BINDING", failure.exception.code)
|
|
def test_observation_overlapping_component_identity_fails_closed(self) -> None:
|
begin = self.begin()
|
card, item = self.bound_card_item()
|
other_card, other_item = self.bound_card_item("457")
|
other_card["identifiers"]["dynamic_id"] = "456"
|
other_card["stable_keys"] = [
|
"dynamic:456", "opus:457", "url:https://www.bilibili.com/opus/457",
|
]
|
other_item["dynamic_id"] = "456"
|
for position, row in enumerate((card, other_card)):
|
row["position"] = position
|
evidence = self.evidence(
|
begin, outcome="READABLE", observations=[self.observation([card, other_card])],
|
marker=self.end_marker(), items=[item, other_item],
|
)
|
path = Path(str(begin["evidence_path"]))
|
self.write_json(path, evidence)
|
with self.assertRaises(collector.CollectorError) as failure:
|
collector.run(["--config", str(self.config_path), "refresh-commit", "--input", str(path), "--now", "2026-08-13T10:00:25+08:00"])
|
self.assertEqual("E_EVIDENCE_ITEM_BINDING", failure.exception.code)
|
|
def test_new_text_commits_both_manifests_and_artifact(self) -> None:
|
begin = self.begin()
|
card = {
|
"position": 0,
|
"identifiers": {"dynamic_id": "456", "opus_id": "456", "bvid": None},
|
"stable_keys": ["dynamic:456", "opus:456", "url:https://www.bilibili.com/opus/456"],
|
"published_at": "2026-08-13T01:55:00Z",
|
"content_type": "text",
|
"source_url": "https://www.bilibili.com/opus/456",
|
}
|
item = {
|
"dynamic_id": "456", "opus_id": "456", "bvid": None, "content_type": "text",
|
"published_at": "2026-08-13T01:55:00Z", "title": "离线合成正文",
|
"source_url": "https://www.bilibili.com/opus/456", "body_text": "仅测试文本",
|
"body_complete": True, "duration_seconds": None, "artifacts": [],
|
}
|
observation = self.observation([card])
|
marker = {
|
"observation_ordinal": 0,
|
"kind": "EXACT_END_OF_FEED",
|
"selector_id": "DYNAMIC_FEED_END_TEXT",
|
"normalized_text": "没有更多动态了",
|
"marker_sha256": hashlib.sha256("DYNAMIC_FEED_END_TEXT\n没有更多动态了".encode()).hexdigest(),
|
}
|
code, result = self.commit(begin, self.evidence(begin, outcome="READABLE", observations=[observation], marker=marker, items=[item]))
|
self.assertEqual(0, code)
|
self.assertEqual("NEW_ITEMS_SAVED", result["status"])
|
self.assertEqual(1, result["new_items"])
|
self.assertEqual(1, result["saved_artifacts"])
|
self.assertEqual(1, len(collector.load_manifest(self.state / "manifest.jsonl")))
|
formal = (self.archive / "manifest.jsonl").read_text(encoding="utf-8")
|
self.assertIn('"event_type":"DYNAMIC_CONTENT_SAVED"', formal)
|
artifacts = [path for path in self.archive.glob("*.txt")]
|
self.assertEqual(1, len(artifacts))
|
self.assertEqual("仅测试文本\n", artifacts[0].read_text(encoding="utf-8"))
|
|
slot = json.loads(Path(str(result["run_evidence_path"])).read_text(encoding="utf-8"))
|
self.assertEqual("DUAL_MANIFEST_COMMIT", slot["transaction_receipt"]["business_commit_kind"])
|
self.assertEqual("BUSINESS_COMMITTED", slot["transaction_receipt"]["phase"])
|
|
def test_hourly_slot_and_dead_pending_recovery_are_bounded(self) -> None:
|
begin = self.begin()
|
with self.assertRaises(collector.CollectorError) as busy:
|
collector.run(["--config", str(self.config_path), "refresh-begin", "--now", "2026-08-13T10:00:10+08:00"])
|
self.assertEqual("E_BUSY", busy.exception.code)
|
code, terminal = collector.run(["--config", str(self.config_path), "refresh-begin", "--now", "2026-08-13T10:02:01+08:00"])
|
self.assertEqual(4, code)
|
self.assertEqual("E_EVIDENCE_MISSING_AFTER_DEADLINE", terminal["error_code"])
|
with self.assertRaises(collector.CollectorError) as occupied:
|
collector.run(["--config", str(self.config_path), "refresh-begin", "--now", "2026-08-13T10:02:10+08:00"])
|
self.assertEqual("E_RUN_HOUR_OCCUPIED", occupied.exception.code)
|
|
def test_latest_commit_failure_retains_pending_and_reopen_rebuilds_index(self) -> None:
|
begin = self.begin()
|
observation = self.observation([])
|
evidence = self.evidence(
|
begin, outcome="READABLE", observations=[observation], marker=self.end_marker(), items=[]
|
)
|
path = Path(str(begin["evidence_path"]))
|
self.write_json(path, evidence)
|
config = collector.load_config(self.config_path)
|
pending = refresh._load_pending(config)
|
assert pending is not None
|
refresh.bind_controller_evidence(
|
config, pending, path,
|
transitioned_at=collector.parse_datetime("2026-08-13T10:00:24+08:00", "test transitioned_at"),
|
)
|
original = refresh._write_readback
|
injected = {"done": False}
|
|
def fail_latest(target: Path, payload: bytes, description: str) -> None:
|
if description == "latest run index" and not injected["done"]:
|
injected["done"] = True
|
raise OSError("injected latest failure")
|
original(target, payload, description)
|
|
refresh._write_readback = fail_latest
|
try:
|
with self.assertRaisesRegex(OSError, "injected latest failure"):
|
collector.run(["--config", str(self.config_path), "refresh-commit", "--input", str(path), "--now", "2026-08-13T10:00:25+08:00"])
|
finally:
|
refresh._write_readback = original
|
self.assertTrue(refresh._pending_path(config).is_file())
|
self.assertFalse((self.state / "refresh" / "runs" / "latest.json").exists())
|
replay = refresh._recover_or_replay(
|
config, self.config_path, collector.parse_now("2026-08-13T10:00:26+08:00"),
|
allow_final_evidence=True,
|
)
|
self.assertEqual("REFRESH_CONFIRMED_NO_NEW", replay["status"])
|
self.assertTrue((self.state / "refresh" / "runs" / "latest.json").is_file())
|
self.assertFalse(refresh._pending_path(config).exists())
|
|
def test_state_lock_is_kernel_backed_and_owner_metadata_recovers(self) -> None:
|
config = collector.load_config(self.config_path)
|
config.state_dir.mkdir()
|
stale = {
|
"pid": 99999999,
|
"process_created_at": "0",
|
"run_id": None,
|
"acquired_at": "2026-08-13T00:00:00Z",
|
}
|
self.write_json(config.state_dir / ".collector.lock.owner.json", stale)
|
with collector.StateLock(config.lock_path):
|
self.assertEqual(1, config.lock_path.stat().st_size)
|
self.assertTrue((config.state_dir / ".collector.lock.owner.json").is_file())
|
with self.assertRaises(collector.CollectorError) as locked:
|
with collector.StateLock(config.lock_path):
|
pass
|
self.assertEqual("E_STATE_LOCKED", locked.exception.code)
|
self.assertTrue(config.lock_path.is_file())
|
self.assertFalse((config.state_dir / ".collector.lock.owner.json").exists())
|
|
def test_frozen_formal_catalog_shape_is_48_22_16_5_1_read_only(self) -> None:
|
formal = Path(__file__).resolve().parents[3] / "ana-data" / "news-青枫浦上Q" / "manifest.jsonl"
|
if not formal.exists():
|
self.skipTest("formal read-only acceptance snapshot is not present")
|
before = (formal.stat().st_size, hashlib.sha256(formal.read_bytes()).hexdigest())
|
config = collector.load_config(self.config_path)
|
archive = formal.parent
|
rebound = collector.CollectorConfig(
|
config.creator_name, config.creator_uid, config.creator_dynamic_url, config.timezone_name,
|
config.window_hours, config.state_dir, config.download_dir, config.video_dir,
|
config.minimum_complete_age_seconds, config.title_max_length, config.allowed_source_hosts,
|
config.allowed_video_extensions, config.native_handoff,
|
collector.RefreshConfig(
|
archive, formal, config.refresh.intake_dir, config.refresh.overall_deadline_seconds,
|
config.refresh.refresh_action_timeout_seconds, config.refresh.observation_timeout_seconds,
|
config.refresh.page_internal_settle_timeout_seconds,
|
1, 168, 200, 20, config.refresh.max_image_bytes, config.refresh.max_text_bytes,
|
),
|
)
|
_, _, counts = refresh.load_formal_catalog(rebound)
|
self.assertEqual({"events": 48, "components": 22, "saved": 16, "video": 5, "retryable": 1}, counts)
|
self.assertEqual(before, (formal.stat().st_size, hashlib.sha256(formal.read_bytes()).hexdigest()))
|
|
def test_mixed_history_union_saved_precedence_and_schema2_identity(self) -> None:
|
artifact = self.archive / "saved.txt"
|
artifact.write_bytes(b"saved\n")
|
digest = hashlib.sha256(artifact.read_bytes()).hexdigest()
|
config = collector.load_config(self.config_path)
|
pending = {
|
"schema_version": 1, "creator": config.creator_name, "creator_uid": config.creator_uid,
|
"item_type": "article", "stable_id": "700", "source_url": "https://www.bilibili.com/opus/700",
|
"published_at": "2026-08-12T01:00:00Z", "collected_at": "2026-08-12T02:00:00Z",
|
"status": "CONTENT_ACCESS_PENDING_CHROME",
|
}
|
saved = {
|
**pending, "status": "SAVED", "path": "saved.txt", "bytes": artifact.stat().st_size,
|
"sha256": digest.upper(),
|
}
|
formal = self.archive / "manifest.jsonl"
|
formal.write_bytes(collector.canonical_json_bytes(pending) + collector.canonical_json_bytes(saved))
|
_, tokens, counts = refresh.load_formal_catalog(config)
|
self.assertEqual({"events": 2, "components": 1, "saved": 1, "video": 0, "retryable": 0}, counts)
|
self.assertEqual({"opus:700", "url:https://www.bilibili.com/opus/700"}, tokens)
|
|
def test_schema2_status_entity_and_artifact_identity_fail_closed(self) -> None:
|
artifact = self.archive / "new.txt"
|
artifact.write_bytes(b"new\n")
|
source = "https://www.bilibili.com/opus/701"
|
keys = ["dynamic:701", "opus:701", f"url:{source}"]
|
config = collector.load_config(self.config_path)
|
event = {
|
"schema_version": 2, "event_type": "DYNAMIC_CONTENT_SAVED", "creator": config.creator_name,
|
"creator_uid": config.creator_uid, "entity_id": collector.entity_id_for_keys(keys),
|
"dynamic_id": "701", "opus_id": "701", "bvid": None, "dedupe_keys": keys,
|
"content_type": "text", "published_at": "2026-08-12T01:00:00Z", "title": "saved",
|
"source_url": source, "collected_at": "2026-08-12T02:00:00Z",
|
"artifacts": [{"kind": "text", "sequence": 1, "path": "new.txt", "bytes": 4,
|
"sha256": hashlib.sha256(b"new\n").hexdigest()}],
|
"duration_seconds": None, "page_run_id": "a" * 32, "coverage_complete": True, "status": "SAVED",
|
}
|
mutations = (
|
("status", {**event, "status": "UNKNOWN"}, "E_CATALOG_STATUS"),
|
("integer_uid", {**event, "creator_uid": int(config.creator_uid)}, "E_CATALOG_STATUS"),
|
("entity", {**event, "entity_id": "0" * 24}, "E_CATALOG_IDENTITY_CONFLICT"),
|
("artifact", {**event, "artifacts": [{**event["artifacts"][0], "path": "absent.txt"}]}, "E_CATALOG_ARTIFACT"),
|
)
|
for label, bad, expected in mutations:
|
with self.subTest(label=label):
|
(self.archive / "manifest.jsonl").write_bytes(collector.canonical_json_bytes(bad))
|
with self.assertRaises(collector.CollectorError) as failure:
|
refresh.load_formal_catalog(config)
|
self.assertEqual(expected, failure.exception.code)
|
|
def test_schema1_stable_id_must_be_one_to_thirty_two_ascii_digits(self) -> None:
|
config = collector.load_config(self.config_path)
|
base = {
|
"schema_version": 1, "creator": config.creator_name, "creator_uid": config.creator_uid,
|
"item_type": "article", "source_url": "https://www.bilibili.com/opus/123",
|
"published_at": "2026-08-12T01:00:00Z", "collected_at": "2026-08-12T02:00:00Z",
|
"status": "CONTENT_ACCESS_PENDING_CHROME",
|
}
|
for stable_id in ("1" * 33, "123"):
|
with self.subTest(stable_id=stable_id):
|
row = {**base, "stable_id": stable_id}
|
(self.archive / "manifest.jsonl").write_bytes(collector.canonical_json_bytes(row))
|
with self.assertRaises(collector.CollectorError) as failure:
|
refresh.load_formal_catalog(config)
|
self.assertEqual("E_CATALOG_IDENTITY_CONFLICT", failure.exception.code)
|
|
def test_pending_before_started_rebuilds_same_run_and_third_slot_is_untouched(self) -> None:
|
config = collector.load_config(self.config_path)
|
original = refresh._ensure_started_slot
|
refresh._ensure_started_slot = lambda *_args, **_kwargs: (_ for _ in ()).throw(OSError("after pending"))
|
try:
|
with self.assertRaisesRegex(OSError, "after pending"):
|
refresh.refresh_begin(config, self.config_path, collector.parse_now(NOW))
|
finally:
|
refresh._ensure_started_slot = original
|
pending = refresh._load_pending(config)
|
self.assertIsNotNone(pending)
|
run_id = pending["run_id"]
|
result = refresh.refresh_begin(config, self.config_path, collector.parse_now("2026-08-13T10:00:01+08:00"))
|
self.assertEqual(run_id, result["run_id"])
|
self.assertEqual("RUN_STARTED", json.loads(Path(str(result["run_evidence_path"])).read_text(encoding="utf-8"))["run_state"])
|
|
Path(str(result["run_evidence_path"])).write_bytes(b'{"third":true}')
|
before = Path(str(result["run_evidence_path"])).read_bytes()
|
with self.assertRaises(collector.CollectorError) as failure:
|
refresh._recover_or_replay(config, self.config_path, collector.parse_now("2026-08-13T10:00:02+08:00"))
|
self.assertEqual("E_RECOVERY_AMBIGUOUS", failure.exception.code)
|
self.assertEqual(before, Path(str(result["run_evidence_path"])).read_bytes())
|
|
def test_terminal_never_overwrites_third_slot_and_recorded_replay_binds_plan(self) -> None:
|
begin = self.begin()
|
slot_path = Path(str(begin["run_evidence_path"]))
|
third = b'{"third":"terminal"}'
|
evidence = self.evidence(begin, outcome="READABLE", observations=[self.observation([])], marker=self.end_marker())
|
original_write_pending = refresh._write_pending
|
injected = {"done": False}
|
|
def inject_third(config, pending, *, create=False):
|
original_write_pending(config, pending, create=create)
|
if pending.get("planned_terminal") is not None and not injected["done"]:
|
injected["done"] = True
|
slot_path.write_bytes(third)
|
|
refresh._write_pending = inject_third
|
try:
|
with self.assertRaises(collector.CollectorError) as failure:
|
self.commit(begin, evidence)
|
finally:
|
refresh._write_pending = original_write_pending
|
self.assertEqual("E_RUN_EVIDENCE_COMMIT", failure.exception.code)
|
self.assertEqual(third, slot_path.read_bytes())
|
|
config = collector.load_config(self.config_path)
|
pending = refresh._load_pending(config)
|
self.assertIsNotNone(pending)
|
self.assertIsNotNone(pending["planned_terminal"])
|
slot, payload = refresh._terminal_slot_from_planned(config, pending, pending["planned_terminal"])
|
slot_path.write_bytes(payload)
|
pending["phase"] = "TERMINAL_RECORDED"
|
pending["planned_terminal"]["transaction_receipt"]["intent_sha256"] = "0" * 64
|
refresh._write_pending(config, pending)
|
before = slot_path.read_bytes()
|
with self.assertRaises(collector.CollectorError) as replay_failure:
|
refresh._recover_or_replay(config, self.config_path, collector.parse_now("2026-08-13T10:00:26+08:00"))
|
self.assertEqual("E_RUN_EVIDENCE_COMMIT", replay_failure.exception.code)
|
self.assertEqual(before, slot_path.read_bytes())
|
|
def test_mixed_history_cross_line_identity_conflict_fails_closed(self) -> None:
|
config = collector.load_config(self.config_path)
|
first = {
|
"schema_version": 1, "creator": config.creator_name, "creator_uid": config.creator_uid,
|
"item_type": "video", "stable_id": "BV1HA3o6oEJJ",
|
"source_url": "https://www.bilibili.com/video/BV1HA3o6oEJJ",
|
"published_at": "2026-08-12T01:00:00Z", "collected_at": "2026-08-12T02:00:00Z",
|
"status": "VIDEO_DOWNLOAD_PENDING_EXTENSION",
|
}
|
second = {**first, "stable_id": "BV1DVMX6XEPq"}
|
(self.archive / "manifest.jsonl").write_bytes(
|
collector.canonical_json_bytes(first) + collector.canonical_json_bytes(second)
|
)
|
with self.assertRaises(collector.CollectorError) as failure:
|
refresh.load_formal_catalog(config)
|
self.assertEqual("E_CATALOG_IDENTITY_CONFLICT", failure.exception.code)
|
|
def test_dead_formal_lock_takeover_uses_archive_quarantine_and_current_identity(self) -> None:
|
config = collector.load_config(self.config_path)
|
config.state_dir.mkdir()
|
pending = {
|
"schema_version": 2,
|
"run_id": "a" * 32,
|
"owner_nonce": "b" * 32,
|
"phase": "TRANSACTION_INTENT",
|
"task_id": refresh.TASK_ID,
|
"creator_uid": config.creator_uid,
|
"creator_dynamic_url": config.creator_dynamic_url,
|
"started_at": "2026-08-13T02:00:00Z",
|
"deadline_at": "2026-08-13T02:00:30Z",
|
"window_start": "2026-08-10T02:00:00Z",
|
"window_end": "2026-08-13T02:00:00Z",
|
"config_sha256": hashlib.sha256(self.config_path.read_bytes()).hexdigest(),
|
"state_manifest_preimage": refresh._identity(config.manifest_path).as_dict(),
|
"formal_manifest_preimage": refresh._identity(config.refresh.formal_manifest).as_dict(),
|
"evidence_path": str(config.state_dir / "refresh" / "incoming" / ("a" * 32 + ".json")),
|
"intake_root": str(config.refresh.intake_dir / ("a" * 32)),
|
"evidence_identity": {"bytes": 2, "sha256": hashlib.sha256(b"{}").hexdigest()},
|
"planned_terminal": None,
|
"transaction_identity": {
|
"transaction_id": "c" * 32,
|
"state_preimage": refresh._identity(config.manifest_path).as_dict(),
|
"formal_preimage": refresh._identity(config.refresh.formal_manifest).as_dict(),
|
"state_candidate": refresh._identity(config.manifest_path).as_dict(),
|
"formal_candidate": refresh._identity(config.refresh.formal_manifest).as_dict(),
|
"created_artifacts": [],
|
"input_item_count": 0,
|
"new_item_count": 0,
|
"formal_lock_claim": None,
|
},
|
"last_transition_at": "2026-08-13T02:00:00Z",
|
}
|
old_record = {
|
"schema_version": 1,
|
"task_id": refresh.TASK_ID,
|
"run_id": pending["run_id"],
|
"owner_nonce": pending["owner_nonce"],
|
"transaction_id": pending["transaction_identity"]["transaction_id"],
|
"recovery_generation": 0,
|
"holder_pid": 99999999,
|
"holder_process_created_at": "0",
|
"lock_created_at": "2026-08-13T02:00:00Z",
|
}
|
old_payload = collector.canonical_json_bytes(old_record, newline=False)
|
old_claim = {
|
"claim_state": "HELD",
|
"lock_path": str(refresh._formal_lock_path(config)),
|
"lock_bytes": len(old_payload),
|
"lock_sha256": hashlib.sha256(old_payload).hexdigest(),
|
"lock_record": old_record,
|
}
|
pending["transaction_identity"]["formal_lock_claim"] = old_claim
|
refresh._write_pending(config, pending, create=True)
|
refresh._create_new(refresh._formal_lock_path(config), old_payload)
|
refresh._recovery_formal_lock(config, pending, collector.parse_now("2026-08-13T10:00:25+08:00"))
|
claim = pending["transaction_identity"]["formal_lock_claim"]
|
self.assertEqual("HELD", claim["claim_state"])
|
self.assertEqual(1, claim["lock_record"]["recovery_generation"])
|
self.assertEqual(os.getpid(), claim["lock_record"]["holder_pid"])
|
takeover = pending["transaction_identity"]["takeover"]
|
self.assertEqual("NEW_LOCK_HELD", takeover["phase"])
|
quarantine = Path(takeover["quarantine_path"])
|
self.assertEqual(config.refresh.archive_dir, quarantine.parent.parent)
|
self.assertEqual(old_payload, quarantine.read_bytes())
|
self.assertEqual(collector.canonical_json_bytes(claim["lock_record"], newline=False), refresh._formal_lock_path(config).read_bytes())
|
lock_path = refresh._formal_lock_path(config)
|
lock_path.unlink()
|
claim["claim_state"] = "PLANNED"
|
claim["lock_record"]["holder_pid"] = 99999999
|
claim["lock_record"]["holder_process_created_at"] = "0"
|
dead_payload = collector.canonical_json_bytes(claim["lock_record"], newline=False)
|
claim["lock_bytes"] = len(dead_payload)
|
claim["lock_sha256"] = hashlib.sha256(dead_payload).hexdigest()
|
takeover["phase"] = "CREATE_PLANNED"
|
prior_generation = takeover["next_generation"]
|
prior_attempt = takeover["claim_attempt"]
|
refresh._write_pending(config, pending)
|
refresh._recovery_formal_lock(config, pending, collector.parse_now("2026-08-13T10:00:26+08:00"))
|
claim = pending["transaction_identity"]["formal_lock_claim"]
|
takeover = pending["transaction_identity"]["takeover"]
|
self.assertEqual(prior_generation, claim["lock_record"]["recovery_generation"])
|
self.assertEqual(prior_attempt + 1, takeover["claim_attempt"])
|
self.assertEqual(os.getpid(), claim["lock_record"]["holder_pid"])
|
|
lock_path.unlink()
|
claim["claim_state"] = "PLANNED"
|
claim["lock_record"]["holder_pid"] = 99999999
|
claim["lock_record"]["holder_process_created_at"] = "0"
|
dead_payload = collector.canonical_json_bytes(claim["lock_record"], newline=False)
|
claim["lock_bytes"] = len(dead_payload)
|
claim["lock_sha256"] = hashlib.sha256(dead_payload).hexdigest()
|
takeover["phase"] = "CREATE_PLANNED"
|
refresh._write_pending(config, pending)
|
refresh._create_new(lock_path, dead_payload)
|
original_rename = refresh._rename_no_overwrite
|
rename_calls: list[tuple[Path, Path]] = []
|
|
def fail_after_intent(source: Path, target: Path) -> None:
|
rename_calls.append((source, target))
|
raise OSError("crash before failed-claim rename")
|
|
refresh._rename_no_overwrite = fail_after_intent
|
try:
|
with self.assertRaisesRegex(OSError, "crash before failed-claim rename"):
|
refresh._recovery_formal_lock(config, pending, collector.parse_now("2026-08-13T10:00:27+08:00"))
|
finally:
|
refresh._rename_no_overwrite = original_rename
|
persisted = refresh._load_pending(config)
|
intent = persisted["transaction_identity"]["takeover"]["failed_claim_intent"]
|
self.assertEqual(str(rename_calls[0][1]), intent["path"])
|
self.assertEqual(dead_payload, lock_path.read_bytes())
|
self.assertFalse(Path(intent["path"]).exists())
|
|
refresh._recovery_formal_lock(config, persisted, collector.parse_now("2026-08-13T10:00:27+08:00"))
|
pending = persisted
|
takeover = pending["transaction_identity"]["takeover"]
|
self.assertEqual(1, len(takeover["failed_claims"]))
|
failed = takeover["failed_claims"][0]
|
self.assertEqual(dead_payload, Path(failed["path"]).read_bytes())
|
self.assertTrue((refresh._quarantine_root(config) / ".owner.json").is_file())
|
refresh._formal_lock_path(config).unlink()
|
refresh._release_formal_lock(config, pending)
|
self.assertFalse(refresh._formal_lock_path(config).exists())
|
self.assertFalse(quarantine.exists())
|
self.assertFalse(refresh._quarantine_root(config).exists())
|
|
def test_quarantine_rename_is_no_overwrite_and_empty_cleanup_reopens(self) -> None:
|
config = collector.load_config(self.config_path)
|
source = refresh._formal_lock_path(config)
|
root = refresh._quarantine_root(config)
|
root.mkdir()
|
target = root / "third.json"
|
source.write_bytes(b"owned")
|
target.write_bytes(b"third")
|
with self.assertRaises(collector.CollectorError) as collision:
|
refresh._rename_no_overwrite(source, target)
|
self.assertEqual("E_RECOVERY_AMBIGUOUS", collision.exception.code)
|
self.assertEqual(b"owned", source.read_bytes())
|
self.assertEqual(b"third", target.read_bytes())
|
|
target.unlink()
|
flushed: list[Path] = []
|
original_fsync = refresh._fsync_directory
|
refresh._fsync_directory = lambda path: flushed.append(Path(path))
|
try:
|
refresh._rename_no_overwrite(source, target)
|
finally:
|
refresh._fsync_directory = original_fsync
|
self.assertFalse(source.exists())
|
self.assertEqual(b"owned", target.read_bytes())
|
self.assertIn(self.archive, flushed)
|
self.assertIn(root, flushed)
|
|
target.unlink()
|
pending = {
|
"run_id": "a" * 32,
|
"owner_nonce": "b" * 32,
|
"transaction_identity": {
|
"formal_lock_claim": {
|
"lock_bytes": 0, "lock_sha256": hashlib.sha256(b"").hexdigest()
|
},
|
"takeover": {"old_lock_never_created": True, "failed_claims": []},
|
},
|
}
|
refresh._release_formal_lock(config, pending)
|
self.assertFalse(root.exists())
|
|
def test_formal_claim_requires_exact_readback_and_parent_durability_before_held(self) -> None:
|
config = collector.load_config(self.config_path)
|
config.state_dir.mkdir()
|
now = collector.parse_now("2026-08-13T10:00:25+08:00")
|
formal_before = refresh._identity(config.refresh.formal_manifest).as_dict()
|
|
def prepare() -> tuple[dict[str, object], dict[str, object], Path]:
|
pending = self.formal_pending(config)
|
claim = refresh._new_formal_claim(config, pending, 1, now)
|
pending["transaction_identity"]["formal_lock_claim"] = claim
|
pending["transaction_identity"]["takeover"] = {
|
"phase": "CREATE_PLANNED",
|
"old_lock_sha256": hashlib.sha256(b"").hexdigest(),
|
"old_lock_hex": "",
|
"quarantine_path": str(refresh._quarantine_root(config) / "never-created.json"),
|
"next_generation": 1,
|
"takeover_reason": "TEST",
|
"planned_at": collector.canonical_datetime(now),
|
"claim_attempt": 1,
|
"old_lock_never_created": True,
|
}
|
refresh._write_pending(config, pending, create=True)
|
return pending, claim, refresh._formal_lock_path(config)
|
|
pending, _, lock = prepare()
|
original_create = refresh._create_new
|
refresh._create_new = lambda path, payload: path.write_bytes(b"third-party-drift")
|
try:
|
with self.assertRaises(collector.CollectorError) as drift:
|
refresh._recovery_formal_lock(config, pending, now)
|
finally:
|
refresh._create_new = original_create
|
self.assertEqual("E_FORMAL_LOCK", drift.exception.code)
|
self.assertEqual(b"third-party-drift", lock.read_bytes())
|
persisted = refresh._load_pending(config)
|
self.assertEqual("CREATE_PLANNED", persisted["transaction_identity"]["takeover"]["phase"])
|
self.assertEqual("PLANNED", persisted["transaction_identity"]["formal_lock_claim"]["claim_state"])
|
with self.assertRaises(collector.CollectorError) as replay_drift:
|
refresh._recovery_formal_lock(config, persisted, now)
|
self.assertEqual("E_RECOVERY_AMBIGUOUS", replay_drift.exception.code)
|
self.assertEqual(b"third-party-drift", lock.read_bytes())
|
lock.unlink()
|
refresh._pending_path(config).unlink()
|
|
pending, _, lock = prepare()
|
original_read = Path.read_bytes
|
|
def fail_lock_read(path: Path) -> bytes:
|
if path == lock:
|
raise OSError("injected formal lock readback failure")
|
return original_read(path)
|
|
Path.read_bytes = fail_lock_read
|
try:
|
with self.assertRaisesRegex(OSError, "readback failure"):
|
refresh._recovery_formal_lock(config, pending, now)
|
finally:
|
Path.read_bytes = original_read
|
self.assertEqual("CREATE_PLANNED", refresh._load_pending(config)["transaction_identity"]["takeover"]["phase"])
|
lock.unlink()
|
refresh._pending_path(config).unlink()
|
|
pending, claim, lock = prepare()
|
original_fsync = refresh._fsync_directory
|
|
def fail_parent_fsync(path: Path) -> None:
|
if Path(path) == self.archive:
|
raise OSError("injected parent directory fsync failure")
|
original_fsync(path)
|
|
refresh._fsync_directory = fail_parent_fsync
|
try:
|
with self.assertRaisesRegex(OSError, "parent directory fsync failure"):
|
refresh._recovery_formal_lock(config, pending, now)
|
finally:
|
refresh._fsync_directory = original_fsync
|
persisted = refresh._load_pending(config)
|
self.assertEqual("CREATE_PLANNED", persisted["transaction_identity"]["takeover"]["phase"])
|
self.assertEqual("PLANNED", persisted["transaction_identity"]["formal_lock_claim"]["claim_state"])
|
self.assertEqual(refresh._claim_payload(claim), lock.read_bytes())
|
refresh._recovery_formal_lock(config, persisted, now)
|
replayed = refresh._load_pending(config)
|
self.assertEqual("NEW_LOCK_HELD", replayed["transaction_identity"]["takeover"]["phase"])
|
self.assertEqual("HELD", replayed["transaction_identity"]["formal_lock_claim"]["claim_state"])
|
self.assertEqual(formal_before, refresh._identity(config.refresh.formal_manifest).as_dict())
|
refresh._release_formal_lock(config, replayed)
|
|
def test_formal_claim_real_child_kill_reopens_create_and_new_lock_held(self) -> None:
|
config = collector.load_config(self.config_path)
|
config.state_dir.mkdir()
|
formal_before = refresh._identity(config.refresh.formal_manifest).as_dict()
|
script = r'''import hashlib, sys, time
|
from pathlib import Path
|
sys.path.insert(0, sys.argv[1])
|
import bili_dynamic_collector as core
|
import bili_dynamic_refresh as refresh
|
config = core.load_config(Path(sys.argv[2]))
|
marker = Path(sys.argv[3])
|
mode = sys.argv[4]
|
pending = refresh._load_pending(config)
|
now = core.parse_now("2026-08-13T10:00:25+08:00")
|
claim = refresh._new_formal_claim(config, pending, 1, now)
|
pending["transaction_identity"]["formal_lock_claim"] = claim
|
pending["transaction_identity"]["takeover"] = {
|
"phase": "CREATE_PLANNED", "old_lock_sha256": hashlib.sha256(b"").hexdigest(),
|
"old_lock_hex": "", "quarantine_path": str(refresh._quarantine_root(config) / "never-created.json"),
|
"next_generation": 1, "takeover_reason": "CHILD_KILL_TEST",
|
"planned_at": core.canonical_datetime(now), "claim_attempt": 1, "old_lock_never_created": True,
|
}
|
refresh._write_pending(config, pending)
|
if mode == "CREATE_PLANNED":
|
original = refresh._fsync_directory
|
def stop_at_parent(path):
|
if Path(path) == config.refresh.archive_dir and refresh._formal_lock_path(config).exists():
|
marker.write_text("created", encoding="utf-8")
|
time.sleep(120)
|
original(path)
|
refresh._fsync_directory = stop_at_parent
|
refresh._recovery_formal_lock(config, pending, now)
|
marker.write_text("held", encoding="utf-8")
|
time.sleep(120)
|
'''
|
for mode, marker_text in (("CREATE_PLANNED", "created"), ("NEW_LOCK_HELD", "held")):
|
with self.subTest(mode=mode):
|
pending = self.formal_pending(config)
|
refresh._write_pending(config, pending, create=True)
|
marker = self.root / f"{mode}.marker"
|
process = subprocess.Popen(
|
[sys.executable, "-B", "-c", script, str(PROJECT_DEV), str(self.config_path), str(marker), mode],
|
stdout=subprocess.DEVNULL,
|
stderr=subprocess.DEVNULL,
|
)
|
deadline = time.monotonic() + 15
|
while time.monotonic() < deadline and not marker.exists() and process.poll() is None:
|
time.sleep(0.05)
|
self.assertTrue(marker.exists(), f"child failed before {mode} marker; exit={process.poll()}")
|
self.assertEqual(marker_text, marker.read_text(encoding="utf-8"))
|
process.kill()
|
process.wait(timeout=10)
|
child_pid = process.pid
|
if os.name == "nt":
|
process._handle.Close()
|
del process
|
gc.collect()
|
dead_deadline = time.monotonic() + 10
|
while time.monotonic() < dead_deadline:
|
try:
|
collector.process_created_at(child_pid)
|
except ProcessLookupError:
|
break
|
time.sleep(0.05)
|
else:
|
self.fail(f"child process identity remained queryable after kill: {child_pid}")
|
persisted = refresh._load_pending(config)
|
expected_phase = "CREATE_PLANNED" if mode == "CREATE_PLANNED" else "NEW_LOCK_HELD"
|
self.assertEqual(expected_phase, persisted["transaction_identity"]["takeover"]["phase"])
|
self.assertEqual(formal_before, refresh._identity(config.refresh.formal_manifest).as_dict())
|
refresh._recovery_formal_lock(config, persisted, collector.parse_now("2026-08-13T10:00:26+08:00"))
|
recovered = refresh._load_pending(config)
|
self.assertEqual("NEW_LOCK_HELD", recovered["transaction_identity"]["takeover"]["phase"])
|
self.assertEqual("HELD", recovered["transaction_identity"]["formal_lock_claim"]["claim_state"])
|
self.assertEqual(formal_before, refresh._identity(config.refresh.formal_manifest).as_dict())
|
refresh._release_formal_lock(config, recovered)
|
refresh._pending_path(config).unlink()
|
|
def test_runtime_v2_config_and_begin_bind_exact_contract(self) -> None:
|
begin = self.begin()
|
self.assertEqual("bili-supported-chrome-visible-runtime-v2", begin["runtime_contract_id"])
|
self.assertEqual(120, begin["overall_deadline_seconds"])
|
self.assertEqual(35, begin["refresh_action_timeout_seconds"])
|
self.assertEqual(45, begin["observation_timeout_seconds"])
|
self.assertEqual(15, begin["page_internal_settle_timeout_seconds"])
|
pending = refresh._load_pending(collector.load_config(self.config_path))
|
self.assertEqual(3, pending["schema_version"])
|
self.assertEqual(begin["runtime_contract_sha256"], pending["runtime_contract"]["contract_sha256"])
|
|
config = json.loads(self.config_path.read_text(encoding="utf-8"))
|
config["refresh"]["timeout_seconds"] = 30
|
legacy = self.root / "legacy.json"
|
self.write_json(legacy, config)
|
with self.assertRaises(collector.CollectorError) as invalid:
|
collector.load_config(legacy)
|
self.assertEqual("E_CONFIG", invalid.exception.code)
|
|
def test_runtime_pre_dispatch_error_is_zero_action_unreadable(self) -> None:
|
begin = self.begin()
|
evidence = self.evidence(
|
begin, outcome="UNREADABLE_ERROR", observations=None, marker=None,
|
action_outcome="PRE_DISPATCH_ERROR", observation_outcome="NOT_ATTEMPTED", refresh_count=0,
|
)
|
self.set_runtime_diagnostic(evidence, "ACTION_PRE_DISPATCH")
|
before = refresh._identity(self.archive / "manifest.jsonl").as_dict()
|
code, result = self.commit(begin, evidence)
|
self.assertEqual(4, code)
|
self.assertEqual("REFRESH_FAILED_PAGE_UNREADABLE", result["status"])
|
self.assertEqual(0, result["refresh_count"])
|
self.assertEqual(before, refresh._identity(self.archive / "manifest.jsonl").as_dict())
|
|
def test_runtime_pre_dispatch_claimed_observation_fails_closed(self) -> None:
|
for observation_outcome in ("READABLE", "ERROR"):
|
with self.subTest(observation_outcome=observation_outcome):
|
begin = self.begin()
|
evidence = self.evidence(
|
begin, outcome="UNREADABLE_ERROR", observations=None, marker=None,
|
action_outcome="PRE_DISPATCH_ERROR", observation_outcome="NOT_ATTEMPTED", refresh_count=0,
|
)
|
self.set_runtime_diagnostic(evidence, "ACTION_PRE_DISPATCH")
|
evidence["runtime_observation"]["observation_outcome"] = observation_outcome
|
evidence["runtime_observation"]["observation_count"] = 1
|
path = Path(str(begin["evidence_path"]))
|
self.write_json(path, evidence)
|
with self.assertRaises(collector.CollectorError) as invalid:
|
collector.run(["--config", str(self.config_path), "refresh-commit", "--input", str(path), "--now", "2026-08-13T10:00:25+08:00"])
|
self.assertEqual("E_EVIDENCE_SCHEMA", invalid.exception.code)
|
refresh._pending_path(collector.load_config(self.config_path)).unlink()
|
Path(str(begin["run_evidence_path"])).unlink()
|
|
def test_runtime_post_dispatch_error_matrix(self) -> None:
|
begin = self.begin()
|
readable_empty = self.evidence(
|
begin, outcome="READABLE", observations=[self.observation([])], marker=self.end_marker(),
|
action_outcome="POST_DISPATCH_ERROR",
|
)
|
self.set_runtime_diagnostic(readable_empty, "ACTION_POST_DISPATCH")
|
code, result = self.commit(begin, readable_empty)
|
self.assertEqual(4, code)
|
self.assertEqual("PARTIAL_DISCOVERY_UNCONFIRMED", result["status"])
|
self.assertFalse(result["no_new_confirmed"])
|
|
def test_runtime_post_dispatch_read_error_and_deadline_exhausted(self) -> None:
|
cases = (
|
("ERROR", "UNREADABLE_ERROR", "ACTION_POST_DISPATCH", 1),
|
("DEADLINE_EXHAUSTED", "UNREADABLE_TIMEOUT", "ACTION_POST_DISPATCH", 0),
|
)
|
for observation_outcome, page_outcome, code_name, observation_count in cases:
|
with self.subTest(observation_outcome=observation_outcome):
|
begin = self.begin()
|
evidence = self.evidence(
|
begin, outcome=page_outcome, observations=None, marker=None,
|
action_outcome="POST_DISPATCH_ERROR", observation_outcome=observation_outcome,
|
observation_count=observation_count,
|
)
|
self.set_runtime_diagnostic(evidence, code_name)
|
if observation_count == 0:
|
evidence["runtime_observation"]["observation_elapsed_ms"] = 0
|
evidence["read_finished_at"] = evidence["refresh_finished_at"]
|
code, result = self.commit(begin, evidence)
|
self.assertEqual(4, code)
|
self.assertEqual("REFRESH_FAILED_PAGE_UNREADABLE", result["status"])
|
if observation_outcome == "ERROR":
|
Path(str(result["run_evidence_path"])).unlink()
|
(self.state / "refresh" / "runs" / "latest.json").unlink()
|
|
def test_runtime_timeout_cannot_confirm_no_new(self) -> None:
|
begin = self.begin()
|
evidence = self.evidence(
|
begin, outcome="READABLE", observations=[self.observation([])], marker=self.end_marker(),
|
action_outcome="TIMEOUT",
|
)
|
self.set_runtime_diagnostic(evidence, "ACTION_TIMEOUT")
|
code, result = self.commit(begin, evidence)
|
self.assertEqual(4, code)
|
self.assertEqual("PARTIAL_DISCOVERY_UNCONFIRMED", result["status"])
|
self.assertFalse(result["no_new_confirmed"])
|
|
def test_runtime_supported_control_latency_is_inside_new_budgets(self) -> None:
|
begin = self.begin()
|
evidence = self.evidence(
|
begin, outcome="READABLE", observations=[self.observation([])], marker=self.end_marker(),
|
)
|
evidence["refresh_finished_at"] = "2026-08-13T02:00:22.500000Z"
|
evidence["read_finished_at"] = "2026-08-13T02:00:44Z"
|
evidence["runtime_observation"]["refresh_action_elapsed_ms"] = 21500
|
evidence["runtime_observation"]["observation_elapsed_ms"] = 21500
|
evidence["controller_attestation"]["monotonic_action_finished_ms"] = 22500
|
evidence["controller_attestation"]["monotonic_observation_started_ms"] = 22500
|
evidence["controller_attestation"]["monotonic_observation_finished_ms"] = 44000
|
evidence["controller_attestation"]["monotonic_evidence_write_started_ms"] = 44000
|
code, result = self.commit(begin, evidence)
|
self.assertEqual(0, code)
|
self.assertEqual("REFRESH_CONFIRMED_NO_NEW", result["status"])
|
|
def test_runtime_post_dispatch_error_readable_new_saves_without_coverage(self) -> None:
|
begin = self.begin()
|
card, item = self.bound_card_item("987654")
|
evidence = self.evidence(
|
begin, outcome="READABLE", observations=[self.observation([card])], marker=self.end_marker(),
|
items=[item], action_outcome="POST_DISPATCH_ERROR",
|
)
|
self.set_runtime_diagnostic(evidence, "ACTION_POST_DISPATCH")
|
code, result = self.commit(begin, evidence)
|
self.assertEqual(0, code)
|
self.assertEqual("NEW_ITEMS_SAVED", result["status"])
|
self.assertFalse(result["coverage_complete"])
|
|
def test_reviewer_trusted_controller_calls_once_and_caller_forgery_fails(self) -> None:
|
begin = self.begin()
|
calls = {"tabs": 0, "reload": 0, "goto": 0, "evaluate": 0}
|
|
class Fake:
|
def open_tabs(inner) -> list[dict[str, object]]:
|
calls["tabs"] += 1
|
return [{"url": "https://space.bilibili.com/1420210197/dynamic"}]
|
|
def reload(inner, tab: dict[str, object], timeout_seconds: int) -> None:
|
calls["reload"] += 1
|
|
def goto(inner, url: str, timeout_seconds: int) -> dict[str, object]:
|
calls["goto"] += 1
|
return {"url": url}
|
|
def evaluate(inner, tab: dict[str, object], source: str, timeout_seconds: int) -> dict[str, object]:
|
calls["evaluate"] += 1
|
return {
|
"schema_version": 1,
|
"ready_state": "complete", "visibility_state": "visible",
|
"final_url": "https://space.bilibili.com/1420210197/dynamic",
|
"page_title": "青枫浦上Q个人动态-青枫浦上Q动态记录-哔哩哔哩视频",
|
"creator": {"uid": "1420210197", "name": "青枫浦上Q", "profile_url": "https://space.bilibili.com/1420210197"},
|
"cards": [], "unparsed_nodes": [], "terminal_marker_text": "已经到底了", "limit_hit": "NONE",
|
}
|
|
ticks = iter([0.0, 0.001, 0.011, 0.012, 0.022, 0.023, 0.024])
|
path = controller.run_once(
|
collector.load_config(self.config_path), begin, Fake(), monotonic=lambda: next(ticks),
|
wall_now=lambda: collector.parse_datetime("2026-08-13T10:00:24+08:00", "wall"),
|
)
|
self.assertEqual({"tabs": 1, "reload": 1, "goto": 0, "evaluate": 1}, calls)
|
code, result = collector.run(["--config", str(self.config_path), "refresh-commit", "--input", str(path), "--now", "2026-08-13T10:00:25+08:00"])
|
self.assertEqual((0, "REFRESH_CONFIRMED_NO_NEW"), (code, result["status"]))
|
|
Path(str(result["run_evidence_path"])).unlink()
|
(self.state / "refresh" / "runs" / "latest.json").unlink()
|
begin = self.begin()
|
forged = self.evidence(begin, outcome="READABLE", observations=[self.observation([])], marker=self.end_marker())
|
forged["controller_attestation"]["binding_sha256"] = "0" * 64
|
path = Path(str(begin["evidence_path"]))
|
path.parent.mkdir(parents=True, exist_ok=True)
|
path.write_text(json.dumps(forged), encoding="utf-8")
|
with self.assertRaises(collector.CollectorError) as failure:
|
collector.run(["--config", str(self.config_path), "refresh-commit", "--input", str(path), "--now", "2026-08-13T10:00:25+08:00"])
|
self.assertIn(failure.exception.code, {"E_CONTROLLER_ATTESTATION", "E_CONTROLLER_REQUIRED"})
|
|
def test_reviewer_real_js_fixture_adapts_to_accepted_schema(self) -> None:
|
runner = PROJECT_DEV / "test" / "fixtures" / "bili_dynamic_collector" / "page_extract_fixture_runner.js"
|
expected = {
|
"new": (0, "NEW_ITEMS_SAVED", None),
|
"empty": (0, "REFRESH_CONFIRMED_NO_NEW", None),
|
"unparsed": (4, "PARTIAL_DISCOVERY_UNCONFIRMED", "E_COVERAGE_INCOMPLETE"),
|
"identity": (3, "REFRESH_BLOCKED_AUTH_OR_ACCESS", "E_CREATOR_MISMATCH"),
|
}
|
for fixture_case, wanted in expected.items():
|
with self.subTest(fixture_case=fixture_case):
|
self.reset_runtime_fixture()
|
completed = subprocess.run(
|
["node", str(runner), str(refresh.EXTRACTOR_SOURCE), fixture_case],
|
capture_output=True, check=True, timeout=10,
|
)
|
raw = json.loads(completed.stdout.decode("utf-8"))
|
begin = self.begin()
|
calls = {"reload": 0, "evaluate": 0}
|
|
class Fake:
|
def open_tabs(inner) -> list[dict[str, object]]:
|
return [{"url": "https://space.bilibili.com/1420210197/dynamic"}]
|
|
def reload(inner, tab: dict[str, object], timeout_seconds: int) -> None:
|
calls["reload"] += 1
|
|
def goto(inner, url: str, timeout_seconds: int) -> dict[str, object]:
|
raise AssertionError("exact fixture tab must use reload")
|
|
def evaluate(inner, tab: dict[str, object], source: str, timeout_seconds: int) -> dict[str, object]:
|
calls["evaluate"] += 1
|
self.assertIn("projectInfoCollectVisibleDynamicNodes({page_internal_settle_timeout_ms:15000})", source)
|
return raw
|
|
ticks = iter([0.0, 0.001, 0.011, 0.012, 0.022, 0.023, 0.024])
|
path = controller.run_once(
|
collector.load_config(self.config_path), begin, Fake(),
|
monotonic=lambda: next(ticks),
|
wall_now=lambda: collector.parse_datetime("2026-08-13T10:00:24+08:00", "wall"),
|
)
|
code, result = collector.run([
|
"--config", str(self.config_path), "refresh-commit", "--input", str(path),
|
"--now", "2026-08-13T10:00:25+08:00",
|
])
|
self.assertEqual(wanted, (code, result["status"], result["error_code"]))
|
self.assertEqual({"reload": 1, "evaluate": 1}, calls)
|
|
self.reset_runtime_fixture()
|
begin = self.begin()
|
|
class AccessBlocked:
|
def open_tabs(inner) -> list[dict[str, object]]:
|
return [{"url": "https://space.bilibili.com/1420210197/dynamic"}]
|
|
def reload(inner, tab: dict[str, object], timeout_seconds: int) -> None:
|
return None
|
|
def goto(inner, url: str, timeout_seconds: int) -> dict[str, object]:
|
raise AssertionError("exact fixture tab must use reload")
|
|
def evaluate(inner, tab: dict[str, object], source: str, timeout_seconds: int) -> dict[str, object]:
|
raise PermissionError("synthetic access interstitial")
|
|
ticks = iter([0.0, 0.001, 0.011, 0.012, 0.022, 0.023, 0.024])
|
path = controller.run_once(
|
collector.load_config(self.config_path), begin, AccessBlocked(),
|
monotonic=lambda: next(ticks),
|
wall_now=lambda: collector.parse_datetime("2026-08-13T10:00:24+08:00", "wall"),
|
)
|
code, result = collector.run([
|
"--config", str(self.config_path), "refresh-commit", "--input", str(path),
|
"--now", "2026-08-13T10:00:25+08:00",
|
])
|
self.assertEqual((3, "REFRESH_BLOCKED_AUTH_OR_ACCESS", "E_ACCESS_BLOCKED"), (code, result["status"], result["error_code"]))
|
|
def test_reviewer_deadline_plus_epsilon_never_binds_or_saves(self) -> None:
|
begin = self.begin()
|
evidence = self.evidence(begin, outcome="READABLE", observations=[self.observation([])], marker=self.end_marker())
|
path = Path(str(begin["evidence_path"]))
|
self.write_json(path, evidence)
|
before = refresh._identity(self.archive / "manifest.jsonl").as_dict()
|
code, result = collector.run(["--config", str(self.config_path), "refresh-commit", "--input", str(path), "--now", "2026-08-13T10:02:01+08:00"])
|
self.assertEqual((4, "E_OVERALL_DEADLINE"), (code, result["error_code"]))
|
self.assertEqual(before, refresh._identity(self.archive / "manifest.jsonl").as_dict())
|
|
def test_reviewer_controller_commit_crossing_deadline_removes_unbound_evidence(self) -> None:
|
begin = self.begin()
|
|
class Fake:
|
def open_tabs(inner) -> list[dict[str, object]]:
|
return [{"url": "https://space.bilibili.com/1420210197/dynamic"}]
|
|
def reload(inner, tab: dict[str, object], timeout_seconds: int) -> None:
|
return None
|
|
def goto(inner, url: str, timeout_seconds: int) -> dict[str, object]:
|
raise AssertionError("exact fixture tab must use reload")
|
|
def evaluate(inner, tab: dict[str, object], source: str, timeout_seconds: int) -> dict[str, object]:
|
return {
|
"schema_version": 1, "ready_state": "complete", "visibility_state": "visible",
|
"final_url": "https://space.bilibili.com/1420210197/dynamic",
|
"page_title": "青枫浦上Q个人动态-青枫浦上Q动态记录-哔哩哔哩视频",
|
"creator": {"uid": "1420210197", "name": "青枫浦上Q", "profile_url": "https://space.bilibili.com/1420210197"},
|
"cards": [], "unparsed_nodes": [], "terminal_marker_text": "已经到底了", "limit_hit": "NONE",
|
}
|
|
ticks = iter([0.0, 0.001, 0.011, 0.012, 0.022, 119.999, 120.001])
|
with self.assertRaises(collector.CollectorError) as failure:
|
controller.run_once(
|
collector.load_config(self.config_path), begin, Fake(),
|
monotonic=lambda: next(ticks),
|
wall_now=lambda: collector.parse_datetime("2026-08-13T10:00:24+08:00", "wall"),
|
)
|
self.assertEqual("E_OVERALL_DEADLINE", failure.exception.code)
|
self.assertFalse(Path(str(begin["evidence_path"])).exists())
|
pending = refresh._load_pending(collector.load_config(self.config_path))
|
self.assertIsNotNone(pending)
|
self.assertEqual("AWAITING_EVIDENCE", pending["phase"])
|
self.assertFalse((self.archive / "manifest.jsonl").exists())
|
|
def test_reviewer_legacy_pending_never_requests_browser_action(self) -> None:
|
for slot_present in (False, True):
|
for after_deadline in (False, True):
|
for entrypoint in ("begin", "commit", "controller"):
|
with self.subTest(slot_present=slot_present, after_deadline=after_deadline, entrypoint=entrypoint):
|
self.reset_runtime_fixture()
|
begin = self.begin()
|
config = collector.load_config(self.config_path)
|
pending = refresh._load_pending(config)
|
assert pending is not None
|
pending["schema_version"] = refresh.LEGACY_PENDING_SCHEMA
|
pending.pop("runtime_contract")
|
pending.pop("controller_capability")
|
refresh._write_pending(config, pending)
|
if not slot_present:
|
Path(str(begin["run_evidence_path"])).unlink()
|
now_text = "2026-08-13T10:02:01+08:00" if after_deadline else "2026-08-13T10:00:30+08:00"
|
if entrypoint == "controller":
|
class NeverCalled:
|
def open_tabs(inner) -> list[dict[str, object]]:
|
raise AssertionError("legacy pending must not touch the browser adapter")
|
|
with self.assertRaises(collector.CollectorError) as failure:
|
controller.run_once(config, begin, NeverCalled())
|
self.assertEqual("E_CONTROLLER_PENDING", failure.exception.code)
|
continue
|
if entrypoint == "begin":
|
invoke = lambda: refresh.refresh_begin(config, self.config_path, collector.parse_datetime(now_text, "now"))
|
else:
|
invoke = lambda: refresh.refresh_commit(config, self.config_path, Path(str(begin["evidence_path"])), collector.parse_datetime(now_text, "now"))
|
if after_deadline:
|
result = invoke()
|
self.assertEqual(("REFRESH_FAILED_PAGE_UNREADABLE", "E_LEGACY_EVIDENCE_MISSING_AFTER_DEADLINE"), (result["status"], result["error_code"]))
|
self.assertNotEqual("BROWSER_REFRESH_REQUIRED", result["status"])
|
else:
|
with self.assertRaises(collector.CollectorError) as failure:
|
invoke()
|
self.assertEqual("E_LEGACY_RECOVERY_ONLY", failure.exception.code)
|
self.assertFalse((self.archive / "manifest.jsonl").exists())
|
|
def test_reviewer_identity_mismatch_is_durable_blocked_terminal(self) -> None:
|
identities = {
|
"uid": lambda evidence: evidence["creator"].__setitem__("uid", "1"),
|
"name": lambda evidence: evidence["creator"].__setitem__("name", "其他用户"),
|
"profile": lambda evidence: evidence["creator"].__setitem__("profile_url", "https://space.bilibili.com/1"),
|
"final_url": lambda evidence: evidence.__setitem__("final_url", "https://space.bilibili.com/1/dynamic"),
|
}
|
outcomes = {
|
"CONFIRMED": "NONE",
|
"TIMEOUT": "ACTION_TIMEOUT",
|
"POST_DISPATCH_ERROR": "ACTION_POST_DISPATCH",
|
}
|
for identity_name, mutate in identities.items():
|
for action_outcome, diagnostic in outcomes.items():
|
with self.subTest(identity=identity_name, action_outcome=action_outcome):
|
self.reset_runtime_fixture()
|
begin = self.begin()
|
evidence = self.evidence(
|
begin, outcome="READABLE", observations=[self.observation([])],
|
marker=self.end_marker(), action_outcome=action_outcome,
|
)
|
self.set_runtime_diagnostic(evidence, diagnostic)
|
mutate(evidence)
|
code, result = self.commit(begin, evidence)
|
self.assertEqual((3, "REFRESH_BLOCKED_AUTH_OR_ACCESS", "E_CREATOR_MISMATCH"), (code, result["status"], result["error_code"]))
|
slot = json.loads(Path(str(result["run_evidence_path"])).read_text(encoding="utf-8"))
|
self.assertEqual("TERMINAL", slot["run_state"])
|
self.assertFalse(refresh._pending_path(collector.load_config(self.config_path)).exists())
|
self.assertFalse(result["no_new_confirmed"])
|
self.assertFalse((self.archive / "manifest.jsonl").exists())
|
|
|
if __name__ == "__main__":
|
unittest.main()
|