"""Synthetic capture-to-delivery recovery; no accepted real source adapter.""" import hashlib from pathlib import Path import tempfile import unittest from unittest.mock import patch from arr_ingestion.repository import InMemoryIngestionRepository from arr_ingestion.service import IngestionService from arr_ingestion.validation import DeliveryValidator from arr_processing.policy import load_processor_policy from arr_storage.filesystem import FilesystemObjectBackend from arr_storage.store import ManagedObjectStore from arr_web.arr_download_executor import CapturedARRExecutor from integrations.ohip import collect_arr_source as source, processing_handoff as handoff from integrations.ohip.rate_info import RateInfoReader from tests.test_arr_web_capture_executor import FixtureAdapter, FixtureValidator from tests.test_arr_web_download_handoff import REQUEST_ID, DAY from tests.test_ohip_day_capture import DayService, HOTEL from tests.test_ohip_processing_handoff import CountingProcessor class CaptureRecoveryTests(unittest.TestCase): def test_freeze_interruption_reuses_capture_and_package_then_commits_once(self): with tempfile.TemporaryDirectory() as temporary: root = Path(temporary) policy = load_processor_policy(Path(__file__).resolve().parents[1]) repository = InMemoryIngestionRepository() store = ManagedObjectStore(FilesystemObjectBackend(root / "objects", create=True)) ingestion = IngestionService(DeliveryValidator(store, policy), repository) processor = CountingProcessor(policy) adapter, validator, transport = FixtureAdapter(), FixtureValidator(), DayService() factories = [] def factory(archive, hotel): factories.append(hotel) return (source.Reader(archive, hotel, transport, sleep=lambda _: None), RateInfoReader(archive, hotel, transport, sleep=lambda _: None)) def executor(): return CapturedARRExecutor(root=root / "executor", hotel_id=HOTEL, adapter_contract="synthetic-only/v1", adapter=adapter, mapping_validator=validator, reader_factory=factory, policy=policy, object_store=store, repository=repository, ingestion=ingestion, processor=processor) def run(): return executor().execute(request_id=REQUEST_ID, from_date=DAY, to_date=DAY, report_stage=lambda _: None) def hashes(folder): return {str(p.relative_to(folder)): hashlib.sha256(p.read_bytes()).hexdigest() for p in folder.rglob("*") if p.is_file()} publish = handoff.atomic_json def interrupt_ready(path, *args, **kwargs): if path.name == "ready.json": raise OSError("synthetic storage interruption") return publish(path, *args, **kwargs) with patch.object(handoff, "atomic_json", side_effect=interrupt_ready), self.assertRaises(OSError): run() self.assertFalse(repository._jobs) self.assertFalse((root / "executor/requests" / REQUEST_ID / "prepared.json").exists()) captures = hashes(root / "executor/captures") package = next((root / "executor/handoffs").iterdir()) / "frozen" frozen = hashes(package) self.assertEqual(run().status, "succeeded") self.assertEqual(run().status, "succeeded") self.assertEqual(hashes(root / "executor/captures"), captures) self.assertEqual(hashes(package), frozen) self.assertEqual((len(factories), processor.calls), (1, 1)) self.assertEqual((adapter.calls, validator.calls), (2, 2)) self.assertEqual((len(repository._jobs), len(repository._callbacks), len(repository._versions)), (1, 1, 1)) if __name__ == "__main__": unittest.main()