Capability sources, add default source under wf.

This commit is contained in:
lda
2026-05-10 16:50:53 +07:00 Verified
parent f13bad5754
commit 4999cf38bb
17 changed files with 1739 additions and 129 deletions
+65
View File
@@ -0,0 +1,65 @@
from __future__ import annotations
from dataclasses import dataclass
from .service.capability_sources import (
CapabilityBuckets,
CapabilitySource,
SourcePermissions,
SourceVisibility,
)
ADMIN_SOURCE_ID = "wf.admin"
@dataclass(frozen=True, slots=True)
class AdminTool:
name: str
description: str
handler_name: str
mutates_config: bool = False
mutates_auth: bool = False
def admin_source() -> CapabilitySource:
"""Return metadata for privileged broker administration tools."""
tools = {
tool.name: tool
for tool in (
AdminTool(
name="wf.admin.list_sources",
description="List broker capability sources.",
handler_name="list_sources",
),
AdminTool(
name="wf.admin.disable_source",
description="Disable a broker capability source.",
handler_name="disable_source",
mutates_config=True,
),
AdminTool(
name="wf.admin.enable_source",
description="Enable a broker capability source.",
handler_name="enable_source",
mutates_config=True,
),
)
}
return CapabilitySource(
id=ADMIN_SOURCE_ID,
kind="system",
capabilities=CapabilityBuckets(tools=tools),
visibility=SourceVisibility(
planner=False,
mcp_client=False,
admin_dashboard=True,
),
permissions=SourcePermissions(
safe_for_workflow=False,
calls_upstream=False,
mutates_config=True,
mutates_auth=True,
),
description="Privileged broker administration capabilities.",
)
+41 -6
View File
@@ -4,7 +4,24 @@ from typing import Any, Protocol
from pydantic import BaseModel, Field
from wf_authoring import NodeReturn, NodeSpec, node, runtime_error
from wf_authoring import (
NodeReturn,
NodeSpec,
coalesce,
constant,
default_if_none,
first_item,
first_item_maybe,
first_item_or_none,
is_empty,
last_item,
last_item_or_none,
length,
node,
pick_key,
runtime_error,
truthy,
)
from .sources import SpecSource
from .specs import qualify_spec
@@ -16,6 +33,24 @@ MCP_SOURCE_ID = "wf.mcp"
"""Internal source id for broker MCP utility node specs."""
AUTHORING_STD_SPECS: tuple[NodeSpec[Any, Any], ...] = (
coalesce,
default_if_none,
constant,
pick_key,
truthy,
runtime_error,
first_item,
first_item_or_none,
first_item_maybe,
last_item,
last_item_or_none,
length,
is_empty,
)
"""Existing authoring ops that are also exposed through the workflow stdlib."""
class ToolCaller(Protocol):
"""Small service boundary needed by the broker-local MCP utility nodes."""
@@ -56,11 +91,8 @@ class McpCallToolOutput(BaseModel):
def builtin_specs() -> dict[str, NodeSpec[Any, Any]]:
"""Return built-in NodeSpecs available to raw broker workflow plans."""
specs = [
node(
runtime_error,
name="runtime_error",
description="Fail the current workflow branch with a runtime error.",
)
node(spec, name=spec.name.removeprefix("authoring."))
for spec in AUTHORING_STD_SPECS
]
qualified_specs = [qualify_spec(BUILTIN_CONNECTION_ID, spec) for spec in specs]
return {spec.name: spec for spec in qualified_specs}
@@ -96,12 +128,15 @@ def builtin_sources(service: ToolCaller) -> dict[str, SpecSource]:
id=BUILTIN_CONNECTION_ID,
kind="system",
specs=builtin_specs(),
mcp_client_visible=True,
safe_for_workflow=True,
description="Workflow standard-library nodes.",
),
MCP_SOURCE_ID: SpecSource(
id=MCP_SOURCE_ID,
kind="system",
specs=mcp_specs(service),
calls_upstream=True,
description="Broker MCP utility nodes.",
),
}
@@ -0,0 +1,65 @@
from __future__ import annotations
from dataclasses import dataclass, field
from typing import Any, Literal
from wf_authoring import NodeSpec
SourceKind = Literal["system", "connection"]
@dataclass(frozen=True, slots=True)
class SourceVisibility:
planner: bool = False
mcp_client: bool = False
admin_dashboard: bool = True
@dataclass(frozen=True, slots=True)
class SourcePermissions:
safe_for_workflow: bool = False
calls_upstream: bool = False
mutates_config: bool = False
mutates_auth: bool = False
@dataclass(slots=True)
class CapabilityBuckets:
tools: dict[str, Any] = field(default_factory=dict)
node_specs: dict[str, NodeSpec[Any, Any]] = field(default_factory=dict)
prompts: dict[str, Any] = field(default_factory=dict)
resources: dict[str, Any] = field(default_factory=dict)
@dataclass(slots=True)
class CapabilitySource:
id: str
kind: SourceKind
capabilities: CapabilityBuckets = field(default_factory=CapabilityBuckets)
enabled: bool = True
visibility: SourceVisibility = field(default_factory=SourceVisibility)
permissions: SourcePermissions = field(default_factory=SourcePermissions)
description: str | None = None
def as_status(self) -> dict[str, Any]:
return {
"id": self.id,
"kind": self.kind,
"enabled": self.enabled,
"visibility": {
"planner": self.visibility.planner,
"mcp_client": self.visibility.mcp_client,
"admin_dashboard": self.visibility.admin_dashboard,
},
"permissions": {
"safe_for_workflow": self.permissions.safe_for_workflow,
"calls_upstream": self.permissions.calls_upstream,
"mutates_config": self.permissions.mutates_config,
"mutates_auth": self.permissions.mutates_auth,
},
"description": self.description,
"tool_count": len(self.capabilities.tools),
"node_spec_count": len(self.capabilities.node_specs),
"prompt_count": len(self.capabilities.prompts),
"resource_count": len(self.capabilities.resources),
}
+168 -25
View File
@@ -4,7 +4,9 @@ import time
from dataclasses import dataclass, field
from typing import Any
from wf_authoring import NodeSpec, build_async_registry
from pydantic import BaseModel
from wf_authoring import NodeReturn, NodeSpec, build_async_registry
from wf_core import NodeUse, Workflow, execute_workflow_async
from ...connections import ConnectionRegistry, parse_connection_id, qualify_node_name
@@ -19,12 +21,21 @@ from ...models import (
)
from ...sdk import BackendAdapter
from ...shared.errors import error_payload
from ...shared.names import RESERVED_CONNECTION_IDS
from ...storage import Store
from ...workflow.wrappers import _model_from_schema
from ..catalog import CombinedCatalog, snapshot_from_specs
from ..discovery import discover_connection_capabilities, specs_from_discovered_tools
from ..events import McpEvent, make_event
from ..admin_capabilities import admin_source
from .adapters import require_adapter
from .builtins import builtin_sources
from .capability_sources import (
CapabilityBuckets,
CapabilitySource,
SourcePermissions,
SourceVisibility,
)
from .sources import SpecSource
from .specs import get_qualified_spec, qualify_spec
@@ -35,7 +46,7 @@ class WfMcpService:
default_catalog_max_age_seconds: int = 300
connections: ConnectionRegistry = field(default_factory=ConnectionRegistry)
adapters: dict[str, BackendAdapter] = field(default_factory=dict)
spec_sources: dict[str, SpecSource] = field(default_factory=dict)
capability_sources: dict[str, CapabilitySource] = field(default_factory=dict)
events: list[McpEvent] = field(default_factory=list)
include_builtin_specs: bool = True
@@ -44,15 +55,46 @@ class WfMcpService:
if self.include_builtin_specs:
for source in builtin_sources(self).values():
self.register_spec_source(source)
self.register_capability_source(admin_source())
@property
def spec_sources(self) -> dict[str, SpecSource]:
"""Compatibility view of node-spec capability sources."""
return {
source.id: SpecSource(
id=source.id,
kind=source.kind,
specs=dict(source.capabilities.node_specs),
visible=source.enabled and source.visibility.planner,
mcp_client_visible=source.enabled and source.visibility.mcp_client,
admin_dashboard_visible=(
source.enabled and source.visibility.admin_dashboard
),
safe_for_workflow=source.permissions.safe_for_workflow,
calls_upstream=source.permissions.calls_upstream,
mutates_config=source.permissions.mutates_config,
mutates_auth=source.permissions.mutates_auth,
description=source.description,
)
for source in self.capability_sources.values()
if source.capabilities.node_specs
}
@property
def specs_by_connection(self) -> dict[str, dict[str, NodeSpec[Any, Any]]]:
"""Compatibility view of source specs keyed by source id."""
return {source.id: source.specs for source in self.spec_sources.values()}
return {
source.id: dict(source.capabilities.node_specs)
for source in self.capability_sources.values()
if source.capabilities.node_specs
}
def register_connection(self, connection: ConnectionConfig) -> None:
parse_connection_id(connection.id)
if connection.id in RESERVED_CONNECTION_IDS:
raise ValueError(f"connection id {connection.id!r} is reserved by wf-mcp")
self.connections.register(connection)
self._hydrate_connection_source_from_snapshot(connection)
self._record_event(
make_event(
"connection_registered",
@@ -90,13 +132,23 @@ class WfMcpService:
)
for spec in specs
}
self.register_spec_source(
SpecSource(
id=connection_id,
kind="connection",
specs=qualified_specs,
description=f"Specs discovered or registered for {connection_id}.",
)
existing_source = self.capability_sources.get(connection_id)
if existing_source is not None:
# Catalog refreshes replace discovered specs, not operator policy.
existing_source.capabilities.node_specs = qualified_specs
else:
self.register_spec_source(
SpecSource(
id=connection_id,
kind="connection",
specs=qualified_specs,
enabled=self.connections.get(connection_id).enabled,
mcp_client_visible=True,
calls_upstream=True,
description=(
f"Specs discovered or registered for {connection_id}."
),
)
)
snapshot = snapshot_from_specs(
connection_id,
@@ -123,21 +175,43 @@ class WfMcpService:
def get_planner_catalog(self) -> CombinedCatalog:
"""Return all planner-visible specs, including broker-local sources."""
snapshots = dict(self.get_catalog().snapshots)
snapshots: dict[str, CatalogSnapshot] = {}
fetched_at_epoch_ms = int(time.time() * 1000)
for source in self.spec_sources.values():
if not source.visible or source.kind != "system":
for source in self.capability_sources.values():
if not source.enabled or not source.visibility.planner:
continue
stored_snapshot = self.store.load_catalog(source.id)
snapshots[source.id] = snapshot_from_specs(
source.id,
specs=source.specs,
specs=source.capabilities.node_specs,
tool_display_names={
entry.local_name: entry.title
for entry in stored_snapshot.nodes
}
if stored_snapshot is not None
else None,
metadata={
"kind": source.kind,
"description": source.description,
},
fetched_at_epoch_ms=fetched_at_epoch_ms,
max_age_seconds=self.default_catalog_max_age_seconds,
}
if stored_snapshot is None
else stored_snapshot.metadata,
fetched_at_epoch_ms=(
stored_snapshot.fetched_at_epoch_ms
if stored_snapshot is not None
else fetched_at_epoch_ms
),
max_age_seconds=(
stored_snapshot.max_age_seconds
if stored_snapshot is not None
else self.default_catalog_max_age_seconds
),
)
if stored_snapshot is not None:
# Connection resources/prompts are discovered by the backend catalog,
# while planner node visibility is governed by capability sources.
snapshots[source.id].resources = list(stored_snapshot.resources)
snapshots[source.id].prompts = list(stored_snapshot.prompts)
return CombinedCatalog(snapshots=snapshots)
def list_spec_sources(self) -> list[dict[str, Any]]:
@@ -145,9 +219,12 @@ class WfMcpService:
return [
source.as_status()
for source in sorted(
self.spec_sources.values(),
self.capability_sources.values(),
key=lambda source: source.id,
)
if source.capabilities.node_specs
and source.enabled
and source.visibility.planner
]
def list_available_specs(self) -> list[CatalogNodeEntry]:
@@ -407,10 +484,7 @@ class WfMcpService:
)
snapshot = snapshot_from_specs(
connection_id,
specs=self.spec_sources.get(
connection_id,
SpecSource(id=connection_id, kind="connection"),
).specs,
specs=self.specs_by_connection.get(connection_id, {}),
tool_display_names={
tool.name: tool.title for tool in capabilities.tools
},
@@ -495,12 +569,81 @@ class WfMcpService:
def list_events(self) -> list[McpEvent]:
return list(self.events)
def register_capability_source(self, source: CapabilitySource) -> None:
"""Register a capability source as canonical service state."""
self.capability_sources[source.id] = source
def register_spec_source(self, source: SpecSource) -> None:
"""Register a planner source."""
self.spec_sources[source.id] = source
"""Register a legacy spec source through the capability model."""
self.register_capability_source(source.as_capability_source())
def _hydrate_connection_source_from_snapshot(
self,
connection: ConnectionConfig,
) -> None:
"""Restore planner-visible connection specs from a stored catalog snapshot."""
if connection.id in self.capability_sources:
return
snapshot = self.store.load_catalog(connection.id)
if snapshot is None or not snapshot.nodes:
return
specs = {
entry.qualified_name: self._spec_from_snapshot_entry(entry)
for entry in snapshot.nodes
}
self.register_capability_source(
CapabilitySource(
id=connection.id,
kind="connection",
enabled=connection.enabled,
capabilities=CapabilityBuckets(node_specs=specs),
visibility=SourceVisibility(
planner=True,
mcp_client=True,
admin_dashboard=True,
),
permissions=SourcePermissions(calls_upstream=True),
description=f"Specs restored from catalog for {connection.id}.",
)
)
def _spec_from_snapshot_entry(
self,
entry: CatalogNodeEntry,
) -> NodeSpec[Any, Any]:
"""Rebuild an executable tool wrapper from a stored catalog node entry."""
model_prefix = entry.qualified_name.replace(".", "_").replace("-", "_")
input_model = _model_from_schema(f"{model_prefix}_Input", entry.input_schema)
output_model = _model_from_schema(f"{model_prefix}_Output", entry.output_schema)
async def invoke_tool(payload: BaseModel) -> NodeReturn[BaseModel]:
result = await self.call_tool(
entry.connection_id,
entry.local_name,
arguments=payload.model_dump(),
)
return NodeReturn(
outcome=result["outcome"],
output=output_model.model_validate(result["output"]),
)
return NodeSpec(
name=entry.qualified_name,
input_model=input_model,
output_model=output_model,
outcomes=entry.outcomes,
fn=invoke_tool,
description=entry.description,
is_async=True,
accepts_context=False,
input_schema_contract=entry.input_schema,
output_schema_contract=entry.output_schema,
)
def _get_qualified_spec(self, qualified_name: str) -> NodeSpec[Any, Any]:
return get_qualified_spec(self.spec_sources, qualified_name)
return get_qualified_spec(self.capability_sources, qualified_name)
def _record_event(self, event: McpEvent) -> None:
self.events.append(event)
+40 -14
View File
@@ -1,33 +1,59 @@
from __future__ import annotations
from dataclasses import dataclass, field
from typing import Any, Literal
from typing import Any
from wf_authoring import NodeSpec
SpecSourceKind = Literal["connection", "system"]
from .capability_sources import (
CapabilityBuckets,
CapabilitySource,
SourceKind,
SourcePermissions,
SourceVisibility,
)
@dataclass(slots=True)
class SpecSource:
"""Planner-visible collection of workflow node specs.
"""Compatibility wrapper for planner node specs.
Connection sources come from proxied MCP servers. System sources are local broker
capabilities, such as workflow stdlib nodes or service-bound MCP control nodes.
Visibility and permissions stay explicit so callers do not infer source
semantics from the legacy ``kind`` field during the capability-source move.
"""
id: str
kind: SpecSourceKind
kind: SourceKind
specs: dict[str, NodeSpec[Any, Any]] = field(default_factory=dict)
enabled: bool = True
visible: bool = True
mcp_client_visible: bool = False
admin_dashboard_visible: bool = True
safe_for_workflow: bool = False
calls_upstream: bool = False
mutates_config: bool = False
mutates_auth: bool = False
description: str | None = None
def as_capability_source(self) -> CapabilitySource:
return CapabilitySource(
id=self.id,
kind=self.kind,
capabilities=CapabilityBuckets(node_specs=dict(self.specs)),
enabled=self.enabled,
visibility=SourceVisibility(
planner=self.visible,
mcp_client=self.mcp_client_visible,
admin_dashboard=self.admin_dashboard_visible,
),
permissions=SourcePermissions(
safe_for_workflow=self.safe_for_workflow,
calls_upstream=self.calls_upstream,
mutates_config=self.mutates_config,
mutates_auth=self.mutates_auth,
),
description=self.description,
)
def as_status(self) -> dict[str, Any]:
"""Return a compact payload suitable for UI and debugging surfaces."""
return {
"id": self.id,
"kind": self.kind,
"visible": self.visible,
"description": self.description,
"spec_count": len(self.specs),
}
return self.as_capability_source().as_status()
+11 -10
View File
@@ -1,13 +1,12 @@
from __future__ import annotations
from typing import Any
from collections.abc import Mapping
from typing import Any
from wf_authoring import NodeSpec
from ...connections import qualify_node_name
from .sources import SpecSource
from .capability_sources import CapabilitySource
def qualify_spec(connection_id: str, spec: NodeSpec[Any, Any]) -> NodeSpec[Any, Any]:
@@ -27,12 +26,14 @@ def qualify_spec(connection_id: str, spec: NodeSpec[Any, Any]) -> NodeSpec[Any,
def get_qualified_spec(
spec_sources: Mapping[str, SpecSource],
sources: Mapping[str, CapabilitySource],
qualified_name: str,
) -> NodeSpec[Any, Any]:
"""Resolve a namespaced node spec from planner-visible sources."""
source_id, _ = qualified_name.rsplit(".", 1)
source = spec_sources.get(source_id)
if source is None or qualified_name not in source.specs:
raise KeyError(f"unknown qualified node {qualified_name!r}")
return source.specs[qualified_name]
"""Resolve a namespaced node spec from enabled planner-visible sources."""
for source in sources.values():
if not source.enabled or not source.visibility.planner:
continue
spec = source.capabilities.node_specs.get(qualified_name)
if spec is not None:
return spec
raise KeyError(f"unknown qualified node {qualified_name!r}")
+3
View File
@@ -12,6 +12,9 @@ from .service import WfMcpService
def register_broker_tools(server: FastMCP, service: WfMcpService) -> None:
"""Register broker tool handlers on a FastMCP server."""
# These MCP tool names are compatibility exports. Their capability metadata
# belongs to the wf.admin source; future admin-enabled servers can project
# dotted wf.admin.* names from that source.
@server.tool()
async def list_connections() -> list[dict[str, Any]]:
return [