526 lines
27 KiB
Python
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()
|