feat(operations): add runtime work status contract

This commit is contained in:
2026-08-19 20:52:44 +02:00
parent 94c94fefb4
commit 4cf2bfeb3e
6 changed files with 562 additions and 2 deletions
+10
View File
@@ -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
+8 -1
View File
@@ -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
+83 -1
View File
@@ -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
+29
View File
@@ -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:
+258
View File
@@ -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
+174
View File
@@ -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()