Fix production ARR reads and add audited source field review

This commit is contained in:
Wyndham ARR committed 2026-10-08 17:01:54 +08:00
1 parent 2417b1a41c
commit fe8735aac2
26 files changed
+2417 -157

No files matched your search

+39 -7
View File
@@ -18,6 +18,18 @@ class DirectARRExecutor:
"data_executor_store_must_be_outside_repository")
self.source, self.policy, self.processor = source, policy, processor
self.object_store, self.repository, self.ingestion = object_store, repository, ingestion
from arr_web.arr_data_review import DataFieldReviews
self.data_reviews = DataFieldReviews(root=self.root / "data-reviews", policy=policy,
context=self.source_identity())
def get_data_review(self, request_id):
return self.data_reviews.get(request_id)
def update_data_review_item(self, request_id, item_id, revision, value, actor):
return self.data_reviews.update(request_id, item_id, revision, value, actor)
def finalize_data_review(self, request_id, revision, actor):
return self.data_reviews.finalize(request_id, revision, actor)
def source_identity(self):
return {"version": "arr-direct-data-executor/v1", "source_version": SOURCE_VERSION,
@@ -53,14 +65,34 @@ class DirectARRExecutor:
== ("data-" + request_id, self.source.hotel_id, from_date.isoformat(), SOURCE_VERSION),
"data_checkpoint_binding_mismatch")
else:
report_stage("downloading")
collected = self.source.fetch(from_date.isoformat(), request_id)
if collected.get("collection_complete") is not True:
return DownloadOutcome("failed", retryable=True)
payload = protected_read(Path(collected["data_path"]), 100 * 1024 * 1024)
require(handoff._hash(payload) == collected["data_sha256"], "data_source_changed")
collected_file = directory / "collected.json"
if collected_file.exists():
checkpoint = strict_json(protected_read(collected_file, 65536))
require(fingerprint(checkpoint["identity"]) == fingerprint(identity), "data_collection_checkpoint_conflict")
payload, original_manifest = self.data_reviews.original(request_id)
require(handoff._hash(payload) == checkpoint["data_sha256"]
and original_manifest == checkpoint["manifest_sha256"], "data_collection_checkpoint_changed")
else:
report_stage("downloading")
collected = self.source.fetch(from_date.isoformat(), request_id)
if collected.get("collection_complete") is not True:
return DownloadOutcome("failed", retryable=True)
payload = protected_read(Path(collected["data_path"]), 100 * 1024 * 1024)
require(handoff._hash(payload) == collected["data_sha256"], "data_source_changed")
original_manifest = collected["manifest_sha256"]
review = self.data_reviews.prepare(request_id, payload, original_manifest,
from_date.isoformat())
if not collected_file.exists():
atomic_json(collected_file, {"identity": identity, "data_sha256": handoff._hash(payload),
"manifest_sha256": original_manifest}, replace=False)
manifest_sha256 = original_manifest
if review:
completed = self.data_reviews.payload(request_id)
if completed is None:
return DownloadOutcome("needs_data_review")
payload, manifest_sha256 = completed
binding = handoff.CaptureBinding("data-" + request_id, self.source.hotel_id,
from_date.isoformat(), collected["manifest_sha256"], SOURCE_VERSION)
from_date.isoformat(), manifest_sha256, SOURCE_VERSION)
report_stage("processing")
result = handoff.prepare(self.root / "handoffs", binding, "ARR.json", payload, self.policy,
processor=self.processor, source_role="source_data")