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

553 lines
31 KiB
Python

"""Date-driven OHIP data input. No XML, pricing, Finance writes or live startup calls.
The output is acquisition data: ordered reservations, fifteen field observations,
multi-value notes/packages and role-preserving related facts. Business filtering
and display choices belong to the consuming data processor.
"""
from __future__ import annotations
import argparse
from collections import Counter
import copy
from datetime import date
import hashlib
import json
import os
from pathlib import Path
import re
import stat
import time
from . import collect_arr_source as base
from . import profile_summary as profiles
from . import rate_info, room_calendar_evidence as calendar, source_fields as fields
from .audit_arr_capture import protected_read
from .capture_job import atomic_json, job_lock, private_directory, sync_directory
from .data_client import DataReader, document, json_bytes, typed_id
VERSION = "arr-ohip-data/v1"
FIELDS = ("BLOCK_CODE", "ADULTS", "CHILDREN", "COMPANY_NAME", "CONFIRMATION_NO", "DISP_ROOM_NO",
"EFFECTIVE_RATE_AMOUNT", "FULL_NAME", "RES_COMMENT", "NO_OF_ROOMS", "PRODUCTS", "RATE_CODE",
"ROOM_CATEGORY_LABEL", "ARRIVAL", "DEPARTURE")
OPTIONAL_FIELDS = {"BLOCK_CODE", "RES_COMMENT", "PRODUCTS", "ROOM_CATEGORY_LABEL"}
PROJECT_ROOT = Path(__file__).resolve().parents[2]
MAX_DATA_BYTES = 25 * 1024 * 1024
require = base.require
class FieldGap(base.CollectionError):
def __init__(self, code, state="missing"):
super().__init__(code)
self.state = state
def value(item):
empty = item == "" or item == [] or item == () or (type(item) is str and not item.strip())
return {"state": "empty" if empty else "available", "value": item}
def observe(select):
try:
return value(select())
except FieldGap as error:
return {"state": error.state, "value": None, "reason": str(error)}
except base.CollectionError as error:
code = str(error)
failed = code in {"http_failure", "http_permission_denied", "transport_retry_exhausted", "request_limit_exceeded",
"response_too_large", "operation_mismatch", "hotel_mismatch", "invalid_data_envelope",
"missing_oracle_request_id", "secret_in_response", "invalid_json", "upstream_warning_or_error"}
return {"state": "failed" if failed else "missing", "value": None, "reason": code}
def _string(obj, key, code):
require(type(obj) is dict and type(obj.get(key)) is str, code)
return obj[key]
def _day(detail, report_date):
index = fields.single_day_rate_index(detail, report_date)
return detail["roomStay"]["roomRates"][index]
def _profile(reader, identity):
result = reader.query("getProfile", identity=identity)
require(typed_id(result.data.get("profileIdList"), "Profile") == identity, "related_profile_identity_mismatch")
profile = result.data.get("profileDetails")
require(type(profile) is dict, "missing_profile_details")
return profile
def _name(detail, reader):
guests = detail.get("reservationGuests")
require(type(guests) is list and all(type(g) is dict for g in guests), "missing_reservation_guests")
require(all("primary" not in g or type(g["primary"]) is bool for g in guests), "invalid_primary_guest")
primary = [g for g in guests if g.get("primary") is True]
require(len(primary) == 1, "ambiguous_primary_guest")
info = primary[0].get("profileInfo")
require(type(info) is dict, "missing_primary_guest_profile")
identity = typed_id(info.get("profileIdList"), "Profile")
working = detail
try:
guest = profiles.primary_guest(working, include_prefix=True)
except base.CollectionError:
# Retrieve the identified guest's profile when the reservation embeds no
# full primary-name structure. No name search or surname-based join.
profile = _profile(reader, identity)
working = copy.deepcopy(detail)
target = next(g for g in working["reservationGuests"] if g.get("primary") is True)
target["profileInfo"]["profile"] = profile
guest = profiles.primary_guest(working, include_prefix=True)
summary = reader.query(profiles.OPERATION, identity=identity)
inspected = profiles.inspect_summary(summary.raw, guest, reader.hotel_id, include_prefix=True)
require(inspected["name_type"] in (None, "Primary"), "profile_nonprimary_summary_name")
require(all(v in {"equal", "both_absent"} for v in inspected["component_comparisons"].values()),
"profile_name_components_mismatch")
require(any(guest.name.get(k, "").strip() for k in ("surname", "givenName", "middleName")),
"profile_primary_name_unverifiable")
require(type(inspected["full_name"]) is str, "profile_full_name_unavailable")
return inspected["full_name"]
def _companies(detail, reader, related):
container = detail.get("reservationProfiles")
require(type(container) is dict and type(container.get("reservationProfile")) is list,
"missing_reservation_profiles")
top = container["reservationProfile"]
require(all(type(p) is dict for p in top), "invalid_reservation_profiles")
roles = {"Company", "TravelAgent", "Source", "Group"}
rows = []
for item in top:
require(type(item.get("reservationProfileType")) is str, "missing_profile_role")
role = item["reservationProfileType"]
if role not in roles:
continue
identity = typed_id(item.get("profileIdList"), "Profile")
profile = item.get("profile", {})
require(type(profile) is dict, "invalid_associated_profile")
company = profile.get("company", {})
require(type(company) is dict, "invalid_company_profile")
if type(company.get("companyName")) is not str:
profile = _profile(reader, identity)
company = profile.get("company")
name = _string(company, "companyName", "missing_company_name")
rows.append({"role": role, "profile_id": identity, "name": name})
related["associated_profiles"] = rows
# Compare an explicit arrival-day association set when it is provided.
day = _day(detail, reader.options.arrival_date)
if "stayProfiles" in day:
day_profiles = day["stayProfiles"]
require(type(day_profiles) is list and all(type(p) is dict for p in day_profiles), "invalid_day_profiles")
day_ids = [(p.get("reservationProfileType"), typed_id(p.get("profileIdList"), "Profile"))
for p in day_profiles if p.get("reservationProfileType") in roles]
require(Counter(day_ids) == Counter((r["role"], r["profile_id"]) for r in rows), "company_association_changed")
names = {(r["role"], r["profile_id"]): r["name"] for r in rows}
for p in day_profiles:
if p.get("reservationProfileType") not in roles:
continue
day_profile = p.get("profile", {})
require(type(day_profile) is dict and type(day_profile.get("company", {})) is dict, "invalid_day_company")
day_name = day_profile.get("company", {}).get("companyName")
if day_name is not None:
require(type(day_name) is str and names[(p["reservationProfileType"], typed_id(p["profileIdList"], "Profile"))] == day_name,
"company_name_changed")
eligible = [row["name"] for row in rows if row["role"] in {"Company", "TravelAgent"}]
names = set(eligible)
if len(names) > 1:
raise FieldGap("multiple_company_or_travel_agent_names", "ambiguous")
if not eligible and rows:
raise FieldGap("only_source_or_group_profiles_available", "ambiguous")
return eligible[0] if eligible else ""
def _block(search, detail, reader):
try:
return fields.agreed_block_code(search, detail, reader.options.arrival_date)
except base.CollectionError:
pass
day = _day(detail, reader.options.arrival_date)
block = day.get("reservationBlock")
require(type(block) is dict, "missing_reservation_block")
require("hotelId" not in block or block["hotelId"] == reader.hotel_id, "block_association_hotel_mismatch")
identifiers = block.get("blockIdList")
if identifiers == []:
search_block = search.get("roomStay", {}).get("reservationBlock")
require(type(search_block) is dict and search_block.get("blockIdList") == [], "block_association_changed")
return ""
identity = typed_id(identifiers, "Block")
# An explicit code disagreement must remain visible, not be overwritten by
# today's block details. Only the genuinely missing-code case is supplemented.
require(not any(i.get("type") == "BlockCode" for i in identifiers), "block_code_disagreement")
searched = search.get("roomStay", {}).get("reservationBlock")
require(type(searched) is dict and typed_id(searched.get("blockIdList"), "Block") == identity,
"block_association_changed")
require("hotelId" not in searched or searched["hotelId"] == reader.hotel_id, "block_association_hotel_mismatch")
result = reader.query("getBlock", identity=identity)
collection = result.data.get("blocks")
require(type(collection) is dict and collection.get("hasMore", False) is False, "incomplete_block_response")
items = collection.get("blockInfo")
require(type(items) is list and len(items) == 1 and type(items[0]) is dict, "ambiguous_block_response")
item = items[0].get("block")
require(type(item) is dict and typed_id(item.get("blockIdList"), "Block") == identity
and item.get("hotelId") == reader.hotel_id, "block_response_identity_mismatch")
code = _string(item.get("blockDetails"), "blockCode", "missing_block_code")
earlier_codes = [i.get("id") for i in searched["blockIdList"] if i.get("type") == "BlockCode"]
require(not earlier_codes or earlier_codes == [code], "block_code_changed")
return code
def _packages(detail, reader, related):
packages = detail.get("reservationPackages")
require(type(packages) is list and all(type(p) is dict for p in packages), "missing_reservation_packages")
result = []
related["reservation_packages"] = result
for item in packages:
code = _string(item, "packageCode", "missing_package_code")
require(0 < len(code) <= 20, "invalid_package_code")
entry = {"package": copy.deepcopy(item)}
result.append(entry)
if "scheduleList" in item:
require(type(item["scheduleList"]) is list and all(type(s) is dict for s in item["scheduleList"]),
"invalid_package_schedule")
continue
dates = fields.stay_dates(detail, reader.options.arrival_date)
fetched = reader.query("getPackage", identity=base.reservation_id(detail), params={
"productCode": code, "reservationTimeSpanStartDate": dates.arrival.isoformat(),
"reservationTimeSpanEndDate": dates.departure.isoformat()})
container = fetched.data.get("reservationPackages")
require(type(container) is dict and type(container.get("reservationPackage")) is list,
"missing_package_details")
supplemental = container["reservationPackage"]
require(bool(supplemental) and all(type(p) is dict and p.get("packageCode") == code for p in supplemental),
"package_response_code_mismatch")
if "reservationIdList" in fetched.data:
require(typed_id(fetched.data["reservationIdList"], "Reservation") == base.reservation_id(detail),
"package_response_reservation_mismatch")
# Keep all group/source/schedule variants. Combining them into a display
# string or selecting package consumption dates is a later business step.
entry["details"] = supplemental
return result
def _rate(detail, reader, related):
response = reader.query(rate_info.POST, identity=base.reservation_id(detail))
# The day is fixed in the request and the response is bound to that request.
selected = rate_info.day_rate_candidate(response.data)
day = _day(detail, reader.options.arrival_date)
rate_container = day.get("rates", {})
require(type(rate_container) is dict, "invalid_day_rates")
rates = rate_container.get("rate", [])
require(type(rates) is list, "invalid_day_rates")
currencies = {r["base"]["currencyCode"] for r in rates if type(r) is dict and type(r.get("base")) is dict
and type(r["base"].get("currencyCode")) is str}
if currencies:
require(currencies == {selected.currency}, "rate_currency_mismatch")
related["currency"] = selected.currency
return str(selected.effective_rate)
def _calendar(reader, identities):
params = {"startDate": reader.options.arrival_date, "endDate": reader.options.arrival_date,
"includeRoomMoveHistory": "true", "showRoomMoveSegments": "true"}
collected = {identity: [] for identity in identities}
room_ids, total, page_size = set(), None, None
for _ in range(reader.options.max_pages):
response = reader.query("getRoomCalendar", params=params)
inspected = calendar.inspect(response.request, 200, response.raw, hotel_id=reader.hotel_id,
report_date=reader.options.arrival_date, reservation_ids=identities)
require(inspected["room_collection_state"] in {"present", "explicit_empty"}, "calendar_missing_rooms")
meta = inspected["pagination"]
require("totalRooms" in meta, "calendar_missing_total")
if total is None:
total, page_size = meta["totalRooms"], meta.get("recordsPerPage")
require(total <= 10000, "calendar_room_limit_exceeded")
require(meta["totalRooms"] == total and meta.get("recordsPerPage") == page_size, "calendar_totals_changed")
rooms = response.data["roomCalendar"]["room"]
for room in rooms:
number = _string(room, "roomId", "calendar_missing_room_identity")
require(number not in room_ids, "calendar_duplicate_room")
room_ids.add(number)
for target in inspected["targets"]:
collected[target["reservation_id"]].extend(target["occurrences"])
require(len(room_ids) <= total, "calendar_inconsistent_total")
if len(room_ids) == total:
return collected
require(bool(rooms) and type(page_size) is int and len(rooms) == page_size
and type(meta.get("pageIndex")) is int, "calendar_incomplete_page")
# First request omits pageIndex. Later requests follow the server's actual
# initial page and size, never assume a zero- or one-based origin.
params = {**params, "recordsPerPage": page_size, "pageIndex": meta["pageIndex"] + 1}
raise base.CollectionError("calendar_page_limit_exceeded")
def _history_room(occurrences, report_date):
if not occurrences:
raise FieldGap("reservation_not_found_in_room_history")
candidates = []
for entry in occurrences:
if entry["schedule_category"].get("value") != "Reservation":
raise FieldGap("room_history_category_unknown")
start, end = entry["segment_start"].get("value"), entry["segment_end"].get("value")
if not start or not end:
start, end = entry["schedule_start"].get("value"), entry["schedule_end"].get("value")
if not start or not end:
raise FieldGap("room_history_dates_missing")
if not start[:10] <= report_date <= end[:10]:
continue
room = entry["room"].get("value")
require(type(room) is str, "room_history_room_missing")
candidates.append(room)
for move in entry["moves"]:
moved_on = move["date"].get("value")
# Move timestamps use a different documented zone. Do not infer a
# conversion or a winning room from conflicting move facts.
if moved_on:
for side in ("from_room", "to_room"):
other = move[side].get("value")
if other and other != room:
raise FieldGap("multiple_rooms_in_history", "ambiguous")
names = set(candidates)
if len(names) != 1:
raise FieldGap("multiple_rooms_on_report_date" if names else "room_history_outside_report_date",
"ambiguous" if names else "missing")
return candidates[0]
def _record(search, detail, reader, sequence):
day = reader.options.arrival_date
related = {}
dates = lambda: fields.stay_dates(detail, day)
counts = lambda: fields.agreed_guest_counts(detail, day)
values = {
"BLOCK_CODE": observe(lambda: _block(search, detail, reader)),
"ADULTS": observe(lambda: counts().adults),
"CHILDREN": observe(lambda: counts().children),
"COMPANY_NAME": observe(lambda: _companies(detail, reader, related)),
"CONFIRMATION_NO": observe(lambda: fields.confirmation_no(detail)),
"DISP_ROOM_NO": observe(lambda: fields.agreed_room_number(search, detail, day)),
"EFFECTIVE_RATE_AMOUNT": observe(lambda: _rate(detail, reader, related)),
"FULL_NAME": observe(lambda: _name(detail, reader)),
"RES_COMMENT": observe(lambda: list(fields.reservation_gen_notes(detail))),
"NO_OF_ROOMS": observe(lambda: fields.arrival_number_of_units(detail, day)),
"PRODUCTS": observe(lambda: _packages(detail, reader, related)),
"RATE_CODE": observe(lambda: fields.arrival_rate_code(detail, day)),
"ROOM_CATEGORY_LABEL": observe(lambda: fields.agreed_room_type(search, detail, day)),
"ARRIVAL": observe(lambda: dates().arrival.isoformat()),
"DEPARTURE": observe(lambda: dates().departure.isoformat()),
}
return {"source_sequence": sequence, "reservation_id": base.reservation_id(detail),
"fields": values, "related": related}
def collect(reader):
"""Collect all source rows; return no Finance/business success claims."""
options = reader.options
rows, fatal = [], None
reader.archive.write("capture.json", {"version": VERSION, "options": vars(options),
"service_url": base.SERVICE, "application_id": base.APPLICATION,
"source_kind": reader.source_kind, "max_requests": reader.max_requests})
complete = False
try:
searches = base.search_day(reader, options)
search_refs = list(reader.used)
for index, search in enumerate(searches, 1):
before = len(reader.used)
identity = base.reservation_id(search)
response = reader.query(base.DETAIL, identity=identity)
container = response.data.get("reservations")
details = container.get("reservation") if type(container) is dict else None
require(type(details) is list and len(details) == 1, "ambiguous_detail_count")
detail = details[0]
require(base.validate_row(detail, options) == identity, "detail_identity_mismatch")
for key in ("reservationStatus", "lastModifyDateTime"):
require(type(search.get(key)) is str and bool(search[key]) and search[key] == detail.get(key),
"search_detail_state_mismatch")
# Trace is intentionally neither requested nor validated by this new
# source. COMMENT counts still guard completeness of requested notes.
indicators = search.get("reservationIndicators", [])
require(type(indicators) is list and all(type(i) is dict for i in indicators), "invalid_reservation_indicators")
comment_counts = [i.get("count") for i in indicators if i.get("indicatorName") == "COMMENT"]
require(len(comment_counts) <= 1, "ambiguous_comment_count")
if comment_counts:
comments = detail.get("comments")
require(type(comments) is list and base.integer(comment_counts[0], "invalid_comment_count") == len(comments),
"comment_count_mismatch")
row = _record(search, detail, reader, index)
row["sources"] = list(dict.fromkeys(search_refs + reader.used[before:]))
rows.append(row)
pending = [r for r in rows if r["fields"]["DISP_ROOM_NO"]["state"] not in {"available", "empty"}]
if pending:
before = len(reader.used)
history = None
def fetch_history():
nonlocal history
history = _calendar(reader, [r["reservation_id"] for r in pending])
return True
fetched = observe(fetch_history)
for row in pending:
if history is None:
row["fields"]["DISP_ROOM_NO"] = fetched
else:
occurrences = history[row["reservation_id"]]
row["related"]["room_history"] = occurrences
row["fields"]["DISP_ROOM_NO"] = observe(lambda: _history_room(occurrences, options.arrival_date))
row["sources"] = list(dict.fromkeys(row["sources"] + reader.used[before:]))
require(base.search_day(reader, options) == searches, "source_changed_during_collection")
complete = True
except base.CollectionError as error:
fatal = str(error)
except Exception:
fatal = "unexpected_collection_failure"
states = Counter(f["state"] for row in rows for f in row["fields"].values())
ready = complete and bool(rows) and all(
f["state"] == "available" or (key in OPTIONAL_FIELDS and f["state"] == "empty")
for row in rows for key, f in row["fields"].items())
failed = fatal is not None or states["failed"] > 0
status = "failed" if failed else "collected" if ready else "collected_with_gaps"
payload = {"version": VERSION, "hotel_id": options.hotel_id, "report_date": options.arrival_date,
"source_kind": reader.source_kind,
"status": status, "collection_complete": complete and not failed, "input_complete": ready and not failed,
"field_names": list(FIELDS), "trace_included": False, "atomic_snapshot": False, "records": rows}
raw = json_bytes(payload)
require(len(raw) <= MAX_DATA_BYTES, "data_output_too_large")
reader.archive.write("arr-data.json", raw)
summary = {"version": VERSION, "status": status, "collection_complete": payload["collection_complete"],
"source_kind": reader.source_kind,
"input_complete": payload["input_complete"], "records": len(rows), "field_states": dict(states),
"http_attempts": reader.request_count, "error": fatal, "finance_ready": False,
"data_sha256": hashlib.sha256(raw).hexdigest()}
reader.archive.write("result.json", {**summary, "files": list(reader.archive.files)})
return {**summary, "manifest_sha256": reader.archive.files[-1]["sha256"]}
class ARRDataSource:
"""Reusable application entrypoint; completed request IDs replay without HTTP.
Pass a transport factory for local tests. For real runs a protected application
credential is loaded only once an uncached fetch is explicitly requested.
"""
def __init__(self, root: Path, hotel_id: str, *, credential_file: Path | None = None,
transport_factory=None, page_size=100, max_pages=100, max_records=10000,
max_requests=25000, sleep=time.sleep):
root = Path(root).absolute()
require(not root.resolve().is_relative_to(PROJECT_ROOT), "data_store_must_be_outside_repository")
base.Options("2000-01-01", hotel_id, page_size, max_pages, max_records).validate()
require(type(max_requests) is int and 1 <= max_requests <= 100000, "invalid_request_limit")
require((credential_file is None) != (transport_factory is None), "exactly_one_transport_source_required")
require(transport_factory is None or callable(transport_factory), "invalid_transport_factory")
self.root, self.hotel_id = root, hotel_id
self.credential_file, self.transport_factory = credential_file, transport_factory
self.page_size, self.max_pages, self.max_records = page_size, max_pages, max_records
self.max_requests, self.sleep = max_requests, sleep
def fetch(self, report_date: str, request_id: str) -> dict:
options = base.Options(report_date, self.hotel_id, self.page_size, self.max_pages, self.max_records)
options.validate()
require(type(request_id) is str and re.fullmatch(r"[0-9a-f]{32}", request_id), "invalid_data_request_id")
identity = {"version": VERSION, "options": vars(options), "max_requests": self.max_requests,
"request_id": request_id, "service_url": base.SERVICE, "application_id": base.APPLICATION,
"source_kind": "ohip_platform" if self.transport_factory is None else "test_transport"}
private_directory(self.root)
folder = self.root / request_id
private_directory(folder)
sync_directory(self.root)
with job_lock(folder):
context = folder / "request.json"
if not context.exists():
require(not any(p.name.startswith("attempt-") for p in folder.iterdir()), "orphan_data_request")
atomic_json(context, identity, replace=False)
require(json_bytes(document(protected_read(context, 65536))) == json_bytes(identity), "data_request_conflict")
pointer = folder / "ready.json"
if pointer.exists():
return self._replay(folder, document(protected_read(pointer, 65536)), identity)
attempts = [p for p in folder.iterdir() if re.fullmatch(r"attempt-[0-9]{4}", p.name)]
attempt = max([int(p.name[-4:]) for p in attempts], default=0) + 1
require(attempt <= 9999, "attempt_limit_exceeded")
key = ""
if self.transport_factory is None:
key = base.load_key(Path(self.credential_file))
transport = base.HTTPTransport(key)
else:
transport = self.transport_factory()
archive = base.Archive(folder / f"attempt-{attempt:04d}")
reader = DataReader(archive, options, transport, key=key, sleep=self.sleep, max_requests=self.max_requests,
source_kind=identity["source_kind"])
summary = collect(reader)
summary.update(data_path=str(archive.path / "arr-data.json"), attempt=attempt)
if summary["status"] != "failed":
atomic_json(pointer, {"attempt": archive.path.name, "manifest_sha256": summary["manifest_sha256"],
"data_sha256": summary["data_sha256"]}, replace=False)
return summary
@staticmethod
def _replay(folder, pointer, identity):
require(set(pointer) == {"attempt", "manifest_sha256", "data_sha256"}
and type(pointer["attempt"]) is str and re.fullmatch(r"attempt-[0-9]{4}", pointer["attempt"]),
"invalid_data_checkpoint")
directory = folder / pointer["attempt"]
info = directory.lstat()
require(stat.S_ISDIR(info.st_mode) and info.st_uid == os.getuid()
and stat.S_IMODE(info.st_mode) == 0o700, "unsafe_data_directory")
raw = protected_read(directory / "result.json", base.MAX_MANIFEST_BYTES)
require(hashlib.sha256(raw).hexdigest() == pointer["manifest_sha256"], "data_manifest_changed")
result = document(raw)
require(result.get("version") == VERSION and result.get("status") in {"collected", "collected_with_gaps"},
"invalid_completed_data")
files = result.get("files")
require(type(files) is list, "invalid_data_inventory")
seen, size = set(), 0
for item in files:
name = item.get("name")
require(type(name) is str and re.fullmatch(r"[a-z0-9_.-]+", name)
and name not in {".", "..", "result.json"} and name not in seen, "invalid_data_inventory")
seen.add(name)
data = protected_read(directory / name, max(base.MAX_RESPONSE_BYTES, MAX_DATA_BYTES))
size += len(data)
require(size <= base.MAX_ARCHIVE_BYTES and len(data) == item.get("bytes")
and hashlib.sha256(data).hexdigest() == item.get("sha256"), "data_archive_changed")
require({p.name for p in directory.iterdir()} == seen | {"result.json"}, "data_inventory_changed")
capture = document(protected_read(directory / "capture.json", 65536))
require(set(capture) == {"version", "options", "service_url", "application_id", "max_requests", "source_kind"}
and json_bytes(capture) == json_bytes({k: identity[k] for k in capture}), "data_capture_context_mismatch")
data = protected_read(directory / "arr-data.json", MAX_DATA_BYTES)
require("arr-data.json" in seen and hashlib.sha256(data).hexdigest() == pointer["data_sha256"] == result["data_sha256"],
"data_output_changed")
payload = document(data)
require(payload.get("hotel_id") == identity["options"]["hotel_id"]
and payload.get("source_kind") == result.get("source_kind") == identity["source_kind"]
and payload.get("report_date") == identity["options"]["arrival_date"]
and payload.get("version") == VERSION and payload.get("status") == result["status"],
"data_output_context_mismatch")
return {k: v for k, v in result.items() if k != "files"} | {
"manifest_sha256": pointer["manifest_sha256"], "data_path": str(directory / "arr-data.json"),
"attempt": int(pointer["attempt"][-4:])}
def main(argv=None):
parser = argparse.ArgumentParser(description="Fetch the 15 ARR input fields as private JSON; no XML or Finance writes.")
parser.add_argument("--report-date", required=True)
parser.add_argument("--request-id", required=True, help="32 lowercase hex characters; reuse to replay a completed request")
parser.add_argument("--hotel-id", required=True, help="Expected platform-selected hotel; cannot change platform routing")
parser.add_argument("--credential-file", required=True, type=Path)
parser.add_argument("--output-root", required=True, type=Path)
parser.add_argument("--page-size", type=int, default=100)
parser.add_argument("--max-pages", type=int, default=100)
parser.add_argument("--max-records", type=int, default=10000)
parser.add_argument("--max-requests", type=int, default=25000)
args = parser.parse_args(argv)
try:
source = ARRDataSource(args.output_root, args.hotel_id, credential_file=args.credential_file,
page_size=args.page_size, max_pages=args.max_pages, max_records=args.max_records,
max_requests=args.max_requests)
summary = source.fetch(args.report_date, args.request_id)
except base.CollectionError as error:
summary = {"status": "failed", "error": str(error), "finance_ready": False}
except Exception:
summary = {"status": "failed", "error": "data_source_setup_failed", "finance_ready": False}
print(json.dumps(summary, ensure_ascii=False, allow_nan=False))
return 0 if summary["status"] == "collected" else 2 if summary["status"] == "collected_with_gaps" else 1
if __name__ == "__main__":
raise SystemExit(main())