refactor: define workflow api surface protocol
This commit is contained in:
@@ -18,6 +18,14 @@ from .next_actions import NextActionPatchExample, NextActionTool, NextActions
|
|||||||
from .refs import WorkflowSurfaceCapabilityId, parse_workflow_surface_capability_id
|
from .refs import WorkflowSurfaceCapabilityId, parse_workflow_surface_capability_id
|
||||||
from .runs import WorkflowRunApi
|
from .runs import WorkflowRunApi
|
||||||
from .service import WorkflowApi
|
from .service import WorkflowApi
|
||||||
|
from .surface import (
|
||||||
|
WorkflowApiSurface,
|
||||||
|
WorkflowArtifactSurface,
|
||||||
|
WorkflowCapabilitySurface,
|
||||||
|
WorkflowDeploymentSurface,
|
||||||
|
WorkflowDraftSurface,
|
||||||
|
WorkflowRunSurface,
|
||||||
|
)
|
||||||
from .wrapper_hints import (
|
from .wrapper_hints import (
|
||||||
MissingDecision,
|
MissingDecision,
|
||||||
MissingDecisionKind,
|
MissingDecisionKind,
|
||||||
@@ -64,15 +72,21 @@ __all__ = [
|
|||||||
"RuntimeDependencies",
|
"RuntimeDependencies",
|
||||||
"TraceRange",
|
"TraceRange",
|
||||||
"WorkflowApi",
|
"WorkflowApi",
|
||||||
|
"WorkflowApiSurface",
|
||||||
"WorkflowArtifactApi",
|
"WorkflowArtifactApi",
|
||||||
|
"WorkflowArtifactSurface",
|
||||||
"WorkflowCapabilityApi",
|
"WorkflowCapabilityApi",
|
||||||
|
"WorkflowCapabilitySurface",
|
||||||
"WorkflowDeploymentApi",
|
"WorkflowDeploymentApi",
|
||||||
|
"WorkflowDeploymentSurface",
|
||||||
"WorkflowDraftApi",
|
"WorkflowDraftApi",
|
||||||
|
"WorkflowDraftSurface",
|
||||||
"WorkflowEventRecorder",
|
"WorkflowEventRecorder",
|
||||||
"WorkflowLiveSourceChecker",
|
"WorkflowLiveSourceChecker",
|
||||||
"WorkflowOperationContext",
|
"WorkflowOperationContext",
|
||||||
"WorkflowRuntimeRunner",
|
"WorkflowRuntimeRunner",
|
||||||
"WorkflowRunApi",
|
"WorkflowRunApi",
|
||||||
|
"WorkflowRunSurface",
|
||||||
"WorkflowSpecProvider",
|
"WorkflowSpecProvider",
|
||||||
"WorkflowSurfaceCapabilityId",
|
"WorkflowSurfaceCapabilityId",
|
||||||
"WrapperAuthoringHints",
|
"WrapperAuthoringHints",
|
||||||
|
|||||||
@@ -0,0 +1,200 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from collections.abc import Sequence
|
||||||
|
from typing import Any, Protocol
|
||||||
|
|
||||||
|
from wf_artifacts import ArtifactKind
|
||||||
|
|
||||||
|
from .runs import TraceRangeLike
|
||||||
|
|
||||||
|
|
||||||
|
class WorkflowCapabilitySurface(Protocol):
|
||||||
|
"""Capability discovery methods exposed by workflow frontends."""
|
||||||
|
|
||||||
|
async def list_capabilities(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
query: str | None = None,
|
||||||
|
source_id: str | None = None,
|
||||||
|
cursor: str | None = None,
|
||||||
|
limit: int = 50,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
async def inspect_capability(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
qualified_name: str,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
|
||||||
|
class WorkflowDraftSurface(Protocol):
|
||||||
|
"""Draft workspace methods exposed by workflow frontends.
|
||||||
|
|
||||||
|
This protocol intentionally describes the transport-facing workflow surface,
|
||||||
|
not every same-process authoring helper on ``WorkflowApi``.
|
||||||
|
"""
|
||||||
|
|
||||||
|
async def list_draft_workspaces(self) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
async def get_draft_workspace(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
workspace_id: str,
|
||||||
|
include_draft: bool = False,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
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[Any] | None = None,
|
||||||
|
output: Sequence[Any] | None = None,
|
||||||
|
input_map: dict[str, str] | None = None,
|
||||||
|
output_map: dict[str, str] | None = None,
|
||||||
|
error_message_source: Any | None = None,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
async def patch_draft_workspace(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
workspace_id: str,
|
||||||
|
revision: int,
|
||||||
|
patch: list[dict[str, Any]],
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
async def validate_draft_workspace(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
workspace_id: str,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
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]: ...
|
||||||
|
|
||||||
|
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]: ...
|
||||||
|
|
||||||
|
|
||||||
|
class WorkflowArtifactSurface(Protocol):
|
||||||
|
"""Artifact catalog methods exposed by workflow frontends."""
|
||||||
|
|
||||||
|
async def list_artifacts(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
query: str | None = None,
|
||||||
|
kind: ArtifactKind | None = None,
|
||||||
|
cursor: str | None = None,
|
||||||
|
limit: int = 50,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
async def inspect_artifact(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
artifact_id: str,
|
||||||
|
version: int,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
|
||||||
|
class WorkflowDeploymentSurface(Protocol):
|
||||||
|
"""Deployment methods exposed by workflow frontends."""
|
||||||
|
|
||||||
|
async def list_deployments(self) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
async def inspect_deployment(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
deployment_id: str,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
async def save_deployment(
|
||||||
|
self,
|
||||||
|
deployment: dict[str, Any],
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
async def delete_deployment(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
deployment_id: str,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
async def validate_deployment(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
deployment_id: str,
|
||||||
|
live_check: bool = False,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
|
||||||
|
class WorkflowRunSurface(Protocol):
|
||||||
|
"""Run lifecycle methods exposed by workflow frontends."""
|
||||||
|
|
||||||
|
async def run_deployment(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
deployment_id: str,
|
||||||
|
workflow_input: dict[str, Any],
|
||||||
|
trace_range: TraceRangeLike | None = None,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
async def inspect_run(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
run_id: str,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
async def read_run_trace(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
run_id: str,
|
||||||
|
trace_range: TraceRangeLike,
|
||||||
|
) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
|
||||||
|
class WorkflowApiSurface(
|
||||||
|
WorkflowCapabilitySurface,
|
||||||
|
WorkflowDraftSurface,
|
||||||
|
WorkflowArtifactSurface,
|
||||||
|
WorkflowDeploymentSurface,
|
||||||
|
WorkflowRunSurface,
|
||||||
|
Protocol,
|
||||||
|
):
|
||||||
|
"""Public workflow operation surface shared by local and remote adapters."""
|
||||||
|
|
||||||
|
|
||||||
|
__all__ = [
|
||||||
|
"WorkflowApiSurface",
|
||||||
|
"WorkflowArtifactSurface",
|
||||||
|
"WorkflowCapabilitySurface",
|
||||||
|
"WorkflowDeploymentSurface",
|
||||||
|
"WorkflowDraftSurface",
|
||||||
|
"WorkflowRunSurface",
|
||||||
|
]
|
||||||
@@ -7,7 +7,7 @@ import json
|
|||||||
import typer
|
import typer
|
||||||
from pydantic import ValidationError
|
from pydantic import ValidationError
|
||||||
|
|
||||||
from wf_api import WorkflowApi
|
from wf_api import WorkflowApi, WorkflowApiSurface
|
||||||
from wf_config import (
|
from wf_config import (
|
||||||
FilesystemStoreConfig,
|
FilesystemStoreConfig,
|
||||||
LocalTargetConfig,
|
LocalTargetConfig,
|
||||||
@@ -27,7 +27,7 @@ class CliContext:
|
|||||||
|
|
||||||
config_path: Path
|
config_path: Path
|
||||||
service: WfMcpService | None
|
service: WfMcpService | None
|
||||||
handlers: "WorkflowApi | RpcWorkflowApiClient"
|
handlers: WorkflowApiSurface
|
||||||
|
|
||||||
|
|
||||||
@dataclass(frozen=True)
|
@dataclass(frozen=True)
|
||||||
|
|||||||
@@ -7,7 +7,7 @@ from uuid import uuid4
|
|||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
|
|
||||||
from wf_api.models import TraceRange
|
from wf_api.runs import TraceRangeLike
|
||||||
|
|
||||||
|
|
||||||
@dataclass(slots=True)
|
@dataclass(slots=True)
|
||||||
@@ -270,7 +270,7 @@ class RpcWorkflowApiClient:
|
|||||||
*,
|
*,
|
||||||
deployment_id: str,
|
deployment_id: str,
|
||||||
workflow_input: dict[str, Any],
|
workflow_input: dict[str, Any],
|
||||||
trace_range: TraceRange | None = None,
|
trace_range: TraceRangeLike | None = None,
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
return await self._call(
|
return await self._call(
|
||||||
"workflow.runs.start",
|
"workflow.runs.start",
|
||||||
@@ -288,7 +288,7 @@ class RpcWorkflowApiClient:
|
|||||||
self,
|
self,
|
||||||
*,
|
*,
|
||||||
run_id: str,
|
run_id: str,
|
||||||
trace_range: TraceRange,
|
trace_range: TraceRangeLike,
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
return await self._call(
|
return await self._call(
|
||||||
"workflow.runs.trace",
|
"workflow.runs.trace",
|
||||||
@@ -299,7 +299,9 @@ class RpcWorkflowApiClient:
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def _trace_range_payload(trace_range: TraceRange | None) -> dict[str, int] | None:
|
def _trace_range_payload(
|
||||||
|
trace_range: TraceRangeLike | None,
|
||||||
|
) -> dict[str, int] | None:
|
||||||
if trace_range is None:
|
if trace_range is None:
|
||||||
return None
|
return None
|
||||||
return {"start": trace_range.start, "limit": trace_range.limit}
|
return {"start": trace_range.start, "limit": trace_range.limit}
|
||||||
|
|||||||
@@ -0,0 +1,24 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from wf_api import WorkflowApi
|
||||||
|
from wf_mcp.broker.service import WfMcpService
|
||||||
|
from wf_mcp.broker.service.workflow_operation_context import context_from_service
|
||||||
|
from wf_mcp.storage import FileStore
|
||||||
|
from wf_transport_rpc_http import RpcWorkflowApiClient
|
||||||
|
|
||||||
|
|
||||||
|
def test_workflow_api_satisfies_surface_protocol(tmp_path) -> None:
|
||||||
|
from wf_api.surface import WorkflowApiSurface
|
||||||
|
|
||||||
|
service = WfMcpService(store=FileStore(tmp_path / "mcp"))
|
||||||
|
api: WorkflowApiSurface = WorkflowApi(context_from_service(service))
|
||||||
|
|
||||||
|
assert api is not None
|
||||||
|
|
||||||
|
|
||||||
|
def test_rpc_workflow_client_satisfies_surface_protocol() -> None:
|
||||||
|
from wf_api.surface import WorkflowApiSurface
|
||||||
|
|
||||||
|
api: WorkflowApiSurface = RpcWorkflowApiClient("http://example.test/rpc")
|
||||||
|
|
||||||
|
assert api is not None
|
||||||
Reference in New Issue
Block a user