from __future__ import annotations from dataclasses import replace import hashlib import json import multiprocessing import os from pathlib import Path import tempfile import unittest from unittest.mock import patch from arr_ingestion.contracts import IngestionError from arr_ingestion.repository import InMemoryIngestionRepository from arr_ingestion.postgres import DatabaseConfig, PostgresIngestionRepository from arr_ingestion.service import IngestionService from arr_ingestion.validation import DeliveryValidator from arr_processing.local import LocalDailyProcessor from arr_processing.policy import load_processor_policy from arr_storage.contracts import ObjectKeyPolicy from arr_storage.filesystem import FilesystemObjectBackend from arr_storage.store import ManagedObjectStore from integrations.ohip import processing_handoff as handoff from integrations.ohip.capture_job import job_lock from tests.test_arr_opera_daily_ingest import reservation, success_xml, xml_document PROJECT = Path(__file__).resolve().parents[1] BINDING = handoff.CaptureBinding("synthetic-arr-001", "OHIPSB02", "2026-07-27", "a" * 64, "synthetic-xml-fixture/v1") def hold_lock(directory, entered, release): with job_lock(Path(directory)): entered.set() if not release.wait(15): raise RuntimeError("test lock wait expired") def exit_after_freeze(root): publish = handoff.atomic_json def exit_before_ready(path, *args, **kwargs): if path.name == "ready.json": os._exit(73) # No exception handlers, context cleanup or lock release code. return publish(path, *args, **kwargs) handoff.atomic_json = exit_before_ready handoff.prepare(Path(root) / "handoffs", BINDING, "ARR.XML", success_xml().encode(), load_processor_policy(PROJECT)) class CountingProcessor: def __init__(self, policy): self.delegate = LocalDailyProcessor(policy) self.calls = 0 def run(self, *args, **kwargs): self.calls += 1 return self.delegate.run(*args, **kwargs) class FrozenHandoffTests(unittest.TestCase): @classmethod def setUpClass(cls): cls.policy = load_processor_policy(PROJECT) def setUp(self): self.temporary = tempfile.TemporaryDirectory() self.addCleanup(self.temporary.cleanup) self.root = Path(self.temporary.name) self.repository = InMemoryIngestionRepository() self.store = ManagedObjectStore(FilesystemObjectBackend(self.root / "remote-objects", create=True)) self.service = IngestionService(DeliveryValidator(self.store, self.policy), self.repository) self.processor = CountingProcessor(self.policy) def prepare(self, payload=None, binding=BINDING, **kwargs): return handoff.prepare(self.root / "handoffs", binding, "ARR.XML", payload if payload is not None else success_xml().encode(), self.policy, processor=kwargs.get("processor", self.processor)) def deliver(self, prepared): return handoff.deliver(Path(prepared["directory"]), prepared["manifest_sha256"], object_store=self.store, repository=self.repository, service=self.service) def test_prepare_has_no_external_effects_and_reuses_exact_bytes(self): prepared = self.prepare() frozen = Path(prepared["directory"]) / "frozen" before = {str(path.relative_to(frozen)): hashlib.sha256(path.read_bytes()).hexdigest() for path in frozen.rglob("*") if path.is_file()} again = self.prepare() after = {str(path.relative_to(frozen)): hashlib.sha256(path.read_bytes()).hexdigest() for path in frozen.rglob("*") if path.is_file()} self.assertTrue(again["reused"]) self.assertEqual(prepared["manifest_sha256"], again["manifest_sha256"]) self.assertEqual(before, after) self.assertEqual(self.processor.calls, 1) self.assertFalse(self.repository._jobs) self.assertFalse(list((self.root / "remote-objects").rglob("*.xml"))) for path in frozen.rglob("*"): self.assertEqual(path.stat().st_mode & 0o777, 0o700 if path.is_dir() else 0o600) def test_exact_delivery_replay_keeps_one_job_callback_and_version(self): prepared = self.prepare() first = self.deliver(prepared) second = self.deliver(prepared) self.assertEqual(first, second) self.assertEqual(first["ingestion_status"], "committed") self.assertEqual(first["version_no"], 1) self.assertEqual((len(self.repository._jobs), len(self.repository._callbacks), len(self.repository._versions)), (1, 1, 1)) self.assertEqual(self.processor.calls, 1) self.assertFalse(first["source_mapping_verified"]) def test_expected_binding_and_policy_checked_before_any_destination_write(self): prepared = self.prepare() directory = Path(prepared["directory"]) changed_bindings = [replace(BINDING, **change) for change in ( {"batch_id": "another-batch"}, {"hotel_id": "OTHER"}, {"arrival_date": "2026-09-16"}, {"manifest_sha256": "b" * 64}, {"adapter_contract": "synthetic/v99"}, )] with patch.object(self.store, "upload_committed", side_effect=AssertionError("unexpected write")) as upload: for binding in changed_bindings: with self.subTest(binding_field=binding), self.assertRaisesRegex( handoff.source.CollectionError, "handoff_expected_binding_mismatch" ): handoff.deliver(directory, prepared["manifest_sha256"], object_store=self.store, repository=self.repository, service=self.service, expected_binding=binding, expected_policy=self.policy) for policy in (replace(self.policy, processor_version="other-version"), replace(self.policy, rule_set_sha256="c" * 64)): with self.assertRaisesRegex(handoff.source.CollectionError, "handoff_expected_policy_mismatch"): handoff.deliver(directory, prepared["manifest_sha256"], object_store=self.store, repository=self.repository, service=self.service, expected_binding=BINDING, expected_policy=policy) upload.assert_not_called() self.assertFalse(self.repository._jobs) receipt = handoff.deliver(directory, prepared["manifest_sha256"], object_store=self.store, repository=self.repository, service=self.service, expected_binding=BINDING, expected_policy=self.policy) self.assertEqual(receipt["ingestion_status"], "committed") self.assertEqual(self.deliver(prepared), receipt) # Legacy pin-only reuse stays compatible. self.assertEqual(len(self.repository._versions), 1) def test_lost_commit_acknowledgement_replays_without_failing_job(self): prepared = self.prepare() ingest = self.service.ingest def lost_ack(raw): ingest(raw) raise IngestionError("DATABASE_WRITE_FAILED", "synthetic lost acknowledgement") with patch.object(self.service, "ingest", side_effect=lost_ack): with self.assertRaises(IngestionError): self.deliver(prepared) self.assertEqual(self.repository._jobs[prepared["job_id"]].status, "succeeded") self.assertFalse((Path(prepared["directory"]) / "receipt.json").exists()) receipt = self.deliver(prepared) self.assertEqual(receipt["version_no"], 1) self.assertEqual(len(self.repository._versions), 1) def test_receipt_write_crash_after_commit_is_recoverable(self): prepared = self.prepare() with patch.object(handoff, "atomic_json", side_effect=KeyboardInterrupt): with self.assertRaises(KeyboardInterrupt): self.deliver(prepared) self.assertEqual(len(self.repository._versions), 1) self.assertEqual(self.deliver(prepared)["version_no"], 1) self.assertEqual(len(self.repository._callbacks), 1) def test_lost_registration_acknowledgement_preserves_same_job(self): prepared = self.prepare() register = self.repository.register_job def lost_ack(registration): register(registration) raise IngestionError("DATABASE_WRITE_FAILED", "synthetic registration acknowledgement") with patch.object(self.repository, "register_job", side_effect=lost_ack): with self.assertRaises(IngestionError): self.deliver(prepared) self.assertEqual(len(self.repository._jobs), 1) self.assertFalse(self.repository._versions) self.assertEqual(self.deliver(prepared)["version_no"], 1) self.assertEqual(len(self.repository._jobs), 1) def test_failure_before_commit_retries_frozen_delivery(self): prepared = self.prepare() with patch.object(self.service, "ingest", side_effect=IngestionError("DATABASE_UNAVAILABLE", "synthetic failure")): with self.assertRaises(IngestionError): self.deliver(prepared) self.assertFalse(self.repository._callbacks) self.assertEqual(self.repository._jobs[prepared["job_id"]].status, "queued") self.assertEqual(self.deliver(prepared)["ingestion_status"], "committed") self.assertEqual(self.processor.calls, 1) def test_partial_object_upload_retries_same_artifacts(self): prepared = self.prepare() upload = self.store.upload_committed def fail_result(**kwargs): if kwargs["role"] == "result_json": raise IngestionError("OBJECT_STORE_UNAVAILABLE", "synthetic upload failure") return upload(**kwargs) with patch.object(self.store, "upload_committed", side_effect=fail_result): with self.assertRaises(IngestionError): self.deliver(prepared) self.assertFalse(self.repository._jobs) self.assertEqual(self.deliver(prepared)["ingestion_status"], "committed") self.assertEqual(self.processor.calls, 1) def test_review_replay_creates_one_case_without_finance(self): prepared = self.prepare(xml_document(reservation(1, rate_code="GRPA1", rate_amount="1800")).encode()) first = self.deliver(prepared) second = self.deliver(prepared) self.assertEqual(first, second) self.assertEqual(first["ingestion_status"], "recorded_review") self.assertEqual(len(self.repository._reviews), 1) self.assertFalse(self.repository._versions) def test_business_failure_is_recorded_once_not_reprocessed(self): prepared = self.prepare(xml_document(reservation(1, departure="2026-07-26")).encode()) self.assertEqual(self.deliver(prepared)["ingestion_status"], "recorded_failure") self.assertEqual(self.deliver(prepared)["ingestion_status"], "recorded_failure") self.assertEqual(len(self.repository._callbacks), 1) self.assertEqual(self.processor.calls, 1) self.assertFalse(self.repository._versions) def test_changed_source_capture_adapter_or_rules_conflicts_without_new_job(self): prepared = self.prepare() for change in ({"manifest_sha256": "b" * 64}, {"adapter_contract": "different/v1"}, {"arrival_date": "2026-07-28"}): with self.subTest(change=change), self.assertRaisesRegex(source_error(), "handoff_request_conflict"): self.prepare(binding=replace(BINDING, **change)) with self.assertRaisesRegex(source_error(), "handoff_request_conflict"): self.prepare(success_xml().replace("", "\n").encode()) with self.assertRaisesRegex(source_error(), "handoff_request_conflict"): handoff.prepare(self.root / "handoffs", BINDING, "ARR.XML", success_xml().encode(), replace(self.policy, rule_set_sha256="b" * 64)) self.assertEqual(self.processor.calls, 1) self.assertEqual(self.deliver(prepared)["version_no"], 1) def test_changed_artifact_or_pin_refused_before_upload(self): prepared = self.prepare() with patch.object(self.store, "upload_committed") as upload: with self.assertRaisesRegex(source_error(), "prepared_pin_mismatch"): self.deliver({**prepared, "manifest_sha256": "b" * 64}) source_path = Path(prepared["directory"]) / "frozen/artifacts/source.xml" source_path.write_bytes(source_path.read_bytes() + b"\n") with self.assertRaisesRegex(source_error(), "prepared_artifact_changed"): self.deliver(prepared) upload.assert_not_called() def test_persisted_package_identity_rejects_numeric_type_changes(self): prepared = self.prepare() directory = Path(prepared["directory"]) path = directory / "identity.json" original = json.loads(path.read_bytes()) path.write_bytes(handoff.source.json_bytes(dict(reversed(list(original.items()))))) self.assertTrue(self.prepare()["reused"]) changed = dict(original, source_byte_size=float(original["source_byte_size"])) path.write_bytes(handoff.source.json_bytes(changed)) before = {str(p): p.read_bytes() for p in directory.rglob("*") if p.is_file()} with self.assertRaisesRegex(source_error(), "handoff_request_conflict"): self.prepare() with patch.object(self.store, "upload_committed", side_effect=AssertionError("unexpected write")) as upload, \ self.assertRaisesRegex(source_error(), "prepared_identity_changed"): self.deliver(prepared) upload.assert_not_called() self.assertEqual(before, {str(p): p.read_bytes() for p in directory.rglob("*") if p.is_file()}) self.assertFalse(self.repository._jobs) self.assertEqual(self.processor.calls, 1) def test_manifest_changes_and_missing_ready_pin_refused(self): prepared = self.prepare() directory = Path(prepared["directory"]) manifest = directory / "frozen/manifest.json" manifest.write_bytes(manifest.read_bytes() + b"\n") with self.assertRaisesRegex(source_error(), "prepared_manifest_changed"): self.prepare() (directory / "ready.json").unlink() (directory / "freeze-intent.json").unlink() # Legacy package has no pre-freeze pin. with self.assertRaisesRegex(source_error(), "uncommitted_handoff_package"): self.prepare() def test_symlink_artifact_refused_before_upload(self): prepared = self.prepare() source_path = Path(prepared["directory"]) / "frozen/artifacts/source.xml" original = self.root / "original.xml" source_path.rename(original) source_path.symlink_to(original) with patch.object(self.store, "upload_committed") as upload: with self.assertRaises(OSError): self.deliver(prepared) upload.assert_not_called() def test_preparation_failure_can_retry_without_external_side_effect(self): with patch.object(self.processor, "run", side_effect=KeyboardInterrupt): with self.assertRaises(KeyboardInterrupt): self.prepare() self.assertFalse(self.repository._jobs) self.assertEqual(self.prepare()["job_id"], BINDING.job_id) def test_date_mismatch_rejected_before_freezing(self): with self.assertRaisesRegex(source_error(), "handoff_date_mismatch"): self.prepare(binding=replace(BINDING, arrival_date="2026-07-28")) self.assertFalse(self.repository._jobs) self.assertFalse(list((self.root / "handoffs").rglob("ready.json"))) def test_legacy_unpinned_package_is_not_promoted(self): atomic = handoff.atomic_json def crash_on_ready(path, *args, **kwargs): if path.name == "ready.json": raise KeyboardInterrupt return atomic(path, *args, **kwargs) with patch.object(handoff, "atomic_json", side_effect=crash_on_ready): with self.assertRaises(KeyboardInterrupt): self.prepare() self.assertFalse(self.repository._jobs) (self.root / "handoffs" / BINDING.job_id / "freeze-intent.json").unlink() with self.assertRaisesRegex(source_error(), "uncommitted_handoff_package"): self.prepare() def test_interrupted_ready_publication_recovers_exact_frozen_bytes(self): publish = handoff.atomic_json def interrupt_ready(path, *args, **kwargs): if path.name == "ready.json": raise KeyboardInterrupt return publish(path, *args, **kwargs) with patch.object(handoff, "atomic_json", side_effect=interrupt_ready): with self.assertRaises(KeyboardInterrupt): self.prepare() directory = self.root / "handoffs" / BINDING.job_id frozen = directory / "frozen" before = {str(p.relative_to(frozen)): p.read_bytes() for p in frozen.rglob("*") if p.is_file()} self.assertFalse(self.repository._jobs) with patch.object(self.processor, "run", side_effect=AssertionError("must not reprocess")): prepared = self.prepare() self.assertTrue(prepared["reused"]) self.assertEqual(before, {str(p.relative_to(frozen)): p.read_bytes() for p in frozen.rglob("*") if p.is_file()}) self.assertEqual(self.deliver(prepared)["version_no"], 1) self.assertEqual(self.deliver(prepared)["version_no"], 1) self.assertEqual((len(self.repository._jobs), len(self.repository._callbacks), len(self.repository._versions)), (1, 1, 1)) def interrupt_preparation(self, boundary): publish = handoff.atomic_json def interrupt(path, *args, **kwargs): if path.name == boundary: if boundary == "freeze-intent.json": publish(path, *args, **kwargs) raise KeyboardInterrupt return publish(path, *args, **kwargs) with patch.object(handoff, "atomic_json", side_effect=interrupt): with self.assertRaises(KeyboardInterrupt): self.prepare() return self.root / "handoffs" / BINDING.job_id def test_pre_rename_interruption_reprepares_without_external_effects(self): directory = self.interrupt_preparation("freeze-intent.json") self.assertTrue((directory / "freeze-intent.json").exists()) self.assertFalse((directory / "frozen").exists()) self.assertFalse(self.repository._jobs) prepared = self.prepare() self.assertFalse(prepared["reused"]) self.assertEqual(self.processor.calls, 2) self.assertFalse(self.repository._jobs) self.assertEqual(self.deliver(prepared)["version_no"], 1) def test_pending_package_cannot_be_delivered_before_verified_recovery(self): directory = self.interrupt_preparation("ready.json") intent = json.loads((directory / "freeze-intent.json").read_bytes()) with patch.object(self.store, "upload_committed") as upload: with self.assertRaises(FileNotFoundError): self.deliver({"directory": str(directory), "manifest_sha256": intent["manifest_sha256"]}) upload.assert_not_called() self.assertFalse(self.repository._jobs) def test_pending_artifact_corruption_is_rejected_without_reprocessing(self): directory = self.interrupt_preparation("ready.json") artifact = directory / "frozen/artifacts/source.xml" artifact.write_bytes(artifact.read_bytes() + b"\n") with patch.object(self.processor, "run", side_effect=AssertionError("must not overwrite corrupt package")): with self.assertRaisesRegex(source_error(), "prepared_artifact_changed"): self.prepare() self.assertFalse((directory / "ready.json").exists()) self.assertFalse(self.repository._jobs) def test_pending_manifest_cannot_supply_its_own_replacement_pin(self): directory = self.interrupt_preparation("ready.json") manifest = directory / "frozen/manifest.json" manifest.write_bytes(manifest.read_bytes() + b"\n") with self.assertRaisesRegex(source_error(), "prepared_manifest_changed"): self.prepare() self.assertFalse((directory / "ready.json").exists()) self.assertEqual(self.processor.calls, 1) def test_pending_intent_must_match_identity_and_pin(self): directory = self.interrupt_preparation("ready.json") path = directory / "freeze-intent.json" original = json.loads(path.read_bytes()) for changes, code in [({"identity_sha256": "0" * 64}, "freeze_intent_context_mismatch"), ({"manifest_sha256": "invalid"}, "invalid_freeze_intent_pin"), ({"manifest_sha256": "0" * 64}, "prepared_manifest_changed"), ({"extra": True}, "freeze_intent_context_mismatch")]: with self.subTest(changes=changes): path.write_bytes(handoff.source.json_bytes({**original, **changes})) with self.assertRaisesRegex(source_error(), code): self.prepare() self.assertFalse((directory / "ready.json").exists()) path.write_bytes(handoff.source.json_bytes(original)) self.assertTrue(self.prepare()["reused"]) self.assertEqual(self.processor.calls, 1) def test_pending_recovery_can_itself_be_interrupted_and_retried(self): directory = self.interrupt_preparation("ready.json") before = {p.name: p.read_bytes() for p in (directory / "frozen/artifacts").iterdir()} self.interrupt_preparation("ready.json") self.assertTrue(self.prepare()["reused"]) self.assertEqual(before, {p.name: p.read_bytes() for p in (directory / "frozen/artifacts").iterdir()}) self.assertEqual(self.processor.calls, 1) def test_committed_legacy_package_without_intent_remains_readable(self): prepared = self.prepare() (Path(prepared["directory"]) / "freeze-intent.json").unlink() self.assertTrue(self.prepare()["reused"]) self.assertEqual(self.deliver(prepared)["version_no"], 1) def test_real_process_exit_after_freeze_recovers_without_processor(self): child = multiprocessing.get_context("spawn").Process(target=exit_after_freeze, args=(str(self.root),)) child.start() try: child.join(30) self.assertEqual(child.exitcode, 73) finally: if child.is_alive(): child.terminate() child.join(5) directory = self.root / "handoffs" / BINDING.job_id self.assertFalse((directory / "ready.json").exists()) pin = json.loads((directory / "freeze-intent.json").read_bytes())["manifest_sha256"] with patch.object(self.processor, "run", side_effect=AssertionError("must not reprocess")): prepared = self.prepare() self.assertTrue(prepared["reused"]) self.assertEqual(prepared["manifest_sha256"], pin) self.assertEqual(self.deliver(prepared)["version_no"], 1) self.assertEqual(self.deliver(prepared)["version_no"], 1) self.assertEqual(len(self.repository._versions), 1) def test_processor_metadata_tampering_is_rejected_by_independent_validation(self): run = self.processor.run def altered_output(*args, **kwargs): output = run(*args, **kwargs) path = output.artifacts["structured_result_json"] document = json.loads(path.read_text()) document["processor_version"] = "unapproved" path.write_text(json.dumps(document)) return output with patch.object(self.processor, "run", side_effect=altered_output): with self.assertRaises(IngestionError): self.prepare() self.assertFalse(list((self.root / "handoffs").rglob("ready.json"))) self.assertFalse(self.repository._jobs) def test_changed_local_receipt_does_not_override_database_result(self): prepared = self.prepare() expected = self.deliver(prepared) (Path(prepared["directory"]) / "receipt.json").write_text('{"version_no": 999}') self.assertEqual(self.deliver(prepared), expected) self.assertEqual(len(self.repository._versions), 1) def test_postgres_exact_replay_precedes_terminal_check_and_has_no_writes(self): # Exercise real Postgres repository SQL/control flow with recorded rows; # no connection or real database is used. Transaction/race acceptance is separate. prepared = self.prepare() self.deliver(prepared) raw = (Path(prepared["directory"]) / "frozen/delivery.json").read_bytes() verified = DeliveryValidator(self.store, self.policy).validate(raw) class ReplayCursor: def __init__(self, mismatch=False): self.statements = [] self.mismatch = mismatch def execute(self, query, params): assert query.count("%s") == len(params) self.statements.append(" ".join(query.split())) def fetchone(self): if len(self.statements) == 1: return (1, 2, verified.envelope.processor_version, verified.envelope.rule_set_sha256, "accepted", 3, "succeeded") return (4, "f" * 64 if self.mismatch else verified.envelope_sha256, "committed", 5, None, verified.envelope.business_date, 1, 1) repository = PostgresIngestionRepository(DatabaseConfig("postgresql://unused")) cursor = ReplayCursor() outcome = repository._commit(cursor, verified) self.assertEqual(outcome.daily_version_id, 5) self.assertTrue(all(query.startswith("SELECT ") for query in cursor.statements)) with self.assertRaisesRegex(IngestionError, "identifier conflicts"): repository._commit(ReplayCursor(mismatch=True), verified) def test_another_process_holding_batch_blocks_submission(self): prepared = self.prepare() context = multiprocessing.get_context("spawn") entered, release = context.Event(), context.Event() child = context.Process(target=hold_lock, args=(prepared["directory"], entered, release)) child.start() try: self.assertTrue(entered.wait(10)) with self.assertRaisesRegex(source_error(), "batch_busy"): self.deliver(prepared) self.assertFalse(self.repository._jobs) finally: release.set() child.join(15) if child.is_alive(): child.terminate() child.join() self.assertEqual(child.exitcode, 0) self.assertEqual(self.deliver(prepared)["version_no"], 1) def test_other_store_namespace_rejected_before_any_upload(self): prepared = self.prepare() self.store = ManagedObjectStore(FilesystemObjectBackend(self.root / "other", create=True), ObjectKeyPolicy("other")) with patch.object(self.store, "upload_committed") as upload: with self.assertRaises(IngestionError): self.deliver(prepared) upload.assert_not_called() def source_error(): return handoff.source.CollectionError if __name__ == "__main__": unittest.main()