124 lines
6.4 KiB
Python
124 lines
6.4 KiB
Python
"""v3 executor with synthetic profiles/XML and isolated local ingestion only."""
|
|
import json
|
|
from unittest.mock import patch
|
|
|
|
from arr_web import arr_download_executor as execution
|
|
from integrations.ohip import collect_arr_source as source, prepare_arr_source
|
|
from integrations.ohip import profile_reader, audit_arr_named_day
|
|
from integrations.ohip.rate_info import RateInfoReader
|
|
from tests import test_arr_web_capture_executor as base
|
|
from tests.test_ohip_named_day import NamedService
|
|
|
|
|
|
class NamedCaptureExecutorTests(base.CaptureExecutorTests):
|
|
# Inherit the real orchestration/recovery tests with a v3 factory; those
|
|
# assertions exercise frozen handoff and local DB behavior, not a mock return.
|
|
def setUp(self):
|
|
super().setUp()
|
|
self.transport = NamedService()
|
|
|
|
def factory(self, archive, hotel):
|
|
self.factory_calls += 1
|
|
return (source.Reader(archive, hotel, self.transport, sleep=lambda _: None),
|
|
RateInfoReader(archive, hotel, self.transport, sleep=lambda _: None),
|
|
profile_reader.ProfileSummaryReader(archive, hotel, self.transport,
|
|
max_profiles=3, sleep=lambda _: None))
|
|
|
|
def make_executor(self, **changes):
|
|
return super().make_executor(**({'capture_version': 'v3', 'max_profiles': 3} | changes))
|
|
|
|
def test_named_archive_and_budget_reach_both_source_checks(self):
|
|
seen = []
|
|
def inspect(archive):
|
|
self.assertIs(type(archive), audit_arr_named_day.VerifiedArchive)
|
|
self.assertEqual(archive.options.max_profiles, 3)
|
|
self.assertEqual(archive.result['valid_name_candidates'], 3)
|
|
seen.append(archive.pin)
|
|
original_adapter, original_validator = self.adapter.adapt, self.validator.validate
|
|
def adapt(archive):
|
|
inspect(archive)
|
|
return original_adapter(archive)
|
|
def validate(archive, payload):
|
|
inspect(archive)
|
|
return original_validator(archive, payload)
|
|
with patch.object(self.adapter, 'adapt', side_effect=adapt), \
|
|
patch.object(self.validator, 'validate', side_effect=validate):
|
|
self.assertEqual(self.run_executor().status, 'succeeded')
|
|
checkpoint = json.loads(self.checkpoint().read_bytes())
|
|
identity = checkpoint['identity']
|
|
self.assertEqual(identity['version'], execution.NAMED_VERSION)
|
|
self.assertEqual(identity['capture']['source_contract'], 'arr-api-date-candidates/v3')
|
|
self.assertEqual(seen, [checkpoint['binding']['manifest_sha256']] * 2)
|
|
self.assertEqual(len(self.transport.profile_calls), 3)
|
|
|
|
def test_profile_budget_is_explicit_validated_and_frozen(self):
|
|
for value in (None, True, False, 0, -1, 1.5, '3', 10001):
|
|
with self.subTest(value=value), self.assertRaisesRegex(source.CollectionError, 'invalid_profile_limit'):
|
|
self.make_executor(max_profiles=value)
|
|
for changes in ({'capture_version': 'v2'}, {'capture_version': 'v4'}):
|
|
with self.assertRaises(source.CollectionError):
|
|
self.make_executor(**changes)
|
|
self.assertFalse(self.executor.root.exists())
|
|
self.run_executor()
|
|
for changes in ({'max_profiles': 2}, {'capture_version': 'v2', 'max_profiles': None}):
|
|
with self.assertRaisesRegex(source.CollectionError, 'executor_request_conflict'):
|
|
self.run_executor(self.make_executor(**changes))
|
|
self.assertEqual(self.factory_calls, 1)
|
|
self.assertEqual(len(self.repository._versions), 1)
|
|
|
|
def test_profile_budget_cannot_exceed_reservation_budget(self):
|
|
with self.assertRaisesRegex(source.CollectionError, 'invalid_profile_limit'):
|
|
self.make_executor(max_profiles=3, max_records=2)
|
|
self.assertFalse(self.executor.root.exists())
|
|
self.assertEqual(self.factory_calls, 0)
|
|
|
|
def test_profile_failure_never_calls_mapping_processor_or_ingestion(self):
|
|
transport = self.transport
|
|
self.transport = lambda method, path, raw: (503, {}, b'{}') if path.endswith('profiles/searches') else transport(method, path, raw)
|
|
failed = self.run_executor()
|
|
self.assertEqual((failed.status, failed.retryable), ('failed', True))
|
|
self.assertEqual((self.adapter.calls, self.validator.calls, self.processor.calls), (0, 0, 0))
|
|
self.assertFalse(self.repository._jobs)
|
|
self.assertFalse(self.checkpoint().exists())
|
|
|
|
def test_preparation_diagnostics_cannot_be_used_as_adapter_or_validator(self):
|
|
for changes in ({'adapter': prepare_arr_source}, {'mapping_validator': prepare_arr_source}):
|
|
with self.assertRaisesRegex(ValueError, 'accepted source adapter'):
|
|
self.make_executor(**changes)
|
|
self.assertFalse(self.executor.root.exists())
|
|
|
|
def test_boolean_integer_alias_in_request_identity_refused_before_reuse(self):
|
|
# Refuse mapping before creating Finance state, but persist capture identity.
|
|
self.validator.refuse = True
|
|
with self.assertRaises(ValueError):
|
|
self.run_executor()
|
|
request = self.executor.root / 'requests' / base.REQUEST_ID / 'request.json'
|
|
document = json.loads(request.read_bytes())
|
|
# Capture recheck_after_details is a real boolean in nested options.
|
|
def change_boolean(node):
|
|
if type(node) is dict:
|
|
for key, value in node.items():
|
|
if type(value) is bool:
|
|
node[key] = int(value)
|
|
return True
|
|
if change_boolean(value):
|
|
return True
|
|
return False
|
|
if not change_boolean(document):
|
|
# Protocol flags need not contain a bool: the page count is numeric.
|
|
def change_numeric(node):
|
|
if type(node) is dict:
|
|
for key, value in node.items():
|
|
if type(value) is int:
|
|
node[key] = float(value)
|
|
return True
|
|
if change_numeric(value):
|
|
return True
|
|
return False
|
|
self.assertTrue(change_numeric(document))
|
|
request.write_text(json.dumps(document))
|
|
with self.assertRaisesRegex(source.CollectionError, 'executor_request_conflict'):
|
|
self.run_executor()
|
|
self.assertEqual(self.factory_calls, 1)
|
|
self.assertFalse(self.repository._jobs)
|