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

165 lines
9.5 KiB
Python

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