189 lines
9.1 KiB
Python
189 lines
9.1 KiB
Python
"""Extract private, unaccepted ARR field evidence from a pinned v2/v3 capture.
|
|
|
|
This is an offline first stage of source adaptation, not an XML adapter. It
|
|
does not select report rows, resolve field alternatives, or implement the Web
|
|
adapter protocol. stdout contains aggregate verification results only.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
from collections import Counter
|
|
from decimal import Decimal
|
|
import json
|
|
from pathlib import Path
|
|
import re
|
|
|
|
from . import audit_arr_day as audit
|
|
from . import audit_arr_named_day as named_audit
|
|
from . import collect_arr_source as capture
|
|
from .source_facts_contract import (BLOCKERS, CONTEXT_PATHS, CONTRACT_SHA256, FIELDS, FLAGS,
|
|
MAX_FACTS_BYTES, VERSION, atom, canonical, kind, definitions)
|
|
|
|
|
|
def _document(archive, name):
|
|
raw = archive.read(name)
|
|
capture.strict_json(raw)
|
|
return json.loads(raw, parse_float=Decimal)
|
|
|
|
|
|
def _reference(archive, filename, pointer, request_name=None):
|
|
ref = {"file": filename, "sha256": archive.inventory[filename]["sha256"], "pointer": pointer}
|
|
if request_name is not None:
|
|
request = archive.document(request_name)
|
|
ref["request"] = {"file": request_name, "sha256": archive.inventory[request_name]["sha256"],
|
|
**{key: request[key] for key in ("operation_id", "method", "path", "body", "attempt")}}
|
|
return ref
|
|
|
|
|
|
def _bindings(archive):
|
|
searches, details = [], []
|
|
search_finished = False
|
|
for name in archive.inventory:
|
|
if not re.fullmatch(r"request-[0-9]{6}\.json", name):
|
|
continue
|
|
label = name[:-5]
|
|
if archive.document(label + ".meta.json").get("http_status") != 200:
|
|
continue
|
|
request = archive.document(name)
|
|
response = label + ".response.bin"
|
|
data = _document(archive, response)["data"]["reservations"]
|
|
if request["operation_id"] == capture.SEARCH and not search_finished:
|
|
for i, row in enumerate(data["reservationInfo"]):
|
|
searches.append((_reference(archive, response, f"/data/reservations/reservationInfo/{i}", name), row))
|
|
search_finished = not data.get("hasMore", False)
|
|
elif request["operation_id"] == capture.DETAIL:
|
|
for i, row in enumerate(data["reservation"]):
|
|
details.append((_reference(archive, response, f"/data/reservations/reservation/{i}", name), row))
|
|
return searches, details
|
|
|
|
|
|
def _walk(value, steps, pointer):
|
|
if not steps:
|
|
return [{"pointer": pointer, **atom(value)}]
|
|
step, tail = steps[0], steps[1:]
|
|
if value is None:
|
|
return [{"pointer": pointer, "state": "null", "kind": "null"}]
|
|
if step == "*":
|
|
if not isinstance(value, list):
|
|
return [{"pointer": pointer, "state": "invalid_container", "kind": kind(value)}]
|
|
if not value:
|
|
return [{"pointer": pointer, "state": "empty_collection", "kind": "array"}]
|
|
result = []
|
|
for i, item in enumerate(value):
|
|
result.extend(_walk(item, tail, pointer + "/" + str(i)))
|
|
return result
|
|
if not isinstance(value, dict):
|
|
return [{"pointer": pointer, "state": "invalid_container", "kind": kind(value)}]
|
|
child_pointer = pointer + "/" + step.replace("~", "~0").replace("/", "~1")
|
|
if step not in value:
|
|
return [{"pointer": child_pointer, "state": "missing"}]
|
|
return _walk(value[step], tail, child_pointer)
|
|
|
|
|
|
def _fields(specs, sources):
|
|
result = {}
|
|
for field, variants in specs.items():
|
|
result[field] = {"mapping_state": "unresolved", "variants": []}
|
|
for source, variant, path in variants:
|
|
reference, value = sources[source]
|
|
result[field]["variants"].append({"source": source, "variant": variant, "path": path,
|
|
"observations": (_walk(value, path.split("/"), reference["pointer"]) if reference is not None
|
|
else [{"state": "source_not_acquired"}])})
|
|
return result
|
|
|
|
|
|
def build_source_facts(archive: audit.VerifiedArchive) -> bytes:
|
|
# Reopen against the externally supplied pin rather than trusting mutated
|
|
# Python attributes on a previously constructed archive instance.
|
|
named = type(archive) is named_audit.VerifiedArchive
|
|
protocol = named_audit if named else audit
|
|
version, contract_sha256, fields, context_paths, blockers = definitions(named)
|
|
archive = protocol.VerifiedArchive(archive.directory, archive.pin)
|
|
_, _, assessments = protocol.replay(archive)
|
|
searches, details = _bindings(archive)
|
|
capture.require(len(searches) == len(details) == len(assessments["records"]), "facts_source_count_mismatch")
|
|
records = []
|
|
for i, (search, detail, assessment) in enumerate(zip(searches, details, assessments["records"])):
|
|
identity = capture.reservation_id(detail[1])
|
|
capture.require(capture.reservation_id(search[1]) == identity == assessment["reservation_id"],
|
|
"facts_source_identity_mismatch")
|
|
rate_request = assessment["rate_request"]
|
|
rate_file = rate_request[:-5] + ".response.bin"
|
|
sources = {"search": search, "detail": detail,
|
|
"rate": (_reference(archive, rate_file, "/data", rate_request), _document(archive, rate_file)["data"]),
|
|
"assessment": (_reference(archive, "rate-assessments.json", f"/records/{i}"), assessment)}
|
|
if named:
|
|
profile_request = assessment["profile"].get("profile_request")
|
|
if profile_request is None:
|
|
sources["profile"] = (None, None)
|
|
else:
|
|
profile_file = profile_request[:-5] + ".response.bin"
|
|
sources["profile"] = (_reference(archive, profile_file,
|
|
"/data/profileSummaries/profileInfo/0", profile_request),
|
|
_document(archive, profile_file)["data"]["profileSummaries"]["profileInfo"][0])
|
|
records.append({"capture_sequence": i + 1, "reservation_id": identity,
|
|
"sources": {name: pair[0] for name, pair in sources.items()},
|
|
"fields": _fields(fields, sources), "context": _fields(context_paths, sources)})
|
|
result = {"version": version, "contract_sha256": contract_sha256, "status": "unaccepted_field_evidence",
|
|
"capture_manifest_sha256": archive.pin, "hotel_id": archive.options.hotel_id,
|
|
"from_date": archive.options.from_date, "to_date": archive.options.to_date, "rate_date": archive.options.rate_date,
|
|
"record_order": "capture_order_not_report_order", **FLAGS,
|
|
"business_blockers": list(blockers), "records": records}
|
|
raw = canonical(result)
|
|
capture.require(len(raw) <= MAX_FACTS_BYTES, "facts_byte_budget_exceeded")
|
|
return raw
|
|
|
|
|
|
def summarize_facts(document):
|
|
# Fixed contract keys and state labels only; never print values/IDs/paths.
|
|
counts = {name: Counter() for name in FIELDS}
|
|
for record in document["records"]:
|
|
for name in FIELDS:
|
|
field = record["fields"][name]
|
|
if not field["variants"]:
|
|
counts[name]["no_proven_source_path"] += 1
|
|
for variant in field["variants"]:
|
|
for observation in variant["observations"]:
|
|
counts[name][observation["state"]] += 1
|
|
return {"records": len(document["records"]), "field_observation_states": counts,
|
|
"record_order": "capture_order_not_report_order", **FLAGS}
|
|
|
|
|
|
def export_source_facts(capture_dir: Path, capture_pin: str, output_dir: Path, *, capture_version="v2") -> dict:
|
|
from .validate_source_facts import verify_source_facts
|
|
|
|
repository = Path(__file__).resolve().parents[2]
|
|
capture.require(not output_dir.resolve().is_relative_to(repository), "output_must_be_outside_repository")
|
|
capture.require(capture_version in ("v2", "v3"), "unsupported_facts_capture_version")
|
|
protocol = named_audit if capture_version == "v3" else audit
|
|
archive = protocol.VerifiedArchive(capture_dir, capture_pin)
|
|
raw = build_source_facts(archive)
|
|
verification = verify_source_facts(archive, raw)
|
|
# Only create an output directory after independent verification succeeds.
|
|
sink = capture.Archive(output_dir)
|
|
sink.write("source-facts.json", raw)
|
|
sink.write("verification.json", verification)
|
|
return {**verification, **summarize_facts(json.loads(raw))}
|
|
|
|
|
|
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("--capture-version", choices=("v2", "v3"), default="v2")
|
|
args = parser.parse_args(argv)
|
|
try:
|
|
result = export_source_facts(args.capture_dir, args.capture_sha256, args.output_dir,
|
|
capture_version=args.capture_version)
|
|
except capture.CollectionError as error:
|
|
result = {"status": "failed", "error": str(error), **FLAGS}
|
|
except Exception:
|
|
result = {"status": "failed", "error": "facts_setup_or_validation_failed", **FLAGS}
|
|
print(json.dumps(result, ensure_ascii=False, sort_keys=True))
|
|
return 0 if result.get("facts_verified") is True else 1
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|