165 lines
9.5 KiB
Python
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"])
|