feat: add mcp backed workflow server

This commit is contained in:
lda
2026-06-04 23:17:02 +07:00 Unverified
parent 0d34174a84
commit 3b8a1fe6dc
10 changed files with 862 additions and 15 deletions
+97
View File
@@ -0,0 +1,97 @@
from __future__ import annotations
import ast
import pytest
from wf_mcp.broker.config import build_service_from_config
from wf_mcp.broker.server import (
build_workflow_server_from_config,
workflow_server_from_service,
)
from wf_mcp.broker.service.core import WfMcpService
from wf_mcp.models import BrokerConfig, ConnectionConfig
from wf_mcp.source_registry import (
FileSourceRegistryStore,
McpSourceRegistryEntry,
SourceRegistryFile,
)
from wf_mcp.storage import FileStore
from wf_server import WorkflowServer
def _registry_entry(source_id: str) -> McpSourceRegistryEntry:
return McpSourceRegistryEntry.model_validate(
{
"id": source_id,
"kind": "mcp",
"enabled": True,
"provider": "demo",
"account": "registry",
"transport": {"kind": "stdio", "command": "demo-server"},
}
)
def test_wf_server_package_stays_mcp_free() -> None:
path = "src/wf_server/context.py"
tree = ast.parse(open(path, encoding="utf-8").read(), filename=path)
violations: list[str] = []
for node in ast.walk(tree):
if isinstance(node, ast.ImportFrom) and node.module:
if node.module.startswith("wf_mcp"):
violations.append(f"{node.lineno}: from {node.module} import ...")
elif isinstance(node, ast.Import):
for alias in node.names:
if alias.name.startswith("wf_mcp"):
violations.append(f"{node.lineno}: import {alias.name}")
assert violations == []
def test_workflow_server_from_service_wires_neutral_surfaces(tmp_path) -> None:
config = BrokerConfig(
store_root=tmp_path / "store",
connections=[
ConnectionConfig(id="demo.default", server="demo", account="default")
],
)
service = build_service_from_config(config)
server = workflow_server_from_service(
service,
config=config,
source_registry_store=FileSourceRegistryStore(config.store_root),
)
assert isinstance(server, WorkflowServer)
assert server.config.store_root == config.store_root
assert server.api.context is server.context
assert server.source_registry_admin is not None
assert server.admin.connections is service.connection_service
assert server.admin.events is service.events
def test_build_workflow_server_from_config_exposes_registry_admin(tmp_path) -> None:
config = BrokerConfig(store_root=tmp_path / "store", connections=[])
FileSourceRegistryStore(config.store_root).save_registry(
SourceRegistryFile(sources=[_registry_entry("demo.registry")])
)
server = build_workflow_server_from_config(config)
assert server.source_registry_admin is not None
assert "demo.registry" in server.context.specs.capability_sources
def test_workflow_server_from_service_rejects_missing_stores(tmp_path) -> None:
config = BrokerConfig(store_root=tmp_path / "store", connections=[])
service = WfMcpService(store=FileStore(config.store_root))
with pytest.raises(ValueError, match="requires workflow stores"):
workflow_server_from_service(
service,
config=config,
source_registry_store=FileSourceRegistryStore(config.store_root),
)
@@ -0,0 +1,90 @@
from __future__ import annotations
import httpx
from wf_mcp.broker.server import build_workflow_server_from_config
from wf_mcp.models import BrokerConfig, ConnectionConfig
from wf_mcp.source_registry import (
FileSourceRegistryStore,
McpSourceRegistryEntry,
SourceRegistryFile,
)
from wf_transport_rpc_http import RpcWorkflowApiClient, create_rpc_app
def _registry_entry(source_id: str, *, enabled: bool = True) -> McpSourceRegistryEntry:
return McpSourceRegistryEntry.model_validate(
{
"id": source_id,
"kind": "mcp",
"enabled": enabled,
"provider": "demo",
"account": "registry",
"transport": {"kind": "stdio", "command": "demo-server"},
}
)
async def _rpc(client: httpx.AsyncClient, method: str, params: dict) -> dict:
response = await client.post(
"/rpc",
json={"jsonrpc": "2.0", "id": "test", "method": method, "params": params},
)
assert response.status_code == 200
return response.json()
async def test_mcp_backed_rpc_lists_and_mutates_source_registry(tmp_path) -> None:
config = BrokerConfig(store_root=tmp_path / "store", connections=[])
FileSourceRegistryStore(config.store_root).save_registry(
SourceRegistryFile(sources=[_registry_entry("demo.registry")])
)
server = build_workflow_server_from_config(config)
app = create_rpc_app(server)
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(
transport=transport, base_url="http://test"
) as http_client:
client = RpcWorkflowApiClient(
url="http://test/rpc",
http_client=http_client,
)
listed = await client.list_registry_entries(limit=10)
disabled = await client.disable_registry_entry(source_id="demo.registry")
inspected = await client.inspect_registry_entry(source_id="demo.registry")
assert listed["entries"][0]["id"] == "demo.registry"
assert disabled["entry"]["enabled"] is False
assert inspected["entry"]["enabled"] is False
async def test_mcp_backed_rpc_reports_connections_and_events(tmp_path) -> None:
config = BrokerConfig(
store_root=tmp_path / "store",
connections=[
ConnectionConfig(
id="demo.default",
server="demo",
account="default",
)
],
)
server = build_workflow_server_from_config(config)
app = create_rpc_app(server)
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(
transport=transport, base_url="http://test"
) as http_client:
connections = await _rpc(
http_client, "workflow.admin.connections.list", {"limit": 20}
)
events = await _rpc(http_client, "workflow.admin.events.list", {"limit": 20})
assert connections["result"]["connections"][0]["id"] == "demo.default"
assert any(
event["kind"] == "connection_registered"
for event in events["result"]["events"]
)