Add governed institutional capability contracts
This commit is contained in:
@@ -53,7 +53,9 @@ from govoplan_core.core.postbox import (
|
||||
)
|
||||
from govoplan_core.core.workflows import (
|
||||
CAPABILITY_WORKFLOW_RUNTIME_WORKER,
|
||||
CAPABILITY_WORKFLOW_TRIGGER_DISPATCHER,
|
||||
WorkflowRuntimeWorker,
|
||||
WorkflowTriggerDispatcher,
|
||||
)
|
||||
from govoplan_core.core.registry import PlatformRegistry
|
||||
from govoplan_core.core.runtime import configure_runtime
|
||||
@@ -385,6 +387,18 @@ def _workflow_runtime_worker(
|
||||
return capability
|
||||
|
||||
|
||||
def _workflow_trigger_dispatcher(
|
||||
registry: PlatformRegistry | None = None,
|
||||
) -> WorkflowTriggerDispatcher | None:
|
||||
registry = registry or _platform_registry()
|
||||
if not registry.has_capability(CAPABILITY_WORKFLOW_TRIGGER_DISPATCHER):
|
||||
return None
|
||||
capability = registry.require_capability(CAPABILITY_WORKFLOW_TRIGGER_DISPATCHER)
|
||||
if not isinstance(capability, WorkflowTriggerDispatcher):
|
||||
raise RuntimeError("Workflow trigger dispatcher capability is invalid")
|
||||
return capability
|
||||
|
||||
|
||||
def _postbox_routing_provider(
|
||||
registry: PlatformRegistry | None = None,
|
||||
) -> PostboxRoutingProvider | None:
|
||||
@@ -764,7 +778,8 @@ def dispatch_platform_events(self, limit: int = 100):
|
||||
"observer_failed": 0,
|
||||
}
|
||||
dataflow_dispatcher = _dataflow_trigger_dispatcher(registry)
|
||||
consumers = ()
|
||||
workflow_dispatcher = _workflow_trigger_dispatcher(registry)
|
||||
consumers: list[DurableEventConsumer] = []
|
||||
if dataflow_dispatcher is not None:
|
||||
|
||||
def deliver_to_dataflow(
|
||||
@@ -776,19 +791,38 @@ def dispatch_platform_events(self, limit: int = 100):
|
||||
event=event,
|
||||
)
|
||||
|
||||
consumers = (
|
||||
consumers.append(
|
||||
DurableEventConsumer(
|
||||
consumer_id="dataflow.event-triggers.v1",
|
||||
event_types=frozenset({"*"}),
|
||||
classifications=frozenset({"public", "internal"}),
|
||||
handler=deliver_to_dataflow,
|
||||
),
|
||||
)
|
||||
)
|
||||
if workflow_dispatcher is not None:
|
||||
|
||||
def deliver_to_workflow(
|
||||
event: PlatformEvent,
|
||||
_delivery_key: str,
|
||||
) -> None:
|
||||
workflow_dispatcher.ingest_event(
|
||||
session,
|
||||
event=event,
|
||||
)
|
||||
|
||||
consumers.append(
|
||||
DurableEventConsumer(
|
||||
consumer_id="workflow.event-triggers.v1",
|
||||
event_types=frozenset({"*"}),
|
||||
classifications=frozenset({"public", "internal"}),
|
||||
handler=deliver_to_workflow,
|
||||
)
|
||||
)
|
||||
|
||||
result = dict(
|
||||
outbox.dispatch_pending(
|
||||
session,
|
||||
consumers=consumers,
|
||||
consumers=tuple(consumers),
|
||||
observer=publish_platform_event,
|
||||
limit=limit,
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user