553 lines
31 KiB
Python
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())
|