Files
ARR-2.0-0918/tests/test_ohip_processing_handoff.py

526 lines
27 KiB
Python

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("<RES_DETAIL>", "<RES_DETAIL>\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()