the WorkflowApiBackend collapses

This commit is contained in:
lda
2026-06-02 11:48:55 +07:00 Verified
parent 09995d86d5
commit 55d45226f4
19 changed files with 984 additions and 1385 deletions
+11 -482
View File
@@ -1,497 +1,26 @@
from __future__ import annotations
from collections.abc import Sequence
from typing import TYPE_CHECKING, Any
from typing import TYPE_CHECKING
from wf_artifacts import (
ArtifactKind,
)
from wf_core.models.steps import (
InputBinding,
OutputBinding,
)
from wf_core.paths import GraphSourcePath
from wf_api.artifacts import WorkflowArtifactApi
from wf_api.capabilities import WorkflowCapabilityApi
from wf_api.deployments import WorkflowDeploymentApi
from wf_api.drafts import WorkflowDraftApi
from wf_api.listing import paged_list_payload
from wf_api.models import RawWorkflowPlan
from wf_api.runs import WorkflowRunApi
from wf_api import WorkflowApi
from ..broker.service.workflow_operation_context import context_from_service
from .models import TraceRange
if TYPE_CHECKING:
from ..broker.service import WfMcpService
class WorkflowSurfaceHandlers:
"""Reusable implementation behind MCP workflow artifact tools."""
class WorkflowSurfaceHandlers(WorkflowApi):
"""Compatibility wrapper for old wf_mcp.workflow_surface imports.
New code should construct `WorkflowApi(context_from_service(service))`
directly. This shim keeps tests and legacy broker artifact tools working
while the MCP surface is migrated.
"""
def __init__(self, service: WfMcpService) -> None:
self.service = service
context = context_from_service(service)
self._capabilities = WorkflowCapabilityApi(context)
self._drafts = WorkflowDraftApi(context)
self._artifacts = WorkflowArtifactApi(context)
self._deployments = WorkflowDeploymentApi(context)
self._runs = WorkflowRunApi(context)
super().__init__(context_from_service(service))
async def list_artifacts(
self,
*,
query: str | None = None,
kind: ArtifactKind | None = None,
cursor: str | None = None,
limit: int = 50,
) -> dict[str, Any]:
"""Return compact paged saved artifact summaries.
Saved artifacts can contain full raw workflow plans, so list results
deliberately stay summary-only. Use inspect/run tools for detail.
"""
if self.service.artifact_store is None:
return paged_list_payload("nodes", [], cursor=cursor, limit=limit)
return await self._artifacts.list_artifacts(
query=query,
kind=kind,
cursor=cursor,
limit=limit,
)
async def list_capabilities(
self,
*,
query: str | None = None,
source_id: str | None = None,
cursor: str | None = None,
limit: int = 50,
) -> dict[str, Any]:
"""Return compact paged planner-visible workflow capability summaries."""
return await self._capabilities.list_capabilities(
query=query,
source_id=source_id,
cursor=cursor,
limit=limit,
)
async def inspect_capability(self, *, qualified_name: str) -> dict[str, Any]:
"""Return one planner-visible workflow capability contract."""
return await self._capabilities.inspect_capability(
qualified_name=qualified_name,
)
async def call_capability(
self,
*,
qualified_name: str,
payload: dict[str, Any],
deployment_id: str | None = None,
) -> dict[str, Any]:
"""Execute one planner-visible workflow capability for authoring tests."""
return await self._capabilities.call_capability(
qualified_name=qualified_name,
payload=payload,
deployment_id=deployment_id,
)
async def save_artifact(self, artifact: dict[str, Any]) -> dict[str, Any]:
if self.service.artifact_store is None:
raise KeyError("workflow artifact store is not configured")
return await self._artifacts.save_artifact(artifact)
async def create_artifact_from_plan(
self,
*,
artifact_id: str,
version: int,
title: str,
plan: RawWorkflowPlan | dict[str, Any],
outcomes: Sequence[str],
kind: ArtifactKind = "workflow",
description: str | None = None,
required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None,
) -> dict[str, Any]:
if self.service.artifact_store is None:
raise KeyError("workflow artifact store is not configured")
return await self._artifacts.create_artifact_from_plan(
artifact_id=artifact_id,
version=version,
title=title,
plan=plan,
outcomes=outcomes,
kind=kind,
description=description,
required_capabilities=required_capabilities,
source_bindings=source_bindings,
created_from_catalog_version=created_from_catalog_version,
)
async def validate_draft(self, *, draft: dict[str, Any]) -> dict[str, Any]:
return await self._drafts.validate_draft(draft=draft)
async def compile_draft(self, *, draft: dict[str, Any]) -> dict[str, Any]:
return await self._drafts.compile_draft(draft=draft)
async def create_artifact_from_draft(
self,
*,
artifact_id: str,
version: int,
title: str,
draft: dict[str, Any],
outcomes: Sequence[str],
kind: ArtifactKind = "workflow",
description: str | None = None,
required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None,
) -> dict[str, Any]:
if self.service.artifact_store is None:
raise KeyError("workflow artifact store is not configured")
return await self._artifacts.create_artifact_from_draft(
artifact_id=artifact_id,
version=version,
title=title,
draft=draft,
outcomes=outcomes,
kind=kind,
description=description,
required_capabilities=required_capabilities,
source_bindings=source_bindings,
created_from_catalog_version=created_from_catalog_version,
)
async def patch_draft(
self,
*,
draft: dict[str, Any],
patch: list[dict[str, Any]],
) -> dict[str, Any]:
return await self._drafts.patch_draft(draft=draft, patch=patch)
async def list_draft_workspaces(self) -> dict[str, Any]:
"""Return compact summaries for stored draft workspaces."""
return await self._drafts.list_draft_workspaces()
async def create_draft_workspace(
self,
*,
workspace_id: str,
draft: dict[str, Any],
title: str | None = None,
) -> dict[str, Any]:
return await self._drafts.create_draft_workspace(
workspace_id=workspace_id,
draft=draft,
title=title,
)
async def get_draft_workspace(
self,
*,
workspace_id: str,
include_draft: bool = False,
) -> dict[str, Any]:
return await self._drafts.get_draft_workspace(
workspace_id=workspace_id,
include_draft=include_draft,
)
async def delete_draft_workspace(self, *, workspace_id: str) -> dict[str, Any]:
return await self._drafts.delete_draft_workspace(workspace_id=workspace_id)
async def validate_draft_workspace(self, *, workspace_id: str) -> dict[str, Any]:
"""Refresh stored validation status without changing draft revision."""
return await self._drafts.validate_draft_workspace(workspace_id=workspace_id)
async def patch_draft_workspace(
self,
*,
workspace_id: str,
revision: int,
patch: list[dict[str, Any]],
) -> dict[str, Any]:
return await self._drafts.patch_draft_workspace(
workspace_id=workspace_id,
revision=revision,
patch=patch,
)
async def set_draft_name(
self,
*,
workspace_id: str,
revision: int,
name: str,
) -> dict[str, Any]:
return await self._drafts.set_draft_name(
workspace_id=workspace_id,
revision=revision,
name=name,
)
async def set_draft_route(
self,
*,
workspace_id: str,
revision: int,
step_id: str,
outcome: str,
target: str,
) -> dict[str, Any]:
return await self._drafts.set_draft_route(
workspace_id=workspace_id,
revision=revision,
step_id=step_id,
outcome=outcome,
target=target,
)
async def set_step_input_map(
self,
*,
workspace_id: str,
revision: int,
step_id: str,
input_map: dict[str, str],
) -> dict[str, Any]:
return await self._drafts.set_step_input_map(
workspace_id=workspace_id,
revision=revision,
step_id=step_id,
input_map=input_map,
)
async def set_step_output_map(
self,
*,
workspace_id: str,
revision: int,
step_id: str,
output_map: dict[str, str],
) -> dict[str, Any]:
return await self._drafts.set_step_output_map(
workspace_id=workspace_id,
revision=revision,
step_id=step_id,
output_map=output_map,
)
async def create_minimal_draft_workspace(
self,
*,
workspace_id: str,
name: str,
capability_name: str,
input_schema: dict[str, Any],
state_schema: dict[str, Any],
output_schema: dict[str, Any],
input: Sequence[InputBinding] | None = None,
output: Sequence[OutputBinding] | None = None,
input_map: dict[str, str] | None = None,
output_map: dict[str, str] | None = None,
error_message_source: str | GraphSourcePath | None = None,
title: str | None = None,
) -> dict[str, Any]:
"""Bootstrap the smallest patchable draft around one workflow capability."""
return await self._drafts.create_minimal_draft_workspace(
workspace_id=workspace_id,
name=name,
capability_name=capability_name,
input_schema=input_schema,
state_schema=state_schema,
output_schema=output_schema,
input=input,
output=output,
input_map=input_map,
output_map=output_map,
error_message_source=error_message_source,
title=title,
)
async def create_draft_workspace_from_capability(
self,
*,
workspace_id: str,
capability_name: str,
name: str | None = None,
title: str | None = None,
input_schema: dict[str, Any] | None = None,
state_schema: dict[str, Any] | None = None,
output_schema: dict[str, Any] | None = None,
input: Sequence[InputBinding] | None = None,
output: Sequence[OutputBinding] | None = None,
input_map: dict[str, str] | None = None,
output_map: dict[str, str] | None = None,
error_message_source: str | GraphSourcePath | None = None,
) -> dict[str, Any]:
"""Create a patchable draft workspace from inspect_capability hints."""
return await self._capabilities.create_draft_workspace_from_capability(
workspace_id=workspace_id,
capability_name=capability_name,
name=name,
title=title,
input_schema=input_schema,
state_schema=state_schema,
output_schema=output_schema,
input=input,
output=output,
input_map=input_map,
output_map=output_map,
error_message_source=error_message_source,
)
async def create_artifact_from_workspace(
self,
*,
workspace_id: str,
artifact_id: str,
version: int,
title: str,
outcomes: Sequence[str],
kind: ArtifactKind = "workflow",
description: str | None = None,
required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None,
) -> dict[str, Any]:
if self.service.artifact_store is None:
raise KeyError("workflow artifact store is not configured")
return await self._artifacts.create_artifact_from_workspace(
workspace_id=workspace_id,
artifact_id=artifact_id,
version=version,
title=title,
outcomes=outcomes,
kind=kind,
description=description,
required_capabilities=required_capabilities,
source_bindings=source_bindings,
created_from_catalog_version=created_from_catalog_version,
)
async def create_wrapper_from_workspace(
self,
*,
workspace_id: str,
artifact_id: str,
version: int,
title: str,
outcomes: Sequence[str],
description: str | None = None,
required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None,
) -> dict[str, Any]:
"""Save the current draft workspace as a callable wrapper artifact."""
if self.service.artifact_store is None:
raise KeyError("workflow artifact store is not configured")
return await self._artifacts.create_wrapper_from_workspace(
workspace_id=workspace_id,
artifact_id=artifact_id,
version=version,
title=title,
outcomes=outcomes,
description=description,
required_capabilities=required_capabilities,
source_bindings=source_bindings,
created_from_catalog_version=created_from_catalog_version,
)
async def inspect_artifact(
self, *, artifact_id: str, version: int
) -> dict[str, Any]:
if self.service.artifact_store is None:
raise KeyError("workflow artifact store is not configured")
return await self._artifacts.inspect_artifact(
artifact_id=artifact_id,
version=version,
)
async def list_deployments(self) -> dict[str, Any]:
if self.service.artifact_store is None:
return {"deployments": []}
return await self._deployments.list_deployments()
async def inspect_deployment(self, *, deployment_id: str) -> dict[str, Any]:
if self.service.artifact_store is None:
raise KeyError("workflow artifact store is not configured")
return await self._deployments.inspect_deployment(
deployment_id=deployment_id,
)
async def save_deployment(self, deployment: dict[str, Any]) -> dict[str, Any]:
if self.service.artifact_store is None:
raise KeyError("workflow artifact store is not configured")
return await self._deployments.save_deployment(deployment)
async def delete_deployment(self, *, deployment_id: str) -> dict[str, Any]:
"""Delete one mutable deployment environment binding."""
if self.service.artifact_store is None:
raise KeyError("workflow artifact store is not configured")
return await self._deployments.delete_deployment(
deployment_id=deployment_id,
)
async def validate_deployment(
self,
*,
deployment_id: str,
live_check: bool = False,
) -> dict[str, Any]:
if self.service.artifact_store is None:
raise KeyError("workflow artifact store is not configured")
return await self._deployments.validate_deployment(
deployment_id=deployment_id,
live_check=live_check,
)
async def run_deployment(
self,
*,
deployment_id: str,
workflow_input: dict[str, Any],
trace_range: TraceRange | None = None,
) -> dict[str, Any]:
return await self._runs.run_deployment(
deployment_id=deployment_id,
workflow_input=workflow_input,
trace_range=trace_range,
)
async def resume_run(
self,
*,
run_id: str,
resume_payload: dict[str, Any],
resume_outcome: str = "submitted",
trace_range: TraceRange | None = None,
) -> dict[str, Any]:
"""Resume one durable interrupted deployment run."""
return await self._runs.resume_run(
run_id=run_id,
resume_payload=resume_payload,
resume_outcome=resume_outcome,
trace_range=trace_range,
)
async def inspect_run(self, *, run_id: str) -> dict[str, Any]:
"""Return one durable stopped-run summary without debug trace entries."""
return await self._runs.inspect_run(run_id=run_id)
async def read_run_trace(
self,
*,
run_id: str,
trace_range: TraceRange,
) -> dict[str, Any]:
"""Return only a caller-bounded debug trace slice from a stopped run."""
return await self._runs.read_run_trace(
run_id=run_id,
trace_range=trace_range,
)
__all__ = ["WorkflowSurfaceHandlers"]
+5 -13
View File
@@ -8,9 +8,8 @@ from pydantic import Field
from wf_artifacts import ArtifactKind
from wf_artifacts.models import RequiredCapability
from wf_api import WorkflowApi
from wf_api.backend import TraceRange as ApiTraceRange
from wf_mcp.broker.service import WfMcpService
from wf_mcp.broker.service.workflow_api_backend import WfMcpWorkflowApiBackend
from wf_mcp.broker.service.workflow_operation_context import context_from_service
from .models import (
CallCapabilityResult,
CreateArtifactFromWorkspaceRequest,
@@ -37,10 +36,7 @@ from .models import (
def register_workflow_tools(server: FastMCP[Any], service: WfMcpService) -> None:
"""Register stable workflow tools on the public MCP server surface."""
handlers = WorkflowApi(WfMcpWorkflowApiBackend(service))
def _to_api_trace_range(tr: TraceRange) -> ApiTraceRange:
return ApiTraceRange(start=tr.start, limit=tr.limit)
handlers = WorkflowApi(context_from_service(service))
@server.tool(
name="wf.workflow.list_artifacts",
@@ -669,9 +665,7 @@ def register_workflow_tools(server: FastMCP[Any], service: WfMcpService) -> None
await handlers.run_deployment(
deployment_id=deployment_id,
workflow_input=workflow_input,
trace_range=_to_api_trace_range(trace_range)
if trace_range is not None
else None,
trace_range=trace_range,
)
)
@@ -703,9 +697,7 @@ def register_workflow_tools(server: FastMCP[Any], service: WfMcpService) -> None
run_id=run_id,
resume_payload=resume_payload,
resume_outcome=resume_outcome,
trace_range=_to_api_trace_range(trace_range)
if trace_range is not None
else None,
trace_range=trace_range,
)
)
@@ -742,6 +734,6 @@ def register_workflow_tools(server: FastMCP[Any], service: WfMcpService) -> None
return RunDeploymentResult.model_validate(
await handlers.read_run_trace(
run_id=run_id,
trace_range=_to_api_trace_range(trace_range),
trace_range=trace_range,
)
)