Files
ARR-2.0-0918/arr_web/arr_data_executor.py

76 lines
5.1 KiB
Python

"""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)