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

111 lines
7.3 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
from arr_web.arr_data_review import DataFieldReviews
self.data_reviews = DataFieldReviews(root=self.root / "data-reviews", policy=policy,
context=self.source_identity())
def get_data_review(self, request_id):
return self.data_reviews.get(request_id)
def get_acquisition_progress(self, *, request_id, report_date):
return self.source.get_acquisition_progress(request_id=request_id, report_date=report_date)
def update_data_review_item(self, request_id, item_id, revision, value, actor):
return self.data_reviews.update(request_id, item_id, revision, value, actor)
def finalize_data_review(self, request_id, revision, actor):
return self.data_reviews.finalize(request_id, revision, actor)
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:
collected_file = directory / "collected.json"
if collected_file.exists():
checkpoint = strict_json(protected_read(collected_file, 65536))
require(fingerprint(checkpoint["identity"]) == fingerprint(identity), "data_collection_checkpoint_conflict")
payload, original_manifest = self.data_reviews.original(request_id)
require(handoff._hash(payload) == checkpoint["data_sha256"]
and original_manifest == checkpoint["manifest_sha256"], "data_collection_checkpoint_changed")
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")
original_manifest = collected["manifest_sha256"]
review = self.data_reviews.prepare(request_id, payload, original_manifest,
from_date.isoformat())
if not collected_file.exists():
atomic_json(collected_file, {"identity": identity, "data_sha256": handoff._hash(payload),
"manifest_sha256": original_manifest}, replace=False)
manifest_sha256 = original_manifest
if review:
completed = self.data_reviews.payload(request_id)
if completed is None:
return DownloadOutcome("needs_data_review")
payload, manifest_sha256 = completed
binding = handoff.CaptureBinding("data-" + request_id, self.source.hotel_id,
from_date.isoformat(), 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)