Files
ARR-2.0-0918/integrations/ohip/profile_supplement.py

263 lines
13 KiB
Python

"""Acquire primary-name candidates independently of an immutable complete v2 day.
Only missing profile summaries are requested. An exhausted lookup closes this
attempt's network path; all rows remain represented as candidates or gaps. Each
retry has its own archive, and reused responses require the pinned prior chain.
Nothing here accepts ARR display semantics or changes the main capture.
"""
from __future__ import annotations
import argparse
from collections import Counter
import hashlib
import json
import os
from pathlib import Path
import re
import stat
import urllib.error
from . import audit_arr_capture as integrity, audit_arr_day as base_audit
from . import collect_arr_source as source, profile_reader, profile_summary
VERSION = "arr-profile-supplement/v1"
FLAGS = {"finance_ready": False, "report_equivalence_verified": False, "complete_arr_output": False,
"atomic_snapshot": False}
FILES = re.compile(r"(?:capture\.json|name-assessments\.json|profile-[0-9]{6}\.(?:json|meta\.json|response\.bin))")
LOOKUP_FAILURES = {"profile_http_failure", "profile_transport_retry_exhausted", "profile_limit_exceeded"}
require = source.require
def load_base(directory, pin):
archive = base_audit.VerifiedArchive(Path(directory), pin)
_, details, _ = base_audit.replay(archive)
return archive, details
def binding(base):
return {"base_manifest_sha256": base.pin, "hotel_id": base.options.hotel_id,
**{key: getattr(base.options, key) for key in ("from_date", "to_date", "rate_date")}}
def _config(base, max_profiles, prior_pins):
require(type(max_profiles) is int and 1 <= max_profiles <= base.options.max_records, "invalid_profile_limit")
require(len(prior_pins) < 100, "supplement_chain_limit_exceeded")
return {"version": VERSION, **binding(base), "max_profiles": max_profiles, "prior_pins": prior_pins,
"service_url": source.SERVICE, "application_id": source.APPLICATION,
"operations": [profile_summary.OPERATION]}
class _CachedReader:
def __init__(self, reader, cache, round_index):
self.reader, self.hotel_id = reader, reader.hotel_id
self.cache, self.round_index, self.closed = dict(cache), round_index, False
self.reused = 0
def read(self, profile_id):
if profile_id in self.cache:
raw, origin, request = self.cache[profile_id]
self.reused += int(origin < self.round_index)
ref = request if origin == self.round_index else f"prior-{origin:04d}/{request}"
return raw, ref
if self.closed:
raise source.CollectionError("profile_lookup_deferred")
try:
raw, request = self.reader.read(profile_id)
except source.CollectionError as error:
if str(error) in LOOKUP_FAILURES:
self.closed = True
raise
self.cache[profile_id] = (raw, self.round_index, request)
return raw, request
def _collect(base, details, sink, reader, max_profiles, prior_pins, cache):
config = _config(base, max_profiles, prior_pins)
require(reader.archive is sink and reader.hotel_id == base.options.hotel_id
and reader.request_count == 0 and reader.unique_profiles == 0
and reader.max_profiles == max_profiles, "profile_reader_context_mismatch")
sink.write("capture.json", config)
cached = _CachedReader(reader, cache, len(prior_pins))
records, issues, reusable = [], Counter(), set()
result = {"version": VERSION, **binding(base), **FLAGS, "status": "failed",
"main_capture_complete": True, "all_names_valid": False, "records": len(details)}
try:
for detail in details:
try:
profile = profile_reader.assess(detail, cached)
except source.CollectionError as error:
if str(error) not in LOOKUP_FAILURES | {"profile_lookup_deferred"}:
raise
# Error bodies and IDs stay in the private originals. Failed
# requests never masquerade as successful field provenance.
profile = {"error": str(error)}
records.append({"reservation_id": source.reservation_id(detail), "profile": profile})
if "error" in profile:
issues[profile["error"]] += 1
else:
reusable.add(profile["profile_id"])
result.update(status="name_gaps" if issues else "complete_name_candidates", all_names_valid=not issues)
except source.CollectionError as error:
result["error"] = str(error)
except Exception:
result["error"] = "unexpected_supplement_failure"
result.update(assessed_records=len(records), valid_name_candidates=sum("error" not in row["profile"] for row in records),
profile_issues=dict(sorted(issues.items())), profile_http_attempts=reader.request_count,
newly_acquired_profiles=reader.unique_profiles, reused_name_records=cached.reused)
sink.write("name-assessments.json", {"records": records})
sink.write("result.json", {**result, "files": list(sink.files)})
result["manifest_sha256"] = sink.files[-1]["sha256"]
# A200 without a usable name is a retryable gap, not a permanent cache hit.
return result, records, {key: value for key, value in cached.cache.items() if key in reusable}
class _Archive:
"""Verify supplement bytes; protocol replay below authenticates their meaning."""
read = integrity.VerifiedArchive.read
document = integrity.VerifiedArchive.document
def __init__(self, directory, pin):
require(type(pin) is str and re.fullmatch(r"[0-9a-f]{64}", pin), "invalid_manifest_pin")
self.directory, self.pin = Path(directory), pin
info = self.directory.lstat()
require(stat.S_ISDIR(info.st_mode) and info.st_uid == os.getuid()
and stat.S_IMODE(info.st_mode) == 0o700, "unsafe_archive_directory")
raw = integrity.protected_read(self.directory / "result.json", source.MAX_MANIFEST_BYTES)
require(hashlib.sha256(raw).hexdigest() == pin, "manifest_hash_mismatch")
self.result = source.strict_json(raw)
require(self.result.get("version") == VERSION and self.result.get("status") in
{"name_gaps", "complete_name_candidates"}, "unusable_profile_supplement")
entries = self.result.get("files")
require(type(entries) is list and 2 <= len(entries) <= 90002, "invalid_file_inventory")
self.inventory, total = {}, 0
for entry in entries:
require(type(entry) is dict and set(entry) == {"name", "bytes", "sha256"}, "invalid_file_entry")
name = entry["name"]
require(type(name) is str and FILES.fullmatch(name), "unsafe_inventory_name")
require(name not in self.inventory, "duplicate_inventory_name")
require(type(entry["bytes"]) is int and 0 <= entry["bytes"] <= source.MAX_RESPONSE_BYTES + 1,
"invalid_inventory_size")
require(type(entry["sha256"]) is str and re.fullmatch(r"[0-9a-f]{64}", entry["sha256"]),
"invalid_inventory_hash")
self.inventory[name] = entry
total += entry["bytes"]
require(total <= source.MAX_ARCHIVE_BYTES, "archive_byte_budget_exceeded")
require({p.name for p in self.directory.iterdir()} == set(self.inventory) | {"result.json"},
"archive_file_set_mismatch")
for name in self.inventory:
self.read(name)
self.capture = self.document("capture.json")
def _replay(base, details, archive, prior_pins, cache):
config = archive.capture
require(base_audit._same_json(config, _config(base, config.get("max_profiles"), prior_pins)),
"supplement_context_mismatch")
sink = base_audit.MemorySink()
counter, consumed = 0, {"capture.json", "name-assessments.json"}
def transport(method, path, body):
nonlocal counter
counter += 1
label = f"profile-{counter:06d}"
request_name, meta_name, response_name = label + ".json", label + ".meta.json", label + ".response.bin"
request, meta = archive.document(request_name), archive.document(meta_name)
consumed.update((request_name, meta_name))
require(request.get("operation_id") == profile_summary.OPERATION and request.get("method") == method
and request.get("path") == path and base_audit._same_json(request.get("body"), source.strict_json(body))
and type(request.get("attempt")) is int
and request["attempt"] == sink.latest_request["attempt"], "archived_request_mismatch")
if meta.get("error") == "transport_failure":
require(response_name not in archive.inventory, "unexpected_transport_response")
raise urllib.error.URLError("archived_transport_failure")
raw = archive.read(response_name)
consumed.add(response_name)
require(type(meta.get("http_status")) is int and meta.get("oversized") is False, "invalid_response_metadata")
headers = {}
if meta.get("retry_after") is not None:
require(type(meta["retry_after"]) is str, "invalid_archived_retry_after")
headers["retry-after"] = meta["retry_after"]
return meta["http_status"], headers, raw
reader = profile_reader.ProfileSummaryReader(sink, base.options.hotel_id, transport,
max_profiles=config["max_profiles"], sleep=lambda _: None)
result, records, cache = _collect(base, details, sink, reader, config["max_profiles"], prior_pins, cache)
require(result["status"] != "failed", "supplement_protocol_replay_failed")
require(consumed == set(archive.inventory), "unconsumed_archive_files")
require([item["name"] for item in sink.files[:-1]] == list(archive.inventory), "archive_sequence_mismatch")
actual = {key: value for key, value in result.items() if key != "manifest_sha256"}
expected = {key: value for key, value in archive.result.items() if key != "files"}
require(base_audit._same_json(actual, expected), "supplement_summary_mismatch")
require(base_audit._same_json({"records": records}, archive.document("name-assessments.json")),
"supplement_assessment_mismatch")
return records, cache
def verify_chain(base, supplements):
"""Reopen and replay every supplied pin in order; no networking or credentials."""
require(type(base) is base_audit.VerifiedArchive, "supplement_requires_v2_base")
base, details = load_base(base.directory, base.pin)
require(len(supplements) <= 100, "supplement_chain_limit_exceeded")
pins, cache, records = [], {}, []
for directory, pin in supplements:
require(pin not in pins, "duplicate_supplement_pin")
archive = _Archive(directory, pin)
records, cache = _replay(base, details, archive, pins, cache)
pins.append(pin)
return records, cache
def run(base_directory, base_pin, output_directory, *, max_profiles, reader_factory, previous=()):
"""New immutable attempt. Credentials are loaded by the factory after validation."""
base, details = load_base(base_directory, base_pin)
pins = [pin for _, pin in previous]
_config(base, max_profiles, pins)
_, cache = verify_chain(base, previous)
output_directory = Path(output_directory)
require(not output_directory.resolve().is_relative_to(Path(__file__).resolve().parents[2]),
"output_must_be_outside_repository")
sink = source.Archive(output_directory)
reader = reader_factory(sink, base.options.hotel_id, max_profiles)
result, _, _ = _collect(base, details, sink, reader, max_profiles, pins, cache)
if result["status"] != "failed":
verify_chain(base, [*previous, (output_directory, result["manifest_sha256"])])
return result
def main(argv=None):
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--capture-dir", type=Path, required=True)
parser.add_argument("--capture-sha256", required=True)
parser.add_argument("--output-dir", type=Path, required=True)
parser.add_argument("--max-profiles", type=int, required=True)
parser.add_argument("--credential-file", type=Path, required=True)
parser.add_argument("--previous", nargs=2, action="append", default=[], metavar=("DIRECTORY", "SHA256"))
args = parser.parse_args(argv)
def factory(sink, hotel, limit):
# Lazily load the key: a fully reusable supplement needs no credentials.
transport = None
def request(method, path, body):
nonlocal transport
if transport is None:
key = source.load_key(args.credential_file)
transport = source.HTTPTransport(key, timeout=40)
reader._key = key
return transport(method, path, body)
reader = profile_reader.ProfileSummaryReader(sink, hotel, request, max_profiles=limit)
return reader
try:
result = run(args.capture_dir, args.capture_sha256, args.output_dir,
max_profiles=args.max_profiles, reader_factory=factory, previous=args.previous)
except source.CollectionError as error:
result = {"status": "failed", "error": str(error), **FLAGS}
except Exception:
result = {"status": "failed", "error": "profile_supplement_failed", **FLAGS}
print(json.dumps(result, ensure_ascii=False, sort_keys=True))
return 0 if result["status"] in {"name_gaps", "complete_name_candidates"} else 1
if __name__ == "__main__":
raise SystemExit(main())