"""Explicit-date capture → accepted adaptation → frozen processing orchestration. No production source adapter/validator, credentials or runtime enablement lives here. Both source functions must be supplied after independent source acceptance. The default Web application continues to use UnavailableARRDownloads. """ from __future__ import annotations import os from dataclasses import replace from datetime import date from pathlib import Path from typing import Callable, Protocol from arr_ingestion.repository import IngestionRepository from arr_ingestion.service import IngestionService from arr_ingestion.validation import ProcessorPolicy from arr_processing.local import LocalDailyProcessor from arr_storage.store import ManagedObjectStore from arr_web.arr_download_handoff import outcome_from_handoff from arr_web.arr_downloads import DownloadOutcome, validate_request_id from integrations.ohip import audit_arr_day as audit from integrations.ohip import capture_day_job as capture from integrations.ohip import collect_arr_day as day from integrations.ohip import audit_arr_named_day as named_audit from integrations.ohip import capture_named_day_job as named_capture from integrations.ohip import collect_arr_named_day as named_day from integrations.ohip import collect_arr_source as source from integrations.ohip import processing_handoff as handoff from integrations.ohip.audit_arr_capture import protected_read from integrations.ohip.capture_job import atomic_json, job_lock, private_directory, sync_directory VERSION = "arr-web-capture-executor/v1" NAMED_VERSION = "arr-web-capture-executor/v2" PROJECT_ROOT = Path(__file__).resolve().parents[1] class ARRSourceAdapter(Protocol): def adapt(self, archive: audit.VerifiedArchive | named_audit.VerifiedArchive) -> bytes: """Return deterministic XML under the separately accepted source contract.""" ... class ARRMappingValidator(Protocol): def validate(self, archive: audit.VerifiedArchive | named_audit.VerifiedArchive, payload: bytes) -> None: """Independently verify complete source mapping; raise on any mismatch.""" ... def _document(path: Path) -> dict: return source.strict_json(protected_read(path, 1024 * 1024)) class CapturedARRExecutor: def __init__( self, *, root: Path, hotel_id: str, adapter_contract: str, adapter: ARRSourceAdapter, mapping_validator: ARRMappingValidator, reader_factory: Callable, policy: ProcessorPolicy, object_store: ManagedObjectStore, repository: IngestionRepository, ingestion: IngestionService, processor: LocalDailyProcessor | None = None, capture_version: str = "v2", max_profiles: int | None = None, page_size: int = 100, max_pages: int = 100, max_records: int = 10000, ) -> None: # Merely possessing API credentials or a complete capture is insufficient. if not callable(getattr(adapter, "adapt", None)) or not callable(getattr(mapping_validator, "validate", None)): raise ValueError("an accepted source adapter and mapping validator are required") if not callable(reader_factory): raise ValueError("a capture reader factory is required") source.require(capture_version in ("v2", "v3"), "unsupported_executor_capture_version") named = capture_version == "v3" if not named: source.require(max_profiles is None, "profile_limit_requires_named_capture") # Validate acquisition bounds before creating state or loading credentials. options_type = named_day.Options if named else day.Options capture_options = options_type("2000-01-01", "2000-01-01", "2000-01-01", hotel_id, page_size=page_size, max_pages=max_pages, max_records=max_records, **({"max_profiles": max_profiles} if named else {})) capture_options.validate() handoff.CaptureBinding("validation", hotel_id, "2000-01-01", "0" * 64, adapter_contract).validate() root = Path(root).absolute() source.require(not root.resolve().is_relative_to(PROJECT_ROOT), "executor_store_must_be_outside_repository") self.root, self.hotel_id, self.adapter_contract = root, hotel_id, adapter_contract self.adapter, self.mapping_validator, self.reader_factory = adapter, mapping_validator, reader_factory self.policy, self.object_store, self.repository, self.ingestion = policy, object_store, repository, ingestion self.processor = processor self.capture_version, self.max_profiles = capture_version, max_profiles self._capture_options = capture_options def execute( self, *, request_id: str, from_date: date, to_date: date, report_stage: Callable[[str], None], ) -> DownloadOutcome: validate_request_id(request_id) if type(from_date) is not date or type(to_date) is not date or from_date != to_date: raise ValueError("explicit equal report dates are required") named = self.capture_version == "v3" capture_api, archive_api = (named_capture, named_audit) if named else (capture, audit) capture_protocol = named_capture.NamedDayProtocol if named else capture.DayProtocol options = replace(self._capture_options, from_date=from_date.isoformat(), to_date=to_date.isoformat(), rate_date=from_date.isoformat()) options.validate() batch_id = "web-" + request_id identity = { "version": NAMED_VERSION if named else VERSION, "request_id": request_id, "batch_id": batch_id, "capture": capture_protocol.identity(batch_id, options), "adapter_contract": self.adapter_contract, "processor_version": self.policy.processor_version, "rule_set_sha256": self.policy.rule_set_sha256, } for folder in (self.root, self.root / "requests", self.root / "captures", self.root / "handoffs"): private_directory(folder) directory = self.root / "requests" / request_id private_directory(directory) sync_directory(self.root.parent) sync_directory(self.root) sync_directory(directory.parent) with job_lock(directory): request_file, prepared_file = directory / "request.json", directory / "prepared.json" if not os.path.lexists(request_file): source.require(not os.path.lexists(prepared_file), "orphan_prepared_checkpoint") atomic_json(request_file, identity, replace=False) source.require(audit._same_json(_document(request_file), identity), "executor_request_conflict") if os.path.lexists(prepared_file): prepared = _document(prepared_file) source.require(audit._same_json(prepared.get("identity"), identity), "prepared_checkpoint_context_mismatch") binding = handoff.CaptureBinding(**prepared["binding"]) binding.validate() source.require((binding.batch_id, binding.hotel_id, binding.arrival_date, binding.adapter_contract) == (batch_id, self.hotel_id, from_date.isoformat(), self.adapter_contract), "prepared_checkpoint_binding_mismatch") # No capture or adaptation after a package has been pinned. report_stage("processing") else: report_stage("downloading") captured = capture_api.run_batch(self.root / "captures", batch_id, options, self.reader_factory) if captured.get("candidate_capture_complete") is not True: if captured.get("status") == "capture_failed": return DownloadOutcome("failed", retryable=True) raise RuntimeError("capture batch unavailable") source.require(captured.get("batch_id") == batch_id, "captured_batch_mismatch") archive = archive_api.VerifiedArchive(Path(captured["capture_dir"]), captured["manifest_sha256"]) source.require(archive.options == options, "captured_dates_or_hotel_mismatch") archive_api.replay(archive) payload = self.adapter.adapt(archive) source.require(type(payload) is bytes, "invalid_adapted_source") validation = self.mapping_validator.validate(archive, payload) source.require(validation is None, "invalid_mapping_validation_result") binding = handoff.CaptureBinding(batch_id, self.hotel_id, from_date.isoformat(), captured["manifest_sha256"], self.adapter_contract) report_stage("processing") result = handoff.prepare(self.root / "handoffs", binding, "ARR.XML", payload, self.policy, processor=self.processor) prepared = {"identity": identity, "binding": vars(binding), "manifest_sha256": result["manifest_sha256"]} # Flush the package reference BEFORE the first possible Finance/OSS write. atomic_json(prepared_file, prepared, replace=False) receipt = handoff.deliver(self.root / "handoffs" / binding.job_id, prepared["manifest_sha256"], object_store=self.object_store, repository=self.repository, service=self.ingestion, expected_binding=binding, expected_policy=self.policy) return outcome_from_handoff(receipt, job_id=binding.job_id, report_date=from_date, manifest_sha256=prepared["manifest_sha256"])