from __future__ import annotations
|
|
import hashlib
|
import json
|
import os
|
import shutil
|
import subprocess
|
import sys
|
import tempfile
|
import unittest
|
from unittest import mock
|
from datetime import datetime, timezone
|
from pathlib import Path
|
|
|
PROJECT_DEV = Path(__file__).resolve().parents[1]
|
PROJECT_ROOT = PROJECT_DEV.parents[1]
|
sys.path.insert(0, str(PROJECT_DEV))
|
|
import bili_half_hour_pipeline as product
|
|
|
NOW = datetime(2026, 8, 29, 8, 0, tzinfo=timezone.utc)
|
UID = "1420210197"
|
BVID = "BV1Q541167Qg"
|
|
|
def canonical(value: object) -> bytes:
|
return (json.dumps(value, ensure_ascii=True, sort_keys=True, separators=(",", ":")) + "\n").encode("ascii")
|
|
|
class PipelineFixture:
|
def __init__(self, root: Path) -> None:
|
self.root = root
|
self.archive = root / "ana-data" / "news-fixture"
|
self.state = root / "dev" / "tmp" / "pipeline"
|
self.archive.mkdir(parents=True)
|
self.formal = self.archive / "manifest.jsonl"
|
self.handoffs = self.archive / "video-processing-handoffs.jsonl"
|
self.formal.write_bytes(canonical({
|
"schema_version": 1, "stable_id": "OLD", "item_type": "text", "status": "SAVED",
|
"creator_uid": UID, "path": "old.txt", "bytes": 3, "sha256": hashlib.sha256(b"old").hexdigest().upper()
|
}))
|
(self.archive / "old.txt").write_bytes(b"old")
|
self.handoffs.write_bytes(canonical({
|
"schema_version": 1, "type": "media-processing-handoff", "status": "READY",
|
"handoff_id": "HANDOFF-OLD", "queue_job_id": "0" * 64, "creator_uid": UID,
|
"bvid": "BV1Q541167Qf", "source_url": "https://www.bilibili.com/video/BV1Q541167Qf",
|
"media_path": "F:/video/old.mkv", "mapping_path": "F:/video/old.download.json",
|
"bytes": 1, "sha256": "A" * 64, "duration_seconds": 1.0,
|
"video_codec": "av1", "audio_codec": "aac", "created_at": NOW.isoformat()
|
}))
|
self.config_path = root / "pipeline.json"
|
self.config_path.write_bytes(canonical({
|
"schema_version": 1,
|
"task_id": product.TASK_ID,
|
"interval_minutes": 30,
|
"creator": {"uid": UID, "name": "fixture", "dynamic_url": f"https://space.bilibili.com/{UID}/dynamic"},
|
"paths": {
|
"project_root": str(root), "archive_root": str(self.archive),
|
"formal_manifest": str(self.formal), "processing_handoffs": str(self.handoffs),
|
"state_dir": str(self.state), "video_root": str(root / "external-video")
|
},
|
"downstream": {
|
"video_downloader_thread_id": "019fcc5d-798f-7ea1-8325-3a4d1f2dc5a5",
|
"media_processor_thread_id": "019fb7a4-bdfd-79f2-bd6b-e67e2b7d8efd",
|
"minutes_thread_id": "019fae88-ef98-7f83-b964-dcd4c5f842ef",
|
"reply_thread_id": "019fbcbb-bed7-7c90-83ab-f50610f80d3a"
|
},
|
"git": {
|
"remote": "origin",
|
"branch": "main",
|
"allowed_extensions": [".txt", ".md", ".json", ".jsonl", ".srt", ".pdf", ".png", ".jpg"],
|
"allowed_docs": ["ana-data/news-fixture/目录导读.md"],
|
}
|
}))
|
(self.archive / "目录导读.md").write_bytes(b"# fixture\n")
|
self.config = product.load_config(self.config_path)
|
|
def append_formal(self, value: dict[str, object]) -> None:
|
with self.formal.open("ab") as stream:
|
stream.write(canonical(value))
|
|
def append_handoff(self, value: dict[str, object]) -> None:
|
with self.handoffs.open("ab") as stream:
|
stream.write(canonical(value))
|
|
|
def prepare_migration_fixture(root: Path) -> tuple[PipelineFixture, dict[str, Path], dict[str, bytes], Path]:
|
fixture = PipelineFixture(root)
|
(fixture.root / "external-video").mkdir()
|
fixture.append_formal({
|
"schema_version": 1,
|
"stable_id": BVID,
|
"item_type": "video",
|
"status": product.VIDEO_COMPLETE,
|
"creator": "fixture",
|
"title": "A / canonical: title?",
|
"published_at": "2026-08-30T12:34:56+08:00",
|
})
|
transcript_dir = fixture.archive / f"{BVID}.transcript"
|
minutes_dir = fixture.archive / f"{BVID}.minutes"
|
transcript_dir.mkdir()
|
minutes_dir.mkdir()
|
payloads = {
|
"transcript_txt": b"transcript\n",
|
"transcript_srt": b"1\n00:00:00,000 --> 00:00:01,000\ntext\n",
|
"transcript_json": b'{"segments":[]}\n',
|
"minutes_md": b"# minutes\n",
|
"minutes_pdf": b"%PDF-fixture\n",
|
}
|
legacy_paths = {
|
"transcript_txt": transcript_dir / f"{BVID}.txt",
|
"transcript_srt": transcript_dir / f"{BVID}.srt",
|
"transcript_json": transcript_dir / f"{BVID}.json",
|
"minutes_md": minutes_dir / f"{BVID}.md",
|
"minutes_pdf": minutes_dir / f"{BVID}.pdf",
|
}
|
for kind, path in legacy_paths.items():
|
path.write_bytes(payloads[kind])
|
flac = transcript_dir / f"{BVID}.audio.flac"
|
flac.write_bytes(b"fLaC-fixture")
|
subprocess.run(["git", "init", "-b", "main"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "config", "user.name", "Fixture"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "config", "user.email", "fixture@example.invalid"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "add", "."], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "commit", "-m", "legacy baseline"], cwd=fixture.root, check=True, capture_output=True)
|
payloads["minutes_pdf"] = b"%PDF-modified-worktree-fixture\n"
|
legacy_paths["minutes_pdf"].write_bytes(payloads["minutes_pdf"])
|
(fixture.archive / "目录导读.md").write_bytes(b"# canonical naming guide\n")
|
product.initialize(fixture.config, NOW)
|
return fixture, legacy_paths, payloads, flac
|
|
|
class HalfHourPipelineTests(unittest.TestCase):
|
@staticmethod
|
def _new_handoff(fixture: PipelineFixture, *, source_url: str | None = None) -> None:
|
fixture.append_handoff({
|
"schema_version": 1, "type": "media-processing-handoff", "status": "READY",
|
"handoff_id": "HANDOFF-NEW", "queue_job_id": "1" * 64, "creator_uid": UID, "bvid": BVID,
|
"source_url": source_url or f"https://www.bilibili.com/video/{BVID}", "media_path": "F:/video/new.mkv",
|
"mapping_path": "F:/video/new.download.json", "bytes": 100, "sha256": "B" * 64,
|
"duration_seconds": 60.0, "video_codec": "av1", "audio_codec": "aac", "created_at": NOW.isoformat(),
|
})
|
|
def test_external_schema_integers_are_strict_and_fail_before_mutation(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
original = json.loads(fixture.config_path.read_text(encoding="ascii"))
|
for field, bad_values in {
|
"schema_version": [True, 1.0, "1", None],
|
"interval_minutes": [True, 30.0, "30", None],
|
}.items():
|
for bad in bad_values:
|
candidate = dict(original)
|
candidate[field] = bad
|
fixture.config_path.write_bytes(canonical(candidate))
|
with self.assertRaises(product.PipelineError) as rejected:
|
product.load_config(fixture.config_path)
|
self.assertEqual("E_CONFIG", rejected.exception.code)
|
self.assertFalse(fixture.state.exists())
|
fixture.config_path.write_bytes(canonical(original))
|
config = product.load_config(fixture.config_path)
|
product.initialize(config, NOW)
|
state_before = config.state_path.read_bytes()
|
state = json.loads(state_before)
|
state["schema_version"] = True
|
config.state_path.write_bytes(canonical(state))
|
with self.assertRaises(product.PipelineError) as rejected:
|
product.begin(config, NOW)
|
self.assertEqual("E_STATE", rejected.exception.code)
|
self.assertFalse(config.runs_path.exists())
|
|
def test_processing_handoff_accepts_exact_schema_aliases_and_rejects_schema_drift(self) -> None:
|
def handoff(schema_key: str | None, schema_value: object = 1) -> dict[str, object]:
|
value: dict[str, object] = {
|
"type": "media-processing-handoff", "status": "READY",
|
"handoff_id": "HANDOFF-SCHEMA", "queue_job_id": "1" * 64,
|
"creator_uid": UID, "bvid": BVID, "source_url": f"https://www.bilibili.com/video/{BVID}",
|
"media_path": "F:/video/schema.mkv", "mapping_path": "F:/video/schema.download.json",
|
"bytes": 100, "sha256": "B" * 64, "duration_seconds": 60.0,
|
"video_codec": "av1", "audio_codec": "aac", "created_at": NOW.isoformat(),
|
}
|
if schema_key is not None:
|
value[schema_key] = schema_value
|
return value
|
|
positive_sha = (
|
("schema", "b" * 64),
|
("schema_version", "B" * 64),
|
("schema_version", "Aa" * 32),
|
)
|
for schema_key, sha256 in positive_sha:
|
with self.subTest(positive=(schema_key, sha256[:2])), tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
row = handoff(schema_key)
|
row["sha256"] = sha256
|
fixture.append_handoff(row)
|
created = product.reconcile(fixture.config, NOW)["created_outbox_ids"]
|
self.assertEqual(1, len(created))
|
item = product.pending(fixture.config)["items"][0]
|
self.assertEqual(("VIDEO_TRANSCRIPTION_READY", BVID), (item["kind"], item["payload"]["bvid"]))
|
self.assertEqual(sha256.upper(), item["payload"]["sha256"])
|
|
invalid: list[tuple[str, dict[str, object]]] = [
|
("bool", handoff("schema", True)),
|
("float", handoff("schema", 1.0)),
|
("string", handoff("schema", "1")),
|
("null", handoff("schema", None)),
|
("wrong_integer", handoff("schema", 2)),
|
("missing", handoff(None)),
|
]
|
invalid_sha = handoff("schema")
|
invalid_sha["sha256"] = "G" * 64
|
invalid.append(("invalid_sha", invalid_sha))
|
unicode_expansion_sha = handoff("schema")
|
unicode_expansion_sha["sha256"] = "0" * 62 + "\ufb00"
|
invalid.append(("unicode_expansion_sha", unicode_expansion_sha))
|
dual = handoff("schema")
|
dual["schema_version"] = 1
|
invalid.append(("dual", dual))
|
extra = handoff("schema")
|
extra["schema_extra"] = 1
|
invalid.append(("extra", extra))
|
for label, row in invalid:
|
with self.subTest(negative=label), tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
fixture.append_handoff(row)
|
state_before = fixture.config.state_path.read_bytes()
|
with self.assertRaises(product.PipelineError) as rejected:
|
product.reconcile(fixture.config, NOW)
|
self.assertEqual("E_SOURCE_BINDING", rejected.exception.code)
|
self.assertEqual(state_before, fixture.config.state_path.read_bytes())
|
self.assertFalse(fixture.config.outbox_path.exists())
|
|
def test_outbox_history_schema_tamper_fails_without_append(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
fixture.append_handoff({
|
"schema_version": 1, "type": "media-processing-handoff", "status": "READY",
|
"handoff_id": "HANDOFF-NEW", "queue_job_id": "1" * 64, "creator_uid": UID, "bvid": BVID,
|
"source_url": f"https://www.bilibili.com/video/{BVID}", "media_path": "F:/video/new.mkv",
|
"mapping_path": "F:/video/new.download.json", "bytes": 100, "sha256": "B" * 64,
|
"duration_seconds": 60.0, "video_codec": "av1", "audio_codec": "aac", "created_at": NOW.isoformat(),
|
})
|
product.reconcile(fixture.config, NOW)
|
rows = [json.loads(line) for line in fixture.config.outbox_path.read_text(encoding="ascii").splitlines()]
|
rows[0]["schema_version"] = True
|
fixture.config.outbox_path.write_bytes(b"".join(canonical(row) for row in rows))
|
before = fixture.config.outbox_path.read_bytes()
|
with self.assertRaises(product.PipelineError) as rejected:
|
product.pending(fixture.config)
|
self.assertEqual("E_OUTBOX", rejected.exception.code)
|
self.assertEqual(before, fixture.config.outbox_path.read_bytes())
|
|
def test_baseline_nonoverlap_reconcile_and_exact_once_outbox(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
initialized = product.initialize(fixture.config, NOW)
|
self.assertEqual(("INITIALIZED", 1, 1), (
|
initialized["status"], initialized["baseline"]["formal"]["lines"], initialized["baseline"]["handoff"]["lines"]
|
))
|
self.assertEqual("RUN_STARTED", product.begin(fixture.config, NOW)["status"])
|
self.assertEqual("RUN_RESUMED", product.begin(fixture.config, NOW)["status"])
|
with self.assertRaises(product.PipelineError) as active:
|
product.begin(fixture.config, datetime(2026, 8, 29, 8, 30, tzinfo=timezone.utc))
|
self.assertEqual("E_RUN_ACTIVE", active.exception.code)
|
|
body = b"new text\n"
|
image = b"\x89PNG\r\nfixture"
|
(fixture.archive / "new.txt").write_bytes(body)
|
(fixture.archive / "new.png").write_bytes(image)
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": "123456", "item_type": "text", "status": "SAVED",
|
"creator_uid": UID, "path": "new.txt", "bytes": len(body), "sha256": hashlib.sha256(body).hexdigest().upper(),
|
"image_path": "new.png", "image_bytes": len(image), "image_sha256": hashlib.sha256(image).hexdigest().upper()
|
})
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": BVID, "item_type": "video", "status": "VIDEO_DOWNLOAD_PENDING_EXTENSION",
|
"creator_uid": UID, "source_url": f"https://www.bilibili.com/video/{BVID}", "title": "fixture video",
|
"published_at": NOW.isoformat(), "expected_duration_seconds": 60
|
})
|
fixture.append_handoff({
|
"schema_version": 1, "type": "media-processing-handoff", "status": "READY",
|
"handoff_id": "HANDOFF-NEW", "queue_job_id": "1" * 64, "creator_uid": UID, "bvid": BVID,
|
"source_url": f"https://www.bilibili.com/video/{BVID}", "media_path": "F:/video/new.mkv",
|
"mapping_path": "F:/video/new.download.json", "bytes": 100, "sha256": "B" * 64,
|
"duration_seconds": 60.0, "video_codec": "av1", "audio_codec": "aac", "created_at": NOW.isoformat()
|
})
|
result = product.reconcile(fixture.config, NOW)
|
self.assertEqual(3, len(result["created_outbox_ids"]))
|
self.assertEqual(0, len(product.reconcile(fixture.config, NOW)["created_outbox_ids"]))
|
pending = product.pending(fixture.config)
|
self.assertEqual(3, pending["count"])
|
self.assertEqual(
|
{"GIT_DELIVERY_READY", "VIDEO_DOWNLOAD_READY", "VIDEO_TRANSCRIPTION_READY"},
|
{item["kind"] for item in pending["items"]},
|
)
|
video_item = next(item for item in pending["items"] if item["kind"] == "VIDEO_DOWNLOAD_READY")
|
self.assertEqual(
|
fixture.config.video_downloader_thread_id,
|
product.dispatch_intent(fixture.config, video_item["outbox_id"], NOW)["target_thread_id"],
|
)
|
serialized = fixture.config.outbox_path.read_text(encoding="ascii").lower()
|
for secret in ("cookie", "sessdata", "localstorage", "profile"):
|
self.assertNotIn(secret, serialized)
|
|
def test_dispatch_intent_is_restart_safe_and_observed_exact_once(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
fixture.append_handoff({
|
"schema_version": 1, "type": "media-processing-handoff", "status": "READY",
|
"handoff_id": "HANDOFF-NEW", "queue_job_id": "1" * 64, "creator_uid": UID, "bvid": BVID,
|
"source_url": f"https://www.bilibili.com/video/{BVID}", "media_path": "F:/video/new.mkv",
|
"mapping_path": "F:/video/new.download.json", "bytes": 100, "sha256": "B" * 64,
|
"duration_seconds": 60.0, "video_codec": "av1", "audio_codec": "aac", "created_at": NOW.isoformat()
|
})
|
outbox_id = product.reconcile(fixture.config, NOW)["created_outbox_ids"][0]
|
intent = product.dispatch_intent(fixture.config, outbox_id, NOW)
|
self.assertEqual(("DISPATCH_INTENT_DURABLE", fixture.config.media_thread_id), (intent["status"], intent["target_thread_id"]))
|
self.assertEqual(outbox_id, intent["envelope"]["outbox_id"])
|
self.assertEqual("DISPATCH_INTENT", product.pending(fixture.config)["items"][0]["delivery_state"])
|
resumed = product.dispatch_intent(fixture.config, outbox_id, NOW)
|
self.assertEqual(("DISPATCH_INTENT_RESUMED", intent["envelope"]), (resumed["status"], resumed["envelope"]))
|
observed = product.observe_dispatch(fixture.config, outbox_id, "DELIVERY-1", NOW)
|
self.assertEqual("DISPATCH_OBSERVED", observed["status"])
|
self.assertEqual("DISPATCH_ALREADY_OBSERVED", product.observe_dispatch(fixture.config, outbox_id, "DELIVERY-1", NOW)["status"])
|
self.assertEqual(0, product.pending(fixture.config)["count"])
|
with self.assertRaises(product.PipelineError) as conflict:
|
product.observe_dispatch(fixture.config, outbox_id, "DELIVERY-2", NOW)
|
self.assertEqual("E_DISPATCH_RECEIPT", conflict.exception.code)
|
|
def test_transcript_and_minutes_receipts_create_restart_safe_downstream(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
fixture.append_handoff({
|
"schema_version": 1, "type": "media-processing-handoff", "status": "READY",
|
"handoff_id": "HANDOFF-NEW", "queue_job_id": "1" * 64, "creator_uid": UID, "bvid": BVID,
|
"source_url": f"https://www.bilibili.com/video/{BVID}", "media_path": "F:/video/new.mkv",
|
"mapping_path": "F:/video/new.download.json", "bytes": 100, "sha256": "B" * 64,
|
"duration_seconds": 60.0, "video_codec": "av1", "audio_codec": "aac", "created_at": NOW.isoformat()
|
})
|
source_id = product.reconcile(fixture.config, NOW)["created_outbox_ids"][0]
|
product.dispatch_intent(fixture.config, source_id, NOW)
|
product.observe_dispatch(fixture.config, source_id, "TRANSCRIPT-DISPATCH-1", NOW)
|
transcript = fixture.archive / "new.transcript.txt"
|
transcript_srt = fixture.archive / "new.transcript.srt"
|
transcript_json = fixture.archive / "new.transcript.json"
|
transcript.write_bytes(b"transcript\n")
|
transcript_srt.write_bytes(b"1\n00:00:00,000 --> 00:00:01,000\ntranscript\n")
|
transcript_json.write_bytes(b'{"segments":[]}\n')
|
receipt = fixture.root / "transcript-receipt.json"
|
receipt.write_bytes(canonical({
|
"schema_version": 1, "type": "TRANSCRIPTION_COMPLETE", "outbox_id": source_id,
|
"stable_id": BVID, "terminal_id": "TRANSCRIPT-1", "status": "COMPLETE",
|
"files": [
|
{"path": str(transcript), "bytes": transcript.stat().st_size,
|
"sha256": hashlib.sha256(transcript.read_bytes()).hexdigest().upper(), "kind": "transcript_txt"},
|
{"path": str(transcript_srt), "bytes": transcript_srt.stat().st_size,
|
"sha256": hashlib.sha256(transcript_srt.read_bytes()).hexdigest().upper(), "kind": "transcript_srt"},
|
{"path": str(transcript_json), "bytes": transcript_json.stat().st_size,
|
"sha256": hashlib.sha256(transcript_json.read_bytes()).hexdigest().upper(), "kind": "transcript_json"},
|
],
|
"created_at": NOW.isoformat()
|
}))
|
committed = product.ingest_receipt(fixture.config, receipt, NOW)
|
self.assertEqual(("RECEIPT_COMMITTED", 2), (committed["status"], len(committed["created_outbox_ids"])))
|
replay = product.ingest_receipt(fixture.config, receipt, NOW)
|
self.assertEqual("RECEIPT_ALREADY_COMMITTED", replay["status"])
|
terminal_before = fixture.config.terminals_path.read_bytes()
|
outbox_before = fixture.config.outbox_path.read_bytes()
|
conflicting = json.loads(receipt.read_text(encoding="ascii"))
|
conflicting["terminal_id"] = "TRANSCRIPT-CONFLICT"
|
receipt.write_bytes(canonical(conflicting))
|
with self.assertRaises(product.PipelineError) as conflict:
|
product.ingest_receipt(fixture.config, receipt, NOW)
|
self.assertEqual("E_RECEIPT_CONFLICT", conflict.exception.code)
|
self.assertEqual(terminal_before, fixture.config.terminals_path.read_bytes())
|
self.assertEqual(outbox_before, fixture.config.outbox_path.read_bytes())
|
items = product.pending(fixture.config)["items"]
|
minutes_item = next(item for item in items if item["kind"] == "MINUTES_READY")
|
product.dispatch_intent(fixture.config, minutes_item["outbox_id"], NOW)
|
product.observe_dispatch(fixture.config, minutes_item["outbox_id"], "MINUTES-DISPATCH-1", NOW)
|
minutes = fixture.archive / "new.minutes.md"
|
minutes_pdf = fixture.archive / "new.minutes.pdf"
|
minutes.write_bytes(b"# minutes\n")
|
minutes_pdf.write_bytes(b"%PDF-fixture\n")
|
minutes_receipt = fixture.root / "minutes-receipt.json"
|
minutes_receipt.write_bytes(canonical({
|
"schema_version": 1, "type": "MINUTES_COMPLETE", "outbox_id": minutes_item["outbox_id"],
|
"stable_id": BVID, "terminal_id": "MINUTES-1", "status": "COMPLETE",
|
"files": [
|
{"path": str(minutes), "bytes": minutes.stat().st_size,
|
"sha256": hashlib.sha256(minutes.read_bytes()).hexdigest().upper(), "kind": "minutes_md"},
|
{"path": str(minutes_pdf), "bytes": minutes_pdf.stat().st_size,
|
"sha256": hashlib.sha256(minutes_pdf.read_bytes()).hexdigest().upper(), "kind": "minutes_pdf"},
|
],
|
"created_at": NOW.isoformat()
|
}))
|
result = product.ingest_receipt(fixture.config, minutes_receipt, NOW)
|
self.assertEqual(1, len(result["created_outbox_ids"]))
|
git_items = [
|
row for row in product._outbox_rows(fixture.config)
|
if row.get("event") == "CREATED" and row.get("kind") == "GIT_DELIVERY_READY"
|
]
|
self.assertEqual(2, len(git_items))
|
self.assertEqual([3, 2], [len(product._git_artifacts(fixture.config, row["payload"]["files"])[0]) for row in git_items])
|
mismatched = dict(git_items[0]["payload"]["files"][0])
|
mismatched["kind"] = "minutes_pdf"
|
with self.assertRaises(product.PipelineError) as kind_drift:
|
product._git_artifacts(fixture.config, [mismatched])
|
self.assertEqual("E_GIT_SCOPE", kind_drift.exception.code)
|
|
def test_git_delivery_uses_exact_paths_and_nonforce_push(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
(fixture.root / ".git").mkdir()
|
(fixture.root / ".git" / "index").write_bytes(b"USER-INDEX")
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
artifact = fixture.archive / "new.txt"
|
artifact.write_bytes(b"content\n")
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": "123", "item_type": "text", "status": "SAVED", "creator_uid": UID,
|
"path": "new.txt", "bytes": artifact.stat().st_size, "sha256": hashlib.sha256(artifact.read_bytes()).hexdigest().upper()
|
})
|
outbox_id = product.reconcile(fixture.config, NOW)["created_outbox_ids"][0]
|
commands: list[list[str]] = []
|
rel = artifact.relative_to(fixture.root).as_posix()
|
|
def runner(command: list[str], **_: object) -> subprocess.CompletedProcess[str]:
|
commands.append(command)
|
if command[:5] == ["git", "diff", "--cached", "--name-only", "-z"] and command[5] != "--":
|
return subprocess.CompletedProcess(command, 0, b"", b"")
|
if command[:6] == ["git", "diff", "--cached", "--name-only", "-z", "--"]:
|
return subprocess.CompletedProcess(command, 0, rel.encode("utf-8") + b"\0", b"")
|
if command[:4] == ["git", "diff", "--cached", "--name-only"]:
|
return subprocess.CompletedProcess(command, 0, "", "")
|
if command[:3] == ["git", "rev-parse", "HEAD"]:
|
return subprocess.CompletedProcess(command, 0, "a" * 40 + "\n", "")
|
if command[:2] == ["git", "hash-object"]:
|
return subprocess.CompletedProcess(command, 0, "d" * 40 + "\n", "")
|
if command[:2] == ["git", "write-tree"]:
|
return subprocess.CompletedProcess(command, 0, "b" * 40 + "\n", "")
|
if command[:2] == ["git", "commit-tree"]:
|
return subprocess.CompletedProcess(command, 0, "c" * 40 + "\n", "")
|
return subprocess.CompletedProcess(command, 0, "", "")
|
|
result = product.git_deliver(fixture.config, outbox_id, NOW, runner=runner)
|
self.assertEqual("GIT_PUSHED", result["status"])
|
self.assertIn(["git", "update-index", "--add", "--cacheinfo", f"100644,{'d' * 40},{rel}"], commands)
|
self.assertIn(["git", "push", "origin", "c" * 40 + ":refs/heads/main"], commands)
|
flattened = "\n".join(" ".join(command) for command in commands)
|
self.assertNotIn("git add .", flattened)
|
self.assertNotIn("--force", flattened)
|
|
def test_git_delivery_accepts_exact_published_head_without_git_writes(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
for command in (
|
["git", "init", "-b", "main"],
|
["git", "config", "user.name", "Fixture"],
|
["git", "config", "user.email", "fixture@example.invalid"],
|
["git", "add", "."],
|
["git", "commit", "-m", "fixture baseline"],
|
):
|
subprocess.run(command, cwd=fixture.root, check=True, capture_output=True)
|
remote = fixture.root / "remote.git"
|
subprocess.run(["git", "init", "--bare", str(remote)], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "remote", "add", "origin", str(remote)], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "push", "-u", "origin", "main"], cwd=fixture.root, check=True, capture_output=True)
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
artifact = fixture.archive / "已交付.txt"
|
artifact.write_bytes(b"already published\n")
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": "published", "item_type": "text", "status": "SAVED",
|
"creator_uid": UID, "path": artifact.name, "bytes": artifact.stat().st_size,
|
"sha256": hashlib.sha256(artifact.read_bytes()).hexdigest().upper(),
|
})
|
subprocess.run(
|
["git", "add", "--", artifact.relative_to(fixture.root).as_posix(), fixture.formal.relative_to(fixture.root).as_posix()],
|
cwd=fixture.root, check=True, capture_output=True,
|
)
|
subprocess.run(["git", "commit", "-m", "accepted data delivery"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "push", "origin", "main"], cwd=fixture.root, check=True, capture_output=True)
|
head_before = subprocess.run(
|
["git", "rev-parse", "HEAD"], cwd=fixture.root, check=True, text=True, capture_output=True,
|
).stdout.strip()
|
shared_index = fixture.root / ".git" / "index"
|
index_before = shared_index.read_bytes()
|
outbox_id = product.reconcile(fixture.config, NOW)["created_outbox_ids"][0]
|
commands: list[list[str]] = []
|
|
def runner(command: list[str], **kwargs: object) -> subprocess.CompletedProcess[Any]:
|
commands.append(command)
|
return subprocess.run(command, **kwargs)
|
|
result = product.git_deliver(fixture.config, outbox_id, NOW, runner=runner)
|
self.assertEqual(("GIT_NO_CHANGES", head_before), (result["status"], result["commit_sha"]))
|
self.assertEqual(index_before, shared_index.read_bytes())
|
self.assertEqual(head_before, subprocess.run(
|
["git", "rev-parse", "HEAD"], cwd=fixture.root, check=True, text=True, capture_output=True,
|
).stdout.strip())
|
self.assertEqual(head_before, subprocess.run(
|
["git", "rev-parse", "refs/remotes/origin/main"], cwd=fixture.root, check=True, text=True, capture_output=True,
|
).stdout.strip())
|
forbidden = {"hash-object", "update-index", "commit-tree", "update-ref", "push"}
|
self.assertFalse(any(len(command) > 1 and command[1] in forbidden for command in commands))
|
self.assertEqual(0, len(list(fixture.config.state_dir.glob(f"git-index-{outbox_id}*"))))
|
events = [json.loads(line) for line in fixture.config.outbox_path.read_text(encoding="ascii").splitlines()]
|
complete = [row for row in events if row.get("outbox_id") == outbox_id and row.get("event") == "COMPLETE"]
|
self.assertEqual(["NO_CHANGES"], [row["result"] for row in complete])
|
|
local_only = fixture.archive / "local-only.txt"
|
local_only.write_bytes(b"not published\n")
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": "local-only", "item_type": "text", "status": "SAVED",
|
"creator_uid": UID, "path": local_only.name, "bytes": local_only.stat().st_size,
|
"sha256": hashlib.sha256(local_only.read_bytes()).hexdigest().upper(),
|
})
|
subprocess.run(
|
["git", "add", "--", local_only.relative_to(fixture.root).as_posix(), fixture.formal.relative_to(fixture.root).as_posix()],
|
cwd=fixture.root, check=True, capture_output=True,
|
)
|
subprocess.run(["git", "commit", "-m", "local only"], cwd=fixture.root, check=True, capture_output=True)
|
later = datetime(2026, 8, 29, 8, 30, tzinfo=timezone.utc)
|
local_outbox = product.reconcile(fixture.config, later)["created_outbox_ids"][0]
|
before = fixture.config.outbox_path.read_bytes()
|
with self.assertRaises(product.PipelineError) as unpublished:
|
product.git_deliver(fixture.config, local_outbox, later)
|
self.assertEqual("E_GIT_PUSH", unpublished.exception.code)
|
self.assertEqual(before, fixture.config.outbox_path.read_bytes())
|
|
def test_git_no_changes_rejects_terminal_boundary_artifact_and_ref_drift(self) -> None:
|
for mode, expected_code in (
|
("artifact_after_read", "E_GIT_SCOPE"),
|
("head_after_capture", "E_GIT_PUSH"),
|
("remote_after_read", "E_GIT_PUSH"),
|
):
|
with self.subTest(mode=mode), tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
for command in (
|
["git", "init", "-b", "main"],
|
["git", "config", "user.name", "Fixture"],
|
["git", "config", "user.email", "fixture@example.invalid"],
|
["git", "add", "."],
|
["git", "commit", "-m", "fixture baseline"],
|
):
|
subprocess.run(command, cwd=fixture.root, check=True, capture_output=True)
|
remote = fixture.root / "remote.git"
|
subprocess.run(["git", "init", "--bare", str(remote)], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "remote", "add", "origin", str(remote)], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "push", "-u", "origin", "main"], cwd=fixture.root, check=True, capture_output=True)
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
artifact = fixture.archive / "终态竞态.txt"
|
artifact.write_bytes(b"published artifact\n")
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": f"race-{mode}", "item_type": "text", "status": "SAVED",
|
"creator_uid": UID, "path": artifact.name, "bytes": artifact.stat().st_size,
|
"sha256": hashlib.sha256(artifact.read_bytes()).hexdigest().upper(),
|
})
|
subprocess.run(
|
["git", "add", "--", artifact.relative_to(fixture.root).as_posix(), fixture.formal.relative_to(fixture.root).as_posix()],
|
cwd=fixture.root, check=True, capture_output=True,
|
)
|
subprocess.run(["git", "commit", "-m", "accepted published artifact"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "push", "origin", "main"], cwd=fixture.root, check=True, capture_output=True)
|
outbox_id = product.reconcile(fixture.config, NOW)["created_outbox_ids"][0]
|
outbox_before = fixture.config.outbox_path.read_bytes()
|
injected = False
|
|
def runner(command: list[str], **kwargs: object) -> subprocess.CompletedProcess[Any]:
|
nonlocal injected
|
result = subprocess.run(command, **kwargs)
|
if not injected and mode in {"artifact_after_read", "head_after_capture"} and command[:2] == ["git", "show"]:
|
injected = True
|
if mode == "artifact_after_read":
|
identity = artifact.stat()
|
payload = artifact.read_bytes()
|
artifact.write_bytes(b"X" + payload[1:])
|
os.utime(artifact, ns=(identity.st_atime_ns, identity.st_mtime_ns))
|
else:
|
old_head = subprocess.run(
|
["git", "rev-parse", "HEAD"], cwd=fixture.root, check=True, text=True, capture_output=True,
|
).stdout.strip()
|
tree = subprocess.run(
|
["git", "rev-parse", "HEAD^{tree}"], cwd=fixture.root, check=True, text=True, capture_output=True,
|
).stdout.strip()
|
env = dict(os.environ)
|
env.update({
|
"GIT_AUTHOR_NAME": "Race", "GIT_COMMITTER_NAME": "Race",
|
"GIT_AUTHOR_EMAIL": "race@example.invalid", "GIT_COMMITTER_EMAIL": "race@example.invalid",
|
})
|
new_head = subprocess.run(
|
["git", "commit-tree", tree, "-p", old_head, "-m", "same tree local race"],
|
cwd=fixture.root, check=True, text=True, capture_output=True, env=env,
|
).stdout.strip()
|
subprocess.run(["git", "update-ref", "HEAD", new_head, old_head], cwd=fixture.root, check=True, capture_output=True)
|
elif (
|
not injected
|
and mode == "remote_after_read"
|
and command[:3] == ["git", "rev-parse", "refs/remotes/origin/main"]
|
):
|
injected = True
|
old_remote = result.stdout.strip()
|
parent = subprocess.run(
|
["git", "rev-parse", "HEAD^"], cwd=fixture.root, check=True, text=True, capture_output=True,
|
).stdout.strip()
|
subprocess.run(
|
["git", "update-ref", "refs/remotes/origin/main", parent, old_remote],
|
cwd=fixture.root, check=True, capture_output=True,
|
)
|
return result
|
|
with self.assertRaises(product.PipelineError) as drift:
|
product.git_deliver(fixture.config, outbox_id, NOW, runner=runner)
|
self.assertTrue(injected)
|
self.assertEqual(expected_code, drift.exception.code)
|
self.assertEqual(outbox_before, fixture.config.outbox_path.read_bytes())
|
self.assertEqual(0, len(list(fixture.config.state_dir.glob(f"git-index-{outbox_id}*"))))
|
self.assertEqual(0, len(list((fixture.root / ".git").rglob("*.lock"))))
|
|
def test_append_only_rewrite_and_unsafe_artifact_fail_closed(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
fixture.formal.write_bytes(b"")
|
with self.assertRaises(product.PipelineError) as rewritten:
|
product.reconcile(fixture.config, NOW)
|
self.assertEqual("E_HISTORY_REWRITE", rewritten.exception.code)
|
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
original = fixture.formal.read_bytes()
|
replacement = original.replace(b'"stable_id":"OLD"', b'"stable_id":"NEW"')
|
self.assertEqual(len(original), len(replacement))
|
fixture.formal.write_bytes(replacement)
|
with self.assertRaises(product.PipelineError) as prefix:
|
product.reconcile(fixture.config, NOW)
|
self.assertEqual("E_HISTORY_REWRITE", prefix.exception.code)
|
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
outside = fixture.root / "outside.txt"
|
outside.write_bytes(b"outside")
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": "123", "item_type": "text", "status": "SAVED", "creator_uid": UID,
|
"path": str(outside), "bytes": outside.stat().st_size, "sha256": hashlib.sha256(outside.read_bytes()).hexdigest().upper()
|
})
|
with self.assertRaises(product.PipelineError) as escaped:
|
product.reconcile(fixture.config, NOW)
|
self.assertEqual("E_ARTIFACT", escaped.exception.code)
|
|
def test_git_commit_journal_recovers_push_without_second_commit(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
(fixture.root / ".git").mkdir()
|
(fixture.root / ".git" / "index").write_bytes(b"USER-INDEX")
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
artifact = fixture.archive / "recover.txt"
|
artifact.write_bytes(b"recover\n")
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": "recover", "item_type": "text", "status": "SAVED",
|
"creator_uid": UID, "path": "recover.txt", "bytes": artifact.stat().st_size,
|
"sha256": hashlib.sha256(artifact.read_bytes()).hexdigest().upper(),
|
})
|
outbox_id = product.reconcile(fixture.config, NOW)["created_outbox_ids"][0]
|
rel = artifact.relative_to(fixture.root).as_posix()
|
commands: list[list[str]] = []
|
push_count = 0
|
|
def runner(command: list[str], **_: object) -> subprocess.CompletedProcess[str]:
|
nonlocal push_count
|
commands.append(command)
|
if command[:5] == ["git", "diff", "--cached", "--name-only", "-z"] and command[5] != "--":
|
return subprocess.CompletedProcess(command, 0, b"", b"")
|
if command[:6] == ["git", "diff", "--cached", "--name-only", "-z", "--"]:
|
return subprocess.CompletedProcess(command, 0, rel.encode("utf-8") + b"\0", b"")
|
if command[:4] == ["git", "diff", "--cached", "--name-only"]:
|
return subprocess.CompletedProcess(command, 0, "", "")
|
if command[:3] == ["git", "rev-parse", "HEAD"]:
|
return subprocess.CompletedProcess(command, 0, ("a" if push_count == 0 else "c") * 40 + "\n", "")
|
if command[:2] == ["git", "hash-object"]:
|
return subprocess.CompletedProcess(command, 0, "d" * 40 + "\n", "")
|
if command[:2] == ["git", "write-tree"]:
|
return subprocess.CompletedProcess(command, 0, "b" * 40 + "\n", "")
|
if command[:2] == ["git", "commit-tree"]:
|
return subprocess.CompletedProcess(command, 0, "c" * 40 + "\n", "")
|
if command[:2] == ["git", "push"]:
|
push_count += 1
|
return subprocess.CompletedProcess(command, 1 if push_count == 1 else 0, "", "")
|
return subprocess.CompletedProcess(command, 0, "", "")
|
|
with self.assertRaises(product.PipelineError) as first:
|
product.git_deliver(fixture.config, outbox_id, NOW, runner=runner)
|
self.assertEqual("E_GIT_PUSH", first.exception.code)
|
self.assertEqual("GIT_PUSHED", product.git_deliver(fixture.config, outbox_id, NOW, runner=runner)["status"])
|
self.assertEqual(1, sum(1 for command in commands if command[:2] == ["git", "commit-tree"]))
|
self.assertEqual(1, sum(1 for command in commands if command[:2] == ["git", "update-index"]))
|
|
def test_payload_source_rebinding_and_signed_url_reject_before_dispatch(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": BVID, "item_type": "video", "status": "VIDEO_DOWNLOAD_PENDING_EXTENSION",
|
"creator_uid": UID, "source_url": f"https://www.bilibili.com/video/{BVID}?token=SYNTHETIC_SECRET",
|
"title": "fixture", "published_at": NOW.isoformat(), "expected_duration_seconds": 60,
|
})
|
with self.assertRaises(product.PipelineError) as secret:
|
product.reconcile(fixture.config, NOW)
|
self.assertEqual("E_SECRET_FIELD", secret.exception.code)
|
self.assertFalse(fixture.config.outbox_path.exists())
|
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
self._new_handoff(fixture)
|
outbox_id = product.reconcile(fixture.config, NOW)["created_outbox_ids"][0]
|
rows = [json.loads(line) for line in fixture.config.outbox_path.read_text(encoding="ascii").splitlines()]
|
rows[0]["payload"]["bvid"] = "BV1Q541167Qf"
|
fixture.config.outbox_path.write_bytes(b"".join(canonical(row) for row in rows))
|
before = fixture.config.outbox_path.read_bytes()
|
with self.assertRaises(product.PipelineError) as tampered:
|
product.dispatch_intent(fixture.config, outbox_id, NOW)
|
self.assertEqual("E_OUTBOX", tampered.exception.code)
|
self.assertEqual(before, fixture.config.outbox_path.read_bytes())
|
|
def test_legacy_audit_handoff_identity_key_is_narrowly_projected(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": BVID, "item_type": "video",
|
"status": "VIDEO_DOWNLOAD_PENDING_EXTENSION", "creator_uid": UID,
|
"source_url": f"https://www.bilibili.com/video/{BVID}",
|
"title": "fixture", "published_at": NOW.isoformat(),
|
"expected_duration_seconds": 60,
|
"runtime_authorization_handoff_id": "HANDOFF-LEGACY-AUDIT-IDENTITY",
|
})
|
outbox_id = product.reconcile(fixture.config, NOW)["created_outbox_ids"][0]
|
pending = product.pending(fixture.config)
|
self.assertEqual((1, outbox_id), (pending["count"], pending["items"][0]["outbox_id"]))
|
serialized = fixture.config.outbox_path.read_text(encoding="ascii").lower()
|
self.assertNotIn("authorization", serialized)
|
self.assertEqual(
|
"DISPATCH_INTENT_DURABLE",
|
product.dispatch_intent(fixture.config, outbox_id, NOW)["status"],
|
)
|
|
def test_receipt_requires_observed_and_late_observed_is_mutation_zero(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
self._new_handoff(fixture)
|
outbox_id = product.reconcile(fixture.config, NOW)["created_outbox_ids"][0]
|
product.dispatch_intent(fixture.config, outbox_id, NOW)
|
transcript = fixture.archive / "ordered.transcript.txt"
|
transcript.write_bytes(b"ordered\n")
|
receipt = fixture.root / "receipt.json"
|
receipt.write_bytes(canonical({
|
"schema_version": 1, "type": "TRANSCRIPTION_COMPLETE", "outbox_id": outbox_id,
|
"stable_id": BVID, "terminal_id": "TRANSCRIPT-ORDER", "status": "COMPLETE",
|
"files": [{"path": str(transcript), "bytes": transcript.stat().st_size,
|
"sha256": hashlib.sha256(transcript.read_bytes()).hexdigest().upper(), "kind": "transcript"}],
|
"created_at": NOW.isoformat(),
|
}))
|
outbox_before = fixture.config.outbox_path.read_bytes()
|
with self.assertRaises(product.PipelineError) as early:
|
product.ingest_receipt(fixture.config, receipt, NOW)
|
self.assertEqual("E_RECEIPT", early.exception.code)
|
self.assertEqual(outbox_before, fixture.config.outbox_path.read_bytes())
|
self.assertFalse(fixture.config.terminals_path.exists())
|
product._append(fixture.config.outbox_path, {
|
"schema_version": 1, "event": "COMPLETE", "outbox_id": outbox_id,
|
"result": "TRANSCRIPTION_COMPLETE", "terminal_id": "TRANSCRIPT-ORDER",
|
"receipt_sha256": hashlib.sha256(receipt.read_bytes()).hexdigest().upper(), "completed_at": NOW.isoformat(),
|
})
|
malformed = fixture.config.outbox_path.read_bytes()
|
with self.assertRaises(product.PipelineError) as late:
|
product.observe_dispatch(fixture.config, outbox_id, "DELIVERY-LATE", NOW)
|
self.assertEqual("E_OUTBOX", late.exception.code)
|
self.assertEqual(malformed, fixture.config.outbox_path.read_bytes())
|
|
def test_git_commit_failure_uses_only_task_index_and_retry_is_deterministic(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
(fixture.root / ".git").mkdir()
|
shared_index = fixture.root / ".git" / "index"
|
shared_index.write_bytes(b"USER-INDEX")
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
artifact = fixture.archive / "failure.txt"
|
artifact.write_bytes(b"failure\n")
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": "failure", "item_type": "text", "status": "SAVED", "creator_uid": UID,
|
"path": "failure.txt", "bytes": artifact.stat().st_size, "sha256": hashlib.sha256(artifact.read_bytes()).hexdigest().upper(),
|
})
|
outbox_id = product.reconcile(fixture.config, NOW)["created_outbox_ids"][0]
|
rel = artifact.relative_to(fixture.root).as_posix()
|
calls: list[tuple[list[str], dict[str, object]]] = []
|
commit_attempt = 0
|
|
def runner(command: list[str], **kwargs: object) -> subprocess.CompletedProcess[str]:
|
nonlocal commit_attempt
|
calls.append((command, kwargs))
|
if command[:5] == ["git", "diff", "--cached", "--name-only", "-z"] and command[5] != "--":
|
return subprocess.CompletedProcess(command, 0, b"", b"")
|
if command[:6] == ["git", "diff", "--cached", "--name-only", "-z", "--"]:
|
return subprocess.CompletedProcess(command, 0, rel.encode("utf-8") + b"\0", b"")
|
if command[:4] == ["git", "diff", "--cached", "--name-only"]:
|
return subprocess.CompletedProcess(command, 0, "", "")
|
if command[:3] == ["git", "rev-parse", "HEAD"]:
|
return subprocess.CompletedProcess(command, 0, "a" * 40 + "\n", "")
|
if command[:2] == ["git", "hash-object"]:
|
return subprocess.CompletedProcess(command, 0, "d" * 40 + "\n", "")
|
if command[:2] == ["git", "write-tree"]:
|
return subprocess.CompletedProcess(command, 0, "b" * 40 + "\n", "")
|
if command[:2] == ["git", "commit-tree"]:
|
commit_attempt += 1
|
return subprocess.CompletedProcess(command, 1 if commit_attempt == 1 else 0, "" if commit_attempt == 1 else "c" * 40 + "\n", "")
|
return subprocess.CompletedProcess(command, 0, "", "")
|
|
with self.assertRaises(product.PipelineError) as failed:
|
product.git_deliver(fixture.config, outbox_id, NOW, runner=runner)
|
self.assertEqual("E_GIT_COMMIT", failed.exception.code)
|
self.assertEqual(b"USER-INDEX", shared_index.read_bytes())
|
self.assertEqual(1, fixture.config.outbox_path.read_text(encoding="ascii").count('"event":"GIT_COMMIT_INTENT"'))
|
self.assertEqual("GIT_PUSHED", product.git_deliver(fixture.config, outbox_id, NOW, runner=runner)["status"])
|
self.assertEqual(b"USER-INDEX", shared_index.read_bytes())
|
self.assertEqual(1, sum(command[:2] == ["git", "update-index"] for command, _ in calls))
|
task_envs = [kwargs["env"] for command, kwargs in calls if command[:2] == ["git", "update-index"]]
|
self.assertTrue(all(str(env["GIT_INDEX_FILE"]).startswith(str(fixture.config.state_dir)) for env in task_envs))
|
|
def test_real_git_commit_failure_preserves_shared_index_and_exact_retry_pushes(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
for command in (
|
["git", "init", "-b", "main"],
|
["git", "config", "user.name", "Fixture"],
|
["git", "config", "user.email", "fixture@example.invalid"],
|
["git", "add", "."],
|
["git", "commit", "-m", "fixture baseline"],
|
):
|
subprocess.run(command, cwd=fixture.root, check=True, capture_output=True)
|
remote = fixture.root / "remote.git"
|
subprocess.run(["git", "init", "--bare", str(remote)], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "remote", "add", "origin", str(remote)], cwd=fixture.root, check=True, capture_output=True)
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
artifact = fixture.archive / "真实-交付.txt"
|
artifact.write_bytes(b"real failure\n")
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": "real-failure", "item_type": "text", "status": "SAVED", "creator_uid": UID,
|
"path": artifact.name, "bytes": artifact.stat().st_size,
|
"sha256": hashlib.sha256(artifact.read_bytes()).hexdigest().upper(),
|
})
|
outbox_id = product.reconcile(fixture.config, NOW)["created_outbox_ids"][0]
|
shared_index = fixture.root / ".git" / "index"
|
index_before = shared_index.read_bytes()
|
failed_once = False
|
|
def runner(command: list[str], **kwargs: object) -> subprocess.CompletedProcess[str]:
|
nonlocal failed_once
|
if command[:2] == ["git", "commit-tree"] and not failed_once:
|
failed_once = True
|
return subprocess.CompletedProcess(command, 1, "", "injected")
|
return subprocess.run(command, **kwargs)
|
|
with self.assertRaises(product.PipelineError) as failure:
|
product.git_deliver(fixture.config, outbox_id, NOW, runner=runner)
|
self.assertEqual("E_GIT_COMMIT", failure.exception.code)
|
self.assertEqual(index_before, shared_index.read_bytes())
|
staged = subprocess.run(["git", "diff", "--cached", "--name-only"], cwd=fixture.root, check=True, text=True, capture_output=True)
|
self.assertEqual("", staged.stdout)
|
result = product.git_deliver(fixture.config, outbox_id, NOW, runner=runner)
|
self.assertEqual("GIT_PUSHED", result["status"])
|
self.assertEqual(index_before, shared_index.read_bytes())
|
remote_head = subprocess.run(
|
["git", "--git-dir", str(remote), "rev-parse", "refs/heads/main"], check=True, text=True, capture_output=True,
|
).stdout.strip()
|
self.assertEqual(result["commit_sha"], remote_head)
|
|
def test_real_git_batch_freezes_shared_index_across_moving_head_and_rejects_drift(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
for command in (
|
["git", "init", "-b", "main"],
|
["git", "config", "user.name", "Fixture"],
|
["git", "config", "user.email", "fixture@example.invalid"],
|
["git", "add", "."],
|
["git", "commit", "-m", "fixture baseline"],
|
):
|
subprocess.run(command, cwd=fixture.root, check=True, capture_output=True)
|
remote = fixture.root / "remote.git"
|
subprocess.run(["git", "init", "--bare", str(remote)], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "remote", "add", "origin", str(remote)], cwd=fixture.root, check=True, capture_output=True)
|
user_file = fixture.root / "user-staged.txt"
|
user_file.write_bytes(b"user staged preimage\n")
|
subprocess.run(["git", "add", user_file.name], cwd=fixture.root, check=True, capture_output=True)
|
shared_index = fixture.root / ".git" / "index"
|
index_before = shared_index.read_bytes()
|
baseline_head = subprocess.run(
|
["git", "rev-parse", "HEAD"], cwd=fixture.root, check=True, text=True, capture_output=True,
|
).stdout.strip()
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
for stable_id in ("batch-one", "batch-two"):
|
artifact = fixture.archive / f"{stable_id}.txt"
|
artifact.write_bytes((stable_id + "\n").encode("ascii"))
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": stable_id, "item_type": "text", "status": "SAVED",
|
"creator_uid": UID, "path": artifact.name, "bytes": artifact.stat().st_size,
|
"sha256": hashlib.sha256(artifact.read_bytes()).hexdigest().upper(),
|
})
|
outbox_ids = product.reconcile(fixture.config, NOW)["created_outbox_ids"]
|
self.assertEqual(2, len(outbox_ids))
|
first = product.git_deliver(fixture.config, outbox_ids[0], NOW)
|
self.assertEqual(index_before, shared_index.read_bytes())
|
interpreted_after_head_move = subprocess.run(
|
["git", "diff", "--cached", "--name-only"], cwd=fixture.root, check=True, text=True, capture_output=True,
|
).stdout.splitlines()
|
self.assertIn("user-staged.txt", interpreted_after_head_move)
|
guards = list(fixture.config.state_dir.glob("git-shared-index-guard-*.json"))
|
self.assertEqual(1, len(guards))
|
guards[0].unlink()
|
outbox_before_recovery = fixture.config.outbox_path.read_bytes()
|
with self.assertRaises(product.PipelineError) as wrong_frozen_bytes:
|
product.recover_git_index_guard(
|
fixture.config,
|
outbox_ids[1],
|
baseline_head,
|
len(index_before),
|
"0" * 64,
|
)
|
self.assertEqual("E_GIT_INDEX_DIRTY", wrong_frozen_bytes.exception.code)
|
self.assertEqual(outbox_before_recovery, fixture.config.outbox_path.read_bytes())
|
self.assertEqual([], list(fixture.config.state_dir.glob("git-shared-index-guard-*.json")))
|
recovered = product.recover_git_index_guard(
|
fixture.config,
|
outbox_ids[1],
|
baseline_head,
|
len(index_before),
|
hashlib.sha256(index_before).hexdigest().upper(),
|
)
|
self.assertEqual("GIT_INDEX_GUARD_RECOVERED", recovered["status"])
|
second = product.git_deliver(fixture.config, outbox_ids[1], NOW)
|
self.assertEqual(index_before, shared_index.read_bytes())
|
self.assertNotEqual(first["commit_sha"], second["commit_sha"])
|
self.assertEqual([], list(fixture.config.state_dir.glob("git-index-*")))
|
remote_head = subprocess.run(
|
["git", "--git-dir", str(remote), "rev-parse", "refs/heads/main"], check=True, text=True, capture_output=True,
|
).stdout.strip()
|
self.assertEqual(second["commit_sha"], remote_head)
|
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
for command in (
|
["git", "init", "-b", "main"],
|
["git", "config", "user.name", "Fixture"],
|
["git", "config", "user.email", "fixture@example.invalid"],
|
["git", "add", "."],
|
["git", "commit", "-m", "fixture baseline"],
|
):
|
subprocess.run(command, cwd=fixture.root, check=True, capture_output=True)
|
remote = fixture.root / "remote.git"
|
subprocess.run(["git", "init", "--bare", str(remote)], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "remote", "add", "origin", str(remote)], cwd=fixture.root, check=True, capture_output=True)
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
for stable_id in ("drift-one", "drift-two"):
|
artifact = fixture.archive / f"{stable_id}.txt"
|
artifact.write_bytes((stable_id + "\n").encode("ascii"))
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": stable_id, "item_type": "text", "status": "SAVED",
|
"creator_uid": UID, "path": artifact.name, "bytes": artifact.stat().st_size,
|
"sha256": hashlib.sha256(artifact.read_bytes()).hexdigest().upper(),
|
})
|
outbox_ids = product.reconcile(fixture.config, NOW)["created_outbox_ids"]
|
first = product.git_deliver(fixture.config, outbox_ids[0], NOW)
|
user_file = fixture.root / "late-user-stage.txt"
|
user_file.write_bytes(b"late user stage\n")
|
subprocess.run(["git", "add", user_file.name], cwd=fixture.root, check=True, capture_output=True)
|
outbox_before = fixture.config.outbox_path.read_bytes()
|
with self.assertRaises(product.PipelineError) as drift:
|
product.git_deliver(fixture.config, outbox_ids[1], NOW)
|
self.assertEqual("E_GIT_INDEX_DIRTY", drift.exception.code)
|
self.assertEqual(outbox_before, fixture.config.outbox_path.read_bytes())
|
self.assertEqual(first["commit_sha"], subprocess.run(
|
["git", "rev-parse", "HEAD"], cwd=fixture.root, check=True, text=True, capture_output=True,
|
).stdout.strip())
|
self.assertEqual([], list(fixture.config.state_dir.glob("git-index-*")))
|
|
def test_existing_git_guard_recovery_revalidates_guard_and_shared_index(self) -> None:
|
baseline_head = "a" * 40
|
|
def prepare(raw: str) -> tuple[PipelineFixture, Path, Path, str, bytes, dict[str, bytes]]:
|
fixture = PipelineFixture(Path(raw))
|
(fixture.root / ".git").mkdir()
|
shared_index = fixture.root / ".git" / "index"
|
index_bytes = b"INDEX-A"
|
shared_index.write_bytes(index_bytes)
|
product.initialize(fixture.config, NOW)
|
product.begin(fixture.config, NOW)
|
artifact = fixture.archive / "guarded.txt"
|
artifact.write_bytes(b"guarded\n")
|
fixture.append_formal({
|
"schema_version": 1, "stable_id": "guarded", "item_type": "text", "status": "SAVED",
|
"creator_uid": UID, "path": artifact.name, "bytes": artifact.stat().st_size,
|
"sha256": hashlib.sha256(artifact.read_bytes()).hexdigest().upper(),
|
})
|
outbox_id = product.reconcile(fixture.config, NOW)["created_outbox_ids"][0]
|
staged = {"value": b"user-staged.txt\0"}
|
|
def runner(command: list[str], **_: object) -> subprocess.CompletedProcess[object]:
|
if command[:3] == ["git", "rev-parse", "HEAD"]:
|
return subprocess.CompletedProcess(command, 0, baseline_head + "\n", "")
|
if command[:5] == ["git", "diff", "--cached", "--name-only", "-z"]:
|
return subprocess.CompletedProcess(command, 0, staged["value"], b"")
|
return subprocess.CompletedProcess(command, 0, "", "")
|
|
rows = product._outbox_rows(fixture.config)
|
item = next(row for row in rows if row.get("event") == "CREATED" and row.get("outbox_id") == outbox_id)
|
guard = product._ensure_git_index_guard(fixture.config, rows, item, runner)
|
guard_path = fixture.config.git_index_guard_path(guard["batch_id"])
|
fixture.guard_runner = runner # type: ignore[attr-defined]
|
return fixture, shared_index, guard_path, outbox_id, index_bytes, staged
|
|
with tempfile.TemporaryDirectory() as raw:
|
fixture, _, _, outbox_id, index_bytes, _ = prepare(raw)
|
result = product.recover_git_index_guard(
|
fixture.config, outbox_id, baseline_head, len(index_bytes), hashlib.sha256(index_bytes).hexdigest().upper(),
|
runner=fixture.guard_runner, # type: ignore[attr-defined]
|
)
|
self.assertEqual("GIT_INDEX_GUARD_ALREADY_DURABLE", result["status"])
|
|
for case in ("bytes", "staged_paths", "file_identity", "guard_field"):
|
with self.subTest(case=case), tempfile.TemporaryDirectory() as raw:
|
fixture, shared_index, guard_path, outbox_id, index_bytes, staged = prepare(raw)
|
if case == "bytes":
|
shared_index.write_bytes(b"INDEX-B")
|
elif case == "staged_paths":
|
staged["value"] = b"other-user-staged.txt\0"
|
elif case == "file_identity":
|
replacement = shared_index.with_name("index-replacement")
|
replacement.write_bytes(index_bytes)
|
os.replace(replacement, shared_index)
|
else:
|
guard = json.loads(guard_path.read_text(encoding="ascii"))
|
guard["outbox_ids"] = ["F" * 64]
|
guard_path.write_bytes(canonical(guard))
|
outbox_before = fixture.config.outbox_path.read_bytes()
|
with self.assertRaises(product.PipelineError) as rejected:
|
product.recover_git_index_guard(
|
fixture.config, outbox_id, baseline_head, len(index_bytes),
|
hashlib.sha256(index_bytes).hexdigest().upper(),
|
runner=fixture.guard_runner, # type: ignore[attr-defined]
|
)
|
self.assertEqual("E_GIT_INDEX_DIRTY", rejected.exception.code)
|
self.assertEqual(outbox_before, fixture.config.outbox_path.read_bytes())
|
|
for broken in (False, True):
|
with self.subTest(symlink_broken=broken), tempfile.TemporaryDirectory() as raw:
|
fixture, _, guard_path, outbox_id, index_bytes, _ = prepare(raw)
|
target = guard_path.with_name("guard-link-target.json")
|
if not broken:
|
target.write_bytes(guard_path.read_bytes())
|
guard_path.unlink()
|
try:
|
os.symlink(target, guard_path)
|
except OSError:
|
continue
|
outbox_before = fixture.config.outbox_path.read_bytes()
|
with self.assertRaises(product.PipelineError) as rejected:
|
product.recover_git_index_guard(
|
fixture.config, outbox_id, baseline_head, len(index_bytes),
|
hashlib.sha256(index_bytes).hexdigest().upper(),
|
runner=fixture.guard_runner, # type: ignore[attr-defined]
|
)
|
self.assertEqual("E_GIT_INDEX_DIRTY", rejected.exception.code)
|
self.assertEqual(outbox_before, fixture.config.outbox_path.read_bytes())
|
|
def test_full_chain_reparse_and_identity_drift_fail_closed(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
artifact = fixture.archive / "stable.txt"
|
artifact.write_bytes(b"stable\n")
|
digest = hashlib.sha256(artifact.read_bytes()).hexdigest().upper()
|
original = product._strict_chain
|
calls = 0
|
|
def drifting(root: Path, target: Path, *, final_file: bool) -> tuple[tuple[str, tuple[int, ...]], ...]:
|
nonlocal calls
|
calls += 1
|
value = original(root, target, final_file=final_file)
|
if calls == 2:
|
mutated = list(value)
|
name, identity = mutated[-2]
|
mutated[-2] = (name, identity[:-1] + (identity[-1] + 1,))
|
return tuple(mutated)
|
return value
|
|
with mock.patch.object(product, "_strict_chain", side_effect=drifting):
|
with self.assertRaises(product.PipelineError) as drift:
|
product._artifact("stable.txt", artifact.stat().st_size, digest, fixture.config)
|
self.assertEqual("E_ARTIFACT", drift.exception.code)
|
outside = fixture.root / "outside.txt"
|
outside.write_bytes(b"outside")
|
link = fixture.archive / "linked.txt"
|
try:
|
os.symlink(outside, link)
|
except OSError:
|
return
|
with self.assertRaises(product.PipelineError) as reparse:
|
product._artifact("linked.txt", outside.stat().st_size, hashlib.sha256(outside.read_bytes()).hexdigest().upper(), fixture.config)
|
self.assertEqual("E_ARTIFACT", reparse.exception.code)
|
|
def test_finish_allows_next_slot_and_cli_is_ascii_only(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
product.initialize(fixture.config, NOW)
|
first = product.begin(fixture.config, NOW)
|
finished = product.finish(fixture.config, "COMPLETE", NOW)
|
self.assertEqual(first["run"]["run_id"], finished["run_id"])
|
second = product.begin(fixture.config, datetime(2026, 8, 29, 8, 30, tzinfo=timezone.utc))
|
self.assertNotEqual(first["run"]["run_id"], second["run"]["run_id"])
|
code, result = product.run(["--config", str(fixture.config_path), "pending"])
|
self.assertEqual((0, "PENDING"), (code, result["status"]))
|
product._canonical(result).decode("ascii")
|
|
def test_git_preflight_detects_index_regression_without_mutating_shared_index(self) -> None:
|
class Interrupted(BaseException):
|
pass
|
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
subprocess.run(["git", "init", "-b", "main"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "config", "user.name", "Fixture"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "config", "user.email", "fixture@example.invalid"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "add", "."], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "commit", "-m", "baseline"], cwd=fixture.root, check=True, capture_output=True)
|
relative = (fixture.archive / "old.txt").relative_to(fixture.root).as_posix()
|
subprocess.run(["git", "rm", "--cached", "--", relative], cwd=fixture.root, check=True, capture_output=True)
|
index_path = fixture.root / ".git" / "index"
|
lock_path = fixture.root / ".git" / "index.lock"
|
lock_path.write_bytes(b"")
|
before = index_path.read_bytes()
|
before_stat = os.stat(index_path)
|
result = product.git_preflight(fixture.config)
|
self.assertFalse(result["index_matches_head"])
|
self.assertEqual([relative], result["archive_staged_delete_present"])
|
self.assertEqual((0, hashlib.sha256(b"").hexdigest().upper()), (
|
result["stale_index_lock"]["bytes"], result["stale_index_lock"]["sha256"]
|
))
|
self.assertEqual(before, index_path.read_bytes())
|
self.assertEqual((before_stat.st_dev, before_stat.st_ino), (os.stat(index_path).st_dev, os.stat(index_path).st_ino))
|
|
def fail_runner(*_args: object, **_kwargs: object) -> object:
|
raise RuntimeError("synthetic read failure")
|
|
with self.assertRaises(RuntimeError):
|
product.git_preflight(fixture.config, runner=fail_runner) # type: ignore[arg-type]
|
self.assertEqual(before, index_path.read_bytes())
|
|
def interrupt_runner(*_args: object, **_kwargs: object) -> object:
|
raise Interrupted()
|
|
with self.assertRaises(Interrupted):
|
product.git_preflight(fixture.config, runner=interrupt_runner) # type: ignore[arg-type]
|
self.assertEqual(before, index_path.read_bytes())
|
|
def test_git_preflight_uses_one_held_index_during_transient_path_swap(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
subprocess.run(["git", "init", "-b", "main"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "config", "user.name", "Fixture"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "config", "user.email", "fixture@example.invalid"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "add", "."], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "commit", "-m", "baseline"], cwd=fixture.root, check=True, capture_output=True)
|
index_path = fixture.root / ".git" / "index"
|
alternate = fixture.root / ".git" / "alternate-index"
|
parked = fixture.root / ".git" / "original-index-held-test"
|
shutil.copyfile(index_path, alternate)
|
relative = (fixture.archive / "old.txt").relative_to(fixture.root).as_posix()
|
alternate_env = {**os.environ, "GIT_INDEX_FILE": str(alternate)}
|
subprocess.run(
|
["git", "rm", "--cached", "--", relative],
|
cwd=fixture.root,
|
env=alternate_env,
|
check=True,
|
capture_output=True,
|
)
|
alternate_staged = subprocess.run(
|
["git", "diff", "--cached", "--name-only", "HEAD", "--"],
|
cwd=fixture.root,
|
env=alternate_env,
|
check=True,
|
text=True,
|
capture_output=True,
|
).stdout.splitlines()
|
self.assertEqual([relative], alternate_staged)
|
before = index_path.read_bytes()
|
before_stat = os.stat(index_path)
|
race = {"attempted": False, "performed": False, "denied": False}
|
|
def swap_runner(command: list[str], **kwargs: object) -> subprocess.CompletedProcess[object]:
|
if not race["attempted"] and command[1:2] == ["diff"]:
|
race["attempted"] = True
|
try:
|
os.replace(index_path, parked)
|
os.replace(alternate, index_path)
|
race["performed"] = True
|
except OSError:
|
race["denied"] = True
|
try:
|
return subprocess.run(command, **kwargs) # type: ignore[arg-type]
|
finally:
|
if race["performed"]:
|
os.replace(index_path, alternate)
|
os.replace(parked, index_path)
|
return subprocess.run(command, **kwargs) # type: ignore[arg-type]
|
|
result = product.git_preflight(fixture.config, runner=swap_runner) # type: ignore[arg-type]
|
self.assertTrue(race["attempted"])
|
self.assertTrue(race["performed"] or race["denied"])
|
self.assertTrue(result["index_matches_head"])
|
self.assertEqual([], result["staged_paths"])
|
self.assertEqual([], result["archive_staged_delete_present"])
|
self.assertEqual(before, index_path.read_bytes())
|
self.assertEqual(
|
(before_stat.st_dev, before_stat.st_ino),
|
(os.stat(index_path).st_dev, os.stat(index_path).st_ino),
|
)
|
|
def test_migration_restarts_from_one_durable_report_at_every_write_boundary(self) -> None:
|
class Interrupted(BaseException):
|
pass
|
|
boundaries = [
|
*(f"move-{index}" for index in range(0, 7)),
|
"rmdir-1",
|
"rmdir-2",
|
"outbox",
|
]
|
for boundary in boundaries:
|
with self.subTest(boundary=boundary), tempfile.TemporaryDirectory() as raw:
|
fixture, legacy_paths, payloads, flac = prepare_migration_fixture(Path(raw))
|
initial_plan = product.plan_video_artifact_migration(fixture.config, NOW)
|
original_move = product._move_relocation_file
|
original_rmdir = Path.rmdir
|
original_append = product._append_outbox
|
move_calls = 0
|
rmdir_calls = 0
|
|
def interrupted_move(*args: object, **kwargs: object) -> None:
|
nonlocal move_calls
|
if boundary == "move-0" and move_calls == 0:
|
raise Interrupted()
|
original_move(*args, **kwargs)
|
move_calls += 1
|
if boundary == f"move-{move_calls}":
|
raise Interrupted()
|
|
def interrupted_rmdir(path: Path) -> None:
|
nonlocal rmdir_calls
|
original_rmdir(path)
|
if path.name in {f"{BVID}.transcript", f"{BVID}.minutes"}:
|
rmdir_calls += 1
|
if boundary == f"rmdir-{rmdir_calls}":
|
raise Interrupted()
|
|
def interrupted_append(*args: object, **kwargs: object) -> str:
|
value = original_append(*args, **kwargs)
|
if boundary == "outbox":
|
raise Interrupted()
|
return value
|
|
with (
|
mock.patch.object(product, "_move_relocation_file", side_effect=interrupted_move),
|
mock.patch.object(Path, "rmdir", new=interrupted_rmdir),
|
mock.patch.object(product, "_append_outbox", side_effect=interrupted_append),
|
):
|
with self.assertRaises(Interrupted):
|
product.migrate_video_artifacts(fixture.config, NOW)
|
reports = list(fixture.config.relocation_root.glob("*.json"))
|
self.assertEqual(1, len(reports))
|
self.assertEqual(initial_plan, product._read_relocation_report(fixture.config, reports[0]))
|
|
completed = product.migrate_video_artifacts(
|
fixture.config,
|
datetime(2026, 8, 29, 8, 1, tzinfo=timezone.utc),
|
)
|
self.assertEqual(initial_plan["batch_id"], completed["batch_id"])
|
self.assertEqual(initial_plan["created_at"], product._read_relocation_report(fixture.config, reports[0])["created_at"])
|
rows = product._outbox_rows(fixture.config)
|
created = [
|
row for row in rows
|
if row.get("event") == "CREATED" and row.get("outbox_id") == completed["git_outbox_id"]
|
]
|
self.assertEqual(1, len(created))
|
for alias in initial_plan["items"][0]["aliases"]:
|
self.assertFalse((fixture.root / alias["old_path"]).exists())
|
self.assertEqual(payloads[alias["kind"]], (fixture.root / alias["new_path"]).read_bytes())
|
target_flac = fixture.root / "external-video" / "intermediate" / "transcription" / f"{BVID}.audio.flac"
|
self.assertFalse(flac.exists())
|
self.assertEqual(b"fLaC-fixture", target_flac.read_bytes())
|
self.assertFalse((fixture.archive / f"{BVID}.transcript").exists())
|
self.assertFalse((fixture.archive / f"{BVID}.minutes").exists())
|
self.assertTrue(all(not path.exists() for path in legacy_paths.values()))
|
|
def test_relocation_source_replacement_cannot_delete_replacement_bytes(self) -> None:
|
source_tmp_root = PROJECT_ROOT / "dev" / "tmp"
|
source_tmp_root.mkdir(parents=True, exist_ok=True)
|
target_parents = [source_tmp_root]
|
if Path("F:/").is_dir():
|
target_parents.append(Path("F:/"))
|
for target_parent in target_parents:
|
with (
|
self.subTest(target_parent=str(target_parent)),
|
tempfile.TemporaryDirectory(dir=source_tmp_root) as source_raw,
|
tempfile.TemporaryDirectory(dir=target_parent) as target_raw,
|
):
|
source_root = Path(source_raw)
|
target_root = Path(target_raw)
|
source = source_root / "source.bin"
|
target = target_root / "target.bin"
|
parked = source_root / "parked-original.bin"
|
source.write_bytes(b"AAAA")
|
original_delete = product._delete_held_relocation_source
|
race = {"swapped": False, "denied": False}
|
|
def replace_before_delete(*args: object, **kwargs: object) -> None:
|
try:
|
os.replace(source, parked)
|
source.write_bytes(b"BBBB")
|
race["swapped"] = True
|
except OSError:
|
race["denied"] = True
|
original_delete(*args, **kwargs)
|
|
if os.name == "nt":
|
with mock.patch.object(product, "_delete_held_relocation_source", side_effect=replace_before_delete):
|
product._move_relocation_file(
|
source_root,
|
target_root,
|
source,
|
target,
|
4,
|
hashlib.sha256(b"AAAA").hexdigest().upper(),
|
)
|
self.assertTrue(race["denied"])
|
self.assertFalse(race["swapped"])
|
self.assertFalse(source.exists())
|
else:
|
with (
|
mock.patch.object(product, "_delete_held_relocation_source", side_effect=replace_before_delete),
|
self.assertRaises(product.PipelineError) as rejected,
|
):
|
product._move_relocation_file(
|
source_root,
|
target_root,
|
source,
|
target,
|
4,
|
hashlib.sha256(b"AAAA").hexdigest().upper(),
|
)
|
self.assertEqual("E_RELOCATION", rejected.exception.code)
|
self.assertTrue(race["swapped"])
|
self.assertEqual(b"BBBB", source.read_bytes())
|
self.assertEqual(b"AAAA", parked.read_bytes())
|
self.assertEqual(b"AAAA", target.read_bytes())
|
|
def test_relocation_report_is_rebound_to_the_exact_external_plan(self) -> None:
|
mutations = (
|
"missing_intermediate", "missing_remove", "cross_path", "extra_alias",
|
"omitted_legacy_id", "omitted_canonical_id",
|
)
|
for mutation in mutations:
|
with self.subTest(mutation=mutation), tempfile.TemporaryDirectory() as raw:
|
fixture, legacy_paths, payloads, flac = prepare_migration_fixture(Path(raw))
|
report = json.loads(json.dumps(product.plan_video_artifact_migration(fixture.config, NOW)))
|
if mutation == "missing_intermediate":
|
report["items"][0]["intermediates"] = []
|
elif mutation == "missing_remove":
|
report["remove_paths"] = report["remove_paths"][1:]
|
elif mutation == "cross_path":
|
report["items"][0]["aliases"][0]["old_path"] = fixture.formal.relative_to(fixture.root).as_posix()
|
elif mutation == "extra_alias":
|
report["items"][0]["aliases"].append(dict(report["items"][0]["aliases"][0]))
|
else:
|
omitted = "BV1Q541167Qf"
|
omitted_title = "Omitted legacy item"
|
omitted_published = "2026-08-31T01:02:03+08:00"
|
fixture.append_formal({
|
"schema_version": 1,
|
"stable_id": omitted,
|
"item_type": "video",
|
"status": product.VIDEO_COMPLETE,
|
"creator": "fixture",
|
"title": omitted_title,
|
"published_at": omitted_published,
|
})
|
omitted_base = omitted if mutation == "omitted_legacy_id" else product._canonical_video_base(
|
omitted,
|
omitted_title,
|
omitted_published,
|
)
|
omitted_dir = fixture.archive / f"{omitted_base}.transcript"
|
omitted_dir.mkdir()
|
(omitted_dir / f"{omitted_base}.txt").write_bytes(b"omitted\n")
|
(omitted_dir / f"{omitted_base}.srt").write_bytes(b"1\n00:00:00,000 --> 00:00:01,000\nomitted\n")
|
(omitted_dir / f"{omitted_base}.json").write_bytes(b'{"segments":[]}\n')
|
report["batch_id"] = product._relocation_batch_id(report)
|
fixture.config.relocation_root.mkdir(parents=True)
|
forged_path = fixture.config.relocation_root / f"{report['batch_id']}.json"
|
forged_path.write_bytes(canonical(report))
|
state_before = fixture.config.state_path.read_bytes()
|
formal_before = fixture.formal.read_bytes()
|
with self.assertRaises(product.PipelineError) as rejected:
|
product.migrate_video_artifacts(fixture.config, NOW)
|
self.assertEqual("E_RELOCATION", rejected.exception.code)
|
self.assertEqual(state_before, fixture.config.state_path.read_bytes())
|
self.assertEqual(formal_before, fixture.formal.read_bytes())
|
self.assertFalse(fixture.config.outbox_path.exists())
|
self.assertEqual(b"fLaC-fixture", flac.read_bytes())
|
for kind, path in legacy_paths.items():
|
self.assertEqual(payloads[kind], path.read_bytes())
|
canonical_base = product._canonical_video_base(BVID, "A / canonical: title?", "2026-08-30T12:34:56+08:00")
|
self.assertFalse((fixture.archive / f"{canonical_base}.transcript").exists())
|
self.assertFalse((fixture.archive / f"{canonical_base}.minutes").exists())
|
|
def test_video_artifact_migration_is_byte_exact_alias_resolvable_and_git_scoped(self) -> None:
|
with tempfile.TemporaryDirectory() as raw:
|
fixture = PipelineFixture(Path(raw))
|
(fixture.root / "external-video").mkdir()
|
stable_id = BVID
|
title = "A / canonical: title?"
|
published_at = "2026-08-30T12:34:56+08:00"
|
fixture.append_formal({
|
"schema_version": 1,
|
"stable_id": stable_id,
|
"item_type": "video",
|
"status": product.VIDEO_COMPLETE,
|
"creator_uid": UID,
|
"title": title,
|
"published_at": published_at,
|
})
|
transcript_dir = fixture.archive / f"{stable_id}.transcript"
|
minutes_dir = fixture.archive / f"{stable_id}.minutes"
|
transcript_dir.mkdir()
|
minutes_dir.mkdir()
|
payloads = {
|
"transcript_txt": b"transcript\n",
|
"transcript_srt": b"1\n00:00:00,000 --> 00:00:01,000\ntext\n",
|
"transcript_json": b'{"segments":[]}\n',
|
"minutes_md": b"# minutes\n",
|
"minutes_pdf": b"%PDF-fixture\n",
|
}
|
legacy_paths = {
|
"transcript_txt": transcript_dir / f"{stable_id}.txt",
|
"transcript_srt": transcript_dir / f"{stable_id}.srt",
|
"transcript_json": transcript_dir / f"{stable_id}.json",
|
"minutes_md": minutes_dir / f"{stable_id}.md",
|
"minutes_pdf": minutes_dir / f"{stable_id}.pdf",
|
}
|
for kind, path in legacy_paths.items():
|
path.write_bytes(payloads[kind])
|
flac = transcript_dir / f"{stable_id}.audio.flac"
|
flac.write_bytes(b"fLaC-fixture")
|
|
subprocess.run(["git", "init", "-b", "main"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "config", "user.name", "Fixture"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "config", "user.email", "fixture@example.invalid"], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "add", "."], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "commit", "-m", "legacy baseline"], cwd=fixture.root, check=True, capture_output=True)
|
remote = fixture.root / "remote.git"
|
subprocess.run(["git", "init", "--bare", str(remote)], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "remote", "add", "origin", str(remote)], cwd=fixture.root, check=True, capture_output=True)
|
subprocess.run(["git", "push", "-u", "origin", "main"], cwd=fixture.root, check=True, capture_output=True)
|
(fixture.archive / "目录导读.md").write_bytes(b"# canonical naming guide\n")
|
product.initialize(fixture.config, NOW)
|
index_path = fixture.root / ".git" / "index"
|
index_before = index_path.read_bytes()
|
|
plan = product.plan_video_artifact_migration(fixture.config, NOW)
|
expected_base = "20260830-123456_video_A _ canonical_ title__BV1Q541167Qg"
|
self.assertEqual(expected_base, plan["items"][0]["canonical_base"])
|
result = product.migrate_video_artifacts(fixture.config, NOW)
|
self.assertEqual(("VIDEO_ARTIFACTS_MIGRATED", 1, 5, 1), (
|
result["status"], result["item_count"], result["public_file_count"], result["intermediate_count"]
|
))
|
self.assertEqual(index_before, index_path.read_bytes())
|
for alias in plan["items"][0]["aliases"]:
|
old_path = fixture.root / alias["old_path"]
|
new_path = fixture.root / alias["new_path"]
|
self.assertFalse(old_path.exists())
|
self.assertEqual(payloads[alias["kind"]], new_path.read_bytes())
|
rebound = product._artifact(alias["old_path"], alias["bytes"], alias["sha256"], fixture.config)
|
self.assertEqual(alias["old_path"], rebound["path"])
|
external_flac = fixture.root / "external-video" / "intermediate" / "transcription" / f"{stable_id}.audio.flac"
|
self.assertEqual(b"fLaC-fixture", external_flac.read_bytes())
|
self.assertFalse(flac.exists())
|
report_path = fixture.root / result["report"]["path"]
|
self.assertTrue(report_path.is_file())
|
rerun = product.migrate_video_artifacts(fixture.config, NOW)
|
self.assertEqual(result["git_outbox_id"], rerun["git_outbox_id"])
|
|
delivered = product.git_deliver(fixture.config, result["git_outbox_id"], NOW)
|
self.assertEqual("GIT_PUSHED", delivered["status"])
|
self.assertEqual(index_before, index_path.read_bytes())
|
changed = subprocess.run(
|
["git", "diff-tree", "--no-commit-id", "--name-status", "-r", "HEAD"],
|
cwd=fixture.root,
|
check=True,
|
text=True,
|
capture_output=True,
|
).stdout
|
for alias in plan["items"][0]["aliases"]:
|
self.assertIn(alias["old_path"], changed)
|
self.assertIn(alias["new_path"], changed)
|
self.assertNotIn(".flac", changed.lower())
|
local_head = subprocess.run(["git", "rev-parse", "HEAD"], cwd=fixture.root, check=True, text=True, capture_output=True).stdout.strip()
|
remote_head = subprocess.run(
|
["git", "--git-dir", str(remote), "rev-parse", "refs/heads/main"],
|
check=True,
|
text=True,
|
capture_output=True,
|
).stdout.strip()
|
self.assertEqual(local_head, remote_head)
|
|
|
if __name__ == "__main__":
|
unittest.main()
|