import asyncio import base64 import binascii import hashlib import hmac import html import json import logging import re import secrets import threading import time import uuid from datetime import datetime, timedelta, timezone from urllib.parse import parse_qsl, quote, urlencode, urlsplit from starlette.concurrency import run_in_threadpool from starlette.responses import Response, StreamingResponse from starlette.routing import Match, Route from managed_files import ( ManagedFileAccessError, ManagedFileDownload, ManagedFileIdentity, ManagedFileListing, ManagedFileMutation, ManagedFileOperation, ManagedFileRootRegistry, parse_managed_relative_path, ) from scanner_db import ( RuntimeControlRevisionConflictError, RuntimeControlTransitionError, RuntimeOperationIdentityConflictError, ) from runtime_document import ( MAX_CONFIG_DOCUMENT_BYTES, MAX_SECRETS_DOCUMENT_BYTES, RuntimeDocumentError, load_yaml_document, ) ADMIN_PREFIX = '/admin-internal' EDGE_MARKER_HEADER = 'x-truf-admin-edge' OPERATOR_HEADER = 'x-truf-admin-operator' OPERATOR_RE = re.compile(r'^[A-Za-z0-9_.-]{1,64}$') SHA256_RE = re.compile(r'^[0-9a-f]{64}$') DEFAULT_MAX_BODY_BYTES = 8 * 1024 DEFAULT_SNAPSHOT_LIMIT = 200 DEFAULT_REQUEUE_LIMIT = 100 DEFAULT_OVERVIEW_TIMEOUT_SECONDS = 5 DEFAULT_AUDIT_PAGE_LIMIT = 50 DEFAULT_OPERATION_PAGE_LIMIT = 50 DEFAULT_WORKER_PAGE_LIMIT = 25 WORKER_PAGE_LIMITS = (25, 50, 100) WORKER_FILTER_FIELDS = frozenset(( 'source', 'worker', 'assignment', 'scan', 'phase', 'category', 'code', 'retryable', 'window', 'limit', 'details', 'diagnostic_offset', 'metric_offset', )) WORKER_FILTER_WINDOWS = { '24h': timedelta(hours=24), '7d': timedelta(days=7), '30d': timedelta(days=30), '90d': timedelta(days=90), 'all': None, } KEY_RE = re.compile(r'^[A-Za-z0-9][A-Za-z0-9_.:@-]{0,127}$') QUEUE_STATUSES = ( 'pending', 'deferred', 'in_progress', 'done', 'failed', 'quarantined', 'cold', ) CORE_PRODUCERS = ('gitlab', 'dockerhub', 'huggingface') PRODUCER_IDS = tuple(f'discovery-producer:{source}' for source in CORE_PRODUCERS) PRODUCER_ACTIONS = ('start', 'stop', 'restart', 'pause', 'resume', 'set-interval') MAX_PRODUCER_INTERVAL_SECONDS = 365 * 24 * 60 * 60 PIPELINE_SOURCE_IDS = ('result-ingester', 'jsonl-projector', 'janitor', 'worker-api') MANAGED_SOURCE_ACTIONS = { **{source_id: PRODUCER_ACTIONS for source_id in PRODUCER_IDS}, **{ source_id: ( 'start', 'stop', 'restart', 'pause', 'resume', 'set-restart', 'set-restart-delay', ) for source_id in PIPELINE_SOURCE_IDS }, 'keychecks': ( 'start', 'stop', 'restart', 'pause', 'resume', 'once', 'set-mode', 'set-interval', 'set-restart', 'set-restart-delay', ), 'docker-shadow': ('start', 'stop'), 'dashboard': ('start', 'stop', 'restart'), } ALL_MANAGED_SOURCE_ACTIONS = frozenset(( 'start', 'stop', 'restart', 'pause', 'resume', 'once', 'set-mode', 'set-interval', 'set-restart', 'set-restart-delay', )) MAX_MANAGED_SOURCE_DELAY_SECONDS = 365 * 24 * 60 * 60 MAX_MANAGED_SOURCE_LOG_LINES = 5000 SECURITY_HEADERS = { 'Cache-Control': 'no-store', 'Referrer-Policy': 'same-origin', 'Content-Security-Policy': ( "default-src 'none'; style-src 'self'; script-src 'self'; form-action 'self'; " "base-uri 'none'; frame-ancestors 'none'" ), 'X-Content-Type-Options': 'nosniff', 'Strict-Transport-Security': 'max-age=31536000; includeSubDomains', } logger = logging.getLogger(__name__) class AdminAPIError(RuntimeError): def __init__(self, status_code, message): super().__init__(message) self.status_code = int(status_code) def _validate_key(value, label): value = str(value or '') if not KEY_RE.fullmatch(value): raise AdminAPIError(400, f'{label} is invalid') return value def _validate_cap(value): value = '' if value is None else str(value) if not re.fullmatch(r'0|[1-9][0-9]{0,4}', value): raise AdminAPIError(400, 'active assignment cap is invalid') cap = int(value) if cap > 10000: raise AdminAPIError(400, 'active assignment cap is invalid') return cap def _validate_origin(origin): origin = str(origin or '') if not 1 <= len(origin) <= 512 or any(character.isspace() for character in origin): raise ValueError('admin origin must be an exact HTTPS origin') try: parsed = urlsplit(origin) parsed_port = parsed.port except ValueError as exc: raise ValueError('admin origin must be an exact HTTPS origin') from exc if ( parsed.scheme != 'https' or not parsed.hostname or parsed.username is not None or parsed.password is not None or parsed.path or parsed.query or parsed.fragment or parsed_port is not None and not 1 <= parsed_port <= 65535 ): raise ValueError('admin origin must be an exact HTTPS origin') return origin class AdminService: def __init__( self, db_url, origin, edge_marker, *, db_factory, max_body_bytes=DEFAULT_MAX_BODY_BYTES, snapshot_limit=DEFAULT_SNAPSHOT_LIMIT, requeue_limit=DEFAULT_REQUEUE_LIMIT, supervisor_metadata=None, runtime_snapshot_provider=None, source_action_provider=None, dashboard_action_provider=None, source_log_provider=None, package_compatibility_provider=None, runtime_config_path=None, runtime_config_provider=None, document_loader=None, candidate_preview_provider=None, candidate_save_provider=None, candidate_verify_provider=None, runtime_apply_provider=None, managed_file_roots=None, ): self.db_url = str(db_url or '') self.origin = _validate_origin(origin) self.edge_marker = str(edge_marker or '') if ( not 32 <= len(self.edge_marker) <= 512 or any(character.isspace() for character in self.edge_marker) ): raise ValueError('admin edge marker must be a 32..512 character secret') self.db_factory = db_factory self.max_body_bytes = int(max_body_bytes) self.snapshot_limit = int(snapshot_limit) self.requeue_limit = int(requeue_limit) self.supervisor_metadata = dict(supervisor_metadata or {}) self.runtime_snapshot_provider = runtime_snapshot_provider self.source_action_provider = source_action_provider self.dashboard_action_provider = dashboard_action_provider self.source_log_provider = source_log_provider self.package_compatibility_provider = package_compatibility_provider self.runtime_config_path = str(runtime_config_path or '') self.runtime_config_provider = runtime_config_provider self.document_loader = document_loader self.candidate_preview_provider = candidate_preview_provider self.candidate_save_provider = candidate_save_provider self.candidate_verify_provider = candidate_verify_provider self.runtime_apply_provider = runtime_apply_provider if managed_file_roots is None: managed_file_roots = ManagedFileRootRegistry() if not isinstance(managed_file_roots, ManagedFileRootRegistry): raise ValueError('managed file root registry is invalid') self.managed_file_roots = managed_file_roots self._managed_file_operation_lock = threading.Lock() self._queue_snapshot_lock = threading.Lock() self._queue_snapshot_cache = None self._queue_snapshot_retry_at = 0.0 if not 1024 <= self.max_body_bytes <= 64 * 1024: raise ValueError('admin body byte limit must be between 1024 and 65536') if not 1 <= self.snapshot_limit <= 500: raise ValueError('admin snapshot limit must be between 1 and 500') if not 1 <= self.requeue_limit <= 500: raise ValueError('admin requeue limit must be between 1 and 500') self.csrf_token = secrets.token_urlsafe(48) def _call(self, method_name, *args, **kwargs): db = self.db_factory(db_url=self.db_url, initialize=False) if not db.enabled: db.close() raise RuntimeError('admin PostgreSQL connection is unavailable') try: return getattr(db, method_name)(*args, **kwargs) finally: db.close() def _worker_page_limit(self, limit): limit = int(limit) if limit not in WORKER_PAGE_LIMITS: raise ValueError('worker page limit is out of range') return min(limit, self.snapshot_limit) def snapshot(self, filters=None, *, limit=DEFAULT_WORKER_PAGE_LIMIT): return self._call( 'admin_remote_worker_snapshot', self._worker_page_limit(limit), filters=dict(filters or {}), ) def diagnostic_groups( self, filters=None, *, occurrence_offset=0, limit=DEFAULT_WORKER_PAGE_LIMIT, ): return self._call( 'admin_worker_diagnostic_groups', self._worker_page_limit(limit), filters=dict(filters or {}), occurrence_offset=occurrence_offset, ) def duration_metrics( self, filters=None, *, offset=0, limit=DEFAULT_WORKER_PAGE_LIMIT, ): filters = dict(filters or {}) if 'since' not in filters: filters['since'] = ( datetime.now(timezone.utc) - WORKER_FILTER_WINDOWS['30d'] ).isoformat(timespec='seconds') return self._call( 'admin_worker_duration_metrics', self._worker_page_limit(limit), filters=filters, offset=offset, ) def assignment_detail(self, reservation_id): detail = self._call( 'admin_worker_assignment_detail', reservation_id, event_limit=self.snapshot_limit, diagnostic_limit=self.snapshot_limit, ) if detail is None: raise AdminAPIError(404, 'worker assignment was not found') try: source = detail['assignment']['source'] detail['current_effective_policy'] = next( row for row in self.deadline_policy()['rows'] if row['source'] == source ) except (KeyError, StopIteration, TypeError, RuntimeError): detail['current_effective_policy'] = None return detail def diagnostic_envelope(self, reservation_id, diagnostic_uid): envelope = self._call( 'admin_worker_diagnostic_envelope', reservation_id, diagnostic_uid, ) if envelope is None: raise AdminAPIError(404, 'worker diagnostic was not found') return envelope @staticmethod def _deadline_policy(config): try: worker = config['supervisor']['worker_api'] sources = config['sources'] enabled = worker.get('sources') or ('gitlab', 'dockerhub', 'huggingface') fallback = int(worker['assignment_ttl_seconds']) overrides = dict(worker.get('assignment_ttl_seconds_by_source') or {}) upload = int(worker['bundle_body_timeout_seconds']) rows = [] for source in enabled: scan = int(sources[source]['timeout']) assignment = int(overrides.get(source, fallback)) required = scan + upload + 60 rows.append({ 'source': source, 'scan_deadline_seconds': scan, 'upload_deadline_seconds': upload, 'assignment_deadline_seconds': assignment, 'assignment_policy_source': ( 'source override' if source in overrides else 'global fallback' ), 'required_minimum_seconds': required, 'handoff_margin_seconds': 60, 'relationship': ( f'{assignment} >= {scan} + {upload} + 60' ), 'valid': assignment >= required, 'future_assignments_only': True, }) return {'rows': rows, 'future_assignments_only': True} except (KeyError, TypeError, ValueError, OverflowError) as exc: raise RuntimeError('runtime deadline policy is unavailable') from exc def deadline_policy(self): if not self.runtime_config_path: raise RuntimeError('runtime deadline policy is unavailable') provider = self.runtime_config_provider if provider is None: from runtime_document_io import load_managed_runtime_config provider = load_managed_runtime_config loaded = provider(self.runtime_config_path) config = getattr(loaded, 'config', None) if not isinstance(config, dict): raise RuntimeError('runtime deadline policy is unavailable') return self._deadline_policy(config) def deadline_policy_from_text(self, document_text): try: config = load_yaml_document( document_text.encode('utf-8', errors='strict'), max_bytes=MAX_CONFIG_DOCUMENT_BYTES, ) except (RuntimeDocumentError, UnicodeError) as exc: raise RuntimeError('runtime deadline policy is unavailable') from exc return self._deadline_policy(config) def runtime_document_observability(self, document_text): return { 'policy': self._overview_component( 'candidate deadline policy', lambda: self.deadline_policy_from_text(document_text), ), 'metrics': self._overview_component( 'worker duration metrics', self.duration_metrics, ), } def _package_compatibility_snapshot(self): provider = self.package_compatibility_provider if provider is None: raise RuntimeError('package compatibility is unavailable') snapshot = provider() if not isinstance(snapshot, dict) or set(snapshot) != { 'profiles', 'required_capabilities', }: raise RuntimeError('package compatibility is malformed') profiles = snapshot['profiles'] required = snapshot['required_capabilities'] if ( not isinstance(profiles, list) or not 1 <= len(profiles) <= 16 or not isinstance(required, list) or not 1 <= len(required) <= 16 ): raise RuntimeError('package compatibility is malformed') capability_fields = {'source', 'platform', 'planning_kind'} def capability(item): if not isinstance(item, dict) or set(item) != capability_fields: raise RuntimeError('package compatibility is malformed') value = {key: item[key] for key in capability_fields} if any( type(part) is not str or not re.fullmatch(r'[a-z][a-z0-9_]{0,63}', part) for part in value.values() ): raise RuntimeError('package compatibility is malformed') return value profile_fields = { 'profile_name', 'protocol_version', 'bundle_format_version', 'platform_tag', 'code_manifest_sha256', 'detector_policy_sha256', 'sources', 'capabilities', } normalized_profiles = [] for item in profiles: if not isinstance(item, dict) or set(item) != profile_fields: raise RuntimeError('package compatibility is malformed') if ( type(item['profile_name']) is not str or not KEY_RE.fullmatch(item['profile_name']) or type(item['protocol_version']) is not int or item['protocol_version'] < 1 or type(item['bundle_format_version']) is not int or item['bundle_format_version'] < 1 or type(item['platform_tag']) is not str or not 1 <= len(item['platform_tag']) <= 128 or not re.fullmatch(r'[a-f0-9]{64}', item['code_manifest_sha256']) or not re.fullmatch(r'[a-f0-9]{64}', item['detector_policy_sha256']) or not isinstance(item['sources'], list) or not 1 <= len(item['sources']) <= 16 or len(item['sources']) != len(set(item['sources'])) or any( type(source) is not str or not re.fullmatch(r'[a-z][a-z0-9_]{0,63}', source) for source in item['sources'] ) or not isinstance(item['capabilities'], list) or not 1 <= len(item['capabilities']) <= 16 ): raise RuntimeError('package compatibility is malformed') normalized_profiles.append({ key: item[key] for key in profile_fields - {'sources', 'capabilities'} } | { 'sources': list(item['sources']), 'capabilities': [capability(value) for value in item['capabilities']], }) return { 'profiles': normalized_profiles, 'required_capabilities': [capability(item) for item in required], } def _runtime_snapshot(self): if not self.supervisor_metadata: raise RuntimeError('Supervisor metadata is unavailable') provider = self.runtime_snapshot_provider if provider is None: from supervisor import get_runtime_snapshot provider = get_runtime_snapshot snapshot = provider( self.supervisor_metadata, timeout=DEFAULT_OVERVIEW_TIMEOUT_SECONDS, ) if not ( isinstance(snapshot, dict) and isinstance(snapshot.get('runtime'), dict) and isinstance(snapshot.get('postgres'), dict) and isinstance(snapshot.get('pipeline'), dict) and isinstance(snapshot.get('sources'), list) and all(isinstance(item, dict) for item in snapshot['sources']) ): raise RuntimeError('Supervisor snapshot is malformed') return snapshot def _queue_snapshot(self): def copied(snapshot): return dict( snapshot, counts=dict(snapshot.get('counts') or {}), truncated_statuses=list(snapshot.get('truncated_statuses') or []), ) with self._queue_snapshot_lock: now = time.monotonic() if now < self._queue_snapshot_retry_at: if self._queue_snapshot_cache is None: raise RuntimeError('queue snapshot is unavailable') snapshot = copied(self._queue_snapshot_cache) snapshot.update({ 'degraded': True, 'stale': True, 'reason': 'bounded_count_retry_backoff', 'retry_after_sec': max(1, int(self._queue_snapshot_retry_at - now)), }) return snapshot snapshot = self._call('admin_target_queue_health', CORE_PRODUCERS) if not isinstance(snapshot, dict) or not isinstance(snapshot.get('counts'), dict): raise RuntimeError('queue snapshot is malformed') counts = snapshot['counts'] if any( status not in QUEUE_STATUSES or type(value) is not int or value < 0 for status, value in counts.items() ): raise RuntimeError('queue snapshot is malformed') truncated = snapshot.get('truncated_statuses') if not isinstance(truncated, list) or any( type(status) is not str or status not in QUEUE_STATUSES for status in truncated ): raise RuntimeError('queue snapshot is malformed') retry_after = snapshot.get('retry_after_sec') if type(retry_after) is not int or not 0 <= retry_after <= 3600: raise RuntimeError('queue snapshot is malformed') if snapshot.get('stale') is True: self._queue_snapshot_retry_at = now + retry_after if self._queue_snapshot_cache is not None: cached = copied(self._queue_snapshot_cache) cached.update({ 'degraded': True, 'stale': True, 'reason': 'bounded_count_query_failed', 'retry_after_sec': retry_after, }) return cached if not counts: raise RuntimeError('queue snapshot is unavailable') return copied(snapshot) self._queue_snapshot_cache = copied(snapshot) self._queue_snapshot_retry_at = 0.0 return copied(snapshot) def _control_snapshot(self): snapshot = self._call('runtime_drain_progress') if not isinstance(snapshot, dict): raise RuntimeError('control snapshot is malformed') if ( type(snapshot.get('revision')) is not int or snapshot['revision'] < 0 or any( type(snapshot.get(key)) is not bool for key in ( 'discovery_paused', 'dispatch_paused', 'effective_discovery_paused', 'effective_dispatch_paused', ) ) or snapshot.get('drain_state') not in ('normal', 'draining', 'drained') or any( type(snapshot.get(key)) is not int or snapshot[key] < 0 for key in ( 'live_remote_assignments', 'precommit_result_bundles', 'blocker_count', ) ) or type(snapshot.get('actor')) is not str or type(snapshot.get('updated_at')) is not str ): raise RuntimeError('control snapshot is malformed') drain_active = snapshot['drain_state'] != 'normal' if ( snapshot['blocker_count'] != ( snapshot['live_remote_assignments'] + snapshot['precommit_result_bundles'] ) or snapshot['effective_discovery_paused'] is not ( snapshot['discovery_paused'] or drain_active ) or snapshot['effective_dispatch_paused'] is not ( snapshot['dispatch_paused'] or drain_active ) ): raise RuntimeError('control snapshot is malformed') return snapshot def _recent_operations(self): rows = self._call('recent_runtime_operations', self.snapshot_limit) fields = ( 'operation_id', 'actor', 'action', 'target_ref', 'status', 'safe_category', 'safe_detail', 'requested_at', 'completed_at', 'updated_at', ) if not isinstance(rows, list) or len(rows) > self.snapshot_limit: raise RuntimeError('operation snapshot is malformed') result = [] for row in rows: if not isinstance(row, dict) or any( row.get(key) is not None and type(row.get(key)) is not str for key in fields ): raise RuntimeError('operation snapshot is malformed') result.append({key: row.get(key) for key in fields}) return result def operation_page(self, before=None): before_updated_at = before_operation_id = None if before is not None: before_updated_at, before_operation_id = before rows = self._call( 'recent_runtime_operations', DEFAULT_OPERATION_PAGE_LIMIT + 1, before_updated_at=before_updated_at, before_operation_id=before_operation_id, ) fields = ( 'operation_id', 'actor', 'action', 'target_ref', 'status', 'safe_category', 'safe_detail', 'requested_at', 'completed_at', 'updated_at', ) if not isinstance(rows, list) or len(rows) > DEFAULT_OPERATION_PAGE_LIMIT + 1: raise RuntimeError('operation page is malformed') operations = [] for row in rows[:DEFAULT_OPERATION_PAGE_LIMIT]: if not isinstance(row, dict) or any( row.get(key) is not None and type(row.get(key)) is not str for key in fields ): raise RuntimeError('operation page is malformed') if not row.get('operation_id') or not row.get('updated_at'): raise RuntimeError('operation page is malformed') operations.append({key: row.get(key) for key in fields}) next_before = None if len(rows) > DEFAULT_OPERATION_PAGE_LIMIT and operations: last = operations[-1] next_before = (last['updated_at'], last['operation_id']) return {'operations': operations, 'next_before': next_before} def operation_status(self, operation_id): operation = self._call('runtime_operation', operation_id) if operation is None: raise AdminAPIError(404, 'Not Found') return operation def audit_page(self, before_event_id=None): return self._call( 'runtime_audit_events', before_event_id=before_event_id, limit=min(DEFAULT_AUDIT_PAGE_LIMIT, self.snapshot_limit), ) @staticmethod def _managed_file_traversal(traversal): if traversal is None: raise AdminAPIError(503, 'managed files are unavailable') return traversal def list_managed_files(self, traversal, root_id, relative_path=None): traversal = self._managed_file_traversal(traversal) root = self.managed_file_roots.get(root_id) if root is None: raise AdminAPIError(404, 'managed file target was not found') if not root.permissions.allow_list: raise AdminAPIError(403, 'managed file operation is not allowed') if relative_path is not None: try: parse_managed_relative_path(relative_path, root.limits) except ManagedFileAccessError as exc: _managed_file_error(exc) try: listing = traversal.list_directory(root_id, relative_path) except ManagedFileAccessError as exc: _managed_file_error(exc) if not isinstance(listing, ManagedFileListing): raise RuntimeError('managed file listing is malformed') return listing def download_managed_file(self, traversal, root_id, relative_path): traversal = self._managed_file_traversal(traversal) self._managed_file_root(root_id, relative_path, ManagedFileOperation.READ) try: download = traversal.download_file(root_id, relative_path) except ManagedFileAccessError as exc: _managed_file_error(exc) if not isinstance(download, ManagedFileDownload): raise RuntimeError('managed file download is malformed') return download @staticmethod def _overview_component(name, callback): try: value = callback() if value is None: raise RuntimeError('overview component is unavailable') return {'available': True, 'value': value} except Exception as exc: logger.warning('Admin overview component %s is unavailable: %s', name, type(exc).__name__) return {'available': False, 'value': None} def overview(self): return { 'runtime': self._overview_component('runtime', self._runtime_snapshot), 'queue': self._overview_component( 'queue', self._queue_snapshot, ), 'control': self._overview_component( 'control', self._control_snapshot, ), 'operations': self._overview_component( 'operations', self._recent_operations, ), } def search_snapshot(self): return { 'runtime': self._overview_component('runtime', self._runtime_snapshot), 'control': self._overview_component('control', self._control_snapshot), } def workers_dispatch_snapshot( self, filters=None, *, diagnostic_occurrence_offset=0, metric_offset=0, page_limit=DEFAULT_WORKER_PAGE_LIMIT, include_diagnostics=False, include_metrics=False, ): filters = dict(filters or {}) return { 'workers': self._overview_component( 'workers', lambda: self.snapshot(filters, limit=page_limit), ), 'diagnostics': ( self._overview_component( 'diagnostics', lambda: self.diagnostic_groups( filters, occurrence_offset=diagnostic_occurrence_offset, limit=page_limit, ), ) if include_diagnostics else {'available': False, 'value': None} ), 'metrics': ( self._overview_component( 'metrics', lambda: self.duration_metrics( filters, offset=metric_offset, limit=page_limit, ), ) if include_metrics else {'available': False, 'value': None} ), 'policy': self._overview_component('policy', self.deadline_policy), 'control': self._overview_component('control', self._control_snapshot), 'packages': self._overview_component( 'packages', self._package_compatibility_snapshot, ), } def set_dispatch_paused(self, paused, expected_revision, actor, operation_id): if type(paused) is not bool: raise AdminAPIError(400, 'dispatch control state is invalid') if type(expected_revision) is not int or expected_revision < 0: raise AdminAPIError(400, 'control revision is invalid') try: return self._call( 'set_runtime_dispatch_paused', paused, expected_revision=expected_revision, actor=actor, operation_id=operation_id, ) except ( RuntimeControlRevisionConflictError, RuntimeControlTransitionError, RuntimeOperationIdentityConflictError, ) as exc: raise AdminAPIError(409, 'dispatch control changed; refresh and retry') from exc def start_drain(self, expected_revision, actor, operation_id): return self._drain_action( 'start_runtime_drain', expected_revision, actor, operation_id, ) def cancel_drain(self, expected_revision, actor, operation_id): return self._drain_action( 'cancel_runtime_drain', expected_revision, actor, operation_id, ) def _drain_action(self, method, expected_revision, actor, operation_id): if type(expected_revision) is not int or expected_revision < 0: raise AdminAPIError(400, 'control revision is invalid') try: return self._call( method, expected_revision=expected_revision, actor=actor, operation_id=operation_id, ) except ( RuntimeControlRevisionConflictError, RuntimeControlTransitionError, RuntimeOperationIdentityConflictError, ) as exc: raise AdminAPIError(409, 'drain control changed; refresh and retry') from exc def set_discovery_paused(self, paused, expected_revision, actor, operation_id): if type(paused) is not bool: raise AdminAPIError(400, 'discovery control state is invalid') if type(expected_revision) is not int or expected_revision < 0: raise AdminAPIError(400, 'control revision is invalid') try: return self._call( 'set_runtime_discovery_paused', paused, expected_revision=expected_revision, actor=actor, operation_id=operation_id, ) except ( RuntimeControlRevisionConflictError, RuntimeControlTransitionError, RuntimeOperationIdentityConflictError, ) as exc: raise AdminAPIError(409, 'discovery control changed; refresh and retry') from exc def _complete_producer_operation(self, operation_id, *, succeeded, outcome=None): last_error = None for attempt in range(3): try: return self._call( 'complete_runtime_source_operation', operation_id, succeeded=succeeded, outcome=outcome, ) except RuntimeOperationIdentityConflictError: raise except Exception as exc: last_error = exc if attempt < 2: time.sleep(0.05) raise last_error def _managed_source_entry(self, source_id): if type(source_id) is not str or not KEY_RE.fullmatch(source_id) or source_id == 'all': raise AdminAPIError(400, 'managed source action is invalid') snapshot = self._runtime_snapshot() matches = [ item for item in snapshot.get('sources', []) if isinstance(item, dict) and item.get('id') == source_id ] if len(matches) != 1: raise AdminAPIError(400, 'managed source action is invalid') allowed = matches[0].get('allowed_actions') if ( not isinstance(allowed, list) or len(allowed) != len(set(allowed)) or any(action not in ALL_MANAGED_SOURCE_ACTIONS for action in allowed) ): raise AdminAPIError(502, 'managed source state is invalid') return matches[0] def managed_source_action( self, source_id, source_action, actor, operation_id, *, interval_seconds=None, mode=None, restart_enabled=None, restart_delay_seconds=None, ): if source_id == 'dashboard': allowed_actions = MANAGED_SOURCE_ACTIONS['dashboard'] else: allowed_actions = self._managed_source_entry(source_id).get('allowed_actions') if source_action not in allowed_actions: raise AdminAPIError(400, 'managed source action is invalid') if source_action == 'set-interval': if ( type(interval_seconds) is not int or not 1 <= interval_seconds <= MAX_PRODUCER_INTERVAL_SECONDS ): raise AdminAPIError(400, 'managed source interval is invalid') elif interval_seconds is not None: raise AdminAPIError(400, 'managed source action is invalid') if source_action == 'set-mode': if mode not in ('loop', 'once', 'repeat') or ( source_id == 'keychecks' and mode == 'loop' ): raise AdminAPIError(400, 'managed source mode is invalid') elif mode is not None: raise AdminAPIError(400, 'managed source action is invalid') if source_action == 'set-restart': if type(restart_enabled) is not bool: raise AdminAPIError(400, 'managed source restart setting is invalid') elif restart_enabled is not None: raise AdminAPIError(400, 'managed source action is invalid') if source_action == 'set-restart-delay': if ( type(restart_delay_seconds) is not int or not 1 <= restart_delay_seconds <= MAX_MANAGED_SOURCE_DELAY_SECONDS ): raise AdminAPIError(400, 'managed source restart delay is invalid') elif restart_delay_seconds is not None: raise AdminAPIError(400, 'managed source action is invalid') try: operation_parameters = {} if mode is not None: operation_parameters['mode'] = mode if restart_enabled is not None: operation_parameters['restart_enabled'] = restart_enabled if restart_delay_seconds is not None: operation_parameters['restart_delay_seconds'] = restart_delay_seconds created = self._call( 'create_runtime_source_operation', operation_id=operation_id, actor=actor, source_id=source_id, source_action=source_action, interval_seconds=interval_seconds, **operation_parameters, ) except RuntimeOperationIdentityConflictError as exc: raise AdminAPIError(409, 'operation ID is already bound to another request') from exc if created.get('replayed') is True: if created.get('status') == 'succeeded': resulting = created.get('resulting_identity') or {} return {'outcome': resulting.get('outcome', 'completed')} if created.get('status') == 'running': raise AdminAPIError(409, 'managed source action is already in progress') raise AdminAPIError(409, 'managed source action already completed with failure') parameters = {} if interval_seconds is not None: parameters['interval_seconds'] = interval_seconds if mode is not None: parameters['mode'] = mode if restart_enabled is not None: parameters['restart_enabled'] = restart_enabled if restart_delay_seconds is not None: parameters['restart_delay_seconds'] = restart_delay_seconds try: if source_id == 'dashboard': provider = self.dashboard_action_provider if provider is None: from supervisor import send_dashboard_action provider = send_dashboard_action result = provider( self.supervisor_metadata, source_action, timeout=60, ) else: provider = self.source_action_provider if provider is None: from supervisor import send_managed_source_action provider = send_managed_source_action result = provider( self.supervisor_metadata, source_id, source_action, timeout=60, **parameters, ) if ( not isinstance(result, dict) or result.get('outcome') not in ('completed', 'dependency-blocked') ): raise RuntimeError('managed source action response is invalid') except Exception: try: self._complete_producer_operation(operation_id, succeeded=False) except Exception as completion_error: logger.error( 'Managed source operation failure reconciliation failed: %s', type(completion_error).__name__, ) raise AdminAPIError(502, 'managed source action failed') from None try: self._complete_producer_operation( operation_id, succeeded=True, outcome=result['outcome'], ) except RuntimeOperationIdentityConflictError as exc: raise AdminAPIError(409, 'managed source operation completion conflicted') from exc return result def producer_action( self, source_id, source_action, actor, operation_id, *, interval_seconds=None, ): if source_id not in PRODUCER_IDS or source_action not in PRODUCER_ACTIONS: raise AdminAPIError(400, 'producer action is invalid') return self.managed_source_action( source_id, source_action, actor, operation_id, interval_seconds=interval_seconds, ) def managed_source_log(self, source_id, line_count): self._managed_source_entry(source_id) if type(line_count) is not int or not 1 <= line_count <= MAX_MANAGED_SOURCE_LOG_LINES: raise AdminAPIError(400, 'managed source log line count is invalid') provider = self.source_log_provider if provider is None: from supervisor import send_managed_source_log_tail provider = send_managed_source_log_tail try: result = provider( self.supervisor_metadata, source_id, line_count, timeout=60, ) except Exception: raise AdminAPIError(502, 'managed source log tail failed') from None if ( not isinstance(result, dict) or result.get('source_id') != source_id or result.get('line_count') != len(result.get('lines') or []) or not isinstance(result.get('lines'), list) or any(type(line) is not str for line in result['lines']) or type(result.get('response_truncated')) is not bool ): raise AdminAPIError(502, 'managed source log tail failed') return { 'source_id': source_id, 'line_count': result['line_count'], 'lines': list(result['lines']), 'response_truncated': result['response_truncated'], } @staticmethod def _runtime_document_error(exc): status = 409 if exc.category == 'reference' and exc.path == 'revision' else 400 document = exc.document or 'document' path = exc.path or 'root' raise AdminAPIError( status, f'{document} validation failed ({exc.category}) at {path}', ) from exc def runtime_document_editor(self, document): if document not in ('config', 'secrets') or not self.runtime_config_path: raise AdminAPIError(503, 'runtime document editor is unavailable') provider = self.document_loader if provider is None: from runtime_document_io import load_managed_runtime_editor_document provider = load_managed_runtime_editor_document try: editor = provider(self.runtime_config_path, document) except RuntimeDocumentError as exc: self._runtime_document_error(exc) if ( getattr(editor, 'document', None) != document or getattr(editor, 'source', None) not in ('active', 'candidate') or type(getattr(editor, 'text', None)) is not str or getattr(editor, 'state', None) is None ): raise AdminAPIError(503, 'runtime document editor is unavailable') return editor def preview_runtime_document(self, document, document_text): payload = None provider = self.candidate_preview_provider if provider is None: from runtime_document_io import preview_managed_runtime_candidate provider = preview_managed_runtime_candidate try: payload = document_text.encode('utf-8', errors='strict') document_text = None return provider( self.runtime_config_path, document, payload, ) except RuntimeDocumentError as exc: self._runtime_document_error(exc) finally: document_text = payload = None def save_runtime_document_candidate( self, document, document_text, actor, operation_id, *, expected_hashes, ): candidate_bytes = None try: candidate_bytes = document_text.encode('utf-8', errors='strict') document_text = None return self._save_runtime_document_candidate_bytes( document, candidate_bytes, actor, operation_id, expected_hashes=expected_hashes, ) finally: document_text = candidate_bytes = None def _save_runtime_document_candidate_bytes( self, document, candidate_bytes, actor, operation_id, *, expected_hashes, ): preview = operation = revision = current = selected = None def complete_success(candidate_sha256, written): for attempt in range(3): try: return self._call( 'complete_runtime_document_operation', operation_id, succeeded=True, candidate_sha256=candidate_sha256, written=written, ) except Exception: if attempt == 2: raise AdminAPIError( 503, 'runtime document completion is pending', ) from None try: action = f'runtime.{document}.save' selected_key = f'candidate_{document}' other_key = ( 'candidate_secrets' if document == 'config' else 'candidate_config' ) proposed_sha256 = hashlib.sha256(candidate_bytes).hexdigest() proposed_bytes = len(candidate_bytes) operation = self._call('runtime_operation', operation_id) if operation is not None: identity = operation.get('expected_identity') or {} selected_identity_key = f'{selected_key}_sha256' other_identity_key = f'{other_key}_sha256' submitted_selected = expected_hashes[selected_key] identity_matches = ( operation.get('actor') == actor and operation.get('action') == action and operation.get('target_kind') == 'runtime-document' and operation.get('target_ref') == document and hmac.compare_digest( identity.get('active_config_sha256', ''), expected_hashes['active_config'], ) and hmac.compare_digest( identity.get('active_secrets_sha256', ''), expected_hashes['active_secrets'], ) and hmac.compare_digest( identity.get(other_identity_key, ''), expected_hashes[other_key], ) and any( hmac.compare_digest(identity.get(key, ''), submitted_selected) for key in (selected_identity_key, 'candidate_after_sha256') ) and hmac.compare_digest( identity.get('candidate_after_sha256', ''), proposed_sha256, ) and identity.get('candidate_after_bytes') == proposed_bytes ) if not identity_matches: raise AdminAPIError( 409, 'operation ID is already bound to another request', ) try: operation = self._call( 'create_runtime_document_operation', operation_id=operation_id, actor=actor, action=action, active_config_sha256=identity['active_config_sha256'], active_secrets_sha256=identity['active_secrets_sha256'], candidate_config_sha256=identity['candidate_config_sha256'], candidate_secrets_sha256=identity['candidate_secrets_sha256'], candidate_after_sha256=identity['candidate_after_sha256'], candidate_before_bytes=identity['candidate_before_bytes'], candidate_after_bytes=identity['candidate_after_bytes'], candidate_before_present=identity['candidate_before_present'], ) except RuntimeOperationIdentityConflictError as exc: raise AdminAPIError( 409, 'operation ID is already bound to another request', ) from exc if operation.get('status') == 'succeeded': return operation if operation.get('status') != 'running': raise AdminAPIError(409, 'runtime document save already has another result') current = self.runtime_document_editor(document) current_hashes = _runtime_document_hashes(current.state) if any( not hmac.compare_digest(current_hashes[key], identity[f'{key}_sha256']) for key in ('active_config', 'active_secrets', other_key) ): raise AdminAPIError(409, 'runtime document revision changed') selected = ( current.state.candidate_config if document == 'config' else current.state.candidate_secrets ) if selected.present and ( hmac.compare_digest(selected.sha256, identity['candidate_after_sha256']) and selected.byte_count == identity['candidate_after_bytes'] ): return complete_success( selected.sha256, written=( not identity['candidate_before_present'] or identity[selected_identity_key] != selected.sha256 or identity['candidate_before_bytes'] != selected.byte_count ), ) if not ( hmac.compare_digest(selected.sha256, identity[selected_identity_key]) and selected.byte_count == identity['candidate_before_bytes'] ): raise AdminAPIError(409, 'runtime document revision changed') expected_hashes = { 'active_config': identity['active_config_sha256'], 'active_secrets': identity['active_secrets_sha256'], 'candidate_config': identity['candidate_config_sha256'], 'candidate_secrets': identity['candidate_secrets_sha256'], } else: provider = self.candidate_preview_provider if provider is None: from runtime_document_io import preview_managed_runtime_candidate provider = preview_managed_runtime_candidate try: preview = provider(self.runtime_config_path, document, candidate_bytes) except RuntimeDocumentError as exc: self._runtime_document_error(exc) current_hashes = _runtime_document_hashes(preview.state) if any( not hmac.compare_digest(current_hashes[key], expected_hashes[key]) for key in current_hashes ): raise AdminAPIError(409, 'runtime document revision changed') selected = ( preview.state.candidate_config if document == 'config' else preview.state.candidate_secrets ) try: operation = self._call( 'create_runtime_document_operation', operation_id=operation_id, actor=actor, action=action, active_config_sha256=expected_hashes['active_config'], active_secrets_sha256=expected_hashes['active_secrets'], candidate_config_sha256=expected_hashes['candidate_config'], candidate_secrets_sha256=expected_hashes['candidate_secrets'], candidate_after_sha256=preview.proposed.sha256, candidate_before_bytes=selected.byte_count, candidate_after_bytes=preview.proposed.byte_count, candidate_before_present=selected.present, ) except RuntimeOperationIdentityConflictError as exc: raise AdminAPIError( 409, 'operation ID is already bound to another request', ) from exc provider = self.candidate_save_provider if provider is None: from runtime_document_io import save_managed_runtime_candidate provider = save_managed_runtime_candidate try: revision = provider( self.runtime_config_path, document, candidate_bytes, expected_active_config_sha256=expected_hashes['active_config'], expected_active_secrets_sha256=expected_hashes['active_secrets'], expected_candidate_config_sha256=expected_hashes['candidate_config'], expected_candidate_secrets_sha256=expected_hashes['candidate_secrets'], ) except RuntimeDocumentError as exc: try: current = self.runtime_document_editor(document) except AdminAPIError: current = None if current is not None: current_hashes = _runtime_document_hashes(current.state) selected = ( current.state.candidate_config if document == 'config' else current.state.candidate_secrets ) other_key = ( 'candidate_secrets' if document == 'config' else 'candidate_config' ) if ( all( hmac.compare_digest( current_hashes[key], expected_hashes[key], ) for key in ('active_config', 'active_secrets', other_key) ) and selected.present and hmac.compare_digest( selected.sha256, operation['expected_identity'][ 'candidate_after_sha256' ], ) and selected.byte_count == operation['expected_identity'][ 'candidate_after_bytes' ] ): identity = operation['expected_identity'] return complete_success( selected.sha256, written=( not identity['candidate_before_present'] or identity[f'candidate_{document}_sha256'] != selected.sha256 or identity['candidate_before_bytes'] != selected.byte_count ), ) for _ in range(3): try: self._call( 'complete_runtime_document_operation', operation_id, succeeded=False, ) break except Exception: continue self._runtime_document_error(exc) complete_success(revision.proposed.sha256, revision.written) return revision finally: candidate_bytes = preview = operation = revision = current = selected = None complete_success = None def request_runtime_apply( self, action, actor, operation_id, *, expected_hashes, ): if action not in ('apply-config', 'apply-secrets', 'apply-both'): raise AdminAPIError(400, 'runtime apply action is invalid') candidate_config = ( expected_hashes['candidate_config'] if action in ('apply-config', 'apply-both') else None ) candidate_secrets = ( expected_hashes['candidate_secrets'] if action in ('apply-secrets', 'apply-both') else None ) expected_identity = { 'active_config_sha256': expected_hashes['active_config'], 'active_secrets_sha256': expected_hashes['active_secrets'], 'candidate_config_sha256': candidate_config, 'candidate_secrets_sha256': candidate_secrets, } verified = None operation = self._call('runtime_operation', operation_id) if operation is None: if self.runtime_apply_provider is None: raise AdminAPIError(503, 'runtime apply agent is unavailable') verify = self.candidate_verify_provider if verify is None: from runtime_document_io import verify_managed_runtime_candidates verify = verify_managed_runtime_candidates try: verified = verify( self.runtime_config_path, action, expected_active_config_sha256=expected_hashes['active_config'], expected_active_secrets_sha256=expected_hashes['active_secrets'], expected_candidate_config_sha256=candidate_config, expected_candidate_secrets_sha256=candidate_secrets, ) except RuntimeDocumentError as exc: self._runtime_document_error(exc) try: operation = self._call( 'create_runtime_operation', operation_id=operation_id, actor=actor, action=action, expected_identity=expected_identity, ) except RuntimeOperationIdentityConflictError as exc: raise AdminAPIError(409, 'operation ID is already bound to another request') from exc if operation.get('status') in ('succeeded', 'failed', 'rolled_back', 'failed_hold'): if operation.get('status') == 'succeeded': return {'operation': operation, 'verification': verified} raise AdminAPIError(409, 'runtime apply operation already has a terminal result') if self.runtime_apply_provider is None: raise AdminAPIError(503, 'runtime apply agent is unavailable') try: self.runtime_apply_provider( operation_id=operation_id, action=action, active_config_sha256=expected_identity['active_config_sha256'], active_secrets_sha256=expected_identity['active_secrets_sha256'], candidate_config_sha256=expected_identity['candidate_config_sha256'], candidate_secrets_sha256=expected_identity['candidate_secrets_sha256'], ) except Exception: raise AdminAPIError(502, 'runtime apply dispatch failed') from None return {'operation': operation, 'verification': verified} def _managed_file_root(self, root_id, relative_path, operation): root = self.managed_file_roots.get(root_id) if root is None: raise AdminAPIError(404, 'managed file target was not found') if not root.permissions.allows(operation): raise AdminAPIError(403, 'managed file operation is not allowed') try: parse_managed_relative_path(relative_path, root.limits) except ManagedFileAccessError as exc: _managed_file_error(exc) return root def _complete_managed_file_operation( self, operation_id, *, succeeded, mutation=None, before_sha256=None, before_byte_count=None, after_sha256=None, after_byte_count=None, written=None, ): if mutation is not None: before_sha256 = mutation.before.sha256 if mutation.before else None before_byte_count = mutation.before.byte_count if mutation.before else None after_sha256 = mutation.after.sha256 if mutation.after else None after_byte_count = mutation.after.byte_count if mutation.after else None written = mutation.written arguments = { 'succeeded': succeeded, 'before_sha256': before_sha256, 'before_byte_count': before_byte_count, 'after_sha256': after_sha256, 'after_byte_count': after_byte_count, 'written': written, } if not succeeded: arguments = {'succeeded': False} last_error = None for attempt in range(3): try: return self._call( 'complete_runtime_managed_file_operation', operation_id, **arguments, ) except RuntimeOperationIdentityConflictError: raise except Exception as exc: last_error = exc if attempt < 2: time.sleep(0.05) raise last_error @staticmethod def _managed_file_observed_identity( traversal, root_id, relative_path, operation_type, *, require_private_sha256=None, ): try: return traversal.mutation_file_identity( root_id, relative_path, operation_type, require_private_sha256=require_private_sha256, ) except ManagedFileAccessError as exc: if exc.category == 'not_found': return None _managed_file_error(exc) def _complete_observed_managed_file_operation( self, traversal, *, action, root_id, relative_path, operation_type, operation_id, expected_sha256, proposed_sha256, proposed_byte_count, ): current = self._managed_file_observed_identity( traversal, root_id, relative_path, operation_type, require_private_sha256=( proposed_sha256 if action in ('files.create', 'files.replace') else None ), ) if action == 'files.create': if current is None: return None if ( current.sha256 != proposed_sha256 or current.byte_count != proposed_byte_count ): raise AdminAPIError(409, 'managed file state changed concurrently') return self._complete_managed_file_operation( operation_id, succeeded=True, after_sha256=current.sha256, after_byte_count=current.byte_count, written=True, ) if action == 'files.replace': if current is None: raise AdminAPIError(409, 'managed file state changed concurrently') if ( current.sha256 == proposed_sha256 and current.byte_count == proposed_byte_count ): no_write = hmac.compare_digest(expected_sha256, proposed_sha256) return self._complete_managed_file_operation( operation_id, succeeded=True, before_sha256=expected_sha256, before_byte_count=current.byte_count if no_write else None, after_sha256=current.sha256, after_byte_count=current.byte_count, written=not no_write, ) if not hmac.compare_digest(current.sha256, expected_sha256): raise AdminAPIError(409, 'managed file state changed concurrently') return None if current is None: return self._complete_managed_file_operation( operation_id, succeeded=True, before_sha256=expected_sha256, before_byte_count=None, after_sha256=None, after_byte_count=None, written=True, ) if not hmac.compare_digest(current.sha256, expected_sha256): raise AdminAPIError(409, 'managed file state changed concurrently') return None def _managed_file_mutation_locked_payload( self, traversal, *, action, root_id, relative_path, actor, operation_id, payload_box, expected_sha256=None, ): operation_type = ( ManagedFileOperation.DELETE if action == 'files.delete' else ManagedFileOperation.CREATE_REPLACE ) proposed_sha256 = proposed_byte_count = None if action in ('files.create', 'files.replace'): if len(payload_box) != 1 or type(payload_box[0]) is not bytes: raise AdminAPIError(400, 'managed file content is invalid') proposed_sha256 = hashlib.sha256(payload_box[0]).hexdigest() proposed_byte_count = len(payload_box[0]) if ( action in ('files.replace', 'files.delete') and not SHA256_RE.fullmatch(str(expected_sha256 or '')) ): raise AdminAPIError(400, 'managed file revision is invalid') existing = self._call('runtime_operation', operation_id) if existing is None: traversal = self._managed_file_traversal(traversal) self._managed_file_root(root_id, relative_path, operation_type) current = self._managed_file_observed_identity( traversal, root_id, relative_path, operation_type, ) if action == 'files.create': if current is not None: raise AdminAPIError(409, 'managed file state changed concurrently') elif ( current is None or not hmac.compare_digest(current.sha256, expected_sha256) ): raise AdminAPIError(409, 'managed file state changed concurrently') try: created = self._call( 'create_runtime_managed_file_operation', operation_id=operation_id, actor=actor, action=action, root_id=root_id, relative_path=relative_path, expected_sha256=expected_sha256, proposed_sha256=proposed_sha256, proposed_byte_count=proposed_byte_count, ) except RuntimeOperationIdentityConflictError as exc: raise AdminAPIError( 409, 'operation ID is already bound to another request', ) from exc if created.get('status') == 'succeeded': return created if created.get('status') != 'running': raise AdminAPIError(409, 'managed file operation already has a terminal result') traversal = self._managed_file_traversal(traversal) self._managed_file_root(root_id, relative_path, operation_type) if created.get('replayed') is True: completed = self._complete_observed_managed_file_operation( traversal, action=action, root_id=root_id, relative_path=relative_path, operation_type=operation_type, operation_id=operation_id, expected_sha256=expected_sha256, proposed_sha256=proposed_sha256, proposed_byte_count=proposed_byte_count, ) if completed is not None: return completed try: if action == 'files.delete': mutation = traversal.delete_file( root_id, relative_path, expected_sha256=expected_sha256, ) else: mutation = traversal.create_replace_file( root_id, relative_path, payload_box[0], expected_sha256=( expected_sha256 if action == 'files.replace' else None ), ) except ManagedFileAccessError as exc: if ( created.get('replayed') is True and exc.category in ('hash_conflict', 'concurrent_change', 'not_found') ): try: completed = self._complete_observed_managed_file_operation( traversal, action=action, root_id=root_id, relative_path=relative_path, operation_type=operation_type, operation_id=operation_id, expected_sha256=expected_sha256, proposed_sha256=proposed_sha256, proposed_byte_count=proposed_byte_count, ) if completed is not None: return completed except AdminAPIError as state_error: if state_error.status_code >= 500: raise try: self._complete_managed_file_operation( operation_id, succeeded=False, ) except Exception as completion_error: logger.error( 'Managed file failure reconciliation failed: %s', type(completion_error).__name__, ) _managed_file_error(exc) try: return self._complete_managed_file_operation( operation_id, succeeded=True, mutation=mutation, ) except RuntimeOperationIdentityConflictError as exc: raise AdminAPIError(409, 'managed file completion conflicted') from exc except Exception as exc: raise AdminAPIError(503, 'managed file completion is pending') from exc finally: mutation = None def _managed_file_mutation( self, traversal, *, action, root_id, relative_path, actor, operation_id, payload=None, expected_sha256=None, ): payload_box = [payload] if payload is not None else [] payload = None execution_db = None execution_acquired = False failure = None try: execution_db = self.db_factory( db_url=self.db_url, initialize=False, ) if not execution_db.enabled: raise AdminAPIError(503, 'managed files are unavailable') try: execution_db.acquire_runtime_managed_file_execution(operation_id) except Exception as exc: raise AdminAPIError(503, 'managed files are unavailable') from exc execution_acquired = True with self._managed_file_operation_lock: return self._managed_file_mutation_locked_payload( traversal, action=action, root_id=root_id, relative_path=relative_path, actor=actor, operation_id=operation_id, payload_box=payload_box, expected_sha256=expected_sha256, ) except BaseException as exc: failure = exc raise finally: release_error = None try: if execution_acquired: try: execution_db.release_runtime_managed_file_execution( operation_id, ) except BaseException as exc: release_error = exc finally: try: if execution_db is not None: execution_db.close() except BaseException as exc: if ( release_error is None or isinstance(release_error, Exception) and not isinstance(exc, Exception) ): release_error = exc finally: payload_box.clear() payload = payload_box = execution_db = None if release_error is not None: if failure is None: if not isinstance(release_error, Exception): raise release_error raise AdminAPIError( 503, 'managed file completion is pending', ) from release_error logger.error( 'Managed file execution lock release failed: %s', type(release_error).__name__, ) def create_managed_file( self, traversal, root_id, relative_path, payload, actor, operation_id, ): try: return self._managed_file_mutation( traversal, action='files.create', root_id=root_id, relative_path=relative_path, payload=payload, actor=actor, operation_id=operation_id, ) finally: payload = None def replace_managed_file( self, traversal, root_id, relative_path, payload, expected_sha256, actor, operation_id, ): try: return self._managed_file_mutation( traversal, action='files.replace', root_id=root_id, relative_path=relative_path, payload=payload, expected_sha256=expected_sha256, actor=actor, operation_id=operation_id, ) finally: payload = None def delete_managed_file( self, traversal, root_id, relative_path, expected_sha256, actor, operation_id, ): return self._managed_file_mutation( traversal, action='files.delete', root_id=root_id, relative_path=relative_path, expected_sha256=expected_sha256, actor=actor, operation_id=operation_id, ) def _complete_worker_admin_operation( self, operation_id, *, succeeded, affected_count=None, ): last_error = None for attempt in range(3): try: return self._call( 'complete_runtime_worker_admin_operation', operation_id, succeeded=succeeded, affected_count=affected_count, ) except RuntimeOperationIdentityConflictError: raise except Exception as exc: last_error = exc if attempt < 2: time.sleep(0.05) raise last_error def _worker_admin_mutation( self, *, action, target_ref, parameters, actor, operation_id, callback, success, affected_count=None, ): request_bytes = json.dumps( {'action': action, 'target_ref': target_ref, 'parameters': parameters}, sort_keys=True, separators=(',', ':'), ensure_ascii=True, ).encode('ascii') request_sha256 = hashlib.sha256(request_bytes).hexdigest() try: created = self._call( 'create_runtime_worker_admin_operation', operation_id=operation_id, actor=actor, action=action, target_ref=target_ref, request_sha256=request_sha256, ) except RuntimeOperationIdentityConflictError as exc: raise AdminAPIError(409, 'operation ID is already bound to another request') from exc if created.get('replayed') is True: if created.get('status') == 'succeeded': raise AdminAPIError(409, 'worker administration action already completed') if created.get('status') == 'running': raise AdminAPIError(409, 'worker administration action is already in progress') raise AdminAPIError(409, 'worker administration action already failed') try: result = callback() applied = success(result) except Exception: try: self._complete_worker_admin_operation(operation_id, succeeded=False) except Exception as completion_error: logger.error( 'Worker admin failure reconciliation failed: %s', type(completion_error).__name__, ) raise try: self._complete_worker_admin_operation( operation_id, succeeded=applied, affected_count=( affected_count(result) if applied and affected_count else 1 if applied else None ), ) except RuntimeOperationIdentityConflictError as exc: raise AdminAPIError(409, 'worker administration completion conflicted') from exc return result def create_user(self, user_key, cap, actor, operation_id): user_key = _validate_key(user_key, 'user key') cap = _validate_cap(cap) return self._worker_admin_mutation( action='workers.user.create', target_ref=user_key, parameters={'active_assignment_cap': cap}, actor=actor, operation_id=operation_id, callback=lambda: self._call('create_remote_worker_user', user_key, cap), success=bool, ) def set_user_cap(self, user_key, cap, actor, operation_id): user_key = _validate_key(user_key, 'user key') cap = _validate_cap(cap) return self._worker_admin_mutation( action='workers.user.set-cap', target_ref=user_key, parameters={'active_assignment_cap': cap}, actor=actor, operation_id=operation_id, callback=lambda: self._call('set_remote_worker_user_cap', user_key, cap), success=bool, ) def set_user_disabled(self, user_key, disabled, actor, operation_id): user_key = _validate_key(user_key, 'user key') disabled = disabled is True return self._worker_admin_mutation( action='workers.user.disable' if disabled else 'workers.user.enable', target_ref=user_key, parameters={}, actor=actor, operation_id=operation_id, callback=lambda: self._call( 'set_remote_worker_user_disabled', user_key, disabled, ), success=bool, ) def issue_device( self, user_key, device_key, actor, operation_id, *, rotate=False, ): user_key = _validate_key(user_key, 'user key') device_key = _validate_key(device_key, 'device key') rotate = rotate is True token_holder = {} def issue(): token = secrets.token_urlsafe(48) token_holder['token'] = token token_sha256 = hashlib.sha256(token.encode('ascii')).hexdigest() return self._call( 'issue_remote_worker_device', user_key, device_key, token_sha256, rotate=rotate, ) result = self._worker_admin_mutation( action='workers.device.rotate' if rotate else 'workers.device.issue', target_ref=device_key, parameters={'user_key': user_key}, actor=actor, operation_id=operation_id, callback=issue, success=bool, ) return result, token_holder.get('token') if result else None def set_device_revoked(self, device_key, revoked, actor, operation_id): device_key = _validate_key(device_key, 'device key') revoked = revoked is True return self._worker_admin_mutation( action='workers.device.revoke' if revoked else 'workers.device.unrevoke', target_ref=device_key, parameters={}, actor=actor, operation_id=operation_id, callback=lambda: self._call( 'set_remote_worker_device_revoked', device_key, revoked, ), success=bool, ) def requeue(self, queue_ids, actor, operation_id): if not queue_ids or len(queue_ids) > self.requeue_limit: raise AdminAPIError(400, 'queue ID selection exceeds its bound') selection_sha256 = hashlib.sha256( ','.join(str(value) for value in queue_ids).encode('ascii') ).hexdigest() return self._worker_admin_mutation( action='workers.queue.requeue', target_ref='deferred-queue', parameters={ 'selection_sha256': selection_sha256, 'item_count': len(queue_ids), }, actor=actor, operation_id=operation_id, callback=lambda: self._call( 'admin_requeue_deferred_targets', queue_ids, max_items=self.requeue_limit, ), success=lambda result: type(result) is int and result >= 0, affected_count=lambda result: result, ) def discard_source_queue(self, source, actor, operation_id): source = _validate_key(source, 'source') if source not in CORE_PRODUCERS: raise AdminAPIError(400, 'source is not a managed producer') return self._worker_admin_mutation( action='workers.queue.discard-source', target_ref=source, parameters={'statuses': ['pending', 'deferred', 'cold']}, actor=actor, operation_id=operation_id, callback=lambda: self._call( 'admin_discard_queued_source', source, ), success=lambda result: type(result) is int and result >= 0, affected_count=lambda result: result, ) def _secure_response(body, status_code=200, media_type='text/plain'): response = Response(body, status_code=status_code, media_type=media_type) response.headers.update(SECURITY_HEADERS) return response def _secure_json_response(value, *, filename=None): body = json.dumps( value, ensure_ascii=True, sort_keys=True, separators=(',', ':'), allow_nan=False, ) headers = dict(SECURITY_HEADERS) if filename: headers['Content-Disposition'] = ( "attachment; filename*=UTF-8''" + quote(filename, safe='') ) return Response(body, media_type='application/json', headers=headers) def _diagnostic_json_response(envelope, diagnostic_uid): body = envelope['canonical_json'] headers = dict(SECURITY_HEADERS) headers.update({ 'Content-Disposition': ( "attachment; filename*=UTF-8''" + quote(f'{diagnostic_uid}.json', safe='') ), 'ETag': f'"{envelope["sha256"]}"', 'Content-Length': str(len(body.encode('ascii'))), }) return Response(body, media_type='application/json', headers=headers) def _single_header(request, name): values = request.headers.getlist(name) return values[0] if len(values) == 1 else None def _trusted_operator(request, service): marker = _single_header(request, EDGE_MARKER_HEADER) operator = _single_header(request, OPERATOR_HEADER) if marker is None or operator is None or not OPERATOR_RE.fullmatch(operator): return None try: if not hmac.compare_digest(marker, service.edge_marker): return None except TypeError: return None return operator async def _form_fields(request, service, expected, *, body_limit=None): limit = service.max_body_bytes if body_limit is None else int(body_limit) origin = _single_header(request, 'origin') if origin is None or not hmac.compare_digest(origin, service.origin): raise AdminAPIError(403, 'mutation authorization failed') content_type = _single_header(request, 'content-type') if content_type is None or content_type.split(';', 1)[0].strip().lower() != ( 'application/x-www-form-urlencoded' ): raise AdminAPIError(415, 'urlencoded form body is required') content_length = _single_header(request, 'content-length') if content_length is not None: try: declared_length = int(content_length) except (TypeError, ValueError, OverflowError): raise AdminAPIError(400, 'Content-Length is invalid') from None if declared_length < 0 or declared_length > limit: raise AdminAPIError(413, 'form body exceeds its byte bound') body = bytearray() chunk = encoded = pairs = fields = key = value = result = None try: async for chunk in request.stream(): if len(body) + len(chunk) > limit: raise AdminAPIError(413, 'form body exceeds its byte bound') body.extend(chunk) encoded = bytes(body).decode('ascii', errors='strict') pairs = parse_qsl( encoded, keep_blank_values=True, strict_parsing=True, encoding='utf-8', errors='strict', max_num_fields=10, ) fields = {} for key, value in pairs: if key in fields: raise AdminAPIError(400, 'form fields must not be repeated') fields[key] = value if set(fields) != set(expected): raise AdminAPIError(400, 'form shape is invalid') if not hmac.compare_digest(fields['csrf_token'], service.csrf_token): raise AdminAPIError(403, 'mutation authorization failed') result = fields fields = None return result except (UnicodeDecodeError, UnicodeEncodeError, ValueError) as exc: raise AdminAPIError(400, 'form body is invalid') from exc finally: body.clear() if isinstance(fields, dict): fields.clear() chunk = encoded = pairs = fields = key = value = result = None def _parse_queue_ids(value, limit): value = str(value or '') if not re.fullmatch(r'[1-9][0-9]*(?:,[1-9][0-9]*)*', value): raise AdminAPIError(400, 'queue IDs must be comma-separated positive integers') parts = value.split(',') if any(len(item) > 19 for item in parts): raise AdminAPIError(400, 'queue ID selection exceeds its bound') ids = [int(item) for item in parts] if ( any(item > 9223372036854775807 for item in ids) or len(ids) > limit or len(ids) != len(set(ids)) ): raise AdminAPIError(400, 'queue ID selection exceeds its bound') return ids def _parse_revision(value): value = str(value or '') if not re.fullmatch(r'0|[1-9][0-9]{0,18}', value): raise AdminAPIError(400, 'control revision is invalid') revision = int(value) if revision > 9223372036854775806: raise AdminAPIError(400, 'control revision is invalid') return revision def _parse_producer_interval(value): value = str(value or '') if not re.fullmatch(r'[1-9][0-9]{0,7}', value): raise AdminAPIError(400, 'producer interval is invalid') interval = int(value) if interval > MAX_PRODUCER_INTERVAL_SECONDS: raise AdminAPIError(400, 'producer interval is invalid') return interval def _parse_managed_source_delay(value, label): value = str(value or '') if not re.fullmatch(r'[1-9][0-9]{0,7}', value): raise AdminAPIError(400, f'{label} is invalid') delay = int(value) if delay > MAX_MANAGED_SOURCE_DELAY_SECONDS: raise AdminAPIError(400, f'{label} is invalid') return delay def _parse_operation_id(value): value = str(value or '') try: parsed = uuid.UUID(value) except (AttributeError, TypeError, ValueError) as exc: raise AdminAPIError(400, 'operation ID is invalid') from exc if parsed.int == 0 or str(parsed) != value: raise AdminAPIError(400, 'operation ID is invalid') return value def _parse_audit_cursor(value): value = str(value or '') if not re.fullmatch(r'[1-9][0-9]{0,18}', value): raise AdminAPIError(400, 'audit cursor is invalid') cursor = int(value) if cursor > 9223372036854775807: raise AdminAPIError(400, 'audit cursor is invalid') return cursor def _operation_cursor(updated_at, operation_id): payload = json.dumps( [updated_at, operation_id], ensure_ascii=True, separators=(',', ':'), ).encode('ascii') return base64.urlsafe_b64encode(payload).rstrip(b'=').decode('ascii') def _parse_operation_cursor(value): value = str(value or '') if not re.fullmatch(r'[A-Za-z0-9_-]{1,256}', value): raise AdminAPIError(400, 'operation cursor is invalid') try: payload = base64.urlsafe_b64decode(value + '=' * (-len(value) % 4)) parts = json.loads(payload.decode('ascii')) if ( not isinstance(parts, list) or len(parts) != 2 or any(type(part) is not str for part in parts) ): raise ValueError updated_at, operation_id = parts parsed_at = datetime.fromisoformat(updated_at.replace('Z', '+00:00')) if parsed_at.tzinfo is None or not updated_at: raise ValueError _parse_operation_id(operation_id) if _operation_cursor(updated_at, operation_id) != value: raise ValueError return updated_at, operation_id except (AdminAPIError, binascii.Error, UnicodeDecodeError, ValueError) as exc: raise AdminAPIError(400, 'operation cursor is invalid') from exc def _query_fields(request, allowed): fields = {} for key, value in request.query_params.multi_items(): if key not in allowed or key in fields: raise AdminAPIError(400, 'query shape is invalid') fields[key] = value return fields def _parse_worker_filters(query): selected = {key: str(value or '') for key, value in query.items()} filters = {} patterns = { 'source': r'[a-z0-9][a-z0-9_.-]{0,63}', 'phase': r'[a-z][a-z0-9_]{0,63}', 'category': r'[a-z][a-z0-9_]{0,127}', 'code': r'[A-Za-z0-9][A-Za-z0-9._:-]{0,255}', } for name, pattern in patterns.items(): if selected.get(name): if re.fullmatch(pattern, selected[name]) is None: raise AdminAPIError(400, 'worker filter is invalid') filters[name] = selected[name] if selected.get('worker'): worker = selected['worker'].strip() if not 1 <= len(worker) <= 128 or '\x00' in worker: raise AdminAPIError(400, 'worker filter is invalid') selected['worker'] = worker filters['worker'] = worker if selected.get('assignment'): if selected['assignment'] not in { 'accepted', 'prebundle_failed', 'expired', 'unfinished', }: raise AdminAPIError(400, 'worker filter is invalid') filters['assignment_outcome'] = selected['assignment'] if selected.get('scan'): if selected['scan'] not in { 'clean', 'found', 'degraded', 'error', 'skipped', 'unavailable', }: raise AdminAPIError(400, 'worker filter is invalid') filters['scan_outcome'] = selected['scan'] if selected.get('retryable'): if selected['retryable'] not in {'true', 'false'}: raise AdminAPIError(400, 'worker filter is invalid') filters['retryable'] = selected['retryable'] == 'true' for name in ('diagnostic_offset', 'metric_offset'): offset = selected.get(name) or '0' if re.fullmatch(r'0|[1-9][0-9]{0,18}', offset) is None: raise AdminAPIError(400, 'worker filter is invalid') if int(offset) > 9223372036854775807: raise AdminAPIError(400, 'worker filter is invalid') selected[name] = offset window = selected.get('window') or '30d' selected['window'] = window if window not in WORKER_FILTER_WINDOWS: raise AdminAPIError(400, 'worker filter is invalid') delta = WORKER_FILTER_WINDOWS[window] if delta is not None: filters['since'] = ( datetime.now(timezone.utc) - delta ).isoformat(timespec='seconds') page_limit = selected.get('limit') or str(DEFAULT_WORKER_PAGE_LIMIT) if page_limit not in {str(value) for value in WORKER_PAGE_LIMITS}: raise AdminAPIError(400, 'worker filter is invalid') selected['limit'] = page_limit details = selected.get('details') or 'assignments' if details not in {'assignments', 'diagnostics', 'metrics', 'all'}: raise AdminAPIError(400, 'worker filter is invalid') selected['details'] = details return filters, selected def _parse_assignment_id(value): value = str(value or '') if re.fullmatch(r'[1-9][0-9]{0,18}', value) is None: raise AdminAPIError(404, 'worker assignment was not found') result = int(value) if result > 9223372036854775807: raise AdminAPIError(404, 'worker assignment was not found') return result def _managed_file_error(exc): status = { 'invalid_path': 400, 'invalid_hash': 400, 'invalid_content': 400, 'operation_not_allowed': 403, 'unknown_root': 404, 'not_found': 404, 'unsafe_target': 404, 'hash_conflict': 409, 'concurrent_change': 409, 'limit_exceeded': 413, 'root_unavailable': 503, 'filesystem_unavailable': 503, 'closed': 503, 'durability_uncertain': 503, 'download_busy': 503, }.get(getattr(exc, 'category', None), 503) messages = { 400: 'managed file request is invalid', 403: 'managed file operation is not allowed', 404: 'managed file target was not found', 409: 'managed file state changed concurrently', 413: 'managed file limit was exceeded', 503: 'managed files are unavailable', } raise AdminAPIError(status, messages[status]) from exc def _parse_managed_file_hash(value): value = str(value or '') if not SHA256_RE.fullmatch(value): raise AdminAPIError(400, 'managed file revision is invalid') return value def _parse_managed_file_content(value): raw = decoded = None try: raw = str(value or '').encode('ascii', errors='strict') decoded = base64.b64decode(raw, altchars=b'-_', validate=True) if base64.urlsafe_b64encode(decoded) != raw: raise ValueError('noncanonical') return decoded except (binascii.Error, UnicodeEncodeError, ValueError) as exc: raise AdminAPIError(400, 'managed file content is invalid') from exc finally: value = raw = decoded = None def _managed_file_listing_location(root_id, relative_path): parent = relative_path.rpartition('/')[0] or None query = {'root_id': root_id} if parent is not None: query['relative_path'] = parent return '../files?' + urlencode(query) def _parse_document_hashes(fields, *, require_config=True, require_secrets=True): names = ('active_config', 'active_secrets') if require_config: names += ('candidate_config',) if require_secrets: names += ('candidate_secrets',) hashes = {} for name in names: value = str(fields.get(f'expected_{name}_sha256') or '') if not SHA256_RE.fullmatch(value): raise AdminAPIError(400, 'runtime document revision is invalid') hashes[name] = value return hashes def _form(action, title, csrf_token, fields, relative_root, hidden=(), *, values=None): values = dict(values or {}) controls = [] for field in fields: name, label, input_type = field[:3] bounds = '' if len(field) == 5: bounds = ( f' min="{html.escape(str(field[3]))}"' f' max="{html.escape(str(field[4]))}" step="1"' ) value = values.get(name) value_attribute = ( f' value="{html.escape(str(value))}"' if value is not None else '' ) controls.append( f'' ) return ( f'
' f'

{html.escape(title)}

' f'' + ''.join( f'' for name, value in hidden ) + ''.join(controls) + f'
' ) def _table(columns, rows): header = ''.join(f'{html.escape(label)}' for key, label in columns) rendered_rows = [] for row in rows: rendered_rows.append('' + ''.join( f'{html.escape(str(row.get(key) if row.get(key) is not None else ""))}' for key, _label in columns ) + '') return f'{header}{"".join(rendered_rows)}
' def _admin_navigation(active, relative_root): links = ( ('workers', relative_root + '/', 'Workers / Dispatch'), ('overview', relative_root + '/overview', 'Overview'), ('search', relative_root + '/search', 'Search'), ('supervisor', relative_root + '/supervisor', 'Supervisor'), ('logs', relative_root + '/logs', 'Logs'), ('config', relative_root + '/config', 'Config'), ('secrets', relative_root + '/secrets', 'Secrets'), ('files', relative_root + '/files', 'Files'), ('operations', relative_root + '/operations', 'Operations'), ('audit', relative_root + '/audit', 'Audit'), ) return '' def _page_shell(title, subtitle, active, content, relative_root, script=None): script_tag = ( f'' if script else '' ) return f''' {html.escape(title)}{script_tag}

{html.escape(title)}

{html.escape(subtitle)}

{_admin_navigation(active, relative_root)}
{content}
''' def _identity_json(value): if value is None: return '' return json.dumps( value, ensure_ascii=True, sort_keys=True, separators=(',', ':'), allow_nan=False, ) def _operation_link(operation_id, relative_root): operation_id = html.escape(str(operation_id)) href = html.escape(f'{relative_root}/operations/{operation_id}') return f'{operation_id}' def _render_operations_table(operations, relative_root): headers = ( 'Operation', 'Actor', 'Action', 'Target', 'Status', 'Category', 'Requested', 'Completed', 'Updated', ) rows = [] for operation in operations: cells = [ _operation_link(operation.get('operation_id', ''), relative_root), ] cells.extend( html.escape(str(operation.get(key) or '')) for key in ( 'actor', 'action', 'target_ref', 'status', 'safe_category', 'requested_at', 'completed_at', 'updated_at', ) ) rows.append('' + ''.join(f'{cell}' for cell in cells) + '') header = ''.join(f'{html.escape(label)}' for label in headers) return f'{header}{"".join(rows)}
' def _render_operations_page(page, relative_root='.'): operations = list(page.get('operations') or []) pagination = f'Newest operations' next_before = page.get('next_before') if next_before is not None: pagination += ( ' ' ) content = ( '

Recent durable operations

' '

Newest first. Open an operation to inspect its persisted status after a runtime restart.

' + _render_operations_table(operations, relative_root) + f'

{pagination}

' ) return _page_shell( 'Operations', 'Bounded recent operation status from durable storage.', 'operations', content, relative_root, ) def _render_operation_page(operation, relative_root='..'): fields = ( ('operation_id', 'Operation'), ('actor', 'Actor'), ('action', 'Action'), ('target_kind', 'Target kind'), ('target_ref', 'Target'), ('status', 'Status'), ('safe_category', 'Category'), ('safe_detail', 'Detail'), ('expected_revision', 'Expected revision'), ('resulting_revision', 'Resulting revision'), ('agent_state', 'Agent state'), ('agent_result_sha256', 'Agent result SHA-256'), ('requested_at', 'Requested'), ('started_at', 'Started'), ('completed_at', 'Completed'), ('agent_reconciled_at', 'Agent reconciled'), ('updated_at', 'Updated'), ) rows = ({'field': label, 'value': operation.get(key)} for key, label in fields) expected = html.escape(_identity_json(operation.get('expected_identity'))) resulting = html.escape(_identity_json(operation.get('resulting_identity'))) content = ( '

Durable status

' + _table((('field', 'Field'), ('value', 'Value')), rows) + '

Expected identity

' + f'
{expected}
' + '

Resulting identity

' + f'
{resulting}
' ) return _page_shell( 'Operation status', 'Persisted state remains queryable across runtime restarts.', 'operations', content, relative_root, ) def _render_audit_page(page, relative_root='.'): events = list(page.get('events') or []) headers = ( 'Event', 'Operation', 'Actor', 'Action', 'Target', 'Time', 'Result', 'Category', 'Before identity', 'After identity', 'Before bytes', 'After bytes', ) rows = [] for event in events: target = f'{event.get("target_kind") or ""}:{event.get("target_ref") or ""}' cells = [ html.escape(str(event.get('id') or '')), _operation_link(event.get('operation_id', ''), relative_root), html.escape(str(event.get('actor') or '')), html.escape(str(event.get('action') or '')), html.escape(target), html.escape(str(event.get('created_at') or '')), html.escape(str(event.get('result') or '')), html.escape(str(event.get('safe_category') or '')), '' + html.escape(_identity_json(event.get('before_identity'))) + '', '' + html.escape(_identity_json(event.get('after_identity'))) + '', html.escape(str(event.get('before_bytes') if event.get('before_bytes') is not None else '')), html.escape(str(event.get('after_bytes') if event.get('after_bytes') is not None else '')), ] rows.append('' + ''.join(f'{cell}' for cell in cells) + '') header = ''.join(f'{html.escape(label)}' for label in headers) table = f'{header}{"".join(rows)}
' next_before = page.get('next_before_event_id') pagination = f'Newest events' if next_before is not None: pagination += ( ' ' ) content = ( '

Append-only events

' '

Newest first. Each page is bounded and uses a stable event cursor.

' + table + f'

{pagination}

' ) return _page_shell( 'Audit', 'Content-free accepted and terminal operation evidence.', 'audit', content, relative_root, ) def _managed_file_query(root_id, relative_path=None): values = {'root_id': root_id} if relative_path is not None: values['relative_path'] = relative_path return urlencode(values) def _render_managed_file_forms(service, root, relative_path): forms = [] root_id = html.escape(root.root_id) def hidden(operation_id): return ( f'' f'' f'' ) path_value = '' if relative_path is None else relative_path + '/' path_input = ( '' ) expected_input = ( '' ) content_input = ( '' ) if root.permissions.allow_create_replace: forms.append( '
' '

Create file

' + hidden(str(uuid.uuid4())) + path_input + content_input + '
' ) forms.append( '
' '

Replace file

' + hidden(str(uuid.uuid4())) + path_input + expected_input + content_input + '
' ) if root.permissions.allow_delete: forms.append( '
' '

Delete file

' + hidden(str(uuid.uuid4())) + path_input + expected_input + '
' ) if not forms: return '

This root is read-only.

' return ( '

Mutation content is canonical URL-safe Base64 and is bounded by the ' f'{service.max_body_bytes}-byte admin form limit.

' '
' + ''.join(forms) + '
' ) def _render_files_page(service, *, root=None, relative_path=None, listing=None): root_rows = [] for configured in service.managed_file_roots.roots: root_rows.append({ 'root': configured.root_id, 'list': configured.permissions.allow_list, 'read': configured.permissions.allow_read, 'create_replace': configured.permissions.allow_create_replace, 'delete': configured.permissions.allow_delete, 'path_bytes': configured.limits.max_relative_path_bytes, 'entries': configured.limits.max_listing_entries, 'file_bytes': configured.limits.max_file_bytes, }) roots = _table(( ('root', 'Logical root'), ('list', 'List'), ('read', 'Read'), ('create_replace', 'Create/replace'), ('delete', 'Delete'), ('path_bytes', 'Path bytes'), ('entries', 'Listing entries'), ('file_bytes', 'File bytes'), ), root_rows) root_links = ''.join( '
  • ' + html.escape(configured.root_id) + '
  • ' for configured in service.managed_file_roots.roots ) content = ( '

    Logical roots

    ' '

    Only configured logical roots are exposed. Host paths are never accepted.

    ' + roots + '
    ' ) if root is not None: location = root.root_id + (f'/{relative_path}' if relative_path else '') entries = [] if listing is not None: for entry in listing.entries: child_path = ( f'{relative_path}/{entry.name}' if relative_path else entry.name ) if entry.kind == 'directory': name = ( '' + html.escape(entry.name) + '/' ) elif root.permissions.allow_read: name = ( '' + html.escape(entry.name) + '' ) else: name = html.escape(entry.name) entries.append( '' + name + '' + html.escape(entry.kind) + '' + html.escape('' if entry.byte_count is None else str(entry.byte_count)) + '' ) listing_block = '

    Listing is not permitted for this root.

    ' if listing is not None: listing_block = ( '' '' '' + ''.join(entries) + '
    NameKindBytes
    ' ) parent = '' if relative_path: parent_path = relative_path.rpartition('/')[0] or None parent = ( '

    Parent directory

    ' ) content += ( '

    ' + html.escape(location) + '

    ' + parent + listing_block + '

    Typed mutations

    ' + _render_managed_file_forms(service, root, relative_path) + '
    ' ) return _page_shell( 'Files', 'Bounded descriptor-safe access by logical root.', 'files', content, '.', ) def _filter_select(name, label, values, selected, *, include_empty=True): options = [''] if include_empty else [] for value, title in values: current = ' selected' if selected.get(name) == value else '' options.append( f'' ) return ( f'' ) def _render_worker_filters(selected, relative_root): text_fields = ( ('source', 'Source'), ('worker', 'Worker / device'), ('phase', 'Phase'), ('category', 'Category'), ('code', 'Stable code'), ) controls = ''.join( f'' for name, label in text_fields ) controls += _filter_select('assignment', 'Assignment outcome', ( ('accepted', 'Accepted'), ('prebundle_failed', 'Prebundle failed'), ('expired', 'Expired'), ('unfinished', 'Unfinished'), ), selected) controls += _filter_select('scan', 'Scan outcome', ( ('clean', 'Clean'), ('found', 'Found'), ('degraded', 'Degraded'), ('error', 'Error'), ('skipped', 'Skipped'), ('unavailable', 'Unavailable'), ), selected) controls += _filter_select('retryable', 'Retryability', ( ('true', 'Retryable'), ('false', 'Not retryable'), ), selected) controls += _filter_select('window', 'Time window', ( ('24h', 'Last 24 hours'), ('7d', 'Last 7 days'), ('30d', 'Last 30 days'), ('90d', 'Last 90 days'), ('all', 'All'), ), selected, include_empty=False) controls += _filter_select('limit', 'Rows per heavy section', ( ('25', '25 rows'), ('50', '50 rows'), ('100', '100 rows'), ), selected, include_empty=False) controls += _filter_select('details', 'Load detail sections', ( ('assignments', 'Assignments only (fast)'), ('diagnostics', 'Assignments + diagnostics'), ('metrics', 'Assignments + durations'), ('all', 'All detail sections'), ), selected, include_empty=False) return ( f'
    {controls}
    ' '
    ' f'Clear filters' '
    ' ) def _render_assignment_rows(assignments, relative_root): headers = ( 'Assignment', 'Source / target', 'Worker', 'Assignment outcome', 'Scan outcome', 'Diagnostics', 'Phase / progress', 'Deadlines', 'Slot / cap', 'Package identity', 'Ingestion / projection', ) rows = [] for row in assignments: reservation_id = int(row['reservation_id']) terminal = row.get('finished_at') is not None age_authority = 'assignment resolution' if terminal else 'current time' deadline_label = 'at resolution' if terminal else 'remaining' diagnostic_count = int(row.get('diagnostic_count') or 0) protocol2 = str(row.get('protocol_version') or '') == '2' if diagnostic_count: diagnostic_state = ( f'current stored: {diagnostic_count}: ' f'{row.get("primary_diagnostic") or "none"}' ) elif protocol2: diagnostic_state = 'current protocol-2: 0 diagnostics observed' elif row.get('diagnostic_projection_version') == 1: diagnostic_state = 'current projection: 0 diagnostics observed' else: diagnostic_state = 'legacy/unavailable' if row.get('active_phase'): phase = ( f"latest persisted phase: {row['active_phase']}" if terminal else f"current phase: {row['active_phase']}" ) phase_age = ( f"{row['phase_age_seconds']}s" if row.get('phase_age_seconds') is not None else 'unavailable' ) progress_age = ( f"{row['last_progress_age_seconds']}s" if row.get('last_progress_age_seconds') is not None else 'unavailable' ) elif protocol2: phase = 'unavailable (no persisted progress)' phase_age = progress_age = 'unavailable (no persisted progress)' else: phase = phase_age = progress_age = 'legacy/unavailable' warning = row.get('scan_warning_summary') warning_class = row.get('scan_warning_class') if warning: warning_prefix = f'persisted warning [{warning_class}]' if warning_class else 'persisted warning' scan_outcome = f"{row.get('scan_outcome') or 'unavailable'}; {warning_prefix}: {warning}" elif str(row.get('scan_outcome') or '').lower() == 'degraded': scan_outcome = 'degraded; persisted warning detail unavailable' else: scan_outcome = str(row.get('scan_outcome') or 'unavailable') package = '
    '.join(html.escape(str(value)) for value in ( f"protocol {row.get('protocol_version') or 'unavailable'} / bundle {row.get('bundle_format_version') or 'unavailable'}", f"platform {row.get('platform_tag') or 'unavailable'}", f"code {row.get('code_manifest_sha256') or 'unavailable'}", f"detector {row.get('detector_policy_sha256') or 'unavailable'}", )) deadlines = '
    '.join(html.escape(str(value)) for value in ( f"scan: {row.get('scan_deadline_at') or 'legacy/unavailable'} ({row.get('scan_remaining_seconds') if row.get('scan_remaining_seconds') is not None else 'unavailable'}s {deadline_label})", ( f"result upload: {row['remote_result_upload_body_timeout_seconds']}s " '(persisted at assignment issuance)' if row.get('remote_result_upload_body_timeout_seconds') is not None else 'result upload: legacy/unavailable (not persisted)' ), f"assignment: {row.get('assignment_deadline_at') or 'unavailable'} ({row.get('assignment_remaining_seconds') if row.get('assignment_remaining_seconds') is not None else 'unavailable'}s {deadline_label})", )) cells = ( ('assignment', f'{reservation_id}'), ('source_target', html.escape(f"{row.get('source') or ''}: {row.get('target') or ''}")), ('worker', html.escape(f"{row.get('user_key') or ''} / {row.get('device_key') or ''}")), ('assignment_outcome', html.escape(str(row.get('assignment_outcome') or 'unavailable'))), ('scan_outcome', html.escape(scan_outcome)), ('diagnostics', html.escape(diagnostic_state)), ('progress', html.escape( f'{phase}; phase age {phase_age}; progress age {progress_age}; ' f'age authority: {age_authority}' )), ('deadlines', deadlines), ('slot_cap', html.escape( f"{row.get('slot_id') if row.get('slot_id') is not None else 'unavailable'} / {row.get('active_assignment_cap') or 0}" )), ('package', package), ('pipeline', html.escape( f"{row.get('ingestion_state') or 'unavailable'} / {row.get('projection_state') or 'unavailable'}" )), ) rows.append( f'' + ''.join( f'{cell}' for field, cell in cells ) + '' ) header = ''.join(f'{html.escape(value)}' for value in headers) return f'{header}{"".join(rows)}
    ' def _render_diagnostic_groups(groups, relative_root): if not groups: return '

    No diagnostic occurrences match the selected filters.

    ' rendered = [] for group in groups: occurrences = ''.join( '
  • ' f'' f'assignment {item["reservation_id"]}: ' f'{html.escape(item["occurred_at"])}; {html.escape(item["assignment_outcome"])} / ' f'{html.escape(item["scan_outcome"])}; {html.escape(item["phase"])}; ' f'{html.escape(item["category"] + "/" + item["code"])}; ' f'retryable={html.escape(str(item["retryable"]).lower())}' '
  • ' for item in group['occurrences'] ) rendered.append( '
    ' f'

    {html.escape(group["fingerprint"])}

    ' f'

    {group["count"]} total filtered occurrences across ' f'{group["affected_assignment_count"]} assignments; ' f'{group["page_occurrence_count"]} occurrences on this page. ' f'Page assignments: ' f'{html.escape(", ".join(str(value) for value in group["affected_assignments"]))}.

    ' f'
      {occurrences}
    ' ) return ''.join(rendered) def _render_diagnostic_group_status(snapshot, selected, relative_root): matched = int(snapshot.get('matched_occurrence_count') or 0) page_count = int(snapshot.get('page_occurrence_count') or 0) offset = int(snapshot.get('occurrence_offset') or 0) limit = int(snapshot.get('occurrence_limit') or 0) status = ( f'

    Matched occurrences: {matched}. Showing {page_count} occurrences ' f'at deterministic occurrence offset {offset} with page limit {limit}. ' 'Fingerprint counts and affected-assignment counts cover the full filtered set.

    ' ) def link(label, target_offset): query = { key: value for key, value in selected.items() if value and key != 'diagnostic_offset' } query['diagnostic_offset'] = str(target_offset) return ( f'' f'{html.escape(label)}' ) links = [] previous_offset = snapshot.get('previous_occurrence_offset') if previous_offset is not None: links.append(link('Previous diagnostic occurrences', int(previous_offset))) next_offset = snapshot.get('next_occurrence_offset') if next_offset is not None: links.append(link('Next diagnostic occurrences', int(next_offset))) if links: status += '
    ' + ''.join(links) + '
    ' return status def _render_duration_metrics(metrics): rows = [] for metric in metrics: sufficient = bool(metric.get('sufficient')) percentile = lambda name: ( f"{metric[name]:.3f}" if sufficient else 'insufficient' ) rows.append({ **metric, 'p50': percentile('p50_seconds'), 'p95': percentile('p95_seconds'), 'p99': percentile('p99_seconds'), 'sample_label': ( str(metric['sample_count']) if sufficient else f"{metric['sample_count']} (minimum {metric['minimum_sample_count']})" ), }) return _table(( ('source', 'Source'), ('phase', 'Phase'), ('outcome', 'Outcome'), ('p50', 'p50 seconds'), ('p95', 'p95 seconds'), ('p99', 'p99 seconds'), ('sample_label', 'Samples'), ), rows) def _render_metric_status(snapshot, selected, relative_root): total = int(snapshot.get('total_group_count') or 0) page_count = int(snapshot.get('page_group_count') or 0) offset = int(snapshot.get('metric_offset') or 0) limit = int(snapshot.get('metric_limit') or 0) status = ( f'

    Total source / phase / outcome groups: {total}. Showing ' f'{page_count} groups at offset {offset} with page limit {limit}.

    ' ) def link(label, target_offset): query = { key: value for key, value in selected.items() if value and key != 'metric_offset' } query['metric_offset'] = str(target_offset) return ( f'' f'{html.escape(label)}' ) links = [] previous_offset = snapshot.get('previous_metric_offset') if previous_offset is not None: links.append(link('Previous duration groups', int(previous_offset))) next_offset = snapshot.get('next_metric_offset') if next_offset is not None: links.append(link('Next duration groups', int(next_offset))) if links: status += '
    ' + ''.join(links) + '
    ' return status def _render_deadline_policy(policy): return ( '

    Future assignments only. Saving or applying a policy ' 'does not alter deadlines on existing assignments.

    ' + _table(( ('source', 'Source'), ('scan_deadline_seconds', 'Scan deadline seconds'), ('upload_deadline_seconds', 'Upload deadline seconds'), ('assignment_deadline_seconds', 'Assignment deadline seconds'), ('assignment_policy_source', 'Assignment value source'), ('handoff_margin_seconds', 'Handoff margin seconds'), ('required_minimum_seconds', 'Required minimum'), ('relationship', 'Validation relationship'), ('valid', 'Valid'), ), policy.get('rows') or []) ) def _render_workers_page(service, snapshot, notice='', issued_token=None, relative_root='.'): snapshot = dict(snapshot or {}) worker_snapshot = _component_value(snapshot, 'workers', dict) diagnostic_snapshot = _component_value(snapshot, 'diagnostics', dict) metric_snapshot = _component_value(snapshot, 'metrics', dict) policy = _component_value(snapshot, 'policy', dict) control = _component_value(snapshot, 'control', dict) packages = _component_value(snapshot, 'packages', dict) selected_filters = dict(snapshot.get('selected_filters') or {}) users = list((worker_snapshot or {}).get('users') or []) workers = list((worker_snapshot or {}).get('workers') or []) assignments = list((worker_snapshot or {}).get('assignments') or []) deferred = list((worker_snapshot or {}).get('deferred_queue') or []) token_block = '' if issued_token: token_block = ( '

    Device token

    ' '

    Shown once. Store it now; only its SHA-256 digest was persisted.

    ' f'{html.escape(issued_token)}
    ' ) notice_block = f'

    {html.escape(notice)}

    ' if notice else '' forms = ''.join(( _form('/users/create', 'Create user', service.csrf_token, ( ('user_key', 'User key', 'text'), ('active_assignment_cap', 'Assignment cap', 'number', 0, 10000), ), relative_root, hidden=(('operation_id', str(uuid.uuid4())),)), _form('/users/cap', 'Update user cap', service.csrf_token, ( ('user_key', 'User key', 'text'), ('active_assignment_cap', 'Assignment cap', 'number', 0, 10000), ), relative_root, hidden=(('operation_id', str(uuid.uuid4())),)), _form('/users/disable', 'Disable user', service.csrf_token, ( ('user_key', 'User key', 'text'), ), relative_root, hidden=(('operation_id', str(uuid.uuid4())),)), _form('/users/enable', 'Enable user', service.csrf_token, ( ('user_key', 'User key', 'text'), ), relative_root, hidden=(('operation_id', str(uuid.uuid4())),)), _form('/devices/issue', 'Issue device token', service.csrf_token, ( ('user_key', 'User key', 'text'), ('device_key', 'Device key', 'text'), ), relative_root, hidden=(('operation_id', str(uuid.uuid4())),)), _form('/devices/rotate', 'Rotate device token', service.csrf_token, ( ('user_key', 'User key', 'text'), ('device_key', 'Device key', 'text'), ), relative_root, hidden=(('operation_id', str(uuid.uuid4())),)), _form('/devices/revoke', 'Revoke device', service.csrf_token, ( ('device_key', 'Device key', 'text'), ), relative_root, hidden=(('operation_id', str(uuid.uuid4())),)), _form('/devices/unrevoke', 'Unrevoke device', service.csrf_token, ( ('device_key', 'Device key', 'text'), ), relative_root, hidden=(('operation_id', str(uuid.uuid4())),)), _form('/queue/requeue', 'Requeue deferred targets', service.csrf_token, ( ('queue_ids', 'Queue IDs, comma separated', 'text'), ), relative_root, hidden=(('operation_id', str(uuid.uuid4())),)), _form( '/queue/discard-source', 'Discard stale unassigned source backlog', service.csrf_token, ( ('source', 'Source', 'text'), ('confirm_source', 'Type source again to confirm', 'text'), ), relative_root, hidden=(('operation_id', str(uuid.uuid4())),), ), )) if worker_snapshot is None: worker_note = '

    Worker administration snapshot unavailable.

    ' else: worker_note = '' if control is None: dispatch_content = '

    Dispatch control unavailable.

    ' dispatch_forms = '' else: revision = control['revision'] dispatch_content = _table(( ('revision', 'Revision'), ('dispatch_paused', 'Explicit pause'), ('effective_dispatch_paused', 'Effective pause'), ('drain_state', 'Drain state'), ('live_remote_assignments', 'Live assignments'), ('precommit_result_bundles', 'Pre-commit bundles'), ('blocker_count', 'Drain blockers'), ('actor', 'Last actor'), ('updated_at', 'Updated'), ), (control,)) hidden = (('expected_revision', revision),) dispatch_forms = '
    ' + ''.join(( _form( '/dispatch/pause', 'Pause new assignments', service.csrf_token, (), relative_root, hidden=hidden + (('operation_id', str(uuid.uuid4())),), ), _form( '/dispatch/resume', 'Resume new assignments', service.csrf_token, (), relative_root, hidden=hidden + (('operation_id', str(uuid.uuid4())),), ), _form( '/dispatch/drain/start', 'Start drain', service.csrf_token, (), relative_root, hidden=hidden + (('operation_id', str(uuid.uuid4())),), ), _form( '/dispatch/drain/cancel', 'Cancel drain', service.csrf_token, (), relative_root, hidden=hidden + (('operation_id', str(uuid.uuid4())),), ), )) + '
    ' if packages is None: package_content = '

    Package compatibility unavailable.

    ' else: profile_rows = [] capability_rows = [] for profile in packages['profiles']: profile_rows.append({ 'profile_name': profile['profile_name'], 'protocol_version': profile['protocol_version'], 'bundle_format_version': profile['bundle_format_version'], 'platform_tag': profile['platform_tag'], 'code_manifest_sha256': profile['code_manifest_sha256'], 'detector_policy_sha256': profile['detector_policy_sha256'], 'sources': ', '.join(profile['sources']), }) capability_rows.extend({ 'profile_name': profile['profile_name'], **capability, } for capability in profile['capabilities']) required_rows = packages['required_capabilities'] package_content = ( '

    Trusted package profiles

    ' + _table(( ('profile_name', 'Profile'), ('protocol_version', 'Protocol'), ('bundle_format_version', 'Bundle format'), ('platform_tag', 'Platform tag'), ('sources', 'Sources'), ('code_manifest_sha256', 'Code manifest SHA-256'), ('detector_policy_sha256', 'Detector policy SHA-256'), ), profile_rows) + '

    Profile capabilities

    ' + _table(( ('profile_name', 'Profile'), ('source', 'Source'), ('platform', 'Worker platform'), ('planning_kind', 'Planning kind'), ), capability_rows) + '

    Required capabilities

    ' + _table(( ('source', 'Source'), ('platform', 'Worker platform'), ('planning_kind', 'Planning kind'), ), required_rows) ) assignment_table = _render_assignment_rows(assignments, relative_root) loaded_details = selected_filters.get('details') or 'assignments' if loaded_details not in {'diagnostics', 'all'}: diagnostic_content = ( '

    Not loaded. Select Assignments + diagnostics or All detail ' 'sections in the filter above.

    ' ) elif diagnostic_snapshot is None: diagnostic_content = '

    Diagnostic grouping unavailable.

    ' else: diagnostic_content = _render_diagnostic_group_status( diagnostic_snapshot, selected_filters, relative_root, ) + _render_diagnostic_groups( diagnostic_snapshot.get('groups') or [], relative_root, ) if loaded_details not in {'metrics', 'all'}: metrics_content = ( '

    Not loaded. Select Assignments + durations or All detail ' 'sections in the filter above.

    ' ) elif metric_snapshot is None: metrics_content = '

    Duration metrics unavailable.

    ' else: metrics_content = _render_metric_status( metric_snapshot, selected_filters, relative_root, ) + _render_duration_metrics(metric_snapshot.get('metrics') or []) policy_content = ( '

    Effective deadline policy unavailable.

    ' if policy is None else _render_deadline_policy(policy) ) for worker in workers: active = [ item for item in assignments if item.get('device_key') == worker.get('device_key') and item.get('assignment_outcome') == 'unfinished' ] worker['current_phases'] = worker.get('current_phases') or ( ', '.join(sorted({ str(item.get('active_phase') or 'legacy/unavailable') for item in active })) or 'none' ) progress_ages = [ int(item['last_progress_age_seconds']) for item in active if item.get('last_progress_age_seconds') is not None ] if worker.get('latest_progress_age_seconds') is None: worker['latest_progress_age_seconds'] = ( min(progress_ages) if progress_ages else None ) worker['activity'] = ( 'active progress' if worker.get('latest_progress_age_seconds') is not None else 'active, progress unavailable' if worker.get('active_slot_count') else 'no active assignment' ) worker_scope = ( '

    Observability scope: Filters apply independently to the compatible ' 'fields in Assignments, Repeated diagnostics, Observed durations, Effective deadline ' 'policy, and Deferred queue. Users, Workers, Dispatch, Package compatibility, and Typed ' 'operations remain global.

    ' ) content = f'''{notice_block}{token_block}

    Filters

    Every filter is independent. Diagnostic filters do not imply an assignment or scan outcome. Expensive diagnostic and duration sections are queried only when selected.

    {worker_scope}{_render_worker_filters(selected_filters, relative_root)}

    Dispatch and drain

    Pausing or draining blocks new assignment commits. Authenticated status, terminal reports, uploads, receipt replay, expiry, ingestion, projection, and maintenance remain available.

    {dispatch_content}{dispatch_forms}

    Package compatibility

    Protocol-1 workers are completion-only; new claims require a matching trusted protocol-2 package profile.

    {package_content}

    Users

    {worker_note}{_table((('user_key', 'User'), ('active_assignment_cap', 'Cap'), ('disabled', 'Disabled')), users)}

    Workers

    {_table((('device_key', 'Device'), ('user_key', 'User'), ('active_assignment_cap', 'Cap'), ('active_slot_count', 'Active slots'), ('activity', 'Observed activity'), ('current_phases', 'Current phases'), ('latest_progress_age_seconds', 'Latest progress age seconds'), ('last_contact_at', 'Last API contact'), ('known_reasons', 'Known idle / backoff reason'), ('pending_local_recovery', 'Pending local recovery'), ('active_package_identity', 'Active package identity'), ('completed_count', 'Accepted'), ('failed_count', 'Prebundle failed'), ('expired_count', 'Expired'), ('revoked', 'Revoked')), workers)}

    Assignments

    Assignment transport, scan outcome, diagnostics, progress, deadlines, slot/cap, package identity, ingestion, and projection are separate authoritative fields.

    {assignment_table}

    Repeated diagnostics

    Grouping is deterministic and every bounded occurrence remains linked below its fingerprint.

    {diagnostic_content}

    Observed durations

    Percentiles are observations only and never change policy automatically. Rows below five samples are explicitly insufficient.

    {metrics_content}

    Effective deadline policy

    {policy_content}

    Deferred queue

    {_table((('queue_id', 'Queue'), ('source', 'Source'), ('target', 'Target'), ('available_after', 'Available after')), deferred)}

    Typed operations

    Discarding source backlog suspends only never-issued pending, deferred, and cold rows. Active, previously issued, and completed targets are preserved. If discovery sees a discarded target again, the same queue row is reactivated with fresh discovery data.

    {forms}
    ''' return _page_shell( 'Workers / Dispatch', 'Authoritative records only; no liveness inference.', 'workers', content, relative_root, ) def _render_material(name, material): if not isinstance(material, dict): return f'

    {html.escape(name)}

    Not captured.

    ' fields = ({ 'encoding': material.get('encoding'), 'original_size': material.get('original_size'), 'stored_size': material.get('stored_size'), 'sha256': material.get('sha256'), 'truncated': material.get('truncated'), 'state': 'truncated transformation' if material.get('truncated') else 'complete', },) head = str(material.get('head') or '') tail = material.get('tail') stored = html.escape(head) if tail is not None: omitted = max( 0, int(material.get('original_size') or 0) - int(material.get('stored_size') or 0), ) stored += ( f'\n[... {omitted} original bytes omitted by the stored transformation ...]\n' + html.escape(str(tail)) ) return ( f'

    {html.escape(name)}

    ' + _table(( ('state', 'State'), ('encoding', 'Stored encoding'), ('original_size', 'Original bytes'), ('stored_size', 'Stored bytes'), ('sha256', 'Original SHA-256'), ('truncated', 'Truncated'), ), fields) + f'
    {stored}
    ' ) def _render_assignment_detail_page(detail, relative_root='..'): assignment = detail['assignment'] reservation_id = int(assignment['reservation_id']) overview = ({ **assignment, **detail['deadlines'], 'diagnostic_availability': detail['diagnostics']['availability'], 'historical_upload_timeout': ( f"{detail['deadlines']['upload_timeout_seconds']}s " '(persisted at assignment issuance)' if detail['deadlines'].get('upload_timeout_seconds') is not None else 'legacy/unavailable (not persisted for this assignment)' ), },) timeline = _table(( ('timestamp', 'Timestamp'), ('kind', 'Kind'), ('label', 'Event'), ('sequence', 'Sequence'), ('received_at', 'Server received'), ), detail['timeline']) durations = _table(( ('phase', 'Phase'), ('duration_seconds', 'Duration seconds'), ('outcome', 'Outcome'), ('complete', 'Complete'), ('ended_at', 'Observed through'), ('authority', 'End authority'), ), detail['durations']) progress = detail['progress'] progress_summary = _table(( ('availability', 'Availability'), ('current_phase', 'Current / latest phase'), ('phase_started_at', 'Phase started'), ('phase_age_seconds', 'Phase age seconds'), ('last_progress_at', 'Last progress'), ('last_progress_age_seconds', 'Last progress age seconds'), ('age_authority', 'Age authority'), ('total_event_count', 'Total persisted events'), ('omitted_older_event_count', 'Older events omitted by bound'), ), (progress,)) scan = _table(tuple((name, label) for name, label in ( ('available', 'Available'), ('target_scan_id', 'Target scan'), ('status', 'Status'), ('started_at', 'Started'), ('ended_at', 'Ended'), ('duration_seconds', 'Duration seconds'), ('findings_count', 'Findings'), ('verified_findings_count', 'Verified findings'), ('error_count', 'Errors'), ('skipped_reason', 'Skipped reason'), ('first_error_summary', 'First error summary'), ('warning_class', 'First persisted warning class'), ('warning_summary', 'First persisted warning'), )), (detail['scan'],)) transport = _table(( ('receipt_id', 'Receipt'), ('bundle_state', 'Bundle state'), ('bundle_ready_at', 'Bundle received'), ('bundle_committed_at', 'Ingestion committed'), ('bundle_acknowledged_at', 'Bundle acknowledged'), ('queue_status', 'Queue status'), ('queue_settled_at', 'Queue settled'), ('projection_status', 'Projection status'), ('projection_completed_at', 'Projection completed'), ), (detail['transport'],)) package = _table(tuple( (key, key.replace('_', ' ').title()) for key in detail['package'] ), (detail['package'],)) current_policy = detail.get('current_effective_policy') current_policy_html = ( _render_deadline_policy({'rows': [current_policy]}) if current_policy else '

    Current effective source policy unavailable.

    ' ) diagnostics = [] for row in detail['diagnostics']['items']: envelope = row['diagnostic'] uid = str(row['diagnostic_uid']) target_id = f'diagnostic-json-{uid}' download = f'./{reservation_id}/diagnostics/{uid}.json' materials = [] http_context = envelope.get('http') or {} process_context = envelope.get('process') or {} if http_context: materials.append(_render_material('HTTP body', http_context.get('body'))) materials.append(_render_material('HTTP headers', http_context.get('headers'))) if process_context: materials.append(_render_material('Process stdout', process_context.get('stdout'))) materials.append(_render_material('Process stderr', process_context.get('stderr'))) if not materials: materials.append('

    No body or process-log material was captured.

    ') summary = ({ 'diagnostic_uid': uid, 'phase': row.get('phase'), 'kind': row.get('kind'), 'category': row.get('category'), 'code': row.get('code'), 'retryable': bool(row.get('retryable')), 'summary': row.get('summary'), 'occurred_at': row.get('occurred_at'), 'received_at': row.get('received_at'), 'canonical_sha256': row.get('envelope_sha256'), 'transformation': ( 'schema validation plus canonical ASCII JSON with sorted keys and compact separators' ), },) diagnostics.append( f'
    ' f'

    {html.escape(str(row.get("category")))} / ' f'{html.escape(str(row.get("code")))}

    ' + _table(tuple((key, key.replace('_', ' ').title()) for key in summary[0]), summary) + '
    ' f'' f'Download canonical JSON' '
    ' f'' + ''.join(materials) + '
    ' ) if not diagnostics: diagnostics.append( '

    No canonical diagnostic envelope is stored. ' f'Availability: {html.escape(detail["diagnostics"]["availability"])}.

    ' ) legacy = detail['legacy_evidence'] legacy_rows = legacy.get('errors') or [] legacy_content = ( '

    These fields are legacy evidence, not a synthesized diagnostic envelope.

    ' + _table(( ('id', 'Error'), ('category', 'Category'), ('summary', 'Summary'), ('created_at', 'Created'), ), legacy_rows) + ''.join( '

    Exact stored raw_error

    ' f'
    {html.escape(str(row.get("raw_error") or ""))}
    ' for row in legacy_rows ) ) if legacy.get('available') else ( '

    No legacy error evidence is stored. No envelope has been invented.

    ' ) truncation_notes = [] for name in ('progress', 'diagnostics'): if detail[name].get('truncated'): if name == 'progress': truncation_notes.append( f'progress retained the latest events chronologically; ' f'{detail[name].get("omitted_older_event_count", 0)} older events omitted' ) else: truncation_notes.append(f'{name} reached the server result bound') if legacy.get('truncated'): truncation_notes.append('legacy evidence reached the server result bound') truncation = ( '

    ' + html.escape('; '.join(truncation_notes)) + '.

    ' if truncation_notes else '' ) content = f'''

    Assignment

    Machine-readable detail

    {_table((('reservation_id', 'Reservation'), ('queue_id', 'Queue'), ('source', 'Source'), ('target', 'Target'), ('user', 'User'), ('worker', 'Worker / device'), ('active_assignment_cap', 'Configured cap'), ('assignment_outcome', 'Assignment outcome'), ('scan_outcome', 'Scan outcome'), ('issued_at', 'Issued'), ('resolved_at', 'Resolved'), ('assignment_code', 'Assignment code'), ('assignment_detail', 'Assignment detail'), ('diagnostic_availability', 'Diagnostic availability'), ('scan_deadline_at', 'Scan deadline'), ('historical_upload_timeout', 'Historical result-upload timeout'), ('assignment_deadline_at', 'Assignment deadline')), overview)}

    Progress authority

    {progress_summary}

    Ordered timeline

    {truncation}{timeline}

    Duration breakdown

    {durations}

    Transport, receipt, ingestion, settlement, projection

    {transport}

    Scan summary

    {scan}

    Package identity

    {package}

    Current effective source policy

    This is current policy context; the stored assignment deadline above remains immutable.

    {current_policy_html}

    Diagnostics

    {''.join(diagnostics)}

    Legacy evidence

    {legacy_content}
    ''' return _page_shell( f'Worker assignment {reservation_id}', 'Authoritative timestamps and exact stored diagnostic material.', 'workers', content, relative_root, script='../admin.js', ) def _component_value(snapshot, name, expected_type): component = snapshot.get(name) if isinstance(snapshot, dict) else None if not isinstance(component, dict) or component.get('available') is not True: return None value = component.get('value') return value if isinstance(value, expected_type) else None def _safe_count(value): return value if type(value) is int and value >= 0 else '' def _render_overview_page(snapshot, relative_root='.'): runtime = _component_value(snapshot, 'runtime', dict) queue = _component_value(snapshot, 'queue', dict) control = _component_value(snapshot, 'control', dict) operations = _component_value(snapshot, 'operations', list) if runtime is None: runtime_content = '

    Runtime snapshot unavailable.

    ' producers = [] pipeline_workers = [] pipeline = {} else: runtime_state = runtime.get('runtime') if isinstance(runtime.get('runtime'), dict) else {} postgres = runtime.get('postgres') if isinstance(runtime.get('postgres'), dict) else {} sources = runtime.get('sources') if isinstance(runtime.get('sources'), list) else [] by_id = { item.get('id'): item for item in sources if isinstance(item, dict) and isinstance(item.get('id'), str) } runtime_content = _table(( ('component', 'Component'), ('state', 'State'), ('ready', 'Ready'), ('safe_error_category', 'Safe error'), ), ( { 'component': 'Supervisor', 'state': runtime_state.get('phase', ''), 'ready': ( runtime_state.get('phase') == 'ACTIVE' and runtime_state.get('start_gate_open') is True and runtime_state.get('shutdown_requested') is False and runtime_state.get('runtime_failed') is False ), 'safe_error_category': 'runtime_failed' if runtime_state.get('runtime_failed') else '', }, { 'component': 'PostgreSQL', 'state': postgres.get('state', ''), 'ready': postgres.get('ready', False), 'safe_error_category': postgres.get('safe_error_category', ''), }, )) producers = [] for source in CORE_PRODUCERS: item = by_id.get(f'discovery-producer:{source}', {}) cycle = item.get('last_cycle_result') if isinstance(item.get('last_cycle_result'), dict) else {} producers.append({ 'source': source, 'role': item.get('role', 'unavailable'), 'lifecycle_state': item.get('lifecycle_state', 'unavailable'), 'last_cycle': cycle.get('status', ''), 'last_success': item.get('last_successful_discovery_at', ''), 'next_run': item.get('next_scheduled_run_at', ''), 'safe_error_category': item.get('safe_error_category', ''), }) pipeline_workers = [] for source_id in PIPELINE_SOURCE_IDS: item = by_id.get(source_id, {}) pipeline_workers.append({ 'id': source_id, 'role': item.get('role', 'unavailable'), 'lifecycle_state': item.get('lifecycle_state', 'unavailable'), 'process_state': item.get('process_state', 'unavailable'), 'desired_state': item.get('desired_state', ''), 'safe_error_category': item.get('safe_error_category', ''), }) pipeline = runtime.get('pipeline') if isinstance(runtime.get('pipeline'), dict) else {} pipeline_rows = ({ 'metric': 'Ingester lease', 'state': pipeline.get('ingester_state', '') if type(pipeline.get('ingester_state')) is str else '', 'ready': pipeline.get('ingester_ready', '') if type(pipeline.get('ingester_ready')) is bool else '', }, { 'metric': 'Projector lease', 'state': pipeline.get('projector_state', '') if type(pipeline.get('projector_state')) is str else '', 'ready': pipeline.get('projector_ready', '') if type(pipeline.get('projector_ready')) is bool else '', }, { 'metric': 'Final cutover', 'state': ( 'valid' if pipeline.get('cutover_ready') is True else 'invalid' if pipeline.get('cutover_ready') is False else '' ), 'ready': pipeline.get('cutover_ready', '') if type(pipeline.get('cutover_ready')) is bool else '', }) queue_rows = [] if queue is not None and isinstance(queue.get('counts'), dict): truncated = set(queue.get('truncated_statuses') or ()) queue_rows = [ { 'status': status, 'count': _safe_count(queue['counts'].get(status, 0)), 'truncated': status in truncated, } for status in QUEUE_STATUSES ] queue_note = ( f'Bounded sample: {html.escape(str(queue.get("sampled_rows", "")))} rows; ' f'degraded={html.escape(str(bool(queue.get("degraded"))))}; ' f'stale={html.escape(str(bool(queue.get("stale"))))}.' ) else: queue_note = 'Queue snapshot unavailable.' if control is None: control_rows = [] control_note = 'Control snapshot unavailable.' else: control_rows = ({ 'revision': control.get('revision', ''), 'discovery_paused': control.get('effective_discovery_paused', ''), 'dispatch_paused': control.get('effective_dispatch_paused', ''), 'drain_state': control.get('drain_state', ''), 'assignments': _safe_count(control.get('live_remote_assignments')), 'bundles': _safe_count(control.get('precommit_result_bundles')), 'updated_at': control.get('updated_at', ''), },) control_note = '' operation_rows = [] if operations is not None: for item in operations: if not isinstance(item, dict): continue operation_rows.append({ key: item.get(key, '') for key in ( 'operation_id', 'actor', 'action', 'target_ref', 'status', 'safe_category', 'safe_detail', 'requested_at', 'completed_at', 'updated_at', ) }) operations_note = '' else: operations_note = 'Recent operations unavailable.' content = f'''

    Runtime health

    {runtime_content}

    Discovery producers

    {_table((('source', 'Source'), ('role', 'Role'), ('lifecycle_state', 'Lifecycle'), ('last_cycle', 'Last cycle'), ('last_success', 'Last success'), ('next_run', 'Next run'), ('safe_error_category', 'Safe error')), producers)}

    Pipeline workers

    {_table((('id', 'ID'), ('role', 'Role'), ('lifecycle_state', 'Lifecycle'), ('process_state', 'Process'), ('desired_state', 'Desired'), ('safe_error_category', 'Safe error')), pipeline_workers)}

    Pipeline authority

    {_table((('metric', 'Metric'), ('state', 'State'), ('ready', 'Ready')), pipeline_rows)}

    Queue status

    {queue_note}

    {_table((('status', 'Status'), ('count', 'Bounded count'), ('truncated', 'Truncated')), queue_rows)}

    Assignments, bundles and control

    {control_note}

    {_table((('revision', 'Revision'), ('discovery_paused', 'Discovery paused'), ('dispatch_paused', 'Dispatch paused'), ('drain_state', 'Drain'), ('assignments', 'Active assignments'), ('bundles', 'Pre-commit bundles'), ('updated_at', 'Updated')), control_rows)}

    Bundle items: {html.escape(str(_safe_count(pipeline.get('bundle_items'))))}; bundle bytes: {html.escape(str(_safe_count(pipeline.get('bundle_bytes'))))}.

    Recent operation outcomes

    {operations_note}

    {_table((('operation_id', 'Operation'), ('actor', 'Actor'), ('action', 'Action'), ('target_ref', 'Target'), ('status', 'Status'), ('safe_category', 'Category'), ('safe_detail', 'Detail'), ('requested_at', 'Requested'), ('completed_at', 'Completed'), ('updated_at', 'Updated')), operation_rows)}
    ''' return _page_shell( 'Runtime overview', 'Bounded structured health; unavailable components degrade independently.', 'overview', content, relative_root, ) def _render_search_page(service, snapshot, relative_root='.'): runtime = _component_value(snapshot, 'runtime', dict) control = _component_value(snapshot, 'control', dict) if control is None: control_content = '

    Discovery control unavailable.

    ' control_forms = '' else: revision = control.get('revision', '') control_content = _table(( ('revision', 'Revision'), ('discovery_paused', 'Explicit pause'), ('effective_discovery_paused', 'Effective pause'), ('drain_state', 'Drain'), ('actor', 'Last actor'), ('updated_at', 'Updated'), ), ({ 'revision': revision, 'discovery_paused': control.get('discovery_paused', ''), 'effective_discovery_paused': control.get('effective_discovery_paused', ''), 'drain_state': control.get('drain_state', ''), 'actor': control.get('actor', ''), 'updated_at': control.get('updated_at', ''), },)) control_forms = '
    ' + ''.join(( _form( '/search/discovery/pause', 'Pause discovery', service.csrf_token, (), relative_root, hidden=( ('expected_revision', revision), ('operation_id', str(uuid.uuid4())), ), ), _form( '/search/discovery/resume', 'Resume discovery', service.csrf_token, (), relative_root, hidden=( ('expected_revision', revision), ('operation_id', str(uuid.uuid4())), ), ), )) + '
    ' if runtime is None: producer_content = '

    Producer state unavailable.

    ' producer_forms = '' else: sources = runtime.get('sources') if isinstance(runtime.get('sources'), list) else [] by_id = { item.get('id'): item for item in sources if isinstance(item, dict) and item.get('id') in PRODUCER_IDS } producer_rows = [] forms = [] for source_id in PRODUCER_IDS: item = by_id.get(source_id, {}) cycle = item.get('last_cycle_result') if isinstance( item.get('last_cycle_result'), dict, ) else {} producer_rows.append({ 'id': source_id, 'lifecycle_state': item.get('lifecycle_state', 'unavailable'), 'desired_state': item.get('desired_state', ''), 'process_state': item.get('process_state', ''), 'interval_seconds': item.get('interval_seconds', ''), 'restart_enabled': item.get('restart_enabled', ''), 'restart_count': item.get('restart_count', ''), 'last_cycle': cycle.get('status', ''), 'fetched': cycle.get('fetched_count', ''), 'queued_new': cycle.get('queued_new_count', ''), 'queued_updated': cycle.get('queued_updated_count', ''), 'last_success': item.get('last_successful_discovery_at', ''), 'next_run': item.get('next_scheduled_run_at', ''), 'safe_error_category': item.get('safe_error_category', ''), }) label = source_id.split(':', 1)[1] for action in ('start', 'stop', 'restart', 'pause', 'resume'): forms.append(_form( f'/search/producers/{action}', f'{action.title()} {label}', service.csrf_token, (), relative_root, hidden=( ('source_id', source_id), ('operation_id', str(uuid.uuid4())), ), )) forms.append(_form( '/search/producers/interval', f'Set {label} interval', service.csrf_token, (( 'interval_seconds', 'Interval seconds', 'number', 1, MAX_PRODUCER_INTERVAL_SECONDS, ),), relative_root, hidden=( ('source_id', source_id), ('operation_id', str(uuid.uuid4())), ), values={'interval_seconds': item.get('interval_seconds')}, )) producer_content = _table(( ('id', 'Producer'), ('lifecycle_state', 'Lifecycle'), ('desired_state', 'Desired'), ('process_state', 'Process'), ('interval_seconds', 'Interval seconds'), ('restart_enabled', 'Restart'), ('restart_count', 'Restarts'), ('last_cycle', 'Last cycle'), ('fetched', 'Fetched'), ('queued_new', 'Queued new'), ('queued_updated', 'Queued updated'), ('last_success', 'Last success'), ('next_run', 'Next run'), ('safe_error_category', 'Safe error'), ), producer_rows) producer_forms = '
    ' + ''.join(forms) + '
    ' content = f'''

    Persistent discovery control

    Stored in PostgreSQL and preserved across runtime restarts.

    {control_content}{control_forms}

    Discovery producers

    Lifecycle and interval changes affect the current Supervisor runtime only.

    {producer_content}{producer_forms}
    ''' return _page_shell( 'Search operations', 'Persistent discovery authority and exact producer controls.', 'search', content, relative_root, ) def _render_supervisor_page(service, snapshot, relative_root='.'): if not isinstance(snapshot, dict): return _page_shell( 'Supervisor', 'Structured managed-process state and typed controls.', 'supervisor', '

    Supervisor state unavailable.

    ', relative_root, ) runtime = snapshot.get('runtime') if isinstance(snapshot.get('runtime'), dict) else {} postgres = snapshot.get('postgres') if isinstance(snapshot.get('postgres'), dict) else {} dashboard = snapshot.get('dashboard') if isinstance(snapshot.get('dashboard'), dict) else {} sources = [ item for item in snapshot.get('sources', []) if isinstance(item, dict) and type(item.get('id')) is str and KEY_RE.fullmatch(item['id']) and isinstance(item.get('allowed_actions'), list) and all(action in ALL_MANAGED_SOURCE_ACTIONS for action in item['allowed_actions']) ] content = '

    Runtime

    ' + _table(( ('phase', 'Phase'), ('pid', 'PID'), ('start_gate_open', 'Start gate'), ('shutdown_requested', 'Shutdown requested'), ('runtime_failed', 'Failed'), ('postgres_state', 'PostgreSQL'), ('postgres_ready', 'PostgreSQL ready'), ('postgres_error', 'PostgreSQL safe error'), ), ({ 'phase': runtime.get('phase', ''), 'pid': runtime.get('pid', ''), 'start_gate_open': runtime.get('start_gate_open', ''), 'shutdown_requested': runtime.get('shutdown_requested', ''), 'runtime_failed': runtime.get('runtime_failed', ''), 'postgres_state': postgres.get('state', ''), 'postgres_ready': postgres.get('ready', ''), 'postgres_error': postgres.get('safe_error_category', ''), },)) + '
    ' content += '

    Dashboard

    ' + _table(( ('status', 'Status'), ('desired_state', 'Desired state'), ('healthy', 'Healthy'), ('pid', 'PID'), ('safe_error_category', 'Safe error'), ), (dashboard,)) dashboard_forms = ''.join(_form( f'/supervisor/dashboard/{action}', f'{action.title()} dashboard', service.csrf_token, (), relative_root, hidden=(('operation_id', str(uuid.uuid4())),), ) for action in ('start', 'stop', 'restart')) dashboard_summary = 'Dashboard controls' if dashboard.get('status'): dashboard_summary += f' - {dashboard["status"]}' content += ( '
    ' + html.escape(dashboard_summary) + '
    ' + dashboard_forms + '
    ' ) content += '

    Managed sources

    ' + _table(( ('id', 'ID'), ('role', 'Role'), ('lifecycle_state', 'Lifecycle'), ('desired_state', 'Desired'), ('process_state', 'Process'), ('pid', 'PID'), ('mode', 'Mode'), ('interval_seconds', 'Interval'), ('restart_enabled', 'Restart'), ('restart_delay_seconds', 'Restart delay'), ('restart_count', 'Restarts'), ('restart_streak', 'Restart streak'), ('safe_error_category', 'Safe error'), ), sources) panels = [] for source in sources: source_id = source['id'] allowed = source.get('allowed_actions') if not isinstance(allowed, list): allowed = [] forms = [] for action in ('start', 'stop', 'restart', 'pause', 'resume', 'once'): if action in allowed: forms.append(_form( f'/supervisor/sources/{action}', f'{action.title()} {source_id}', service.csrf_token, (), relative_root, hidden=( ('source_id', source_id), ('operation_id', str(uuid.uuid4())), ), )) if 'set-mode' in allowed: for mode in ('loop', 'once', 'repeat'): if source_id == 'keychecks' and mode == 'loop': continue forms.append(_form( f'/supervisor/sources/mode/{mode}', f'Set {source_id} mode {mode}', service.csrf_token, (), relative_root, hidden=( ('source_id', source_id), ('operation_id', str(uuid.uuid4())), ), )) if 'set-interval' in allowed: forms.append(_form( '/supervisor/sources/interval', f'Set {source_id} interval', service.csrf_token, (( 'interval_seconds', 'Interval seconds', 'number', 1, MAX_MANAGED_SOURCE_DELAY_SECONDS, ),), relative_root, hidden=( ('source_id', source_id), ('operation_id', str(uuid.uuid4())), ), values={'interval_seconds': source.get('interval_seconds')}, )) if 'set-restart' in allowed: for enabled, label in ((True, 'Enable'), (False, 'Disable')): route = 'enable' if enabled else 'disable' forms.append(_form( f'/supervisor/sources/restart-policy/{route}', f'{label} {source_id} restart', service.csrf_token, (), relative_root, hidden=( ('source_id', source_id), ('operation_id', str(uuid.uuid4())), ), )) if 'set-restart-delay' in allowed: forms.append(_form( '/supervisor/sources/restart-delay', f'Set {source_id} restart delay', service.csrf_token, (( 'restart_delay_seconds', 'Restart delay seconds', 'number', 1, MAX_MANAGED_SOURCE_DELAY_SECONDS, ),), relative_root, hidden=( ('source_id', source_id), ('operation_id', str(uuid.uuid4())), ), values={ 'restart_delay_seconds': source.get('restart_delay_seconds'), }, )) summary_parts = [source_id] for key in ('lifecycle_state', 'process_state'): value = source.get(key) if value: summary_parts.append(str(value)) panels.append( '
    ' + html.escape(' - '.join(summary_parts)) + '
    ' + ''.join(forms) + '
    ' ) content += '
    ' + ''.join(panels) + '
    ' return _page_shell( 'Supervisor', 'Structured managed-process state and typed controls.', 'supervisor', content, relative_root, ) def _render_logs_page(service, snapshot, tail=None, relative_root='.'): if snapshot is None: content = '

    Managed source list unavailable.

    ' sources = [] else: sources = [ item for item in snapshot.get('sources', []) if isinstance(item, dict) and type(item.get('id')) is str and KEY_RE.fullmatch(item['id']) and isinstance(item.get('allowed_actions'), list) and all(action in ALL_MANAGED_SOURCE_ACTIONS for action in item['allowed_actions']) ] content = ( '

    Managed source logs

    ' '

    Only the active bounded log for an allowlisted source is available.

    ' '
    ' + ''.join(_form( '/logs/tail', f'Tail {item["id"]}', service.csrf_token, (( 'line_count', 'Newest lines', 'number', 1, MAX_MANAGED_SOURCE_LOG_LINES, ),), relative_root, hidden=(('source_id', item['id']),), values={'line_count': 40}, ) for item in sources) + '
    ' ) if isinstance(tail, dict): lines = '\n'.join(tail.get('lines') or []) suffix = ' Response was truncated.' if tail.get('response_truncated') else '' content += ( '

    Tail result

    ' + html.escape( f'{tail.get("source_id", "")} ยท {tail.get("line_count", 0)} lines.{suffix}' ) + '

    ' + html.escape(lines) + '
    ' ) return _page_shell( 'Managed logs', 'Bounded allowlisted log tails; no path or command input.', 'logs', content, relative_root, ) def _runtime_document_hashes(state): return { 'active_config': state.active_config.sha256, 'active_secrets': state.active_secrets.sha256, 'candidate_config': state.candidate_config.sha256, 'candidate_secrets': state.candidate_secrets.sha256, } def _render_runtime_document_page( service, document, editor, *, document_text=None, preview=None, notice='', relative_root='.', operation_id=None, observability=None, ): state = preview.state if preview is not None else editor.state text = editor.text if document_text is None else document_text hashes = _runtime_document_hashes(state) hidden = ''.join( f'' for name, value in hashes.items() ) operation_id = str(uuid.uuid4()) if operation_id is None else operation_id editor_form = ( f'
    ' f'' f'{hidden}' f'' '
    ' f'
    ' '
    ' ) identity_rows = ({ 'active_config': hashes['active_config'], 'active_secrets': hashes['active_secrets'], 'candidate_config': hashes['candidate_config'], 'candidate_secrets': hashes['candidate_secrets'], 'config_candidate_present': state.candidate_config.present, 'secrets_candidate_present': state.candidate_secrets.present, },) content = ( (f'

    {html.escape(notice)}

    ' if notice else '') + f'

    {html.escape(document.title())} candidate

    ' + f'

    Editing source: {html.escape(editor.source)}. Save stages a private candidate and never activates it.

    ' + _table(( ('active_config', 'Active config SHA-256'), ('active_secrets', 'Active secrets SHA-256'), ('candidate_config', 'Candidate config SHA-256'), ('candidate_secrets', 'Candidate secrets SHA-256'), ('config_candidate_present', 'Config candidate'), ('secrets_candidate_present', 'Secrets candidate'), ), identity_rows) + editor_form + '
    ' ) if preview is not None: diff = preview.diff if hasattr(diff, 'entries'): rows = [ { 'path': entry.path, 'change': entry.change, 'before': entry.before, 'after': entry.after, 'redacted': entry.value_redacted, } for entry in diff.entries ] diff_html = _table(( ('path', 'Path'), ('change', 'Change'), ('before', 'Before'), ('after', 'After'), ('redacted', 'Redacted'), ), rows) diff_html += ( f'

    Truncated: {html.escape(str(diff.truncated))}; ' f'format-only change: {html.escape(str(diff.format_only_changed))}.

    ' ) else: fields = ( 'document_changed', 'semantic_changed', 'pools_before', 'pools_after', 'pools_added', 'pools_removed', 'pools_changed', 'entries_before', 'entries_after', 'entries_added', 'entries_removed', 'pools_reordered', 'usernames_added', 'usernames_removed', 'usernames_changed', 'tokens_changed', ) diff_html = _table(tuple((field, field.replace('_', ' ').title()) for field in fields), ({ field: getattr(diff, field) for field in fields },)) content += '

    Validated preview

    ' + diff_html + '
    ' if document == 'config': observability = dict(observability or {}) policy = _component_value(observability, 'policy', dict) metrics = _component_value(observability, 'metrics', dict) policy_html = ( '

    Effective deadline policy unavailable.

    ' if policy is None else _render_deadline_policy(policy) ) metrics_html = ( '

    Duration metrics unavailable.

    ' if metrics is None else ( f'

    Total source / phase / outcome groups: ' f'{int(metrics.get("total_group_count") or 0)}. ' + ( 'Additional groups are available through paginated Workers observability.

    ' if metrics.get('has_next') else '

    ' ) + _render_duration_metrics(metrics.get('metrics') or []) ) ) candidate_label = ( 'Validated candidate effective values' if preview is not None else f'{editor.source.title()} editor effective values' ) content += ( f'

    {html.escape(candidate_label)}

    {policy_html}' '

    Observed source / phase / outcome durations

    ' '

    p50/p95/p99 values include sample counts and never mutate policy.

    ' f'{metrics_html}
    ' ) relevant_candidate = ( state.candidate_config if document == 'config' else state.candidate_secrets ) relevant_active = ( state.active_config if document == 'config' else state.active_secrets ) can_apply_document = bool( relevant_candidate.present and relevant_candidate.sha256 != relevant_active.sha256 ) can_apply_both = bool( state.candidate_config.present and state.candidate_secrets.present and ( state.candidate_config.sha256 != state.active_config.sha256 or state.candidate_secrets.sha256 != state.active_secrets.sha256 ) ) apply_hidden = [ ('operation_id', str(uuid.uuid4())), ('expected_active_config_sha256', hashes['active_config']), ('expected_active_secrets_sha256', hashes['active_secrets']), (f'expected_candidate_{document}_sha256', relevant_candidate.sha256), ] apply_form = _form( f'/{document}/apply', f'Apply {document} candidate', service.csrf_token, (), relative_root, hidden=tuple(apply_hidden), ) if can_apply_document else '' both_form = _form( '/runtime/apply-both', 'Apply both candidates', service.csrf_token, (), relative_root, hidden=( ('operation_id', str(uuid.uuid4())), ('expected_active_config_sha256', hashes['active_config']), ('expected_active_secrets_sha256', hashes['active_secrets']), ('expected_candidate_config_sha256', hashes['candidate_config']), ('expected_candidate_secrets_sha256', hashes['candidate_secrets']), ), ) if can_apply_both else '' apply_controls = apply_form + both_form apply_state = ( f'
    {apply_controls}
    ' if apply_controls else '

    No staged candidate differs from the active document. ' 'Save a changed candidate before applying.

    ' ) content += ( '

    Apply

    Apply requires the fixed host agent. ' 'The operation is hash-bound and durable before dispatch.

    ' f'{apply_state}
    ' ) return _page_shell( f'{document.title()} editor', 'Plaintext is rendered only inside this protected no-store editor.', document, content, relative_root, ) def _redirect(location): response = _secure_response('', status_code=303) response.headers['Location'] = location return response class _ManagedFileStreamingResponse(StreamingResponse): def __init__(self, snapshot, *args, **kwargs): self._managed_file_snapshot = snapshot super().__init__(*args, **kwargs) async def __call__(self, scope, receive, send): try: await super().__call__(scope, receive, send) finally: self._managed_file_snapshot.close() def _managed_file_download_response(download, relative_path): leaf_name = relative_path.rsplit('/', 1)[-1] headers = dict(SECURITY_HEADERS) headers['Content-Disposition'] = ( "attachment; filename*=UTF-8''" + quote(leaf_name, safe='') ) headers['Content-Length'] = str(download.identity.byte_count) headers['ETag'] = f'"{download.identity.sha256}"' if download.snapshot is not None: try: return _ManagedFileStreamingResponse( download.snapshot, download.snapshot.chunks(), media_type='application/octet-stream', headers=headers, ) except BaseException: download.snapshot.close() raise return Response( download.content, media_type='application/octet-stream', headers=headers, ) async def _download_managed_file( service, traversal, root_id, relative_path): task = asyncio.create_task(asyncio.to_thread( service.download_managed_file, traversal, root_id, relative_path, )) try: return await asyncio.shield(task) except asyncio.CancelledError: def close_snapshot(completed): try: download = completed.result() except BaseException: return if download.snapshot is not None: download.snapshot.close() task.add_done_callback(close_snapshot) raise ADMIN_CSS = ''' :root { color-scheme: light; font-family: ui-monospace, Consolas, monospace; background: #f3f0e8; color: #18211d; } body { margin: 0; } header, main { box-sizing: border-box; width: 100%; max-width: 92rem; min-width: 0; margin: 0 auto; padding: 1.25rem; } header { border-bottom: 4px solid #18211d; } nav { display: flex; gap: .5rem; flex-wrap: wrap; } nav a { color: inherit; border: 1px solid #6d746f; padding: .4rem .65rem; text-decoration: none; background: #fffdf6; } nav a[aria-current="page"] { background: #18211d; color: #fffdf6; } section, article { box-sizing: border-box; max-width: 100%; min-width: 0; } section { margin: 1.5rem 0; overflow-x: auto; } table { width: 100%; border-collapse: collapse; background: #fffdf6; } th, td { border: 1px solid #9b9a8e; padding: .45rem; text-align: left; overflow-wrap: anywhere; } .forms { display: grid; grid-template-columns: repeat(auto-fit, minmax(16rem, 1fr)); gap: .75rem; } .control-panels { display: grid; gap: .75rem; margin-top: 1rem; } .control-panel { border: 1px solid #6d746f; background: #e4eee8; } .control-panel > summary { cursor: pointer; font-weight: bold; padding: .8rem; } .control-panel[open] > summary { border-bottom: 1px solid #6d746f; } .control-panel > .forms { padding: .8rem; } form { border: 1px solid #6d746f; background: #fffdf6; padding: .8rem; } label, input, select, button { display: block; box-sizing: border-box; width: 100%; margin: .45rem 0; } input, select, button, textarea { font: inherit; padding: .45rem; box-sizing: border-box; } textarea { display: block; width: 100%; resize: vertical; white-space: pre; overflow: auto; } .button-row { display: grid; grid-template-columns: repeat(auto-fit, minmax(10rem, 1fr)); gap: .5rem; } button { background: #174c3c; color: white; border: 0; cursor: pointer; } .button-link { display: block; box-sizing: border-box; margin: .45rem 0; padding: .45rem; text-align: center; background: #fffdf6; border: 1px solid #174c3c; color: #174c3c; text-decoration: none; } .filter-grid { display: grid; grid-template-columns: repeat(auto-fit, minmax(12rem, 1fr)); gap: .5rem; } pre { max-height: 40rem; overflow: auto; white-space: pre-wrap; overflow-wrap: anywhere; background: #18211d; color: #fffdf6; padding: .8rem; } .diagnostic, .diagnostic-group, .material { border: 1px solid #6d746f; padding: .8rem; margin: .8rem 0; background: #fffdf6; } .diagnostic-group code { overflow-wrap: anywhere; } .canonical-json { background: #f3f0e8; min-height: 12rem; } .issued { border: 3px solid #8b2f24; padding: 1rem; background: #fff4df; } .issued code { overflow-wrap: anywhere; } .notice { border-left: 4px solid #174c3c; padding: .7rem; background: #e4eee8; } .unavailable { color: #8b2f24; font-weight: bold; } @media (max-width: 48rem) { header, main { padding: .75rem; } table { width: max-content; min-width: 100%; } th, td { overflow-wrap: normal; word-break: normal; } } '''.strip() ADMIN_JS = ''' document.addEventListener('click', async (event) => { const button = event.target.closest('[data-copy-target]'); if (!button) return; const target = document.getElementById(button.dataset.copyTarget); if (!target) return; await navigator.clipboard.writeText(target.value); button.textContent = 'Copied canonical JSON'; }); '''.strip() async def _dispatch_runtime_document(request, service, document, action): fields = hashes = editor = preview = document_text = operation_id = None try: if action in ('preview', 'save'): expected = { 'csrf_token', 'document_text', 'operation_id', 'expected_active_config_sha256', 'expected_active_secrets_sha256', 'expected_candidate_config_sha256', 'expected_candidate_secrets_sha256', } max_document = ( MAX_CONFIG_DOCUMENT_BYTES if document == 'config' else MAX_SECRETS_DOCUMENT_BYTES ) fields = await _form_fields( request, service, expected, body_limit=(max_document * 3) + 4096, ) document_text = fields.pop('document_text') hashes = _parse_document_hashes(fields) operation_id = _parse_operation_id(fields['operation_id']) editor = await asyncio.to_thread(service.runtime_document_editor, document) try: if action == 'preview': preview = await asyncio.to_thread( service.preview_runtime_document, document, document_text, ) observability = None if document == 'config': observability = await asyncio.to_thread( service.runtime_document_observability, document_text, ) return _secure_response( _render_runtime_document_page( service, document, editor, document_text=document_text, preview=preview, notice='Candidate is valid; nothing was saved.', relative_root='..', operation_id=operation_id, observability=observability, ), media_type='text/html', ) await asyncio.to_thread( service.save_runtime_document_candidate, document, document_text, request.state.admin_actor, operation_id, expected_hashes=hashes, ) except AdminAPIError as exc: if action == 'save': editor = await asyncio.to_thread( service.runtime_document_editor, document, ) observability = None if document == 'config': observability = await asyncio.to_thread( service.runtime_document_observability, document_text, ) return _secure_response( _render_runtime_document_page( service, document, editor, document_text=document_text, notice=str(exc), relative_root='..', operation_id=operation_id, observability=observability, ), status_code=exc.status_code, media_type='text/html', ) return _redirect(f'../{document}') apply_action = f'apply-{document}' if document != 'both' else 'apply-both' candidate_fields = ( {'expected_candidate_config_sha256'} if apply_action == 'apply-config' else {'expected_candidate_secrets_sha256'} if apply_action == 'apply-secrets' else {'expected_candidate_config_sha256', 'expected_candidate_secrets_sha256'} ) fields = await _form_fields( request, service, { 'csrf_token', 'operation_id', 'expected_active_config_sha256', 'expected_active_secrets_sha256', *candidate_fields, }, ) hashes = _parse_document_hashes( fields, require_config=apply_action in ('apply-config', 'apply-both'), require_secrets=apply_action in ('apply-secrets', 'apply-both'), ) operation_id = _parse_operation_id(fields['operation_id']) await asyncio.to_thread( service.request_runtime_apply, apply_action, request.state.admin_actor, operation_id, expected_hashes=hashes, ) return _redirect(f'../operations/{operation_id}') finally: if isinstance(fields, dict): fields.clear() fields = hashes = editor = preview = document_text = operation_id = None async def _dispatch_managed_file_mutation(request, service, action): fields = payload = encoded_content = operation_id = expected_sha256 = None try: expected = {'csrf_token', 'operation_id', 'root_id', 'relative_path'} if action in ('create', 'replace'): expected.add('content_base64') if action in ('replace', 'delete'): expected.add('expected_sha256') fields = await _form_fields(request, service, expected) operation_id = _parse_operation_id(fields['operation_id']) if action in ('replace', 'delete'): expected_sha256 = _parse_managed_file_hash(fields['expected_sha256']) if action in ('create', 'replace'): encoded_content = fields.pop('content_base64') payload = _parse_managed_file_content(encoded_content) encoded_content = None traversal = getattr(request.app.state, 'managed_file_traversal', None) if action == 'create': await asyncio.to_thread( service.create_managed_file, traversal, fields['root_id'], fields['relative_path'], payload, request.state.admin_actor, operation_id, ) elif action == 'replace': await asyncio.to_thread( service.replace_managed_file, traversal, fields['root_id'], fields['relative_path'], payload, expected_sha256, request.state.admin_actor, operation_id, ) else: await asyncio.to_thread( service.delete_managed_file, traversal, fields['root_id'], fields['relative_path'], expected_sha256, request.state.admin_actor, operation_id, ) return _redirect(_managed_file_listing_location( fields['root_id'], fields['relative_path'], )) finally: if isinstance(fields, dict): fields.clear() fields = payload = encoded_content = operation_id = expected_sha256 = None async def _dispatch(request, service): path = str(request.scope.get('path') or '') raw_path = request.scope.get('raw_path') try: canonical_raw_path = path.encode('ascii') except UnicodeEncodeError: canonical_raw_path = None if type(raw_path) is not bytes or raw_path != canonical_raw_path: raise AdminAPIError(404, 'page not found') method = request.method.upper() allowed_query = set() if method == 'GET' and path == f'{ADMIN_PREFIX}/audit': allowed_query = {'before'} elif method == 'GET' and path == f'{ADMIN_PREFIX}/operations': allowed_query = {'before'} elif method == 'GET' and path in (ADMIN_PREFIX, f'{ADMIN_PREFIX}/'): allowed_query = WORKER_FILTER_FIELDS elif method == 'GET' and path in ( f'{ADMIN_PREFIX}/files', f'{ADMIN_PREFIX}/files/download', ): allowed_query = {'root_id', 'relative_path'} query = _query_fields(request, allowed_query) if method == 'GET': if path in (ADMIN_PREFIX, f'{ADMIN_PREFIX}/'): filters, selected = _parse_worker_filters(query) snapshot = await asyncio.to_thread( service.workers_dispatch_snapshot, filters, diagnostic_occurrence_offset=int(selected['diagnostic_offset']), metric_offset=int(selected['metric_offset']), page_limit=int(selected['limit']), include_diagnostics=selected['details'] in {'diagnostics', 'all'}, include_metrics=selected['details'] in {'metrics', 'all'}, ) snapshot['selected_filters'] = selected return _secure_response(_render_workers_page(service, snapshot), media_type='text/html') if path == f'{ADMIN_PREFIX}/overview': snapshot = await asyncio.to_thread(service.overview) return _secure_response( _render_overview_page(snapshot), media_type='text/html', ) if path == f'{ADMIN_PREFIX}/search': snapshot = await asyncio.to_thread(service.search_snapshot) return _secure_response( _render_search_page(service, snapshot), media_type='text/html', ) if path in (f'{ADMIN_PREFIX}/supervisor', f'{ADMIN_PREFIX}/logs'): try: snapshot = await asyncio.to_thread(service._runtime_snapshot) except Exception: snapshot = None renderer = _render_supervisor_page if path.endswith('/supervisor') else _render_logs_page return _secure_response( renderer(service, snapshot), media_type='text/html', ) if path in (f'{ADMIN_PREFIX}/config', f'{ADMIN_PREFIX}/secrets'): document = path.rsplit('/', 1)[-1] editor = await asyncio.to_thread(service.runtime_document_editor, document) observability = None if document == 'config': observability = await asyncio.to_thread( service.runtime_document_observability, editor.text, ) return _secure_response( _render_runtime_document_page( service, document, editor, observability=observability, ), media_type='text/html', ) if path == f'{ADMIN_PREFIX}/files': if set(query) not in (set(), {'root_id'}, {'root_id', 'relative_path'}): raise AdminAPIError(400, 'query shape is invalid') root = listing = None relative_path = query.get('relative_path') if 'root_id' in query: root = service.managed_file_roots.get(query['root_id']) if root is None: raise AdminAPIError(404, 'managed file target was not found') if relative_path == '': raise AdminAPIError(400, 'managed file request is invalid') if relative_path is not None: try: parse_managed_relative_path(relative_path, root.limits) except ManagedFileAccessError as exc: _managed_file_error(exc) if root.permissions.allow_list: listing = await asyncio.to_thread( service.list_managed_files, getattr(request.app.state, 'managed_file_traversal', None), root.root_id, relative_path, ) return _secure_response( _render_files_page( service, root=root, relative_path=relative_path, listing=listing, ), media_type='text/html', ) if path == f'{ADMIN_PREFIX}/files/download': if set(query) != {'root_id', 'relative_path'} or not query['relative_path']: raise AdminAPIError(400, 'query shape is invalid') download = await _download_managed_file( service, getattr(request.app.state, 'managed_file_traversal', None), query['root_id'], query['relative_path'], ) return _managed_file_download_response(download, query['relative_path']) if path == f'{ADMIN_PREFIX}/operations': before = ( _parse_operation_cursor(query['before']) if 'before' in query else None ) page = await asyncio.to_thread(service.operation_page, before) return _secure_response( _render_operations_page(page), media_type='text/html', ) if path.startswith(f'{ADMIN_PREFIX}/operations/'): try: operation_id = _parse_operation_id( path[len(f'{ADMIN_PREFIX}/operations/'):], ) except AdminAPIError as exc: raise AdminAPIError(404, 'Not Found') from exc operation = await asyncio.to_thread( service.operation_status, operation_id, ) return _secure_response( _render_operation_page(operation), media_type='text/html', ) if path == f'{ADMIN_PREFIX}/audit': before_event_id = ( _parse_audit_cursor(query['before']) if 'before' in query else None ) page = await asyncio.to_thread( service.audit_page, before_event_id, ) return _secure_response( _render_audit_page(page), media_type='text/html', ) assignment_path = path[len(f'{ADMIN_PREFIX}/assignments/'):] match = re.fullmatch(r'([1-9][0-9]{0,18})(\.json)?', assignment_path) if match: reservation_id = _parse_assignment_id(match.group(1)) detail = await asyncio.to_thread( service.assignment_detail, reservation_id, ) if match.group(2): return _secure_json_response(detail) return _secure_response( _render_assignment_detail_page(detail), media_type='text/html', ) match = re.fullmatch( r'([1-9][0-9]{0,18})/diagnostics/([a-f0-9]{64})\.json', assignment_path, ) if match: reservation_id = _parse_assignment_id(match.group(1)) diagnostic_uid = match.group(2) envelope = await asyncio.to_thread( service.diagnostic_envelope, reservation_id, diagnostic_uid, ) return _diagnostic_json_response(envelope, diagnostic_uid) if path == f'{ADMIN_PREFIX}/admin.css': return _secure_response(ADMIN_CSS, media_type='text/css') if path == f'{ADMIN_PREFIX}/admin.js': return _secure_response(ADMIN_JS, media_type='text/javascript') raise AdminAPIError(404, 'Not Found') if method != 'POST': raise AdminAPIError(405, 'Method Not Allowed') notice = '' issued_token = None document_routes = { f'{ADMIN_PREFIX}/{document}/{action}': (document, action) for document in ('config', 'secrets') for action in ('preview', 'save', 'apply') } document_routes[f'{ADMIN_PREFIX}/runtime/apply-both'] = ('both', 'apply') if path in document_routes: document, action = document_routes[path] return await _dispatch_runtime_document( request, service, document, action, ) managed_file_routes = { f'{ADMIN_PREFIX}/files/{action}': action for action in ('create', 'replace', 'delete') } if path in managed_file_routes: return await _dispatch_managed_file_mutation( request, service, managed_file_routes[path], ) dispatch_routes = { f'{ADMIN_PREFIX}/dispatch/pause': ('dispatch', True), f'{ADMIN_PREFIX}/dispatch/resume': ('dispatch', False), f'{ADMIN_PREFIX}/dispatch/drain/start': ('drain', 'start'), f'{ADMIN_PREFIX}/dispatch/drain/cancel': ('drain', 'cancel'), } if path in dispatch_routes: fields = await _form_fields( request, service, {'csrf_token', 'expected_revision', 'operation_id'}, ) kind, value = dispatch_routes[path] arguments = ( _parse_revision(fields['expected_revision']), request.state.admin_actor, _parse_operation_id(fields['operation_id']), ) if kind == 'dispatch': await asyncio.to_thread( service.set_dispatch_paused, value, *arguments, ) else: operation = service.start_drain if value == 'start' else service.cancel_drain await asyncio.to_thread(operation, *arguments) return _redirect('../' if kind == 'dispatch' else '../../') if path in ( f'{ADMIN_PREFIX}/search/discovery/pause', f'{ADMIN_PREFIX}/search/discovery/resume', ): fields = await _form_fields( request, service, {'csrf_token', 'expected_revision', 'operation_id'}, ) await asyncio.to_thread( service.set_discovery_paused, path.endswith('/pause'), _parse_revision(fields['expected_revision']), request.state.admin_actor, _parse_operation_id(fields['operation_id']), ) return _redirect('../../search') producer_routes = { f'{ADMIN_PREFIX}/search/producers/{action}': action for action in ('start', 'stop', 'restart', 'pause', 'resume') } if path in producer_routes: fields = await _form_fields( request, service, {'csrf_token', 'source_id', 'operation_id'}, ) await asyncio.to_thread( service.producer_action, fields['source_id'], producer_routes[path], request.state.admin_actor, _parse_operation_id(fields['operation_id']), ) return _redirect('../../search') if path == f'{ADMIN_PREFIX}/search/producers/interval': fields = await _form_fields( request, service, { 'csrf_token', 'source_id', 'interval_seconds', 'operation_id', }, ) await asyncio.to_thread( service.producer_action, fields['source_id'], 'set-interval', request.state.admin_actor, _parse_operation_id(fields['operation_id']), interval_seconds=_parse_producer_interval(fields['interval_seconds']), ) return _redirect('../../search') source_routes = { f'{ADMIN_PREFIX}/supervisor/sources/{action}': (action, {}, '../../supervisor') for action in ('start', 'stop', 'restart', 'pause', 'resume', 'once') } source_routes.update({ f'{ADMIN_PREFIX}/supervisor/sources/mode/{mode}': ( 'set-mode', {'mode': mode}, '../../../supervisor', ) for mode in ('loop', 'once', 'repeat') }) source_routes.update({ f'{ADMIN_PREFIX}/supervisor/sources/restart-policy/{value}': ( 'set-restart', {'restart_enabled': value == 'enable'}, '../../../supervisor', ) for value in ('enable', 'disable') }) if path in source_routes: fields = await _form_fields( request, service, {'csrf_token', 'source_id', 'operation_id'}, ) source_action, parameters, redirect = source_routes[path] await asyncio.to_thread( service.managed_source_action, fields['source_id'], source_action, request.state.admin_actor, _parse_operation_id(fields['operation_id']), **parameters, ) return _redirect(redirect) if path in ( f'{ADMIN_PREFIX}/supervisor/sources/interval', f'{ADMIN_PREFIX}/supervisor/sources/restart-delay', ): field_name = 'interval_seconds' if path.endswith('/interval') else 'restart_delay_seconds' fields = await _form_fields( request, service, {'csrf_token', 'source_id', 'operation_id', field_name}, ) delay = _parse_managed_source_delay(fields[field_name], field_name.replace('_', ' ')) await asyncio.to_thread( service.managed_source_action, fields['source_id'], 'set-interval' if field_name == 'interval_seconds' else 'set-restart-delay', request.state.admin_actor, _parse_operation_id(fields['operation_id']), **{field_name: delay}, ) return _redirect('../../supervisor') dashboard_routes = { f'{ADMIN_PREFIX}/supervisor/dashboard/{action}': action for action in ('start', 'stop', 'restart') } if path in dashboard_routes: fields = await _form_fields( request, service, {'csrf_token', 'operation_id'}, ) await asyncio.to_thread( service.managed_source_action, 'dashboard', dashboard_routes[path], request.state.admin_actor, _parse_operation_id(fields['operation_id']), ) return _redirect('../../supervisor') if path == f'{ADMIN_PREFIX}/logs/tail': fields = await _form_fields( request, service, {'csrf_token', 'source_id', 'line_count'}, ) line_count = _parse_managed_source_delay(fields['line_count'], 'log line count') if line_count > MAX_MANAGED_SOURCE_LOG_LINES: raise AdminAPIError(400, 'log line count is invalid') tail = await asyncio.to_thread( service.managed_source_log, fields['source_id'], line_count, ) try: snapshot = await asyncio.to_thread(service._runtime_snapshot) except Exception: snapshot = None return _secure_response( _render_logs_page(service, snapshot, tail=tail, relative_root='..'), media_type='text/html', ) if path in (f'{ADMIN_PREFIX}/users/create', f'{ADMIN_PREFIX}/users/cap'): fields = await _form_fields( request, service, { 'csrf_token', 'user_key', 'active_assignment_cap', 'operation_id', }, ) operation = service.create_user if path.endswith('/create') else service.set_user_cap result = await run_in_threadpool( operation, fields['user_key'], fields['active_assignment_cap'], request.state.admin_actor, _parse_operation_id(fields['operation_id']), ) if not result: raise AdminAPIError(409 if path.endswith('/create') else 404, 'user operation was not applied') return _redirect('../') elif path in (f'{ADMIN_PREFIX}/users/disable', f'{ADMIN_PREFIX}/users/enable'): fields = await _form_fields( request, service, {'csrf_token', 'user_key', 'operation_id'}, ) disabled = path.endswith('/disable') result = await run_in_threadpool( service.set_user_disabled, fields['user_key'], disabled, request.state.admin_actor, _parse_operation_id(fields['operation_id']), ) if not result: raise AdminAPIError(404, 'user was not found') return _redirect('../') elif path in (f'{ADMIN_PREFIX}/devices/issue', f'{ADMIN_PREFIX}/devices/rotate'): fields = await _form_fields( request, service, { 'csrf_token', 'user_key', 'device_key', 'operation_id', }, ) rotate = path.endswith('/rotate') result, issued_token = await run_in_threadpool( service.issue_device, fields['user_key'], fields['device_key'], request.state.admin_actor, _parse_operation_id(fields['operation_id']), rotate=rotate, ) if not result: raise AdminAPIError(409 if not rotate else 404, 'device token operation was not applied') notice = 'Device token rotated.' if rotate else 'Device token issued.' elif path in (f'{ADMIN_PREFIX}/devices/revoke', f'{ADMIN_PREFIX}/devices/unrevoke'): fields = await _form_fields( request, service, {'csrf_token', 'device_key', 'operation_id'}, ) revoked = path.endswith('/revoke') result = await run_in_threadpool( service.set_device_revoked, fields['device_key'], revoked, request.state.admin_actor, _parse_operation_id(fields['operation_id']), ) if not result: raise AdminAPIError(404, 'device was not found') return _redirect('../') elif path == f'{ADMIN_PREFIX}/queue/requeue': fields = await _form_fields( request, service, {'csrf_token', 'queue_ids', 'operation_id'}, ) queue_ids = _parse_queue_ids(fields['queue_ids'], service.requeue_limit) count = await run_in_threadpool( service.requeue, queue_ids, request.state.admin_actor, _parse_operation_id(fields['operation_id']), ) if type(count) is not int or count < 0: raise AdminAPIError(500, 'queue operation returned an invalid result') return _redirect('../') elif path == f'{ADMIN_PREFIX}/queue/discard-source': fields = await _form_fields( request, service, {'csrf_token', 'source', 'confirm_source', 'operation_id'}, ) source = _validate_key(fields['source'], 'source') if fields['confirm_source'] != source: raise AdminAPIError(400, 'source confirmation does not match') count = await run_in_threadpool( service.discard_source_queue, source, request.state.admin_actor, _parse_operation_id(fields['operation_id']), ) if type(count) is not int or count < 0: raise AdminAPIError(500, 'queue operation returned an invalid result') return _redirect('../') else: raise AdminAPIError(404, 'Not Found') snapshot = await asyncio.to_thread(service.workers_dispatch_snapshot) return _secure_response( _render_workers_page( service, snapshot, notice=notice, issued_token=issued_token, relative_root='..', ), media_type='text/html', ) async def admin_endpoint(request): service = request.app.state.admin_service actor = _trusted_operator(request, service) if actor is None: return _secure_response('Not Found', status_code=404) request.state.admin_actor = actor try: return await _dispatch(request, service) except AdminAPIError as exc: return _secure_response(str(exc), status_code=exc.status_code) except Exception as exc: logger.error('Admin request failed: %s', type(exc).__name__) return _secure_response('Internal Server Error', status_code=500) class _AdminRoute(Route): def matches(self, scope): match, child_scope = super().matches(scope) if match == Match.PARTIAL and scope.get('type') == 'http': return Match.FULL, child_scope return match, child_scope async def handle(self, scope, receive, send): await self.app(scope, receive, send) def admin_routes(): methods = ['GET', 'POST', 'PUT', 'PATCH', 'DELETE', 'OPTIONS', 'HEAD', 'TRACE'] return [ _AdminRoute(ADMIN_PREFIX, admin_endpoint, methods=methods), _AdminRoute(f'{ADMIN_PREFIX}/{{admin_path:path}}', admin_endpoint, methods=methods), ]