380 lines
20 KiB
Python
380 lines
20 KiB
Python
"""Saved zero-item source reviews continue locally without a human confirmation."""
|
|
import hashlib
|
|
import json
|
|
from pathlib import Path
|
|
import sqlite3
|
|
import tempfile
|
|
import threading
|
|
import unittest
|
|
from unittest.mock import Mock, patch
|
|
|
|
from arr_processing.policy import load_processor_policy
|
|
from arr_web import arr_data_review as review_module
|
|
from arr_web.arr_data_review import DataFieldReviews, EMPTY_REVIEW_ACTOR
|
|
from arr_web.arr_downloads import DownloadOutcome, PersistentARRDownloads
|
|
from integrations.ohip.collect_arr_source import CollectionError
|
|
from tests import test_arr_data_review as review_fixtures
|
|
from tests.test_arr_data_review import (
|
|
ACTOR, CONTEXT, DAY, PROJECT, REANALYSIS_MANIFEST, REANALYSIS_POLICY,
|
|
REQUEST, SOURCE_MANIFEST, gap, optional_omission_source,
|
|
raw, source_document,
|
|
)
|
|
|
|
|
|
def omission_source():
|
|
document = source_document()
|
|
for field, reason in (("BLOCK_CODE", "missing_reservation_block"),
|
|
("PRODUCTS", "missing_reservation_packages")):
|
|
gap(document, field)
|
|
document["records"][0]["fields"][field]["reason"] = reason
|
|
return document
|
|
|
|
|
|
class EmptyReviewRecoveryTests(unittest.TestCase):
|
|
@classmethod
|
|
def setUpClass(cls):
|
|
cls.policy = load_processor_policy(PROJECT)
|
|
|
|
def setUp(self):
|
|
temporary = tempfile.TemporaryDirectory()
|
|
self.addCleanup(temporary.cleanup)
|
|
self.root = Path(temporary.name) / "reviews"
|
|
self.service = DataFieldReviews(root=self.root, policy=self.policy, context=CONTEXT)
|
|
self.document = omission_source()
|
|
|
|
def prepare(self):
|
|
return self.service.prepare(REQUEST, raw(self.document), SOURCE_MANIFEST, DAY)
|
|
|
|
def reanalyze(self, review):
|
|
return self.service.apply_source_reanalysis(REQUEST, raw(optional_omission_source(self.document)),
|
|
REANALYSIS_MANIFEST, REANALYSIS_POLICY, expected_revision=review["revision"])
|
|
|
|
def test_zero_item_reanalysis_freezes_under_system_actor_without_original_or_decision_changes(self):
|
|
review = self.reanalyze(self.prepare())
|
|
original = self.service.original(REQUEST)
|
|
self.assertEqual((review["total_count"], review["pending_count"]), (0, 0))
|
|
self.assertTrue(self.service.recover_ready(REQUEST))
|
|
payload = self.service.payload(REQUEST)
|
|
audit = json.loads(payload[0])["manual_data_review"]["manifest"]
|
|
self.assertEqual(audit["changes"], [])
|
|
self.assertEqual([event["action"] for event in audit["events"]], ["source_reanalysis", "finalize"])
|
|
self.assertEqual(audit["events"][-1]["actor"], EMPTY_REVIEW_ACTOR)
|
|
self.assertTrue(self.service.recover_ready(REQUEST))
|
|
self.assertEqual(self.service.payload(REQUEST), payload)
|
|
self.assertEqual(self.service.original(REQUEST), original)
|
|
|
|
def test_fully_filled_manual_items_still_require_explicit_confirmation(self):
|
|
review = self.prepare()
|
|
saved = self.service.update(REQUEST, "1:BLOCK_CODE", review["revision"], "STAFF-VERIFIED", ACTOR)
|
|
review = self.reanalyze(saved)
|
|
self.assertEqual((review["total_count"], review["pending_count"]), (1, 0))
|
|
self.assertFalse(self.service.recover_ready(REQUEST))
|
|
self.assertIsNone(self.service.payload(REQUEST))
|
|
self.assertEqual(self.service.get(REQUEST), review)
|
|
self.service.finalize(REQUEST, review["revision"], ACTOR)
|
|
self.assertTrue(self.service.recover_ready(REQUEST), "an explicit saved freeze may resume")
|
|
audit = json.loads(self.service.payload(REQUEST)[0])["manual_data_review"]["manifest"]
|
|
self.assertEqual(audit["events"][-1]["actor"], ACTOR)
|
|
|
|
def test_pending_failed_or_ambiguous_observations_do_not_auto_continue(self):
|
|
for state in ("missing", "failed", "ambiguous"):
|
|
with self.subTest(state=state):
|
|
document = omission_source()
|
|
gap(document, "DISP_ROOM_NO", state=state)
|
|
key = {"missing": "b", "failed": "c", "ambiguous": "d"}[state] * 32
|
|
review = self.service.prepare(key, raw(document), SOURCE_MANIFEST, DAY)
|
|
review = self.service.apply_source_reanalysis(key, raw(optional_omission_source(document)),
|
|
REANALYSIS_MANIFEST, REANALYSIS_POLICY, expected_revision=review["revision"])
|
|
self.assertGreater(review["pending_count"], 0)
|
|
self.assertFalse(self.service.recover_ready(key))
|
|
self.assertIsNone(self.service.payload(key))
|
|
|
|
def test_complete_ordinary_prepare_still_needs_no_review_or_derived_source(self):
|
|
self.document = source_document()
|
|
self.assertIsNone(self.prepare())
|
|
self.assertIsNone(self.service.payload(REQUEST))
|
|
|
|
def test_cancelled_and_pm_hidden_history_decisions_freeze_without_overlay_and_remain_audited(self):
|
|
for kind in ("cancelled", "pm"):
|
|
with self.subTest(kind=kind):
|
|
document = source_document(2)
|
|
gap(document, "BLOCK_CODE")
|
|
if kind == "cancelled":
|
|
document["records"][0]["reservation_status"] = "Cancelled"
|
|
else:
|
|
document["records"][0]["fields"]["ROOM_CATEGORY_LABEL"] = {"state": "available", "value": "PM"}
|
|
key = ("b" if kind == "cancelled" else "c") * 32
|
|
old_rule = "is_cancelled_record" if kind == "cancelled" else "is_pm_record"
|
|
# Represent an audited legacy decision saved before the approved
|
|
# exclusion; the original source and rule identity stay pinned.
|
|
with patch.object(self.service.rules, old_rule, return_value=False):
|
|
review = self.service.prepare(key, raw(document), SOURCE_MANIFEST, DAY)
|
|
self.service.update(key, "1:BLOCK_CODE", review["revision"], "PRESERVED-" + kind, ACTOR)
|
|
original = self.service.original(key)
|
|
review = self.service.prepare(key, raw(document), SOURCE_MANIFEST, DAY)
|
|
self.assertEqual((review["total_count"], review["pending_count"]), (0, 0))
|
|
self.assertTrue(self.service.recover_ready(key))
|
|
audit = json.loads(self.service.payload(key)[0])["manual_data_review"]["manifest"]
|
|
self.assertEqual(audit["changes"][0]["value"], "PRESERVED-" + kind)
|
|
self.assertEqual(audit["changes"][0]["actor"], ACTOR)
|
|
self.assertEqual(audit["events"][-1]["actor"], EMPTY_REVIEW_ACTOR)
|
|
self.assertNotIn("source_reanalysis", audit)
|
|
self.assertNotIn("reservation_status_evidence", audit)
|
|
self.assertEqual(self.service.original(key), original)
|
|
|
|
def test_failed_original_source_is_not_automatically_frozen(self):
|
|
self.document["status"] = "failed"
|
|
review = self.reanalyze(self.prepare())
|
|
self.assertEqual(review["total_count"], 0)
|
|
self.assertFalse(self.service.recover_ready(REQUEST))
|
|
self.assertIsNone(self.service.payload(REQUEST))
|
|
|
|
def test_invalid_receipt_blocks_freeze_without_state_change(self):
|
|
review = self.reanalyze(self.prepare())
|
|
pointer = (self.root / REQUEST / "state.json").read_bytes()
|
|
reference = next((self.root / REQUEST).glob("source-reanalysis-*.json"))
|
|
reference.write_bytes(reference.read_bytes() + b" ")
|
|
with self.assertRaises(CollectionError):
|
|
self.service.recover_ready(REQUEST)
|
|
self.assertEqual((self.root / REQUEST / "state.json").read_bytes(), pointer)
|
|
self.assertFalse((self.root / REQUEST / "finalize-intent.json").exists())
|
|
|
|
def test_original_hash_or_context_conflict_never_auto_freezes(self):
|
|
self.reanalyze(self.prepare())
|
|
path = self.root / REQUEST / "original.json"
|
|
original = path.read_bytes()
|
|
path.write_bytes(original + b" ")
|
|
with self.assertRaises(CollectionError):
|
|
self.service.recover_ready(REQUEST)
|
|
path.write_bytes(original)
|
|
changed = DataFieldReviews(root=self.root, policy=self.policy, context={**CONTEXT, "hotel_id": "OTHER"})
|
|
with self.assertRaises(CollectionError):
|
|
changed.recover_ready(REQUEST)
|
|
self.assertFalse((self.root / REQUEST / "finalize-intent.json").exists())
|
|
|
|
def test_interrupted_freeze_replays_same_bytes_and_one_system_event(self):
|
|
review = self.reanalyze(self.prepare())
|
|
publish = review_module.atomic_json
|
|
def interrupt(path, document, **kwargs):
|
|
if path.name == "state.json":
|
|
raise OSError("synthetic freeze pointer interruption")
|
|
return publish(path, document, **kwargs)
|
|
with patch.object(review_module, "atomic_json", side_effect=interrupt), self.assertRaises(OSError):
|
|
self.service.recover_ready(REQUEST)
|
|
frozen = (self.root / REQUEST / "reviewed-source.json").read_bytes()
|
|
self.assertEqual(self.service.get(REQUEST), review)
|
|
restored = DataFieldReviews(root=self.root, policy=self.policy, context=CONTEXT)
|
|
self.assertTrue(restored.recover_ready(REQUEST))
|
|
self.assertEqual(restored.payload(REQUEST)[0], frozen)
|
|
audit = json.loads(frozen)["manual_data_review"]["manifest"]
|
|
self.assertEqual(sum(event["action"] == "finalize" for event in audit["events"]), 1)
|
|
self.assertEqual(audit["events"][-1]["actor"], EMPTY_REVIEW_ACTOR)
|
|
|
|
def test_saved_status_zero_items_preserves_prior_decisions_in_audit(self):
|
|
self.document = source_document(2)
|
|
gap(self.document, "BLOCK_CODE")
|
|
review = self.prepare()
|
|
review = self.service.update(REQUEST, "1:BLOCK_CODE", review["revision"], "PRESERVED", ACTOR)
|
|
evidence = {"version": "arr-reservation-status-evidence/v1", "policy_id": review_module.STATUS_POLICY,
|
|
"original_sha256": hashlib.sha256(raw(self.document)).hexdigest(), "source_manifest_sha256": SOURCE_MANIFEST,
|
|
"hotel_id": CONTEXT["hotel_id"], "report_date": DAY,
|
|
"records": [{"source_sequence": row["source_sequence"], "reservation_id": row["reservation_id"],
|
|
"reservation_status": status, "sources": [{"file": "synthetic.response.bin", "sha256": "f" * 64}]}
|
|
for row, status in zip(self.document["records"], ["Cancelled", "InHouse"])]}
|
|
with patch("integrations.ohip.reservation_status.saved_status_evidence", return_value=evidence):
|
|
review = self.service.apply_saved_status_evidence(REQUEST, Path("synthetic-unused"),
|
|
expected_revision=review["revision"])
|
|
self.assertEqual((review["total_count"], review["pending_count"]), (0, 0))
|
|
original = self.service.original(REQUEST)
|
|
self.assertTrue(self.service.recover_ready(REQUEST))
|
|
audit = json.loads(self.service.payload(REQUEST)[0])["manual_data_review"]["manifest"]
|
|
self.assertEqual(audit["changes"][0]["value"], "PRESERVED")
|
|
self.assertEqual(audit["changes"][0]["actor"], ACTOR)
|
|
self.assertEqual(audit["events"][-1]["actor"], EMPTY_REVIEW_ACTOR)
|
|
self.assertEqual(self.service.original(REQUEST), original)
|
|
|
|
|
|
class EmptyReviewExecutorTests(unittest.TestCase):
|
|
@classmethod
|
|
def setUpClass(cls):
|
|
cls.policy = load_processor_policy(PROJECT)
|
|
|
|
setUp = review_fixtures.DirectFieldReviewTests.setUp
|
|
configure_executor = review_fixtures.DirectFieldReviewTests.configure_executor
|
|
execute = review_fixtures.DirectFieldReviewTests.execute
|
|
wait_task = review_fixtures.DirectFieldReviewTests.wait_task
|
|
|
|
def fixture(self):
|
|
document = omission_source()
|
|
payload = raw(document)
|
|
path = self.files / "synthetic-collected.json"
|
|
path.write_bytes(payload)
|
|
path.chmod(0o600)
|
|
fetched = {"collection_complete": True, "data_path": str(path),
|
|
"data_sha256": hashlib.sha256(payload).hexdigest(), "manifest_sha256": SOURCE_MANIFEST}
|
|
return document, fetched
|
|
|
|
def reanalyze(self, document):
|
|
review = self.executor.get_data_review(REQUEST)
|
|
return self.executor.data_reviews.apply_source_reanalysis(REQUEST, raw(optional_omission_source(document)),
|
|
REANALYSIS_MANIFEST, REANALYSIS_POLICY, expected_revision=review["revision"])
|
|
|
|
def test_executor_zero_item_saved_review_continues_same_source_with_one_finance_version(self):
|
|
document, fetched = self.fixture()
|
|
with patch.object(self.source, "fetch", return_value=fetched) as fetch:
|
|
self.assertEqual(self.execute().status, "needs_data_review")
|
|
self.reanalyze(document)
|
|
original = self.executor.data_reviews.original(REQUEST)
|
|
self.assertEqual(self.execute().status, "succeeded")
|
|
self.assertEqual(self.execute().status, "succeeded")
|
|
self.assertEqual(fetch.call_count, 1)
|
|
self.assertEqual(len(self.repository._versions), 1)
|
|
self.assertEqual(self.executor.data_reviews.original(REQUEST), original)
|
|
|
|
def test_missing_or_changed_saved_checkpoint_never_auto_freezes_or_refetches(self):
|
|
document, fetched = self.fixture()
|
|
with patch.object(self.source, "fetch", return_value=fetched):
|
|
self.assertEqual(self.execute().status, "needs_data_review")
|
|
self.reanalyze(document)
|
|
path = self.files / "executor" / "requests" / REQUEST / "collected.json"
|
|
saved = path.read_bytes()
|
|
with patch.object(self.source, "fetch", side_effect=AssertionError("must not refetch")) as fetch:
|
|
path.unlink()
|
|
with self.assertRaises(FileNotFoundError):
|
|
self.executor.recover_data_review(REQUEST)
|
|
path.write_bytes(saved)
|
|
path.chmod(0o600)
|
|
changed = json.loads(saved)
|
|
changed["data_sha256"] = "0" * 64
|
|
path.write_bytes(raw(changed))
|
|
with self.assertRaises(CollectionError):
|
|
self.executor.recover_data_review(REQUEST)
|
|
fetch.assert_not_called()
|
|
self.assertEqual(self.executor.get_data_review(REQUEST)["status"], "editing")
|
|
self.assertIsNone(self.executor.data_reviews.payload(REQUEST))
|
|
self.assertEqual(len(self.repository._versions), 0)
|
|
|
|
def test_startup_recovers_without_browser_get_or_another_fetch_and_restart_is_idempotent(self):
|
|
document, fetched = self.fixture()
|
|
with patch.object(self.source, "fetch", return_value=fetched):
|
|
queue = PersistentARRDownloads(self.files / "queue", self.executor)
|
|
try:
|
|
queue.create(DAY, REQUEST)
|
|
self.assertEqual(self.wait_task(queue)["status"], "needs_data_review")
|
|
finally:
|
|
queue.close(wait=True)
|
|
self.reanalyze(document)
|
|
completed = threading.Event()
|
|
execute = self.executor.execute
|
|
def run(**kwargs):
|
|
try:
|
|
return execute(**kwargs)
|
|
finally:
|
|
completed.set()
|
|
with patch.object(self.source, "fetch", side_effect=AssertionError("must replay saved checkpoint")), \
|
|
patch.object(self.executor, "execute", side_effect=run):
|
|
queue = PersistentARRDownloads(self.files / "queue", self.executor)
|
|
try:
|
|
self.assertTrue(completed.wait(30), "startup must recover without any browser read")
|
|
self.assertEqual(self.wait_task(queue)["status"], "succeeded")
|
|
finally:
|
|
queue.close(wait=True)
|
|
queue = PersistentARRDownloads(self.files / "queue", self.executor)
|
|
try:
|
|
self.assertEqual(queue.get(REQUEST)["status"], "succeeded")
|
|
finally:
|
|
queue.close(wait=True)
|
|
self.assertEqual(len(self.repository._versions), 1)
|
|
|
|
def test_successful_freeze_before_lost_queue_update_recovers_on_next_startup(self):
|
|
document, fetched = self.fixture()
|
|
with patch.object(self.source, "fetch", return_value=fetched):
|
|
queue = PersistentARRDownloads(self.files / "queue", self.executor)
|
|
try:
|
|
queue.create(DAY, REQUEST)
|
|
self.assertEqual(self.wait_task(queue)["status"], "needs_data_review")
|
|
finally:
|
|
queue.close(wait=True)
|
|
self.reanalyze(document)
|
|
recover = self.executor.recover_data_review
|
|
recovered = threading.Event()
|
|
def lost_update(request_id):
|
|
recover(request_id)
|
|
recovered.set()
|
|
raise OSError("synthetic loss before queue update")
|
|
with patch.object(self.executor, "recover_data_review", side_effect=lost_update):
|
|
queue = PersistentARRDownloads(self.files / "queue", self.executor)
|
|
try:
|
|
self.assertTrue(recovered.wait(10))
|
|
with sqlite3.connect(self.files / "queue" / "arr-downloads.sqlite3") as db:
|
|
self.assertEqual(db.execute("SELECT status FROM downloads WHERE request_id=?", (REQUEST,)).fetchone()[0],
|
|
"needs_data_review")
|
|
self.assertTrue(queue.ready)
|
|
finally:
|
|
queue.close(wait=True)
|
|
frozen = self.executor.data_reviews.payload(REQUEST)
|
|
with patch.object(self.source, "fetch", side_effect=AssertionError("must replay saved checkpoint")):
|
|
queue = PersistentARRDownloads(self.files / "queue", self.executor)
|
|
try:
|
|
self.assertEqual(self.wait_task(queue)["status"], "succeeded")
|
|
finally:
|
|
queue.close(wait=True)
|
|
self.assertEqual(self.executor.data_reviews.payload(REQUEST), frozen)
|
|
self.assertEqual(len(self.repository._versions), 1)
|
|
|
|
def test_authorized_wrapper_recovery_is_local_without_access_refresh(self):
|
|
from arr_web.local_ohip import AuthorizedExecutor
|
|
executor, access = Mock(), Mock()
|
|
executor.recover_data_review.return_value = True
|
|
self.assertTrue(AuthorizedExecutor(executor, access).recover_data_review(REQUEST))
|
|
executor.recover_data_review.assert_called_once_with(REQUEST)
|
|
access.require_ready.assert_not_called()
|
|
|
|
|
|
class EmptyReviewQueueRecoveryTests(unittest.TestCase):
|
|
def test_one_bad_review_cannot_kill_worker_and_get_never_runs_new_recovery(self):
|
|
temporary = tempfile.TemporaryDirectory()
|
|
self.addCleanup(temporary.cleanup)
|
|
root = Path(temporary.name)
|
|
entered = threading.Event()
|
|
executor = Mock()
|
|
executor.source_identity.return_value = {"kind": "synthetic-empty-review-recovery"}
|
|
executor.get_acquisition_progress.return_value = None
|
|
executor.get_data_review.return_value = {"status": "editing"}
|
|
executor.recover_data_review.return_value = False
|
|
queue = PersistentARRDownloads(root, executor)
|
|
queue.close(wait=True)
|
|
now = "2026-10-08T00:00:00+00:00"
|
|
with sqlite3.connect(root / "arr-downloads.sqlite3") as db:
|
|
db.executemany("INSERT INTO downloads(request_id,report_date,status,created_at,updated_at) VALUES(?,?,?,?,?)",
|
|
[("a" * 32, DAY, "needs_data_review", now, now),
|
|
("b" * 32, "2026-09-16", "needs_data_review", now, now),
|
|
("c" * 32, "2026-09-17", "needs_data_review", now, now)])
|
|
def recover(request_id):
|
|
if request_id == "a" * 32:
|
|
raise CollectionError("synthetic conflicting evidence")
|
|
return request_id == "c" * 32
|
|
def execute(**kwargs):
|
|
entered.set()
|
|
return DownloadOutcome("succeeded", "arrjob-synthetic-recovered")
|
|
executor.recover_data_review.side_effect = recover
|
|
executor.execute.side_effect = execute
|
|
queue = PersistentARRDownloads(root, executor)
|
|
try:
|
|
self.assertTrue(entered.wait(5), "startup resumes the good task without a browser read")
|
|
queue.close(wait=True)
|
|
self.assertEqual(queue.get("c" * 32)["status"], "succeeded")
|
|
before = executor.recover_data_review.call_count
|
|
self.assertEqual(queue.get("a" * 32)["status"], "needs_data_review")
|
|
self.assertEqual(queue.get("b" * 32)["status"], "needs_data_review")
|
|
queue.latest()
|
|
queue.pending_reviews()
|
|
self.assertEqual(executor.recover_data_review.call_count, before)
|
|
self.assertEqual(executor.execute.call_count, 1)
|
|
finally:
|
|
queue.close(wait=True)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|