import copy import hashlib import os import sys import tempfile import unittest from dataclasses import FrozenInstanceError from datetime import datetime, timedelta, timezone from types import SimpleNamespace APP_DIR = os.path.abspath(os.path.join(os.path.dirname(__file__), '..', 'app')) if APP_DIR not in sys.path: sys.path.insert(0, APP_DIR) from console_runner import _GitResolutionFailure from result_bundle import FORMAT_VERSION, ResultBundleReader, bundle_ready_path from runtime_security import ensure_private_directory from scan_execution import PROTOCOL_VERSION, ScanExecutionError, remote_execution_snapshot_sha256 from worker_assignment import ( ASSIGNMENT_SOURCE_ADAPTERS, CORE_ASSIGNMENT_SOURCE_ADAPTERS, LEGACY_GITHUB_ASSIGNMENT_ADAPTER, PROTOCOL1_NEW_CLAIM_SOURCES, SUPPORTED_REMOTE_GIT_SOURCES, RemoteAssignmentBuilder, RemoteGitAssignmentBuilder, assignment_source_adapter, ) from worker_package import worker_package_build_compatibility from worker_contracts import ( DIAGNOSTIC_PROJECTION_VERSION, ordered_diagnostic_uid_set_sha256, ) from lifecycle_authority import ( GIT_MANIFEST_NAME, REMOTE_WORKER_CODE_AUTHORITY_FILES, TRUFFLEHOG_MANIFEST_NAME, ) from scanner_db import ( ScanEventConflictError, canonical_git_scan_plan_bytes, validate_result_git_scan_plan, ) TOKEN_SHA256 = 'd' * 64 class AssignmentSourceAdapterTests(unittest.TestCase): def test_core_registry_declares_exact_distributed_capabilities(self): self.assertEqual( tuple(CORE_ASSIGNMENT_SOURCE_ADAPTERS), ('gitlab', 'dockerhub', 'huggingface'), ) self.assertEqual( { source: adapter.package_capability.as_dict() for source, adapter in CORE_ASSIGNMENT_SOURCE_ADAPTERS.items() }, { 'gitlab': { 'source': 'gitlab', 'platform': 'gitlab', 'planning_kind': 'exact_git_v1', }, 'dockerhub': { 'source': 'dockerhub', 'platform': 'docker', 'planning_kind': 'docker_direct_v1', }, 'huggingface': { 'source': 'huggingface', 'platform': 'huggingface', 'planning_kind': 'huggingface_space_v1', }, }, ) self.assertNotIn('github', CORE_ASSIGNMENT_SOURCE_ADAPTERS) self.assertIs(ASSIGNMENT_SOURCE_ADAPTERS['github'], LEGACY_GITHUB_ASSIGNMENT_ADAPTER) def test_registry_is_immutable_and_unknown_sources_fail_closed(self): with self.assertRaises(TypeError): CORE_ASSIGNMENT_SOURCE_ADAPTERS['github'] = LEGACY_GITHUB_ASSIGNMENT_ADAPTER with self.assertRaises(FrozenInstanceError): LEGACY_GITHUB_ASSIGNMENT_ADAPTER.worker_platform = 'gitlab' with self.assertRaisesRegex(ValueError, 'unsupported'): assignment_source_adapter('docker') with self.assertRaisesRegex(ValueError, 'unsupported'): assignment_source_adapter('unknown') def test_direct_adapters_are_available_only_with_credential_free_args(self): for source, platform in (('dockerhub', 'docker'), ('huggingface', 'huggingface')): adapter = assignment_source_adapter(source) self.assertEqual(adapter.worker_platform, platform) self.assertIsNotNone(adapter.assignment_flow) adapter.validate_source_args(SimpleNamespace(platform=platform)) with self.assertRaisesRegex(ValueError, 'credential-free'): adapter.validate_source_args(SimpleNamespace( platform=platform, token='private-token', )) self.assertTrue(callable(adapter.snapshot_validator)) def test_protocol1_compatibility_aliases_remain_stable(self): self.assertIs(RemoteGitAssignmentBuilder, RemoteAssignmentBuilder) self.assertIs(SUPPORTED_REMOTE_GIT_SOURCES, PROTOCOL1_NEW_CLAIM_SOURCES) self.assertEqual(PROTOCOL1_NEW_CLAIM_SOURCES, {'github', 'gitlab'}) class _FakeDB: enabled = True def __init__(self, owner): self.owner = owner def set_application_name(self, value): self.owner.applications.append(value) def remote_assignment_status(self, reservation_id, device_id, token_sha256): if token_sha256 != TOKEN_SHA256: raise AssertionError('remote status lost its credential fence') return { 'reservation_id': reservation_id, 'state': 'scanning', 'expires_at': '2099-01-01T00:00:00.000000', } def reconcile_remote_assignment_request(self, request_id, device_id, token_sha256): if token_sha256 != TOKEN_SHA256 or device_id != 11: raise AssertionError('remote request recovery lost its credential fence') self.owner.reconcile_calls.append(request_id) return self.owner.reconciled_by_request.get(request_id, self.owner.reconciled) def remote_bound_git_scan_plan( self, reservation_id, device_id, claim_lease_token, token_sha256, ): if token_sha256 != TOKEN_SHA256: raise AssertionError('remote plan recovery lost its credential fence') self.owner.plan_recovery_calls.append( (reservation_id, device_id, claim_lease_token), ) return self.owner.bound_plan def mark_result_bundle_ready( self, reservation_id, metadata, remote_acceptance=None, bundle_capacity_bytes=None, ): if remote_acceptance.get('token_sha256') != TOKEN_SHA256: raise AssertionError('remote acceptance lost its credential fence') if ( metadata.get('effective_diagnostic_projection_version') != DIAGNOSTIC_PROJECTION_VERSION or isinstance(metadata.get('effective_diagnostic_count'), bool) or not isinstance(metadata.get('effective_diagnostic_count'), int) or metadata['effective_diagnostic_count'] < 0 or not isinstance(metadata.get('effective_diagnostic_uids_sha256'), str) or len(metadata['effective_diagnostic_uids_sha256']) != 64 ): raise AssertionError('remote acceptance lost diagnostic projection authority') if self.owner.bundle_root is None: raise AssertionError('fake DB lacks bundle authority for remote acceptance') diagnostics = ResultBundleReader(os.path.join( self.owner.bundle_root, *metadata['relative_path'].split('/'), )).effective_diagnostics() if ( metadata['effective_diagnostic_count'] != len(diagnostics) or metadata['effective_diagnostic_uids_sha256'] != ordered_diagnostic_uid_set_sha256(diagnostics) ): raise AssertionError('remote diagnostic projection authority is incorrect') if self.owner.accept_failures: self.owner.accept_failures -= 1 raise RuntimeError('synthetic acceptance failure') self.owner.accepted.append((reservation_id, metadata, remote_acceptance)) return {'receipt_id': 'f' * 64, 'reservation_id': reservation_id} def close(self): pass class _DBFactory: def __init__(self): self.applications = [] self.accepted = [] self.bound_plan = None self.plan_recovery_calls = [] self.reconciled = None self.reconciled_by_request = {} self.reconcile_calls = [] self.accept_failures = 0 self.bundle_root = None def __call__(self, **_kwargs): return _FakeDB(self) def _source_args(root, platform='github'): policy_path = os.path.join(root, 'server-policy.yaml') if not os.path.exists(policy_path): with open(policy_path, 'wb') as handle: handle.write(b'detectors: []\n') return SimpleNamespace( platform=platform, exact_git_planning_enabled=True, workers=1, timeout=30, save_dir=root, detectors='', exclude_detectors='', drop_detectors='', no_verification=False, trufflehog_config=policy_path, token='task-token', scan_full_history=False, max_depth=25, max_commit_age_days=0, commit_lookup_pages=3, skip_if_commit_lookup_fails=True, result_bundle_max_event_bytes=1 << 20, result_bundle_max_items=20, result_bundle_max_total_bytes=32 << 20, projection_backlog_max_items=20, projection_backlog_max_bytes=32 << 20, projection_backlog_headroom_bytes=2 << 20, keycheck_queue_max_items=100, keycheck_queue_max_bytes=8 << 20, pipeline_quarantine_max_items=20, pipeline_quarantine_max_bytes=8 << 20, keycheck_candidates_per_event=50, keycheck_candidate_bytes_per_event=1 << 20, target_retry_max_attempts=3, target_retry_base_delay_sec=60, target_retry_max_delay_sec=600, target_timeout_retry_delay_sec=300, max_active_scans=1, admission_resolution_attempts=2, admission_resolution_seconds=1, admission_resolution_retry_delay_sec=0.01, target_claim_order='oldest', ) def _direct_source_args(root, source): args = _source_args(root, { 'dockerhub': 'docker', 'huggingface': 'huggingface', }[source]) args.exact_git_planning_enabled = False args.token = '' args.docker_username = '' args.docker_token = '' args.auth_name = None return args def _package_manifest(args, sources=None): with open(args.trufflehog_config, 'rb') as handle: policy_digest = hashlib.sha256(handle.read()).hexdigest() files = {} for name in REMOTE_WORKER_CODE_AUTHORITY_FILES: files[name] = {'path': f'app/{name}', 'sha256': '1' * 64} source_values = list(sources or [args.platform]) capability_values = { 'github': { 'source': 'github', 'platform': 'github', 'planning_kind': 'exact_git_v1', }, 'gitlab': { 'source': 'gitlab', 'platform': 'gitlab', 'planning_kind': 'exact_git_v1', }, 'dockerhub': { 'source': 'dockerhub', 'platform': 'docker', 'planning_kind': 'docker_direct_v1', }, 'huggingface': { 'source': 'huggingface', 'platform': 'huggingface', 'planning_kind': 'huggingface_space_v1', }, } return { 'schema': 3, 'protocol_version': PROTOCOL_VERSION, 'bundle_format_version': FORMAT_VERSION, 'platform_tag': 'windows-x86_64', 'capabilities': [capability_values[source] for source in source_values], 'app_root': 'app', 'files': files, 'executables': { TRUFFLEHOG_MANIFEST_NAME: {'path': 'bin/trufflehog.exe', 'sha256': '2' * 64}, GIT_MANIFEST_NAME: {'path': 'runtime/git/cmd/git.exe', 'sha256': '3' * 64}, }, 'assets': { 'detector_policy': {'path': 'app/detectors.yaml', 'sha256': policy_digest}, }, 'runtime_trees': { 'git': {'path': 'runtime/git', 'sha256': '4' * 64, 'file_count': 1}, 'python': {'path': 'runtime/python', 'sha256': '5' * 64, 'file_count': 1}, }, } class RemoteGitAssignmentBuilderTests(unittest.TestCase): def _builder(self, root, admission, planner, compatibility=None, db_factory=None): args = _source_args(root) package_manifest = compatibility or _package_manifest(args) return RemoteGitAssignmentBuilder( 'postgresql://scanner@example/db', root, {'github': args}, {'windows-fixture': {'package_manifest': package_manifest}}, 'supervisor-instance', db_factory=db_factory or _DBFactory(), admission=admission, planner=planner, ) def _multi_builder( self, root, admission, planner, compatibility=None, db_factory=None, **builder_options, ): github = _source_args(root, 'github') gitlab = _source_args(root, 'gitlab') package_manifest = compatibility or _package_manifest( github, ['github', 'gitlab'], ) return RemoteAssignmentBuilder( 'postgresql://scanner@example/db', root, {'github': github, 'gitlab': gitlab}, {'windows-fixture': {'package_manifest': package_manifest}}, 'supervisor-instance', db_factory=db_factory or _DBFactory(), admission=admission, planner=planner, **builder_options, ) @staticmethod def _claim(args, kwargs): producer = args[3] source = args[1] platform = args[2] issued_at = datetime(2026, 9, 17, tzinfo=timezone.utc) expires_at = issued_at + timedelta(seconds=kwargs['lease_seconds']) return { 'reservation_id': 41, 'reservation_token': kwargs['reservation_token'], 'bundle_id': kwargs['bundle_id'], 'scan_event_id': kwargs['scan_event_id'], 'queue_id': 9, 'claim_lease_token': 'lease-token', 'attempts': 1, 'declared_bundle_bytes': args[5], 'ready_relative_path': f"ready/{kwargs['bundle_id'][:2]}/{kwargs['bundle_id']}.trb", 'source': source, 'platform': platform, 'query': 'q', 'target': f'https://{source}.com/example/project.git', 'normalized_target': f'https://{source}.com/example/project', 'run_id': None, 'cycle_id': None, 'producer_instance_id': args[4], 'producer_pid': producer['pid'], 'producer_creation_time': producer['creation_time'], 'producer_executable': producer['executable'], 'assignment_kind': 'remote', 'remote_user_id': kwargs['remote_assignment']['user_id'], 'remote_device_id': kwargs['remote_assignment']['device_id'], 'remote_effective_config_sha256': kwargs['remote_assignment']['effective_config_sha256'], 'remote_client_compat_sha256': kwargs['remote_assignment']['client_compat_sha256'], 'remote_issued_at': issued_at.isoformat(timespec='seconds'), 'remote_expires_at': expires_at.isoformat(timespec='seconds'), 'remote_result_upload_body_timeout_seconds': kwargs[ 'remote_assignment' ]['result_upload_body_timeout_seconds'], } def test_compatibility_snapshot_is_bounded_deterministic_and_secret_free(self): with tempfile.TemporaryDirectory() as root: ensure_private_directory(root, reject_reparse=True) builder = self._multi_builder( root, lambda *_args, **_kwargs: None, lambda *_args, **_kwargs: None, ) snapshot = builder.compatibility_snapshot() self.assertEqual(set(snapshot), {'profiles', 'required_capabilities'}) self.assertEqual(len(snapshot['profiles']), 1) profile = snapshot['profiles'][0] self.assertEqual(profile['profile_name'], 'windows-fixture') self.assertEqual(profile['sources'], ['github', 'gitlab']) self.assertEqual(profile['capabilities'], sorted( profile['capabilities'], key=lambda item: ( item['source'], item['platform'], item['planning_kind'], ), )) self.assertEqual(snapshot['required_capabilities'], sorted( snapshot['required_capabilities'], key=lambda item: item['source'], )) rendered = repr(snapshot).lower() for forbidden in ( 'package_manifest', 'credential', 'auth_entry', 'token', 'file_inventory', str(root).lower(), ): self.assertNotIn(forbidden, rendered) def test_claim_uses_stable_remote_admission_identity(self): calls = [] def admission(*args, **kwargs): calls.append((args, kwargs)) return SimpleNamespace(claim=self._claim(args, kwargs)) with tempfile.TemporaryDirectory() as root: ensure_private_directory(root, reject_reparse=True) package_manifest = _package_manifest(_source_args(root)) compatibility = worker_package_build_compatibility(package_manifest) builder = self._builder( root, admission, lambda _args, _url, _source, _claim, _kwargs, **_options: {'version': 1}, package_manifest, ) builder.source_args['github'].target_claim_order = 'newest' builder.source_args['github'].drop_detectors = 'Generic' identity = { 'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256, } request_id = '1' * 32 first = builder(identity, request_id, compatibility) second = builder(identity, request_id, compatibility) self.assertEqual(first['reservation'], second['reservation']) self.assertEqual(first['deadlines'], { 'target_scan_timeout_seconds': 30, 'result_upload_body_timeout_seconds': 1800, 'assignment_ttl_seconds': 86400, 'assignment_issued_at': '2026-09-17T00:00:00+00:00', 'assignment_deadline_at': '2026-09-18T00:00:00+00:00', }) self.assertEqual(first['scan_kwargs']['git_plan'], {'version': 1}) self.assertNotIn('token', first['event_scan_options']) self.assertEqual(first['scan_policy']['drop_detectors'], ['generic']) self.assertEqual(first['scan_policy']['result_bundle_max_event_bytes'], 1 << 20) self.assertEqual(calls[0][1]['lease_seconds'], 86400) self.assertEqual(calls[0][0][5], 1 << 20) self.assertEqual(calls[0][0][6], 2 << 20) self.assertEqual(calls[0][1]['reserved_bundle_bytes'], 2 << 20) self.assertEqual(calls[0][1]['remote_max_active'], 50) self.assertEqual(calls[0][1]['claim_order'], 'newest') self.assertEqual( calls[0][1]['capacity_limits']['projection_headroom_bytes'], 2 << 20, ) self.assertEqual(calls[0][0][1:3], ('github', 'github')) self.assertEqual(calls[0][1]['reservation_token'], request_id) self.assertEqual(calls[0][1]['remote_assignment']['user_id'], 7) self.assertEqual(calls[0][1]['remote_assignment']['device_id'], 11) self.assertEqual( calls[0][1]['remote_assignment']['token_sha256'], TOKEN_SHA256, ) self.assertEqual( calls[0][1]['remote_assignment'][ 'result_upload_body_timeout_seconds' ], 1800, ) snapshot = calls[0][1]['remote_assignment']['execution_snapshot'] self.assertNotIn('token', snapshot['execution']['scan_kwargs']) self.assertEqual(snapshot['credential_ref'], {'source': 'github', 'auth_entry': ''}) self.assertEqual(calls[0][1]['bundle_id'], calls[1][1]['bundle_id']) self.assertEqual(calls[0][1]['scan_event_id'], calls[1][1]['scan_event_id']) def test_source_override_selects_lease_and_global_fallback(self): calls = [] def admission(*args, **kwargs): calls.append((args, kwargs)) return SimpleNamespace(claim=self._claim(args, kwargs)) with tempfile.TemporaryDirectory() as root: ensure_private_directory(root, reject_reparse=True) manifest = _package_manifest( _source_args(root), ['github', 'gitlab'], ) compatibility = worker_package_build_compatibility(manifest) builder = self._multi_builder( root, admission, lambda *_args, **_kwargs: {'version': 1}, manifest, assignment_ttl_seconds=3600, assignment_ttl_seconds_by_source={'gitlab': 7200}, result_upload_body_timeout_seconds=900, ) identity = { 'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256, } github = builder(identity, '0' * 32, compatibility) gitlab = builder(identity, '1' * 32, compatibility) self.assertEqual( [(args[1], kwargs['lease_seconds']) for args, kwargs in calls], [('github', 3600), ('gitlab', 7200)], ) self.assertEqual(github['deadlines']['assignment_ttl_seconds'], 3600) self.assertEqual(gitlab['deadlines']['assignment_ttl_seconds'], 7200) self.assertEqual( gitlab['deadlines']['result_upload_body_timeout_seconds'], 900, ) self.assertEqual([ kwargs['remote_assignment']['result_upload_body_timeout_seconds'] for _args, kwargs in calls ], [900, 900]) def test_every_worker_build_mismatch_does_not_admit(self): with tempfile.TemporaryDirectory() as root: package_manifest = _package_manifest(_source_args(root)) compatibility = worker_package_build_compatibility(package_manifest) admitted = [] builder = self._builder( root, lambda *args, **kwargs: admitted.append((args, kwargs)), lambda *_args, **_kwargs: {'version': 1}, package_manifest, ) for field in compatibility: bad = dict(compatibility) if field == 'platform_tag': bad[field] = 'linux-x86_64' elif field in {'protocol_version', 'bundle_format_version'}: bad[field] += 1 else: bad[field] = 'c' * 64 with self.subTest(field=field): self.assertEqual(builder({ 'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256, }, '2' * 32, bad), { 'no_assignment': {'reason': 'compatibility'}, }) self.assertEqual(admitted, []) def test_protocol1_build_is_rejected_before_admission(self): admitted = [] with tempfile.TemporaryDirectory() as root: package_manifest = _package_manifest(_source_args(root)) supplied = worker_package_build_compatibility(package_manifest) supplied['protocol_version'] = 1 builder = self._builder( root, lambda *args, **kwargs: admitted.append((args, kwargs)), lambda *_args, **_kwargs: self.fail('protocol mismatch must not plan'), package_manifest, ) self.assertEqual(builder({ 'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256, }, '2' * 32, supplied), { 'no_assignment': {'reason': 'compatibility'}, }) self.assertEqual(admitted, []) def test_protocol1_reconciliation_preserves_snapshot_without_readmission(self): identity = {'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256} for source in ('github', 'gitlab'): with self.subTest(source=source), tempfile.TemporaryDirectory() as root: args = _source_args(root, source) manifest = _package_manifest(args, [source]) compatibility = worker_package_build_compatibility(manifest) captured = [] def admission(*call_args, **kwargs): claim = self._claim(call_args, kwargs) captured.append(claim) return SimpleNamespace(claim=claim) initial = RemoteAssignmentBuilder( 'postgresql://scanner@example/db', root, {source: args}, {'windows-fixture': {'package_manifest': manifest}}, 'supervisor-instance', db_factory=_DBFactory(), admission=admission, planner=lambda *_args, **_kwargs: {'version': 1}, credential_refs={source: ''}, )(identity, '2' * 32, compatibility) legacy_snapshot = copy.deepcopy(initial['execution_snapshot']) legacy_snapshot['compatibility']['protocol_version'] = 1 legacy_sha256 = remote_execution_snapshot_sha256(legacy_snapshot) legacy_build = dict(compatibility) legacy_build['protocol_version'] = 1 replay_db = _DBFactory() replay_db.reconciled = { 'state': 'committed', 'claim': captured[0], 'execution_snapshot': legacy_snapshot, 'execution_snapshot_sha256': legacy_sha256, 'execution_plan': initial['execution_plan'], 'receipt': None, } replay = RemoteAssignmentBuilder( 'postgresql://scanner@example/db', root, {source: args}, {'windows-fixture': {'package_manifest': manifest}}, 'supervisor-instance', db_factory=replay_db, admission=lambda *_args, **_kwargs: self.fail('must not readmit'), planner=lambda *_args, **_kwargs: self.fail('must not replan'), credential_refs={source: ''}, assignment_ttl_seconds=600, result_upload_body_timeout_seconds=300, )(identity, '2' * 32, legacy_build) self.assertEqual(replay['execution_snapshot'], legacy_snapshot) self.assertEqual(replay['execution_snapshot_sha256'], legacy_sha256) self.assertEqual(replay['execution_plan'], initial['execution_plan']) self.assertEqual( { name: replay['deadlines'][name] for name in ( 'assignment_ttl_seconds', 'assignment_issued_at', 'assignment_deadline_at', ) }, { name: initial['deadlines'][name] for name in ( 'assignment_ttl_seconds', 'assignment_issued_at', 'assignment_deadline_at', ) }, ) self.assertEqual( replay['deadlines']['result_upload_body_timeout_seconds'], 1800, ) self.assertEqual(replay_db.plan_recovery_calls, []) legacy_db = _DBFactory() legacy_claim = dict(captured[0]) legacy_claim.pop('remote_result_upload_body_timeout_seconds') legacy_db.reconciled = { 'state': 'committed', 'claim': legacy_claim, 'execution_snapshot': legacy_snapshot, 'execution_snapshot_sha256': legacy_sha256, 'execution_plan': initial['execution_plan'], 'receipt': None, } legacy_replay = RemoteAssignmentBuilder( 'postgresql://scanner@example/db', root, {source: args}, {'windows-fixture': {'package_manifest': manifest}}, 'supervisor-instance', db_factory=legacy_db, admission=lambda *_args, **_kwargs: self.fail('must not readmit'), planner=lambda *_args, **_kwargs: self.fail('must not replan'), credential_refs={source: ''}, result_upload_body_timeout_seconds=300, )(identity, '2' * 32, legacy_build) self.assertEqual( legacy_replay['deadlines'][ 'result_upload_body_timeout_seconds' ], 300, ) unavailable = RemoteAssignmentBuilder( 'postgresql://scanner@example/db', root, {source: args}, {'windows-fixture': {'package_manifest': manifest}}, 'supervisor-instance', db_factory=replay_db, admission=lambda *_args, **_kwargs: self.fail('must not readmit'), planner=lambda *_args, **_kwargs: self.fail('must not replan'), credential_refs={source: 'different'}, ) with self.assertRaisesRegex(RuntimeError, 'credential reference'): unavailable(identity, '2' * 32, legacy_build) def test_profile_source_requires_exact_package_capability(self): with tempfile.TemporaryDirectory() as root: github = _source_args(root, 'github') gitlab_only = _package_manifest(github, ['gitlab']) with self.assertRaisesRegex(ValueError, 'invalid sources'): RemoteAssignmentBuilder( 'postgresql://scanner@example/db', root, {'github': github}, {'windows-fixture': { 'package_manifest': gitlab_only, 'sources': ['github'], }}, 'supervisor-instance', admission=lambda *_args, **_kwargs: None, planner=lambda *_args, **_kwargs: None, ) def test_clean_no_work_admission_does_not_plan(self): with tempfile.TemporaryDirectory() as root: package_manifest = _package_manifest(_source_args(root)) compatibility = worker_package_build_compatibility(package_manifest) builder = self._builder( root, lambda *_args, **_kwargs: SimpleNamespace( claim=None, reason='no_claimable_target', ), lambda *_args, **_kwargs: self.fail('no-work admission must not plan'), package_manifest, ) result = builder( {'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256}, 'a' * 32, compatibility, ) self.assertEqual(result, { 'no_assignment': {'reason': 'empty_queue'}, }) def test_multisource_no_work_falls_back_in_deterministic_wraparound_order(self): calls = [] def admission(*args, **kwargs): calls.append((args, kwargs)) claim = None if args[1] == 'github' else self._claim(args, kwargs) return SimpleNamespace( claim=claim, reason='no_claimable_target' if claim is None else None, ) identity = {'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256} request_id = '2' * 32 with tempfile.TemporaryDirectory() as root: manifest = _package_manifest( _source_args(root), ['github', 'gitlab'], ) compatibility = worker_package_build_compatibility(manifest) result = self._multi_builder( root, admission, lambda *_args, **_kwargs: {'version': 1}, manifest, )(identity, request_id, compatibility) self.assertEqual([call[0][1] for call in calls], ['github', 'gitlab']) self.assertEqual(result['reservation']['source'], 'gitlab') tokens = [call[1]['reservation_token'] for call in calls] self.assertEqual(len(set(tokens)), 2) self.assertNotIn(request_id, tokens) for _args, kwargs in calls: self.assertEqual(len(kwargs['reservation_token']), 32) self.assertEqual( kwargs['remote_assignment']['execution_snapshot'][ 'credential_ref' ]['source'], _args[1], ) def test_multisource_rotation_wraps_and_stops_after_first_claim(self): calls = [] def admission(*args, **kwargs): calls.append((args, kwargs)) return SimpleNamespace(claim=self._claim(args, kwargs)) identity = {'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256} with tempfile.TemporaryDirectory() as root: manifest = _package_manifest( _source_args(root), ['github', 'gitlab'], ) compatibility = worker_package_build_compatibility(manifest) result = self._multi_builder( root, admission, lambda *_args, **_kwargs: {'version': 1}, manifest, )(identity, '1' * 32, compatibility) self.assertEqual([call[0][1] for call in calls], ['gitlab']) self.assertEqual(result['reservation']['source'], 'gitlab') def test_multisource_retry_recovers_same_fallback_without_readmission(self): calls = [] def admission(*args, **kwargs): calls.append((args, kwargs)) if args[1] == 'github': return SimpleNamespace( claim=None, reason='no_claimable_target', ) return SimpleNamespace(claim=self._claim(args, kwargs)) identity = {'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256} request_id = '2' * 32 with tempfile.TemporaryDirectory() as root: manifest = _package_manifest( _source_args(root), ['github', 'gitlab'], ) compatibility = worker_package_build_compatibility(manifest) first = self._multi_builder( root, admission, lambda *_args, **_kwargs: {'version': 1}, manifest, )(identity, request_id, compatibility) github_token = calls[0][1]['reservation_token'] gitlab_token = calls[1][1]['reservation_token'] snapshot = calls[1][1]['remote_assignment']['execution_snapshot'] db_factory = _DBFactory() db_factory.reconciled_by_request = { github_token: {'state': 'aborted', 'receipt': None}, gitlab_token: { 'state': 'committed', 'claim': first['reservation'], 'execution_snapshot': snapshot, 'execution_snapshot_sha256': remote_execution_snapshot_sha256( snapshot, ), 'execution_plan': { 'kind': 'exact_git_v1', 'execution_target': first['reservation']['target'], 'bound_plan': first['scan_kwargs']['git_plan'], }, 'git_plan': first['scan_kwargs']['git_plan'], 'receipt': None, }, } replay = self._multi_builder( root, lambda *_args, **_kwargs: self.fail('must not readmit'), lambda *_args, **_kwargs: self.fail('must not replan'), manifest, db_factory, )(identity, request_id, compatibility) self.assertEqual(replay, first) self.assertEqual( db_factory.reconcile_calls, [request_id, github_token, gitlab_token], ) def test_replayed_claim_reuses_bound_plan_without_provider_resolution(self): db_factory = _DBFactory() db_factory.bound_plan = {'version': 1, 'head_sha': 'a' * 40} def admission(*args, **kwargs): return SimpleNamespace(claim=self._claim(args, kwargs)) with tempfile.TemporaryDirectory() as root: package_manifest = _package_manifest(_source_args(root)) compatibility = worker_package_build_compatibility(package_manifest) builder = self._builder( root, admission, lambda *_args, **_kwargs: self.fail( 'provider resolution must not be repeated' ), package_manifest, db_factory, ) result = builder( { 'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256, }, '4' * 32, compatibility, ) self.assertEqual(result['scan_kwargs']['git_plan'], db_factory.bound_plan) self.assertEqual(db_factory.plan_recovery_calls, [(41, 11, 'lease-token')]) def test_direct_claim_and_replay_never_use_git_planning_or_credentials(self): identity = {'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256} digest = 'sha256:' + ('a' * 64) cases = ( ('dockerhub', f'docker.io/library/alpine@{digest}'), ('huggingface', 'Owner/Space'), ) for source, target in cases: with self.subTest(source=source), tempfile.TemporaryDirectory() as root: args = _direct_source_args(root, source) manifest = _package_manifest(args, [source]) compatibility = worker_package_build_compatibility(manifest) captured = [] def admission(*call_args, **kwargs): claim = self._claim(call_args, kwargs) claim['target'] = target claim['normalized_target'] = ( target.lower() if source == 'huggingface' else target ) captured.append(claim) return SimpleNamespace(claim=claim) first_db = _DBFactory() first = RemoteAssignmentBuilder( 'postgresql://scanner@example/db', root, {source: args}, {'windows-fixture': {'package_manifest': manifest}}, 'supervisor-instance', db_factory=first_db, admission=admission, planner=lambda *_args, **_kwargs: self.fail( 'direct assignment must not invoke the Git planner' ), credential_refs={source: ''}, )(identity, 'c' * 32, compatibility) self.assertEqual(first['scan_kwargs'], first['event_scan_options']) self.assertEqual(first['execution_plan']['bound_plan'], None) self.assertNotIn('token', first['scan_kwargs']) self.assertNotIn('git_plan', first['scan_kwargs']) self.assertEqual(first_db.plan_recovery_calls, []) replay_db = _DBFactory() replay_db.reconciled = { 'state': 'committed', 'claim': captured[0], 'execution_snapshot': first['execution_snapshot'], 'execution_snapshot_sha256': first['execution_snapshot_sha256'], 'execution_plan': first['execution_plan'], 'receipt': None, } replay = RemoteAssignmentBuilder( 'postgresql://scanner@example/db', root, {source: args}, {'windows-fixture': {'package_manifest': manifest}}, 'supervisor-instance', db_factory=replay_db, admission=lambda *_args, **_kwargs: self.fail('must not readmit'), planner=lambda *_args, **_kwargs: self.fail('must not replan'), credential_refs={source: ''}, )(identity, 'c' * 32, compatibility) self.assertEqual(replay, first) self.assertEqual(replay_db.plan_recovery_calls, []) def test_lost_reply_uses_durable_snapshot_after_live_profile_changes(self): captured = [] def admission(*args, **kwargs): claim = self._claim(args, kwargs) captured.append((claim, kwargs['remote_assignment']['execution_snapshot'])) return SimpleNamespace(claim=claim) identity = {'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256} request_id = '5' * 32 with tempfile.TemporaryDirectory() as root: original_manifest = _package_manifest(_source_args(root)) original_build = worker_package_build_compatibility(original_manifest) first = self._builder( root, admission, lambda *_args, **_kwargs: {'version': 1, 'head_sha': 'a' * 40}, original_manifest, )(identity, request_id, original_build) claim, snapshot = captured[0] changed_manifest = _package_manifest(_source_args(root)) first_name = next(iter(changed_manifest['files'])) changed_manifest['files'][first_name]['sha256'] = '9' * 64 db_factory = _DBFactory() db_factory.reconciled = { 'state': 'committed', 'claim': claim, 'execution_snapshot': snapshot, 'execution_snapshot_sha256': remote_execution_snapshot_sha256(snapshot), 'git_plan': first['scan_kwargs']['git_plan'], 'receipt': None, } replay_builder = self._builder( root, lambda *_args, **_kwargs: self.fail('must not readmit'), lambda *_args, **_kwargs: self.fail('must not replan'), changed_manifest, db_factory, ) replay_builder.source_args['github'].drop_detectors = 'changed-live-policy' replay = replay_builder(identity, request_id, original_build) self.assertEqual(replay, first) self.assertEqual(db_factory.reconcile_calls, [request_id]) def test_lost_reply_consumes_generic_execution_plan_without_replanning(self): captured = [] def admission(*args, **kwargs): claim = self._claim(args, kwargs) captured.append((claim, kwargs['remote_assignment']['execution_snapshot'])) return SimpleNamespace(claim=claim) identity = {'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256} request_id = 'a' * 32 with tempfile.TemporaryDirectory() as root: manifest = _package_manifest(_source_args(root)) compatibility = worker_package_build_compatibility(manifest) first = self._builder( root, admission, lambda *_args, **_kwargs: {'version': 1}, manifest, )(identity, request_id, compatibility) claim, snapshot = captured[0] db_factory = _DBFactory() db_factory.reconciled = { 'state': 'committed', 'claim': claim, 'execution_snapshot': snapshot, 'execution_snapshot_sha256': remote_execution_snapshot_sha256(snapshot), 'execution_plan': { 'kind': 'exact_git_v1', 'execution_target': claim['target'], 'bound_plan': first['scan_kwargs']['git_plan'], }, 'git_plan': first['scan_kwargs']['git_plan'], 'receipt': None, } replay = self._builder( root, lambda *_args, **_kwargs: self.fail('must not readmit'), lambda *_args, **_kwargs: self.fail('must not replan'), manifest, db_factory, )(identity, request_id, compatibility) self.assertEqual(replay, first) self.assertEqual(db_factory.plan_recovery_calls, []) def test_recovered_execution_plan_kind_mismatch_fails_closed(self): captured = [] def admission(*args, **kwargs): claim = self._claim(args, kwargs) captured.append((claim, kwargs['remote_assignment']['execution_snapshot'])) return SimpleNamespace(claim=claim) identity = {'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256} with tempfile.TemporaryDirectory() as root: manifest = _package_manifest(_source_args(root)) compatibility = worker_package_build_compatibility(manifest) self._builder( root, admission, lambda *_args, **_kwargs: {'version': 1}, manifest, )(identity, 'b' * 32, compatibility) claim, snapshot = captured[0] db_factory = _DBFactory() db_factory.reconciled = { 'state': 'committed', 'claim': claim, 'execution_snapshot': snapshot, 'execution_snapshot_sha256': remote_execution_snapshot_sha256(snapshot), 'execution_plan': { 'kind': 'docker_direct_v1', 'execution_target': claim['target'], 'bound_plan': None, }, 'receipt': None, } builder = self._builder( root, lambda *_args, **_kwargs: self.fail('must not readmit'), lambda *_args, **_kwargs: self.fail('must not replan'), manifest, db_factory, ) with self.assertRaisesRegex(ScanEventConflictError, 'execution plan'): builder(identity, 'b' * 32, compatibility) self.assertEqual(db_factory.plan_recovery_calls, []) def test_lost_reply_returns_an_explicit_terminal_resolution(self): receipt = { 'receipt_id': '6' * 64, 'resolution': 'expired', 'reservation_id': 41, 'bundle_id': '7' * 32, 'scan_event_id': '8' * 32, 'resolved_at': '2026-09-18T00:00:00+00:00', } db_factory = _DBFactory() db_factory.reconciled = { 'state': 'committed', 'receipt': receipt, } identity = {'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256} with tempfile.TemporaryDirectory() as root: package_manifest = _package_manifest(_source_args(root)) compatibility = worker_package_build_compatibility(package_manifest) result = self._builder( root, lambda *_args, **_kwargs: self.fail('must not readmit'), lambda *_args, **_kwargs: self.fail('must not replan'), package_manifest, db_factory, )(identity, '6' * 32, compatibility) self.assertEqual(result, {'resolution': receipt}) self.assertEqual(db_factory.reconcile_calls, ['6' * 32]) def test_planning_failure_is_published_and_accepted_centrally(self): db_factory = _DBFactory() def admission(*args, **kwargs): return SimpleNamespace(claim=self._claim(args, kwargs)) def planner(*_args, **_kwargs): raise _GitResolutionFailure(ValueError('bad target'), invalid_target=True) with tempfile.TemporaryDirectory() as root: ensure_private_directory(root, reject_reparse=True) db_factory.bundle_root = root package_manifest = _package_manifest(_source_args(root)) compatibility = worker_package_build_compatibility(package_manifest) builder = self._builder(root, admission, planner, package_manifest, db_factory) result = builder( { 'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256, }, '3' * 32, compatibility, ) self.assertIsNone(result) self.assertEqual(len(db_factory.accepted), 1) accepted = db_factory.accepted[0] self.assertEqual(accepted[0], 41) self.assertRegex(accepted[2]['payload_sha256'], r'^[0-9a-f]{64}$') published = os.path.join(root, accepted[1]['relative_path']) metadata = ResultBundleReader(published, max_event_bytes=1 << 20).validate() self.assertEqual(metadata.error_count, 1) effective_diagnostics = ResultBundleReader( published, max_event_bytes=1 << 20, ).effective_diagnostics() self.assertEqual( accepted[1]['effective_diagnostic_projection_version'], DIAGNOSTIC_PROJECTION_VERSION, ) self.assertEqual( accepted[1]['effective_diagnostic_count'], len(effective_diagnostics), ) self.assertEqual( accepted[1]['effective_diagnostic_uids_sha256'], ordered_diagnostic_uid_set_sha256(effective_diagnostics), ) def test_planning_failure_recovery_replays_the_same_projection_authority(self): db_factory = _DBFactory() db_factory.accept_failures = 1 captured = [] def admission(*args, **kwargs): claim = self._claim(args, kwargs) captured.append((claim, kwargs['remote_assignment']['execution_snapshot'])) return SimpleNamespace(claim=claim) def planner(*_args, **_kwargs): raise _GitResolutionFailure(ValueError('bad target'), invalid_target=True) identity = {'user_id': 7, 'device_id': 11, 'token_sha256': TOKEN_SHA256} request_id = 'e' * 32 with tempfile.TemporaryDirectory() as root: ensure_private_directory(root, reject_reparse=True) db_factory.bundle_root = root package_manifest = _package_manifest(_source_args(root)) compatibility = worker_package_build_compatibility(package_manifest) builder = self._builder( root, admission, planner, package_manifest, db_factory, ) with self.assertRaisesRegex(RuntimeError, 'synthetic acceptance'): builder(identity, request_id, compatibility) claim, snapshot = captured[0] db_factory.reconciled = { 'state': 'committed', 'claim': claim, 'execution_snapshot': snapshot, 'execution_snapshot_sha256': remote_execution_snapshot_sha256(snapshot), 'git_plan': None, 'receipt': None, } replay = self._builder( root, lambda *_args, **_kwargs: self.fail('must not readmit'), lambda *_args, **_kwargs: self.fail('must not replan'), package_manifest, db_factory, )(identity, request_id, compatibility) self.assertIsNone(replay) self.assertEqual(len(db_factory.accepted), 1) published = bundle_ready_path(root, claim['bundle_id']) effective_diagnostics = ResultBundleReader( published, max_event_bytes=1 << 20, ).effective_diagnostics() accepted_metadata = db_factory.accepted[0][1] self.assertEqual( accepted_metadata['effective_diagnostic_count'], len(effective_diagnostics), ) self.assertEqual( accepted_metadata['effective_diagnostic_uids_sha256'], ordered_diagnostic_uid_set_sha256(effective_diagnostics), ) def test_remote_result_plan_must_match_bound_canonical_plan(self): plan = {'version': 1, 'head_sha': 'a' * 40} encoded = canonical_git_scan_plan_bytes(plan) reservation = { 'git_scan_plan_json': encoded.decode('ascii'), 'git_scan_plan_sha256': hashlib.sha256(encoded).hexdigest(), } self.assertEqual( validate_result_git_scan_plan(reservation, {'git_scan_plan': dict(plan)}), plan, ) with self.assertRaises(ScanEventConflictError): validate_result_git_scan_plan( reservation, {'git_scan_plan': {'version': 1, 'head_sha': 'b' * 40}}, ) if __name__ == '__main__': unittest.main()