feat: add remote workflow lifecycle rpc methods

This commit is contained in:
lda
2026-06-03 11:18:32 +07:00 Verified
parent 9155dae8c8
commit 826b28682f
12 changed files with 820 additions and 12 deletions
+22
View File
@@ -4,12 +4,22 @@ from .app import create_rpc_app
from .client import RpcWorkflowApiClient
from .errors import WorkflowRpcError
from .models import (
CreateArtifactFromWorkspaceParams,
CreateDraftFromCapabilityParams,
CreateWrapperFromWorkspaceParams,
DeleteDeploymentParams,
GetDraftWorkspaceParams,
HealthParams,
InspectArtifactParams,
InspectCapabilityParams,
InspectDeploymentParams,
InspectRunParams,
ListArtifactsParams,
ListCapabilitiesParams,
ListDeploymentsParams,
ListDraftWorkspacesParams,
PatchDraftParams,
PatchDraftWorkspaceParams,
ReadRunTraceParams,
ResumeRunParams,
SaveArtifactParams,
@@ -18,15 +28,26 @@ from .models import (
TraceRangeParams,
ValidateDeploymentParams,
ValidateDraftParams,
ValidateDraftWorkspaceParams,
)
__all__ = [
"CreateArtifactFromWorkspaceParams",
"CreateDraftFromCapabilityParams",
"CreateWrapperFromWorkspaceParams",
"DeleteDeploymentParams",
"GetDraftWorkspaceParams",
"HealthParams",
"InspectArtifactParams",
"InspectCapabilityParams",
"InspectDeploymentParams",
"InspectRunParams",
"ListArtifactsParams",
"ListCapabilitiesParams",
"ListDeploymentsParams",
"ListDraftWorkspacesParams",
"PatchDraftParams",
"PatchDraftWorkspaceParams",
"ReadRunTraceParams",
"ResumeRunParams",
"SaveArtifactParams",
@@ -35,6 +56,7 @@ __all__ = [
"TraceRangeParams",
"ValidateDeploymentParams",
"ValidateDraftParams",
"ValidateDraftWorkspaceParams",
"WorkflowRpcError",
"create_rpc_app",
"RpcWorkflowApiClient",
+187
View File
@@ -10,11 +10,21 @@ from wf_server import WorkflowServer
from .errors import WorkflowRpcError, raise_workflow_rpc_error
from .models import (
CreateArtifactFromWorkspaceParams,
CreateDraftFromCapabilityParams,
CreateWrapperFromWorkspaceParams,
DeleteDeploymentParams,
GetDraftWorkspaceParams,
InspectArtifactParams,
InspectCapabilityParams,
InspectDeploymentParams,
InspectRunParams,
ListArtifactsParams,
ListCapabilitiesParams,
ListDeploymentsParams,
ListDraftWorkspacesParams,
PatchDraftParams,
PatchDraftWorkspaceParams,
ReadRunTraceParams,
ResumeRunParams,
SaveArtifactParams,
@@ -22,6 +32,7 @@ from .models import (
StartRunParams,
ValidateDeploymentParams,
ValidateDraftParams,
ValidateDraftWorkspaceParams,
)
@@ -116,6 +127,125 @@ def create_rpc_app(server: WorkflowServer, *, rpc_path: str = "/rpc") -> jsonrpc
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(name="workflow.draft_workspaces.list", errors=[WorkflowRpcError])
async def workflow_draft_workspaces_list(
params: ListDraftWorkspacesParams = Body(
default_factory=ListDraftWorkspacesParams
),
) -> dict[str, Any]:
try:
return await server.api.list_draft_workspaces()
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(name="workflow.draft_workspaces.get", errors=[WorkflowRpcError])
async def workflow_draft_workspaces_get(
params: GetDraftWorkspaceParams = Params(...), # type: ignore[reportArgumentType],
) -> dict[str, Any]:
try:
return await server.api.get_draft_workspace(
workspace_id=params.workspace_id,
include_draft=params.include_draft,
)
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(
name="workflow.draft_workspaces.create_from_capability",
errors=[WorkflowRpcError],
)
async def workflow_draft_workspaces_create_from_capability(
params: CreateDraftFromCapabilityParams = Params(...), # type: ignore[reportArgumentType],
) -> dict[str, Any]:
try:
return await server.api.create_draft_workspace_from_capability(
workspace_id=params.workspace_id,
capability_name=params.capability_name,
name=params.name,
title=params.title,
input_schema=params.input_schema,
state_schema=params.state_schema,
output_schema=params.output_schema,
input=params.input,
output=params.output,
input_map=params.input_map,
output_map=params.output_map,
error_message_source=params.error_message_source,
)
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(
name="workflow.draft_workspaces.patch", errors=[WorkflowRpcError]
)
async def workflow_draft_workspaces_patch(
params: PatchDraftWorkspaceParams = Params(...), # type: ignore[reportArgumentType],
) -> dict[str, Any]:
try:
return await server.api.patch_draft_workspace(
workspace_id=params.workspace_id,
revision=params.revision,
patch=params.patch,
)
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(
name="workflow.draft_workspaces.validate", errors=[WorkflowRpcError]
)
async def workflow_draft_workspaces_validate(
params: ValidateDraftWorkspaceParams = Params(...), # type: ignore[reportArgumentType],
) -> dict[str, Any]:
try:
return await server.api.validate_draft_workspace(
workspace_id=params.workspace_id,
)
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(
name="workflow.draft_workspaces.create_artifact", errors=[WorkflowRpcError]
)
async def workflow_draft_workspaces_create_artifact(
params: CreateArtifactFromWorkspaceParams = Params(...), # type: ignore[reportArgumentType],
) -> dict[str, Any]:
try:
return await server.api.create_artifact_from_workspace(
workspace_id=params.workspace_id,
artifact_id=params.artifact_id,
version=params.version,
title=params.title,
outcomes=tuple(params.outcomes),
kind=params.kind,
description=params.description,
required_capabilities=params.required_capabilities,
source_bindings=params.source_bindings,
created_from_catalog_version=params.created_from_catalog_version,
)
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(
name="workflow.draft_workspaces.create_wrapper", errors=[WorkflowRpcError]
)
async def workflow_draft_workspaces_create_wrapper(
params: CreateWrapperFromWorkspaceParams = Params(...), # type: ignore[reportArgumentType],
) -> dict[str, Any]:
try:
return await server.api.create_wrapper_from_workspace(
workspace_id=params.workspace_id,
artifact_id=params.artifact_id,
version=params.version,
title=params.title,
outcomes=tuple(params.outcomes),
description=params.description,
required_capabilities=params.required_capabilities,
source_bindings=params.source_bindings,
created_from_catalog_version=params.created_from_catalog_version,
)
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(name="workflow.artifacts.save", errors=[WorkflowRpcError])
async def workflow_artifacts_save(
params: SaveArtifactParams = Params(...), # type: ignore[reportArgumentType],
@@ -146,6 +276,63 @@ def create_rpc_app(server: WorkflowServer, *, rpc_path: str = "/rpc") -> jsonrpc
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(name="workflow.artifacts.list", errors=[WorkflowRpcError])
async def workflow_artifacts_list(
params: ListArtifactsParams = Body(default_factory=ListArtifactsParams),
) -> dict[str, Any]:
try:
return await server.api.list_artifacts(
query=params.query,
kind=params.kind,
cursor=params.cursor,
limit=params.limit,
)
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(name="workflow.artifacts.inspect", errors=[WorkflowRpcError])
async def workflow_artifacts_inspect(
params: InspectArtifactParams = Params(...), # type: ignore[reportArgumentType],
) -> dict[str, Any]:
try:
return await server.api.inspect_artifact(
artifact_id=params.artifact_id,
version=params.version,
)
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(name="workflow.deployments.list", errors=[WorkflowRpcError])
async def workflow_deployments_list(
params: ListDeploymentsParams = Body(default_factory=ListDeploymentsParams),
) -> dict[str, Any]:
try:
return await server.api.list_deployments()
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(name="workflow.deployments.inspect", errors=[WorkflowRpcError])
async def workflow_deployments_inspect(
params: InspectDeploymentParams = Params(...), # type: ignore[reportArgumentType],
) -> dict[str, Any]:
try:
return await server.api.inspect_deployment(
deployment_id=params.deployment_id,
)
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(name="workflow.deployments.delete", errors=[WorkflowRpcError])
async def workflow_deployments_delete(
params: DeleteDeploymentParams = Params(...), # type: ignore[reportArgumentType],
) -> dict[str, Any]:
try:
return await server.api.delete_deployment(
deployment_id=params.deployment_id,
)
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
raise_workflow_rpc_error(exc)
@entrypoint.method(name="workflow.runs.start", errors=[WorkflowRpcError])
async def workflow_runs_start(
params: StartRunParams = Params(...), # type: ignore[reportArgumentType],
+196 -2
View File
@@ -1,7 +1,8 @@
from __future__ import annotations
from collections.abc import Sequence
from dataclasses import dataclass
from typing import Any
from typing import Any, Literal
from uuid import uuid4
import httpx
@@ -14,7 +15,7 @@ class RpcWorkflowApiClient:
"""Small WorkflowApi-compatible adapter for JSON-RPC HTTP targets.
This is intentionally not a full WorkflowApi clone. It implements only the
methods used by the first remote CLI slice.
methods used by CLI commands targeting rpc_http transports.
"""
url: str
@@ -47,6 +48,8 @@ class RpcWorkflowApiClient:
raise RuntimeError("JSON-RPC response result must be an object")
return result
# -- capabilities --
async def list_capabilities(
self,
*,
@@ -71,6 +74,197 @@ class RpcWorkflowApiClient:
{"qualified_name": qualified_name},
)
# -- draft workspaces --
async def list_draft_workspaces(self) -> dict[str, Any]:
return await self._call("workflow.draft_workspaces.list", {})
async def get_draft_workspace(
self,
*,
workspace_id: str,
include_draft: bool = False,
) -> dict[str, Any]:
return await self._call(
"workflow.draft_workspaces.get",
{"workspace_id": workspace_id, "include_draft": include_draft},
)
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]:
return await self._call(
"workflow.draft_workspaces.create_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 patch_draft_workspace(
self,
*,
workspace_id: str,
revision: int,
patch: list[dict[str, Any]],
) -> dict[str, Any]:
return await self._call(
"workflow.draft_workspaces.patch",
{"workspace_id": workspace_id, "revision": revision, "patch": patch},
)
async def validate_draft_workspace(
self,
*,
workspace_id: str,
) -> dict[str, Any]:
return await self._call(
"workflow.draft_workspaces.validate",
{"workspace_id": workspace_id},
)
async def create_artifact_from_workspace(
self,
*,
workspace_id: str,
artifact_id: str,
version: int,
title: str,
outcomes: Sequence[str],
kind: Literal["workflow", "wrapper"] = "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]:
return await self._call(
"workflow.draft_workspaces.create_artifact",
{
"workspace_id": workspace_id,
"artifact_id": artifact_id,
"version": version,
"title": title,
"outcomes": list(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]:
return await self._call(
"workflow.draft_workspaces.create_wrapper",
{
"workspace_id": workspace_id,
"artifact_id": artifact_id,
"version": version,
"title": title,
"outcomes": list(outcomes),
"description": description,
"required_capabilities": required_capabilities,
"source_bindings": source_bindings,
"created_from_catalog_version": created_from_catalog_version,
},
)
# -- artifacts --
async def list_artifacts(
self,
*,
query: str | None = None,
kind: Literal["workflow", "wrapper"] | None = None,
cursor: str | None = None,
limit: int = 50,
) -> dict[str, Any]:
return await self._call(
"workflow.artifacts.list",
{
"query": query,
"kind": kind,
"cursor": cursor,
"limit": limit,
},
)
async def inspect_artifact(
self, *, artifact_id: str, version: int
) -> dict[str, Any]:
return await self._call(
"workflow.artifacts.inspect",
{"artifact_id": artifact_id, "version": version},
)
async def save_artifact(self, artifact: dict[str, Any]) -> dict[str, Any]:
return await self._call("workflow.artifacts.save", {"artifact": artifact})
# -- deployments --
async def list_deployments(self) -> dict[str, Any]:
return await self._call("workflow.deployments.list", {})
async def inspect_deployment(self, *, deployment_id: str) -> dict[str, Any]:
return await self._call(
"workflow.deployments.inspect",
{"deployment_id": deployment_id},
)
async def validate_deployment(
self, *, deployment_id: str, live_check: bool = False
) -> dict[str, Any]:
return await self._call(
"workflow.deployments.validate",
{"deployment_id": deployment_id, "live_check": live_check},
)
async def save_deployment(self, deployment: dict[str, Any]) -> dict[str, Any]:
return await self._call("workflow.deployments.save", {"deployment": deployment})
async def delete_deployment(self, *, deployment_id: str) -> dict[str, Any]:
return await self._call(
"workflow.deployments.delete",
{"deployment_id": deployment_id},
)
# -- runs --
async def run_deployment(
self,
*,
+69 -1
View File
@@ -1,6 +1,6 @@
from __future__ import annotations
from typing import Any
from typing import Any, Literal
from pydantic import BaseModel, ConfigDict, Field
@@ -73,6 +73,74 @@ class SaveDeploymentParams(RpcParamsModel):
deployment: dict[str, Any]
class ListDraftWorkspacesParams(RpcParamsModel):
pass
class GetDraftWorkspaceParams(RpcParamsModel):
workspace_id: str = Field(min_length=1)
include_draft: bool = False
class PatchDraftWorkspaceParams(RpcParamsModel):
workspace_id: str = Field(min_length=1)
revision: int = Field(ge=1)
patch: list[dict[str, Any]]
class ValidateDraftWorkspaceParams(RpcParamsModel):
workspace_id: str = Field(min_length=1)
class CreateArtifactFromWorkspaceParams(RpcParamsModel):
workspace_id: str = Field(min_length=1)
artifact_id: str = Field(min_length=1)
version: int = Field(ge=1)
title: str = Field(min_length=1)
outcomes: list[str]
kind: Literal["workflow", "wrapper"] = "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
class CreateWrapperFromWorkspaceParams(RpcParamsModel):
workspace_id: str = Field(min_length=1)
artifact_id: str = Field(min_length=1)
version: int = Field(ge=1)
title: str = Field(min_length=1)
outcomes: list[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
class ListArtifactsParams(RpcParamsModel):
query: str | None = None
kind: Literal["workflow", "wrapper"] | None = None
cursor: str | None = None
limit: int = Field(default=50, ge=1, le=100)
class InspectArtifactParams(RpcParamsModel):
artifact_id: str = Field(min_length=1)
version: int = Field(ge=1)
class ListDeploymentsParams(RpcParamsModel):
pass
class InspectDeploymentParams(RpcParamsModel):
deployment_id: str = Field(min_length=1)
class DeleteDeploymentParams(RpcParamsModel):
deployment_id: str = Field(min_length=1)
class ValidateDeploymentParams(RpcParamsModel):
deployment_id: str = Field(min_length=1)
live_check: bool = False