Files
truf-server/tests/test_worker_assignment.py
2026-09-30 20:30:56 +03:00

1127 lines
50 KiB
Python

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