263 lines
13 KiB
Python
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())
|