290 lines
11 KiB
Python
290 lines
11 KiB
Python
from __future__ import annotations
|
|
|
|
from pathlib import Path
|
|
|
|
import pytest
|
|
|
|
import wf_api.capabilities as capabilities_module
|
|
from tests.wf_mcp.test_support import echo_tool
|
|
from tests.wf_mcp.workflow_surface.conftest import echo_artifact, failing_tool
|
|
from wf_api.capabilities import WorkflowCapabilityApi
|
|
from wf_artifacts import FileDraftWorkspaceStore, FileWorkflowArtifactStore
|
|
from wf_mcp.broker import WfMcpService
|
|
from wf_mcp.broker.service.workflow_operation_context import context_from_service
|
|
from wf_mcp.models import ConnectionConfig
|
|
from wf_mcp.storage import FileStore
|
|
from wf_mcp.workflow_surface import WorkflowSurfaceHandlers
|
|
|
|
|
|
def _capability_api(
|
|
artifact_store: FileWorkflowArtifactStore,
|
|
*,
|
|
register_echo: bool = False,
|
|
register_failing: bool = False,
|
|
) -> tuple[WorkflowCapabilityApi, WfMcpService]:
|
|
mcp_root = artifact_store.root / "caps_mcp" / str(id(artifact_store))
|
|
service = WfMcpService(
|
|
store=FileStore(mcp_root),
|
|
artifact_store=artifact_store,
|
|
draft_workspace_store=FileDraftWorkspaceStore(mcp_root),
|
|
)
|
|
if register_echo:
|
|
service.register_connection(
|
|
ConnectionConfig(id="demo.personal", server="demo", account="personal")
|
|
)
|
|
service.register_specs("demo.personal", echo_tool)
|
|
if register_failing:
|
|
service.register_connection(
|
|
ConnectionConfig(id="demo.personal", server="demo", account="personal")
|
|
)
|
|
service.register_specs("demo.personal", failing_tool)
|
|
context = context_from_service(service)
|
|
return WorkflowCapabilityApi(context), service
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_capabilities_returns_planner_visible_sources(
|
|
tmp_path: Path,
|
|
) -> None:
|
|
artifact_store = FileWorkflowArtifactStore(tmp_path / "cap_api_list")
|
|
api, _service = _capability_api(artifact_store, register_echo=True)
|
|
|
|
result = await api.list_capabilities()
|
|
|
|
assert result["total"] >= 1
|
|
assert any(
|
|
item["name"] == "demo.personal.echo_tool" for item in result["capabilities"]
|
|
)
|
|
first = next(
|
|
item
|
|
for item in result["capabilities"]
|
|
if item["name"] == "demo.personal.echo_tool"
|
|
)
|
|
assert first["kind"] == "node_spec"
|
|
assert "input_fields" in first
|
|
assert "output_fields" in first
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_capabilities_filters_by_source(tmp_path: Path) -> None:
|
|
artifact_store = FileWorkflowArtifactStore(tmp_path / "cap_api_list_filter")
|
|
api, _service = _capability_api(artifact_store, register_echo=True)
|
|
|
|
result = await api.list_capabilities(source_id="wf.std", query="truthy")
|
|
|
|
assert [item["name"] for item in result["capabilities"]] == ["wf.std.truthy"]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_inspect_capability_returns_detail_with_wrapper_hints(
|
|
tmp_path: Path,
|
|
) -> None:
|
|
artifact_store = FileWorkflowArtifactStore(tmp_path / "cap_api_inspect")
|
|
api, _service = _capability_api(artifact_store, register_echo=True)
|
|
|
|
detail = await api.inspect_capability(qualified_name="demo.personal.echo_tool")
|
|
|
|
assert detail["name"] == "demo.personal.echo_tool"
|
|
assert detail["source_id"] == "demo.personal"
|
|
assert detail["kind"] == "node_spec"
|
|
assert "wrapper_hints" in detail
|
|
hints = detail["wrapper_hints"]
|
|
assert hints["capability_name"] == "demo.personal.echo_tool"
|
|
assert hints["input_map"] == {"input.text": "text"}
|
|
assert hints["output_map"] == {"echoed": "state.echoed"}
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_inspect_capability_raises_on_unknown(tmp_path: Path) -> None:
|
|
artifact_store = FileWorkflowArtifactStore(tmp_path / "cap_api_inspect_unknown")
|
|
api, _service = _capability_api(artifact_store, register_echo=True)
|
|
|
|
with pytest.raises(KeyError, match="no.such.capability"):
|
|
await api.inspect_capability(qualified_name="no.such.capability")
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_call_capability_node_spec_success(tmp_path: Path) -> None:
|
|
artifact_store = FileWorkflowArtifactStore(tmp_path / "cap_api_call")
|
|
api, _service = _capability_api(artifact_store, register_echo=True)
|
|
|
|
result = await api.call_capability(
|
|
qualified_name="demo.personal.echo_tool",
|
|
payload={"text": "hello"},
|
|
)
|
|
|
|
assert result["kind"] == "node_spec"
|
|
assert result["outcome"] == "ok"
|
|
assert result["output"] == {"echoed": "hello"}
|
|
assert result["diagnostics"] == []
|
|
assert result["deployment_id"] is None
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_call_capability_node_spec_failure(tmp_path: Path) -> None:
|
|
artifact_store = FileWorkflowArtifactStore(tmp_path / "cap_api_call_fail")
|
|
api, _service = _capability_api(artifact_store, register_failing=True)
|
|
|
|
result = await api.call_capability(
|
|
qualified_name="demo.personal.failing_tool",
|
|
payload={"message": "boom"},
|
|
)
|
|
|
|
assert result["kind"] == "node_spec"
|
|
assert result["outcome"] == "runtime_error"
|
|
assert result["output"] is None
|
|
assert result["diagnostics"][0]["code"] == "capability_call_failed"
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_list_capabilities_includes_saved_wrapper(tmp_path: Path) -> None:
|
|
artifact_store = FileWorkflowArtifactStore(tmp_path / "cap_api_wrapper_list")
|
|
artifact_store.save_artifact(
|
|
echo_artifact().model_copy(
|
|
update={
|
|
"id": "echo_wrapper",
|
|
"kind": "wrapper",
|
|
"description": "Reusable echo wrapper.",
|
|
}
|
|
)
|
|
)
|
|
api, _service = _capability_api(artifact_store)
|
|
|
|
result = await api.list_capabilities(source_id="workflow", query="echo")
|
|
|
|
names = [item["name"] for item in result["capabilities"]]
|
|
assert names == ["workflow.echo_wrapper.v1"]
|
|
row = result["capabilities"][0]
|
|
assert row["source_id"] == "workflow"
|
|
assert row["kind"] == "wrapper_artifact"
|
|
assert row["artifact_id"] == "echo_wrapper"
|
|
assert row["version"] == 1
|
|
assert row["title"] == "Echo"
|
|
assert row["description"] == "Reusable echo wrapper."
|
|
assert row["outcomes"] == ["completed"]
|
|
assert row["input_fields"] == ["text"]
|
|
assert row["output_fields"] == ["echoed"]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_inspect_capability_saved_wrapper(tmp_path: Path) -> None:
|
|
artifact_store = FileWorkflowArtifactStore(tmp_path / "cap_api_wrapper_inspect")
|
|
artifact_store.save_artifact(
|
|
echo_artifact().model_copy(update={"id": "echo_wrapper", "kind": "wrapper"})
|
|
)
|
|
api, _service = _capability_api(artifact_store)
|
|
|
|
detail = await api.inspect_capability(qualified_name="workflow.echo_wrapper.v1")
|
|
|
|
assert detail["name"] == "workflow.echo_wrapper.v1"
|
|
assert detail["source_id"] == "workflow"
|
|
assert detail["kind"] == "wrapper_artifact"
|
|
assert detail["artifact_id"] == "echo_wrapper"
|
|
assert detail["outcomes"] == ["completed"]
|
|
assert "input_schema" in detail
|
|
assert detail["input_schema"]["properties"]["text"]["type"] == "string"
|
|
assert detail["output_schema"]["properties"]["echoed"]["type"] == "string"
|
|
hints = detail["wrapper_hints"]
|
|
assert hints["capability_name"] == "workflow.echo_wrapper.v1"
|
|
assert hints["declared_outcomes"] == ["completed"]
|
|
assert hints["suggested_wrapper_outcomes"] == ["completed"]
|
|
assert hints["input_map"] == {"input.text": "text"}
|
|
assert hints["output_map"] == {"echoed": "state.echoed"}
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_call_capability_saved_wrapper(tmp_path: Path) -> None:
|
|
artifact_store = FileWorkflowArtifactStore(tmp_path / "cap_api_wrapper_call")
|
|
artifact_store.save_artifact(
|
|
echo_artifact().model_copy(update={"id": "echo_wrapper", "kind": "wrapper"})
|
|
)
|
|
api, service = _capability_api(artifact_store, register_echo=True)
|
|
|
|
result = await api.call_capability(
|
|
qualified_name="workflow.echo_wrapper.v1",
|
|
payload={"text": "hi"},
|
|
)
|
|
|
|
assert result["kind"] == "wrapper_artifact"
|
|
assert result["outcome"] == "completed"
|
|
assert result["diagnostics"] == []
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_create_draft_workspace_from_capability(tmp_path: Path) -> None:
|
|
artifact_store = FileWorkflowArtifactStore(tmp_path / "cap_api_draft_bootstrap")
|
|
api, _service = _capability_api(artifact_store, register_echo=True)
|
|
|
|
result = await api.create_draft_workspace_from_capability(
|
|
workspace_id="echo_ws",
|
|
capability_name="demo.personal.echo_tool",
|
|
)
|
|
|
|
assert result["workspace_id"] == "echo_ws"
|
|
assert result["revision"] == 1
|
|
assert "wrapper_hints" in result
|
|
assert "next_actions" in result
|
|
assert result["wrapper_hints"]["capability_name"] == "demo.personal.echo_tool"
|
|
|
|
fetched = await api.drafts.get_draft_workspace(
|
|
workspace_id="echo_ws", include_draft=True
|
|
)
|
|
|
|
assert fetched["draft"]["steps"]["call"]["use"] == "demo.personal.echo_tool"
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_create_draft_workspace_validates_the_merged_result(
|
|
tmp_path: Path,
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
artifact_store = FileWorkflowArtifactStore(tmp_path / "cap_api_projection")
|
|
api, _service = _capability_api(artifact_store, register_echo=True)
|
|
projected: list[object] = []
|
|
projector = capabilities_module._PROJECT_CREATE_DRAFT_FROM_CAPABILITY
|
|
|
|
def capture_projection(value: object):
|
|
projected.append(value)
|
|
return projector(value)
|
|
|
|
monkeypatch.setattr(
|
|
capabilities_module,
|
|
"_PROJECT_CREATE_DRAFT_FROM_CAPABILITY",
|
|
capture_projection,
|
|
)
|
|
|
|
result = await api.create_draft_workspace_from_capability(
|
|
workspace_id="echo_projected",
|
|
capability_name="demo.personal.echo_tool",
|
|
)
|
|
|
|
assert projected == [result]
|
|
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_handler_delegates_to_capability_api(tmp_path: Path) -> None:
|
|
"""WorkflowSurfaceHandlers methods produce the same result as direct API."""
|
|
artifact_store = FileWorkflowArtifactStore(tmp_path / "cap_api_delegation")
|
|
mcp_root = artifact_store.root / "delegation_mcp"
|
|
service = WfMcpService(
|
|
store=FileStore(mcp_root),
|
|
artifact_store=artifact_store,
|
|
draft_workspace_store=FileDraftWorkspaceStore(mcp_root),
|
|
)
|
|
service.register_connection(
|
|
ConnectionConfig(id="demo.personal", server="demo", account="personal")
|
|
)
|
|
service.register_specs("demo.personal", echo_tool)
|
|
|
|
h = WorkflowSurfaceHandlers(service)
|
|
context = context_from_service(service)
|
|
api = WorkflowCapabilityApi(context)
|
|
|
|
handler_result = await h.inspect_capability(
|
|
qualified_name="demo.personal.echo_tool"
|
|
)
|
|
api_result = await api.inspect_capability(qualified_name="demo.personal.echo_tool")
|
|
|
|
assert handler_result["name"] == api_result["name"]
|
|
assert handler_result["wrapper_hints"] == api_result["wrapper_hints"]
|