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'{html.escape(label)} '
)
return (
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 '' + ''.join(
f'{html.escape(label)} '
for key, href, label in links
) + ' '
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''
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 += (
' Older operations '
)
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' '
)
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''
next_before = page.get('next_before_event_id')
pagination = f'Newest events '
if next_before is not None:
pagination += (
' Older events '
)
content = (
'Append-only events '
'Newest first. Each page is bounded and uses a stable event cursor.
'
+ table + f' '
)
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 = (
'Relative file path'
f' '
)
expected_input = (
'Expected SHA-256'
' '
)
content_input = (
'URL-safe Base64 content'
' '
)
if root.permissions.allow_create_replace:
forms.append(
''
)
forms.append(
''
)
if root.permissions.allow_delete:
forms.append(
''
)
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 = (
'Name '
'Kind Bytes '
'' + ''.join(entries) + '
'
)
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 = ['All '] if include_empty else []
for value, title in values:
current = ' selected' if selected.get(name) == value else ''
options.append(
f'{html.escape(title)} '
)
return (
f'{html.escape(label)}'
+ ''.join(options) + ' '
)
def _render_worker_filters(selected, relative_root):
text_fields = (
('source', 'Source'), ('worker', 'Worker / device'), ('phase', 'Phase'),
('category', 'Category'), ('code', 'Stable code'),
)
controls = ''.join(
f'{html.escape(label)} '
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''
)
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''
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''
+ ''.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}
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''
)
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),
]