Compare commits

...
3 Commits
14 changed files with 541 additions and 38 deletions
+14
View File
@@ -823,6 +823,10 @@ WebUI modules receive only the core route context:
- `auth`
A module should call its own API client and module-owned backend routes. Shared API helpers should live in core only when they are truly platform-level concerns.
For ordinary JSON mutations, use Core's `apiPostJson` and `apiPatchJson`
helpers. They preserve the shared authentication, CSRF, error, and request
invalidation behavior while leaving endpoint types and feature semantics in the
owning module.
Modules can also contribute named UI capabilities for explicit extension
points. Capability values must be narrow, typed contracts, not imports from a
@@ -1005,6 +1009,16 @@ Decision: templates and reporting are separate modules.
and export targets
- report permissions, report execution history, generated report evidence, and
report-specific retention inputs
Cross-module reports use Core's versioned
`reporting.report_provider.<provider-id>` contract. Source modules own
authorization, parameters, source revisions, effective scope, result schema,
and privacy transforms; Reporting owns discovery, validation, governed
execution, provenance, export history, and the global `/reports` route. The
optional `policy.reporting_governance` capability can only tighten execution,
retention, export, and re-identification-risk handling. Reporting exposes
`reporting.retention` so the Policy-owned retention run can minimize expired
provider results without importing Reporting models.
- downstream export handoff to files, dataflow, connectors, or publication
surfaces
+7
View File
@@ -52,6 +52,13 @@ software version, module-composition hash, queues, start time, and heartbeat.
The registration identity includes a process incarnation so a stale process
cannot update a replacement's row.
Worker metadata also records the orchestrator pool and declared concurrency.
Every Celery prefork child disposes the SQLAlchemy pool inherited from its
parent and creates a process-local pool before handling work. Deployment
rendering must therefore budget one database pool for the worker parent and
each child. Ops compares active queue ownership, software versions, and the
order-independent module-composition hash with the graph loaded by the API.
Drain is durable operator intent:
- an API enters not-ready state after observing drain;
+24 -3
View File
@@ -3,9 +3,15 @@ from __future__ import annotations
from datetime import datetime, timedelta, timezone
from importlib.metadata import PackageNotFoundError, version
import logging
import os
from celery import Celery
from celery.signals import heartbeat_sent, worker_ready, worker_shutdown
from celery.signals import (
heartbeat_sent,
worker_process_init,
worker_ready,
worker_shutdown,
)
from govoplan_core.core.campaigns import (
CAPABILITY_CAMPAIGNS_DELIVERY_TASKS,
@@ -196,6 +202,21 @@ def _worker_runtime_identity(sender: object | None = None) -> RuntimeIdentity:
return _worker_identity
def _worker_metadata() -> dict[str, object]:
raw_concurrency = str(os.getenv("CELERY_WORKER_CONCURRENCY") or "").strip()
return {
"process": "celery-worker",
"worker_pool": str(os.getenv("GOVOPLAN_WORKER_POOL") or "default"),
"concurrency": int(raw_concurrency) if raw_concurrency.isdigit() else None,
}
@worker_process_init.connect
def _reset_worker_process_database(**_kwargs) -> None:
# SQLAlchemy pools must not be shared across prefork child processes.
configure_database(settings.database_url, dispose_previous=True)
@worker_ready.connect
def _register_worker_runtime(sender=None, **_kwargs) -> None:
global _worker_consumer
@@ -206,7 +227,7 @@ def _register_worker_runtime(sender=None, **_kwargs) -> None:
node = register_runtime_node(
session,
identity,
metadata={"process": "celery-worker"},
metadata=_worker_metadata(),
)
session.commit()
except Exception:
@@ -227,7 +248,7 @@ def _heartbeat_worker_runtime(sender=None, **_kwargs) -> None:
node = heartbeat_runtime_node(
session,
identity,
metadata={"process": "celery-worker"},
metadata=_worker_metadata(),
)
session.commit()
except Exception: # noqa: BLE001 - uncertain authority must fail closed
@@ -3,7 +3,6 @@ from __future__ import annotations
import base64
from collections.abc import Mapping, Sequence
from dataclasses import dataclass, field
from datetime import UTC, datetime
from pathlib import Path
import json
import os
@@ -21,6 +20,7 @@ from govoplan_core.core.module_package_catalog import (
_is_http_url,
_load_private_key,
_parse_trusted_keys,
_record_catalog_acceptance,
)
from govoplan_core.core.external_references import (
IntegrationMaturity,
@@ -858,33 +858,10 @@ def sign_configuration_package_catalog(*, path: Path, key_id: str, private_key_p
def record_configuration_package_catalog_acceptance(validation: dict[str, object]) -> None:
state_path = _configured_sequence_state_path()
if state_path is None or validation.get("valid") is not True:
return
channel = validation.get("channel")
sequence = validation.get("sequence")
if not isinstance(channel, str) or not isinstance(sequence, int):
return
try:
state = json.loads(state_path.read_text(encoding="utf-8")) if state_path.exists() else {}
except json.JSONDecodeError:
state = {}
if not isinstance(state, dict):
state = {}
channels = state.get("channels")
if not isinstance(channels, dict):
channels = {}
channel_state = channels.get(channel)
if not isinstance(channel_state, dict):
channel_state = {}
channel_state["last_sequence"] = max(int(channel_state.get("last_sequence") or 0), sequence)
channel_state["accepted_at"] = datetime.now(tz=UTC).isoformat().replace("+00:00", "Z")
channel_state["key_id"] = validation.get("key_id")
channel_state["source"] = validation.get("source") or validation.get("path")
channels[channel] = channel_state
state["channels"] = channels
state_path.parent.mkdir(parents=True, exist_ok=True)
state_path.write_text(json.dumps(state, indent=2, sort_keys=True) + "\n", encoding="utf-8")
_record_catalog_acceptance(
validation,
state_path=_configured_sequence_state_path(),
)
def configuration_package_claim_issues(
@@ -244,7 +244,17 @@ def sign_module_package_catalog(
def record_module_package_catalog_acceptance(validation: dict[str, object]) -> None:
state_path = _configured_sequence_state_path()
_record_catalog_acceptance(
validation,
state_path=_configured_sequence_state_path(),
)
def _record_catalog_acceptance(
validation: dict[str, object],
*,
state_path: Path | None,
) -> None:
if state_path is None or validation.get("valid") is not True:
return
channel = validation.get("channel")
+311
View File
@@ -0,0 +1,311 @@
"""Provider-neutral contracts for cross-module governed reports."""
from __future__ import annotations
from collections.abc import Mapping
from dataclasses import dataclass, field
from datetime import datetime
from typing import Literal, Protocol, runtime_checkable
REPORT_PROVIDER_CAPABILITY_PREFIX = "reporting.report_provider."
CAPABILITY_POLICY_REPORTING_GOVERNANCE = "policy.reporting_governance"
CAPABILITY_REPORTING_RETENTION = "reporting.retention"
REPORT_PROVIDER_CONTRACT_VERSION = "1.0"
ReportParameterType = Literal[
"string",
"integer",
"number",
"boolean",
"date",
"datetime",
"reference",
]
ReportFieldType = Literal[
"string",
"integer",
"number",
"boolean",
"date",
"datetime",
"object",
"suppressed_count",
]
ReidentificationRisk = Literal["low", "moderate", "high"]
ReportGovernanceAction = Literal["catalogue", "execute", "export"]
@dataclass(frozen=True, slots=True)
class ReportParameterOption:
value: str
label: str
description: str | None = None
def to_dict(self) -> dict[str, object]:
return {
"value": self.value,
"label": self.label,
"description": self.description,
}
@dataclass(frozen=True, slots=True)
class ReportParameterDescriptor:
key: str
label: str
type: ReportParameterType = "string"
required: bool = False
description: str | None = None
options_from_provider: bool = False
def to_dict(self) -> dict[str, object]:
return {
"key": self.key,
"label": self.label,
"type": self.type,
"required": self.required,
"description": self.description,
"options_from_provider": self.options_from_provider,
}
@dataclass(frozen=True, slots=True)
class ReportResultField:
path: str
label: str
type: ReportFieldType
group: str
nullable: bool = False
sensitive: bool = False
def to_dict(self) -> dict[str, object]:
return {
"path": self.path,
"label": self.label,
"type": self.type,
"group": self.group,
"nullable": self.nullable,
"sensitive": self.sensitive,
}
@dataclass(frozen=True, slots=True)
class ReportPrivacyTransform:
id: str
label: str
required: bool = True
def to_dict(self) -> dict[str, object]:
return {"id": self.id, "label": self.label, "required": self.required}
@dataclass(frozen=True, slots=True)
class ReportDescriptor:
provider_id: str
report_id: str
revision: str
title: str
summary: str
parameters: tuple[ReportParameterDescriptor, ...]
result_schema: tuple[ReportResultField, ...]
privacy_transforms: tuple[ReportPrivacyTransform, ...]
purpose_required: bool = True
audience_scope_required: bool = True
retention_class: str = "report_result"
export_formats: tuple[str, ...] = ("json",)
reidentification_risk: ReidentificationRisk = "moderate"
presentation: Mapping[str, object] = field(default_factory=dict)
contract_version: str = REPORT_PROVIDER_CONTRACT_VERSION
def to_dict(self) -> dict[str, object]:
return {
"contract_version": self.contract_version,
"provider_id": self.provider_id,
"report_id": self.report_id,
"revision": self.revision,
"title": self.title,
"summary": self.summary,
"parameters": [item.to_dict() for item in self.parameters],
"result_schema": [item.to_dict() for item in self.result_schema],
"privacy_transforms": [item.to_dict() for item in self.privacy_transforms],
"purpose_required": self.purpose_required,
"audience_scope_required": self.audience_scope_required,
"retention_class": self.retention_class,
"export_formats": list(self.export_formats),
"reidentification_risk": self.reidentification_risk,
"presentation": dict(self.presentation),
}
@dataclass(frozen=True, slots=True)
class ReportProviderRequest:
report_id: str
parameters: Mapping[str, object]
purpose: str
audience_scope: Mapping[str, object]
@dataclass(frozen=True, slots=True)
class ReportProviderResult:
report_id: str
generated_at: datetime
payload: Mapping[str, object]
source_revisions: tuple[Mapping[str, object], ...]
effective_scope: Mapping[str, object]
applied_privacy_transforms: tuple[str, ...]
provenance: Mapping[str, object]
@runtime_checkable
class ReportProvider(Protocol):
provider_id: str
contract_version: str
def list_reports(
self,
session: object,
principal: object,
) -> tuple[ReportDescriptor, ...]: ...
def parameter_options(
self,
session: object,
principal: object,
*,
report_id: str,
parameter_key: str,
query: str,
limit: int,
) -> tuple[ReportParameterOption, ...]: ...
def execute_report(
self,
session: object,
principal: object,
*,
request: ReportProviderRequest,
) -> ReportProviderResult: ...
def authorize_result(
self,
session: object,
principal: object,
*,
report_id: str,
source_revisions: tuple[Mapping[str, object], ...],
effective_scope: Mapping[str, object],
) -> bool: ...
@dataclass(frozen=True, slots=True)
class ReportingGovernanceRequest:
action: ReportGovernanceAction
tenant_id: str
provider_id: str
report_id: str
purpose: str | None
audience_scope: Mapping[str, object]
retention_class: str
export_format: str | None
reidentification_risk: ReidentificationRisk
declared_privacy_transforms: tuple[str, ...]
applied_privacy_transforms: tuple[str, ...] = ()
@dataclass(frozen=True, slots=True)
class ReportingGovernanceDecision:
allowed: bool
reason: str | None
retention_days: int | None
export_formats: tuple[str, ...]
required_privacy_transforms: tuple[str, ...]
provenance: Mapping[str, object] = field(default_factory=dict)
@runtime_checkable
class ReportingGovernanceProvider(Protocol):
def decide_reporting_action(
self,
session: object,
principal: object,
*,
request: ReportingGovernanceRequest,
) -> ReportingGovernanceDecision: ...
@runtime_checkable
class ReportingRetentionProvider(Protocol):
def apply_retention(
self,
session: object,
*,
dry_run: bool,
now: datetime,
limit: int = 500,
) -> Mapping[str, int]: ...
def report_providers(registry: object | None) -> tuple[tuple[str, ReportProvider], ...]:
if (
registry is None
or not hasattr(registry, "capability_names")
or not hasattr(registry, "capability")
):
return ()
providers: list[tuple[str, ReportProvider]] = []
for capability_name in registry.capability_names():
if not capability_name.startswith(REPORT_PROVIDER_CAPABILITY_PREFIX):
continue
provider_id = capability_name.removeprefix(REPORT_PROVIDER_CAPABILITY_PREFIX)
capability = registry.capability(capability_name)
if not isinstance(capability, ReportProvider):
raise TypeError(f"Invalid report provider capability: {capability_name}")
if capability.provider_id != provider_id:
raise ValueError(
f"Report provider id {capability.provider_id!r} does not match "
f"capability {capability_name!r}"
)
if capability.contract_version != REPORT_PROVIDER_CONTRACT_VERSION:
raise ValueError(
f"Unsupported report provider contract: {capability.contract_version}"
)
providers.append((provider_id, capability))
return tuple(providers)
def reporting_governance_provider(
registry: object | None,
) -> ReportingGovernanceProvider | None:
if (
registry is None
or not hasattr(registry, "has_capability")
or not registry.has_capability(CAPABILITY_POLICY_REPORTING_GOVERNANCE)
):
return None
capability = registry.capability(CAPABILITY_POLICY_REPORTING_GOVERNANCE)
if not isinstance(capability, ReportingGovernanceProvider):
raise TypeError("Invalid reporting governance provider capability")
return capability
__all__ = [
"CAPABILITY_POLICY_REPORTING_GOVERNANCE",
"CAPABILITY_REPORTING_RETENTION",
"REPORT_PROVIDER_CAPABILITY_PREFIX",
"REPORT_PROVIDER_CONTRACT_VERSION",
"ReportDescriptor",
"ReportParameterDescriptor",
"ReportParameterOption",
"ReportPrivacyTransform",
"ReportProvider",
"ReportProviderRequest",
"ReportProviderResult",
"ReportResultField",
"ReportingGovernanceDecision",
"ReportingGovernanceProvider",
"ReportingGovernanceRequest",
"ReportingRetentionProvider",
"report_providers",
"reporting_governance_provider",
]
@@ -166,9 +166,7 @@ def runtime_identity(
for item in str(getattr(settings, "celery_queues", "") or "").split(",")
if item.strip()
)
composition_hash = hashlib.sha256(
json.dumps(sorted(module_ids), separators=(",", ":")).encode("utf-8")
).hexdigest()
composition_hash = runtime_composition_hash(module_ids)
return RuntimeIdentity(
installation_id=installation_id,
node_id=effective_node_id,
@@ -180,6 +178,12 @@ def runtime_identity(
)
def runtime_composition_hash(module_ids: tuple[str, ...]) -> str:
return hashlib.sha256(
json.dumps(sorted(set(module_ids)), separators=(",", ":")).encode("utf-8")
).hexdigest()
def register_runtime_node(
session: Session,
identity: RuntimeIdentity,
+6 -1
View File
@@ -91,7 +91,12 @@ _default_database: DatabaseHandle | None = None
def configure_database(database_url: str, *, engine: Engine | None = None, dispose_previous: bool = False) -> DatabaseHandle:
global _default_database
if engine is None and _default_database is not None and _default_database.database_url == database_url:
if (
engine is None
and not dispose_previous
and _default_database is not None
and _default_database.database_url == database_url
):
return _default_database
previous_database = _default_database
_default_database = DatabaseHandle(database_url, engine=engine)
+24
View File
@@ -73,6 +73,30 @@ class Settings(BaseSettings):
ge=0,
le=86_400,
)
database_connection_limit: int | None = Field(
default=None,
alias="GOVOPLAN_DB_CONNECTION_LIMIT",
ge=10,
le=1_000_000,
)
database_connection_reserve: int = Field(
default=10,
alias="GOVOPLAN_DB_CONNECTION_RESERVE",
ge=1,
le=999_999,
)
database_connection_peak: int | None = Field(
default=None,
alias="GOVOPLAN_DB_CONNECTION_PEAK",
ge=1,
le=1_000_000,
)
database_connection_available: int | None = Field(
default=None,
alias="GOVOPLAN_DB_CONNECTION_AVAILABLE",
ge=1,
le=1_000_000,
)
access_database_url: str | None = Field(default=None, alias="ACCESS_DATABASE_URL")
access_db_schema: str | None = Field(default=None, alias="ACCESS_DB_SCHEMA")
access_table_prefix: str = Field(default="access_", alias="ACCESS_TABLE_PREFIX")
+18 -2
View File
@@ -232,7 +232,15 @@ class ModuleSystemTests(unittest.TestCase):
self.assertTrue(manifests["campaigns"].required_capabilities)
self.assertEqual(
manifests["campaigns"].optional_dependencies,
("files", "mail", "notifications", "addresses", "postbox", "approvals"),
(
"files",
"mail",
"notifications",
"addresses",
"postbox",
"approvals",
"reporting",
),
)
self.assertEqual(manifests["dashboard"].dependencies, ())
self.assertTrue(manifests["dashboard"].required_capabilities)
@@ -3436,7 +3444,15 @@ finally:
"version_max_exclusive": "0.3.0",
}, modules["campaigns"]["requires_interfaces"])
self.assertEqual(
["files", "mail", "notifications", "addresses", "postbox", "approvals"],
[
"files",
"mail",
"notifications",
"addresses",
"postbox",
"approvals",
"reporting",
],
modules["campaigns"]["optional_dependencies"],
)
self.assertEqual("requires_review", modules["files"]["migration_safety"])
+77
View File
@@ -0,0 +1,77 @@
from __future__ import annotations
from govoplan_core.core.modules import ModuleContext, ModuleManifest
from govoplan_core.core.registry import PlatformRegistry
from govoplan_core.core.reporting import (
REPORT_PROVIDER_CAPABILITY_PREFIX,
REPORT_PROVIDER_CONTRACT_VERSION,
ReportProviderRequest,
report_providers,
)
class _Provider:
provider_id = "example"
contract_version = REPORT_PROVIDER_CONTRACT_VERSION
def list_reports(self, session, principal):
return ()
def parameter_options(self, session, principal, **kwargs):
return ()
def execute_report(self, session, principal, *, request: ReportProviderRequest):
raise NotImplementedError
def authorize_result(self, session, principal, **kwargs):
return True
def test_report_provider_discovery_is_prefix_based_and_optional() -> None:
registry = PlatformRegistry()
registry.register(
ModuleManifest(
id="example",
name="Example",
version="1.0",
capability_factories={
REPORT_PROVIDER_CAPABILITY_PREFIX + "example": lambda _context: (
_Provider()
)
},
)
)
registry.configure_capability_context(
ModuleContext(registry=registry, settings=object())
)
providers = report_providers(registry)
assert len(providers) == 1
assert providers[0][0] == "example"
def test_report_provider_id_must_match_capability_name() -> None:
registry = PlatformRegistry()
registry.register(
ModuleManifest(
id="example",
name="Example",
version="1.0",
capability_factories={
REPORT_PROVIDER_CAPABILITY_PREFIX + "other": lambda _context: (
_Provider()
)
},
)
)
registry.configure_capability_context(
ModuleContext(registry=registry, settings=object())
)
try:
report_providers(registry)
except ValueError as exc:
assert "does not match" in str(exc)
else:
raise AssertionError("mismatched provider id was accepted")
+17
View File
@@ -123,3 +123,20 @@ def test_worker_disables_consumers_without_reclaiming_stale_identity(
celery_app._worker_consumer,
celery_app._worker_draining,
) = original_state
def test_worker_child_replaces_inherited_database_pool(monkeypatch) -> None:
from govoplan_core import celery_app
calls: list[tuple[str, bool]] = []
monkeypatch.setattr(
celery_app,
"configure_database",
lambda url, *, dispose_previous=False: calls.append(
(url, dispose_previous)
),
)
celery_app._reset_worker_process_database()
assert calls == [(celery_app.settings.database_url, True)]
+7
View File
@@ -18,6 +18,7 @@ from govoplan_core.core.runtime_coordination import (
list_runtime_nodes,
register_runtime_node,
request_runtime_node_drain,
runtime_composition_hash,
)
from govoplan_core.db.base import Base
@@ -128,3 +129,9 @@ def test_expired_lease_reassignment_increments_fence() -> None:
finally:
session.close()
engine.dispose()
def test_runtime_composition_hash_is_order_and_duplicate_independent() -> None:
assert runtime_composition_hash(("files", "core", "files")) == (
runtime_composition_hash(("core", "files"))
)
+13
View File
@@ -127,6 +127,19 @@ export function apiPostJson<TResponse, TPayload = unknown>(
return apiPost<TResponse>(settings, path, { ...init, body: JSON.stringify(payload) });
}
export function apiPatch<TResponse>(settings: ApiSettings, path: string, init: Omit<RequestInit, "method"> = {}): Promise<TResponse> {
return apiFetch<TResponse>(settings, path, { ...init, method: "PATCH" });
}
export function apiPatchJson<TResponse, TPayload = unknown>(
settings: ApiSettings,
path: string,
payload: TPayload,
init: Omit<RequestInit, "method" | "body"> = {}
): Promise<TResponse> {
return apiPatch<TResponse>(settings, path, { ...init, body: JSON.stringify(payload) });
}
export function loadApiSettings(): ApiSettings {
const storedBaseUrl = loadStoredSetting("baseUrl");
const storedApiKey = sessionStorage.getItem(`${SESSION_STORAGE_KEY}.apiKey`);