refactor: split workflow rpc transport by domain
This commit is contained in:
@@ -3,37 +3,15 @@ from __future__ import annotations
|
|||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
import fastapi_jsonrpc as jsonrpc
|
import fastapi_jsonrpc as jsonrpc
|
||||||
from fastapi import Body
|
|
||||||
from fastapi_jsonrpc import Params
|
|
||||||
|
|
||||||
from wf_server import WorkflowServer
|
from wf_server import WorkflowServer
|
||||||
|
|
||||||
from .errors import WorkflowRpcError, raise_workflow_rpc_error
|
from .errors import WorkflowRpcError
|
||||||
from .models import (
|
from .methods_artifacts import register_methods as register_artifact_methods
|
||||||
CreateArtifactFromWorkspaceParams,
|
from .methods_capabilities import register_methods as register_capability_methods
|
||||||
CreateDraftFromCapabilityParams,
|
from .methods_deployments import register_methods as register_deployment_methods
|
||||||
CreateWrapperFromWorkspaceParams,
|
from .methods_drafts import register_methods as register_draft_methods
|
||||||
DeleteDeploymentParams,
|
from .methods_runs import register_methods as register_run_methods
|
||||||
GetDraftWorkspaceParams,
|
|
||||||
InspectArtifactParams,
|
|
||||||
InspectCapabilityParams,
|
|
||||||
InspectDeploymentParams,
|
|
||||||
InspectRunParams,
|
|
||||||
ListArtifactsParams,
|
|
||||||
ListCapabilitiesParams,
|
|
||||||
ListDeploymentsParams,
|
|
||||||
ListDraftWorkspacesParams,
|
|
||||||
PatchDraftParams,
|
|
||||||
PatchDraftWorkspaceParams,
|
|
||||||
ReadRunTraceParams,
|
|
||||||
ResumeRunParams,
|
|
||||||
SaveArtifactParams,
|
|
||||||
SaveDeploymentParams,
|
|
||||||
StartRunParams,
|
|
||||||
ValidateDeploymentParams,
|
|
||||||
ValidateDraftParams,
|
|
||||||
ValidateDraftWorkspaceParams,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def create_rpc_app(server: WorkflowServer, *, rpc_path: str = "/rpc") -> jsonrpc.API:
|
def create_rpc_app(server: WorkflowServer, *, rpc_path: str = "/rpc") -> jsonrpc.API:
|
||||||
@@ -60,334 +38,11 @@ def create_rpc_app(server: WorkflowServer, *, rpc_path: str = "/rpc") -> jsonrpc
|
|||||||
"store_root": str(server.config.store_root),
|
"store_root": str(server.config.store_root),
|
||||||
}
|
}
|
||||||
|
|
||||||
@entrypoint.method(name="workflow.capabilities.list", errors=[WorkflowRpcError])
|
register_capability_methods(entrypoint, server)
|
||||||
async def workflow_capabilities_list(
|
register_draft_methods(entrypoint, server)
|
||||||
params: ListCapabilitiesParams = Body(default_factory=ListCapabilitiesParams),
|
register_artifact_methods(entrypoint, server)
|
||||||
) -> dict[str, Any]:
|
register_deployment_methods(entrypoint, server)
|
||||||
try:
|
register_run_methods(entrypoint, server)
|
||||||
return await server.api.list_capabilities(
|
|
||||||
query=params.query,
|
|
||||||
source_id=params.source_id,
|
|
||||||
cursor=params.cursor,
|
|
||||||
limit=params.limit,
|
|
||||||
)
|
|
||||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
|
||||||
raise_workflow_rpc_error(exc)
|
|
||||||
|
|
||||||
@entrypoint.method(name="workflow.capabilities.inspect", errors=[WorkflowRpcError])
|
|
||||||
async def workflow_capabilities_inspect(
|
|
||||||
params: InspectCapabilityParams = Params(...), # type: ignore[reportArgumentType]
|
|
||||||
) -> dict[str, Any]:
|
|
||||||
try:
|
|
||||||
return await server.api.inspect_capability(
|
|
||||||
qualified_name=params.qualified_name,
|
|
||||||
)
|
|
||||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
|
||||||
raise_workflow_rpc_error(exc)
|
|
||||||
|
|
||||||
@entrypoint.method(
|
|
||||||
name="workflow.drafts.create_from_capability", errors=[WorkflowRpcError]
|
|
||||||
)
|
|
||||||
async def workflow_drafts_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.drafts.patch", errors=[WorkflowRpcError])
|
|
||||||
async def workflow_drafts_patch(
|
|
||||||
params: PatchDraftParams = Params(...), # type: ignore[reportArgumentType],
|
|
||||||
) -> dict[str, Any]:
|
|
||||||
try:
|
|
||||||
return await server.api.patch_draft(draft=params.draft, patch=params.patch)
|
|
||||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
|
||||||
raise_workflow_rpc_error(exc)
|
|
||||||
|
|
||||||
@entrypoint.method(name="workflow.drafts.validate", errors=[WorkflowRpcError])
|
|
||||||
async def workflow_drafts_validate(
|
|
||||||
params: ValidateDraftParams = Params(...), # type: ignore[reportArgumentType],
|
|
||||||
) -> dict[str, Any]:
|
|
||||||
try:
|
|
||||||
return await server.api.validate_draft(draft=params.draft)
|
|
||||||
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],
|
|
||||||
) -> dict[str, Any]:
|
|
||||||
try:
|
|
||||||
return await server.api.save_artifact(params.artifact)
|
|
||||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
|
||||||
raise_workflow_rpc_error(exc)
|
|
||||||
|
|
||||||
@entrypoint.method(name="workflow.deployments.save", errors=[WorkflowRpcError])
|
|
||||||
async def workflow_deployments_save(
|
|
||||||
params: SaveDeploymentParams = Params(...), # type: ignore[reportArgumentType],
|
|
||||||
) -> dict[str, Any]:
|
|
||||||
try:
|
|
||||||
return await server.api.save_deployment(params.deployment)
|
|
||||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
|
||||||
raise_workflow_rpc_error(exc)
|
|
||||||
|
|
||||||
@entrypoint.method(name="workflow.deployments.validate", errors=[WorkflowRpcError])
|
|
||||||
async def workflow_deployments_validate(
|
|
||||||
params: ValidateDeploymentParams = Params(...), # type: ignore[reportArgumentType],
|
|
||||||
) -> dict[str, Any]:
|
|
||||||
try:
|
|
||||||
return await server.api.validate_deployment(
|
|
||||||
deployment_id=params.deployment_id,
|
|
||||||
live_check=params.live_check,
|
|
||||||
)
|
|
||||||
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],
|
|
||||||
) -> dict[str, Any]:
|
|
||||||
try:
|
|
||||||
return await server.api.run_deployment(
|
|
||||||
deployment_id=params.deployment_id,
|
|
||||||
workflow_input=params.workflow_input,
|
|
||||||
trace_range=(
|
|
||||||
params.trace_range.to_api_trace_range()
|
|
||||||
if params.trace_range is not None
|
|
||||||
else None
|
|
||||||
),
|
|
||||||
)
|
|
||||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
|
||||||
raise_workflow_rpc_error(exc)
|
|
||||||
|
|
||||||
@entrypoint.method(name="workflow.runs.inspect", errors=[WorkflowRpcError])
|
|
||||||
async def workflow_runs_inspect(
|
|
||||||
params: InspectRunParams = Params(...), # type: ignore[reportArgumentType],
|
|
||||||
) -> dict[str, Any]:
|
|
||||||
try:
|
|
||||||
return await server.api.inspect_run(run_id=params.run_id)
|
|
||||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
|
||||||
raise_workflow_rpc_error(exc)
|
|
||||||
|
|
||||||
@entrypoint.method(name="workflow.runs.trace", errors=[WorkflowRpcError])
|
|
||||||
async def workflow_runs_trace(
|
|
||||||
params: ReadRunTraceParams = Params(...), # type: ignore[reportArgumentType],
|
|
||||||
) -> dict[str, Any]:
|
|
||||||
try:
|
|
||||||
return await server.api.read_run_trace(
|
|
||||||
run_id=params.run_id,
|
|
||||||
trace_range=params.trace_range.to_api_trace_range(),
|
|
||||||
)
|
|
||||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
|
||||||
raise_workflow_rpc_error(exc)
|
|
||||||
|
|
||||||
@entrypoint.method(name="workflow.runs.resume", errors=[WorkflowRpcError])
|
|
||||||
async def workflow_runs_resume(
|
|
||||||
params: ResumeRunParams = Params(...), # type: ignore[reportArgumentType],
|
|
||||||
) -> dict[str, Any]:
|
|
||||||
try:
|
|
||||||
return await server.api.resume_run(
|
|
||||||
run_id=params.run_id,
|
|
||||||
resume_payload=params.resume_payload,
|
|
||||||
resume_outcome=params.resume_outcome,
|
|
||||||
trace_range=(
|
|
||||||
params.trace_range.to_api_trace_range()
|
|
||||||
if params.trace_range is not None
|
|
||||||
else None
|
|
||||||
),
|
|
||||||
)
|
|
||||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
|
||||||
raise_workflow_rpc_error(exc)
|
|
||||||
|
|
||||||
app.bind_entrypoint(entrypoint)
|
app.bind_entrypoint(entrypoint)
|
||||||
return app
|
return app
|
||||||
|
|||||||
@@ -1,307 +1,24 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
from collections.abc import Sequence
|
import httpx # noqa: F401 # Backcompat for tests patching client.httpx.AsyncClient.
|
||||||
from dataclasses import dataclass
|
|
||||||
from typing import Any, Literal
|
|
||||||
from uuid import uuid4
|
|
||||||
|
|
||||||
import httpx
|
from .client_artifacts import RpcArtifactClientMixin
|
||||||
|
from .client_base import RpcClientTransport
|
||||||
from wf_api.runs import TraceRangeLike
|
from .client_capabilities import RpcCapabilityClientMixin
|
||||||
|
from .client_deployments import RpcDeploymentClientMixin
|
||||||
|
from .client_drafts import RpcDraftClientMixin
|
||||||
|
from .client_runs import RpcRunClientMixin
|
||||||
|
|
||||||
|
|
||||||
@dataclass(slots=True)
|
class RpcWorkflowApiClient(
|
||||||
class RpcWorkflowApiClient:
|
RpcClientTransport,
|
||||||
"""Small WorkflowApi-compatible adapter for JSON-RPC HTTP targets.
|
RpcCapabilityClientMixin,
|
||||||
|
RpcDraftClientMixin,
|
||||||
This is intentionally not a full WorkflowApi clone. It implements only the
|
RpcArtifactClientMixin,
|
||||||
methods used by CLI commands targeting rpc_http transports.
|
RpcDeploymentClientMixin,
|
||||||
"""
|
RpcRunClientMixin,
|
||||||
|
):
|
||||||
url: str
|
"""WorkflowApiSurface implementation backed by JSON-RPC HTTP calls."""
|
||||||
timeout_seconds: float = 30.0
|
|
||||||
http_client: httpx.AsyncClient | None = None
|
|
||||||
|
|
||||||
async def _call(self, method: str, params: dict[str, Any]) -> dict[str, Any]:
|
|
||||||
request = {
|
|
||||||
"jsonrpc": "2.0",
|
|
||||||
"id": uuid4().hex,
|
|
||||||
"method": method,
|
|
||||||
"params": params,
|
|
||||||
}
|
|
||||||
if self.http_client is None:
|
|
||||||
async with httpx.AsyncClient(timeout=self.timeout_seconds) as client:
|
|
||||||
response = await client.post(self.url, json=request)
|
|
||||||
else:
|
|
||||||
response = await self.http_client.post(self.url, json=request)
|
|
||||||
response.raise_for_status()
|
|
||||||
payload = response.json()
|
|
||||||
if "error" in payload:
|
|
||||||
error = payload["error"]
|
|
||||||
message = error.get("message", "JSON-RPC error")
|
|
||||||
data = error.get("data")
|
|
||||||
if isinstance(data, dict) and data.get("message"):
|
|
||||||
message = f"{message}: {data['message']}"
|
|
||||||
raise RuntimeError(message)
|
|
||||||
result = payload.get("result")
|
|
||||||
if not isinstance(result, dict):
|
|
||||||
raise RuntimeError("JSON-RPC response result must be an object")
|
|
||||||
return result
|
|
||||||
|
|
||||||
# -- capabilities --
|
|
||||||
|
|
||||||
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 await self._call(
|
|
||||||
"workflow.capabilities.list",
|
|
||||||
{
|
|
||||||
"query": query,
|
|
||||||
"source_id": source_id,
|
|
||||||
"cursor": cursor,
|
|
||||||
"limit": limit,
|
|
||||||
},
|
|
||||||
)
|
|
||||||
|
|
||||||
async def inspect_capability(self, *, qualified_name: str) -> dict[str, Any]:
|
|
||||||
return await self._call(
|
|
||||||
"workflow.capabilities.inspect",
|
|
||||||
{"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,
|
|
||||||
*,
|
|
||||||
deployment_id: str,
|
|
||||||
workflow_input: dict[str, Any],
|
|
||||||
trace_range: TraceRangeLike | None = None,
|
|
||||||
) -> dict[str, Any]:
|
|
||||||
return await self._call(
|
|
||||||
"workflow.runs.start",
|
|
||||||
{
|
|
||||||
"deployment_id": deployment_id,
|
|
||||||
"workflow_input": workflow_input,
|
|
||||||
"trace_range": _trace_range_payload(trace_range),
|
|
||||||
},
|
|
||||||
)
|
|
||||||
|
|
||||||
async def inspect_run(self, *, run_id: str) -> dict[str, Any]:
|
|
||||||
return await self._call("workflow.runs.inspect", {"run_id": run_id})
|
|
||||||
|
|
||||||
async def read_run_trace(
|
|
||||||
self,
|
|
||||||
*,
|
|
||||||
run_id: str,
|
|
||||||
trace_range: TraceRangeLike,
|
|
||||||
) -> dict[str, Any]:
|
|
||||||
return await self._call(
|
|
||||||
"workflow.runs.trace",
|
|
||||||
{
|
|
||||||
"run_id": run_id,
|
|
||||||
"trace_range": _trace_range_payload(trace_range),
|
|
||||||
},
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def _trace_range_payload(
|
__all__ = ["RpcWorkflowApiClient"]
|
||||||
trace_range: TraceRangeLike | None,
|
|
||||||
) -> dict[str, int] | None:
|
|
||||||
if trace_range is None:
|
|
||||||
return None
|
|
||||||
return {"start": trace_range.start, "limit": trace_range.limit}
|
|
||||||
|
|||||||
@@ -0,0 +1,38 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Any, Literal
|
||||||
|
|
||||||
|
|
||||||
|
class RpcArtifactClientMixin:
|
||||||
|
"""JSON-RPC implementation of workflow artifact surface methods."""
|
||||||
|
|
||||||
|
async def _call(self, method: str, params: dict[str, Any]) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
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})
|
||||||
@@ -0,0 +1,42 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from dataclasses import dataclass
|
||||||
|
from typing import Any
|
||||||
|
from uuid import uuid4
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(slots=True)
|
||||||
|
class RpcClientTransport:
|
||||||
|
"""Shared JSON-RPC request plumbing for workflow RPC client mixins."""
|
||||||
|
|
||||||
|
url: str
|
||||||
|
timeout_seconds: float = 30.0
|
||||||
|
http_client: httpx.AsyncClient | None = None
|
||||||
|
|
||||||
|
async def _call(self, method: str, params: dict[str, Any]) -> dict[str, Any]:
|
||||||
|
request = {
|
||||||
|
"jsonrpc": "2.0",
|
||||||
|
"id": uuid4().hex,
|
||||||
|
"method": method,
|
||||||
|
"params": params,
|
||||||
|
}
|
||||||
|
if self.http_client is None:
|
||||||
|
async with httpx.AsyncClient(timeout=self.timeout_seconds) as client:
|
||||||
|
response = await client.post(self.url, json=request)
|
||||||
|
else:
|
||||||
|
response = await self.http_client.post(self.url, json=request)
|
||||||
|
response.raise_for_status()
|
||||||
|
payload = response.json()
|
||||||
|
if "error" in payload:
|
||||||
|
error = payload["error"]
|
||||||
|
message = error.get("message", "JSON-RPC error")
|
||||||
|
data = error.get("data")
|
||||||
|
if isinstance(data, dict) and data.get("message"):
|
||||||
|
message = f"{message}: {data['message']}"
|
||||||
|
raise RuntimeError(message)
|
||||||
|
result = payload.get("result")
|
||||||
|
if not isinstance(result, dict):
|
||||||
|
raise RuntimeError("JSON-RPC response result must be an object")
|
||||||
|
return result
|
||||||
@@ -0,0 +1,33 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
|
||||||
|
class RpcCapabilityClientMixin:
|
||||||
|
"""JSON-RPC implementation of workflow capability surface methods."""
|
||||||
|
|
||||||
|
async def _call(self, method: str, params: dict[str, Any]) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
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 await self._call(
|
||||||
|
"workflow.capabilities.list",
|
||||||
|
{
|
||||||
|
"query": query,
|
||||||
|
"source_id": source_id,
|
||||||
|
"cursor": cursor,
|
||||||
|
"limit": limit,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
async def inspect_capability(self, *, qualified_name: str) -> dict[str, Any]:
|
||||||
|
return await self._call(
|
||||||
|
"workflow.capabilities.inspect",
|
||||||
|
{"qualified_name": qualified_name},
|
||||||
|
)
|
||||||
@@ -0,0 +1,35 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
|
||||||
|
class RpcDeploymentClientMixin:
|
||||||
|
"""JSON-RPC implementation of workflow deployment surface methods."""
|
||||||
|
|
||||||
|
async def _call(self, method: str, params: dict[str, Any]) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
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},
|
||||||
|
)
|
||||||
@@ -0,0 +1,138 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from collections.abc import Sequence
|
||||||
|
from typing import Any, Literal
|
||||||
|
|
||||||
|
|
||||||
|
class RpcDraftClientMixin:
|
||||||
|
"""JSON-RPC implementation of workflow draft workspace surface methods."""
|
||||||
|
|
||||||
|
async def _call(self, method: str, params: dict[str, Any]) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
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,
|
||||||
|
},
|
||||||
|
)
|
||||||
@@ -0,0 +1,52 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
from wf_api.runs import TraceRangeLike
|
||||||
|
|
||||||
|
|
||||||
|
class RpcRunClientMixin:
|
||||||
|
"""JSON-RPC implementation of workflow run lifecycle surface methods."""
|
||||||
|
|
||||||
|
async def _call(self, method: str, params: dict[str, Any]) -> dict[str, Any]: ...
|
||||||
|
|
||||||
|
async def run_deployment(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
deployment_id: str,
|
||||||
|
workflow_input: dict[str, Any],
|
||||||
|
trace_range: TraceRangeLike | None = None,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
return await self._call(
|
||||||
|
"workflow.runs.start",
|
||||||
|
{
|
||||||
|
"deployment_id": deployment_id,
|
||||||
|
"workflow_input": workflow_input,
|
||||||
|
"trace_range": _trace_range_payload(trace_range),
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
async def inspect_run(self, *, run_id: str) -> dict[str, Any]:
|
||||||
|
return await self._call("workflow.runs.inspect", {"run_id": run_id})
|
||||||
|
|
||||||
|
async def read_run_trace(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
run_id: str,
|
||||||
|
trace_range: TraceRangeLike,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
return await self._call(
|
||||||
|
"workflow.runs.trace",
|
||||||
|
{
|
||||||
|
"run_id": run_id,
|
||||||
|
"trace_range": _trace_range_payload(trace_range),
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _trace_range_payload(
|
||||||
|
trace_range: TraceRangeLike | None,
|
||||||
|
) -> dict[str, int] | None:
|
||||||
|
if trace_range is None:
|
||||||
|
return None
|
||||||
|
return {"start": trace_range.start, "limit": trace_range.limit}
|
||||||
@@ -0,0 +1,54 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
from fastapi import Body
|
||||||
|
import fastapi_jsonrpc as jsonrpc
|
||||||
|
from fastapi_jsonrpc import Params
|
||||||
|
|
||||||
|
from wf_server import WorkflowServer
|
||||||
|
|
||||||
|
from .errors import WorkflowRpcError, raise_workflow_rpc_error
|
||||||
|
from .models import InspectArtifactParams, ListArtifactsParams, SaveArtifactParams
|
||||||
|
|
||||||
|
|
||||||
|
def register_methods(
|
||||||
|
entrypoint: jsonrpc.Entrypoint,
|
||||||
|
server: WorkflowServer,
|
||||||
|
) -> None:
|
||||||
|
"""Register artifact JSON-RPC methods."""
|
||||||
|
|
||||||
|
@entrypoint.method(name="workflow.artifacts.save", errors=[WorkflowRpcError])
|
||||||
|
async def workflow_artifacts_save(
|
||||||
|
params: SaveArtifactParams = Params(...), # type: ignore[reportArgumentType]
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
try:
|
||||||
|
return await server.api.save_artifact(params.artifact)
|
||||||
|
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)
|
||||||
@@ -0,0 +1,44 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
from fastapi import Body
|
||||||
|
import fastapi_jsonrpc as jsonrpc
|
||||||
|
from fastapi_jsonrpc import Params
|
||||||
|
|
||||||
|
from wf_server import WorkflowServer
|
||||||
|
|
||||||
|
from .errors import WorkflowRpcError, raise_workflow_rpc_error
|
||||||
|
from .models import InspectCapabilityParams, ListCapabilitiesParams
|
||||||
|
|
||||||
|
|
||||||
|
def register_methods(
|
||||||
|
entrypoint: jsonrpc.Entrypoint,
|
||||||
|
server: WorkflowServer,
|
||||||
|
) -> None:
|
||||||
|
"""Register capability discovery JSON-RPC methods."""
|
||||||
|
|
||||||
|
@entrypoint.method(name="workflow.capabilities.list", errors=[WorkflowRpcError])
|
||||||
|
async def workflow_capabilities_list(
|
||||||
|
params: ListCapabilitiesParams = Body(default_factory=ListCapabilitiesParams),
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
try:
|
||||||
|
return await server.api.list_capabilities(
|
||||||
|
query=params.query,
|
||||||
|
source_id=params.source_id,
|
||||||
|
cursor=params.cursor,
|
||||||
|
limit=params.limit,
|
||||||
|
)
|
||||||
|
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||||
|
raise_workflow_rpc_error(exc)
|
||||||
|
|
||||||
|
@entrypoint.method(name="workflow.capabilities.inspect", errors=[WorkflowRpcError])
|
||||||
|
async def workflow_capabilities_inspect(
|
||||||
|
params: InspectCapabilityParams = Params(...), # type: ignore[reportArgumentType]
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
try:
|
||||||
|
return await server.api.inspect_capability(
|
||||||
|
qualified_name=params.qualified_name,
|
||||||
|
)
|
||||||
|
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||||
|
raise_workflow_rpc_error(exc)
|
||||||
@@ -0,0 +1,77 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
from fastapi import Body
|
||||||
|
import fastapi_jsonrpc as jsonrpc
|
||||||
|
from fastapi_jsonrpc import Params
|
||||||
|
|
||||||
|
from wf_server import WorkflowServer
|
||||||
|
|
||||||
|
from .errors import WorkflowRpcError, raise_workflow_rpc_error
|
||||||
|
from .models import (
|
||||||
|
DeleteDeploymentParams,
|
||||||
|
InspectDeploymentParams,
|
||||||
|
ListDeploymentsParams,
|
||||||
|
SaveDeploymentParams,
|
||||||
|
ValidateDeploymentParams,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def register_methods(
|
||||||
|
entrypoint: jsonrpc.Entrypoint,
|
||||||
|
server: WorkflowServer,
|
||||||
|
) -> None:
|
||||||
|
"""Register deployment JSON-RPC methods."""
|
||||||
|
|
||||||
|
@entrypoint.method(name="workflow.deployments.save", errors=[WorkflowRpcError])
|
||||||
|
async def workflow_deployments_save(
|
||||||
|
params: SaveDeploymentParams = Params(...), # type: ignore[reportArgumentType]
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
try:
|
||||||
|
return await server.api.save_deployment(params.deployment)
|
||||||
|
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||||
|
raise_workflow_rpc_error(exc)
|
||||||
|
|
||||||
|
@entrypoint.method(name="workflow.deployments.validate", errors=[WorkflowRpcError])
|
||||||
|
async def workflow_deployments_validate(
|
||||||
|
params: ValidateDeploymentParams = Params(...), # type: ignore[reportArgumentType]
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
try:
|
||||||
|
return await server.api.validate_deployment(
|
||||||
|
deployment_id=params.deployment_id,
|
||||||
|
live_check=params.live_check,
|
||||||
|
)
|
||||||
|
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)
|
||||||
@@ -0,0 +1,184 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
from fastapi import Body
|
||||||
|
import fastapi_jsonrpc as jsonrpc
|
||||||
|
from fastapi_jsonrpc import Params
|
||||||
|
|
||||||
|
from wf_server import WorkflowServer
|
||||||
|
|
||||||
|
from .errors import WorkflowRpcError, raise_workflow_rpc_error
|
||||||
|
from .models import (
|
||||||
|
CreateArtifactFromWorkspaceParams,
|
||||||
|
CreateDraftFromCapabilityParams,
|
||||||
|
CreateWrapperFromWorkspaceParams,
|
||||||
|
GetDraftWorkspaceParams,
|
||||||
|
ListDraftWorkspacesParams,
|
||||||
|
PatchDraftParams,
|
||||||
|
PatchDraftWorkspaceParams,
|
||||||
|
ValidateDraftParams,
|
||||||
|
ValidateDraftWorkspaceParams,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def register_methods(
|
||||||
|
entrypoint: jsonrpc.Entrypoint,
|
||||||
|
server: WorkflowServer,
|
||||||
|
) -> None:
|
||||||
|
"""Register draft and draft-workspace JSON-RPC methods."""
|
||||||
|
|
||||||
|
@entrypoint.method(
|
||||||
|
name="workflow.drafts.create_from_capability", errors=[WorkflowRpcError]
|
||||||
|
)
|
||||||
|
async def workflow_drafts_create_from_capability(
|
||||||
|
params: CreateDraftFromCapabilityParams = Params(...), # type: ignore[reportArgumentType]
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
try:
|
||||||
|
return await _create_from_capability(server, params)
|
||||||
|
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||||
|
raise_workflow_rpc_error(exc)
|
||||||
|
|
||||||
|
@entrypoint.method(name="workflow.drafts.patch", errors=[WorkflowRpcError])
|
||||||
|
async def workflow_drafts_patch(
|
||||||
|
params: PatchDraftParams = Params(...), # type: ignore[reportArgumentType]
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
try:
|
||||||
|
return await server.api.patch_draft(draft=params.draft, patch=params.patch)
|
||||||
|
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||||
|
raise_workflow_rpc_error(exc)
|
||||||
|
|
||||||
|
@entrypoint.method(name="workflow.drafts.validate", errors=[WorkflowRpcError])
|
||||||
|
async def workflow_drafts_validate(
|
||||||
|
params: ValidateDraftParams = Params(...), # type: ignore[reportArgumentType]
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
try:
|
||||||
|
return await server.api.validate_draft(draft=params.draft)
|
||||||
|
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 _create_from_capability(server, params)
|
||||||
|
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)
|
||||||
|
|
||||||
|
|
||||||
|
async def _create_from_capability(
|
||||||
|
server: WorkflowServer,
|
||||||
|
params: CreateDraftFromCapabilityParams,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
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,
|
||||||
|
)
|
||||||
@@ -0,0 +1,79 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
import fastapi_jsonrpc as jsonrpc
|
||||||
|
from fastapi_jsonrpc import Params
|
||||||
|
|
||||||
|
from wf_server import WorkflowServer
|
||||||
|
|
||||||
|
from .errors import WorkflowRpcError, raise_workflow_rpc_error
|
||||||
|
from .models import (
|
||||||
|
InspectRunParams,
|
||||||
|
ReadRunTraceParams,
|
||||||
|
ResumeRunParams,
|
||||||
|
StartRunParams,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def register_methods(
|
||||||
|
entrypoint: jsonrpc.Entrypoint,
|
||||||
|
server: WorkflowServer,
|
||||||
|
) -> None:
|
||||||
|
"""Register run lifecycle JSON-RPC methods."""
|
||||||
|
|
||||||
|
@entrypoint.method(name="workflow.runs.start", errors=[WorkflowRpcError])
|
||||||
|
async def workflow_runs_start(
|
||||||
|
params: StartRunParams = Params(...), # type: ignore[reportArgumentType]
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
try:
|
||||||
|
return await server.api.run_deployment(
|
||||||
|
deployment_id=params.deployment_id,
|
||||||
|
workflow_input=params.workflow_input,
|
||||||
|
trace_range=(
|
||||||
|
params.trace_range.to_api_trace_range()
|
||||||
|
if params.trace_range is not None
|
||||||
|
else None
|
||||||
|
),
|
||||||
|
)
|
||||||
|
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||||
|
raise_workflow_rpc_error(exc)
|
||||||
|
|
||||||
|
@entrypoint.method(name="workflow.runs.inspect", errors=[WorkflowRpcError])
|
||||||
|
async def workflow_runs_inspect(
|
||||||
|
params: InspectRunParams = Params(...), # type: ignore[reportArgumentType]
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
try:
|
||||||
|
return await server.api.inspect_run(run_id=params.run_id)
|
||||||
|
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||||
|
raise_workflow_rpc_error(exc)
|
||||||
|
|
||||||
|
@entrypoint.method(name="workflow.runs.trace", errors=[WorkflowRpcError])
|
||||||
|
async def workflow_runs_trace(
|
||||||
|
params: ReadRunTraceParams = Params(...), # type: ignore[reportArgumentType]
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
try:
|
||||||
|
return await server.api.read_run_trace(
|
||||||
|
run_id=params.run_id,
|
||||||
|
trace_range=params.trace_range.to_api_trace_range(),
|
||||||
|
)
|
||||||
|
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||||
|
raise_workflow_rpc_error(exc)
|
||||||
|
|
||||||
|
@entrypoint.method(name="workflow.runs.resume", errors=[WorkflowRpcError])
|
||||||
|
async def workflow_runs_resume(
|
||||||
|
params: ResumeRunParams = Params(...), # type: ignore[reportArgumentType]
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
try:
|
||||||
|
return await server.api.resume_run(
|
||||||
|
run_id=params.run_id,
|
||||||
|
resume_payload=params.resume_payload,
|
||||||
|
resume_outcome=params.resume_outcome,
|
||||||
|
trace_range=(
|
||||||
|
params.trace_range.to_api_trace_range()
|
||||||
|
if params.trace_range is not None
|
||||||
|
else None
|
||||||
|
),
|
||||||
|
)
|
||||||
|
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||||
|
raise_workflow_rpc_error(exc)
|
||||||
@@ -0,0 +1,24 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import importlib
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
|
||||||
|
def test_rpc_transport_has_domain_method_modules() -> None:
|
||||||
|
for module_name in (
|
||||||
|
"wf_transport_rpc_http.methods_capabilities",
|
||||||
|
"wf_transport_rpc_http.methods_drafts",
|
||||||
|
"wf_transport_rpc_http.methods_artifacts",
|
||||||
|
"wf_transport_rpc_http.methods_deployments",
|
||||||
|
"wf_transport_rpc_http.methods_runs",
|
||||||
|
):
|
||||||
|
module = importlib.import_module(module_name)
|
||||||
|
|
||||||
|
assert hasattr(module, "register_methods")
|
||||||
|
|
||||||
|
|
||||||
|
def test_rpc_transport_client_stays_thin() -> None:
|
||||||
|
client_path = Path("src/wf_transport_rpc_http/client.py")
|
||||||
|
line_count = len(client_path.read_text(encoding="utf-8").splitlines())
|
||||||
|
|
||||||
|
assert line_count < 140
|
||||||
Reference in New Issue
Block a user