Files
lda-wf/tests/wf_api/test_deployment_api.py
T
2026-06-15 22:37:59 +07:00

269 lines
8.7 KiB
Python

"""Tests for wf_api.deployments module."""
from __future__ import annotations
from dataclasses import replace
from pathlib import Path
from typing import Any, cast
import pytest
from tests.wf_mcp.test_support import echo_tool
from wf_api.deployments import WorkflowDeploymentApi
from wf_artifacts import (
FileWorkflowArtifactStore,
RequiredCapability,
WorkflowArtifact,
WorkflowDeployment,
)
from wf_mcp.broker import WfMcpService
from wf_mcp.broker.service.workflow_operation_context import context_from_service
from wf_mcp.capabilities import DiscoveredTool
from wf_mcp.models import AuthRecord, ConnectionConfig
from wf_mcp.sdk import BackendAdapter
from wf_mcp.storage import FileStore
from wf_mcp.workflow_surface import WorkflowSurfaceHandlers
def _echo_artifact() -> WorkflowArtifact:
plan: dict[str, Any] = {
"name": "echo",
"input_schema": {
"type": "object",
"properties": {"text": {"type": "string"}},
"required": ["text"],
},
"state_schema": {"fields": {"echoed": {"type": "string"}}},
"output_schema": {
"type": "object",
"properties": {"echoed": {"type": "string"}},
"required": ["echoed"],
},
"start": "echo",
"nodes": [
{
"id": "echo",
"type": "node",
"node": "demo.personal.echo_tool",
"input": [
{
"path": {"root": "input", "parts": ["text"]},
"target": {"root": "local", "parts": ["text"]},
}
],
"output": [
{
"source": {"root": "local", "parts": ["echoed"]},
"target": {"root": "state", "parts": ["echoed"]},
}
],
}
],
"edges": [{"from": "echo", "outcome": "ok", "to": "__end__"}],
}
return WorkflowArtifact(
id="echo",
version=1,
title="Echo",
input_schema=plan["input_schema"],
output_schema=plan["output_schema"],
outcomes=("completed",),
plan=plan,
required_capabilities={
"demo.echo_tool": RequiredCapability(
ref="demo.echo_tool",
kind="node_spec",
)
},
)
def _deployment_api(
artifact_store: FileWorkflowArtifactStore,
*,
register_echo: bool = False,
) -> tuple[WorkflowDeploymentApi, WfMcpService]:
service = WfMcpService(
store=FileStore(
artifact_store.root / "deployments_mcp" / str(id(artifact_store))
),
artifact_store=artifact_store,
)
if register_echo:
service.register_connection(
ConnectionConfig(id="demo.personal", server="demo", account="personal")
)
service.register_specs("demo.personal", echo_tool)
context = context_from_service(service)
return WorkflowDeploymentApi(context), service
@pytest.mark.asyncio
async def test_save_deployment_stores_and_returns_stable_fields(tmp_path: Path) -> None:
artifact_store = FileWorkflowArtifactStore(tmp_path / "deploy_save")
api, _service = _deployment_api(artifact_store)
result = await api.save_deployment(
WorkflowDeployment(
id="echo.personal",
artifact_id="echo",
artifact_version=1,
bindings=[{"logical_source": "demo", "concrete_source": "demo.personal"}],
).model_dump(mode="json")
)
assert result["saved"] is True
assert result["deployment_id"] == "echo.personal"
assert result["artifact_id"] == "echo"
assert result["artifact_version"] == 1
@pytest.mark.asyncio
async def test_list_deployments_returns_compact_summaries(tmp_path: Path) -> None:
artifact_store = FileWorkflowArtifactStore(tmp_path / "deploy_list")
api, _service = _deployment_api(artifact_store)
artifact_store.save_deployment(
WorkflowDeployment(
id="echo.personal",
artifact_id="echo",
artifact_version=1,
bindings=[{"logical_source": "demo", "concrete_source": "demo.personal"}],
)
)
result = await api.list_deployments()
assert len(result["deployments"]) == 1
assert result["deployments"][0]["id"] == "echo.personal"
assert result["deployments"][0]["binding_count"] == 1
assert "bindings" not in result["deployments"][0]
@pytest.mark.asyncio
async def test_list_deployments_returns_empty_without_artifact_store(
tmp_path: Path,
) -> None:
artifact_store = FileWorkflowArtifactStore(tmp_path / "deploy_no_store")
_api, service = _deployment_api(artifact_store)
context = replace(context_from_service(service), artifact_store=None)
api = WorkflowDeploymentApi(context)
result = await api.list_deployments()
assert result["deployments"] == []
@pytest.mark.asyncio
async def test_delete_deployment_removes_one(tmp_path: Path) -> None:
artifact_store = FileWorkflowArtifactStore(tmp_path / "deploy_delete")
api, _service = _deployment_api(artifact_store)
artifact_store.save_deployment(
WorkflowDeployment(
id="echo.personal",
artifact_id="echo",
artifact_version=1,
bindings=[{"logical_source": "demo", "concrete_source": "demo.personal"}],
)
)
result = await api.delete_deployment(deployment_id="echo.personal")
assert result["deployment_id"] == "echo.personal"
assert result["deleted"] is True
assert artifact_store.list_deployments() == []
@pytest.mark.asyncio
async def test_validate_deployment_returns_runnable_for_valid_binding(
tmp_path: Path,
) -> None:
artifact_store = FileWorkflowArtifactStore(tmp_path / "deploy_validate_runnable")
api, service = _deployment_api(artifact_store, register_echo=True)
artifact_store.save_artifact(_echo_artifact())
artifact_store.save_deployment(
WorkflowDeployment(
id="echo.personal",
artifact_id="echo",
artifact_version=1,
bindings=[{"logical_source": "demo", "concrete_source": "demo.personal"}],
)
)
result = await api.validate_deployment(
deployment_id="echo.personal", live_check=False
)
assert result["status"] == "runnable"
assert result["diagnostics"] == []
class FailingLivenessAdapter:
async def list_tools(
self,
connection: ConnectionConfig,
auth: AuthRecord | None,
) -> list[DiscoveredTool]:
raise OSError("stdio process exited")
@pytest.mark.asyncio
async def test_validate_deployment_live_check_calls_live_checker(
tmp_path: Path,
) -> None:
artifact_store = FileWorkflowArtifactStore(tmp_path / "deploy_validate_live")
api, service = _deployment_api(artifact_store, register_echo=True)
artifact_store.save_artifact(_echo_artifact())
artifact_store.save_deployment(
WorkflowDeployment(
id="echo.personal",
artifact_id="echo",
artifact_version=1,
bindings=[{"logical_source": "demo", "concrete_source": "demo.personal"}],
)
)
service.register_adapter(
"demo",
cast(BackendAdapter, FailingLivenessAdapter()),
)
result = await api.validate_deployment(
deployment_id="echo.personal", live_check=True
)
assert result["status"] == "unrunnable"
assert result["diagnostics"][0]["code"] == "source_unreachable"
@pytest.mark.asyncio
async def test_handler_delegation_for_validate_deployment(tmp_path: Path) -> None:
"""WorkflowSurfaceHandlers.validate_deployment delegates to WorkflowDeploymentApi."""
artifact_store = FileWorkflowArtifactStore(tmp_path / "deploy_delegation")
service = WfMcpService(
store=FileStore(artifact_store.root / "delegation_mcp"),
artifact_store=artifact_store,
)
artifact_store.save_artifact(_echo_artifact())
artifact_store.save_deployment(
WorkflowDeployment(
id="echo.personal",
artifact_id="echo",
artifact_version=1,
bindings=[{"logical_source": "demo", "concrete_source": "demo.personal"}],
)
)
h = WorkflowSurfaceHandlers(service)
context = context_from_service(service)
api = WorkflowDeploymentApi(context)
handler_result = await h.validate_deployment(deployment_id="echo.personal")
api_result = await api.validate_deployment(deployment_id="echo.personal")
assert handler_result["status"] == api_result["status"]
assert len(handler_result["diagnostics"]) == len(api_result["diagnostics"])
if handler_result["diagnostics"]:
assert (
handler_result["diagnostics"][0]["code"]
== api_result["diagnostics"][0]["code"]
)