import csv
|
import hashlib
|
import io
|
import json
|
import os
|
import shutil
|
import subprocess
|
import sys
|
import tempfile
|
import time
|
import unittest
|
from pathlib import Path
|
|
|
PROJECT_ROOT = Path(__file__).resolve().parents[3]
|
MODULE_ROOT = PROJECT_ROOT / "dev" / "ana-dev"
|
PYTHON = sys.executable
|
GUARD_PATH = PROJECT_ROOT / "dev-doc" / "ana-doc" / "开发方案" / "SHARED-CONTENT-PUBLISHER-PRODUCTION-GUARD-V001.json"
|
LEDGER_HEADER = [
|
"release_id", "event_seq", "state", "prior_release_id",
|
"prior_release_set_sha256", "candidate_release_set_sha256", "task_id",
|
"case_id", "batch_id", "run_id", "accepted_audit_id", "operator",
|
"process_identity", "event_time", "recovery_or_rollback_receipt",
|
"attempt_id", "attempt_seq",
|
]
|
|
|
def long(path: Path) -> str:
|
text = os.path.abspath(path)
|
if os.name != "nt" or text.startswith("\\\\?\\"):
|
return text
|
if text.startswith("\\\\"):
|
return "\\\\?\\UNC\\" + text[2:]
|
return "\\\\?\\" + text
|
|
|
def write(path: Path, data: bytes) -> None:
|
os.makedirs(long(path.parent), exist_ok=True)
|
with open(long(path), "wb") as handle:
|
handle.write(data)
|
|
|
def read(path: Path) -> bytes:
|
with open(long(path), "rb") as handle:
|
return handle.read()
|
|
|
def ident(path: Path) -> tuple[int, str]:
|
data = read(path)
|
return len(data), hashlib.sha256(data).hexdigest().upper()
|
|
|
def sha(data: bytes) -> str:
|
return hashlib.sha256(data).hexdigest().upper()
|
|
|
def canonical(value, newline=True) -> bytes:
|
data = json.dumps(value, ensure_ascii=False, separators=(",", ":"), sort_keys=True).encode("utf-8")
|
return data + (b"\n" if newline else b"")
|
|
|
def release_set(data: bytes, formal: str) -> str:
|
descriptor = canonical({"bytes": len(data), "formal_relative_path": formal, "sha256": sha(data)}, newline=False)
|
return sha(data + descriptor)
|
|
|
def csv_bytes(header, rows) -> bytes:
|
stream = io.StringIO(newline="")
|
writer = csv.writer(stream, lineterminator="\n")
|
writer.writerow(header)
|
writer.writerows(rows)
|
return stream.getvalue().encode("utf-8")
|
|
|
def ledger_rows(data: bytes) -> list[dict[str, str]]:
|
reader = csv.DictReader(io.StringIO(data.decode("utf-8-sig"), newline=""))
|
if reader.fieldnames != LEDGER_HEADER:
|
raise AssertionError(reader.fieldnames)
|
return list(reader)
|
|
|
def event_chain_digest(rows: list[dict[str, str]]) -> str:
|
payload = {
|
"columns": LEDGER_HEADER,
|
"rows": [[row[column] for column in LEDGER_HEADER] for row in rows],
|
"schema_version": "SHARED_CONTENT_PUBLISHER_ATTEMPT_EVENT_CHAIN_V1",
|
}
|
return sha(canonical(payload, newline=False))
|
|
|
class Fixture:
|
def __init__(self, root: Path, fault="NONE"):
|
self.root = root
|
self.industry = "DEMO-INDUSTRY"
|
self.task = "TASK-DEMO-PUBLISH"
|
self.case = "CASE-DEMO-PUBLISH"
|
self.batch = "BATCH-DEMO-PUBLISH"
|
self.run_id = "RUN-DEMO-PUBLISH"
|
self.attempt = "ATTEMPT-DEMO-000002"
|
self.audit = "AUDIT-DEMO-PUBLISH-PASS"
|
self.handoff = "HANDOFF-DEMO-PUBLISH-REVIEW"
|
self.prior_release = "b" * 64
|
preimage = "".join(f"{key}={value}\n" for key, value in (
|
("industry", self.industry), ("task_id", self.task), ("case_id", self.case),
|
("batch_id", self.batch), ("run_id", self.run_id), ("canonical_audit_id", self.audit),
|
))
|
self.release = hashlib.sha256(preimage.encode()).hexdigest()
|
self.case_root = root / "ana-data" / "cases" / "Demo案例"
|
self.result_root = root / "ana-data" / "result" / "Demo案例"
|
self.manifest_root = self.case_root / "manifest"
|
self.release_store = self.case_root / "核心文档" / ".releases"
|
self.target = self.release_store / self.release
|
self.staging = self.release_store / (".staging-" + self.attempt)
|
self.candidate = root / "ana-data" / "tmp" / "Demo案例" / "candidate"
|
self.history = self.case_root / "审计包" / "TASK-DEMO" / "attempts" / "ATTEMPT-DEMO-000001"
|
self.receipt = self.case_root / "审计包" / "TASK-DEMO" / "attempts" / self.attempt
|
self.case_index = self.case_root / "当前成果索引.md"
|
self.result_index = self.result_root / "当前成果索引.md"
|
self.ledger = self.manifest_root / "promotion_ledger.csv"
|
self.lock = self.manifest_root / ".promotion.lock"
|
for directory in (self.case_root, self.result_root, self.manifest_root, self.release_store, self.candidate, self.history):
|
os.makedirs(long(directory), exist_ok=True)
|
|
self.prior_index_data = b"# Demo current\n\n- release: prior\n- core: [core](core/current.md)\n"
|
self.result_data = b"# Demo result current\n\n- cases: ../../cases/Demo\xe6\xa1\x88\xe4\xbe\x8b/\xe5\xbd\x93\xe5\x89\x8d\xe6\x88\x90\xe6\x9e\x9c\xe7\xb4\xa2\xe5\xbc\x95.md\n"
|
self.prior_core_data = b"# Prior core\n\nprior-evidence\n"
|
write(self.case_index, self.prior_index_data)
|
write(self.result_index, self.result_data)
|
write(self.case_root / "core" / "current.md", self.prior_core_data)
|
|
prior_rows_data = [
|
("CASES-CURRENT", self.rel(self.case_index), "CASES_CURRENT_INDEX", *ident(self.case_index)),
|
("CORE-PRIOR", self.rel(self.case_root / "core" / "current.md"), "CORE_DOCUMENT", *ident(self.case_root / "core" / "current.md")),
|
("RESULT-CURRENT", self.rel(self.result_index), "RESULT_THIN_CURRENT_INDEX", *ident(self.result_index)),
|
]
|
prior_manifest = csv_bytes(
|
["member_id", "formal_relative_path", "artifact_type", "bytes", "sha256"],
|
prior_rows_data,
|
)
|
self.prior_manifest_path = self.manifest_root / "current_output_manifest.csv"
|
write(self.prior_manifest_path, prior_manifest)
|
self.prior_manifest_formal = self.rel(self.prior_manifest_path)
|
self.prior_release_set = release_set(prior_manifest, self.prior_manifest_formal)
|
self.prior_rows = [
|
{"member_id": row[0], "formal_relative_path": row[1], "artifact_type": row[2], "bytes": row[3], "sha256": row[4]}
|
for row in prior_rows_data
|
]
|
|
state = canonical({"semantic_state": "RECOVERY_REQUIRED", "status": "BLOCKED", "exit_code": 30})
|
self.state_rel = "rollback_current_validation.json"
|
write(self.history / self.state_rel, state)
|
self.history_rows = [self.history_row(self.state_rel)]
|
for index in range(50):
|
rel = f"short/item-{index:02d}.json"
|
write(self.history / rel, canonical({"index": index}))
|
self.history_rows.append(self.history_row(rel))
|
for index in range(12):
|
parts = [f"long-{index:02d}"] + [(f"segment-{level:02d}-" + "x" * 22) for level in range(7)] + ["receipt.json"]
|
rel = "/".join(parts)
|
write(self.history / rel, canonical({"long": index}))
|
self.history_rows.append(self.history_row(rel))
|
assert len(self.history_rows) == 63
|
assert sum(len(os.path.abspath(self.history / row["relative_path"])) >= 260 for row in self.history_rows) >= 12
|
|
self.target_rel = self.rel(self.target)
|
self.core_source_rel = "release/core/research.md"
|
self.current_manifest_source_rel = "release/manifest/current_output_manifest.csv"
|
self.index_source_rel = "entry/当前成果索引.md"
|
candidate_core = b"# Candidate core\n\nsource-evidence-closed\n"
|
candidate_index = (f"# Demo current\n\n- release: {self.release}\n- core: [core](核心文档/.releases/{self.release}/core/research.md)\n- evidence: source-evidence-closed\n").encode("utf-8")
|
write(self.candidate / self.core_source_rel, candidate_core)
|
write(self.candidate / self.index_source_rel, candidate_index)
|
core_id = ident(self.candidate / self.core_source_rel)
|
index_id = ident(self.candidate / self.index_source_rel)
|
result_id = ident(self.result_index)
|
self.current_manifest_rows = [
|
("CORE-CANDIDATE", f"{self.target_rel}/core/research.md", "CORE_DOCUMENT", *core_id),
|
("CASES-CURRENT", self.rel(self.case_index), "CASES_CURRENT_INDEX", *index_id),
|
("RESULT-CURRENT", self.rel(self.result_index), "RESULT_THIN_CURRENT_INDEX", *result_id),
|
]
|
current_manifest = csv_bytes(
|
["member_id", "formal_relative_path", "artifact_type", "bytes", "sha256"],
|
self.current_manifest_rows,
|
)
|
write(self.candidate / self.current_manifest_source_rel, current_manifest)
|
current_manifest_id = ident(self.candidate / self.current_manifest_source_rel)
|
self.current_manifest_formal = f"{self.target_rel}/manifest/current_output_manifest.csv"
|
self.candidate_rows = [
|
self.artifact_row("CORE-CANDIDATE", self.core_source_rel, f"{self.target_rel}/core/research.md", "CORE_DOCUMENT", "RELEASE_MEMBER"),
|
self.artifact_row("CURRENT-MANIFEST", self.current_manifest_source_rel, self.current_manifest_formal, "CURRENT_OUTPUT_MANIFEST", "RELEASE_MEMBER"),
|
self.artifact_row("CASES-CURRENT", self.index_source_rel, self.rel(self.case_index), "CASES_CURRENT_INDEX", "CASE_CURRENT_INDEX"),
|
]
|
candidate_manifest = csv_bytes(
|
["member_id", "candidate_relative_path", "formal_relative_path", "artifact_type", "bytes", "sha256"],
|
[(row["member_id"], "candidate/" + row["relative_path"], row["formal_relative_path"], row["artifact_type"], row["bytes"], row["sha256"]) for row in self.candidate_rows],
|
)
|
self.candidate_manifest_rel = "candidate_output_manifest.csv"
|
write(self.candidate / self.candidate_manifest_rel, candidate_manifest)
|
candidate_manifest_id = ident(self.candidate / self.candidate_manifest_rel)
|
self.candidate_set = release_set(current_manifest, self.current_manifest_formal)
|
|
ledger = csv_bytes(LEDGER_HEADER, [[
|
"old-release", "1", "ROLLED_BACK", self.prior_release, self.prior_release_set, "C" * 64,
|
self.task, self.case, self.batch, self.run_id, "AUDIT-OLD", "tester", "PID-OLD-INSTANCE",
|
"2026-01-01T00:00:00Z", "receipts/old.json", "ATTEMPT-DEMO-000001", "1",
|
]])
|
write(self.ledger, ledger)
|
|
state_id = ident(self.history / self.state_rel)
|
prior_manifest_id = ident(self.prior_manifest_path)
|
self.config = {
|
"schema_version": "SHARED_CONTENT_PUBLISHER_CONFIG_V2",
|
"identity": {
|
"industry": self.industry, "task_id": self.task, "case_id": self.case,
|
"batch_id": self.batch, "run_id": self.run_id, "attempt_id": self.attempt,
|
"canonical_audit_id": self.audit, "review_handoff_id": self.handoff,
|
"release_id": self.release, "prior_release_id": self.prior_release,
|
"candidate_release_set_sha256": self.candidate_set, "operator": "test.operator",
|
},
|
"roots": {
|
"operation_root": str(self.root), "resolved_root": str(self.root),
|
"volume_identity": (self.root.drive or str(os.stat(self.root).st_dev)).upper(),
|
"candidate_root": self.rel(self.candidate), "release_store_root": self.rel(self.release_store),
|
"case_current_index_path": self.rel(self.case_index), "result_current_index_path": self.rel(self.result_index),
|
"ledger_path": self.rel(self.ledger), "lock_path": self.rel(self.lock),
|
"attempt_receipt_root": self.rel(self.receipt), "history_root": self.rel(self.history),
|
"current_state_relative_path": self.state_rel,
|
},
|
"expected": {
|
"ledger": {"bytes": len(ledger), "sha256": sha(ledger), "next_event_seq": 2, "next_attempt_seq": 2, "prior_attempt_id": "ATTEMPT-DEMO-000001", "prior_terminal_state": "ROLLED_BACK"},
|
"candidate": {
|
"manifest_relative_path": self.candidate_manifest_rel, "manifest_bytes": candidate_manifest_id[0], "manifest_sha256": candidate_manifest_id[1],
|
"current_manifest_relative_path": self.current_manifest_source_rel, "current_manifest_bytes": current_manifest_id[0], "current_manifest_sha256": current_manifest_id[1],
|
"manifest_self_formal_relative_path": self.current_manifest_formal, "rows": self.candidate_rows,
|
"link_checks": [{"relative_path": self.index_source_rel, "required_utf8_substrings": ["核心文档/.releases/", "source-evidence-closed"]}],
|
"evidence_checks": [{"relative_path": self.core_source_rel, "required_utf8_substrings": ["source-evidence-closed"]}],
|
},
|
"prior": {
|
"release_id": self.prior_release, "manifest_formal_relative_path": self.prior_manifest_formal,
|
"manifest_bytes": prior_manifest_id[0], "manifest_sha256": prior_manifest_id[1],
|
"manifest_self_formal_relative_path": self.prior_manifest_formal, "release_set_sha256": self.prior_release_set,
|
"rows": self.prior_rows,
|
},
|
"history": {"rows": self.history_rows, "minimum_long_paths": 12, "long_path_threshold": 260},
|
"case_current_index": {"bytes": len(self.prior_index_data), "sha256": sha(self.prior_index_data), "required_utf8_substrings": ["release: prior", "core/current.md"]},
|
"result_current_index": {
|
"member_id": "RESULT-CURRENT", "formal_relative_path": self.rel(self.result_index), "artifact_type": "RESULT_THIN_CURRENT_INDEX",
|
"bytes": result_id[0], "sha256": result_id[1], "required_utf8_substrings": ["../../cases/Demo案例/当前成果索引.md"],
|
},
|
"current_state": {
|
"bytes": state_id[0], "sha256": state_id[1], "semantic_field": "semantic_state", "semantic_value": "RECOVERY_REQUIRED",
|
"status_field": "status", "status_value": "BLOCKED", "exit_code_field": "exit_code", "exit_code_value": 30,
|
},
|
},
|
"commit": {"strategy": "MARKDOWN_CURRENT_INDEX_REPLACE_V1", "target_release_relative_path": self.target_rel, "staging_relative_path": self.rel(self.staging)},
|
"test_control": {"environment": "ISOLATED_TEST", "fault": fault, "race_marker_relative_path": ""},
|
}
|
self.config_path = self.root / "batch-config.json"
|
self.save()
|
|
def rel(self, path: Path) -> str:
|
return os.path.relpath(path, self.root).replace("\\", "/")
|
|
def history_row(self, rel: str):
|
size, digest = ident(self.history / rel)
|
return {"relative_path": rel, "bytes": size, "sha256": digest}
|
|
def artifact_row(self, member, source, formal, artifact_type, role):
|
size, digest = ident(self.candidate / source)
|
return {"member_id": member, "relative_path": source, "formal_relative_path": formal, "artifact_type": artifact_type, "commit_role": role, "bytes": size, "sha256": digest}
|
|
def save(self):
|
write(self.config_path, canonical(self.config))
|
|
def run(self, validate=False, config_path=None):
|
path = config_path or self.config_path
|
env = os.environ.copy()
|
env["PYTHONPATH"] = str(MODULE_ROOT)
|
env["PYTHONUTF8"] = "1"
|
env["PYTHONIOENCODING"] = "utf-8"
|
command = [PYTHON, "-m", "shared_content_publisher", "--config", str(path)]
|
if validate:
|
command.append("--validate-only")
|
proc = subprocess.run(command, cwd=self.root, env=env, text=True, encoding="utf-8", capture_output=True, timeout=40)
|
return proc, json.loads(proc.stdout.strip().splitlines()[-1])
|
|
def popen(self):
|
env = os.environ.copy()
|
env["PYTHONPATH"] = str(MODULE_ROOT)
|
env["PYTHONUTF8"] = "1"
|
env["PYTHONIOENCODING"] = "utf-8"
|
return subprocess.Popen([PYTHON, "-m", "shared_content_publisher", "--config", str(self.config_path)], cwd=self.root, env=env, text=True, encoding="utf-8", stdout=subprocess.PIPE, stderr=subprocess.PIPE)
|
|
def formal_state(self):
|
return {
|
"ledger": ident(self.ledger), "case": ident(self.case_index), "result": ident(self.result_index),
|
"target": self.target.exists(), "receipt": self.receipt.exists(), "lock": self.lock.exists(),
|
}
|
|
def rebuild_candidate_manifests(self, current_rows=None):
|
rows = current_rows if current_rows is not None else self.current_manifest_rows
|
current_data = csv_bytes(["member_id", "formal_relative_path", "artifact_type", "bytes", "sha256"], rows)
|
write(self.candidate / self.current_manifest_source_rel, current_data)
|
for row in self.candidate_rows:
|
if row["relative_path"] == self.current_manifest_source_rel:
|
row["bytes"], row["sha256"] = ident(self.candidate / self.current_manifest_source_rel)
|
candidate_data = csv_bytes(
|
["member_id", "candidate_relative_path", "formal_relative_path", "artifact_type", "bytes", "sha256"],
|
[(row["member_id"], "candidate/" + row["relative_path"], row["formal_relative_path"], row["artifact_type"], row["bytes"], row["sha256"]) for row in self.candidate_rows],
|
)
|
write(self.candidate / self.candidate_manifest_rel, candidate_data)
|
current_id = ident(self.candidate / self.current_manifest_source_rel)
|
manifest_id = ident(self.candidate / self.candidate_manifest_rel)
|
section = self.config["expected"]["candidate"]
|
section["rows"] = self.candidate_rows
|
section["current_manifest_bytes"], section["current_manifest_sha256"] = current_id
|
section["manifest_bytes"], section["manifest_sha256"] = manifest_id
|
self.config["identity"]["candidate_release_set_sha256"] = release_set(current_data, self.current_manifest_formal)
|
self.save()
|
|
|
class SharedPublisherTests(unittest.TestCase):
|
def fixture(self, fault="NONE"):
|
root = Path(tempfile.mkdtemp(prefix="shared-publisher-v2-"))
|
self.addCleanup(lambda: shutil.rmtree(long(root), ignore_errors=True))
|
return Fixture(root, fault=fault)
|
|
def assert_prewrite_failure(self, fx, code):
|
before = fx.formal_state()
|
proc, result = fx.run()
|
self.assertNotEqual(proc.returncode, 0, proc.stdout + proc.stderr)
|
self.assertEqual(result["status"], "FAIL_CLOSED")
|
self.assertEqual(result["error_code"], code)
|
self.assertEqual(fx.formal_state(), before)
|
|
def test_00_production_guard_precedes_isolated_mutation(self):
|
guard = json.loads(read(GUARD_PATH).decode("utf-8"))
|
self.assertEqual(guard["schema_version"], "SHARED_CONTENT_PUBLISHER_PRODUCTION_GUARD_V1")
|
committed_attempt = (
|
PROJECT_ROOT / "ana-data" / "cases" / "农业案例" / "审计包"
|
/ "ANA-AGRICULTURE-PESTICIDE-FERTILIZER-20260818-001" / "BATCH-004"
|
/ "RUN-ANA-AGRICULTURE-PESTICIDE-FERTILIZER-20260818-001-BATCH-004-001"
|
/ "attempts" / "ATTEMPT-B004-000003"
|
)
|
prior_files = committed_attempt / "prior_snapshot" / "files"
|
for row in guard["files"]:
|
live = PROJECT_ROOT / row["relative_path"]
|
snapshot = prior_files / row["relative_path"]
|
if os.path.isfile(long(snapshot)):
|
self.assertEqual(ident(snapshot), (row["bytes"], row["sha256"]))
|
elif row["relative_path"].endswith("/promotion_ledger.csv"):
|
data = read(live)[:row["bytes"]]
|
self.assertEqual((len(data), sha(data)), (row["bytes"], row["sha256"]))
|
else:
|
self.assertEqual(ident(live), (row["bytes"], row["sha256"]))
|
for rel in guard["absent_paths"]:
|
path = PROJECT_ROOT / rel
|
if path == committed_attempt:
|
terminal = json.loads(read(path / "terminal.json").decode("utf-8"))
|
self.assertEqual((terminal["status"], terminal["exit_code"]), ("COMMITTED", 0))
|
else:
|
self.assertFalse(path.exists(), rel)
|
|
def test_01_markdown_commit_and_strict_idempotent_replay(self):
|
fx = self.fixture()
|
check, valid = fx.run(validate=True)
|
self.assertEqual(check.returncode, 0, check.stdout + check.stderr)
|
self.assertEqual(valid["lifecycle"], "FRESH")
|
proc, result = fx.run()
|
self.assertEqual(proc.returncode, 0, proc.stdout + proc.stderr)
|
self.assertEqual(result["status"], "COMMITTED")
|
committed_rows = [row for row in ledger_rows(read(fx.ledger)) if row["attempt_id"] == fx.attempt]
|
self.assertEqual(result["attempt_event_chain_row_count"], 4)
|
self.assertEqual(result["attempt_event_chain_sha256"], event_chain_digest(committed_rows))
|
self.assertEqual(result["event_time_start_utc"], committed_rows[0]["event_time"])
|
self.assertEqual(result["event_time_end_utc"], committed_rows[-1]["event_time"])
|
self.assertEqual(result["event_process_identities"], [result["process_instance_identity"]] * 4)
|
self.assertEqual(read(fx.case_index), read(fx.candidate / fx.index_source_rel))
|
self.assertEqual(ident(fx.result_index), (result["result_current_bytes"], result["result_current_sha256"]))
|
ledger_after = read(fx.ledger)
|
again, replay = fx.run()
|
self.assertEqual(again.returncode, 0, again.stdout + again.stderr)
|
self.assertEqual(replay["status"], "IDEMPOTENT_COMMITTED")
|
self.assertEqual(read(fx.ledger), ledger_after)
|
self.assertFalse(fx.lock.exists())
|
|
def test_02_postcommit_failure_rolls_back_markdown_once(self):
|
fx = self.fixture(fault="POSTCOMMIT_READBACK_FAIL")
|
proc, result = fx.run()
|
self.assertEqual(proc.returncode, 20, proc.stdout + proc.stderr)
|
self.assertEqual(result["status"], "ROLLED_BACK")
|
self.assertEqual(read(fx.case_index), fx.prior_index_data)
|
rollback_rows = [row for row in ledger_rows(read(fx.ledger)) if row["attempt_id"] == fx.attempt]
|
self.assertEqual(result["attempt_event_chain_row_count"], 5)
|
self.assertEqual(result["attempt_event_chain_sha256"], event_chain_digest(rollback_rows))
|
self.assertEqual(result["event_process_identities"], [result["process_instance_identity"]] * 5)
|
ledger_after = read(fx.ledger)
|
again, blocked = fx.run()
|
self.assertEqual(again.returncode, 12, again.stdout + again.stderr)
|
self.assertEqual(blocked["error_code"], "TERMINAL_REPEAT_BLOCKED")
|
self.assertEqual(read(fx.ledger), ledger_after)
|
|
def test_03_locked_ledger_race_is_detected_before_append(self):
|
fx = self.fixture(fault="PAUSE_AFTER_LOCK")
|
marker = fx.root / "race.ready"
|
fx.config["test_control"]["race_marker_relative_path"] = fx.rel(marker)
|
fx.save()
|
process = fx.popen()
|
deadline = time.monotonic() + 10
|
while time.monotonic() < deadline and not fx.lock.exists():
|
time.sleep(0.01)
|
self.assertTrue(fx.lock.exists())
|
race = csv_bytes(LEDGER_HEADER, [[
|
"race", "2", "PREPARED", fx.prior_release, fx.prior_release_set, "D" * 64,
|
fx.task, fx.case, fx.batch, fx.run_id, "AUDIT-RACE", "racer", "PID-RACE-INSTANCE",
|
"2026-01-01T00:00:01Z", "race.json", "ATTEMPT-RACE", "2",
|
]])
|
with open(long(fx.ledger), "ab") as handle:
|
handle.write(race.split(b"\n", 1)[1])
|
write(marker, b"ready\n")
|
stdout, stderr = process.communicate(timeout=20)
|
result = json.loads(stdout.strip().splitlines()[-1])
|
self.assertEqual(process.returncode, 12, stdout + stderr)
|
self.assertEqual(result["error_code"], "LEDGER_IDENTITY_LOCKED")
|
self.assertEqual(read(fx.case_index), fx.prior_index_data)
|
self.assertFalse(fx.target.exists())
|
self.assertFalse(fx.receipt.exists())
|
self.assertFalse(fx.lock.exists())
|
|
def test_04_duplicate_sequence_and_declared_sequence_drift(self):
|
fx = self.fixture()
|
base = read(fx.ledger)
|
duplicate = csv_bytes(LEDGER_HEADER, [[
|
"dup", "1", "ROLLED_BACK", fx.prior_release, fx.prior_release_set, "E" * 64,
|
fx.task, fx.case, fx.batch, fx.run_id, "AUDIT-DUP", "tester", "PID-DUP",
|
"2026-01-01T00:00:02Z", "dup.json", "ATTEMPT-DUP", "2",
|
]]).split(b"\n", 1)[1]
|
write(fx.ledger, base + duplicate)
|
fx.config["expected"]["ledger"]["bytes"] = len(base + duplicate)
|
fx.config["expected"]["ledger"]["sha256"] = sha(base + duplicate)
|
fx.save()
|
self.assert_prewrite_failure(fx, "LEDGER_SEQUENCE_UNIQUE")
|
fx2 = self.fixture()
|
fx2.config["expected"]["ledger"]["next_event_seq"] = 99
|
fx2.save()
|
self.assert_prewrite_failure(fx2, "LEDGER_NEXT_EVENT")
|
|
def test_05_path_overlap_and_junction_ancestor_fail_closed(self):
|
fx = self.fixture()
|
fx.config["roots"]["attempt_receipt_root"] = fx.target_rel
|
fx.save()
|
self.assert_prewrite_failure(fx, "PATH_OVERLAP")
|
|
fx2 = self.fixture()
|
alias = fx2.root.parent / (fx2.root.name + "-junction")
|
created = False
|
if os.name == "nt":
|
made = subprocess.run(["cmd", "/c", "mklink", "/J", str(alias), str(fx2.root)], capture_output=True)
|
created = made.returncode == 0
|
else:
|
os.symlink(fx2.root, alias, target_is_directory=True)
|
created = True
|
if not created:
|
self.skipTest("junction creation unavailable")
|
self.addCleanup(lambda: os.rmdir(long(alias)) if os.path.lexists(long(alias)) else None)
|
fx2.config["roots"]["operation_root"] = str(alias)
|
fx2.config["roots"]["resolved_root"] = str(alias)
|
alias_config = alias / fx2.config_path.name
|
fx2.save()
|
proc, result = fx2.run(config_path=alias_config)
|
self.assertEqual(proc.returncode, 12, proc.stdout + proc.stderr)
|
self.assertIn(result["error_code"], ("PATH_REPARSE", "PATH_ALIAS"))
|
|
def test_06_manifest_omission_and_empty_checks_are_rejected(self):
|
fx = self.fixture()
|
fx.rebuild_candidate_manifests([row for row in fx.current_manifest_rows if row[0] != "RESULT-CURRENT"])
|
self.assert_prewrite_failure(fx, "CURRENT_MANIFEST_COVERAGE")
|
fx2 = self.fixture()
|
fx2.config["expected"]["candidate"]["link_checks"] = []
|
fx2.save()
|
self.assert_prewrite_failure(fx2, "CONFIG_TYPE")
|
fx3 = self.fixture()
|
fx3.config["expected"]["candidate"]["evidence_checks"] = []
|
fx3.save()
|
self.assert_prewrite_failure(fx3, "CONFIG_TYPE")
|
|
def test_07_extra_target_file_forces_rollback(self):
|
fx = self.fixture(fault="EXTRA_TARGET_AFTER_RENAME")
|
proc, result = fx.run()
|
self.assertEqual(proc.returncode, 20, proc.stdout + proc.stderr)
|
self.assertEqual(result["status"], "ROLLED_BACK")
|
self.assertEqual(read(fx.case_index), fx.prior_index_data)
|
self.assertTrue((fx.target / "UNDECLARED.txt").exists())
|
|
def test_08_replay_target_attack_rolls_back_before_false_idempotence(self):
|
fx = self.fixture()
|
proc, _ = fx.run()
|
self.assertEqual(proc.returncode, 0, proc.stdout + proc.stderr)
|
write(fx.target / "UNDECLARED.txt", b"attack\n")
|
again, result = fx.run()
|
self.assertEqual(again.returncode, 20, again.stdout + again.stderr)
|
self.assertEqual(result["status"], "ROLLED_BACK")
|
self.assertEqual(read(fx.case_index), fx.prior_index_data)
|
self.assertFalse(fx.lock.exists())
|
|
def test_09_invalid_config_cannot_use_replay_bypass(self):
|
fx = self.fixture()
|
proc, _ = fx.run()
|
self.assertEqual(proc.returncode, 0, proc.stdout + proc.stderr)
|
fx.config["identity"]["unexpected"] = "attack"
|
fx.save()
|
again, result = fx.run()
|
self.assertEqual(again.returncode, 12, again.stdout + again.stderr)
|
self.assertEqual(result["error_code"], "CONFIG_KEYS")
|
self.assertNotEqual(result["status"], "IDEMPOTENT_COMMITTED")
|
|
def test_10_terminal_and_event_chain_tamper_force_rollback(self):
|
fx = self.fixture()
|
proc, _ = fx.run()
|
self.assertEqual(proc.returncode, 0, proc.stdout + proc.stderr)
|
terminal_path = fx.receipt / "terminal.json"
|
terminal = json.loads(read(terminal_path).decode("utf-8"))
|
terminal["volume_identity"] = "DRIFT"
|
write(terminal_path, canonical(terminal))
|
again, result = fx.run()
|
self.assertEqual(again.returncode, 20, again.stdout + again.stderr)
|
self.assertEqual(result["status"], "ROLLED_BACK")
|
self.assertEqual(read(fx.case_index), fx.prior_index_data)
|
self.assertTrue((fx.receipt / "terminal.json").exists())
|
self.assertTrue((fx.receipt / "recovery_terminal.json").exists())
|
recovery_rows = [row for row in ledger_rows(read(fx.ledger)) if row["attempt_id"] == fx.attempt]
|
self.assertEqual(result["attempt_event_chain_row_count"], 6)
|
self.assertEqual(result["attempt_event_chain_sha256"], event_chain_digest(recovery_rows))
|
self.assertEqual(result["event_process_identities"], [row["process_identity"] for row in recovery_rows])
|
|
fx2 = self.fixture()
|
proc2, _ = fx2.run()
|
self.assertEqual(proc2.returncode, 0, proc2.stdout + proc2.stderr)
|
write(fx2.receipt / "terminal.json", b"{invalid terminal\n")
|
again2, result2 = fx2.run()
|
self.assertEqual(again2.returncode, 20, again2.stdout + again2.stderr)
|
self.assertEqual(result2["status"], "ROLLED_BACK")
|
self.assertEqual(read(fx2.case_index), fx2.prior_index_data)
|
|
def test_11_prior_release_set_and_result_thin_entry_are_bound(self):
|
fx = self.fixture()
|
fx.config["expected"]["prior"]["release_set_sha256"] = "0" * 64
|
fx.save()
|
self.assert_prewrite_failure(fx, "PRIOR_RELEASE_SET")
|
fx2 = self.fixture()
|
write(fx2.result_index, b"# not a thin entry\n")
|
self.assert_prewrite_failure(fx2, "RESULT_INDEX_IDENTITY")
|
|
def test_11a_every_ledger_column_tamper_is_fail_closed(self):
|
replacements = {
|
"release_id": "attacked-release",
|
"event_seq": "999",
|
"state": "ABORTED",
|
"prior_release_id": "a" * 64,
|
"prior_release_set_sha256": "A" * 64,
|
"candidate_release_set_sha256": "B" * 64,
|
"task_id": "TASK-ATTACKED",
|
"case_id": "CASE-ATTACKED",
|
"batch_id": "BATCH-ATTACKED",
|
"run_id": "RUN-ATTACKED",
|
"accepted_audit_id": "AUDIT-ATTACKED",
|
"operator": "ATTACKED-OPERATOR",
|
"process_identity": "PID-ATTACKED-INSTANCE",
|
"event_time": "2099-01-01T00:00:00.000000Z",
|
"recovery_or_rollback_receipt": "receipts/attacked.json",
|
"attempt_id": "ATTEMPT-ATTACKED-999999",
|
"attempt_seq": "999",
|
}
|
self.assertEqual(set(replacements), set(LEDGER_HEADER))
|
for column in LEDGER_HEADER:
|
with self.subTest(column=column):
|
fx = self.fixture()
|
first, committed = fx.run()
|
self.assertEqual(first.returncode, 0, first.stdout + first.stderr)
|
self.assertEqual(committed["status"], "COMMITTED")
|
rows = ledger_rows(read(fx.ledger))
|
rows[-1][column] = replacements[column]
|
write(fx.ledger, csv_bytes(LEDGER_HEADER, [[row[key] for key in LEDGER_HEADER] for row in rows]))
|
replay, result = fx.run()
|
self.assertNotEqual(replay.returncode, 0, replay.stdout + replay.stderr)
|
self.assertEqual(result["status"], "FAIL_CLOSED")
|
self.assertNotEqual(result["status"], "IDEMPOTENT_COMMITTED")
|
self.assertFalse(fx.lock.exists())
|
|
def test_11b_every_state_receipt_and_time_are_bound(self):
|
for column in ("recovery_or_rollback_receipt", "event_time"):
|
for state_index in range(4):
|
with self.subTest(column=column, state_index=state_index):
|
fx = self.fixture()
|
first, committed = fx.run()
|
self.assertEqual(first.returncode, 0, first.stdout + first.stderr)
|
self.assertEqual(committed["status"], "COMMITTED")
|
rows = ledger_rows(read(fx.ledger))
|
attempt_indexes = [index for index, row in enumerate(rows) if row["attempt_id"] == fx.attempt]
|
target = rows[attempt_indexes[state_index]]
|
if column == "recovery_or_rollback_receipt":
|
target[column] = f"receipts/attacked-{state_index}.json"
|
else:
|
target[column] = f"2099-01-01T00:00:0{state_index}.000000Z"
|
write(fx.ledger, csv_bytes(LEDGER_HEADER, [[row[key] for key in LEDGER_HEADER] for row in rows]))
|
replay, result = fx.run()
|
self.assertNotEqual(replay.returncode, 0, replay.stdout + replay.stderr)
|
self.assertEqual(result["status"], "FAIL_CLOSED")
|
self.assertNotEqual(result["status"], "IDEMPOTENT_COMMITTED")
|
self.assertFalse(fx.lock.exists())
|
|
def test_12_candidate_history_root_and_schema_negatives(self):
|
fx = self.fixture()
|
write(fx.candidate / fx.core_source_rel, b"drift\n")
|
self.assert_prewrite_failure(fx, "CANDIDATE_IDENTITY")
|
fx2 = self.fixture()
|
long_row = next(row for row in fx2.history_rows if len(os.path.abspath(fx2.history / row["relative_path"])) >= 260)
|
write(fx2.history / long_row["relative_path"], b"drift\n")
|
self.assert_prewrite_failure(fx2, "HISTORY_IDENTITY")
|
fx3 = self.fixture()
|
fx3.config["roots"]["candidate_root"] = "../escape"
|
fx3.save()
|
self.assert_prewrite_failure(fx3, "PATH_RELATIVE")
|
fx4 = self.fixture()
|
del fx4.config["identity"]["operator"]
|
fx4.save()
|
self.assert_prewrite_failure(fx4, "CONFIG_KEYS")
|
|
def test_13_audit_handoff_domains_remain_closed(self):
|
fx = self.fixture()
|
fx.config["identity"]["canonical_audit_id"], fx.config["identity"]["review_handoff_id"] = fx.config["identity"]["review_handoff_id"], fx.config["identity"]["canonical_audit_id"]
|
fx.save()
|
self.assert_prewrite_failure(fx, "AUDIT_ID_DOMAIN")
|
|
|
if __name__ == "__main__":
|
unittest.main(verbosity=2)
|