"""Date-based direct OHIP data -> immutable processing -> existing Finance/review flow.""" from datetime import date from pathlib import Path from arr_web.arr_download_handoff import outcome_from_handoff from arr_web.arr_downloads import DownloadOutcome, validate_request_id from integrations.ohip.arr_data import ARRDataSource, VERSION as SOURCE_VERSION 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, fingerprint, job_lock, private_directory, sync_directory from integrations.ohip.collect_arr_source import APPLICATION, SERVICE, require, strict_json class DirectARRExecutor: def __init__(self, *, root, source: ARRDataSource, policy, object_store, repository, ingestion, processor=None): self.root = Path(root).absolute() require(not self.root.resolve().is_relative_to(Path(__file__).resolve().parents[1]), "data_executor_store_must_be_outside_repository") self.source, self.policy, self.processor = source, policy, processor self.object_store, self.repository, self.ingestion = object_store, repository, ingestion def source_identity(self): return {"version": "arr-direct-data-executor/v1", "source_version": SOURCE_VERSION, "hotel_id": self.source.hotel_id, "service_url": SERVICE, "application_id": APPLICATION, "source_kind": "ohip_platform" if self.source.transport_factory is None else "test_transport", "page_size": self.source.page_size, "max_pages": self.source.max_pages, "max_records": self.source.max_records, "max_requests": self.source.max_requests, "processor_version": self.policy.processor_version, "rule_set_sha256": self.policy.rule_set_sha256} def execute(self, *, request_id: str, from_date: date, to_date: date, report_stage): validate_request_id(request_id) require(type(from_date) is date and type(to_date) is date and from_date == to_date, "single_report_date_required") identity = {**self.source_identity(), "request_id": request_id, "report_date": from_date.isoformat()} for folder in (self.root, self.root / "requests", self.root / "handoffs"): private_directory(folder) directory = self.root / "requests" / request_id private_directory(directory) sync_directory(directory.parent) with job_lock(directory): request_file, prepared_file = directory / "request.json", directory / "prepared.json" if not request_file.exists(): require(not prepared_file.exists(), "orphan_data_checkpoint") atomic_json(request_file, identity, replace=False) require(fingerprint(strict_json(protected_read(request_file, 65536))) == fingerprint(identity), "data_executor_request_conflict") if prepared_file.exists(): prepared = strict_json(protected_read(prepared_file, 65536)) require(fingerprint(prepared.get("identity")) == fingerprint(identity), "data_checkpoint_conflict") binding = handoff.CaptureBinding(**prepared["binding"]) binding.validate() require((binding.batch_id, binding.hotel_id, binding.arrival_date, binding.adapter_contract) == ("data-" + request_id, self.source.hotel_id, from_date.isoformat(), SOURCE_VERSION), "data_checkpoint_binding_mismatch") else: report_stage("downloading") collected = self.source.fetch(from_date.isoformat(), request_id) if collected.get("collection_complete") is not True: return DownloadOutcome("failed", retryable=True) payload = protected_read(Path(collected["data_path"]), 100 * 1024 * 1024) require(handoff._hash(payload) == collected["data_sha256"], "data_source_changed") binding = handoff.CaptureBinding("data-" + request_id, self.source.hotel_id, from_date.isoformat(), collected["manifest_sha256"], SOURCE_VERSION) report_stage("processing") result = handoff.prepare(self.root / "handoffs", binding, "ARR.json", payload, self.policy, processor=self.processor, source_role="source_data") prepared = {"identity": identity, "binding": vars(binding), "manifest_sha256": result["manifest_sha256"]} atomic_json(prepared_file, prepared, replace=False) report_stage("processing") 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"], expected_version=handoff.DATA_VERSION)