From 4cf2bfeb3e4326420476f36c602b2c8c25f04b3a Mon Sep 17 00:00:00 2001 From: Albrecht Degering Date: Wed, 19 Aug 2026 20:52:44 +0200 Subject: [PATCH] feat(operations): add runtime work status contract --- docs/DEPLOYMENT_OPERATOR_GUIDE.md | 10 + src/govoplan_core/core/modules.py | 9 +- src/govoplan_core/core/operations.py | 84 +++++++- src/govoplan_core/core/registry.py | 29 +++ src/govoplan_core/core/runtime_work.py | 258 +++++++++++++++++++++++++ tests/test_runtime_work.py | 174 +++++++++++++++++ 6 files changed, 562 insertions(+), 2 deletions(-) create mode 100644 src/govoplan_core/core/runtime_work.py create mode 100644 tests/test_runtime_work.py diff --git a/docs/DEPLOYMENT_OPERATOR_GUIDE.md b/docs/DEPLOYMENT_OPERATOR_GUIDE.md index db7af89..ce6a497 100644 --- a/docs/DEPLOYMENT_OPERATOR_GUIDE.md +++ b/docs/DEPLOYMENT_OPERATOR_GUIDE.md @@ -7,6 +7,16 @@ files. ## Runtime Configuration Contract +Worker and queue observability is provider-neutral. Runtime modules register a +bounded `RuntimeWorkStatusProviderRegistration` with Core; the Ops module +projects its sanitized status without importing Celery, Redis, or module job +implementations. Providers must use explicit `null` values for unsupported +queue depth, active/reserved work, failure count, and heartbeat evidence. An +unavailable metric must never be interpreted as zero or as proof of health. +The standard Core adapter reports the configured Celery/Redis runtime and +combines its bounded inspection result with registered worker heartbeat and +stale-threshold evidence. + Self-hosted installability follows the staged approach documented in `SELF_HOSTED_INSTALLABILITY.md`: generate an explicit env template, validate it, run production-like rehearsal with Compose-backed dependencies, then use the diff --git a/src/govoplan_core/core/modules.py b/src/govoplan_core/core/modules.py index 38c1b49..17cfc2a 100644 --- a/src/govoplan_core/core/modules.py +++ b/src/govoplan_core/core/modules.py @@ -15,7 +15,10 @@ from govoplan_core.core.views import ViewSurface if TYPE_CHECKING: from fastapi import APIRouter - from govoplan_core.core.operations import OperationalCheckProviderRegistration + from govoplan_core.core.operations import ( + OperationalCheckProviderRegistration, + RuntimeWorkStatusProviderRegistration, + ) from govoplan_core.core.search import ( SearchProviderRegistration, SearchSourceProviderRegistration, @@ -502,6 +505,10 @@ class ModuleManifest: "OperationalCheckProviderRegistration", ..., ] = () + runtime_work_status_providers: tuple[ + "RuntimeWorkStatusProviderRegistration", + ..., + ] = () architecture: ModuleArchitectureDeclaration | None = None information_governance: ModuleInformationGovernance = field( default_factory=ModuleInformationGovernance diff --git a/src/govoplan_core/core/operations.py b/src/govoplan_core/core/operations.py index 866127c..c9eb04d 100644 --- a/src/govoplan_core/core/operations.py +++ b/src/govoplan_core/core/operations.py @@ -1,7 +1,8 @@ from __future__ import annotations -from collections.abc import Callable, Mapping +from collections.abc import Callable, Mapping, Sequence from dataclasses import dataclass, field +from datetime import datetime from typing import Literal @@ -42,3 +43,84 @@ class OperationalCheckProviderRegistration: provider: OperationalCheckProvider cache_seconds: int = 60 + +RuntimeWorkState = Literal[ + "disabled", + "unconfigured", + "starting", + "healthy", + "idle", + "busy", + "degraded", + "stale", + "unreachable", +] + + +@dataclass(frozen=True, slots=True) +class RuntimeWorkStatusContext: + """Sanitized process evidence supplied to a runtime-work provider.""" + + profile: str + observed_at: datetime + stale_after_seconds: int + runtime_nodes: Sequence[Mapping[str, object]] = () + + +@dataclass(frozen=True, slots=True) +class RuntimeWorkStatus: + """Bounded worker/queue status with explicit unsupported metrics.""" + + provider_id: str + label: str + backend: str + enabled: bool + configured: bool + state: RuntimeWorkState + detail: str + observed_at: datetime + active_workers: int | None = None + last_heartbeat_at: datetime | None = None + queue_depths: Mapping[str, int | None] = field(default_factory=dict) + active_work: int | None = None + reserved_work: int | None = None + failures: int | None = None + stale_after_seconds: int | None = None + guidance: str = "" + + def as_dict(self) -> dict[str, object]: + return { + "provider_id": self.provider_id, + "label": self.label, + "backend": self.backend, + "enabled": self.enabled, + "configured": self.configured, + "state": self.state, + "detail": self.detail, + "observed_at": self.observed_at.isoformat(), + "active_workers": self.active_workers, + "last_heartbeat_at": ( + self.last_heartbeat_at.isoformat() + if self.last_heartbeat_at is not None + else None + ), + "queue_depths": dict(self.queue_depths), + "active_work": self.active_work, + "reserved_work": self.reserved_work, + "failures": self.failures, + "stale_after_seconds": self.stale_after_seconds, + "guidance": self.guidance, + } + + +RuntimeWorkStatusProvider = Callable[[RuntimeWorkStatusContext], RuntimeWorkStatus] + + +@dataclass(frozen=True, slots=True) +class RuntimeWorkStatusProviderRegistration: + """Register one optional provider-neutral worker/queue observation.""" + + module_id: str + provider_id: str + provider: RuntimeWorkStatusProvider + cache_seconds: int = 15 diff --git a/src/govoplan_core/core/registry.py b/src/govoplan_core/core/registry.py index 9f9e0ac..1575187 100644 --- a/src/govoplan_core/core/registry.py +++ b/src/govoplan_core/core/registry.py @@ -640,6 +640,7 @@ class PlatformRegistry: permissions = _collect_manifest_permissions(ordered) _validate_public_frontend_route_uniqueness(ordered) _validate_presentation_catalog(ordered) + _validate_runtime_work_provider_uniqueness(ordered) _validate_interface_closure(ordered) _validate_role_template_scopes(ordered, known_scopes=set(permissions)) @@ -692,6 +693,34 @@ class PlatformRegistry: return (self._manifests[module_id] for module_id in ordered) +def _validate_runtime_work_provider_uniqueness( + manifests: tuple[ModuleManifest, ...], +) -> None: + owners: dict[str, str] = {} + for manifest in manifests: + for registration in manifest.runtime_work_status_providers: + if registration.module_id != manifest.id: + raise RegistryError( + f"Runtime-work provider {registration.provider_id!r} belongs to " + f"{registration.module_id!r}, not module {manifest.id!r}" + ) + if not _QUICK_ACCESS_TOOL_ID_RE.fullmatch(registration.provider_id): + raise RegistryError( + f"Invalid runtime-work provider id: {registration.provider_id!r}" + ) + if registration.cache_seconds < 1: + raise RegistryError( + f"Runtime-work provider {registration.provider_id!r} must cache for at least one second" + ) + previous = owners.get(registration.provider_id) + if previous is not None: + raise RegistryError( + f"Duplicate runtime-work provider {registration.provider_id!r} in modules " + f"{previous!r} and {manifest.id!r}" + ) + owners[registration.provider_id] = manifest.id + + def manifest_view_surfaces(manifest: ModuleManifest) -> tuple[ViewSurface, ...]: frontend = manifest.frontend if frontend is None: diff --git a/src/govoplan_core/core/runtime_work.py b/src/govoplan_core/core/runtime_work.py new file mode 100644 index 0000000..c1f9816 --- /dev/null +++ b/src/govoplan_core/core/runtime_work.py @@ -0,0 +1,258 @@ +from __future__ import annotations + +from collections.abc import Callable, Mapping +from datetime import datetime +from typing import Any + +from govoplan_core.core.operations import ( + RuntimeWorkState, + RuntimeWorkStatus, + RuntimeWorkStatusContext, +) +from govoplan_core.settings import settings + + +QueueDepthReader = Callable[[list[str]], Mapping[str, int | None]] + + +def celery_runtime_work_status(context: RuntimeWorkStatusContext) -> RuntimeWorkStatus: + """Observe the built-in Celery backend without exposing it to Ops.""" + + queues = _configured_queues() + backend_configured = bool(str(settings.redis_url or "").strip()) + if not settings.celery_enabled: + return _status( + context, + enabled=False, + configured=backend_configured, + state="disabled", + detail="Background workers are intentionally disabled.", + active_workers=0, + queue_depths={queue: None for queue in queues}, + active_work=0, + reserved_work=0, + guidance=( + "Synchronous development paths may be used in development. " + "Enable and monitor workers before production queue-backed work." + if context.profile in {"development", "local-dev"} + else "Enable a worker backend before accepting queue-backed work." + ), + ) + if not backend_configured or not queues: + return _status( + context, + enabled=True, + configured=False, + state="unconfigured", + detail="Background workers are enabled but their backend or queue list is not configured.", + queue_depths={queue: None for queue in queues}, + guidance="Configure the broker and an explicit queue list, then start the required worker pools.", + ) + + try: + from govoplan_core.celery_app import celery + + inspector = celery.control.inspect(timeout=0.75) + return collect_celery_runtime_work_status( + context, + inspector=inspector, + queues=queues, + queue_depth_reader=_redis_queue_depths, + ) + except Exception: # noqa: BLE001 - status must isolate and sanitize provider failures. + return _status( + context, + enabled=True, + configured=True, + state="unreachable", + detail="The configured worker backend did not return bounded status evidence.", + queue_depths={queue: None for queue in queues}, + guidance="Verify broker reachability and worker processes; do not infer health from missing metrics.", + ) + + +def collect_celery_runtime_work_status( + context: RuntimeWorkStatusContext, + *, + inspector: Any, + queues: list[str], + queue_depth_reader: QueueDepthReader, +) -> RuntimeWorkStatus: + replies = inspector.ping() or {} + if not isinstance(replies, Mapping) or not replies: + state = "starting" if _fresh_worker_nodes(context) else "unreachable" + return _status( + context, + enabled=True, + configured=True, + state=state, + detail=( + "Worker processes are starting but have not answered the bounded status probe." + if state == "starting" + else "No configured worker answered the bounded status probe." + ), + active_workers=0, + queue_depths={queue: None for queue in queues}, + guidance="Wait for startup or verify worker and broker connectivity.", + ) + + active_queues_by_worker = inspector.active_queues() or {} + active_by_worker = inspector.active() or {} + reserved_by_worker = inspector.reserved() or {} + active_queues = sorted( + { + str(queue.get("name")) + for worker_queues in active_queues_by_worker.values() + if isinstance(worker_queues, list) + for queue in worker_queues + if isinstance(queue, Mapping) and queue.get("name") + } + ) + missing_queues = sorted(set(queues) - set(active_queues)) + active_work = _task_count(active_by_worker) + reserved_work = _task_count(reserved_by_worker) + try: + measured_depths = dict(queue_depth_reader(queues)) + except Exception: # noqa: BLE001 - queue depth remains explicitly unsupported. + measured_depths = {} + queue_depths = { + queue: _bounded_count(measured_depths.get(queue)) for queue in queues + } + known_depth = sum(value for value in queue_depths.values() if value is not None) + worker_nodes = _worker_nodes(context) + stale_nodes = [node for node in worker_nodes if node.get("stale") is True] + latest_heartbeat = _latest_heartbeat(worker_nodes) + + if stale_nodes and len(stale_nodes) >= len(worker_nodes) > 0: + state = "stale" + detail = "All registered worker heartbeats are stale." + guidance = "Restore worker heartbeats or replace the stale worker incarnations." + elif missing_queues or stale_nodes: + state = "degraded" + detail = "Worker status is partial: a queue lacks a consumer or a registered worker is stale." + guidance = "Restore the missing queue consumers and investigate stale worker heartbeats." + elif active_work + reserved_work + known_depth > 0: + state = "busy" + detail = "Workers are processing or waiting to process queued work." + guidance = "Monitor queue age and failures; scale only within configured provider limits." + elif all(value is not None for value in queue_depths.values()): + state = "idle" + detail = "Workers are available and all measured queues are empty." + guidance = "No action is required." + else: + state = "healthy" + detail = "Workers answered and cover all configured queues; queue depth is unavailable." + guidance = "Treat queue depth as unavailable, not empty." + + return _status( + context, + enabled=True, + configured=True, + state=state, + detail=detail, + active_workers=len(replies), + last_heartbeat_at=latest_heartbeat, + queue_depths=queue_depths, + active_work=active_work, + reserved_work=reserved_work, + guidance=guidance, + ) + + +def _status( + context: RuntimeWorkStatusContext, + *, + enabled: bool, + configured: bool, + state: RuntimeWorkState, + detail: str, + active_workers: int | None = None, + last_heartbeat_at: datetime | None = None, + queue_depths: Mapping[str, int | None] | None = None, + active_work: int | None = None, + reserved_work: int | None = None, + guidance: str, +) -> RuntimeWorkStatus: + return RuntimeWorkStatus( + provider_id="core.celery", + label="Background workers", + backend="Celery", + enabled=enabled, + configured=configured, + state=state, + detail=detail, + observed_at=context.observed_at, + active_workers=active_workers, + last_heartbeat_at=last_heartbeat_at, + queue_depths=queue_depths or {}, + active_work=active_work, + reserved_work=reserved_work, + failures=None, + stale_after_seconds=context.stale_after_seconds, + guidance=guidance, + ) + + +def _configured_queues() -> list[str]: + return sorted( + { + item.strip() + for item in str(settings.celery_queues or "").split(",") + if item.strip() + } + ) + + +def _redis_queue_depths(queues: list[str]) -> Mapping[str, int | None]: + from redis import Redis + + client = Redis.from_url( + settings.redis_url, + socket_connect_timeout=0.75, + socket_timeout=0.75, + ) + pipeline = client.pipeline(transaction=False) + for queue in queues: + pipeline.llen(queue) + values = pipeline.execute() + return { + queue: _bounded_count(value) + for queue, value in zip(queues, values, strict=True) + } + + +def _task_count(tasks_by_worker: object) -> int: + if not isinstance(tasks_by_worker, Mapping): + return 0 + return sum( + len(tasks) for tasks in tasks_by_worker.values() if isinstance(tasks, list) + ) + + +def _bounded_count(value: object) -> int | None: + if isinstance(value, bool) or not isinstance(value, int | float): + return None + return max(0, int(value)) + + +def _worker_nodes(context: RuntimeWorkStatusContext) -> list[Mapping[str, object]]: + return [node for node in context.runtime_nodes if node.get("role") == "worker"] + + +def _fresh_worker_nodes(context: RuntimeWorkStatusContext) -> list[Mapping[str, object]]: + return [node for node in _worker_nodes(context) if node.get("stale") is not True] + + +def _latest_heartbeat(nodes: list[Mapping[str, object]]) -> datetime | None: + values: list[datetime] = [] + for node in nodes: + raw = node.get("last_heartbeat_at") + if isinstance(raw, datetime): + values.append(raw) + continue + if isinstance(raw, str): + try: + values.append(datetime.fromisoformat(raw.replace("Z", "+00:00"))) + except ValueError: + continue + return max(values) if values else None diff --git a/tests/test_runtime_work.py b/tests/test_runtime_work.py new file mode 100644 index 0000000..eaf1e0e --- /dev/null +++ b/tests/test_runtime_work.py @@ -0,0 +1,174 @@ +from __future__ import annotations + +from datetime import UTC, datetime, timedelta + +import pytest + +from govoplan_core.core.modules import ModuleManifest +from govoplan_core.core.operations import ( + RuntimeWorkStatusContext, + RuntimeWorkStatusProviderRegistration, +) +from govoplan_core.core.registry import PlatformRegistry, RegistryError +from govoplan_core.core.runtime_work import collect_celery_runtime_work_status + + +class Inspector: + def __init__( + self, + *, + replies: dict[str, object] | None = None, + queues: dict[str, list[dict[str, str]]] | None = None, + active: dict[str, list[object]] | None = None, + reserved: dict[str, list[object]] | None = None, + ) -> None: + self._replies = replies or {} + self._queues = queues or {} + self._active = active or {} + self._reserved = reserved or {} + + def ping(self): + return self._replies + + def active_queues(self): + return self._queues + + def active(self): + return self._active + + def reserved(self): + return self._reserved + + +def _context(*, stale: bool = False, worker: bool = True) -> RuntimeWorkStatusContext: + now = datetime.now(UTC) + nodes = ( + { + "role": "worker", + "stale": stale, + "last_heartbeat_at": (now - timedelta(seconds=12)).isoformat(), + }, + ) if worker else () + return RuntimeWorkStatusContext( + profile="production", + observed_at=now, + stale_after_seconds=60, + runtime_nodes=nodes, + ) + + +def test_unknown_queue_depth_is_healthy_not_inferred_idle() -> None: + status = collect_celery_runtime_work_status( + _context(), + inspector=Inspector( + replies={"worker-1": {"ok": "pong"}}, + queues={"worker-1": [{"name": "mail"}]}, + ), + queues=["mail"], + queue_depth_reader=lambda queues: {}, + ) + + assert status.state == "healthy" + assert status.queue_depths == {"mail": None} + assert status.failures is None + assert status.last_heartbeat_at is not None + + +def test_busy_and_idle_require_measured_evidence() -> None: + busy = collect_celery_runtime_work_status( + _context(), + inspector=Inspector( + replies={"worker-1": {"ok": "pong"}}, + queues={"worker-1": [{"name": "mail"}]}, + active={"worker-1": [{"id": "task-1"}]}, + ), + queues=["mail"], + queue_depth_reader=lambda queues: {"mail": 2}, + ) + idle = collect_celery_runtime_work_status( + _context(), + inspector=Inspector( + replies={"worker-1": {"ok": "pong"}}, + queues={"worker-1": [{"name": "mail"}]}, + ), + queues=["mail"], + queue_depth_reader=lambda queues: {"mail": 0}, + ) + + assert busy.state == "busy" + assert busy.active_work == 1 + assert idle.state == "idle" + + +def test_partial_and_stale_worker_evidence_are_not_healthy() -> None: + degraded = collect_celery_runtime_work_status( + _context(), + inspector=Inspector( + replies={"worker-1": {"ok": "pong"}}, + queues={"worker-1": [{"name": "mail"}]}, + ), + queues=["mail", "calendar"], + queue_depth_reader=lambda queues: {queue: 0 for queue in queues}, + ) + stale = collect_celery_runtime_work_status( + _context(stale=True), + inspector=Inspector( + replies={"worker-1": {"ok": "pong"}}, + queues={"worker-1": [{"name": "mail"}]}, + ), + queues=["mail"], + queue_depth_reader=lambda queues: {"mail": 0}, + ) + + assert degraded.state == "degraded" + assert stale.state == "stale" + + +def test_no_reply_distinguishes_starting_from_unreachable() -> None: + starting = collect_celery_runtime_work_status( + _context(), + inspector=Inspector(), + queues=["mail"], + queue_depth_reader=lambda queues: {}, + ) + unreachable = collect_celery_runtime_work_status( + _context(worker=False), + inspector=Inspector(), + queues=["mail"], + queue_depth_reader=lambda queues: {}, + ) + + assert starting.state == "starting" + assert unreachable.state == "unreachable" + + +def test_registry_rejects_duplicate_runtime_work_provider_ids() -> None: + def provider(context: RuntimeWorkStatusContext): + raise AssertionError(context) + + def registration(module_id: str) -> RuntimeWorkStatusProviderRegistration: + return RuntimeWorkStatusProviderRegistration( + module_id=module_id, + provider_id="example.queue", + provider=provider, + ) + registry = PlatformRegistry() + registry.register( + ModuleManifest( + id="alpha", + name="Alpha", + version="1", + runtime_work_status_providers=(registration("alpha"),), + ) + ) + registry.register( + ModuleManifest( + id="beta", + name="Beta", + version="1", + runtime_work_status_providers=(registration("beta"),), + ) + ) + + with pytest.raises(RegistryError, match="Duplicate runtime-work provider"): + registry.validate()