feat: expose workflow authoring rpc methods
This commit is contained in:
@@ -9,7 +9,16 @@ 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, raise_workflow_rpc_error
|
||||||
from .models import InspectCapabilityParams, ListCapabilitiesParams
|
from .models import (
|
||||||
|
CreateDraftFromCapabilityParams,
|
||||||
|
InspectCapabilityParams,
|
||||||
|
ListCapabilitiesParams,
|
||||||
|
PatchDraftParams,
|
||||||
|
SaveArtifactParams,
|
||||||
|
SaveDeploymentParams,
|
||||||
|
ValidateDeploymentParams,
|
||||||
|
ValidateDraftParams,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def create_rpc_app(server: WorkflowServer) -> jsonrpc.API:
|
def create_rpc_app(server: WorkflowServer) -> jsonrpc.API:
|
||||||
@@ -59,5 +68,75 @@ def create_rpc_app(server: WorkflowServer) -> jsonrpc.API:
|
|||||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||||
raise_workflow_rpc_error(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(...),
|
||||||
|
) -> 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(...),
|
||||||
|
) -> 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(...),
|
||||||
|
) -> 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.artifacts.save", errors=[WorkflowRpcError])
|
||||||
|
async def workflow_artifacts_save(
|
||||||
|
params: SaveArtifactParams = Params(...),
|
||||||
|
) -> 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(...),
|
||||||
|
) -> 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(...),
|
||||||
|
) -> 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)
|
||||||
|
|
||||||
app.bind_entrypoint(entrypoint)
|
app.bind_entrypoint(entrypoint)
|
||||||
return app
|
return app
|
||||||
|
|||||||
@@ -58,3 +58,120 @@ def test_rpc_unknown_method_returns_json_rpc_error(tmp_path) -> None:
|
|||||||
assert payload["error"]["message"] == "Method not found"
|
assert payload["error"]["message"] == "Method not found"
|
||||||
|
|
||||||
asyncio.run(scenario())
|
asyncio.run(scenario())
|
||||||
|
|
||||||
|
|
||||||
|
def test_rpc_draft_artifact_deployment_lifecycle(tmp_path) -> None:
|
||||||
|
async def scenario() -> None:
|
||||||
|
server = build_local_static_workflow_server(tmp_path / "store")
|
||||||
|
app = create_rpc_app(server)
|
||||||
|
transport = httpx.ASGITransport(app=app)
|
||||||
|
async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
|
||||||
|
draft_ws = await _rpc(
|
||||||
|
client,
|
||||||
|
"workflow.drafts.create_from_capability",
|
||||||
|
{
|
||||||
|
"workspace_id": "constant_ws",
|
||||||
|
"capability_name": "wf.std.constant",
|
||||||
|
"name": "constant_workflow",
|
||||||
|
"title": "Constant Workflow",
|
||||||
|
"input_map": {},
|
||||||
|
"output_map": {"value": "state.result"},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
draft = {
|
||||||
|
"name": "rpc_constant",
|
||||||
|
"input_schema": {"type": "object", "properties": {}},
|
||||||
|
"state_schema": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {
|
||||||
|
"result": {"type": "string", "reducer": "wf.std.replace"}
|
||||||
|
},
|
||||||
|
},
|
||||||
|
"output_schema": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {"result": {"type": "string"}},
|
||||||
|
"required": ["result"],
|
||||||
|
},
|
||||||
|
"start": "constant",
|
||||||
|
"steps": {
|
||||||
|
"constant": {
|
||||||
|
"use": "wf.std.constant",
|
||||||
|
"input": [
|
||||||
|
{
|
||||||
|
"value": "hello over rpc",
|
||||||
|
"target": {"root": "local", "parts": ["value"]},
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"output": [
|
||||||
|
{
|
||||||
|
"source": {"root": "local", "parts": ["value"]},
|
||||||
|
"target": {"root": "state", "parts": ["result"]},
|
||||||
|
}
|
||||||
|
],
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"routes": {"constant": {"ok": "__end__"}},
|
||||||
|
"output": [
|
||||||
|
{
|
||||||
|
"path": {"root": "state", "parts": ["result"]},
|
||||||
|
"target": {"root": "local", "parts": ["result"]},
|
||||||
|
}
|
||||||
|
],
|
||||||
|
}
|
||||||
|
|
||||||
|
validate_draft = await _rpc(
|
||||||
|
client,
|
||||||
|
"workflow.drafts.validate",
|
||||||
|
{"draft": draft},
|
||||||
|
)
|
||||||
|
compiled_plan = validate_draft["result"]["compiled_plan"]
|
||||||
|
artifact = await _rpc(
|
||||||
|
client,
|
||||||
|
"workflow.artifacts.save",
|
||||||
|
{
|
||||||
|
"artifact": {
|
||||||
|
"id": "constant_rpc",
|
||||||
|
"version": 1,
|
||||||
|
"kind": "wrapper",
|
||||||
|
"title": "Constant RPC",
|
||||||
|
"input_schema": {"type": "object", "properties": {}},
|
||||||
|
"output_schema": {
|
||||||
|
"type": "object",
|
||||||
|
"properties": {"result": {"type": "string"}},
|
||||||
|
"required": ["result"],
|
||||||
|
},
|
||||||
|
"outcomes": ["ok"],
|
||||||
|
"required_capabilities": {},
|
||||||
|
"source_bindings": {"wf.std": "wf.std"},
|
||||||
|
"plan": compiled_plan,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
deployment = await _rpc(
|
||||||
|
client,
|
||||||
|
"workflow.deployments.save",
|
||||||
|
{
|
||||||
|
"deployment": {
|
||||||
|
"id": "constant_rpc.default",
|
||||||
|
"artifact_id": "constant_rpc",
|
||||||
|
"artifact_version": 1,
|
||||||
|
"bindings": [
|
||||||
|
{"logical_source": "wf.std", "concrete_source": "wf.std"}
|
||||||
|
],
|
||||||
|
},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
validate_deployment = await _rpc(
|
||||||
|
client,
|
||||||
|
"workflow.deployments.validate",
|
||||||
|
{"deployment_id": "constant_rpc.default"},
|
||||||
|
)
|
||||||
|
|
||||||
|
assert draft_ws["result"]["workspace_id"] == "constant_ws"
|
||||||
|
assert validate_draft["result"]["status"] == "valid"
|
||||||
|
assert artifact["result"]["artifact_id"] == "constant_rpc"
|
||||||
|
assert deployment["result"]["deployment_id"] == "constant_rpc.default"
|
||||||
|
assert validate_deployment["result"]["status"] == "runnable"
|
||||||
|
|
||||||
|
asyncio.run(scenario())
|
||||||
|
|||||||
Reference in New Issue
Block a user