| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113 |
- """Database adapters for the fixed WP13 control and runtime gateways."""
- from __future__ import annotations
- import json
- import os
- from typing import Any
- from sqlalchemy import create_engine, text
- from sqlalchemy.pool import NullPool
- from app.config.database_urls import validate_postgresql_url
- from app.core.mcp.governed_invocation import normalize_mcp_invocation
- from .governance import fixture_signature_digest, normalize_plugin_manifest
- from .runtime import BuiltinFixtureDispatcher
- class PluginPersistenceError(ValueError):
- """The persistent gateway rejected a closed plugin request."""
- def _digest(value: Any, label: str) -> str:
- if not isinstance(value, str) or len(value) != 64 or set(value) - set("0123456789abcdef"):
- raise PluginPersistenceError(f"{label}_invalid")
- return value
- class DatabasePluginPlatformService:
- """Uses only fixed SECURITY DEFINER functions; never direct fact SQL."""
- def __init__(self) -> None:
- self._runtime_url = os.environ.get("DATABASE_URL", "")
- self._control_url = os.environ.get("PLUGIN_PLATFORM_CONTROL_DATABASE_URL", os.environ.get("BI_AI_CATALOG_CONTROL_DATABASE_URL", ""))
- runtime = validate_postgresql_url(self._runtime_url, "DATABASE_URL")
- control = validate_postgresql_url(self._control_url, "PLUGIN_PLATFORM_CONTROL_DATABASE_URL")
- if runtime.username == control.username:
- raise RuntimeError("plugin control and runtime identities must differ")
- @staticmethod
- def _scope() -> dict[str, str]:
- tenant = os.environ.get("TRUSTED_PLUGIN_TENANT", "")
- domain = os.environ.get("TRUSTED_PLUGIN_DOMAIN", "")
- if not tenant or not domain or not tenant.isascii() or not domain.isascii():
- raise PluginPersistenceError("trusted_plugin_scope_missing")
- return {"tenant_ref": tenant, "domain_ref": domain}
- def _control(self, payload: dict[str, Any]) -> dict[str, Any]:
- with create_engine(self._control_url, poolclass=NullPool, pool_pre_ping=True).begin() as connection:
- value = connection.execute(text("SELECT public.plugin_platform_control_v5(CAST(:payload AS jsonb))"), {"payload": json.dumps(payload, sort_keys=True, separators=(",", ":"))}).scalar_one()
- return dict(value)
- def _claim(self, payload: dict[str, Any]) -> dict[str, Any]:
- with create_engine(self._control_url, poolclass=NullPool, pool_pre_ping=True).begin() as connection:
- value = connection.execute(text("SELECT public.plugin_platform_issue_claim_v3(CAST(:payload AS jsonb))"), {"payload": json.dumps(payload, sort_keys=True, separators=(",", ":"))}).scalar_one()
- return dict(value)
- def _runtime(self, payload: dict[str, Any]) -> dict[str, Any]:
- with create_engine(self._runtime_url, poolclass=NullPool, pool_pre_ping=True).begin() as connection:
- value = connection.execute(text("SELECT public.plugin_platform_runtime_v3(CAST(:payload AS jsonb))"), {"payload": json.dumps(payload, sort_keys=True, separators=(",", ":"))}).scalar_one()
- return dict(value)
- def register(self, *, manifest: Any, registry_record: Any, actor_ref: str) -> dict[str, Any]:
- normalized = normalize_plugin_manifest(manifest)
- if not isinstance(registry_record, dict) or set(registry_record) != {"artifact_digest", "signature", "sbom_digest", "license_digest", "vulnerability_digest", "provenance_digest"}:
- raise PluginPersistenceError("registry_record_closed")
- signature = registry_record["signature"]
- if not isinstance(signature, dict) or set(signature) != {"trust_store_key_id", "signature_digest"} or signature["trust_store_key_id"] != "local-fixture-key-v1":
- raise PluginPersistenceError("trust_store_denied")
- payload = {
- "action": "register", "plugin_uid": normalized["plugin_uid"], "version": normalized["version"], "actor_ref": actor_ref,
- "plugin_type": normalized["type"], "manifest_digest": normalized["manifest_digest"],
- "artifact_digest": _digest(registry_record["artifact_digest"], "artifact_digest"),
- "signature_digest": _digest(signature["signature_digest"], "signature_digest"),
- "sbom_digest": _digest(registry_record["sbom_digest"], "sbom_digest"), "license_digest": _digest(registry_record["license_digest"], "license_digest"),
- "vulnerability_digest": _digest(registry_record["vulnerability_digest"], "vulnerability_digest"), "provenance_digest": _digest(registry_record["provenance_digest"], "provenance_digest"),
- "capabilities": normalized["capabilities"], "resource": normalized["resource"],
- }
- if payload["artifact_digest"] != normalized["distribution"]["artifact_digest"]:
- raise PluginPersistenceError("artifact_digest_mismatch")
- if payload["signature_digest"] != fixture_signature_digest(payload["artifact_digest"]):
- raise PluginPersistenceError("signature_binding_denied")
- return self._control(payload)
- def review(self, *, plugin_uid: str, version: str, actor_ref: str) -> dict[str, Any]:
- return self._control({"action": "review", "plugin_uid": plugin_uid, "version": version, "actor_ref": actor_ref})
- def issue_approval(self, *, plugin_uid: str, version: str, actor_ref: str, approval_action: str) -> dict[str, Any]:
- return self._control({"action": "issue_approval", "plugin_uid": plugin_uid, "version": version, "actor_ref": actor_ref, "approval_action": approval_action, **self._scope(), "expires_in_seconds": 300})
- def transition(self, *, plugin_uid: str, version: str, actor_ref: str, target_state: str, approval_uid: str, expected_fence: int, incident_uid: str = "") -> dict[str, Any]:
- if not isinstance(expected_fence, int) or expected_fence < 0:
- raise PluginPersistenceError("expected_fence_invalid")
- return self._control({"action": "transition", "plugin_uid": plugin_uid, "version": version, "actor_ref": actor_ref, **self._scope(), "approval_uid": approval_uid, "target_state": target_state, "expected_fence": expected_fence, "incident_uid": incident_uid})
- def invoke(self, *, plugin_uid: str, version: str, actor_ref: str, operation: str, input_digest: str, idempotency_key: str) -> dict[str, Any]:
- _digest(input_digest, "input_digest")
- if not isinstance(idempotency_key, str) or not idempotency_key.isascii() or not 1 <= len(idempotency_key) <= 120:
- raise PluginPersistenceError("idempotency_key_invalid")
- scope = self._scope()
- description = self._control({"action": "describe", "plugin_uid": plugin_uid, "version": version, "actor_ref": actor_ref})
- if description["plugin_type"] == "agent_mcp":
- normalize_mcp_invocation({"interface_type": "mcp", "tool_name": "governed_plugin_fixture", "action": operation, "arguments_digest": input_digest, "evidence_refs": []})
- claim = self._claim({"plugin_uid": plugin_uid, "version": version, "principal_ref": actor_ref, "operation_name": operation, "input_digest": input_digest, "idempotency_key": idempotency_key, **scope})
- queued = self._runtime({"action": "enqueue", "claim_uid": claim["claim_uid"], "idempotency_key": idempotency_key, "input_digest": input_digest, "operation_name": operation})
- lease = self._runtime({"action": "claim", "run_uid": queued["run_uid"], "worker": "fixed-builtin"})
- fixture_manifest = {"distribution": {"kind": "builtin_fixture", "fixture_id": "ENGINEERING_EVIDENCE_ONLY"}, "permissions": {"child_process": False, "file": False, "network": False, "secret": False}, "type": description["plugin_type"], "capabilities": description["capabilities"], "resource": description["resource"], "manifest_digest": description["manifest_digest"]}
- result = BuiltinFixtureDispatcher().execute(fixture_manifest, operation=operation, input_digest=input_digest)
- settled = self._runtime({"action": "settle", "run_uid": queued["run_uid"], "worker": "fixed-builtin", "fence": lease["fence"], "success": True, "output_digest": result["output_digest"], "failure_digest": "0" * 64})
- return {"run_uid": queued["run_uid"], **result, **settled}
- __all__ = ["DatabasePluginPlatformService", "PluginPersistenceError"]
|