Files

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())