419 lines
20 KiB
Python
419 lines
20 KiB
Python
"""Acquire ARR candidate data with two read-only OHIP operations.
|
|
|
|
This utility does not produce a report, select/prioritize rows, compute prices,
|
|
or submit anything to ARR/Finance. A complete capture is not report equivalence.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
from dataclasses import dataclass, field
|
|
from datetime import date, datetime, timezone
|
|
from email.utils import parsedate_to_datetime
|
|
import hashlib
|
|
import json
|
|
import os
|
|
from pathlib import Path
|
|
import re
|
|
import stat
|
|
import sys
|
|
import time
|
|
from typing import Callable
|
|
import urllib.error
|
|
import urllib.parse
|
|
import urllib.request
|
|
|
|
|
|
SERVICE = "https://ohip.nianxx.cn"
|
|
APPLICATION = "caller_zloxalzsQ9pYsztt"
|
|
SEARCH = "searchHotelReservations"
|
|
DETAIL = "getReservation"
|
|
FETCH = ("Reservation", "Comments", "Traces", "DailySummary", "RateInfoDetails",
|
|
"Packages", "InventoryItems", "Shares")
|
|
MAX_RESPONSE_BYTES = 8 * 1024 * 1024
|
|
MAX_ARCHIVE_BYTES = 1024 * 1024 * 1024
|
|
MAX_MANIFEST_BYTES = 32 * 1024 * 1024
|
|
# Leave room for the v2 assessment footer even when a response exceeds budget.
|
|
MAX_FINAL_METADATA_BYTES = MAX_RESPONSE_BYTES + 1
|
|
MANIFEST_SUMMARY_RESERVE = 64 * 1024
|
|
MANIFEST_FOOTER_ENTRY_RESERVE = 512
|
|
RETRY_STATUSES = {429, 500, 502, 503, 504}
|
|
|
|
|
|
class CollectionError(Exception):
|
|
"""Only static, non-sensitive error codes may be raised through this class."""
|
|
|
|
|
|
def require(condition: bool, code: str) -> None:
|
|
if not condition:
|
|
raise CollectionError(code)
|
|
|
|
|
|
def utc_now() -> str:
|
|
return datetime.now(timezone.utc).isoformat()
|
|
|
|
|
|
def json_bytes(value: object) -> bytes:
|
|
return (json.dumps(value, ensure_ascii=False, indent=2, allow_nan=False) + "\n").encode()
|
|
|
|
|
|
def strict_json(raw: bytes) -> dict:
|
|
def unique(pairs):
|
|
result = {}
|
|
for key, value in pairs:
|
|
require(key not in result, "duplicate_json_key")
|
|
result[key] = value
|
|
return result
|
|
|
|
def invalid_constant(_):
|
|
raise CollectionError("invalid_json_number")
|
|
|
|
try:
|
|
value = json.loads(raw, object_pairs_hook=unique, parse_constant=invalid_constant)
|
|
except (ValueError, UnicodeError, RecursionError):
|
|
raise CollectionError("invalid_json") from None
|
|
require(isinstance(value, dict), "json_not_object")
|
|
return value
|
|
|
|
|
|
def load_key(path: Path) -> str:
|
|
"""Read a regular, owner-only credential without exposing its contents."""
|
|
try:
|
|
parent = path.parent.lstat()
|
|
require(stat.S_ISDIR(parent.st_mode) and parent.st_uid == os.getuid()
|
|
and stat.S_IMODE(parent.st_mode) == 0o700, "credential_directory_permissions")
|
|
fd = os.open(path, os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK)
|
|
with os.fdopen(fd, "rb") as source:
|
|
info = os.fstat(source.fileno())
|
|
require(stat.S_ISREG(info.st_mode) and info.st_uid == os.getuid()
|
|
and stat.S_IMODE(info.st_mode) == 0o600, "credential_file_permissions")
|
|
raw = source.read(65537)
|
|
require(len(raw) <= 65536, "credential_too_large")
|
|
document = strict_json(raw)
|
|
require(document.get("version") == "ohip.application-credential/v1"
|
|
and document.get("status") == "ready"
|
|
and document.get("service_url", "").rstrip("/") == SERVICE
|
|
and document.get("application_id") == APPLICATION, "credential_identity_mismatch")
|
|
key = document.get("value")
|
|
require(isinstance(key, str) and bool(re.fullmatch(r"[\x21-\x7e]{1,4096}", key)),
|
|
"invalid_credential_value")
|
|
return key
|
|
except (OSError, TypeError, AttributeError):
|
|
raise CollectionError("credential_unreadable") from None
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class Options:
|
|
arrival_date: str
|
|
hotel_id: str
|
|
page_size: int = 100
|
|
max_pages: int = 100
|
|
max_records: int = 10000
|
|
|
|
def validate(self) -> None:
|
|
try:
|
|
require(date.fromisoformat(self.arrival_date).isoformat() == self.arrival_date,
|
|
"invalid_arrival_date")
|
|
except (TypeError, ValueError):
|
|
raise CollectionError("invalid_arrival_date") from None
|
|
require(isinstance(self.hotel_id, str) and bool(re.fullmatch(r"[A-Za-z0-9_-]{1,64}", self.hotel_id)),
|
|
"invalid_hotel_id")
|
|
for value, maximum in ((self.page_size, 100), (self.max_pages, 100), (self.max_records, 10000)):
|
|
require(type(value) is int and 1 <= value <= maximum, "invalid_collection_limit")
|
|
|
|
|
|
class Archive:
|
|
"""New private directory; exclusive file creation prevents accidental reuse."""
|
|
|
|
def __init__(self, path: Path):
|
|
self.path = path
|
|
self.files: list[dict] = []
|
|
self.inventory_bytes = 0
|
|
self.manifest_inventory_bytes = 0
|
|
path.mkdir(mode=0o700, parents=False, exist_ok=False)
|
|
os.chmod(path, 0o700)
|
|
|
|
def write(self, name: str, value: object) -> None:
|
|
require(bool(re.fullmatch(r"[a-z0-9_.-]+", name)), "invalid_archive_name")
|
|
raw = value if isinstance(value, bytes) else json_bytes(value)
|
|
entry = {"name": name, "bytes": len(raw), "sha256": hashlib.sha256(raw).hexdigest()}
|
|
# An inventory entry gains extra indentation at result.files depth.
|
|
# This conservative bound avoids repeatedly serializing the entire list.
|
|
entry_budget = len(json_bytes(entry)) + 64
|
|
if name == "result.json":
|
|
# The separately pinned manifest is outside its own file inventory.
|
|
require(len(raw) <= MAX_MANIFEST_BYTES, "archive_manifest_too_large")
|
|
elif name == "rate-assessments.json":
|
|
require(len(raw) <= MAX_FINAL_METADATA_BYTES, "archive_footer_too_large")
|
|
require(self.inventory_bytes + len(raw) <= MAX_ARCHIVE_BYTES, "archive_byte_budget_exceeded")
|
|
else:
|
|
require(self.inventory_bytes + len(raw) <= MAX_ARCHIVE_BYTES - MAX_FINAL_METADATA_BYTES,
|
|
"archive_byte_budget_exceeded")
|
|
if name != "result.json":
|
|
reserve = MANIFEST_SUMMARY_RESERVE + (0 if name == "rate-assessments.json" else MANIFEST_FOOTER_ENTRY_RESERVE)
|
|
require(self.manifest_inventory_bytes + entry_budget <= MAX_MANIFEST_BYTES - reserve,
|
|
"archive_manifest_budget_exceeded")
|
|
fd = os.open(self.path / name, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600)
|
|
with os.fdopen(fd, "wb") as output:
|
|
output.write(raw)
|
|
output.flush()
|
|
os.fsync(output.fileno())
|
|
self.files.append(entry)
|
|
if name != "result.json":
|
|
self.inventory_bytes += len(raw)
|
|
self.manifest_inventory_bytes += entry_budget
|
|
|
|
|
|
class NoRedirect(urllib.request.HTTPRedirectHandler):
|
|
def redirect_request(self, req, fp, code, msg, headers, newurl):
|
|
raise CollectionError("redirect_refused")
|
|
|
|
|
|
@dataclass
|
|
class HTTPTransport:
|
|
key: str = field(repr=False)
|
|
timeout: float = 25
|
|
opener: object = field(default_factory=lambda: urllib.request.build_opener(NoRedirect()), repr=False)
|
|
|
|
def __call__(self, method: str, path: str, body: bytes | None):
|
|
request = urllib.request.Request(SERVICE + path, data=body, method=method, headers={
|
|
"X-API-Key": self.key, "Accept": "application/json", "Content-Type": "application/json",
|
|
})
|
|
try:
|
|
response = self.opener.open(request, timeout=self.timeout)
|
|
except urllib.error.HTTPError as error:
|
|
response = error
|
|
with response:
|
|
return response.code, dict(response.headers), response.read(MAX_RESPONSE_BYTES + 1)
|
|
|
|
|
|
def retry_delay(headers: dict, attempt: int) -> float:
|
|
value = headers.get("retry-after")
|
|
if value is None:
|
|
return float(2 ** (attempt - 1))
|
|
try:
|
|
if re.fullmatch(r"[0-9]+", value):
|
|
delay = float(value)
|
|
else:
|
|
when = parsedate_to_datetime(value)
|
|
require(when.tzinfo is not None, "invalid_retry_after")
|
|
delay = max(0.0, (when - datetime.now(timezone.utc)).total_seconds())
|
|
# Stop instead of retrying sooner than the server permits.
|
|
require(0 <= delay <= 30, "retry_after_exceeds_budget")
|
|
return delay
|
|
except (ValueError, TypeError, OverflowError):
|
|
raise CollectionError("invalid_retry_after") from None
|
|
|
|
|
|
def check_warnings(value: object) -> None:
|
|
if isinstance(value, dict):
|
|
for key, child in value.items():
|
|
require(key not in {"warnings", "errors"} or not child, "upstream_warning_or_error")
|
|
check_warnings(child)
|
|
elif isinstance(value, list):
|
|
for child in value:
|
|
check_warnings(child)
|
|
|
|
|
|
class Reader:
|
|
def __init__(self, archive: Archive, hotel_id: str, transport: Callable,
|
|
*, key: str = "", sleep: Callable = time.sleep):
|
|
self.archive, self.hotel_id, self.transport = archive, hotel_id, transport
|
|
self._key, self.sleep = key, sleep
|
|
self.request_count = 0
|
|
|
|
def read(self, operation: str, *, body: dict | None = None, identity: str | None = None) -> dict:
|
|
if operation == SEARCH:
|
|
require(isinstance(body, dict) and identity is None, "invalid_search_request")
|
|
method, path, payload = "POST", "/api/v1/reservations/searches", json_bytes(body)
|
|
elif operation == DETAIL:
|
|
require(body is None and isinstance(identity, str)
|
|
and bool(re.fullmatch(r"[A-Za-z0-9_-]{1,128}", identity)), "invalid_reservation_id")
|
|
query = urllib.parse.urlencode([("fetchInstructions", value) for value in FETCH])
|
|
method, path, payload = "GET", "/api/v1/reservations/" + identity + "?" + query, None
|
|
else:
|
|
raise CollectionError("operation_not_allowed")
|
|
|
|
for attempt in range(1, 4):
|
|
self.request_count += 1
|
|
label = f"request-{self.request_count:06d}"
|
|
self.archive.write(label + ".json", {"operation_id": operation, "method": method,
|
|
"path": path, "body": body, "attempt": attempt, "started_at": utc_now()})
|
|
try:
|
|
status, headers, raw = self.transport(method, path, payload)
|
|
except (urllib.error.URLError, TimeoutError, ConnectionError):
|
|
self.archive.write(label + ".meta.json", {"error": "transport_failure", "ended_at": utc_now()})
|
|
if attempt == 3:
|
|
raise CollectionError("transport_retry_exhausted") from None
|
|
self.sleep(float(2 ** (attempt - 1)))
|
|
continue
|
|
headers = {k.lower(): v for k, v in headers.items()}
|
|
require(not self._key or self._key.encode() not in raw, "secret_in_response")
|
|
self.archive.write(label + ".response.bin", raw)
|
|
self.archive.write(label + ".meta.json", {"http_status": status, "ended_at": utc_now(),
|
|
"oversized": len(raw) > MAX_RESPONSE_BYTES})
|
|
require(len(raw) <= MAX_RESPONSE_BYTES, "response_too_large")
|
|
if status in RETRY_STATUSES and attempt < 3:
|
|
self.sleep(retry_delay(headers, attempt))
|
|
continue
|
|
require(status == 200, "http_failure")
|
|
document = strict_json(raw)
|
|
require(document.get("operation_id") == operation, "operation_mismatch")
|
|
require(document.get("hotel_id") == self.hotel_id, "hotel_mismatch")
|
|
require(isinstance(document.get("data"), dict), "invalid_data_envelope")
|
|
check_warnings(document)
|
|
return document["data"]
|
|
raise CollectionError("retry_exhausted")
|
|
|
|
|
|
def reservation_id(row: dict) -> str:
|
|
require(isinstance(row, dict), "invalid_reservation")
|
|
values = row.get("reservationIdList")
|
|
require(isinstance(values, list) and all(isinstance(value, dict) for value in values), "invalid_identity_list")
|
|
identities = [value.get("id") for value in values if value.get("type") == "Reservation"]
|
|
require(len(identities) == 1 and isinstance(identities[0], str)
|
|
and bool(re.fullmatch(r"[A-Za-z0-9_-]{1,128}", identities[0])), "ambiguous_reservation_identity")
|
|
return identities[0]
|
|
|
|
|
|
def validate_row(row: dict, options: Options) -> str:
|
|
identity = reservation_id(row)
|
|
require(row.get("hotelId") == options.hotel_id, "record_hotel_mismatch")
|
|
stay = row.get("roomStay")
|
|
require(isinstance(stay, dict) and stay.get("arrivalDate") == options.arrival_date,
|
|
"record_arrival_mismatch")
|
|
return identity
|
|
|
|
|
|
def integer(value: object, code: str) -> int:
|
|
require(type(value) is int and value >= 0, code)
|
|
return value
|
|
|
|
|
|
def validate_note_counts(search: dict, detail: dict) -> None:
|
|
"""Check explicit source indicators; absence is not inferred to mean zero."""
|
|
indicators = search.get("reservationIndicators", [])
|
|
require(isinstance(indicators, list) and all(isinstance(item, dict) for item in indicators),
|
|
"invalid_reservation_indicators")
|
|
for name, block in (("COMMENT", "comments"), ("TRACE", "traces")):
|
|
values = detail.get(block, [])
|
|
require(isinstance(values, list) and all(isinstance(item, dict) for item in values),
|
|
"invalid_note_or_trace_array")
|
|
expected = [item.get("count") for item in indicators if item.get("indicatorName") == name]
|
|
require(len(expected) <= 1, "ambiguous_reservation_indicator")
|
|
if expected:
|
|
require(integer(expected[0], "invalid_indicator_count") == len(values), "note_or_trace_count_mismatch")
|
|
|
|
|
|
def search_day(reader: Reader, options: Options) -> list[dict]:
|
|
records, seen = [], set()
|
|
expected_counts = None
|
|
for page_number in range(options.max_pages):
|
|
offset = page_number * options.page_size
|
|
data = reader.read(SEARCH, body={"arrivalStartDate": options.arrival_date,
|
|
"arrivalEndDate": options.arrival_date, "limit": options.page_size, "offset": offset,
|
|
"orderBy": ["ConfirmationNo"], "sortOrder": ["Asc"]})
|
|
page = data.get("reservations")
|
|
require(isinstance(page, dict), "invalid_search_page")
|
|
rows = page.get("reservationInfo")
|
|
require(isinstance(rows, list) and len(rows) <= options.page_size, "invalid_search_rows")
|
|
has_more = page.get("hasMore", False) # Oracle documents absence as all rows fetched.
|
|
require(type(has_more) is bool, "invalid_has_more")
|
|
for name, expected in (("offset", offset + options.page_size), ("limit", options.page_size), ("count", len(rows))):
|
|
if name in page:
|
|
require(integer(page[name], "invalid_pagination_integer") == expected, "pagination_mismatch")
|
|
counts = {name: integer(page[name], "invalid_total") for name in ("totalResults", "totalPages") if name in page}
|
|
if expected_counts is None:
|
|
expected_counts = counts
|
|
require(counts == expected_counts, "pagination_totals_changed")
|
|
if "totalResults" in counts:
|
|
require(counts["totalResults"] <= options.max_records, "record_limit_exceeded")
|
|
for row in rows:
|
|
identity = validate_row(row, options)
|
|
require(identity not in seen, "duplicate_reservation")
|
|
seen.add(identity)
|
|
records.append(row)
|
|
require(len(records) <= options.max_records, "record_limit_exceeded")
|
|
if not has_more:
|
|
if "totalResults" in counts:
|
|
require(counts["totalResults"] == len(records), "incomplete_search_total")
|
|
if "totalPages" in counts and records:
|
|
require(counts["totalPages"] == page_number + 1, "incomplete_search_pages")
|
|
require(bool(records), "empty_source_requires_review")
|
|
return records
|
|
require(len(rows) == options.page_size, "short_nonfinal_page")
|
|
if "totalResults" in counts:
|
|
require(len(records) < counts["totalResults"], "has_more_contradicts_total")
|
|
raise CollectionError("page_limit_exceeded")
|
|
|
|
|
|
def collect(options: Options, archive: Archive, reader: Reader) -> dict:
|
|
"""Fail closed; absence of a successful result.json also means incomplete."""
|
|
options.validate()
|
|
require(reader.archive is archive and reader.hotel_id == options.hotel_id, "reader_context_mismatch")
|
|
result = {"version": "arr-api-capture/v1", "status": "failed", "candidate_capture_complete": False,
|
|
"report_equivalence_verified": False, "finance_ready": False, "atomic_snapshot": False,
|
|
"hotel_id": options.hotel_id, "arrival_date": options.arrival_date,
|
|
"search_records": 0, "verified_details": 0, "search_recheck_equal": False}
|
|
archive.write("capture.json", {**result, "status": "started", "started_at": utc_now(),
|
|
"service_url": SERVICE, "application_id": APPLICATION, "options": vars(options),
|
|
"operations": [SEARCH, DETAIL], "fetch_instructions": FETCH})
|
|
try:
|
|
rows = search_day(reader, options)
|
|
result["search_records"] = len(rows)
|
|
for row in rows:
|
|
identity = reservation_id(row)
|
|
data = reader.read(DETAIL, identity=identity)
|
|
reservations = data.get("reservations")
|
|
require(isinstance(reservations, dict), "invalid_detail_envelope")
|
|
details = reservations.get("reservation")
|
|
require(isinstance(details, list) and len(details) == 1, "ambiguous_detail_count")
|
|
detail = details[0]
|
|
require(validate_row(detail, options) == identity, "detail_identity_mismatch")
|
|
for name in ("reservationStatus", "lastModifyDateTime"):
|
|
require(isinstance(row.get(name), str) and bool(row[name]) and detail.get(name) == row[name],
|
|
"search_detail_state_mismatch")
|
|
validate_note_counts(row, detail)
|
|
result["verified_details"] += 1
|
|
rechecked = search_day(reader, options)
|
|
require(rechecked == rows, "source_changed_during_collection")
|
|
result.update(status="complete_candidate_capture", candidate_capture_complete=True, search_recheck_equal=True)
|
|
except CollectionError as error:
|
|
result["error"] = str(error)
|
|
except Exception:
|
|
# Never expose exception strings, response bodies or credential-bearing requests.
|
|
result["error"] = "unexpected_collection_failure"
|
|
finally:
|
|
result.update(completed_at=utc_now(), http_attempts=reader.request_count)
|
|
archive.write("result.json", {**result, "files": list(archive.files)})
|
|
# Return a separate root pin; result.json cannot contain its own digest.
|
|
result["manifest_sha256"] = archive.files[-1]["sha256"]
|
|
return result
|
|
|
|
|
|
def main(argv: list[str] | None = None) -> int:
|
|
parser = argparse.ArgumentParser(description=__doc__)
|
|
parser.add_argument("--arrival-date", required=True, help="Explicit YYYY-MM-DD; no implicit timezone or business date")
|
|
parser.add_argument("--hotel-id", required=True, help="Expected server-selected hotel; does not change server routing")
|
|
parser.add_argument("--credential-file", required=True, type=Path)
|
|
parser.add_argument("--output-dir", required=True, type=Path, help="New private directory outside the repository")
|
|
args = parser.parse_args(argv)
|
|
try:
|
|
options = Options(args.arrival_date, args.hotel_id)
|
|
options.validate()
|
|
repository = Path(__file__).resolve().parents[2]
|
|
require(not args.output_dir.resolve().is_relative_to(repository), "output_must_be_outside_repository")
|
|
key = load_key(args.credential_file)
|
|
archive = Archive(args.output_dir)
|
|
reader = Reader(archive, options.hotel_id, HTTPTransport(key), key=key)
|
|
result = collect(options, archive, reader)
|
|
except CollectionError as error:
|
|
result = {"status": "failed", "error": str(error), "finance_ready": False}
|
|
except Exception:
|
|
result = {"status": "failed", "error": "local_setup_or_archive_failure", "finance_ready": False}
|
|
print(json.dumps(result, ensure_ascii=False))
|
|
return 0 if result.get("candidate_capture_complete") else 1
|
|
|
|
|
|
if __name__ == "__main__":
|
|
sys.exit(main())
|