list and save artifact/deployments tools hooked
This commit is contained in:
@@ -27,6 +27,9 @@ class WorkflowArtifactStore:
|
||||
def get_deployment(self, deployment_id: str) -> WorkflowDeployment:
|
||||
raise NotImplementedError
|
||||
|
||||
def list_deployments(self) -> list[WorkflowDeployment]:
|
||||
raise NotImplementedError
|
||||
|
||||
|
||||
class FileWorkflowArtifactStore(WorkflowArtifactStore):
|
||||
"""JSON file-backed artifact store for local development and tests."""
|
||||
@@ -89,3 +92,11 @@ class FileWorkflowArtifactStore(WorkflowArtifactStore):
|
||||
if not path.exists():
|
||||
raise KeyError(f"unknown workflow deployment {deployment_id!r}")
|
||||
return WorkflowDeployment.model_validate_json(path.read_text(encoding="utf-8"))
|
||||
|
||||
def list_deployments(self) -> list[WorkflowDeployment]:
|
||||
deployments: list[WorkflowDeployment] = []
|
||||
for path in sorted(self.deployments_dir.glob("*.json")):
|
||||
deployments.append(
|
||||
WorkflowDeployment.model_validate_json(path.read_text(encoding="utf-8"))
|
||||
)
|
||||
return deployments
|
||||
|
||||
@@ -6,6 +6,8 @@ from mcp.server.fastmcp import FastMCP
|
||||
from wf_artifacts import (
|
||||
AvailableCapability,
|
||||
AvailableSource,
|
||||
WorkflowArtifact,
|
||||
WorkflowDeployment,
|
||||
validate_deployment_dependencies,
|
||||
)
|
||||
|
||||
@@ -25,6 +27,18 @@ def register_artifact_tools(server: FastMCP, service: WfMcpService) -> None:
|
||||
]
|
||||
return {"nodes": entries}
|
||||
|
||||
@server.tool()
|
||||
async def save_workflow_artifact(artifact: dict[str, Any]) -> dict[str, Any]:
|
||||
if service.artifact_store is None:
|
||||
raise KeyError("workflow artifact store is not configured")
|
||||
workflow_artifact = WorkflowArtifact.model_validate(artifact)
|
||||
service.artifact_store.save_artifact(workflow_artifact)
|
||||
return {
|
||||
"artifact_id": workflow_artifact.id,
|
||||
"version": workflow_artifact.version,
|
||||
"saved": True,
|
||||
}
|
||||
|
||||
@server.tool()
|
||||
async def inspect_workflow_artifact(
|
||||
artifact_id: str,
|
||||
@@ -35,6 +49,30 @@ def register_artifact_tools(server: FastMCP, service: WfMcpService) -> None:
|
||||
artifact = service.artifact_store.get_artifact(artifact_id, version)
|
||||
return artifact.model_dump(mode="json")
|
||||
|
||||
@server.tool()
|
||||
async def list_workflow_deployments() -> dict[str, Any]:
|
||||
if service.artifact_store is None:
|
||||
return {"deployments": []}
|
||||
return {
|
||||
"deployments": [
|
||||
deployment.model_dump(mode="json")
|
||||
for deployment in service.artifact_store.list_deployments()
|
||||
]
|
||||
}
|
||||
|
||||
@server.tool()
|
||||
async def save_workflow_deployment(deployment: dict[str, Any]) -> dict[str, Any]:
|
||||
if service.artifact_store is None:
|
||||
raise KeyError("workflow artifact store is not configured")
|
||||
workflow_deployment = WorkflowDeployment.model_validate(deployment)
|
||||
service.artifact_store.save_deployment(workflow_deployment)
|
||||
return {
|
||||
"deployment_id": workflow_deployment.id,
|
||||
"artifact_id": workflow_deployment.artifact_id,
|
||||
"artifact_version": workflow_deployment.artifact_version,
|
||||
"saved": True,
|
||||
}
|
||||
|
||||
@server.tool()
|
||||
async def validate_workflow_deployment(deployment_id: str) -> dict[str, Any]:
|
||||
if service.artifact_store is None:
|
||||
|
||||
@@ -59,3 +59,31 @@ def test_file_store_round_trips_deployment(tmp_path) -> None:
|
||||
assert loaded.id == "summarize_docs.personal"
|
||||
assert loaded.artifact_id == "summarize_docs"
|
||||
assert loaded.bindings["context7"] == "context7.personal"
|
||||
|
||||
|
||||
def test_file_store_lists_deployments_in_id_order(tmp_path) -> None:
|
||||
store = FileWorkflowArtifactStore(tmp_path)
|
||||
store.save_deployment(
|
||||
WorkflowDeployment(
|
||||
id="summarize_docs.work",
|
||||
artifact_id="summarize_docs",
|
||||
artifact_version=1,
|
||||
bindings={"context7": "context7.work"},
|
||||
)
|
||||
)
|
||||
store.save_deployment(
|
||||
WorkflowDeployment(
|
||||
id="summarize_docs.personal",
|
||||
artifact_id="summarize_docs",
|
||||
artifact_version=1,
|
||||
bindings={"context7": "context7.personal"},
|
||||
)
|
||||
)
|
||||
|
||||
deployments = store.list_deployments()
|
||||
|
||||
assert [deployment.id for deployment in deployments] == [
|
||||
"summarize_docs.personal",
|
||||
"summarize_docs.work",
|
||||
]
|
||||
assert deployments[0].bindings["context7"] == "context7.personal"
|
||||
|
||||
@@ -258,6 +258,64 @@ def test_broker_validates_workflow_deployment_from_artifact_store() -> None:
|
||||
assert payload["diagnostics"][0]["code"] == "source_missing"
|
||||
|
||||
|
||||
def test_broker_saves_workflow_artifact() -> None:
|
||||
artifact_store = FileWorkflowArtifactStore(
|
||||
local_temp_root() / "broker_save_artifacts"
|
||||
)
|
||||
service = WfMcpService(
|
||||
store=FileStore(local_temp_root() / "broker_save_mcp_store"),
|
||||
artifact_store=artifact_store,
|
||||
)
|
||||
server = create_broker_server(service)
|
||||
|
||||
_content, structured = asyncio.run(
|
||||
server.call_tool(
|
||||
"save_workflow_artifact",
|
||||
{"artifact": _artifact().model_dump(mode="json")},
|
||||
)
|
||||
)
|
||||
payload = cast(dict[str, Any], cast(object, structured))
|
||||
loaded = artifact_store.get_artifact("summarize_docs", 1)
|
||||
|
||||
assert payload["artifact_id"] == "summarize_docs"
|
||||
assert payload["version"] == 1
|
||||
assert loaded.title == "Summarize Docs"
|
||||
|
||||
|
||||
def test_broker_saves_and_lists_workflow_deployments() -> None:
|
||||
artifact_store = FileWorkflowArtifactStore(
|
||||
local_temp_root() / "broker_save_deployments"
|
||||
)
|
||||
service = WfMcpService(
|
||||
store=FileStore(local_temp_root() / "broker_save_deployments_mcp_store"),
|
||||
artifact_store=artifact_store,
|
||||
)
|
||||
server = create_broker_server(service)
|
||||
|
||||
_content, save_structured = asyncio.run(
|
||||
server.call_tool(
|
||||
"save_workflow_deployment",
|
||||
{
|
||||
"deployment": WorkflowDeployment(
|
||||
id="summarize_docs.personal",
|
||||
artifact_id="summarize_docs",
|
||||
artifact_version=1,
|
||||
bindings={"context7": "context7.personal"},
|
||||
).model_dump(mode="json")
|
||||
},
|
||||
)
|
||||
)
|
||||
save_payload = cast(dict[str, Any], cast(object, save_structured))
|
||||
_content, list_structured = asyncio.run(
|
||||
server.call_tool("list_workflow_deployments", {})
|
||||
)
|
||||
list_payload = cast(dict[str, Any], cast(object, list_structured))
|
||||
|
||||
assert save_payload["deployment_id"] == "summarize_docs.personal"
|
||||
assert list_payload["deployments"][0]["id"] == "summarize_docs.personal"
|
||||
assert list_payload["deployments"][0]["bindings"]["context7"] == "context7.personal"
|
||||
|
||||
|
||||
def test_build_service_from_config_uses_store_root_for_artifacts() -> None:
|
||||
store_root = local_temp_root() / "broker_config_artifact_store"
|
||||
service = build_service_from_config(
|
||||
|
||||
Reference in New Issue
Block a user