76 lines
5.1 KiB
Python
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)
|