4556 lines
205 KiB
Python
4556 lines
205 KiB
Python
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'<label>{html.escape(label)}<input name="{html.escape(name)}" '
|
|
f'type="{html.escape(input_type)}"{bounds}{value_attribute} required></label>'
|
|
)
|
|
return (
|
|
f'<form method="post" action="{html.escape(relative_root + action)}">'
|
|
f'<h3>{html.escape(title)}</h3>'
|
|
f'<input type="hidden" name="csrf_token" value="{html.escape(csrf_token)}">'
|
|
+ ''.join(
|
|
f'<input type="hidden" name="{html.escape(name)}" value="{html.escape(str(value))}">'
|
|
for name, value in hidden
|
|
)
|
|
+ ''.join(controls)
|
|
+ f'<button type="submit">{html.escape(title)}</button></form>'
|
|
)
|
|
|
|
|
|
def _table(columns, rows):
|
|
header = ''.join(f'<th scope="col">{html.escape(label)}</th>' for key, label in columns)
|
|
rendered_rows = []
|
|
for row in rows:
|
|
rendered_rows.append('<tr>' + ''.join(
|
|
f'<td>{html.escape(str(row.get(key) if row.get(key) is not None else ""))}</td>'
|
|
for key, _label in columns
|
|
) + '</tr>')
|
|
return f'<table><thead><tr>{header}</tr></thead><tbody>{"".join(rendered_rows)}</tbody></table>'
|
|
|
|
|
|
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 '<nav aria-label="Administration">' + ''.join(
|
|
f'<a href="{html.escape(href)}"'
|
|
f'{" aria-current=\"page\"" if key == active else ""}>{html.escape(label)}</a>'
|
|
for key, href, label in links
|
|
) + '</nav>'
|
|
|
|
|
|
def _page_shell(title, subtitle, active, content, relative_root, script=None):
|
|
script_tag = (
|
|
f'<script src="{html.escape(script)}" defer></script>' if script else ''
|
|
)
|
|
return f'''<!doctype html>
|
|
<html lang="en"><head><meta charset="utf-8"><meta name="viewport" content="width=device-width,initial-scale=1">
|
|
<title>{html.escape(title)}</title><link rel="stylesheet" href="{html.escape(relative_root + '/admin.css')}">{script_tag}</head>
|
|
<body><header><h1>{html.escape(title)}</h1><p>{html.escape(subtitle)}</p>
|
|
{_admin_navigation(active, relative_root)}</header><main>{content}</main></body></html>'''
|
|
|
|
|
|
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'<a href="{href}"><code>{operation_id}</code></a>'
|
|
|
|
|
|
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('<tr>' + ''.join(f'<td>{cell}</td>' for cell in cells) + '</tr>')
|
|
header = ''.join(f'<th scope="col">{html.escape(label)}</th>' for label in headers)
|
|
return f'<table><thead><tr>{header}</tr></thead><tbody>{"".join(rows)}</tbody></table>'
|
|
|
|
|
|
def _render_operations_page(page, relative_root='.'):
|
|
operations = list(page.get('operations') or [])
|
|
pagination = f'<a href="{html.escape(relative_root + "/operations")}">Newest operations</a>'
|
|
next_before = page.get('next_before')
|
|
if next_before is not None:
|
|
pagination += (
|
|
' <a rel="next" href="?before='
|
|
+ html.escape(_operation_cursor(*next_before)) + '">Older operations</a>'
|
|
)
|
|
content = (
|
|
'<section><h2>Recent durable operations</h2>'
|
|
'<p>Newest first. Open an operation to inspect its persisted status after a runtime restart.</p>'
|
|
+ _render_operations_table(operations, relative_root)
|
|
+ f'<p class="pagination">{pagination}</p></section>'
|
|
)
|
|
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 = (
|
|
'<section><h2>Durable status</h2>'
|
|
+ _table((('field', 'Field'), ('value', 'Value')), rows)
|
|
+ '</section><section><h2>Expected identity</h2>'
|
|
+ f'<pre>{expected}</pre></section>'
|
|
+ '<section><h2>Resulting identity</h2>'
|
|
+ f'<pre>{resulting}</pre></section>'
|
|
)
|
|
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 '')),
|
|
'<code>' + html.escape(_identity_json(event.get('before_identity'))) + '</code>',
|
|
'<code>' + html.escape(_identity_json(event.get('after_identity'))) + '</code>',
|
|
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('<tr>' + ''.join(f'<td>{cell}</td>' for cell in cells) + '</tr>')
|
|
header = ''.join(f'<th scope="col">{html.escape(label)}</th>' for label in headers)
|
|
table = f'<table><thead><tr>{header}</tr></thead><tbody>{"".join(rows)}</tbody></table>'
|
|
next_before = page.get('next_before_event_id')
|
|
pagination = f'<a href="{html.escape(relative_root + "/audit")}">Newest events</a>'
|
|
if next_before is not None:
|
|
pagination += (
|
|
' <a rel="next" href="?before='
|
|
+ html.escape(str(next_before)) + '">Older events</a>'
|
|
)
|
|
content = (
|
|
'<section><h2>Append-only events</h2>'
|
|
'<p>Newest first. Each page is bounded and uses a stable event cursor.</p>'
|
|
+ table + f'<p class="pagination">{pagination}</p></section>'
|
|
)
|
|
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'<input type="hidden" name="csrf_token" value="{html.escape(service.csrf_token)}">'
|
|
f'<input type="hidden" name="operation_id" value="{operation_id}">'
|
|
f'<input type="hidden" name="root_id" value="{root_id}">'
|
|
)
|
|
|
|
path_value = '' if relative_path is None else relative_path + '/'
|
|
path_input = (
|
|
'<label>Relative file path'
|
|
f'<input name="relative_path" type="text" value="{html.escape(path_value)}" '
|
|
'autocomplete="off" required></label>'
|
|
)
|
|
expected_input = (
|
|
'<label>Expected SHA-256'
|
|
'<input name="expected_sha256" type="text" minlength="64" maxlength="64" '
|
|
'pattern="[0-9a-f]{64}" autocomplete="off" required></label>'
|
|
)
|
|
content_input = (
|
|
'<label>URL-safe Base64 content'
|
|
'<textarea name="content_base64" rows="8" autocomplete="off" '
|
|
'spellcheck="false"></textarea></label>'
|
|
)
|
|
if root.permissions.allow_create_replace:
|
|
forms.append(
|
|
'<form method="post" action="files/create" autocomplete="off">'
|
|
'<h3>Create file</h3>' + hidden(str(uuid.uuid4())) + path_input
|
|
+ content_input + '<button type="submit">Create</button></form>'
|
|
)
|
|
forms.append(
|
|
'<form method="post" action="files/replace" autocomplete="off">'
|
|
'<h3>Replace file</h3>' + hidden(str(uuid.uuid4())) + path_input
|
|
+ expected_input + content_input
|
|
+ '<button type="submit">Replace</button></form>'
|
|
)
|
|
if root.permissions.allow_delete:
|
|
forms.append(
|
|
'<form method="post" action="files/delete" autocomplete="off">'
|
|
'<h3>Delete file</h3>' + hidden(str(uuid.uuid4())) + path_input
|
|
+ expected_input + '<button type="submit">Delete</button></form>'
|
|
)
|
|
if not forms:
|
|
return '<p>This root is read-only.</p>'
|
|
return (
|
|
'<p>Mutation content is canonical URL-safe Base64 and is bounded by the '
|
|
f'{service.max_body_bytes}-byte admin form limit.</p>'
|
|
'<div class="forms">' + ''.join(forms) + '</div>'
|
|
)
|
|
|
|
|
|
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(
|
|
'<li><a href="?'
|
|
+ html.escape(_managed_file_query(configured.root_id)) + '">'
|
|
+ html.escape(configured.root_id) + '</a></li>'
|
|
for configured in service.managed_file_roots.roots
|
|
)
|
|
content = (
|
|
'<section><h2>Logical roots</h2>'
|
|
'<p>Only configured logical roots are exposed. Host paths are never accepted.</p>'
|
|
+ roots + '<ul>' + root_links + '</ul></section>'
|
|
)
|
|
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 = (
|
|
'<a href="?'
|
|
+ html.escape(_managed_file_query(root.root_id, child_path))
|
|
+ '">' + html.escape(entry.name) + '/</a>'
|
|
)
|
|
elif root.permissions.allow_read:
|
|
name = (
|
|
'<a href="files/download?'
|
|
+ html.escape(_managed_file_query(root.root_id, child_path))
|
|
+ '">' + html.escape(entry.name) + '</a>'
|
|
)
|
|
else:
|
|
name = html.escape(entry.name)
|
|
entries.append(
|
|
'<tr><td>' + name + '</td><td>' + html.escape(entry.kind)
|
|
+ '</td><td>'
|
|
+ html.escape('' if entry.byte_count is None else str(entry.byte_count))
|
|
+ '</td></tr>'
|
|
)
|
|
listing_block = '<p>Listing is not permitted for this root.</p>'
|
|
if listing is not None:
|
|
listing_block = (
|
|
'<table><thead><tr><th scope="col">Name</th>'
|
|
'<th scope="col">Kind</th><th scope="col">Bytes</th></tr></thead>'
|
|
'<tbody>' + ''.join(entries) + '</tbody></table>'
|
|
)
|
|
parent = ''
|
|
if relative_path:
|
|
parent_path = relative_path.rpartition('/')[0] or None
|
|
parent = (
|
|
'<p><a href="?'
|
|
+ html.escape(_managed_file_query(root.root_id, parent_path))
|
|
+ '">Parent directory</a></p>'
|
|
)
|
|
content += (
|
|
'<section><h2>' + html.escape(location) + '</h2>' + parent
|
|
+ listing_block + '</section><section><h2>Typed mutations</h2>'
|
|
+ _render_managed_file_forms(service, root, relative_path)
|
|
+ '</section>'
|
|
)
|
|
return _page_shell(
|
|
'Files', 'Bounded descriptor-safe access by logical root.',
|
|
'files', content, '.',
|
|
)
|
|
|
|
|
|
def _filter_select(name, label, values, selected, *, include_empty=True):
|
|
options = ['<option value="">All</option>'] if include_empty else []
|
|
for value, title in values:
|
|
current = ' selected' if selected.get(name) == value else ''
|
|
options.append(
|
|
f'<option value="{html.escape(value)}"{current}>{html.escape(title)}</option>'
|
|
)
|
|
return (
|
|
f'<label>{html.escape(label)}<select name="{html.escape(name)}">'
|
|
+ ''.join(options) + '</select></label>'
|
|
)
|
|
|
|
|
|
def _render_worker_filters(selected, relative_root):
|
|
text_fields = (
|
|
('source', 'Source'), ('worker', 'Worker / device'), ('phase', 'Phase'),
|
|
('category', 'Category'), ('code', 'Stable code'),
|
|
)
|
|
controls = ''.join(
|
|
f'<label>{html.escape(label)}<input name="{name}" value="'
|
|
f'{html.escape(selected.get(name, ""))}"></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'<form method="get" action="{html.escape(relative_root + "/")}" '
|
|
f'class="filters"><div class="filter-grid">{controls}</div>'
|
|
'<div class="button-row"><button type="submit">Apply independent filters</button>'
|
|
f'<a class="button-link" href="{html.escape(relative_root + "/")}">Clear filters</a>'
|
|
'</div></form>'
|
|
)
|
|
|
|
|
|
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 = '<br>'.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 = '<br>'.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'<a href="{html.escape(relative_root + "/assignments/" + str(reservation_id))}">{reservation_id}</a>'),
|
|
('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'<tr data-reservation-id="{reservation_id}">'
|
|
+ ''.join(
|
|
f'<td data-field="{field}">{cell}</td>' for field, cell in cells
|
|
) + '</tr>'
|
|
)
|
|
header = ''.join(f'<th scope="col">{html.escape(value)}</th>' for value in headers)
|
|
return f'<table><thead><tr>{header}</tr></thead><tbody>{"".join(rows)}</tbody></table>'
|
|
|
|
|
|
def _render_diagnostic_groups(groups, relative_root):
|
|
if not groups:
|
|
return '<p>No diagnostic occurrences match the selected filters.</p>'
|
|
rendered = []
|
|
for group in groups:
|
|
occurrences = ''.join(
|
|
'<li>'
|
|
f'<a href="{html.escape(relative_root + "/assignments/" + str(item["reservation_id"]) + "#diagnostic-" + item["diagnostic_uid"])}">'
|
|
f'assignment {item["reservation_id"]}</a>: '
|
|
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())}'
|
|
'</li>'
|
|
for item in group['occurrences']
|
|
)
|
|
rendered.append(
|
|
'<article class="diagnostic-group">'
|
|
f'<h3><code>{html.escape(group["fingerprint"])}</code></h3>'
|
|
f'<p>{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"]))}.</p>'
|
|
f'<ol>{occurrences}</ol></article>'
|
|
)
|
|
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'<p>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.</p>'
|
|
)
|
|
|
|
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'<a class="button-link" href="{html.escape(relative_root + "/?" + urlencode(query))}">'
|
|
f'{html.escape(label)}</a>'
|
|
)
|
|
|
|
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 += '<div class="button-row">' + ''.join(links) + '</div>'
|
|
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'<p>Total source / phase / outcome groups: {total}. Showing '
|
|
f'{page_count} groups at offset {offset} with page limit {limit}.</p>'
|
|
)
|
|
|
|
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'<a class="button-link" href="{html.escape(relative_root + "/?" + urlencode(query))}">'
|
|
f'{html.escape(label)}</a>'
|
|
)
|
|
|
|
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 += '<div class="button-row">' + ''.join(links) + '</div>'
|
|
return status
|
|
|
|
|
|
def _render_deadline_policy(policy):
|
|
return (
|
|
'<p><strong>Future assignments only.</strong> Saving or applying a policy '
|
|
'does not alter deadlines on existing assignments.</p>'
|
|
+ _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 = (
|
|
'<section class="issued"><h2>Device token</h2>'
|
|
'<p>Shown once. Store it now; only its SHA-256 digest was persisted.</p>'
|
|
f'<code>{html.escape(issued_token)}</code></section>'
|
|
)
|
|
notice_block = f'<p class="notice">{html.escape(notice)}</p>' 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 = '<p class="unavailable">Worker administration snapshot unavailable.</p>'
|
|
else:
|
|
worker_note = ''
|
|
|
|
if control is None:
|
|
dispatch_content = '<p class="unavailable">Dispatch control unavailable.</p>'
|
|
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 = '<div class="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())),),
|
|
),
|
|
)) + '</div>'
|
|
|
|
if packages is None:
|
|
package_content = '<p class="unavailable">Package compatibility unavailable.</p>'
|
|
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 = (
|
|
'<h3>Trusted package profiles</h3>'
|
|
+ _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)
|
|
+ '<h3>Profile capabilities</h3>'
|
|
+ _table((
|
|
('profile_name', 'Profile'), ('source', 'Source'),
|
|
('platform', 'Worker platform'), ('planning_kind', 'Planning kind'),
|
|
), capability_rows)
|
|
+ '<h3>Required capabilities</h3>'
|
|
+ _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 = (
|
|
'<p>Not loaded. Select Assignments + diagnostics or All detail '
|
|
'sections in the filter above.</p>'
|
|
)
|
|
elif diagnostic_snapshot is None:
|
|
diagnostic_content = '<p class="unavailable">Diagnostic grouping unavailable.</p>'
|
|
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 = (
|
|
'<p>Not loaded. Select Assignments + durations or All detail '
|
|
'sections in the filter above.</p>'
|
|
)
|
|
elif metric_snapshot is None:
|
|
metrics_content = '<p class="unavailable">Duration metrics unavailable.</p>'
|
|
else:
|
|
metrics_content = _render_metric_status(
|
|
metric_snapshot, selected_filters, relative_root,
|
|
) + _render_duration_metrics(metric_snapshot.get('metrics') or [])
|
|
policy_content = (
|
|
'<p class="unavailable">Effective deadline policy unavailable.</p>'
|
|
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 = (
|
|
'<p><strong>Observability scope:</strong> 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.</p>'
|
|
)
|
|
content = f'''{notice_block}{token_block}<section><h2>Filters</h2>
|
|
<p>Every filter is independent. Diagnostic filters do not imply an assignment or scan outcome. Expensive diagnostic and duration sections are queried only when selected.</p>{worker_scope}{_render_worker_filters(selected_filters, relative_root)}</section>
|
|
<section><h2>Dispatch and drain</h2>
|
|
<p>Pausing or draining blocks new assignment commits. Authenticated status, terminal reports, uploads, receipt replay, expiry, ingestion, projection, and maintenance remain available.</p>{dispatch_content}{dispatch_forms}</section>
|
|
<section><h2>Package compatibility</h2><p>Protocol-1 workers are completion-only; new claims require a matching trusted protocol-2 package profile.</p>{package_content}</section>
|
|
<section><h2>Users</h2>{worker_note}{_table((('user_key', 'User'), ('active_assignment_cap', 'Cap'), ('disabled', 'Disabled')), users)}</section>
|
|
<section><h2>Workers</h2>{_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)}</section>
|
|
<section><h2>Assignments</h2><p>Assignment transport, scan outcome, diagnostics, progress, deadlines, slot/cap, package identity, ingestion, and projection are separate authoritative fields.</p>{assignment_table}</section>
|
|
<section><h2>Repeated diagnostics</h2><p>Grouping is deterministic and every bounded occurrence remains linked below its fingerprint.</p>{diagnostic_content}</section>
|
|
<section><h2>Observed durations</h2><p>Percentiles are observations only and never change policy automatically. Rows below five samples are explicitly insufficient.</p>{metrics_content}</section>
|
|
<section><h2>Effective deadline policy</h2>{policy_content}</section>
|
|
<section><h2>Deferred queue</h2>{_table((('queue_id', 'Queue'), ('source', 'Source'), ('target', 'Target'), ('available_after', 'Available after')), deferred)}</section>
|
|
<section><h2>Typed operations</h2><p>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.</p><div class="forms">{forms}</div></section>'''
|
|
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'<section><h4>{html.escape(name)}</h4><p>Not captured.</p></section>'
|
|
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'<section class="material"><h4>{html.escape(name)}</h4>'
|
|
+ _table((
|
|
('state', 'State'), ('encoding', 'Stored encoding'),
|
|
('original_size', 'Original bytes'), ('stored_size', 'Stored bytes'),
|
|
('sha256', 'Original SHA-256'), ('truncated', 'Truncated'),
|
|
), fields)
|
|
+ f'<pre>{stored}</pre></section>'
|
|
)
|
|
|
|
|
|
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
|
|
'<p class="unavailable">Current effective source policy unavailable.</p>'
|
|
)
|
|
|
|
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('<p>No body or process-log material was captured.</p>')
|
|
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'<article class="diagnostic" id="diagnostic-{html.escape(uid)}">'
|
|
f'<h3>{html.escape(str(row.get("category")))} / '
|
|
f'{html.escape(str(row.get("code")))}</h3>'
|
|
+ _table(tuple((key, key.replace('_', ' ').title()) for key in summary[0]), summary)
|
|
+ '<div class="button-row">'
|
|
f'<button type="button" data-copy-target="{target_id}">Copy canonical JSON</button>'
|
|
f'<a class="button-link" href="{html.escape(download)}" download>Download canonical JSON</a>'
|
|
'</div>'
|
|
f'<textarea class="canonical-json" id="{target_id}" readonly rows="12">'
|
|
f'{html.escape(row["canonical_envelope_json"])}</textarea>'
|
|
+ ''.join(materials) + '</article>'
|
|
)
|
|
if not diagnostics:
|
|
diagnostics.append(
|
|
'<p class="unavailable">No canonical diagnostic envelope is stored. '
|
|
f'Availability: {html.escape(detail["diagnostics"]["availability"])}.</p>'
|
|
)
|
|
|
|
legacy = detail['legacy_evidence']
|
|
legacy_rows = legacy.get('errors') or []
|
|
legacy_content = (
|
|
'<p>These fields are legacy evidence, not a synthesized diagnostic envelope.</p>'
|
|
+ _table((
|
|
('id', 'Error'), ('category', 'Category'), ('summary', 'Summary'),
|
|
('created_at', 'Created'),
|
|
), legacy_rows)
|
|
+ ''.join(
|
|
'<h4>Exact stored raw_error</h4>'
|
|
f'<pre>{html.escape(str(row.get("raw_error") or ""))}</pre>'
|
|
for row in legacy_rows
|
|
)
|
|
) if legacy.get('available') else (
|
|
'<p>No legacy error evidence is stored. No envelope has been invented.</p>'
|
|
)
|
|
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 = (
|
|
'<p class="unavailable">' + html.escape('; '.join(truncation_notes)) + '.</p>'
|
|
if truncation_notes else ''
|
|
)
|
|
content = f'''
|
|
<section><h2>Assignment</h2>
|
|
<p><a href="./{reservation_id}.json">Machine-readable detail</a></p>
|
|
{_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)}</section>
|
|
<section><h2>Progress authority</h2>{progress_summary}</section>
|
|
<section><h2>Ordered timeline</h2>{truncation}{timeline}</section>
|
|
<section><h2>Duration breakdown</h2>{durations}</section>
|
|
<section><h2>Transport, receipt, ingestion, settlement, projection</h2>{transport}</section>
|
|
<section><h2>Scan summary</h2>{scan}</section>
|
|
<section><h2>Package identity</h2>{package}</section>
|
|
<section><h2>Current effective source policy</h2><p>This is current policy context; the stored assignment deadline above remains immutable.</p>{current_policy_html}</section>
|
|
<section><h2>Diagnostics</h2>{''.join(diagnostics)}</section>
|
|
<section><h2>Legacy evidence</h2>{legacy_content}</section>'''
|
|
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 = '<p class="unavailable">Runtime snapshot unavailable.</p>'
|
|
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 = '<span class="unavailable">Queue snapshot unavailable.</span>'
|
|
|
|
if control is None:
|
|
control_rows = []
|
|
control_note = '<span class="unavailable">Control snapshot unavailable.</span>'
|
|
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 = '<span class="unavailable">Recent operations unavailable.</span>'
|
|
|
|
content = f'''
|
|
<section><h2>Runtime health</h2>{runtime_content}</section>
|
|
<section><h2>Discovery producers</h2>{_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)}</section>
|
|
<section><h2>Pipeline workers</h2>{_table((('id', 'ID'), ('role', 'Role'), ('lifecycle_state', 'Lifecycle'), ('process_state', 'Process'), ('desired_state', 'Desired'), ('safe_error_category', 'Safe error')), pipeline_workers)}</section>
|
|
<section><h2>Pipeline authority</h2>{_table((('metric', 'Metric'), ('state', 'State'), ('ready', 'Ready')), pipeline_rows)}</section>
|
|
<section><h2>Queue status</h2><p>{queue_note}</p>{_table((('status', 'Status'), ('count', 'Bounded count'), ('truncated', 'Truncated')), queue_rows)}</section>
|
|
<section><h2>Assignments, bundles and control</h2><p>{control_note}</p>{_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)}
|
|
<p>Bundle items: {html.escape(str(_safe_count(pipeline.get('bundle_items'))))}; bundle bytes: {html.escape(str(_safe_count(pipeline.get('bundle_bytes'))))}.</p></section>
|
|
<section><h2>Recent operation outcomes</h2><p>{operations_note}</p>{_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)}</section>'''
|
|
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 = '<p class="unavailable">Discovery control unavailable.</p>'
|
|
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 = '<div class="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())),
|
|
),
|
|
),
|
|
)) + '</div>'
|
|
|
|
if runtime is None:
|
|
producer_content = '<p class="unavailable">Producer state unavailable.</p>'
|
|
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 = '<div class="forms">' + ''.join(forms) + '</div>'
|
|
|
|
content = f'''
|
|
<section><h2>Persistent discovery control</h2><p>Stored in PostgreSQL and preserved across runtime restarts.</p>{control_content}{control_forms}</section>
|
|
<section><h2>Discovery producers</h2><p>Lifecycle and interval changes affect the current Supervisor runtime only.</p>{producer_content}{producer_forms}</section>'''
|
|
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', '<p class="unavailable">Supervisor state unavailable.</p>',
|
|
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 = '<section><h2>Runtime</h2>' + _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', ''),
|
|
},)) + '</section>'
|
|
content += '<section><h2>Dashboard</h2>' + _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 += (
|
|
'<details class="control-panel"><summary>'
|
|
+ html.escape(dashboard_summary)
|
|
+ '</summary><div class="forms">'
|
|
+ dashboard_forms
|
|
+ '</div></details></section>'
|
|
)
|
|
content += '<section><h2>Managed sources</h2>' + _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(
|
|
'<details class="control-panel"><summary>'
|
|
+ html.escape(' - '.join(summary_parts))
|
|
+ '</summary><div class="forms">'
|
|
+ ''.join(forms)
|
|
+ '</div></details>'
|
|
)
|
|
content += '<div class="control-panels">' + ''.join(panels) + '</div></section>'
|
|
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 = '<p class="unavailable">Managed source list unavailable.</p>'
|
|
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 = (
|
|
'<section><h2>Managed source logs</h2>'
|
|
'<p>Only the active bounded log for an allowlisted source is available.</p>'
|
|
'<div class="forms">'
|
|
+ ''.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)
|
|
+ '</div></section>'
|
|
)
|
|
if isinstance(tail, dict):
|
|
lines = '\n'.join(tail.get('lines') or [])
|
|
suffix = ' Response was truncated.' if tail.get('response_truncated') else ''
|
|
content += (
|
|
'<section><h2>Tail result</h2><p>'
|
|
+ html.escape(
|
|
f'{tail.get("source_id", "")} · {tail.get("line_count", 0)} lines.{suffix}'
|
|
)
|
|
+ '</p><pre>' + html.escape(lines) + '</pre></section>'
|
|
)
|
|
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'<input type="hidden" name="expected_{name}_sha256" value="{value}">'
|
|
for name, value in hashes.items()
|
|
)
|
|
operation_id = str(uuid.uuid4()) if operation_id is None else operation_id
|
|
editor_form = (
|
|
f'<form method="post" autocomplete="off" action="{html.escape(relative_root + "/" + document + "/preview")}">'
|
|
f'<input type="hidden" name="csrf_token" value="{html.escape(service.csrf_token)}">'
|
|
f'{hidden}<input type="hidden" name="operation_id" value="{html.escape(operation_id)}">'
|
|
f'<label>{html.escape(document.title() + " YAML")}<textarea name="document_text" '
|
|
f'rows="32" spellcheck="false" autocomplete="off" required>{html.escape(text)}</textarea></label>'
|
|
'<div class="button-row"><button type="submit">Preview</button>'
|
|
f'<button type="submit" formaction="{html.escape(relative_root + "/" + document + "/save")}">Save candidate</button></div>'
|
|
'</form>'
|
|
)
|
|
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'<p class="notice">{html.escape(notice)}</p>' if notice else '')
|
|
+ f'<section><h2>{html.escape(document.title())} candidate</h2>'
|
|
+ f'<p>Editing source: {html.escape(editor.source)}. Save stages a private candidate and never activates it.</p>'
|
|
+ _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 + '</section>'
|
|
)
|
|
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'<p>Truncated: {html.escape(str(diff.truncated))}; '
|
|
f'format-only change: {html.escape(str(diff.format_only_changed))}.</p>'
|
|
)
|
|
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 += '<section><h2>Validated preview</h2>' + diff_html + '</section>'
|
|
if document == 'config':
|
|
observability = dict(observability or {})
|
|
policy = _component_value(observability, 'policy', dict)
|
|
metrics = _component_value(observability, 'metrics', dict)
|
|
policy_html = (
|
|
'<p class="unavailable">Effective deadline policy unavailable.</p>'
|
|
if policy is None else _render_deadline_policy(policy)
|
|
)
|
|
metrics_html = (
|
|
'<p class="unavailable">Duration metrics unavailable.</p>'
|
|
if metrics is None else (
|
|
f'<p>Total source / phase / outcome groups: '
|
|
f'{int(metrics.get("total_group_count") or 0)}. '
|
|
+ (
|
|
'Additional groups are available through paginated Workers observability.</p>'
|
|
if metrics.get('has_next') else '</p>'
|
|
)
|
|
+ _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'<section><h2>{html.escape(candidate_label)}</h2>{policy_html}'
|
|
'<h3>Observed source / phase / outcome durations</h3>'
|
|
'<p>p50/p95/p99 values include sample counts and never mutate policy.</p>'
|
|
f'{metrics_html}</section>'
|
|
)
|
|
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'<div class="forms">{apply_controls}</div>' if apply_controls else
|
|
'<p class="unavailable">No staged candidate differs from the active document. '
|
|
'Save a changed candidate before applying.</p>'
|
|
)
|
|
content += (
|
|
'<section><h2>Apply</h2><p>Apply requires the fixed host agent. '
|
|
'The operation is hash-bound and durable before dispatch.</p>'
|
|
f'{apply_state}</section>'
|
|
)
|
|
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),
|
|
]
|