import asyncio import hashlib import html import json import os from pathlib import Path import re import sys import tempfile import threading from types import SimpleNamespace import unittest import uuid from unittest import mock from urllib.parse import urlencode, urljoin ROOT = Path(__file__).resolve().parents[1] APP_DIR = ROOT / 'app' sys.path.insert(0, str(APP_DIR)) from starlette.testclient import TestClient from starlette.datastructures import Headers from starlette.requests import Request from starlette.responses import Response from admin_api import ( ADMIN_PREFIX, EDGE_MARKER_HEADER, OPERATOR_HEADER, SECURITY_HEADERS, AdminAPIError, AdminService, _dispatch_managed_file_mutation, _dispatch, _dispatch_runtime_document, _download_managed_file, _form_fields, _managed_file_download_response, _ManagedFileStreamingResponse, _parse_managed_file_content, _trusted_operator, _validate_cap, ) from managed_files import ( ManagedFileAccessError, ManagedFileDirectoryEntry, ManagedFileDownload, ManagedFileIdentity, ManagedFileLimits, ManagedFileListing, ManagedFileMutation, ManagedFilePermissions, ManagedFileRoot, ManagedFileRootRegistry, ) from scanner_db import ( RuntimeControlRevisionConflictError, RuntimeOperationIdentityConflictError, ) from runtime_document import RuntimeDocumentError from runtime_security import ensure_private_directory from worker_api import WorkerAPIError, create_worker_app from worker_contracts import ( MAX_DIAGNOSTIC_BODY_BYTES, MAX_DIAGNOSTIC_LOG_BYTES, make_body_material, make_log_material, ) ORIGIN = 'https://admin.example.test' MARKER = 'fixture-edge-marker-value-32bytes-minimum' OPERATOR = 'fixture.operator' class FakeWorkerService: def __init__(self, root): self.bundle_root = root self.max_bundle_bytes = 1024 * 1024 self.claim_retry_after_seconds = 5 def authenticate(self, authorization): raise WorkerAPIError(401, 'unauthorized', 'worker credentials are invalid') def reap(self): return [] class AdminValidationTests(unittest.TestCase): def test_assignment_cap_accepts_integer_zero(self): self.assertEqual(_validate_cap(0), 0) class FakeManagedFileTraversal: def __init__(self): self.files = { 'guide.txt': b'fixture managed content', 'nested/report.bin': b'\x00\xffreport', } self.calls = [] self.closed = 0 self.error_category = None self.stream_downloads = False self.last_snapshot = None @staticmethod def _identity(content): return ManagedFileIdentity(hashlib.sha256(content).hexdigest(), len(content)) def _fail(self): if self.error_category: raise ManagedFileAccessError(self.error_category) def list_directory(self, root_id, relative_path=None): self._fail() self.calls.append(('list', root_id, relative_path)) prefix = '' if relative_path is None else relative_path + '/' entries = {} for path, content in self.files.items(): if not path.startswith(prefix): continue remainder = path[len(prefix):] name, separator, _tail = remainder.partition('/') entries[name] = ( ManagedFileDirectoryEntry(name, 'directory', None) if separator else ManagedFileDirectoryEntry(name, 'file', len(content)) ) values = tuple(entries[name] for name in sorted(entries)) return ManagedFileListing(values, sum(len(item.name.encode()) for item in values)) def download_file(self, root_id, relative_path): self._fail() self.calls.append(('download', root_id, relative_path)) if relative_path not in self.files: raise ManagedFileAccessError('not_found') content = self.files[relative_path] if self.stream_downloads: snapshot = FakeManagedFileSnapshot(content) self.last_snapshot = snapshot return ManagedFileDownload( self._identity(content), snapshot=snapshot, ) return ManagedFileDownload(self._identity(content), content) def mutation_file_identity( self, root_id, relative_path, operation, *, require_private_sha256=None): self._fail() self.calls.append(( 'identity', root_id, relative_path, operation.value, require_private_sha256, )) if relative_path not in self.files: raise ManagedFileAccessError('not_found') return self._identity(self.files[relative_path]) def create_replace_file(self, root_id, relative_path, payload, *, expected_sha256): self._fail() self.calls.append(( 'create_replace', root_id, relative_path, expected_sha256, hashlib.sha256(payload).hexdigest(), len(payload), )) before_content = self.files.get(relative_path) if expected_sha256 is None: if before_content is not None: raise ManagedFileAccessError('hash_conflict') elif before_content is None: raise ManagedFileAccessError('not_found') elif self._identity(before_content).sha256 != expected_sha256: raise ManagedFileAccessError('hash_conflict') before = self._identity(before_content) if before_content is not None else None after = self._identity(payload) written = before != after self.files[relative_path] = payload return ManagedFileMutation(before, after, written) def delete_file(self, root_id, relative_path, *, expected_sha256): self._fail() self.calls.append(('delete', root_id, relative_path, expected_sha256)) if relative_path not in self.files: raise ManagedFileAccessError('not_found') before = self._identity(self.files[relative_path]) if before.sha256 != expected_sha256: raise ManagedFileAccessError('hash_conflict') del self.files[relative_path] return ManagedFileMutation(before, None, True) def close(self): self.closed += 1 class FakeManagedFileSnapshot: def __init__(self, content): self.content = content self.closed = 0 def chunks(self): try: for offset in range(0, len(self.content), 3): yield self.content[offset:offset + 3] finally: self.close() def close(self): if not self.closed: self.closed = 1 class RecordingDB: enabled = True calls = [] fail_queue = False source_operations = {} worker_admin_operations = {} document_operations = {} apply_operations = {} managed_file_operations = {} audit_events = [] audit_next_before_event_id = None source_identity_conflict = False completion_failures = 0 control_identity_conflict = False control_revision_conflict = False document_completion_failures = 0 managed_file_completion_failures = 0 managed_file_execution_lock = threading.Lock() def __init__(self, **kwargs): self.managed_file_execution_held = False self.calls.append(('open', kwargs)) def close(self): if self.managed_file_execution_held: self.managed_file_execution_held = False self.managed_file_execution_lock.release() self.calls.append(('close',)) def acquire_runtime_managed_file_execution(self, operation_id): if self.managed_file_execution_held: raise RuntimeError('managed file execution lock is already held') self.managed_file_execution_lock.acquire() self.managed_file_execution_held = True self.calls.append(('managed_file_execution_acquire', operation_id)) return True def release_runtime_managed_file_execution(self, operation_id): if not self.managed_file_execution_held: raise RuntimeError('managed file execution lock is not held') self.managed_file_execution_held = False self.managed_file_execution_lock.release() self.calls.append(('managed_file_execution_release', operation_id)) return True def admin_remote_worker_snapshot(self, limit, filters=None): filters = dict(filters or {}) self.calls.append(('snapshot', limit, filters)) common = { 'queue_id': 5, 'user_key': 'fixture-user', 'active_assignment_cap': 2, 'device_key': 'fixture-device', 'source': 'gitlab', 'target': 'example/project@commit', 'issued_at': '2026-09-20T00:01:00+00:00', 'assignment_deadline_at': '2026-09-21T00:01:00+00:00', 'remote_result_upload_body_timeout_seconds': 1800, 'finished_at': None, 'duration_seconds': None, 'accepted': False, 'ingested': False, 'assignment_code': None, 'target_scan_id': None, 'scan_error_count': None, 'first_error_summary': None, 'skipped_reason': None, 'diagnostic_categories': None, 'diagnostic_codes': None, 'primary_diagnostic': None, 'phase_started_at': None, 'last_progress_at': None, 'last_progress_received_at': None, 'scan_deadline_at': None, 'scan_remaining_seconds': None, 'assignment_remaining_seconds': 3600, 'ingestion_state': None, 'projection_state': None, 'protocol_version': '2', 'bundle_format_version': '2', 'platform_tag': 'linux-x86_64', 'code_manifest_sha256': '1' * 64, 'detector_policy_sha256': '2' * 64, 'effective_config_sha256': '3' * 64, } assignments = [{ **common, 'reservation_id': 4, 'assignment_outcome': 'unfinished', 'scan_outcome': 'unavailable', 'diagnostic_count': 0, 'diagnostic_projection_version': None, 'active_phase': 'scanning', 'phase_age_seconds': 17, 'last_progress_age_seconds': 3, 'slot_id': 0, 'scan_deadline_at': '2026-09-20T00:11:00+00:00', 'scan_remaining_seconds': 500, }, { **common, 'reservation_id': 5, 'target_scan_id': 50, 'assignment_outcome': 'accepted', 'scan_outcome': 'degraded', 'finished_at': '2026-09-20T00:10:00+00:00', 'scan_warning_class': 'detector_timeout', 'scan_warning_summary': 'bounded detector output timed out', 'accepted': True, 'ingested': True, 'diagnostic_count': 2, 'diagnostic_projection_version': 1, 'diagnostic_categories': 'rate_limit, scanner', 'diagnostic_codes': 'provider.rate_limit, scanner.exit', 'primary_diagnostic': 'scanner/scanner.exit', 'active_phase': 'awaiting_receipt', 'phase_age_seconds': 4, 'last_progress_age_seconds': 2, 'slot_id': 1, 'ingestion_state': 'acknowledged', 'projection_state': 'completed', }, { **common, 'reservation_id': 6, 'assignment_outcome': 'prebundle_failed', 'scan_outcome': 'unavailable', 'finished_at': '2026-09-20T00:09:00+00:00', 'diagnostic_count': 1, 'diagnostic_projection_version': 1, 'diagnostic_categories': 'storage', 'diagnostic_codes': 'bundle.fsync', 'primary_diagnostic': 'storage/bundle.fsync', 'active_phase': 'bundling', 'phase_age_seconds': 9, 'last_progress_age_seconds': 8, 'slot_id': 0, }, { **common, 'reservation_id': 7, 'assignment_outcome': 'expired', 'scan_outcome': 'unavailable', 'finished_at': '2026-09-21T00:01:00+00:00', 'diagnostic_count': 1, 'diagnostic_projection_version': 1, 'diagnostic_categories': 'assignment_expired', 'diagnostic_codes': 'assignment.expired', 'primary_diagnostic': 'assignment_expired/assignment.expired', 'active_phase': 'scanning', 'phase_age_seconds': 300, 'last_progress_age_seconds': 280, 'slot_id': 0, }, { **common, 'reservation_id': 8, 'target_scan_id': 80, 'assignment_outcome': 'accepted', 'scan_outcome': 'error', 'finished_at': '2026-09-20T00:08:00+00:00', 'accepted': True, 'diagnostic_count': 0, 'diagnostic_projection_version': None, 'active_phase': None, 'phase_age_seconds': None, 'last_progress_age_seconds': None, 'slot_id': None, 'remote_result_upload_body_timeout_seconds': None, 'protocol_version': '1', }, { **common, 'reservation_id': 9, 'assignment_outcome': 'unfinished', 'scan_outcome': 'unavailable', 'diagnostic_count': 0, 'diagnostic_projection_version': 1, 'active_phase': None, 'phase_age_seconds': None, 'last_progress_age_seconds': None, 'slot_id': None, }] return { 'filters': filters, 'users': [{'user_key': 'fixture-user', 'active_assignment_cap': 2, 'disabled': False}], 'workers': [{ 'device_key': 'fixture-device', 'user_key': 'fixture-user', 'active_assignment_cap': 2, 'active_slot_count': 1, 'current_phases': 'scanning', 'latest_progress_age_seconds': 3, 'known_reasons': None, 'pending_local_recovery': False, 'active_package_identity': 'linux-x86_64:' + ('1' * 12), 'unfinished_count': 1, 'completed_count': 2, 'failed_count': 0, 'expired_count': 0, 'last_contact_at': '2026-09-20T00:02:00+00:00', 'revoked': False, }], 'assignments': assignments, 'deferred_queue': [], } def admin_worker_diagnostic_groups(self, limit, filters=None, *, occurrence_offset=0): filters = dict(filters or {}) self.calls.append(('diagnostic_groups', limit, filters, occurrence_offset)) occurrences = [{ 'diagnostic_uid': 'b' * 64, 'reservation_id': 5, 'target_scan_id': 50, 'source': 'gitlab', 'worker': 'fixture-device', 'user': 'fixture-user', 'assignment_outcome': 'accepted', 'scan_outcome': 'error', 'phase': 'scanning', 'kind': 'provider_http', 'category': 'rate_limit', 'code': 'provider.rate_limit', 'summary': 'rate limited', 'retryable': True, 'occurred_at': '2026-09-20T00:03:00Z', 'received_at': '2026-09-20T00:03:01Z', }, { 'diagnostic_uid': 'c' * 64, 'reservation_id': 6, 'target_scan_id': None, 'source': 'gitlab', 'worker': 'fixture-device', 'user': 'fixture-user', 'assignment_outcome': 'prebundle_failed', 'scan_outcome': 'unavailable', 'phase': 'scanning', 'kind': 'provider_http', 'category': 'rate_limit', 'code': 'provider.rate_limit', 'summary': 'rate limited again', 'retryable': True, 'occurred_at': '2026-09-20T00:02:00Z', 'received_at': '2026-09-20T00:02:01Z', }] return { 'filters': filters, 'occurrence_limit': limit, 'matched_occurrence_count': 201, 'page_occurrence_count': 2, 'occurrence_offset': occurrence_offset, 'has_previous': occurrence_offset > 0, 'has_next': occurrence_offset == 0, 'previous_occurrence_offset': ( max(0, occurrence_offset - limit) if occurrence_offset else None ), 'next_occurrence_offset': limit if occurrence_offset == 0 else None, 'truncated': occurrence_offset == 0, 'groups': [{ 'fingerprint': 'f' * 64, 'count': 2, 'affected_assignment_count': 2, 'page_occurrence_count': 2, 'affected_assignments': [5, 6], 'occurrences': occurrences, }], } def admin_worker_duration_metrics(self, limit, filters=None, *, offset=0): filters = dict(filters or {}) self.calls.append(('duration_metrics', limit, filters, offset)) metrics = [{ 'source': 'gitlab', 'phase': 'scanning', 'outcome': 'error', 'sample_count': 12, 'sufficient': True, 'minimum_sample_count': 5, 'p50_seconds': 10.0, 'p95_seconds': 18.0, 'p99_seconds': 21.0, }, { 'source': 'gitlab', 'phase': 'uploading', 'outcome': 'accepted', 'sample_count': 2, 'sufficient': False, 'minimum_sample_count': 5, 'p50_seconds': 3.0, 'p95_seconds': 4.0, 'p99_seconds': 4.0, }] return { 'filters': filters, 'metrics': metrics, 'total_group_count': 201, 'metric_limit': limit, 'metric_offset': offset, 'page_group_count': len(metrics), 'has_previous': offset > 0, 'has_next': offset == 0, 'previous_metric_offset': max(0, offset - limit) if offset else None, 'next_metric_offset': limit if offset == 0 else None, 'truncated': offset == 0, } @staticmethod def _diagnostic_fixture(reservation_id=5, diagnostic_uid=None): def material_value(material): return { 'encoding': material.encoding.value, 'head': material.head, 'tail': material.tail, 'original_size': material.original_size, 'stored_size': material.stored_size, 'sha256': material.sha256, 'truncated': material.truncated, } material = material_value(make_body_material( b'b' * MAX_DIAGNOSTIC_BODY_BYTES, )) truncated = material_value(make_log_material( b'l' * (MAX_DIAGNOSTIC_LOG_BYTES + 1), )) diagnostic_uid = diagnostic_uid or ('b' * 64) return { 'schema': 1, 'diagnostic_uid': diagnostic_uid, 'occurrence_id': 'fixture-occurrence', 'reservation_id': reservation_id, 'scan_event_id': None, 'slot_id': 0, 'source': 'gitlab', 'phase': 'scanning', 'kind': 'provider_http', 'category': 'rate_limit', 'code': 'provider.rate_limit', 'summary': 'provider returned a limit', 'retryable': True, 'attempt': 1, 'assignment_outcome': None, 'scan_outcome': None, 'occurred_at': '2026-09-20T00:03:00Z', 'captured_at': '2026-09-20T00:03:01Z', 'received_at': None, 'http': { 'operation': 'GET', 'status_code': 429, 'content_type': 'text/plain', 'request_id': 'fixture', 'body': material, 'headers': None, }, 'process': { 'name': 'trufflehog', 'exit_code': 1, 'signal': None, 'timed_out': False, 'stdout': None, 'stderr': truncated, }, 'exception': None, } def admin_worker_assignment_detail(self, reservation_id, **_kwargs): self.calls.append(('assignment_detail', reservation_id, _kwargs)) reservation_id = int(reservation_id) states = { 4: ('unfinished', 'unavailable', None, 'scanning'), 5: ('accepted', 'error', '2026-09-20T00:06:00Z', 'awaiting_receipt'), 6: ('prebundle_failed', 'unavailable', '2026-09-20T00:04:00Z', 'bundling'), 7: ('expired', 'unavailable', '2026-09-20T00:11:00Z', 'scanning'), 8: ('accepted', 'error', '2026-09-20T00:06:00Z', 'awaiting_receipt'), } if reservation_id not in states: return None assignment_outcome, scan_outcome, resolved_at, current_phase = states[reservation_id] legacy = reservation_id == 8 diagnostic_uid = ({5: 'b', 6: 'c', 7: 'd'}.get(reservation_id, 'b')) * 64 envelope = self._diagnostic_fixture(reservation_id, diagnostic_uid) canonical = json.dumps(envelope, ensure_ascii=True, sort_keys=True, separators=(',', ':')) diagnostic_items = [] if reservation_id in {5, 6, 7}: diagnostic_items.append({ 'diagnostic_uid': diagnostic_uid, 'phase': current_phase, 'kind': 'provider_http', 'category': 'rate_limit', 'code': 'provider.rate_limit', 'retryable': True, 'summary': 'provider returned a limit', 'occurred_at': '2026-09-20T00:03:00Z', 'received_at': '2026-09-20T00:03:01Z', 'envelope_sha256': '6' * 64, 'diagnostic': envelope, 'canonical_envelope_json': canonical, }) timeline = [{ 'timestamp': '2026-09-20T00:01:00Z', 'kind': 'assignment', 'label': 'issued', }] if not legacy: timeline.append({ 'timestamp': '2026-09-20T00:02:00Z', 'kind': 'phase', 'label': current_phase, 'sequence': 1, 'received_at': '2026-09-20T00:02:01Z', }) if reservation_id == 5: timeline.extend([ {'timestamp': '2026-09-20T00:04:00Z', 'kind': 'transport', 'label': 'bundle received'}, {'timestamp': '2026-09-20T00:05:00Z', 'kind': 'ingestion', 'label': 'bundle ingested'}, {'timestamp': '2026-09-20T00:05:45Z', 'kind': 'settlement', 'label': 'queue settled'}, {'timestamp': '2026-09-20T00:06:00Z', 'kind': 'projection', 'label': 'projection completed'}, ]) elif reservation_id == 6: timeline.append({ 'timestamp': resolved_at, 'kind': 'receipt', 'label': 'prebundle_report', }) elif reservation_id == 7: timeline.append({ 'timestamp': resolved_at, 'kind': 'receipt', 'label': 'expired', }) return { 'schema': 1, 'assignment': { 'reservation_id': int(reservation_id), 'queue_id': 5, 'source': 'gitlab', 'target': 'example/project@commit', 'user': 'fixture-user', 'worker': 'fixture-device', 'active_assignment_cap': 2, 'assignment_outcome': assignment_outcome, 'scan_outcome': scan_outcome, 'issued_at': '2026-09-20T00:01:00Z', 'resolved_at': resolved_at, 'assignment_code': None, 'assignment_detail': None, }, 'deadlines': { 'assignment_deadline_at': '2026-09-21T00:01:00Z', 'scan_deadline_at': '2026-09-20T00:11:00Z', 'upload_timeout_seconds': None if legacy else 1800, 'upload_timeout_availability': ( 'legacy/unavailable' if legacy else 'persisted at assignment issuance' ), }, 'package': { 'protocol_version': 1 if legacy else 2, 'bundle_format_version': 2, 'platform_tag': 'linux-x86_64', 'code_manifest_sha256': '1' * 64, 'detector_policy_sha256': '2' * 64, 'effective_config_sha256': '3' * 64, }, 'transport': { 'resolution': None, 'receipt_id': ( 'receipt-fixture' if resolved_at else None ), 'bundle_state': 'acknowledged' if reservation_id in {5, 8} else None, 'bundle_ready_at': '2026-09-20T00:04:00Z' if reservation_id in {5, 8} else None, 'bundle_committed_at': '2026-09-20T00:05:00Z' if reservation_id in {5, 8} else None, 'bundle_acknowledged_at': '2026-09-20T00:05:30Z' if reservation_id in {5, 8} else None, 'queue_status': 'done' if resolved_at else 'in_progress', 'queue_settled_at': '2026-09-20T00:05:45Z' if reservation_id in {5, 8} else None, 'projection_status': 'completed' if reservation_id in {5, 8} else None, 'projection_completed_at': '2026-09-20T00:06:00Z' if reservation_id in {5, 8} else None, }, 'scan': { 'available': reservation_id in {5, 8}, 'target_scan_id': reservation_id * 10 if reservation_id in {5, 8} else None, 'status': 'error' if reservation_id in {5, 8} else None, 'started_at': '2026-09-20T00:02:00Z', 'ended_at': '2026-09-20T00:04:00Z', 'duration_seconds': 120.0, 'findings_count': 0, 'verified_findings_count': 0, 'error_count': 1, 'skipped_reason': None, 'first_error_summary': 'legacy summary' if legacy else None, 'warning_class': 'detector_timeout' if reservation_id == 5 else None, 'warning_summary': ( 'bounded detector output timed out' if reservation_id == 5 else None ), }, 'timeline': timeline, 'durations': [{ 'phase': current_phase, 'duration_seconds': 120.0, 'outcome': scan_outcome, 'complete': resolved_at is not None, 'ended_at': resolved_at or '2026-09-20T00:04:00Z', 'authority': ( 'assignment_resolution' if resolved_at else 'database_current_time' ), }], 'progress': { 'available': not legacy, 'availability': ( 'legacy/unavailable' if legacy else 'current' ), 'truncated': False, 'total_event_count': 0 if legacy else 1, 'omitted_older_event_count': 0, 'current_phase': None if legacy else current_phase, 'phase_started_at': None if legacy else '2026-09-20T00:02:00Z', 'phase_age_seconds': None if legacy else 120.0, 'last_progress_at': None if legacy else '2026-09-20T00:03:55Z', 'last_progress_age_seconds': None if legacy else 5.0, 'age_authority': ( None if legacy else 'assignment_resolution' if resolved_at else 'database_current_time' ), 'events': [], }, 'diagnostics': { 'availability': 'legacy/unavailable' if legacy else 'current', 'projection_version': None if legacy or reservation_id == 4 else 1, 'declared_count': None if legacy else len(diagnostic_items), 'truncated': False, 'items': diagnostic_items, }, 'legacy_evidence': { 'available': legacy, 'explicitly_not_an_envelope': True, 'truncated': False, 'errors': [{ 'id': 1, 'category': 'scanner', 'summary': 'legacy summary', 'raw_error': 'exact legacy raw error', 'created_at': '2026-09-20T00:04:00Z', }] if legacy else [], 'first_error_summary': 'legacy summary' if legacy else None, }, } def admin_worker_diagnostic_envelope(self, reservation_id, diagnostic_uid): self.calls.append(('diagnostic_envelope', reservation_id, diagnostic_uid)) if int(reservation_id) != 5 or diagnostic_uid != 'b' * 64: return None envelope = self._diagnostic_fixture(5, 'b' * 64) canonical = json.dumps(envelope, ensure_ascii=True, sort_keys=True, separators=(',', ':')) return {'canonical_json': canonical, 'sha256': '6' * 64} def admin_target_queue_health(self, sources): self.calls.append(('queue_counts', tuple(sources))) if self.fail_queue: if self.fail_queue == 'malformed': return { 'counts': {}, 'degraded': False, 'stale': False, 'truncated_statuses': 1, } if self.fail_queue == 'degraded': return { 'counts': {}, 'degraded': True, 'stale': True, 'reason': 'bounded_count_query_failed', 'retry_after_sec': 300, 'truncated_statuses': [], } raise RuntimeError('sensitive queue failure detail') return { 'counts': {'pending': 3, 'done': 100}, 'degraded': True, 'stale': False, 'reason': 'bounded_status_sample', 'sample_limit_per_status': 100, 'timeout_ms': 5000, 'sampled_rows': 103, 'retry_after_sec': 0, 'truncated_statuses': ['done'], } def runtime_drain_progress(self): self.calls.append(('drain_progress',)) return { 'revision': 7, 'discovery_paused': False, 'dispatch_paused': True, 'drain_state': 'draining', 'effective_discovery_paused': True, 'effective_dispatch_paused': True, 'actor': 'fixture.operator', 'operation_id': '00000000-0000-4000-8000-000000000001', 'created_at': '2026-09-20T00:00:00+00:00', 'updated_at': '2026-09-20T00:01:00+00:00', 'live_remote_assignments': 2, 'precommit_result_bundles': 1, 'blocker_count': 3, } def set_runtime_dispatch_paused( self, paused, *, expected_revision, actor, operation_id, ): if self.control_revision_conflict: raise RuntimeControlRevisionConflictError( expected_revision, {'revision': expected_revision + 1}, ) if self.control_identity_conflict: raise RuntimeOperationIdentityConflictError('fixture conflict') self.calls.append(( 'dispatch_paused', paused, expected_revision, actor, operation_id, )) return {'after': {'revision': expected_revision + 1}} def start_runtime_drain(self, *, expected_revision, actor, operation_id): if self.control_revision_conflict: raise RuntimeControlRevisionConflictError( expected_revision, {'revision': expected_revision + 1}, ) if self.control_identity_conflict: raise RuntimeOperationIdentityConflictError('fixture conflict') self.calls.append(('drain_start', expected_revision, actor, operation_id)) return {'after': {'revision': expected_revision + 1}} def cancel_runtime_drain(self, *, expected_revision, actor, operation_id): if self.control_revision_conflict: raise RuntimeControlRevisionConflictError( expected_revision, {'revision': expected_revision + 1}, ) if self.control_identity_conflict: raise RuntimeOperationIdentityConflictError('fixture conflict') self.calls.append(('drain_cancel', expected_revision, actor, operation_id)) return {'after': {'revision': expected_revision + 1}} def recent_runtime_operations( self, limit, *, before_updated_at=None, before_operation_id=None, ): self.calls.append(( 'recent_operations', limit, before_updated_at, before_operation_id, )) recorded = ( list(type(self).apply_operations.values()) + list(type(self).document_operations.values()) + list(type(self).managed_file_operations.values()) ) if recorded: return [dict(operation) for operation in recorded[:limit]] return [{ 'operation_id': '00000000-0000-4000-8000-000000000002', 'actor': 'fixture.operator', 'action': 'restart', 'target_ref': 'runtime', 'status': 'succeeded', 'safe_category': None, 'safe_detail': None, 'requested_at': '2026-09-20T00:00:00+00:00', 'completed_at': '2026-09-20T00:00:05+00:00', 'updated_at': '2026-09-20T00:00:05+00:00', }] def set_runtime_discovery_paused( self, paused, *, expected_revision, actor, operation_id, ): if self.control_revision_conflict: raise RuntimeControlRevisionConflictError( expected_revision, {'revision': expected_revision + 1}, ) self.calls.append(( 'discovery_paused', paused, expected_revision, actor, operation_id, )) return {'after': {'revision': expected_revision + 1}} def create_runtime_source_operation( self, *, operation_id, actor, source_id, source_action, interval_seconds=None, mode=None, restart_enabled=None, restart_delay_seconds=None, ): if self.source_identity_conflict: raise RuntimeOperationIdentityConflictError('fixture conflict') self.calls.append(( 'source_operation_create', operation_id, actor, source_id, source_action, interval_seconds, mode, restart_enabled, restart_delay_seconds, )) existing = self.source_operations.get(operation_id) if existing is not None: return {**existing, 'replayed': True} operation = { 'operation_id': operation_id, 'status': 'running', 'replayed': False, 'resulting_identity': None, } self.source_operations[operation_id] = operation return dict(operation) def complete_runtime_source_operation( self, operation_id, *, succeeded, outcome=None, ): self.calls.append(( 'source_operation_complete', operation_id, succeeded, outcome, )) if type(self).completion_failures: type(self).completion_failures -= 1 raise RuntimeError('fixture completion failure') operation = { **self.source_operations[operation_id], 'status': 'succeeded' if succeeded else 'failed', 'resulting_identity': ( {'outcome': outcome} if succeeded else None ), } self.source_operations[operation_id] = operation return dict(operation) def create_runtime_worker_admin_operation( self, *, operation_id, actor, action, target_ref, request_sha256, ): self.calls.append(( 'worker_admin_create', operation_id, actor, action, target_ref, request_sha256, )) existing = self.worker_admin_operations.get(operation_id) if existing is not None: return {**existing, 'replayed': True} operation = { 'operation_id': operation_id, 'status': 'running', 'replayed': False, 'resulting_identity': None, } self.worker_admin_operations[operation_id] = operation return dict(operation) def complete_runtime_worker_admin_operation( self, operation_id, *, succeeded, affected_count=None, ): self.calls.append(( 'worker_admin_complete', operation_id, succeeded, affected_count, )) operation = { **self.worker_admin_operations[operation_id], 'status': 'succeeded' if succeeded else 'failed', 'resulting_identity': ( {'outcome': 'completed', 'affected_count': affected_count} if succeeded else None ), } self.worker_admin_operations[operation_id] = operation return dict(operation) def create_runtime_document_operation(self, **kwargs): self.calls.append(('document_operation_create', kwargs)) operation_id = kwargs['operation_id'] document = kwargs['action'].split('.')[1] expected_identity = { 'document': document, 'active_config_sha256': kwargs['active_config_sha256'], 'active_secrets_sha256': kwargs['active_secrets_sha256'], 'candidate_config_sha256': kwargs['candidate_config_sha256'], 'candidate_secrets_sha256': kwargs['candidate_secrets_sha256'], 'candidate_after_sha256': kwargs['candidate_after_sha256'], 'candidate_before_bytes': kwargs['candidate_before_bytes'], 'candidate_after_bytes': kwargs['candidate_after_bytes'], 'candidate_before_present': kwargs['candidate_before_present'], } existing = self.document_operations.get(operation_id) if existing is not None: if not ( existing['actor'] == kwargs['actor'] and existing['action'] == kwargs['action'] and existing['target_ref'] == document and existing['expected_identity'] == expected_identity ): raise RuntimeOperationIdentityConflictError('fixture identity conflict') return {**existing, 'replayed': True} if operation_id in self.apply_operations: raise RuntimeOperationIdentityConflictError('fixture identity conflict') operation = { 'operation_id': operation_id, 'status': 'running', 'replayed': False, 'actor': kwargs['actor'], 'action': kwargs['action'], 'target_kind': 'runtime-document', 'target_ref': document, 'expected_identity': expected_identity, } self.document_operations[operation_id] = operation return dict(operation) def complete_runtime_document_operation( self, operation_id, *, succeeded, candidate_sha256=None, written=None, ): self.calls.append(( 'document_operation_complete', operation_id, succeeded, candidate_sha256, written, )) if type(self).document_completion_failures: type(self).document_completion_failures -= 1 raise RuntimeError('fixture document completion failure') operation = { **self.document_operations[operation_id], 'status': 'succeeded' if succeeded else 'failed', } self.document_operations[operation_id] = operation return dict(operation) def create_runtime_operation(self, **kwargs): self.calls.append(('runtime_operation_create', kwargs)) operation_id = kwargs['operation_id'] existing = self.apply_operations.get(operation_id) if existing is not None: if not ( existing['actor'] == kwargs['actor'] and existing['action'] == kwargs['action'] and existing['expected_identity'] == kwargs['expected_identity'] ): raise RuntimeOperationIdentityConflictError('fixture identity conflict') return {**existing, 'replayed': True} if operation_id in self.document_operations: raise RuntimeOperationIdentityConflictError('fixture identity conflict') operation = { 'operation_id': operation_id, 'status': 'requested', 'replayed': False, 'actor': kwargs['actor'], 'action': kwargs['action'], 'expected_identity': dict(kwargs['expected_identity']), } self.apply_operations[operation_id] = operation return dict(operation) def create_runtime_managed_file_operation(self, **kwargs): self.calls.append(('managed_file_operation_create', dict(kwargs))) operation_id = kwargs['operation_id'] expected_identity = { key: kwargs[key] for key in ( 'root_id', 'relative_path', 'expected_sha256', 'proposed_sha256', 'proposed_byte_count', ) } existing = self.managed_file_operations.get(operation_id) if existing is not None: if not ( existing['actor'] == kwargs['actor'] and existing['action'] == kwargs['action'] and existing['target_ref'] == kwargs['root_id'] and existing['expected_identity'] == expected_identity ): raise RuntimeOperationIdentityConflictError('fixture identity conflict') return {**existing, 'replayed': True} operation = { 'operation_id': operation_id, 'status': 'running', 'replayed': False, 'actor': kwargs['actor'], 'action': kwargs['action'], 'target_kind': 'managed-file', 'target_ref': kwargs['root_id'], 'expected_identity': expected_identity, 'resulting_identity': None, } self.managed_file_operations[operation_id] = operation return dict(operation) def complete_runtime_managed_file_operation(self, operation_id, **kwargs): self.calls.append(( 'managed_file_operation_complete', operation_id, dict(kwargs), )) if type(self).managed_file_completion_failures: type(self).managed_file_completion_failures -= 1 raise RuntimeError('fixture managed file completion failure') operation = self.managed_file_operations[operation_id] resulting = None if kwargs['succeeded']: expected = operation['expected_identity'] resulting = { 'root_id': expected['root_id'], 'relative_path': expected['relative_path'], 'outcome': 'completed', 'before_sha256': kwargs.get('before_sha256'), 'before_byte_count': kwargs.get('before_byte_count'), 'after_sha256': kwargs.get('after_sha256'), 'after_byte_count': kwargs.get('after_byte_count'), 'written': kwargs.get('written'), } operation = { **operation, 'status': 'succeeded' if kwargs['succeeded'] else 'failed', 'resulting_identity': resulting, } self.managed_file_operations[operation_id] = operation return dict(operation) def runtime_operation(self, operation_id): self.calls.append(('runtime_operation', operation_id)) operation = ( self.document_operations.get(operation_id) or self.apply_operations.get(operation_id) or self.managed_file_operations.get(operation_id) ) return None if operation is None else dict(operation) def runtime_audit_events(self, *, before_event_id=None, limit=50): self.calls.append(('runtime_audit_events', before_event_id, limit)) return { 'events': [dict(event) for event in type(self).audit_events], 'next_before_event_id': type(self).audit_next_before_event_id, } def create_remote_worker_user(self, user_key, cap): self.calls.append(('create_user', user_key, cap)) return {'user_key': user_key} def set_remote_worker_user_cap(self, user_key, cap): self.calls.append(('set_cap', user_key, cap)) return True def set_remote_worker_user_disabled(self, user_key, disabled): self.calls.append(('set_disabled', user_key, disabled)) return True def issue_remote_worker_device(self, user_key, device_key, token_sha256, *, rotate=False): self.calls.append(('issue_device', user_key, device_key, token_sha256, rotate)) return {'user_key': user_key, 'device_key': device_key} def set_remote_worker_device_revoked(self, device_key, revoked=True): self.calls.append(('set_revoked', device_key, revoked)) return True def admin_requeue_deferred_targets(self, queue_ids, *, max_items): self.calls.append(('requeue', list(queue_ids), max_items)) return len(queue_ids) def admin_discard_queued_source(self, source): self.calls.append(('discard_source_queue', source)) return 17 class AdminAPITests(unittest.TestCase): def setUp(self): RecordingDB.calls = [] RecordingDB.fail_queue = False RecordingDB.source_operations = {} RecordingDB.worker_admin_operations = {} RecordingDB.document_operations = {} RecordingDB.apply_operations = {} RecordingDB.managed_file_operations = {} RecordingDB.audit_events = [] RecordingDB.audit_next_before_event_id = None RecordingDB.source_identity_conflict = False RecordingDB.completion_failures = 0 RecordingDB.control_identity_conflict = False RecordingDB.control_revision_conflict = False RecordingDB.document_completion_failures = 0 RecordingDB.managed_file_completion_failures = 0 self.temp_dir = tempfile.TemporaryDirectory() root = ensure_private_directory( os.path.join(self.temp_dir.name, 'bundles'), reject_reparse=True, ) self.worker = FakeWorkerService(root) self.runtime_provider = mock.Mock(return_value={ 'snapshot_schema': 2, 'runtime': { 'pid': 123, 'phase': 'ACTIVE', 'manages_postgres': True, 'start_gate_open': True, 'shutdown_requested': False, 'runtime_failed': False, }, 'postgres': { 'state': 'READY', 'ready': True, 'failures': 0, 'safe_error_category': '', }, 'dashboard': { 'status': 'running', 'desired_state': 'running', 'healthy': True, 'pid': 456, 'safe_error_category': '', }, 'sources': [ { 'id': f'discovery-producer:{source}', 'source': source, 'role': 'discovery-producer', 'lifecycle_state': 'waiting', 'process_state': 'stopped', 'desired_state': 'running', 'safe_error_category': '', 'interval_seconds': 3600, 'restart_enabled': True, 'restart_count': 2, 'restart_delay_seconds': 5, 'restart_streak': 0, 'mode': 'repeat', 'pid': None, 'allowed_actions': [ 'start', 'stop', 'restart', 'pause', 'resume', 'set-interval', ], 'last_cycle_result': { 'status': 'completed', 'fetched_count': 4, 'queued_new_count': 2, 'queued_updated_count': 1, }, 'last_successful_discovery_at': '2026-09-20T00:00:00+00:00', 'next_scheduled_run_at': '2026-09-20T01:00:00+00:00', } for source in ('gitlab', 'dockerhub', 'huggingface') ] + [ { 'id': source, 'source': source, 'role': source, 'lifecycle_state': 'running', 'process_state': 'running', 'desired_state': 'running', 'safe_error_category': '', 'mode': 'singleton', 'pid': 789, 'interval_seconds': 0, 'restart_enabled': True, 'restart_delay_seconds': 5, 'restart_count': 0, 'restart_streak': 0, 'allowed_actions': [ 'start', 'stop', 'restart', 'pause', 'resume', 'set-restart', 'set-restart-delay', ], } for source in ('result-ingester', 'jsonl-projector', 'janitor', 'worker-api') ] + [{ 'id': 'keychecks', 'source': 'keychecks', 'role': 'keycheck', 'lifecycle_state': 'waiting', 'process_state': 'stopped', 'desired_state': 'running', 'safe_error_category': '', 'pid': None, 'mode': 'repeat', 'interval_seconds': 3600, 'restart_enabled': True, 'restart_delay_seconds': 5, 'restart_count': 0, 'restart_streak': 0, 'allowed_actions': [ 'start', 'stop', 'restart', 'pause', 'resume', 'once', 'set-mode', 'set-interval', 'set-restart', 'set-restart-delay', ], }, { 'id': 'docker-shadow', 'source': 'docker-shadow', 'role': 'docker-shadow', 'lifecycle_state': 'stopped', 'process_state': 'stopped', 'desired_state': 'stopped', 'safe_error_category': '', 'pid': None, 'mode': 'manual', 'interval_seconds': 0, 'restart_enabled': False, 'restart_delay_seconds': 0, 'restart_count': 0, 'restart_streak': 0, 'allowed_actions': ['start', 'stop'], }], 'pipeline': { 'ingester_ready': True, 'projector_ready': True, 'cutover_ready': True, 'ingester_state': 'ready', 'projector_state': 'ready', 'bundle_items': 1, 'bundle_bytes': 4096, 'projection_items': 0, 'projection_bytes': 0, 'keycheck_items': 0, 'keycheck_bytes': 0, 'quarantine_items': 0, 'quarantine_bytes': 0, }, 'scan_workers': { 'active': 0, 'limit': 2, 'base_active': 0, 'base_limit': 2, 'bonus_active': 0, 'bonus_limit': 0, 'trufflehog': 0, 'sources': {}, }, }) self.source_action_provider = mock.Mock(return_value={ 'source_action': 'restart', 'outcome': 'completed', 'source': {}, }) self.dashboard_action_provider = mock.Mock(return_value={ 'dashboard_action': 'restart', 'outcome': 'completed', 'dashboard': {}, }) self.source_log_provider = mock.Mock(return_value={ 'source_id': 'result-ingester', 'line_count': 2, 'lines': ['line ', 'line two'], 'response_truncated': False, }) self.package_compatibility_provider = mock.Mock(return_value={ 'profiles': [{ 'profile_name': 'linux-x86_64', 'protocol_version': 2, 'bundle_format_version': 2, 'platform_tag': 'linux-x86_64', 'code_manifest_sha256': '1' * 64, 'detector_policy_sha256': '2' * 64, 'sources': ['dockerhub', 'gitlab', 'huggingface'], 'capabilities': [ {'source': 'gitlab', 'platform': 'gitlab', 'planning_kind': 'exact_git_v1'}, {'source': 'dockerhub', 'platform': 'docker', 'planning_kind': 'docker_direct_v1'}, {'source': 'huggingface', 'platform': 'huggingface', 'planning_kind': 'huggingface_space_v1'}, ], }], 'required_capabilities': [ {'source': 'gitlab', 'platform': 'gitlab', 'planning_kind': 'exact_git_v1'}, {'source': 'dockerhub', 'platform': 'docker', 'planning_kind': 'docker_direct_v1'}, {'source': 'huggingface', 'platform': 'huggingface', 'planning_kind': 'huggingface_space_v1'}, ], }) self.document_state = SimpleNamespace( active_config=SimpleNamespace(sha256='a' * 64, byte_count=10, present=True), active_secrets=SimpleNamespace(sha256='b' * 64, byte_count=10, present=True), candidate_config=SimpleNamespace(sha256='c' * 64, byte_count=10, present=True), candidate_secrets=SimpleNamespace(sha256='d' * 64, byte_count=10, present=True), ) self.runtime_config = { 'global': {'project_dir': '/opt/truf/app'}, 'sources': { 'gitlab': {'timeout': 600}, 'dockerhub': {'timeout': 900}, 'huggingface': {'timeout': 300}, }, 'supervisor': {'worker_api': { 'sources': ['gitlab', 'dockerhub', 'huggingface'], 'assignment_ttl_seconds': 86400, 'assignment_ttl_seconds_by_source': {'dockerhub': 90000}, 'bundle_body_timeout_seconds': 1800, }}, } self.runtime_config_provider = mock.Mock(return_value=SimpleNamespace( config=self.runtime_config, )) self.document_loader = mock.Mock(side_effect=lambda _path, document: SimpleNamespace( document=document, source='candidate', text=( 'global:\n project_dir: /opt/truf/app\n' 'sources:\n gitlab:\n timeout: 600\n' ' dockerhub:\n timeout: 900\n' ' huggingface:\n timeout: 300\n' 'supervisor:\n worker_api:\n' ' sources: [gitlab, dockerhub, huggingface]\n' ' assignment_ttl_seconds: 86400\n' ' assignment_ttl_seconds_by_source:\n dockerhub: 90000\n' ' bundle_body_timeout_seconds: 1800\n' if document == 'config' else 'auth_pools:\n fixture:\n - name: fixture\n token: editor-secret\n' ), state=self.document_state, selected=( self.document_state.candidate_config if document == 'config' else self.document_state.candidate_secrets ), )) self.candidate_preview_provider = mock.Mock(side_effect=lambda _path, document, payload: SimpleNamespace( document=document, state=self.document_state, proposed=SimpleNamespace( sha256=hashlib.sha256(payload).hexdigest(), byte_count=len(payload), present=True, ), diff=( SimpleNamespace(entries=(), truncated=False, format_only_changed=False) if document == 'config' else SimpleNamespace(**{ name: value for name, value in ( ('document_changed', True), ('semantic_changed', True), ('pools_before', 1), ('pools_after', 1), ('pools_added', 0), ('pools_removed', 0), ('pools_changed', 1), ('entries_before', 1), ('entries_after', 1), ('entries_added', 0), ('entries_removed', 0), ('pools_reordered', 0), ('usernames_added', 0), ('usernames_removed', 0), ('usernames_changed', 0), ('tokens_changed', 1), ) }) ), )) def save_candidate(_path, document, payload, **_kwargs): proposed = SimpleNamespace( sha256=hashlib.sha256(payload).hexdigest(), byte_count=len(payload), present=True, ) if document == 'config': self.document_state.candidate_config = proposed else: self.document_state.candidate_secrets = proposed return SimpleNamespace( document=document, proposed=proposed, written=True, ) self.candidate_save_provider = mock.Mock(side_effect=save_candidate) self.candidate_verify_provider = mock.Mock(return_value=SimpleNamespace(action='apply-config')) self.runtime_apply_provider = mock.Mock() self.managed_file_traversal = FakeManagedFileTraversal() managed_limits = ManagedFileLimits(1024, 255, 16, 100, 65536, 65536) self.managed_file_registry = ManagedFileRootRegistry(( ManagedFileRoot( 'exports', '/data/managed-files/private-host-path', ManagedFilePermissions(True, True, True, True), managed_limits, ), ManagedFileRoot( 'readonly', '/data/managed-files/private-readonly-path', ManagedFilePermissions(True, True, False, False), managed_limits, ), )) self.admin = AdminService( 'postgresql://fixture', ORIGIN, MARKER, db_factory=RecordingDB, supervisor_metadata={ 'instance_id': 'fixture', 'token': 'supervisor-token-must-not-render', }, runtime_snapshot_provider=self.runtime_provider, source_action_provider=self.source_action_provider, dashboard_action_provider=self.dashboard_action_provider, source_log_provider=self.source_log_provider, package_compatibility_provider=self.package_compatibility_provider, runtime_config_path='/data/config/runtime.yaml', runtime_config_provider=self.runtime_config_provider, document_loader=self.document_loader, candidate_preview_provider=self.candidate_preview_provider, candidate_save_provider=self.candidate_save_provider, candidate_verify_provider=self.candidate_verify_provider, runtime_apply_provider=self.runtime_apply_provider, managed_file_roots=self.managed_file_registry, ) self.traversal_patcher = mock.patch( 'worker_api.ManagedFileTraversal', return_value=self.managed_file_traversal, ) self.traversal_patcher.start() def tearDown(self): self.traversal_patcher.stop() self.temp_dir.cleanup() def _client(self, enabled=True): return TestClient(create_worker_app( self.worker, reaper_interval_seconds=3600, admin_service=self.admin if enabled else None, )) def _headers(self, **extra): return { EDGE_MARKER_HEADER: MARKER, OPERATOR_HEADER: OPERATOR, 'Origin': ORIGIN, **extra, } def _post(self, client, path, fields): fields = dict(fields) if path.startswith(('/users/', '/devices/', '/queue/')): fields.setdefault('operation_id', str(uuid.uuid4())) return client.post( ADMIN_PREFIX + path, data={'csrf_token': self.admin.csrf_token, **fields}, headers=self._headers(), ) def test_managed_file_registry_pages_and_download_are_logical_and_bounded(self): self.assertIs(self.admin.managed_file_roots, self.managed_file_registry) registry = ManagedFileRootRegistry((object(),)) service = AdminService( 'postgresql://fixture', ORIGIN, MARKER, db_factory=RecordingDB, managed_file_roots=registry, ) self.assertIs(service.managed_file_roots, registry) with self.assertRaisesRegex(ValueError, 'registry is invalid'): AdminService( 'postgresql://fixture', ORIGIN, MARKER, db_factory=RecordingDB, managed_file_roots={}, ) with self._client() as client: index = client.get(ADMIN_PREFIX + '/files', headers=self._headers()) listing = client.get( ADMIN_PREFIX + '/files?root_id=exports', headers=self._headers(), ) nested = client.get( ADMIN_PREFIX + '/files?root_id=exports&relative_path=nested', headers=self._headers(), ) download = client.get( ADMIN_PREFIX + '/files/download?root_id=exports&relative_path=nested%2Freport.bin', headers=self._headers(), ) for response in (index, listing, nested, download): self.assertEqual(response.status_code, 200, response.text) self.assertSecurityHeaders(response) self.assertNotIn('/data/managed-files', response.text) self.assertIn('>Files<', index.text) self.assertIn('exports', index.text) self.assertIn('nested/', listing.text) self.assertIn('report.bin', nested.text) self.assertEqual(download.content, b'\x00\xffreport') self.assertEqual(download.headers['content-type'], 'application/octet-stream') self.assertEqual(download.headers['content-length'], '8') self.assertEqual( download.headers['etag'], '"' + hashlib.sha256(b'\x00\xffreport').hexdigest() + '"', ) self.assertIn("filename*=UTF-8''report.bin", download.headers['content-disposition']) def test_managed_file_snapshot_is_streamed_and_closed(self): self.managed_file_traversal.stream_downloads = True with self._client() as client: download = client.get( ADMIN_PREFIX + '/files/download?root_id=exports&relative_path=nested%2Freport.bin', headers=self._headers(), ) self.assertEqual(download.status_code, 200, download.text) self.assertEqual(download.content, b'\x00\xffreport') self.assertEqual(download.headers['content-length'], '8') self.assertEqual(self.managed_file_traversal.last_snapshot.closed, 1) def test_cancelled_managed_file_download_closes_completed_snapshot(self): snapshot = FakeManagedFileSnapshot(b'bounded snapshot') started = threading.Event() release = threading.Event() class BlockingService: @staticmethod def download_managed_file(*_args): started.set() if not release.wait(5): raise RuntimeError('test download timed out') return ManagedFileDownload( ManagedFileIdentity('0' * 64, len(snapshot.content)), snapshot=snapshot, ) async def scenario(): task = asyncio.create_task(_download_managed_file( BlockingService(), object(), 'runtime-results', 'scan_results.jsonl', )) self.assertTrue(await asyncio.to_thread(started.wait, 2)) task.cancel() release.set() with self.assertRaises(asyncio.CancelledError): await task for _ in range(20): if snapshot.closed: break await asyncio.sleep(0.01) asyncio.run(scenario()) self.assertEqual(snapshot.closed, 1) def test_streaming_response_construction_failure_closes_snapshot(self): snapshot = FakeManagedFileSnapshot(b'bounded snapshot') download = ManagedFileDownload( ManagedFileIdentity('0' * 64, len(snapshot.content)), snapshot=snapshot, ) with mock.patch( 'admin_api._ManagedFileStreamingResponse', side_effect=RuntimeError('failed')): with self.assertRaisesRegex(RuntimeError, 'failed'): _managed_file_download_response( download, 'scan_results.jsonl', ) self.assertEqual(snapshot.closed, 1) def test_streaming_send_failure_closes_snapshot(self): snapshot = FakeManagedFileSnapshot(b'bounded snapshot') response = _ManagedFileStreamingResponse( snapshot, snapshot.chunks(), media_type='application/octet-stream', ) async def scenario(): async def receive(): return {'type': 'http.disconnect'} async def send(_message): raise RuntimeError('send failed') with self.assertRaisesRegex(RuntimeError, 'send failed'): await response( { 'type': 'http', 'method': 'GET', 'path': '/', 'headers': [], 'asgi': {'version': '3.0', 'spec_version': '2.4'}, }, receive, send, ) asyncio.run(scenario()) self.assertEqual(snapshot.closed, 1) def test_managed_file_mutations_are_typed_audited_and_content_free(self): created_payload = b'created-binary-\x00-credential-sentinel' replacement_payload = b'replaced-binary-\xff-credential-sentinel' with self._client() as client: create_id = str(uuid.uuid4()) created = self._post(client, '/files/create', { 'operation_id': create_id, 'root_id': 'exports', 'relative_path': 'new.bin', 'content_base64': __import__('base64').urlsafe_b64encode( created_payload, ).decode('ascii'), }) created_hash = hashlib.sha256(created_payload).hexdigest() replace_id = str(uuid.uuid4()) replaced = self._post(client, '/files/replace', { 'operation_id': replace_id, 'root_id': 'exports', 'relative_path': 'new.bin', 'expected_sha256': created_hash, 'content_base64': __import__('base64').urlsafe_b64encode( replacement_payload, ).decode('ascii'), }) replacement_hash = hashlib.sha256(replacement_payload).hexdigest() delete_id = str(uuid.uuid4()) deleted = self._post(client, '/files/delete', { 'operation_id': delete_id, 'root_id': 'exports', 'relative_path': 'new.bin', 'expected_sha256': replacement_hash, }) for response in (created, replaced, deleted): self.assertEqual(response.status_code, 200, response.text) self.assertSecurityHeaders(response) self.assertNotIn('new.bin', self.managed_file_traversal.files) operations = [ call for call in RecordingDB.calls if call[0] == 'managed_file_operation_create' ] self.assertEqual( [call[1]['action'] for call in operations], ['files.create', 'files.replace', 'files.delete'], ) self.assertEqual([call[1]['actor'] for call in operations], [OPERATOR] * 3) self.assertEqual( [call[1]['operation_id'] for call in operations], [create_id, replace_id, delete_id], ) recorded = repr(RecordingDB.calls) self.assertNotIn('credential-sentinel', recorded) self.assertNotIn(__import__('base64').urlsafe_b64encode(created_payload).decode(), recorded) def test_managed_file_create_replay_recovers_after_completion_loss(self): payload = b'durable replay payload' encoded = __import__('base64').urlsafe_b64encode(payload).decode('ascii') operation_id = str(uuid.uuid4()) fields = { 'csrf_token': self.admin.csrf_token, 'operation_id': operation_id, 'root_id': 'exports', 'relative_path': 'replayed.bin', 'content_base64': encoded, } RecordingDB.managed_file_completion_failures = 3 with self._client() as client: first = client.post( ADMIN_PREFIX + '/files/create', data=fields, headers=self._headers(), follow_redirects=False, ) second = client.post( ADMIN_PREFIX + '/files/create', data=fields, headers=self._headers(), follow_redirects=False, ) terminal = client.post( ADMIN_PREFIX + '/files/create', data=fields, headers=self._headers(), follow_redirects=False, ) self.assertEqual(first.status_code, 503) self.assertEqual((second.status_code, terminal.status_code), (303, 303)) self.assertSecurityHeaders(first) self.assertEqual(self.managed_file_traversal.files['replayed.bin'], payload) mutations = [ call for call in self.managed_file_traversal.calls if call[0] == 'create_replace' and call[2] == 'replayed.bin' ] self.assertEqual(len(mutations), 1) self.assertEqual( RecordingDB.managed_file_operations[operation_id]['status'], 'succeeded', ) replay = self.admin.create_managed_file( None, 'exports', 'replayed.bin', payload, OPERATOR, operation_id, ) self.assertEqual(replay['status'], 'succeeded') def test_managed_file_acceptance_precedes_one_concurrent_physical_mutation(self): payload = b'concurrent managed payload' operation_id = str(uuid.uuid4()) original = self.managed_file_traversal.create_replace_file other_admin = AdminService( 'postgresql://fixture', ORIGIN, MARKER, db_factory=RecordingDB, managed_file_roots=self.managed_file_registry, ) def accepted_provider(*args, **kwargs): self.assertEqual( RecordingDB.managed_file_operations[operation_id]['status'], 'running', ) return original(*args, **kwargs) barrier = threading.Barrier(3) results = [] def submit(service): barrier.wait() results.append(service.create_managed_file( self.managed_file_traversal, 'exports', 'concurrent.bin', payload, OPERATOR, operation_id, )) with mock.patch.object( self.managed_file_traversal, 'create_replace_file', side_effect=accepted_provider): threads = [ threading.Thread(target=submit, args=(service,)) for service in (self.admin, other_admin) ] for thread in threads: thread.start() barrier.wait() for thread in threads: thread.join(timeout=5) self.assertTrue(all(not thread.is_alive() for thread in threads)) self.assertEqual([result['status'] for result in results], [ 'succeeded', 'succeeded', ]) mutations = [ call for call in self.managed_file_traversal.calls if call[0] == 'create_replace' and call[2] == 'concurrent.bin' ] self.assertEqual(len(mutations), 1) self.assertEqual( RecordingDB.managed_file_operations[operation_id]['status'], 'succeeded', ) def test_fresh_managed_file_cas_loser_is_not_promoted_to_success(self): payload = b'external matching payload' operation_id = str(uuid.uuid4()) def lose_cas(_root_id, relative_path, value, *, expected_sha256): self.managed_file_traversal.files[relative_path] = value raise ManagedFileAccessError('hash_conflict') with mock.patch.object( self.managed_file_traversal, 'create_replace_file', side_effect=lose_cas): with self.assertRaises(AdminAPIError) as raised: self.admin.create_managed_file( self.managed_file_traversal, 'exports', 'external.bin', payload, OPERATOR, operation_id, ) self.assertEqual(raised.exception.status_code, 409) self.assertEqual( RecordingDB.managed_file_operations[operation_id]['status'], 'failed', ) def assertSecurityHeaders(self, response): for name, expected in SECURITY_HEADERS.items(): self.assertEqual(response.headers.get(name), expected) def _document_fields(self, text, operation_id=None): return { 'document_text': text, 'operation_id': operation_id or str(uuid.uuid4()), 'expected_active_config_sha256': 'a' * 64, 'expected_active_secrets_sha256': 'b' * 64, 'expected_candidate_config_sha256': 'c' * 64, 'expected_candidate_secrets_sha256': 'd' * 64, } def test_runtime_document_pages_preview_save_and_apply_are_exact_and_no_store(self): with self._client() as client: config = client.get(ADMIN_PREFIX + '/config', headers=self._headers()) secrets = client.get(ADMIN_PREFIX + '/secrets', headers=self._headers()) self.assertEqual(config.status_code, 200) self.assertEqual(secrets.status_code, 200) self.assertSecurityHeaders(config) self.assertSecurityHeaders(secrets) self.assertIn('Config YAML', config.text) self.assertNotIn('editor-secret', config.text) self.assertIn('editor-secret', secrets.text) self.assertNotIn('', 'action': 'apply-secrets', 'target_kind': 'runtime-deployment', 'target_ref': 'secrets', 'status': 'failed', 'safe_category': 'startup_failed', 'safe_detail': 'startup_failed', 'expected_revision': None, 'resulting_revision': None, 'expected_identity': expected_identity, 'resulting_identity': None, 'agent_state': 'failed', 'agent_result_sha256': 'e' * 64, 'requested_at': '2026-09-20T00:00:00+00:00', 'started_at': '2026-09-20T00:00:01+00:00', 'completed_at': '2026-09-20T00:00:02+00:00', 'agent_reconciled_at': '2026-09-20T00:00:02+00:00', 'updated_at': '2026-09-20T00:00:02+00:00', } RecordingDB.audit_events = [{ 'id': 41, 'operation_id': operation_id, 'actor': '', 'action': 'apply-secrets', 'target_kind': 'runtime-deployment', 'target_ref': 'secrets', 'result': 'accepted', 'safe_category': None, 'before_identity': expected_identity, 'after_identity': None, 'before_bytes': None, 'after_bytes': None, 'previous_event_id': None, 'previous_event_sha256': None, 'event_sha256': 'f' * 64, 'created_at': '2026-09-20T00:00:00+00:00', }] RecordingDB.audit_next_before_event_id = 41 with self._client() as client: index = client.get( ADMIN_PREFIX + '/operations', headers=self._headers(), ) status = client.get( ADMIN_PREFIX + f'/operations/{operation_id}', headers=self._headers(), ) audit = client.get( ADMIN_PREFIX + '/audit', headers=self._headers(), ) older = client.get( ADMIN_PREFIX + '/audit?before=41', headers=self._headers(), ) for response in (index, status, audit, older): self.assertEqual(response.status_code, 200, response.text) self.assertSecurityHeaders(response) self.assertNotIn('editor-secret', response.text) self.assertNotIn('supervisor-token-must-not-render', response.text) self.assertNotIn(']*>).*?()', r'\1\2', secrets, flags=re.DOTALL | re.IGNORECASE, ) self.assertNotIn('editor-secret', outside_textarea) self.assertIn('
Workers / Dispatch', response.text) self.assertIn('Dispatch and drain', response.text) self.assertIn('Explicit pause', response.text) self.assertIn('Effective pause', response.text) self.assertIn('Drain blockers', response.text) self.assertIn('value="7"', response.text) for action in ( './dispatch/pause', './dispatch/resume', './dispatch/drain/start', './dispatch/drain/cancel', ): self.assertIn(f'action="{action}"', response.text) self.assertEqual(response.text.count('name="operation_id"'), 14) self.assertIn('Authenticated status, terminal reports, uploads', response.text) self.assertIn('Protocol-1 workers are completion-only', response.text) self.assertIn('linux-x86_64', response.text) self.assertIn('exact_git_v1', response.text) self.assertIn('docker_direct_v1', response.text) self.assertIn('huggingface_space_v1', response.text) self.assertIn('fixture-user', response.text) self.assertIn('fixture-device', response.text) self.assertIn('Accepted', response.text) self.assertIn('Ingestion / projection', response.text) snapshot_calls = [call for call in RecordingDB.calls if call[0] == 'snapshot'] self.assertEqual(len(snapshot_calls), 1) self.assertEqual(snapshot_calls[0][1], 25) self.assertFalse(any( call[0] in {'diagnostic_groups', 'duration_metrics'} for call in RecordingDB.calls )) self.assertIn('value="assignments" selected', response.text) self.assertIn('Not loaded. Select Assignments + diagnostics', response.text) self.assertRegex(snapshot_calls[0][2]['since'], r'^\d{4}-\d{2}-\d{2}T') self.assertIn(('drain_progress',), RecordingDB.calls) self.package_compatibility_provider.assert_called_once_with() self.assertNotIn('supervisor-token-must-not-render', response.text) def test_worker_observability_list_separates_final_models_and_keeps_occurrences(self): with self._client() as client: response = client.get( ADMIN_PREFIX + '?details=all', headers=self._headers(), ) self.assertEqual(response.status_code, 200, response.text) for label in ( 'Assignment outcome', 'Scan outcome', 'Diagnostics', 'Phase / progress', 'Deadlines', 'Slot / cap', 'Package identity', 'Ingestion / projection', ): self.assertIn(label, response.text) self.assertIn('active progress', response.text) self.assertIn('Observability scope:', response.text) self.assertEqual(response.text.count('