From 6677cd6ae8279e3cacd19e32894ea27777e257e5 Mon Sep 17 00:00:00 2001 From: James Wong <2421248+jameswnl@users.noreply.github.com> Date: Tue, 29 Sep 2026 18:27:37 -0400 Subject: [PATCH 1/3] Remove /v1/agents/run and use one-step workflows (#55) POST /v1/workflows/run is now the only agent execution endpoint: a one-shot agent invocation is submitted as a one-step workflow using the documented agent/result naming convention with a run-level provider. Intentionally breaking with no compatibility adapter. Removal: - Delete src/app/endpoints/agents.py (manual StepInput construction and get_step_executor dispatch that bypassed workflow middleware, persistence, and transcripts). - Remove AgentRunRequest and its export; document the one-step contract on RunWorkflowRequest instead. - Remove the router registration, the agents OpenAPI tag, and the AGENT_RUN action. Dependency: - Bump lightspeed-cloud-agents a5b30eb -> ffdc8933 to pick up the #268 merge (unified AgentExecutionSpec normalization shared by one-step and multi-step steps on both runners). Tests: - Unit: drop run_agent_handler tests; add one-step forwarding and provider default/override tests to test_workflows_endpoint.py. - Integration: new test_workflows_integration.py pins the one-step contract against real cloud-agents (full input coverage, bare-step defaults, one-step/multi-step normalization equality, secret-value rejection); triage parse test moved over. - E2E: drop agents HTTP/handler files; add a migration-example one-step test to test_workflows_http_e2e.py. Makefile and CI updated. Demos/docs: - Demo script agent-* scenarios become oneshot-* one-step workflows (oneshot-ephemeral gains the allowed_skills allowlist its Landlock demo text always described); demo HTML tabs and both cloud-agents design docs updated; openapi.json regenerated. Verification: unit 3377 passed, integration 264 passed, ruff/black/ docstyle clean, pylint/pyright/mypy report only pre-existing hits in untouched files, live-app route check confirms agents/run is gone. --- .github/workflows/cloud_agents_tests.yaml | 1 - Makefile | 11 +- docs/cloud-agents-demo-curl.sh | 136 +++-- docs/cloud-agents-integration.html | 28 +- .../cloud-agents/e2e-testing-strategy.md | 24 +- .../cloud-agents/integration-architecture.md | 23 +- docs/devel_doc/openapi.json | 201 +------ src/app/endpoints/agents.py | 132 ---- src/app/main.py | 11 +- src/app/routers.py | 5 +- src/models/api/requests/__init__.py | 2 - src/models/api/requests/agents.py | 101 +--- src/models/config.py | 3 - tests/e2e/cloud_agents/conftest.py | 6 +- .../test_agents_run_handler_e2e.py | 311 ---------- .../cloud_agents/test_agents_run_http_e2e.py | 145 ----- .../cloud_agents/test_step_executor_e2e.py | 4 +- .../cloud_agents/test_workflows_http_e2e.py | 74 ++- .../cloud_agents/test_agents_integration.py | 281 --------- .../test_workflows_integration.py | 176 ++++++ tests/unit/app/test_routers.py | 4 +- .../unit/cloud_agents/test_agents_endpoint.py | 563 ------------------ .../cloud_agents/test_workflows_endpoint.py | 105 ++++ uv.lock | 2 +- 24 files changed, 505 insertions(+), 1844 deletions(-) delete mode 100644 src/app/endpoints/agents.py delete mode 100644 tests/e2e/cloud_agents/test_agents_run_handler_e2e.py delete mode 100644 tests/e2e/cloud_agents/test_agents_run_http_e2e.py delete mode 100644 tests/integration/cloud_agents/test_agents_integration.py create mode 100644 tests/integration/cloud_agents/test_workflows_integration.py delete mode 100644 tests/unit/cloud_agents/test_agents_endpoint.py diff --git a/.github/workflows/cloud_agents_tests.yaml b/.github/workflows/cloud_agents_tests.yaml index e1a3ec623..c5d16324e 100644 --- a/.github/workflows/cloud_agents_tests.yaml +++ b/.github/workflows/cloud_agents_tests.yaml @@ -89,7 +89,6 @@ jobs: - name: Run mock-backed e2e tests (spawn=none/local) run: | uv run pytest \ - tests/e2e/cloud_agents/test_agents_run_http_e2e.py \ tests/e2e/cloud_agents/test_workflows_http_e2e.py \ tests/e2e/cloud_agents/test_mock_llm_server.py \ tests/e2e/cloud_agents/test_mock_llm_env.py \ diff --git a/Makefile b/Makefile index b878b82b4..0921787e7 100644 --- a/Makefile +++ b/Makefile @@ -172,7 +172,7 @@ test-e2e-tagged: ## Run e2e tests with E2E_BEHAVE_TAG_EXPR (default: all @cfg_*) test-e2e-tagged-local: ## Same as test-e2e-tagged without script wrapper uv run behave --color --format pretty --tags="$(E2E_BEHAVE_TAG_EXPR)" -D dump_errors=true @tests/e2e/test_list.txt -test-e2e-agents-workflows: ## Real-HTTP e2e tests for /v1/agents/run and /v1/workflows/* (needs OPENAI_API_KEY) +test-e2e-agents-workflows: ## Real-HTTP e2e tests for /v1/workflows/* incl. one-step runs (needs OPENAI_API_KEY) @if [ -z "$$OPENAI_API_KEY" ]; then \ echo "ERROR: OPENAI_API_KEY is not set."; \ exit 1; \ @@ -187,7 +187,7 @@ test-e2e-agents-workflows: ## Real-HTTP e2e tests for /v1/agents/run and /v1/wor fi @echo "Ensuring Postgres is up (needed for workflow run-state/transcript storage)..." $(CONTAINER_RUNTIME) compose -f docker-compose-harness.yaml up -d --wait postgres - uv run pytest tests/e2e/cloud_agents/test_agents_run_http_e2e.py tests/e2e/cloud_agents/test_workflows_http_e2e.py -v + uv run pytest tests/e2e/cloud_agents/test_workflows_http_e2e.py -v test-e2e-agents-workflows-mock: ## Real-HTTP e2e tests for spawn=none/local against a mock LLM (no OPENAI_API_KEY/gateway needed) @if [ -z "$(CONTAINER_RUNTIME)" ]; then \ @@ -201,7 +201,6 @@ test-e2e-agents-workflows-mock: ## Real-HTTP e2e tests for spawn=none/local agai @echo "Ensuring Postgres is up (needed for workflow run-state/transcript storage)..." $(CONTAINER_RUNTIME) compose -f docker-compose-harness.yaml up -d --wait postgres OPENAI_API_KEY=sk-mock-ci-key LIGHTSPEED_E2E_USE_MOCK_LLM=1 uv run pytest \ - tests/e2e/cloud_agents/test_agents_run_http_e2e.py \ tests/e2e/cloud_agents/test_workflows_http_e2e.py \ tests/e2e/cloud_agents/test_mock_llm_server.py \ tests/e2e/cloud_agents/test_mock_llm_env.py \ @@ -243,13 +242,13 @@ demo-agents-workflows-server: ## Start lightspeed-stack ready for the cloud-agen OPENSHELL_GATEWAY_URL="$${OPENSHELL_GATEWAY_URL:-localhost:17670}" \ uv run python src/lightspeed_stack.py -c lightspeed-stack-harness.yaml -demo-agents-workflows: ## Curl-based live demo against a running server (make demo-agents-workflows SCENARIO=agent-none|agent-local|agent-ephemeral|workflow-ephemeral-approval|workflow-none-approval|workflow-local|workflow-ephemeral|discover) +demo-agents-workflows: ## Curl-based live demo against a running server (make demo-agents-workflows SCENARIO=oneshot-none|oneshot-local|oneshot-ephemeral|workflow-ephemeral-approval|workflow-none-approval|workflow-local|workflow-ephemeral|discover) @if ! curl -s -o /dev/null --max-time 2 "$${BASE_URL:-http://localhost:8090}/v1/info" 2>/dev/null; then \ echo "ERROR: no server reachable at $${BASE_URL:-http://localhost:8090}."; \ echo "Start one with: uv run make demo-agents-workflows-server"; \ exit 1; \ fi - ./docs/cloud-agents-demo-curl.sh $(or $(SCENARIO),agent-none) + ./docs/cloud-agents-demo-curl.sh $(or $(SCENARIO),oneshot-none) demo-agents-workflows-all: ## Run every demo scenario against a running server (see demo-agents-workflows-server); does not include discover @if ! curl -s -o /dev/null --max-time 2 "$${BASE_URL:-http://localhost:8090}/v1/info" 2>/dev/null; then \ @@ -258,7 +257,7 @@ demo-agents-workflows-all: ## Run every demo scenario against a running server ( exit 1; \ fi @failed=""; \ - for scenario in agent-none agent-local agent-ephemeral workflow-ephemeral-approval workflow-none-approval workflow-local workflow-ephemeral; do \ + for scenario in oneshot-none oneshot-local oneshot-ephemeral workflow-ephemeral-approval workflow-none-approval workflow-local workflow-ephemeral; do \ echo; \ echo "########## $$scenario ##########"; \ if ! ./docs/cloud-agents-demo-curl.sh "$$scenario"; then \ diff --git a/docs/cloud-agents-demo-curl.sh b/docs/cloud-agents-demo-curl.sh index 9e5437e59..df293900d 100755 --- a/docs/cloud-agents-demo-curl.sh +++ b/docs/cloud-agents-demo-curl.sh @@ -1,41 +1,45 @@ #!/usr/bin/env bash -# Live-demo curl commands covering the full /v1/agents/run and -# /v1/workflows/run endpoint x spawn-mode matrix, matching -# tests/e2e/cloud_agents/test_agents_run_http_e2e.py and -# test_workflows_http_e2e.py. +# Live-demo curl commands covering the /v1/workflows/run spawn-mode +# matrix, matching tests/e2e/cloud_agents/test_workflows_http_e2e.py. +# POST /v1/workflows/run is the only agent execution endpoint: a one-shot +# agent invocation is a one-step workflow (the oneshot-* scenarios below +# use the documented agent/result naming convention with a run-level +# provider). # -# agent-none / agent-ephemeral / workflow-ephemeral-approval illustrate -# the three tabs in docs/cloud-agents-integration.html. The rest -# (agent-local, workflow-none-approval, workflow-local, workflow-ephemeral) -# round out the same matrix for parity with the automated e2e suite -- -# they aren't part of that illustration, just additional scenarios: -# agent-none POST /v1/agents/run, spawn: "none" -# agent-local POST /v1/agents/run, spawn: "local" -# agent-ephemeral POST /v1/agents/run, spawn: "ephemeral", k8s-diag skill + Landlock demo -# workflow-ephemeral-approval POST /v1/workflows/run, spawn: "ephemeral", multi-step + approval -# workflow-none-approval POST /v1/workflows/run, spawn: "none", multi-step + approval -# workflow-local POST /v1/workflows/run, spawn: "local", single step -# workflow-ephemeral POST /v1/workflows/run, spawn: "ephemeral", single step, no approval +# oneshot-none / oneshot-ephemeral / workflow-ephemeral-approval +# illustrate the three tabs in docs/cloud-agents-integration.html. The +# rest (oneshot-local, workflow-none-approval, workflow-local, +# workflow-ephemeral) round out the same matrix for parity with the +# automated e2e suite -- they aren't part of that illustration, just +# additional scenarios: +# oneshot-none one-step workflow, spawn: "none" +# oneshot-local one-step workflow, spawn: "local" +# oneshot-ephemeral one-step workflow, spawn: "ephemeral", k8s-diag skill + Landlock demo +# workflow-ephemeral-approval multi-step + approval, spawn: "ephemeral" +# workflow-none-approval multi-step + approval, spawn: "none" +# workflow-local single step, spawn: "local", no approval +# workflow-ephemeral single step, spawn: "ephemeral", no approval # -# Note: agent-local and workflow-local omit output_schema -- the +# Note: oneshot-local and workflow-local omit output_schema -- the # cloud-agents SubprocessExecutor behind spawn:local has no native # structured-output mode yet (jameswnl/lightspeed-cloud-agents#235), so # it can't reliably guarantee schema-conforming JSON the way spawn:none # and spawn:ephemeral can. # -# MCP server reachability: agent-none and agent-local run in-process on -# this machine, not inside the cluster, so they reach the mock pod-status -# MCP server (deployed via ~/ws/local-infra's ocp-prod-mcp-pod-status-* -# targets) at localhost:8084 -- port-forward it first: +# MCP server reachability: oneshot-none and oneshot-local run in-process +# on this machine, not inside the cluster, so they reach the mock +# pod-status MCP server (deployed via ~/ws/local-infra's +# ocp-prod-mcp-pod-status-* targets) at localhost:8084 -- port-forward +# it first: # oc -n openshell-prod port-forward svc/mcp-pod-status-mock 8084:8084 -# agent-ephemeral runs inside an OpenShell sandbox pod on the cluster, so -# it reaches the same service via in-cluster DNS instead +# oneshot-ephemeral runs inside an OpenShell sandbox pod on the cluster, +# so it reaches the same service via in-cluster DNS instead # (mcp-pod-status-mock:8084), no port-forward needed. # # Usage: -# BASE_URL=http://localhost:8090 ./docs/cloud-agents-demo-curl.sh agent-none -# BASE_URL=http://localhost:8090 ./docs/cloud-agents-demo-curl.sh agent-local -# BASE_URL=http://localhost:8090 ./docs/cloud-agents-demo-curl.sh agent-ephemeral +# BASE_URL=http://localhost:8090 ./docs/cloud-agents-demo-curl.sh oneshot-none +# BASE_URL=http://localhost:8090 ./docs/cloud-agents-demo-curl.sh oneshot-local +# BASE_URL=http://localhost:8090 ./docs/cloud-agents-demo-curl.sh oneshot-ephemeral # BASE_URL=http://localhost:8090 ./docs/cloud-agents-demo-curl.sh workflow-ephemeral-approval # BASE_URL=http://localhost:8090 ./docs/cloud-agents-demo-curl.sh workflow-none-approval # BASE_URL=http://localhost:8090 ./docs/cloud-agents-demo-curl.sh workflow-local @@ -63,33 +67,44 @@ discover() { curl -s "${AUTH_HEADER[@]+"${AUTH_HEADER[@]}"}" "$BASE_URL/v1/mcp-servers" | jq } -run_agent() { - # POST an /v1/agents/run payload, pretty-print the body, and fail on - # HTTP >= 400. $1: title line, $2: spawn mode, $3: prompt, $4: extra - # JSON object merged onto the shared {prompt, spawn, provider, model} - # base -- callers only spell out what differs per spawn mode. - local title="$1" spawn="$2" prompt="$3" extra="{}" +oneshot_workflow_payload() { + # One-step workflow definition shared by the oneshot-* demos: a single + # agent step using the documented agent/result naming convention, plus + # a run-level provider. $1: spawn mode, $2: workflow name, $3: prompt, + # $4: extra JSON object merged onto the step -- callers only spell out + # what differs per spawn mode. + local extra="{}" if [[ $# -ge 4 ]]; then extra="$4" fi - echo "$title" - local payload resp status - payload=$(jq -n \ - --arg prompt "$prompt" \ - --arg spawn "$spawn" \ + jq -n \ + --arg spawn "$1" \ + --arg name "$2" \ + --arg prompt "$3" \ --argjson extra "$extra" \ - '{prompt: $prompt, spawn: $spawn, provider: "openai", model: "gpt-5-mini"} + $extra') - resp=$(curl -s -w '\n%{http_code}' -X POST "$BASE_URL/v1/agents/run" \ - "${AUTH_HEADER[@]+"${AUTH_HEADER[@]}"}" \ - -H "Content-Type: application/json" \ - -d "$payload") - status="${resp##*$'\n'}" - echo "${resp%$'\n'*}" | jq - [[ "$status" -lt 400 ]] + '{ + "definition": { + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": $name}, + "spec": { + "steps": [ + { + "name": "agent", "type": "agent", "spawn": $spawn, + "output_key": "result", + "prompt": $prompt, + "timeout_seconds": 120 + } + $extra + ] + } + }, + "provider": {"name": "openai", "model": "gpt-5-mini"} + }' } -agent_none() { - run_agent "== Agent — In-Process (spawn: none) ==" "none" \ +oneshot_none() { + echo "== One-shot one-step workflow — In-Process (spawn: none) ==" + submit_workflow "$(oneshot_workflow_payload "none" "oneshot-none-demo" \ "Is pod checkout-7f9 healthy?" \ '{ "tools": [], @@ -99,24 +114,29 @@ agent_none() { "properties": { "healthy": {"type": "boolean"}, "reason": {"type": "string"} }, "required": ["healthy", "reason"] } - }' + }')" + finish_workflow "$WF_ID" 60 } -agent_local() { - run_agent "== Agent — Subprocess (spawn: local) ==" "local" \ +oneshot_local() { + echo "== One-shot one-step workflow — Subprocess (spawn: local) ==" + submit_workflow "$(oneshot_workflow_payload "local" "oneshot-local-demo" \ "Check whether pod checkout-7f9 is healthy and say one sentence confirming the result." \ - '{"tools": [], "mcp_servers": [{"name": "kubectl-mcp", "url": "http://localhost:8084/mcp"}]}' + '{"tools": [], "mcp_servers": [{"name": "kubectl-mcp", "url": "http://localhost:8084/mcp"}]}')" + finish_workflow "$WF_ID" 150 } -agent_ephemeral() { +oneshot_ephemeral() { # allowed_skills=["k8s-diag"]: the spawner materializes just that skill # into the sandbox and Landlock-grants /skills/k8s-diag, so the prompt # below demonstrates both sides -- the allowed skill works, and reading # an unlisted skill (/skills/security-audit) is denied at the filesystem # boundary. Requires those skills baked into the sandbox image (/skills). - run_agent "== Agent — OpenShell (spawn: ephemeral, k8s-diag skill + Landlock) ==" "ephemeral" \ + echo "== One-shot one-step workflow — OpenShell (spawn: ephemeral, k8s-diag skill + Landlock) ==" + submit_workflow "$(oneshot_workflow_payload "ephemeral" "oneshot-ephemeral-demo" \ "Use the k8s-diag skill to check whether pod checkout-7f9 is healthy. Then try reading /skills/security-audit/SKILL.md and report whether that read succeeded or was denied, and why." \ - '{"mcp_servers": [{"name": "kubectl-mcp", "url": "http://mcp-pod-status-mock:8084/mcp"}], "provider": "openai", "model": "gpt-5-mini"}' + '{"mcp_servers": [{"name": "kubectl-mcp", "url": "http://mcp-pod-status-mock:8084/mcp"}], "allowed_skills": ["k8s-diag"]}')" + finish_workflow "$WF_ID" 150 } wait_for_status() { @@ -309,15 +329,15 @@ workflow_ephemeral() { case "${1:-}" in discover) discover ;; - agent-none) agent_none ;; - agent-local) agent_local ;; - agent-ephemeral) agent_ephemeral ;; + oneshot-none) oneshot_none ;; + oneshot-local) oneshot_local ;; + oneshot-ephemeral) oneshot_ephemeral ;; workflow-ephemeral-approval) workflow_ephemeral_approval ;; workflow-none-approval) workflow_none_approval ;; workflow-local) workflow_local ;; workflow-ephemeral) workflow_ephemeral ;; *) - echo "Usage: $0 {discover|agent-none|agent-local|agent-ephemeral|workflow-ephemeral-approval|workflow-none-approval|workflow-local|workflow-ephemeral}" >&2 + echo "Usage: $0 {discover|oneshot-none|oneshot-local|oneshot-ephemeral|workflow-ephemeral-approval|workflow-none-approval|workflow-local|workflow-ephemeral}" >&2 exit 1 ;; esac diff --git a/docs/cloud-agents-integration.html b/docs/cloud-agents-integration.html index 361e057b1..a4723823c 100644 --- a/docs/cloud-agents-integration.html +++ b/docs/cloud-agents-integration.html @@ -72,8 +72,8 @@

Cloud Agents Integration — Lightspeed Stack

- - + +
@@ -197,7 +197,7 @@

Everything here reduces to one atomic unit — the agent run — declared on
-

POST /v1/agents/run, spawn: "none" — the agent runs in-process. Same agent loop (LLM + MCP tools + skills) used by every spawn mode.

+

POST /v1/workflows/run with a one-step workflow, spawn: "none" — the agent runs in-process. Same agent loop (LLM + MCP tools + skills) used by every spawn mode.

@@ -215,7 +215,7 @@

POST /v1/agents/run, spawn: "none" — the agent runs in-process. Same agent ClientAPI caller -APIPOST /v1/agents/runvalidates request,dispatches by spawn modespawn: "none" +APIPOST /v1/workflows/runone-step workflow,dispatches by spawn modespawn: "none" SkillsSKILL.md → system prompt @@ -253,7 +253,7 @@

POST /v1/agents/run, spawn: "none" — the agent runs in-process. Same agent
-

Same request, same agent loop, same output shape — POST /v1/agents/run, spawn: "ephemeral". The difference is entirely underneath: a gateway-mediated, hardened sandbox instead of the API process.

+

Same one-step workflow, same agent loop, same output — POST /v1/workflows/run, spawn: "ephemeral". The difference is entirely underneath: a gateway-mediated, hardened sandbox instead of the API process.

@@ -277,7 +277,7 @@

Same request, same agent loop, same output shape — POST /v1/agents/run, sp ClientAPI caller -APIPOST /v1/agents/rundispatches by spawn modesame request shape as Tab 1spawn: "ephemeral" +APIPOST /v1/workflows/rundispatches by spawn modesame one-step shape as Tab 1spawn: "ephemeral" GatewaygRPC: CreateSandboxissues per-sandbox JWT @@ -318,7 +318,7 @@

Same request, same agent loop, same output shape — POST /v1/agents/run, sp
-

POST /v1/workflows/run — a workflow is a graph of steps, and each step is exactly the single-agent-invocation from Tab 2. The workflow runner owns sequencing, context-passing, and approval gates; it does not re-implement the agent loop.

+

POST /v1/workflows/run — a workflow is a graph of steps, and each step is exactly the one-step invocation from Tab 2. The workflow runner owns sequencing, context-passing, and approval gates; it does not re-implement the agent loop.

@@ -454,9 +454,9 @@

POST /v1/workflows/run — a workflow is a graph of steps, and each step is const N = createPlayer('n-', 'tab-none', ['n-skills'], [ { t: 'Client sends the request — params visible in full', on: ['n-f1'], - log: '→ POST /v1/agents/run\n{\n "prompt": "Is pod checkout-7f9 healthy?",\n "spawn": "none",\n "tools": ["get_pod_status"],\n "mcp_servers": ["kubectl-mcp"],\n "skills_paths": ["/skills/k8s-triage"],\n "output_schema": {"healthy": "bool", "reason": "string"}\n}' }, + log: '→ POST /v1/workflows/run\n{\n "definition": {"spec": {"steps": [\n {"name": "agent", "spawn": "none", "output_key": "result",\n "tools": ["get_pod_status"], "mcp_servers": ["kubectl-mcp"]}\n ]}},\n "provider": {"name": "openai", "model": "gpt-4o-mini"}\n}' }, { t: 'lightspeed-stack dispatches by spawn mode', on: ['n-f2'], - log: '→ spawn: "none" → get_step_executor() returns DirectExecutor' }, + log: '→ one-step workflow → runner normalizes → DirectExecutor' }, { t: 'DirectExecutor builds the agent: loads skills + tools + MCP', on: ['n-f3'], show: ['n-skills'], log: '⚙ Agent("openai:gpt-4o-mini")\n system_prompt += SKILL.md content (k8s-triage)\n tools: [get_pod_status] + MCPToolset(kubectl-mcp)' }, { t: 'Agent loop, turn 1 — inference call', on: ['n-f4'], @@ -467,15 +467,15 @@

POST /v1/workflows/run — a workflow is a graph of steps, and each step is log: '→ MCP kubectl-mcp: get_pod_status(...)\n← result: {status:"CrashLoopBackOff", restarts:14}' }, { t: 'Agent loop, turn 2 — inference with tool result', on: ['n-f8'], log: '→ LLM inference (with tool result)\n← structured output: {healthy:false, reason:"CrashLoopBackOff", restarts:14}' }, - { t: 'Result returns to the client', on: ['n-f9', 'n-f10'], - log: '← StepResult {status:"completed", input_tokens:412, output_tokens:38}\n← 200 OK to client' }, + { t: 'Run completes — client reads status and transcript', on: ['n-f9', 'n-f10'], + log: '← 202 Accepted {workflow_id}\n← GET status/transcripts {status:"completed", tokens:{in:412, out:38}}' }, ]); const E = createPlayer('e-', 'tab-eph', ['e-sandbound', 'e-sandlbl', 'e-skills', 'e-agent'], [ { t: 'Client sends the same shape of request — only spawn differs', on: ['e-f1'], - log: '→ POST /v1/agents/run\n{\n "prompt": "Is pod checkout-7f9 healthy?",\n "spawn": "ephemeral",\n "tools": ["get_pod_status"],\n "mcp_servers": ["kubectl-mcp"],\n "skills_image": "registry/skills:v1",\n "sandbox_image": "lightspeed-agentic-sandbox:latest"\n}' }, + log: '→ POST /v1/workflows/run\n{\n "definition": {"spec": {"steps": [\n {"name": "agent", "spawn": "ephemeral", "output_key": "result",\n "allowed_skills": ["k8s-diag"]}\n ]}},\n "provider": {"name": "openai", "model": "gpt-4o-mini"},\n "sandbox_image": "lightspeed-agentic-sandbox:latest"\n}' }, { t: 'lightspeed-stack dispatches by spawn mode', on: ['e-f2'], - log: '→ spawn: "ephemeral" → get_step_executor() returns SandboxExecutor' }, + log: '→ one-step workflow → runner normalizes → SandboxExecutor' }, { t: 'SandboxExecutor builds a spawner from config', on: ['e-f3'], log: '⚙ spawner_factory.build_spawner("openshell", ...) → OpenShellSpawner (gRPC client)' }, { t: 'Spawner asks the gateway to create a sandbox', show: ['e-sandbound', 'e-sandlbl'], @@ -491,7 +491,7 @@

POST /v1/workflows/run — a workflow is a graph of steps, and each step is { t: 'Agent loop, turn 2 — final structured answer', on: ['e-f6', 'e-f7'], log: '← structured output: {healthy:false, reason:"CrashLoopBackOff", restarts:14}' }, { t: 'Transcript pulled, sandbox destroyed, result returned', on: ['e-f5', 'e-f4', 'e-f3', 'e-f2', 'e-f1'], hide: ['e-sandbound', 'e-sandlbl', 'e-skills', 'e-agent'], - log: '→ GET /v1/agent/events (transcript)\n↓ DeleteSandbox — container destroyed\n← 200 OK to client — identical response shape to Tab 1' }, + log: '→ GET /v1/agent/events (transcript)\n↓ DeleteSandbox — container destroyed\n← run completes — same status/transcript APIs as Tab 1' }, ]); const W = createPlayer('w-', 'tab-wf', ['w-step', 'w-sandbound', 'w-sandlbl'], [ diff --git a/docs/design/cloud-agents/e2e-testing-strategy.md b/docs/design/cloud-agents/e2e-testing-strategy.md index ea026d647..50b0e2879 100644 --- a/docs/design/cloud-agents/e2e-testing-strategy.md +++ b/docs/design/cloud-agents/e2e-testing-strategy.md @@ -4,11 +4,13 @@ An audit of `tests/e2e/cloud_agents/` found real coverage gaps across the two cloud-agents HTTP endpoints, `/v1/agents/run` and `/v1/workflows/*`, -crossed with the three spawn modes: +crossed with the three spawn modes. `/v1/agents/run` was later removed +(jameswnl/lightspeed-stack#55) and its coverage migrated to one-step +workflow tests -- the table below is the historical audit record: | Endpoint | none | local | ephemeral | |---|---|---|---| -| `/v1/agents/run` | HTTP-tested | not a valid API value (`AgentRunRequest.spawn: Literal["none","ephemeral"]`) | only handler-direct, bypassing HTTP routing/auth | +| `/v1/agents/run` (removed in #55) | HTTP-tested | not a valid API value (`AgentRunRequest.spawn: Literal["none","ephemeral"]`) | only handler-direct, bypassing HTTP routing/auth | | `/v1/workflows/run` | HTTP-tested | zero coverage | zero coverage | None of `tests/e2e/cloud_agents/` ran in CI (`.github/workflows/cloud_agents_tests.yaml` @@ -17,7 +19,9 @@ a real `OPENAI_API_KEY` against real OpenAI. **Update:** `/v1/agents/run`'s `local` gap above was closed after this audit -- `AgentRunRequest.spawn` now accepts `Literal["none","local","ephemeral"]`, -HTTP-tested in `test_agents_run_http_e2e.py::test_local_spawn_agent_run`. Note +HTTP-tested in `test_agents_run_http_e2e.py::test_local_spawn_agent_run` +(that file was removed in #55; the equivalent coverage is now +`test_workflows_http_e2e.py::test_workflow_with_local_spawn_step`). Note that test omits `output_schema`: the cloud-agents `SubprocessExecutor` behind `spawn:local` has no native structured-output mode yet (jameswnl/lightspeed-cloud-agents#235), so it can't reliably guarantee @@ -45,7 +49,9 @@ missing functionality. `test_workflows_http_e2e.py` gained `test_workflow_with_ephemeral_spawn_step` accordingly, and `test_agents_run_http_e2e.py::TestAgentRunHttpE2E` gained `test_ephemeral_agent_run` to cover the last real gap (agents×ephemeral -had only handler-direct coverage before). +had only handler-direct coverage before). In #55 the agents HTTP file was +removed and one-shot coverage moved into `test_workflows_http_e2e.py` as +one-step workflow tests. (These two files started as one combined `test_agents_workflow_http_e2e.py` and were later split by endpoint as part of a broader @@ -88,7 +94,7 @@ as before; set (what CI does), `OPENAI_BASE_URL` is redirected to the mock for the session. **Scope warning:** this only applies safely to -`test_agents_run_http_e2e.py`/`test_workflows_http_e2e.py` (plus the mock's own self-tests) — +`test_workflows_http_e2e.py` (plus the mock's own self-tests) — every other file in `tests/e2e/cloud_agents/` asserts real-world semantic LLM content (e.g. `"paris" in output`) that only a real model can produce, and fails confusingly (not due to a real bug) if run with the mock @@ -116,7 +122,7 @@ Registered in `pyproject.toml` and applied to every ephemeral-gated test `lightspeed-stack-harness.yaml` had no `spawner:` section at all, which meant `spawn=ephemeral` would 400 against a server started with it — -including `docs/cloud-agents-demo-curl.sh`'s `agent-ephemeral` scenario, +including `docs/cloud-agents-demo-curl.sh`'s `oneshot-ephemeral` scenario, which already documented an ephemeral flow that couldn't have worked. Added: ```yaml @@ -137,7 +143,7 @@ automatically on every push/PR to `harness` (same trigger as the existing unit/integration job): a `postgres:16` service matching the harness config's credentials, `OPENAI_API_KEY=sk-mock-ci-key` + `LIGHTSPEED_E2E_USE_MOCK_LLM=1`, running -`test_agents_run_http_e2e.py` and `test_workflows_http_e2e.py` plus the mock's own self-tests with +`test_workflows_http_e2e.py` plus the mock's own self-tests with `-m "not ephemeral"`. ## File organization by testing layer @@ -152,9 +158,7 @@ name in two different files). Current layout: | File | Layer | Covers | |---|---|---| -| `test_agents_run_http_e2e.py` | real HTTP (`TestClient`) | `/v1/agents/run`, spawn none+local+ephemeral | -| `test_workflows_http_e2e.py` | real HTTP (`TestClient`) | `/v1/workflows/*`, spawn none+local+ephemeral | -| `test_agents_run_handler_e2e.py` | handler-direct (`handler.__wrapped__(...)`) | `/v1/agents/run`, spawn none+local+ephemeral | +| `test_workflows_http_e2e.py` | real HTTP (`TestClient`) | `/v1/workflows/*`, spawn none+local+ephemeral, incl. one-step workflows | | `test_query_direct_handler_e2e.py` | handler-direct | `/v1/query/direct` error paths | | `test_step_executor_e2e.py` | step-executor dispatch (`get_step_executor(...).run(...)`) | single-step execution, spawn none+local+ephemeral | | `test_workflow_definitions_e2e.py` | step-executor dispatch | full workflow-YAML execution, one step-executor call per step | diff --git a/docs/design/cloud-agents/integration-architecture.md b/docs/design/cloud-agents/integration-architecture.md index e6b36c35c..110165986 100644 --- a/docs/design/cloud-agents/integration-architecture.md +++ b/docs/design/cloud-agents/integration-architecture.md @@ -16,7 +16,6 @@ graph TB subgraph new["New Endpoints"] qd["/query/direct"] qds["/query/direct/stream"] - ar["/agents/run"] wf["/workflows/*"] at["/agent-tools"] end @@ -62,8 +61,8 @@ sequenceDiagram participant LLM as OpenAI API participant MCP as MCP Server - Client->>FastAPI: POST /agents/run - FastAPI->>DirectExecutor: run(StepInput) + Client->>FastAPI: POST /v1/workflows/run (one step) + FastAPI->>DirectExecutor: runner normalizes, run(step) DirectExecutor->>Agent: Agent("openai:gpt-4o-mini") Agent->>LLM: inference call LLM-->>Agent: response @@ -75,7 +74,7 @@ sequenceDiagram end Agent-->>DirectExecutor: AgentRunResult DirectExecutor-->>FastAPI: StepResult - FastAPI-->>Client: JSON response + FastAPI-->>Client: 202 Accepted, then status/transcripts ``` **When to use:** Default for all queries and workflow steps. Low latency, no infrastructure needed. @@ -95,8 +94,8 @@ sequenceDiagram participant Agent as pydantic-ai Agent participant LLM as OpenAI API - Client->>FastAPI: POST /agents/run - FastAPI->>SubprocessExec: run(StepInput) + Client->>FastAPI: POST /v1/workflows/run (one step) + FastAPI->>SubprocessExec: runner normalizes, run(step) SubprocessExec->>Child: spawn subprocess Note over Child: Inherits env vars
(OPENAI_API_KEY, etc.) Child->>Agent: create Agent with tools @@ -106,7 +105,7 @@ sequenceDiagram Child-->>SubprocessExec: JSON via stdout Note over Child: Process exits SubprocessExec-->>FastAPI: StepResult - FastAPI-->>Client: JSON response + FastAPI-->>Client: 202 Accepted, then status/transcripts ``` **When to use:** Steps with untrusted tools, crash isolation needed, or resource-intensive operations. @@ -127,8 +126,8 @@ sequenceDiagram participant Container as Sandbox Container participant LLM as OpenAI API - Client->>FastAPI: POST /agents/run - FastAPI->>SandboxExec: run(StepInput) + Client->>FastAPI: POST /v1/workflows/run (one step) + FastAPI->>SandboxExec: runner normalizes, run(step) SandboxExec->>Spawner: spawn(image, env, labels) Spawner->>Gateway: CreateSandbox Note over Gateway: Gateway's own compute driver
(kubernetes or podman) decides
how the sandbox is created --
not something the client sends @@ -151,7 +150,7 @@ sequenceDiagram SandboxExec->>Spawner: destroy(sandbox_name) Spawner->>Gateway: DeleteSandbox SandboxExec-->>FastAPI: StepResult - FastAPI-->>Client: JSON response + FastAPI-->>Client: 202 Accepted, then status/transcripts ``` **When to use:** Full isolation needed, different agent SDK required, or agents with kubectl/filesystem access. @@ -267,8 +266,7 @@ graph TB |---|---|---| | `/v1/query/direct` | POST | Blocking query via DirectExecutor | | `/v1/query/direct/stream` | POST | SSE streaming query | -| `/v1/agents/run` | POST | Single agent execution (any spawn mode) | -| `/v1/workflows/run` | POST | Start a multi-step workflow | +| `/v1/workflows/run` | POST | Start a workflow (multi-step, or one-step for one-shot agent runs) | | `/v1/workflows/{id}` | GET | Get workflow status | | `/v1/workflows/{id}/approve` | POST | Approve a paused step | | `/v1/workflows/{id}/cancel` | POST | Cancel a workflow | @@ -478,7 +476,6 @@ Blue = cloud-agents. Orange = lightspeed-stack. Red = final migration steps. | Component | Multi-pod safe? | Notes | |---|---|---| | `/query/direct` | Yes | Stateless per call | -| `/agents/run` | Yes | Stateless per call | | Conversation state | Yes (PostgreSQL) | Shared database | | Workflow state | Yes (PostgreSQL) | Shared database | | Running workflow tasks | No (in-memory) | Use Temporal for crash recovery | diff --git a/docs/devel_doc/openapi.json b/docs/devel_doc/openapi.json index a1bb31844..813f0c7e0 100644 --- a/docs/devel_doc/openapi.json +++ b/docs/devel_doc/openapi.json @@ -11575,59 +11575,6 @@ } } }, - "/v1/agents/run": { - "post": { - "tags": [ - "agents" - ], - "summary": "Run Agent Handler", - "description": "Execute a single agent with inline parameters.\n\nParameters:\n request: FastAPI request (used by middleware).\n body: Agent execution parameters.\n auth: Authentication tuple (used by middleware).\n\nReturns:\n Agent execution result with status, output, and transcript.", - "operationId": "run_agent_handler_v1_agents_run_post", - "requestBody": { - "content": { - "application/json": { - "schema": { - "$ref": "#/components/schemas/AgentRunRequest" - } - } - }, - "required": true - }, - "responses": { - "200": { - "description": "Agent execution result.", - "content": { - "application/json": { - "schema": { - "additionalProperties": true, - "type": "object", - "title": "Response Run Agent Handler V1 Agents Run Post" - } - } - } - }, - "401": { - "description": "Unauthorized." - }, - "403": { - "description": "Forbidden." - }, - "500": { - "description": "Internal server error." - }, - "422": { - "description": "Validation Error", - "content": { - "application/json": { - "schema": { - "$ref": "#/components/schemas/HTTPValidationError" - } - } - } - } - } - } - }, "/v1/workflows/run": { "post": { "tags": [ @@ -12164,7 +12111,6 @@ "manage_prompts", "read_prompts", "manage_saved_prompts", - "agent_run", "workflow_start", "workflow_view", "workflow_approve", @@ -12570,145 +12516,6 @@ "title": "AgentProvider", "description": "Represents the service provider of an agent." }, - "AgentRunRequest": { - "properties": { - "prompt": { - "type": "string", - "title": "Prompt", - "description": "The task prompt for the agent." - }, - "instructions": { - "anyOf": [ - { - "type": "string" - }, - { - "type": "null" - } - ], - "title": "Instructions", - "description": "System prompt / instructions for the agent." - }, - "model": { - "anyOf": [ - { - "type": "string" - }, - { - "type": "null" - } - ], - "title": "Model", - "description": "Full model ID (e.g. 'openai/gpt-4o'). Falls back to inference.default_model." - }, - "provider": { - "anyOf": [ - { - "type": "string" - }, - { - "type": "null" - } - ], - "title": "Provider", - "description": "Provider name. Used to prefix model if no slash present. Falls back to inference.default_provider." - }, - "spawn": { - "type": "string", - "enum": [ - "none", - "local", - "ephemeral" - ], - "title": "Spawn", - "description": "'none' runs in-process; 'local' spawns a subprocess; 'ephemeral' spawns a container.", - "default": "none" - }, - "sandbox_image": { - "anyOf": [ - { - "type": "string" - }, - { - "type": "null" - } - ], - "title": "Sandbox Image", - "description": "Container image for spawn=ephemeral." - }, - "tools": { - "items": { - "type": "string" - }, - "type": "array", - "title": "Tools", - "description": "Registered tool names for the agent to use." - }, - "mcp_servers": { - "anyOf": [ - { - "items": { - "additionalProperties": true, - "type": "object" - }, - "type": "array" - }, - { - "type": "null" - } - ], - "title": "Mcp Servers", - "description": "MCP server configs: [{name, url, headers}]. Each server is connected via pydantic-ai MCPToolset." - }, - "allowed_skills": { - "anyOf": [ - { - "items": { - "type": "string" - }, - "type": "array" - }, - { - "type": "null" - } - ], - "title": "Allowed Skills", - "description": "Skill allowlist: names of skill subdirectories the agent may use. For spawn=ephemeral the spawner materializes just this subset into the sandbox and Landlock-grants each /skills/ path, so unlisted skills are denied at the filesystem boundary. Omitted means no skills." - }, - "output_schema": { - "anyOf": [ - { - "additionalProperties": true, - "type": "object" - }, - { - "type": "null" - } - ], - "title": "Output Schema", - "description": "JSON Schema for structured output." - }, - "context": { - "anyOf": [ - { - "additionalProperties": true, - "type": "object" - }, - { - "type": "null" - } - ], - "title": "Context", - "description": "Prior context from previous steps or calls." - } - }, - "type": "object", - "required": [ - "prompt" - ], - "title": "AgentRunRequest", - "description": "Request body for POST /v1/agents/run.\n\nAll agent parameters are passed inline \u2014 no registry lookup.\n\nAttributes:\n prompt: The task prompt for the agent.\n instructions: System prompt / instructions.\n model: Full model ID (e.g. \"openai/gpt-4o\").\n provider: Provider name (used with model if no slash in model).\n spawn: Execution mode.\n sandbox_image: Container image for spawn=ephemeral.\n tools: Tool definitions.\n mcp_servers: MCP server names.\n allowed_skills: Skill allowlist (names of skill subdirectories).\n output_schema: JSON Schema for structured output.\n context: Prior context (e.g. from previous steps)." - }, "AgentSkill": { "properties": { "description": { @@ -21707,7 +21514,7 @@ "additionalProperties": true, "type": "object", "title": "Definition", - "description": "Workflow definition with apiVersion, kind, metadata, spec." + "description": "Workflow definition with apiVersion, kind, metadata, spec. A one-step workflow uses a single agent step named 'agent' with output_key 'result'." }, "provider": { "anyOf": [ @@ -21765,7 +21572,7 @@ "definition" ], "title": "RunWorkflowRequest", - "description": "Request body for POST /v1/workflows/run.\n\nAttributes:\n definition: Workflow definition (same schema as cloud-agents YAML).\n provider: Default LLM provider config for all steps.\n sandbox_image: Default sandbox image for ephemeral steps.\n approval_policy: Optional approval policy.\n session_id: Optional caller-provided ID grouping related workflow runs." + "description": "Request body for POST /v1/workflows/run.\n\nA one-shot agent invocation is a one-step workflow: a definition with a\nsingle agent step using the ``agent`` / ``result`` naming convention.\nOne-step workflows take exactly the same path as multi-step workflows\n(normalization, executor selection, status, cancellation, transcripts,\npersistence).\n\nAttributes:\n definition: Workflow definition (same schema as cloud-agents YAML).\n provider: Default LLM provider config for all steps.\n sandbox_image: Default sandbox image for ephemeral steps.\n approval_policy: Optional approval policy.\n session_id: Optional caller-provided ID grouping related workflow runs." }, "SQLiteDatabaseConfiguration": { "properties": { @@ -24065,10 +23872,6 @@ "name": "a2a", "description": "Agent-to-Agent (A2A) protocol." }, - { - "name": "agents", - "description": "Agent execution." - }, { "name": "authorized", "description": "Authorization probe." diff --git a/src/app/endpoints/agents.py b/src/app/endpoints/agents.py deleted file mode 100644 index d9fe2e683..000000000 --- a/src/app/endpoints/agents.py +++ /dev/null @@ -1,132 +0,0 @@ -"""Handler for REST API call to execute an agent.""" - -# pylint: disable=import-outside-toplevel - -from typing import Annotated, Any - -from cloud_agents.workflow.executor.step.base import StepInput, StepMetadata -from cloud_agents.workflow.executor.step.dispatch import get_step_executor -from fastapi import APIRouter, Depends, HTTPException, Request, status - -from authentication import get_auth_dependency -from authentication.interface import AuthTuple -from authorization.middleware import authorize -from configuration import configuration -from log import get_logger -from models.api.requests.agents import AgentRunRequest -from models.config import Action -from utils.endpoints import check_configuration_loaded -from workflow.provider_credentials import credentials_secret_for -from workflow.spawner_factory import build_spawner - -logger = get_logger(__name__) -router = APIRouter(tags=["agents"]) - - -agent_run_responses: dict[int | str, dict[str, Any]] = { - 200: {"description": "Agent execution result."}, - 401: {"description": "Unauthorized."}, - 403: {"description": "Forbidden."}, - 500: {"description": "Internal server error."}, -} - - -@router.post( - "/agents/run", - responses=agent_run_responses, -) -@authorize(Action.AGENT_RUN) -async def run_agent_handler( - request: Request, - body: AgentRunRequest, - auth: Annotated[AuthTuple, Depends(get_auth_dependency())], -) -> dict[str, Any]: - """Execute a single agent with inline parameters. - - Parameters: - request: FastAPI request (used by middleware). - body: Agent execution parameters. - auth: Authentication tuple (used by middleware). - - Returns: - Agent execution result with status, output, and transcript. - """ - _ = request - - check_configuration_loaded(configuration) - - inference = configuration.inference - spawner = None - sandbox_image = None - provider: dict[str, Any] = { - "name": body.provider or inference.default_provider or "", - "model": body.model or inference.default_model or "", - } - if body.spawn == "ephemeral": - spawner_config = configuration.spawner_configuration - if not spawner_config: - raise HTTPException( - status_code=status.HTTP_400_BAD_REQUEST, - detail="spawn=ephemeral requires a spawner configuration.", - ) - spawner = build_spawner(spawner_config) - sandbox_image = body.sandbox_image or spawner_config.sandbox_image - cred_secret = credentials_secret_for(provider["name"]) - if cred_secret: - provider["credentials_secret"] = cred_secret - - executor = get_step_executor( - {"name": "agent-run", "spawn": body.spawn, "prompt": body.prompt}, - spawner=spawner, - ) - - user_id, _, _, _ = auth - - step_input_kwargs: dict[str, Any] = { - "prompt": body.prompt, - "provider": provider, - "system_prompt": body.instructions, - "output_schema": body.output_schema, - "tools": body.tools, - "mcp_servers": body.mcp_servers, - "allowed_skills": body.allowed_skills, - "context": body.context or {}, - "step_name": "agent-run", - "output_key": "result", - # Raw step definition so the ephemeral path (SandboxExecutor -> - # step_runner -> spawner.spawn) sees allowed_skills and can - # materialize just that subset with per-skill Landlock grants. - # Other executors ignore raw_step and read allowed_skills directly. - # mcp_servers names are selected here too: step_runner only injects - # catalog entries (StepInput.mcp_servers) whose names the step lists, - # so without this the sandbox never sees request-level MCP servers. - "raw_step": { - "name": "agent-run", - "prompt": body.prompt, - "output_key": "result", - "allowed_skills": body.allowed_skills, - "mcp_servers": [s["name"] for s in (body.mcp_servers or []) if "name" in s], - }, - "metadata": StepMetadata(user_id=user_id), - } - if sandbox_image: - step_input_kwargs["sandbox_image"] = sandbox_image - step_input = StepInput(**step_input_kwargs) - - result = await executor.run(step_input) - - return { - "status": result.status, - "output": result.output, - "error": result.error, - "transcript": ( - [e if isinstance(e, dict) else e.model_dump() for e in result.transcript] - if result.transcript - else [] - ), - "token_usage": { - "input_tokens": result.input_tokens, - "output_tokens": result.output_tokens, - }, - "duration_ms": result.duration_ms, - } diff --git a/src/app/main.py b/src/app/main.py index 3b51b313d..2e9383ba0 100644 --- a/src/app/main.py +++ b/src/app/main.py @@ -40,7 +40,6 @@ # Global OpenAPI tags so every operation tag is declared (Spectral: operation-tag-defined). _OPENAPI_TAGS: Final[list[dict[str, str]]] = [ {"name": "a2a", "description": "Agent-to-Agent (A2A) protocol."}, - {"name": "agents", "description": "Agent execution."}, {"name": "authorized", "description": "Authorization probe."}, {"name": "config", "description": "Service configuration."}, {"name": "conversations_v1", "description": "Conversations API v1."}, @@ -91,12 +90,10 @@ async def lifespan( # pylint: disable=too-many-branches,too-many-statements,imp # cloud-agents' step/workflow executors emit OTEL spans through a tracer # that only records once a global TracerProvider is set -- otherwise # every span is silently dropped (NoOp tracer) even when - # OTEL_EXPORTER_OTLP_ENDPOINT is set. workflow.executor_factory already - # does this for /v1/workflows/run, but /v1/agents/run (agents.py) uses - # the step executors directly and never triggers it, so its spans were - # never exported. init_tracing is idempotent and a no-op when - # OTEL_EXPORTER_OTLP_ENDPOINT is unset, so calling it once here covers - # every code path. + # OTEL_EXPORTER_OTLP_ENDPOINT is set. init_tracing is idempotent and a + # no-op when OTEL_EXPORTER_OTLP_ENDPOINT is unset, so calling it once + # here covers every code path (including /v1/workflows/run, whose + # executor factory also calls it defensively). init_tracing("workflow-runner") llama_stack_config = configuration.configuration.llama_stack diff --git a/src/app/routers.py b/src/app/routers.py index 97633ba5e..9e653012c 100644 --- a/src/app/routers.py +++ b/src/app/routers.py @@ -6,7 +6,6 @@ # A2A (Agent-to-Agent) protocol support a2a, agent_tools, - agents, authorized, config, conversations_v1, @@ -84,9 +83,9 @@ def include_routers(app: FastAPI) -> None: app.include_router(authorized.router) app.include_router(metrics.router) - # Agent execution and workflow orchestration + # Agent tools and workflow orchestration (POST /v1/workflows/run is + # the only agent execution endpoint; one-shot runs are one-step workflows) app.include_router(agent_tools.router, prefix="/v1") - app.include_router(agents.router, prefix="/v1") app.include_router(workflows.router, prefix="/v1") # A2A (Agent-to-Agent) protocol endpoint diff --git a/src/models/api/requests/__init__.py b/src/models/api/requests/__init__.py index f9dcdeab0..d26c8f2df 100644 --- a/src/models/api/requests/__init__.py +++ b/src/models/api/requests/__init__.py @@ -1,7 +1,6 @@ """Concrete REST API request models grouped by domain.""" from models.api.requests.agents import ( - AgentRunRequest, ApproveWorkflowRequest, RunWorkflowRequest, ) @@ -28,7 +27,6 @@ ) __all__ = [ - "AgentRunRequest", "ApproveWorkflowRequest", "ConversationUpdateRequest", "FeedbackRequest", diff --git a/src/models/api/requests/agents.py b/src/models/api/requests/agents.py index 1383634bc..cfda6a242 100644 --- a/src/models/api/requests/agents.py +++ b/src/models/api/requests/agents.py @@ -1,96 +1,25 @@ -"""Request models for agent and workflow endpoints.""" +"""Request models for workflow endpoints. + +POST /v1/workflows/run is the only agent execution endpoint: a one-shot +agent invocation is submitted as a one-step workflow definition (see the +``RunWorkflowRequest.definition`` schema and the one-step ``agent`` / +``result`` naming convention in cloud-agents). +""" from typing import Any, Literal, Optional from pydantic import BaseModel, Field, field_validator -class AgentRunRequest(BaseModel): - """Request body for POST /v1/agents/run. - - All agent parameters are passed inline — no registry lookup. - - Attributes: - prompt: The task prompt for the agent. - instructions: System prompt / instructions. - model: Full model ID (e.g. "openai/gpt-4o"). - provider: Provider name (used with model if no slash in model). - spawn: Execution mode. - sandbox_image: Container image for spawn=ephemeral. - tools: Tool definitions. - mcp_servers: MCP server names. - allowed_skills: Skill allowlist (names of skill subdirectories). - output_schema: JSON Schema for structured output. - context: Prior context (e.g. from previous steps). - """ - - prompt: str = Field( - ..., - description="The task prompt for the agent.", - ) - - instructions: Optional[str] = Field( - None, - description="System prompt / instructions for the agent.", - ) - - model: Optional[str] = Field( - None, - description="Full model ID (e.g. 'openai/gpt-4o'). " - "Falls back to inference.default_model.", - ) - - provider: Optional[str] = Field( - None, - description="Provider name. Used to prefix model if no slash present. " - "Falls back to inference.default_provider.", - ) - - spawn: Literal["none", "local", "ephemeral"] = Field( - "none", - description="'none' runs in-process; 'local' spawns a subprocess; " - "'ephemeral' spawns a container.", - ) - - sandbox_image: Optional[str] = Field( - None, - description="Container image for spawn=ephemeral.", - ) - - tools: list[str] = Field( - default_factory=list, - description="Registered tool names for the agent to use.", - ) - - mcp_servers: Optional[list[dict[str, Any]]] = Field( - None, - description="MCP server configs: [{name, url, headers}]. " - "Each server is connected via pydantic-ai MCPToolset.", - ) - - allowed_skills: Optional[list[str]] = Field( - None, - description="Skill allowlist: names of skill subdirectories the agent " - "may use. For spawn=ephemeral the spawner materializes just this " - "subset into the sandbox and Landlock-grants each /skills/ " - "path, so unlisted skills are denied at the filesystem boundary. " - "Omitted means no skills.", - ) - - output_schema: Optional[dict[str, Any]] = Field( - None, - description="JSON Schema for structured output.", - ) - - context: Optional[dict[str, Any]] = Field( - None, - description="Prior context from previous steps or calls.", - ) - - class RunWorkflowRequest(BaseModel): """Request body for POST /v1/workflows/run. + A one-shot agent invocation is a one-step workflow: a definition with a + single agent step using the ``agent`` / ``result`` naming convention. + One-step workflows take exactly the same path as multi-step workflows + (normalization, executor selection, status, cancellation, transcripts, + persistence). + Attributes: definition: Workflow definition (same schema as cloud-agents YAML). provider: Default LLM provider config for all steps. @@ -101,7 +30,9 @@ class RunWorkflowRequest(BaseModel): definition: dict[str, Any] = Field( ..., - description="Workflow definition with apiVersion, kind, metadata, spec.", + description="Workflow definition with apiVersion, kind, metadata, spec. " + "A one-step workflow uses a single agent step named 'agent' with " + "output_key 'result'.", ) provider: Optional[dict[str, Any]] = Field( diff --git a/src/models/config.py b/src/models/config.py index 1f1756cf1..a6eb394c8 100644 --- a/src/models/config.py +++ b/src/models/config.py @@ -1316,9 +1316,6 @@ class Action(str, Enum): # User saved prompts (/v1/saved-prompts) MANAGE_SAVED_PROMPTS = "manage_saved_prompts" - # Agent execution (/v1/agents) - AGENT_RUN = "agent_run" - # Workflow operations (/v1/workflows) WORKFLOW_START = "workflow_start" WORKFLOW_VIEW = "workflow_view" diff --git a/tests/e2e/cloud_agents/conftest.py b/tests/e2e/cloud_agents/conftest.py index e5efde093..9a8043776 100644 --- a/tests/e2e/cloud_agents/conftest.py +++ b/tests/e2e/cloud_agents/conftest.py @@ -3,9 +3,9 @@ See mock_llm_env.py for the LIGHTSPEED_E2E_USE_MOCK_LLM gating logic the autouse fixture below wraps. -LIGHTSPEED_E2E_USE_MOCK_LLM is only safe for test_agents_run_http_e2e.py / -test_workflows_http_e2e.py (plus this module's own self-tests) -- those -files' assertions are structural only (status/key-presence), by design, so +LIGHTSPEED_E2E_USE_MOCK_LLM is only safe for +test_workflows_http_e2e.py (plus this module's own self-tests) -- that +file's assertions are structural only (status/key-presence), by design, so a canned response satisfies them. Every other file in this directory asserts real-world semantic content (e.g. "paris" in output) that only a real LLM can produce -- running those with the mock active fails diff --git a/tests/e2e/cloud_agents/test_agents_run_handler_e2e.py b/tests/e2e/cloud_agents/test_agents_run_handler_e2e.py deleted file mode 100644 index 97df3cf00..000000000 --- a/tests/e2e/cloud_agents/test_agents_run_handler_e2e.py +++ /dev/null @@ -1,311 +0,0 @@ -"""Handler-direct e2e tests for POST /v1/agents/run. - -Calls `run_agent_handler.__wrapped__(...)` directly, bypassing the -`@authorize` decorator and FastAPI routing/request-validation entirely -- -see test_agents_run_http_e2e.py for the real-HTTP equivalent of this -file's spawn=none/local/ephemeral coverage, and test_step_executor_e2e.py -for coverage one layer further down (the step-executor dispatch itself, -bypassing this handler too). - -Runs against a real LLM backend (OpenAI via pydantic-ai). Requires -OPENAI_API_KEY. spawn=ephemeral tests additionally require a reachable -OpenShell gateway (OPENSHELL_GATEWAY_URL, default localhost:17670) and -`cloud_agents.spawner.factory` installed, and are marked `ephemeral` so CI -can deselect them with `-m "not ephemeral"`. - -Usage: - uv run pytest tests/e2e/cloud_agents/test_agents_run_handler_e2e.py -v -s -""" - -# pylint: disable=import-outside-toplevel,too-few-public-methods,unused-argument - -from __future__ import annotations - -import importlib.util -import os -from typing import Any - -import pytest -from pydantic import ValidationError - -from app.endpoints.agents import run_agent_handler -from configuration import configuration -from models.api.requests.agents import AgentRunRequest - -from .conftest import AUTH, make_request, skip_if_gateway_unreachable - -pytestmark = pytest.mark.skipif( - not os.environ.get("OPENAI_API_KEY"), - reason="OPENAI_API_KEY not set", -) - -_OPENSHELL_GATEWAY_URL = os.environ.get("OPENSHELL_GATEWAY_URL", "localhost:17670") -_SANDBOX_IMAGE = os.environ.get( - "LIGHTSPEED_SANDBOX_IMAGE", "quay.io/jameswong/lightspeed-agentic-sandbox:latest" -) - -_CONFIG_NONE = { - "name": "e2e-agents-run-handler-test", - "service": { - "host": "localhost", - "port": 8080, - "auth_enabled": False, - "workers": 1, - }, - "llama_stack": { - "use_as_library_client": False, - "url": "http://localhost:8321", - }, - "user_data_collection": {"feedback_enabled": False}, - "authentication": {"module": "noop"}, -} - -_CONFIG_EPHEMERAL = { - **_CONFIG_NONE, - "name": "e2e-agents-run-handler-ephemeral-test", - "spawner": { - "type": "openshell", - "openshell_gateway_url": _OPENSHELL_GATEWAY_URL, - "sandbox_image": _SANDBOX_IMAGE, - }, -} - - -@pytest.fixture(name="e2e_config", scope="module") -def e2e_config_fixture() -> Any: - """Load config for handler-direct tests with no spawner section. - - No Llama Stack needed -- pydantic-ai talks to OpenAI directly. - """ - configuration.init_from_dict(_CONFIG_NONE) - return configuration - - -@pytest.fixture(name="e2e_config_ephemeral") -def e2e_config_ephemeral_fixture() -> Any: - """Load config with a real openshell spawner section. - - Skips (rather than erroring) if the gateway isn't reachable, same as - test_step_executor_e2e.py's openshell_spawner fixture. - - Note: like `e2e_config`, this mutates the global `configuration` - singleton -- harmless in practice (spawn=none tests never read the - spawner section, and `ephemeral` is normally deselected in CI), but - worth knowing if test isolation here ever becomes a real problem. - """ - skip_if_gateway_unreachable() - configuration.init_from_dict(_CONFIG_EPHEMERAL) - return configuration - - -@pytest.fixture(autouse=True) -def reset_spawner_singleton() -> Any: - """Reset the module-level spawner singleton before and after each test.""" - from workflow.spawner_factory import ( - reset_spawner, - ) # pylint: disable=import-outside-toplevel - - reset_spawner() - yield - reset_spawner() - - -class TestAgentRunE2E: - """E2E tests for /v1/agents/run with real LLM calls, spawn=none.""" - - @pytest.mark.asyncio - async def test_simple_text_response(self, e2e_config: Any) -> None: - """Agent returns a text response to a simple prompt.""" - body = AgentRunRequest( - prompt="What is 2 + 2? Reply with just the number.", - provider="openai", - model="gpt-4o-mini", - ) - - result = await run_agent_handler.__wrapped__(make_request(), body, AUTH) - - assert result["status"] == "completed" - assert result["output"] is not None - assert "4" in str(result["output"]) - assert result["token_usage"]["input_tokens"] > 0 - assert result["token_usage"]["output_tokens"] > 0 - assert result["duration_ms"] > 0 - - @pytest.mark.asyncio - async def test_with_instructions(self, e2e_config: Any) -> None: - """Agent follows system instructions.""" - body = AgentRunRequest( - prompt="What is the capital of France?", - provider="openai", - model="gpt-4o-mini", - instructions="You are a geography expert. Answer in exactly one word.", - ) - - result = await run_agent_handler.__wrapped__(make_request(), body, AUTH) - - assert result["status"] == "completed" - summary = str(result["output"]).lower() - assert "paris" in summary - - @pytest.mark.asyncio - async def test_structured_output(self, e2e_config: Any) -> None: - """Agent returns structured JSON matching the output schema.""" - body = AgentRunRequest( - prompt=( - "An alert fired: 'High CPU on node worker-1: 98% for 10 minutes.' " - "Classify the severity and category. Reply with JSON only." - ), - provider="openai", - model="gpt-4o-mini", - instructions=( - "You classify infrastructure alerts. " - "Reply ONLY with a JSON object matching the schema, no markdown." - ), - output_schema={ - "type": "object", - "properties": { - "severity": { - "type": "string", - "enum": ["low", "medium", "high", "critical"], - }, - "category": { - "type": "string", - "enum": ["resource", "network", "storage", "security"], - }, - "summary": {"type": "string"}, - }, - "required": ["severity", "category", "summary"], - }, - ) - - result = await run_agent_handler.__wrapped__(make_request(), body, AUTH) - - assert result["status"] == "completed" - output = result["output"] - assert isinstance(output, dict) - assert len(output) >= 2 - assert any(k in output for k in ("severity", "category", "summary")) - - @pytest.mark.asyncio - async def test_provider_model_resolution(self, e2e_config: Any) -> None: - """Model ID with separate provider field resolves correctly.""" - body = AgentRunRequest( - prompt="Say 'hello' and nothing else.", - provider="openai", - model="gpt-4o-mini", - ) - - result = await run_agent_handler.__wrapped__(make_request(), body, AUTH) - - assert result["status"] == "completed" - assert "hello" in str(result["output"]).lower() - - @pytest.mark.asyncio - async def test_multi_step_context_passing(self, e2e_config: Any) -> None: - """Agent receives prior context and uses it.""" - body = AgentRunRequest( - prompt=( - "The previous analysis found severity=high and category=resource. " - "Based on that, recommend ONE action in a single sentence." - ), - provider="openai", - model="gpt-4o-mini", - instructions="You are an SRE. Be concise.", - context={ - "analysis": { - "status": "completed", - "output": { - "severity": "high", - "category": "resource", - "summary": "Node worker-1 at 98% CPU for 10 min", - }, - } - }, - ) - - result = await run_agent_handler.__wrapped__(make_request(), body, AUTH) - - assert result["status"] == "completed" - assert len(str(result["output"])) > 10 - - @pytest.mark.asyncio - async def test_different_model_produces_response(self, e2e_config: Any) -> None: - """A different model still produces a valid response.""" - body = AgentRunRequest( - prompt="What color is the sky? One word.", - provider="openai", - model="gpt-4o-mini", - instructions="Reply with exactly one word.", - ) - - result = await run_agent_handler.__wrapped__(make_request(), body, AUTH) - - assert result["status"] == "completed" - assert "blue" in str(result["output"]).lower() - - def test_agent_run_missing_prompt_returns_422(self) -> None: - """Missing required 'prompt' field raises a pydantic ValidationError. - - AgentRunRequest is the request model for POST /v1/agents/run -- - FastAPI turns this into a 422 at the real HTTP layer (not exercised - here, no handler/HTTP involved -- pure request-model validation). - """ - with pytest.raises(ValidationError): - AgentRunRequest(provider="openai", model="gpt-4o-mini") # type: ignore[call-arg] - - -class TestAgentRunLocalSpawnE2E: - """POST /v1/agents/run with spawn=local, through the real handler. - - No output_schema here: the cloud-agents SubprocessExecutor behind - spawn=local has no native structured-output mode yet - (jameswnl/lightspeed-cloud-agents#235). - """ - - @pytest.mark.asyncio - async def test_local_spawn_executes_in_subprocess(self, e2e_config: Any) -> None: - """spawn=local runs in a child process and returns real LLM output.""" - body = AgentRunRequest( - prompt="What is 9+9? Reply with just the number.", - provider="openai", - model="gpt-4o-mini", - spawn="local", - ) - - result = await run_agent_handler.__wrapped__(make_request(), body, AUTH) - - assert result["status"] == "completed" - assert "18" in str(result["output"]) - assert result["transcript"] - - -@pytest.mark.ephemeral -@pytest.mark.skipif( - importlib.util.find_spec("cloud_agents.spawner.factory") is None, - reason="cloud_agents.spawner.factory not installed -- " - "re-run: uv pip install -e ~/ws/lightspeed-cloud-agents[local,kubernetes,openshell]", -) -class TestAgentRunEphemeralE2E: - """POST /v1/agents/run with spawn=ephemeral, through the real handler. - - Proves the fix for the spawner=None bug: - https://github.com/jameswnl/lightspeed-stack/issues/23 - """ - - @pytest.mark.asyncio - async def test_ephemeral_spawn_executes_real_sandbox( - self, e2e_config_ephemeral: Any - ) -> None: - """spawn=ephemeral reaches a real OpenShell sandbox and returns output.""" - body = AgentRunRequest( - prompt="What is 9+9? Reply with just the number.", - provider="openai", - model="gpt-4o-mini", - spawn="ephemeral", - ) - - result = await run_agent_handler.__wrapped__(make_request(), body, AUTH) - - assert result["status"] == "completed" - assert "18" in str(result["output"]) - assert result["transcript"] diff --git a/tests/e2e/cloud_agents/test_agents_run_http_e2e.py b/tests/e2e/cloud_agents/test_agents_run_http_e2e.py deleted file mode 100644 index 4df944556..000000000 --- a/tests/e2e/cloud_agents/test_agents_run_http_e2e.py +++ /dev/null @@ -1,145 +0,0 @@ -"""Real HTTP e2e tests for POST /v1/agents/run. - -Companion automated coverage for docs/cloud-agents-demo-curl.sh's -agent-none/agent-local/agent-ephemeral scenarios -- exercises in-process, -subprocess, and ephemeral agent runs through the actual FastAPI app over -real HTTP (real routing, real auth dependency resolution, real -request/response validation), unlike test_agents_run_handler_e2e.py -(calls the handler function directly) or test_step_executor_e2e.py (calls -the step-executor dispatch directly, bypassing the handler and HTTP both). - -Requires OPENAI_API_KEY and a reachable PostgreSQL matching -lightspeed-stack-harness.yaml's database.postgres section -- the shared -`http_client` fixture (conftest.py) enters the app's ASGI lifespan, which -initializes WorkflowStorageFactory against Postgres unconditionally, even -though /v1/agents/run itself doesn't touch workflow storage. Skips cleanly -if either prerequisite is missing. spawn=ephemeral additionally requires a -reachable OpenShell gateway (OPENSHELL_GATEWAY_URL, default -localhost:17670) and is marked `ephemeral` so CI can deselect it with -`-m "not ephemeral"`. - -Set LIGHTSPEED_E2E_USE_MOCK_LLM=1 (see conftest.py) to run the spawn=none -and spawn=local tests here against an in-process mock LLM instead of real -OpenAI -- this is what CI does. spawn=ephemeral still needs a real key and -gateway either way. - -Usage: - cd ~/ws/local-infra && make up # provides Postgres on localhost:5432 - uv run pytest tests/e2e/cloud_agents/test_agents_run_http_e2e.py -v -""" - -from __future__ import annotations - -import os - -import pytest -from fastapi.testclient import TestClient - -from .conftest import postgres_reachable, skip_if_gateway_unreachable - -pytestmark = [ - pytest.mark.skipif( - not os.environ.get("OPENAI_API_KEY"), - reason="OPENAI_API_KEY not set", - ), - pytest.mark.skipif( - not postgres_reachable(), - reason=( - "PostgreSQL not reachable on localhost:5432 " - "(needed by the shared http_client fixture)" - ), - ), -] - - -class TestAgentRunHttpE2E: - """POST /v1/agents/run over real HTTP. - - Mirrors demo-curl agent-none/agent-local/agent-ephemeral. - """ - - def test_in_process_agent_run(self, http_client: TestClient) -> None: - """spawn:none agent run returns a structured, schema-conforming response.""" - response = http_client.post( - "/v1/agents/run", - json={ - "prompt": "Is pod checkout-7f9 healthy? Assume yes, everything is fine.", - "spawn": "none", - "provider": "openai", - "model": "gpt-4o-mini", - "tools": [], - "mcp_servers": None, - "output_schema": { - "type": "object", - "properties": { - "healthy": {"type": "boolean"}, - "reason": {"type": "string"}, - }, - "required": ["healthy", "reason"], - }, - }, - ) - - assert response.status_code == 200 - data = response.json() - assert data["status"] == "completed" - assert isinstance(data["output"], dict) - assert "healthy" in data["output"] - assert "reason" in data["output"] - assert data["token_usage"]["input_tokens"] > 0 - - def test_local_spawn_agent_run(self, http_client: TestClient) -> None: - """spawn:local agent run executes in a subprocess and returns a completed result. - - No output_schema here (unlike test_in_process_agent_run): the - cloud-agents SubprocessExecutor has no native structured-output - mode yet (jameswnl/lightspeed-cloud-agents#235) -- it only embeds - the schema as prompt text, which the deterministic mock LLM used - here doesn't recognize. Matches the existing spawn:local workflow - step test's approach (test_workflow_with_local_spawn_step), which - also avoids output_schema for the same reason. - """ - response = http_client.post( - "/v1/agents/run", - json={ - "prompt": "Say one sentence confirming pod checkout-7f9 is healthy.", - "spawn": "local", - "provider": "openai", - "model": "gpt-4o-mini", - "tools": [], - "mcp_servers": None, - }, - ) - - assert response.status_code == 200 - data = response.json() - assert data["status"] == "completed" - assert data["output"] is not None - assert data["token_usage"]["input_tokens"] > 0 - - @pytest.mark.ephemeral - def test_ephemeral_agent_run(self, http_client: TestClient) -> None: - """spawn:ephemeral agent run reaches a real OpenShell sandbox over HTTP. - - Companion to test_agents_run_handler_e2e.py's handler-direct - coverage -- this goes through real routing, auth dependency - resolution, and request/response validation instead of calling the - handler function directly. - """ - skip_if_gateway_unreachable() - - response = http_client.post( - "/v1/agents/run", - json={ - "prompt": "What is 9+9? Reply with just the number.", - "spawn": "ephemeral", - "provider": "openai", - "model": "gpt-4o-mini", - }, - ) - - assert response.status_code == 200 - data = response.json() - assert data["status"] == "completed" - assert data["output"] is not None - assert data["transcript"] diff --git a/tests/e2e/cloud_agents/test_step_executor_e2e.py b/tests/e2e/cloud_agents/test_step_executor_e2e.py index 043d92e46..a188e6c67 100644 --- a/tests/e2e/cloud_agents/test_step_executor_e2e.py +++ b/tests/e2e/cloud_agents/test_step_executor_e2e.py @@ -2,8 +2,8 @@ Calls `get_step_executor(...).run(...)`/`.run_stream(...)` (or, for ephemeral, `cloud_agents.workflow.core.step_runner.run_step(...)`) -directly -- bypassing the `/v1/agents/run` handler and FastAPI routing -entirely, one layer below test_agents_run_handler_e2e.py. See +directly -- bypassing the workflow engine, the `/v1/workflows/run` +handler, and FastAPI routing entirely. See test_workflow_definitions_e2e.py for the same layer applied to a full multi-step workflow YAML instead of a single step. diff --git a/tests/e2e/cloud_agents/test_workflows_http_e2e.py b/tests/e2e/cloud_agents/test_workflows_http_e2e.py index d415d5f45..3a644f2c2 100644 --- a/tests/e2e/cloud_agents/test_workflows_http_e2e.py +++ b/tests/e2e/cloud_agents/test_workflows_http_e2e.py @@ -1,9 +1,9 @@ """Real HTTP e2e tests for POST /v1/workflows/*. Companion automated coverage for docs/cloud-agents-demo-curl.sh's -workflow-ephemeral-approval / workflow-none-approval / workflow-local / -workflow-ephemeral scenarios -- exercises workflow run/approve/transcripts -across all three spawn +oneshot-none / workflow-ephemeral-approval / workflow-none-approval / +workflow-local / workflow-ephemeral scenarios -- exercises workflow +run/approve/transcripts across all three spawn modes, with and without a human-approval gate, through the actual FastAPI app over real HTTP (real routing, real auth dependency resolution, real request/response validation), unlike test_workflow_definitions_e2e.py @@ -30,6 +30,8 @@ uv run pytest tests/e2e/cloud_agents/test_workflows_http_e2e.py -v """ +# pylint: disable=too-few-public-methods + from __future__ import annotations import os @@ -374,3 +376,69 @@ def test_workflow_with_ephemeral_spawn_step(self, http_client: TestClient) -> No ) assert completed["status"] == "completed" assert "investigate_result" in completed["steps"] + + +class TestOneStepWorkflowHttpE2E: + """A one-step workflow (one-shot agent run) over real HTTP. + + POST /v1/workflows/run is the only agent execution endpoint: the + former standalone agent run is submitted as a single-step workflow + using the documented ``agent`` / ``result`` naming convention, with a + run-level provider. Uses the issue-#55 migration-example shape. + """ + + def test_one_step_workflow_completes(self, http_client: TestClient) -> None: + """One-step spawn:none workflow completes with schema-conforming output.""" + start_response = http_client.post( + "/v1/workflows/run", + json={ + "definition": { + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "one-shot-agent"}, + "spec": { + "steps": [ + { + "name": "agent", + "type": "agent", + "spawn": "none", + "output_key": "result", + "prompt": ( + "Is pod checkout-7f9 healthy? Assume yes, " + "everything is fine." + ), + "output_schema": { + "type": "object", + "properties": { + "healthy": {"type": "boolean"}, + "reason": {"type": "string"}, + }, + "required": ["healthy", "reason"], + }, + "timeout_seconds": 120, + }, + ] + }, + }, + "provider": {"name": "openai", "model": "gpt-4o-mini"}, + }, + ) + + assert start_response.status_code == 202 + workflow_id = start_response.json()["workflow_id"] + assert workflow_id + + completed = wait_for_status( + http_client, + workflow_id, + lambda body: bool(body["is_terminal"]), + timeout_s=150, + ) + assert completed["status"] == "completed" + assert "result" in completed["steps"] + + transcripts_response = http_client.get( + f"/v1/workflows/{workflow_id}/transcripts" + ) + assert transcripts_response.status_code == 200 + assert "result" in transcripts_response.json()["transcripts"] diff --git a/tests/integration/cloud_agents/test_agents_integration.py b/tests/integration/cloud_agents/test_agents_integration.py deleted file mode 100644 index 0e6c9af5f..000000000 --- a/tests/integration/cloud_agents/test_agents_integration.py +++ /dev/null @@ -1,281 +0,0 @@ -"""Integration tests for /v1/agents/run endpoint. - -Tests the full endpoint call chain with mocked step executor, -verifying request parsing, model resolution, and response formatting. -""" - -# pylint: disable=import-outside-toplevel,unused-argument,protected-access,too-few-public-methods,unspecified-encoding - -from __future__ import annotations - -from pathlib import Path -from typing import Any - -import pytest -import yaml -from fastapi import Request -from pytest_mock import MockerFixture - -from app.endpoints.agents import run_agent_handler -from authentication.interface import AuthTuple -from configuration import configuration -from models.api.requests.agents import AgentRunRequest - - -@pytest.fixture(name="agent_config") -def agent_config_fixture() -> Any: - """Load minimal config for agent tests.""" - config_dict = { - "name": "test-agents", - "service": { - "host": "localhost", - "port": 8080, - "auth_enabled": False, - "workers": 1, - }, - "llama_stack": { - "use_as_library_client": False, - "url": "http://localhost:8321", - }, - "user_data_collection": { - "feedback_enabled": False, - }, - "authentication": {"module": "noop"}, - } - configuration.init_from_dict(config_dict) - return configuration - - -def _mock_executor(mocker: MockerFixture, output: Any = None) -> Any: - """Mock the step executor returned by get_step_executor.""" - mock_result = mocker.MagicMock() - mock_result.status = "completed" - mock_result.output = output or {"summary": "test output"} - mock_result.error = None - mock_result.transcript = [ - {"ts": "2024-01-01T00:00:00", "type": "result", "data": {"text": "done"}} - ] - mock_result.input_tokens = 50 - mock_result.output_tokens = 25 - mock_result.duration_ms = 500 - - mock_exec = mocker.AsyncMock() - mock_exec.run.return_value = mock_result - mocker.patch( - "app.endpoints.agents.get_step_executor", - return_value=mock_exec, - ) - return mock_exec - - -class TestAgentRunIntegration: - """Integration tests for POST /v1/agents/run.""" - - @pytest.mark.asyncio - async def test_simple_agent_run( - self, - agent_config: Any, - mock_request_with_auth: Request, - mocker: MockerFixture, - ) -> None: - """Simple agent run returns completed status with output.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - _mock_executor(mocker) - - body = AgentRunRequest( - prompt="Classify this alert", - provider="openai", - model="gpt-4o-mini", - ) - auth: AuthTuple = ("user-1", "testuser", False, "") - - result = await run_agent_handler.__wrapped__(mock_request_with_auth, body, auth) - - assert result["status"] == "completed" - assert result["output"]["summary"] == "test output" - assert result["token_usage"]["input_tokens"] == 50 - assert result["duration_ms"] == 500 - - @pytest.mark.asyncio - async def test_structured_output_passthrough( - self, - agent_config: Any, - mock_request_with_auth: Request, - mocker: MockerFixture, - ) -> None: - """Structured output from executor is passed through.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - _mock_executor( - mocker, - output={"severity": "high", "category": "resource"}, - ) - - body = AgentRunRequest( - prompt="Classify this alert", - provider="openai", - model="gpt-4o-mini", - output_schema={"type": "object"}, - ) - auth: AuthTuple = ("user-1", "testuser", False, "") - - result = await run_agent_handler.__wrapped__(mock_request_with_auth, body, auth) - - assert result["status"] == "completed" - assert result["output"]["severity"] == "high" - - @pytest.mark.asyncio - async def test_transcript_included( - self, - agent_config: Any, - mock_request_with_auth: Request, - mocker: MockerFixture, - ) -> None: - """Transcript events are included in the response.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - _mock_executor(mocker) - - body = AgentRunRequest( - prompt="Check status", - provider="openai", - model="gpt-4o-mini", - ) - auth: AuthTuple = ("user-1", "testuser", False, "") - - result = await run_agent_handler.__wrapped__(mock_request_with_auth, body, auth) - - assert len(result["transcript"]) == 1 - assert result["transcript"][0]["type"] == "result" - - @pytest.mark.asyncio - async def test_step_input_construction( - self, - agent_config: Any, - mock_request_with_auth: Request, - mocker: MockerFixture, - ) -> None: - """StepInput is constructed with correct provider and prompt.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_exec = _mock_executor(mocker) - - body = AgentRunRequest( - prompt="Analyze this", - provider="openai", - model="gpt-4o", - instructions="Be concise", - ) - auth: AuthTuple = ("user-1", "testuser", False, "") - - await run_agent_handler.__wrapped__(mock_request_with_auth, body, auth) - - step_input = mock_exec.run.call_args[0][0] - assert step_input.prompt == "Analyze this" - assert step_input.provider == {"name": "openai", "model": "gpt-4o"} - assert step_input.system_prompt == "Be concise" - - @pytest.mark.asyncio - async def test_mcp_servers_passed_to_step_input( - self, - agent_config: Any, - mock_request_with_auth: Request, - mocker: MockerFixture, - ) -> None: - """MCP server configs are passed through to StepInput.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_exec = _mock_executor(mocker) - - mcp_configs = [ - {"name": "kubectl", "url": "http://mcp-kubectl:8080/sse"}, - { - "name": "github", - "url": "http://mcp-github:8080/sse", - "headers": {"Authorization": "Bearer token123"}, - }, - ] - - body = AgentRunRequest( - prompt="List pods", - provider="openai", - model="gpt-4o-mini", - mcp_servers=mcp_configs, - ) - auth: AuthTuple = ("user-1", "testuser", False, "") - - await run_agent_handler.__wrapped__(mock_request_with_auth, body, auth) - - step_input = mock_exec.run.call_args[0][0] - assert step_input.mcp_servers == mcp_configs - assert len(step_input.mcp_servers) == 2 - assert ( - step_input.mcp_servers[1]["headers"]["Authorization"] == "Bearer token123" - ) - - @pytest.mark.asyncio - async def test_tools_passed_to_step_input( - self, - agent_config: Any, - mock_request_with_auth: Request, - mocker: MockerFixture, - ) -> None: - """Registered tool names are passed through to StepInput.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_exec = _mock_executor(mocker) - - body = AgentRunRequest( - prompt="Check the cluster", - provider="openai", - model="gpt-4o-mini", - tools=["kubectl_get", "http_request"], - ) - auth: AuthTuple = ("user-1", "testuser", False, "") - - await run_agent_handler.__wrapped__(mock_request_with_auth, body, auth) - - step_input = mock_exec.run.call_args[0][0] - assert step_input.tools == ["kubectl_get", "http_request"] - - @pytest.mark.asyncio - async def test_ephemeral_without_spawner_returns_400( - self, - agent_config: Any, - mock_request_with_auth: Request, - mocker: MockerFixture, - ) -> None: - """Requesting spawn=ephemeral without spawner config returns 400.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - from fastapi import HTTPException - - body = AgentRunRequest( - prompt="Fix the issue", - provider="openai", - model="gpt-4o-mini", - spawn="ephemeral", - ) - auth: AuthTuple = ("user-1", "testuser", False, "") - - with pytest.raises(HTTPException) as exc_info: - await run_agent_handler.__wrapped__(mock_request_with_auth, body, auth) - assert exc_info.value.status_code == 400 - - -class TestWorkflowDefinitionExecution: - """Test that cloud-agents workflow definitions are compatible.""" - - def test_triage_classify_definition_parses(self) -> None: - """The triage-classify workflow YAML parses into a valid definition.""" - from cloud_agents.workflow.core.definition import WorkflowDefinition - - wf_path = ( - Path(__file__).resolve().parent.parent.parent.parent - / "lightspeed-cloud-agents" - / "examples" - / "workflow-definitions" - / "triage-classify-workflow.yaml" - ) - if not wf_path.exists(): - pytest.skip("lightspeed-cloud-agents repo not found") - with open(wf_path, encoding="utf-8") as f: - raw = yaml.safe_load(f) - - definition = WorkflowDefinition.model_validate(raw) - assert definition.metadata["name"] == "triage-classify-alerts" - assert len(definition.spec.steps) == 3 diff --git a/tests/integration/cloud_agents/test_workflows_integration.py b/tests/integration/cloud_agents/test_workflows_integration.py new file mode 100644 index 000000000..5aaa7a62d --- /dev/null +++ b/tests/integration/cloud_agents/test_workflows_integration.py @@ -0,0 +1,176 @@ +"""Integration tests for one-step workflows via POST /v1/workflows/run. + +A one-shot agent invocation is a one-step workflow definition submitted to +the normal workflow path. These tests pin the documented one-step contract +(``agent`` / ``result`` naming, full former-standalone input coverage, +identical normalization to multi-step steps) against the real installed +cloud-agents package -- no LLM calls, no mocks of cloud-agents itself. +""" + +# pylint: disable=import-outside-toplevel,too-few-public-methods,unspecified-encoding + +from __future__ import annotations + +from pathlib import Path +from typing import Any + +import pytest +import yaml +from cloud_agents.workflow.core.definition import WorkflowDefinition +from cloud_agents.workflow.core.execution import ( + build_step_input, + normalize_definition, +) +from cloud_agents.workflow.core.validation import validate_definition + + +def _one_step_definition(step: dict[str, Any]) -> dict[str, Any]: + """Wrap a single step dict in the documented one-step definition shape.""" + return { + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "one-shot-agent"}, + "spec": {"steps": [step]}, + } + + +class TestOneStepWorkflowContract: + """The documented one-step workflow shape validates and normalizes.""" + + def test_full_one_step_shape(self) -> None: + """Every former standalone input is expressible on a one-step workflow. + + Covers prompt, instructions, provider/model (run-level), tools, + MCP servers, skills, permissions, spawn mode, sandbox config, + output schema, context, and timeout. + """ + definition = _one_step_definition( + { + "name": "agent", + "type": "agent", + "prompt": "Inspect the cluster", + "output_key": "result", + "instructions": "Be concise.", + "spawn": "ephemeral", + "spawn_config": {"sandbox_image": "custom-sandbox:v2"}, + "tools": ["kubectl_get"], + "mcp_servers": ["cluster"], + "allowed_skills": ["kubernetes"], + "permissions": {"service_account": "agent-runner"}, + "output_schema": {"type": "object"}, + "context": {"cluster": "prod"}, + "timeout_seconds": 120, + } + ) + run_context = { + "provider": {"name": "openai", "model": "gpt-4o"}, + "sandbox_image": "sandbox:latest", + } + + assert validate_definition(definition) == [] + + agent_steps, _ = normalize_definition(definition) + assert len(agent_steps) == 1 + step_input = build_step_input(definition["spec"]["steps"][0], run_context) + assert step_input.prompt == "Inspect the cluster" + assert step_input.provider["name"] == "openai" + assert step_input.provider["model"] == "gpt-4o" + assert step_input.tools == ["kubectl_get"] + assert step_input.allowed_skills == ["kubernetes"] + assert step_input.output_schema == {"type": "object"} + assert step_input.timeout_seconds == 120 + + def test_bare_one_step_defaults(self) -> None: + """A bare single step defaults to name 'agent' / output_key 'result'.""" + definition = _one_step_definition({"prompt": "Inspect the cluster"}) + + assert validate_definition(definition) == [] + + agent_steps, _ = normalize_definition(definition) + assert len(agent_steps) == 1 + assert agent_steps[0].name == "agent" + assert agent_steps[0].output_key == "result" + + def test_one_step_matches_multi_step_normalization(self) -> None: + """A one-step workflow normalizes identically to the same multi-step step. + + This is the core #55 invariant: one-shot runs take exactly the same + path and semantics as a step in a multi-step workflow. + """ + step = { + "name": "agent", + "type": "agent", + "prompt": "Inspect the cluster", + "output_key": "result", + "spawn": "none", + "tools": ["kubectl_get"], + "timeout_seconds": 60, + } + one_step = _one_step_definition(dict(step)) + multi_step = { + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "two-step"}, + "spec": { + "steps": [ + dict(step), + { + "name": "verify", + "type": "agent", + "prompt": "Verify", + "output_key": "verification", + "spawn": "none", + }, + ] + }, + } + + assert validate_definition(one_step) == [] + assert validate_definition(multi_step) == [] + + single_normalized, _ = normalize_definition(one_step) + multi_normalized, _ = normalize_definition(multi_step) + assert single_normalized[0].model_dump() == multi_normalized[0].model_dump() + + def test_secret_value_rejected(self) -> None: + """Inline MCP secret values fail validation, never reaching run state.""" + definition = _one_step_definition( + { + "name": "agent", + "type": "agent", + "prompt": "Inspect the cluster", + "output_key": "result", + "mcp_servers": [ + { + "name": "cluster", + "url": "https://admin:s3cret@example.com/mcp", + } + ], + } + ) + + assert validate_definition(definition) != [] + with pytest.raises(ValueError, match="credentialed URL"): + normalize_definition(definition) + + +class TestWorkflowDefinitionExecution: + """Test that cloud-agents workflow definitions are compatible.""" + + def test_triage_classify_definition_parses(self) -> None: + """The triage-classify workflow YAML parses into a valid definition.""" + wf_path = ( + Path(__file__).resolve().parent.parent.parent.parent + / "lightspeed-cloud-agents" + / "examples" + / "workflow-definitions" + / "triage-classify-workflow.yaml" + ) + if not wf_path.exists(): + pytest.skip("lightspeed-cloud-agents repo not found") + with open(wf_path, encoding="utf-8") as f: + raw = yaml.safe_load(f) + + definition = WorkflowDefinition.model_validate(raw) + assert definition.metadata["name"] == "triage-classify-alerts" + assert len(definition.spec.steps) == 3 diff --git a/tests/unit/app/test_routers.py b/tests/unit/app/test_routers.py index 8e0bfc453..d8c8917f1 100644 --- a/tests/unit/app/test_routers.py +++ b/tests/unit/app/test_routers.py @@ -123,7 +123,7 @@ def test_include_routers() -> None: include_routers(app) # are all routers added? - assert len(app.routers) == 30 + assert len(app.routers) == 29 assert root.router in app.get_routers() assert info.router in app.get_routers() assert models.router in app.get_routers() @@ -166,7 +166,7 @@ def test_check_prefixes() -> None: include_routers(app) # are all routers added? - assert len(app.routers) == 30 + assert len(app.routers) == 29 assert app.get_router_prefix(root.router) == "" assert app.get_router_prefix(info.router) == "/v1" assert app.get_router_prefix(models.router) == "/v1" diff --git a/tests/unit/cloud_agents/test_agents_endpoint.py b/tests/unit/cloud_agents/test_agents_endpoint.py deleted file mode 100644 index b70911d07..000000000 --- a/tests/unit/cloud_agents/test_agents_endpoint.py +++ /dev/null @@ -1,563 +0,0 @@ -"""Unit tests for the /v1/agents/run endpoint.""" - -# pylint: disable=protected-access,import-outside-toplevel,unused-argument - -from __future__ import annotations - -from typing import Any - -import pytest -from pytest_mock import MockerFixture - -from app.endpoints.agents import run_agent_handler -from models.api.requests.agents import AgentRunRequest - - -@pytest.fixture(name="mock_config") -def mock_config_fixture(mocker: MockerFixture) -> Any: - """Mock the configuration singleton.""" - mock_cfg = mocker.patch("app.endpoints.agents.configuration") - mock_cfg.spawner_configuration = None - mock_cfg.inference.default_provider = None - mock_cfg.inference.default_model = None - return mock_cfg - - -@pytest.fixture(name="mock_executor") -def mock_executor_fixture(mocker: MockerFixture) -> Any: - """Mock the step executor via cloud-agents dispatch.""" - mock_result = mocker.MagicMock() - mock_result.status = "completed" - mock_result.output = {"summary": "Done"} - mock_result.error = None - mock_result.transcript = [] - mock_result.input_tokens = 50 - mock_result.output_tokens = 25 - mock_result.duration_ms = 1000 - - mock_exec = mocker.AsyncMock() - mock_exec.run.return_value = mock_result - - mocker.patch( - "app.endpoints.agents.get_step_executor", - return_value=mock_exec, - ) - return mock_exec - - -class TestRunAgentHandler: - """Tests for run_agent_handler.""" - - @pytest.mark.asyncio - async def test_successful_run( - self, - mocker: MockerFixture, - mock_config: Any, - mock_executor: Any, - ) -> None: - """Successful agent run returns result.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - - body = AgentRunRequest( - prompt="Analyze the cluster", - provider="openai", - model="gpt-4o-mini", - ) - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - result = await run_agent_handler.__wrapped__(request, body, auth) - - assert result["status"] == "completed" - assert result["output"] == {"summary": "Done"} - assert result["token_usage"]["input_tokens"] == 50 - - @pytest.mark.asyncio - async def test_passes_provider_and_model( - self, - mocker: MockerFixture, - mock_config: Any, - mock_executor: Any, - ) -> None: - """Provider and model are passed through to step input.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - - body = AgentRunRequest( - prompt="Hello", - provider="openai", - model="gpt-4o-mini", - ) - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - call_args = mock_executor.run.call_args[0][0] - assert call_args.provider["name"] == "openai" - assert call_args.provider["model"] == "gpt-4o-mini" - - @pytest.mark.asyncio - async def test_ephemeral_passes_real_spawner_to_executor( - self, - mocker: MockerFixture, - ) -> None: - """spawn=ephemeral with a spawner config builds and passes a real spawner. - - Regression test: get_step_executor(step_def, spawner=None) used to be - hardcoded regardless of spawner_configuration, silently no-opping - ephemeral spawn through this endpoint. - """ - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_cfg = mocker.patch("app.endpoints.agents.configuration") - spawner_config = mocker.MagicMock() - mock_cfg.spawner_configuration = spawner_config - mock_cfg.inference.default_provider = None - mock_cfg.inference.default_model = None - - fake_spawner = mocker.MagicMock() - mock_build_spawner = mocker.patch( - "app.endpoints.agents.build_spawner", return_value=fake_spawner - ) - - mock_result = mocker.MagicMock() - mock_result.status = "completed" - mock_result.output = {} - mock_result.error = None - mock_result.transcript = [] - mock_result.input_tokens = 0 - mock_result.output_tokens = 0 - mock_result.duration_ms = 0 - mock_exec = mocker.AsyncMock() - mock_exec.run.return_value = mock_result - mock_get_step_executor = mocker.patch( - "app.endpoints.agents.get_step_executor", return_value=mock_exec - ) - - body = AgentRunRequest( - prompt="Fix the issue", - spawn="ephemeral", - ) - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - mock_build_spawner.assert_called_once_with(spawner_config) - mock_get_step_executor.assert_called_once() - _, kwargs = mock_get_step_executor.call_args - assert kwargs["spawner"] is fake_spawner - - @pytest.mark.asyncio - async def test_ephemeral_threads_sandbox_image_and_credentials( - self, - mocker: MockerFixture, - mock_executor: Any, - ) -> None: - """spawn=ephemeral passes sandbox_image and credentials_secret through. - - Regression test: SandboxExecutor.run() reads step_input.sandbox_image - and step_input.provider["credentials_secret"] to know which container - image to use and which env var holds the LLM API key. Without these, - a real sandbox spawns but can't call the LLM (no credentials) and - uses the wrong image. - """ - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_cfg = mocker.patch("app.endpoints.agents.configuration") - spawner_config = mocker.MagicMock() - spawner_config.sandbox_image = "default-sandbox:latest" - mock_cfg.spawner_configuration = spawner_config - mock_cfg.inference.default_provider = None - mock_cfg.inference.default_model = None - mocker.patch( - "app.endpoints.agents.build_spawner", return_value=mocker.MagicMock() - ) - - body = AgentRunRequest( - prompt="Fix the issue", - spawn="ephemeral", - provider="openai", - sandbox_image="custom-sandbox:v2", - ) - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - call_args = mock_executor.run.call_args[0][0] - assert call_args.sandbox_image == "custom-sandbox:v2" - assert call_args.provider["credentials_secret"] == "OPENAI_API_KEY" - - @pytest.mark.asyncio - async def test_ephemeral_unknown_provider_omits_credentials_secret( - self, - mocker: MockerFixture, - mock_executor: Any, - ) -> None: - """spawn=ephemeral with no/unknown provider omits credentials_secret. - - Regression test: a prior version defaulted to "OPENAI_API_KEY" - regardless of provider, which would silently stamp the wrong (or a - nonexistent) env var name onto the sandbox for any non-OpenAI or - misspelled provider. Omitting the key lets the sandbox fail loudly - instead of guessing. - """ - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_cfg = mocker.patch("app.endpoints.agents.configuration") - spawner_config = mocker.MagicMock() - spawner_config.sandbox_image = "default-sandbox:latest" - mock_cfg.spawner_configuration = spawner_config - mock_cfg.inference.default_provider = None - mock_cfg.inference.default_model = None - mocker.patch( - "app.endpoints.agents.build_spawner", return_value=mocker.MagicMock() - ) - - body = AgentRunRequest(prompt="Fix the issue", spawn="ephemeral") - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - call_args = mock_executor.run.call_args[0][0] - assert "credentials_secret" not in call_args.provider - - @pytest.mark.asyncio - async def test_ephemeral_credentials_secret_matches_provider( - self, - mocker: MockerFixture, - mock_executor: Any, - ) -> None: - """credentials_secret uses the env var matching body.provider, not always OpenAI. - - Regression test: hardcoding "OPENAI_API_KEY" regardless of provider - would make the sandbox look for an Anthropic/Gemini/Azure key under - the wrong env var name, so the LLM call inside the sandbox would - fail (or silently pick up an unrelated OpenAI key) for any - non-OpenAI provider. - """ - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_cfg = mocker.patch("app.endpoints.agents.configuration") - spawner_config = mocker.MagicMock() - spawner_config.sandbox_image = "default-sandbox:latest" - mock_cfg.spawner_configuration = spawner_config - mock_cfg.inference.default_provider = None - mock_cfg.inference.default_model = None - mocker.patch( - "app.endpoints.agents.build_spawner", return_value=mocker.MagicMock() - ) - - body = AgentRunRequest( - prompt="Fix the issue", - spawn="ephemeral", - provider="anthropic", - model="claude-sonnet-5", - ) - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - call_args = mock_executor.run.call_args[0][0] - assert call_args.provider["credentials_secret"] == "ANTHROPIC_API_KEY" - - @pytest.mark.asyncio - async def test_ephemeral_falls_back_to_spawner_config_sandbox_image( - self, - mocker: MockerFixture, - mock_executor: Any, - ) -> None: - """Without a request-level sandbox_image, falls back to spawner config.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_cfg = mocker.patch("app.endpoints.agents.configuration") - spawner_config = mocker.MagicMock() - spawner_config.sandbox_image = "default-sandbox:latest" - mock_cfg.spawner_configuration = spawner_config - mock_cfg.inference.default_provider = None - mock_cfg.inference.default_model = None - mocker.patch( - "app.endpoints.agents.build_spawner", return_value=mocker.MagicMock() - ) - - body = AgentRunRequest(prompt="Fix the issue", spawn="ephemeral") - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - call_args = mock_executor.run.call_args[0][0] - assert call_args.sandbox_image == "default-sandbox:latest" - - @pytest.mark.asyncio - async def test_non_ephemeral_has_no_credentials_secret( - self, - mocker: MockerFixture, - mock_config: Any, - mock_executor: Any, - ) -> None: - """spawn=none doesn't set credentials_secret (irrelevant, no sandbox).""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - - body = AgentRunRequest(prompt="Hello", provider="openai", model="gpt-4o-mini") - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - call_args = mock_executor.run.call_args[0][0] - assert "credentials_secret" not in call_args.provider - - @pytest.mark.asyncio - async def test_ephemeral_without_spawner_raises( - self, - mocker: MockerFixture, - mock_config: Any, - ) -> None: - """spawn=ephemeral without spawner config raises 400.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - - body = AgentRunRequest( - prompt="Fix the issue", - spawn="ephemeral", - ) - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - from fastapi import HTTPException - - with pytest.raises(HTTPException) as exc_info: - await run_agent_handler.__wrapped__(request, body, auth) - - assert exc_info.value.status_code == 400 - - @pytest.mark.asyncio - async def test_falls_back_to_configured_default_provider_and_model( - self, - mocker: MockerFixture, - mock_config: Any, - mock_executor: Any, - ) -> None: - """Omitted provider/model fall back to inference.default_provider/model. - - Regression test: the docstring for AgentRunRequest.model promises a - fallback to inference.default_model (mirroring - start_workflow_handler's behavior for RunWorkflowRequest.provider), - but the handler used to hardcode `body.provider or ""` / - `body.model or ""` with no fallback, silently producing - {"name": "", "model": ""} and a confusing "Unknown provider ''" - error deep in cloud_agents even when defaults were configured. - """ - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_config.inference.default_provider = "openai" - mock_config.inference.default_model = "gpt-4o-mini" - - body = AgentRunRequest(prompt="Analyze the cluster") - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - call_args = mock_executor.run.call_args[0][0] - assert call_args.provider["name"] == "openai" - assert call_args.provider["model"] == "gpt-4o-mini" - - @pytest.mark.asyncio - async def test_explicit_provider_and_model_override_configured_defaults( - self, - mocker: MockerFixture, - mock_config: Any, - mock_executor: Any, - ) -> None: - """Explicit provider/model in the request take priority over defaults.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_config.inference.default_provider = "openai" - mock_config.inference.default_model = "gpt-4o-mini" - - body = AgentRunRequest( - prompt="Analyze the cluster", - provider="anthropic", - model="claude-sonnet-5", - ) - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - call_args = mock_executor.run.call_args[0][0] - assert call_args.provider["name"] == "anthropic" - assert call_args.provider["model"] == "claude-sonnet-5" - - @pytest.mark.asyncio - async def test_ephemeral_credentials_secret_uses_default_provider_fallback( - self, - mocker: MockerFixture, - mock_executor: Any, - ) -> None: - """Ephemeral credentials_secret resolution honors the fallback provider. - - Regression test: cred_secret used to be derived from raw - `body.provider or ""`, so a configured default_provider had no - effect on which credentials env var got threaded into the sandbox. - """ - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_cfg = mocker.patch("app.endpoints.agents.configuration") - spawner_config = mocker.MagicMock() - spawner_config.sandbox_image = "default-sandbox:latest" - mock_cfg.spawner_configuration = spawner_config - mock_cfg.inference.default_provider = "anthropic" - mock_cfg.inference.default_model = "claude-sonnet-5" - mocker.patch( - "app.endpoints.agents.build_spawner", return_value=mocker.MagicMock() - ) - - body = AgentRunRequest(prompt="Fix the issue", spawn="ephemeral") - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - call_args = mock_executor.run.call_args[0][0] - assert call_args.provider["name"] == "anthropic" - assert call_args.provider["model"] == "claude-sonnet-5" - assert call_args.provider["credentials_secret"] == "ANTHROPIC_API_KEY" - - @pytest.mark.asyncio - async def test_local_spawn_dispatches_without_spawner( - self, - mocker: MockerFixture, - mock_config: Any, - mock_executor: Any, - ) -> None: - """spawn=local dispatches via get_step_executor with no spawner. - - Mirrors spawn=none: SubprocessExecutor inherits the host process's - env directly, so no spawner/credentials_secret plumbing is needed. - """ - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_build_spawner = mocker.patch("app.endpoints.agents.build_spawner") - mock_get_step_executor = mocker.patch( - "app.endpoints.agents.get_step_executor", return_value=mock_executor - ) - - body = AgentRunRequest(prompt="Fix the issue", spawn="local") - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - mock_build_spawner.assert_not_called() - _, kwargs = mock_get_step_executor.call_args - assert kwargs["spawner"] is None - assert mock_get_step_executor.call_args[0][0]["spawn"] == "local" - call_args = mock_executor.run.call_args[0][0] - assert "credentials_secret" not in call_args.provider - - @pytest.mark.asyncio - async def test_allowed_skills_reach_step_input_and_raw_step( - self, - mocker: MockerFixture, - mock_executor: Any, - ) -> None: - """allowed_skills is threaded to StepInput and the raw step dict. - - DirectExecutor/SubprocessExecutor read step_input.allowed_skills; - SandboxExecutor only forwards step_input.raw_step to step_runner, - which is where spawner.spawn(allowed_skills=...) (per-skill - Landlock grants) comes from. Both must carry the allowlist or the - ephemeral path silently drops it. - """ - mocker.patch("app.endpoints.agents.check_configuration_loaded") - mock_cfg = mocker.patch("app.endpoints.agents.configuration") - spawner_config = mocker.MagicMock() - spawner_config.sandbox_image = "default-sandbox:latest" - mock_cfg.spawner_configuration = spawner_config - mock_cfg.inference.default_provider = None - mock_cfg.inference.default_model = None - mocker.patch( - "app.endpoints.agents.build_spawner", return_value=mocker.MagicMock() - ) - - body = AgentRunRequest( - prompt="Diagnose the pod", - spawn="ephemeral", - provider="openai", - model="gpt-4o-mini", - allowed_skills=["k8s-diag"], - ) - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - call_args = mock_executor.run.call_args[0][0] - assert call_args.allowed_skills == ["k8s-diag"] - assert call_args.raw_step["allowed_skills"] == ["k8s-diag"] - assert call_args.raw_step["name"] == "agent-run" - - @pytest.mark.asyncio - async def test_raw_step_selects_request_mcp_servers_by_name( - self, - mocker: MockerFixture, - mock_config: Any, - mock_executor: Any, - ) -> None: - """raw_step lists request MCP server names for step_runner injection. - - step_runner only injects catalog entries (StepInput.mcp_servers) - whose names appear in step["mcp_servers"] -- without the name list, - LIGHTSPEED_MCP_SERVERS is never set and the sandbox agent has no - MCP tools to call. - """ - mocker.patch("app.endpoints.agents.check_configuration_loaded") - - body = AgentRunRequest( - prompt="Check the pod", - provider="openai", - model="gpt-4o-mini", - mcp_servers=[{"name": "pod-status", "url": "http://x/mcp"}], - ) - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - call_args = mock_executor.run.call_args[0][0] - assert call_args.raw_step["mcp_servers"] == ["pod-status"] - - @pytest.mark.asyncio - async def test_allowed_skills_omitted_means_no_skills( - self, - mocker: MockerFixture, - mock_config: Any, - mock_executor: Any, - ) -> None: - """Omitted allowed_skills stays None (no skills), not empty filtering.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - - body = AgentRunRequest(prompt="Hello", provider="openai", model="gpt-4o-mini") - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - await run_agent_handler.__wrapped__(request, body, auth) - - call_args = mock_executor.run.call_args[0][0] - assert call_args.allowed_skills is None - assert call_args.raw_step["allowed_skills"] is None - - @pytest.mark.asyncio - async def test_local_spawn_result_shape( - self, - mocker: MockerFixture, - mock_config: Any, - mock_executor: Any, - ) -> None: - """spawn=local returns the same result shape as spawn=none.""" - mocker.patch("app.endpoints.agents.check_configuration_loaded") - - body = AgentRunRequest(prompt="Fix the issue", spawn="local") - auth = ("user-1", "testuser", False, "token") - request = mocker.MagicMock() - - result = await run_agent_handler.__wrapped__(request, body, auth) - - assert result["status"] == "completed" - assert result["output"] == {"summary": "Done"} diff --git a/tests/unit/cloud_agents/test_workflows_endpoint.py b/tests/unit/cloud_agents/test_workflows_endpoint.py index e05d10e09..ccf8b316c 100644 --- a/tests/unit/cloud_agents/test_workflows_endpoint.py +++ b/tests/unit/cloud_agents/test_workflows_endpoint.py @@ -149,6 +149,111 @@ async def test_starts_workflow( assert result["status"] == "running" mock_executor.start.assert_called_once() + @pytest.mark.asyncio + async def test_one_step_workflow_definition_forwarded( + self, + mocker: MockerFixture, + mock_config: Any, + mock_executor: Any, + ) -> None: + """A one-step workflow (one-shot agent run) forwards its definition as-is. + + One-step workflows take the same start path as multi-step ones; the + handler does not special-case them. + """ + mocker.patch("app.endpoints.workflows.check_configuration_loaded") + mock_executor.start.return_value = "wf-oneshot" + + definition = { + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "one-shot-agent"}, + "spec": { + "steps": [ + { + "name": "agent", + "type": "agent", + "prompt": "Inspect the cluster", + "output_key": "result", + "spawn": "none", + } + ] + }, + } + body = RunWorkflowRequest( + definition=definition, + provider={"name": "openai", "model": "gpt-4o-mini"}, + ) + auth = ("user-1", "testuser", False, "token") + request = mocker.MagicMock() + + result = await start_workflow_handler.__wrapped__(request, body, auth) + + assert result["workflow_id"] == "wf-oneshot" + workflow_input = mock_executor.start.call_args[0][0] + assert workflow_input["definition"] == definition + assert workflow_input["provider"]["name"] == "openai" + + @pytest.mark.asyncio + async def test_provider_falls_back_to_configured_defaults( + self, + mocker: MockerFixture, + mock_config: Any, + mock_executor: Any, + ) -> None: + """An omitted provider falls back to inference.default_provider/model.""" + mocker.patch("app.endpoints.workflows.check_configuration_loaded") + mock_config.spawner_configuration = None + mock_executor.start.return_value = "wf-abc123" + + body = RunWorkflowRequest( + definition={ + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "test-wf"}, + "spec": {"steps": []}, + }, + ) + auth = ("user-1", "testuser", False, "token") + request = mocker.MagicMock() + + await start_workflow_handler.__wrapped__(request, body, auth) + + workflow_input = mock_executor.start.call_args[0][0] + assert workflow_input["provider"]["name"] == "openai" + assert workflow_input["provider"]["model"] == "gpt-4o" + assert workflow_input["provider"]["credentials_secret"] == "OPENAI_API_KEY" + + @pytest.mark.asyncio + async def test_explicit_provider_overrides_configured_defaults( + self, + mocker: MockerFixture, + mock_config: Any, + mock_executor: Any, + ) -> None: + """An explicit request provider wins over inference defaults.""" + mocker.patch("app.endpoints.workflows.check_configuration_loaded") + mock_config.spawner_configuration = None + mock_executor.start.return_value = "wf-abc123" + + body = RunWorkflowRequest( + definition={ + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "test-wf"}, + "spec": {"steps": []}, + }, + provider={"name": "anthropic", "model": "claude-sonnet-5"}, + ) + auth = ("user-1", "testuser", False, "token") + request = mocker.MagicMock() + + await start_workflow_handler.__wrapped__(request, body, auth) + + workflow_input = mock_executor.start.call_args[0][0] + assert workflow_input["provider"]["name"] == "anthropic" + assert workflow_input["provider"]["model"] == "claude-sonnet-5" + @pytest.mark.asyncio async def test_credentials_secret_added_for_known_provider( self, diff --git a/uv.lock b/uv.lock index dba971861..485e3a960 100644 --- a/uv.lock +++ b/uv.lock @@ -1793,7 +1793,7 @@ wheels = [ [[package]] name = "lightspeed-cloud-agents" version = "0.1.0" -source = { git = "https://github.com/jameswnl/lightspeed-cloud-agents.git?branch=main#a5b30eb832432f78223c9bd1df6126d3a08d5c07" } +source = { git = "https://github.com/jameswnl/lightspeed-cloud-agents.git?branch=main#ffdc8933d5ed6ef487c0ade79267199549839d4d" } dependencies = [ { name = "alembic" }, { name = "asyncpg" }, From 0679d5c0e15bd0b7c2b772e795416596bbd475bd Mon Sep 17 00:00:00 2001 From: James Wong <2421248+jameswnl@users.noreply.github.com> Date: Tue, 29 Sep 2026 19:35:31 -0400 Subject: [PATCH 2/3] Add submission-time validation gate to /v1/workflows/run (#55) Address review feedback on PR #56. LocalWorkflowRunner.start persists the raw definition and provider BEFORE build_graph/normalization runs, and the stack route never called validate_definition -- so a secret-bearing definition returned 202 and landed in workflow state before the async step failed. start_workflow_ handler now rejects such input with 422 before dispatch, mirroring cloud-agents' own local/api.py submission gate: - validate_definition errors -> 422 (secret-bearing MCP, malformed steps) - WorkflowDefinition model errors -> 422 (unknown step types/fields) - secret-shaped run-level credentials_secret -> 422 Tests: four handler-level 422 regression tests (each asserting the executor is never started), an HTTP-level 422 rejection test, and spawn=local/ephemeral one-step HTTP tests for spawn-mode parity. Existing handler tests now use a valid one-step definition instead of an empty steps list. Also pins the bare-step ephemeral spawn default. --- src/app/endpoints/workflows.py | 32 +++ .../cloud_agents/test_workflows_http_e2e.py | 136 +++++++++++ .../test_workflows_integration.py | 7 +- .../cloud_agents/test_workflows_endpoint.py | 230 ++++++++++++++---- 4 files changed, 356 insertions(+), 49 deletions(-) diff --git a/src/app/endpoints/workflows.py b/src/app/endpoints/workflows.py index 17e3fc24d..9a2cefdae 100644 --- a/src/app/endpoints/workflows.py +++ b/src/app/endpoints/workflows.py @@ -4,7 +4,11 @@ from typing import Annotated, Any +from cloud_agents.workflow.core.definition import WorkflowDefinition +from cloud_agents.workflow.core.execution import validate_credential_reference +from cloud_agents.workflow.core.validation import validate_definition from fastapi import APIRouter, Depends, HTTPException, Request, status +from pydantic import ValidationError from authentication import get_auth_dependency from authentication.interface import AuthTuple @@ -87,6 +91,34 @@ async def start_workflow_handler( if cred_secret: provider["credentials_secret"] = cred_secret + # Submission-time gate (issue #55): LocalWorkflowRunner.start persists + # the raw definition and provider BEFORE build_graph/normalization + # runs, so secret-bearing or malformed input must be rejected here -- + # otherwise it returns 202 and lands in workflow state first. Mirrors + # cloud-agents' own local/api.py submission gate. + definition_errors = validate_definition(body.definition) + if definition_errors: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + detail={"validation_errors": definition_errors}, + ) + try: + WorkflowDefinition.model_validate(body.definition) + except ValidationError as exc: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + detail={"validation_errors": [str(exc)]}, + ) from exc + run_credentials_secret = provider.get("credentials_secret") + if run_credentials_secret is not None: + try: + validate_credential_reference(run_credentials_secret) + except ValueError as exc: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + detail={"validation_errors": [f"provider: {exc}"]}, + ) from exc + spawner_config = configuration.spawner_configuration default_sandbox_image = ( spawner_config.sandbox_image # pylint: disable=no-member diff --git a/tests/e2e/cloud_agents/test_workflows_http_e2e.py b/tests/e2e/cloud_agents/test_workflows_http_e2e.py index 3a644f2c2..c8c755971 100644 --- a/tests/e2e/cloud_agents/test_workflows_http_e2e.py +++ b/tests/e2e/cloud_agents/test_workflows_http_e2e.py @@ -442,3 +442,139 @@ def test_one_step_workflow_completes(self, http_client: TestClient) -> None: ) assert transcripts_response.status_code == 200 assert "result" in transcripts_response.json()["transcripts"] + + def test_secret_bearing_definition_rejected_with_422( + self, http_client: TestClient + ) -> None: + """A credentialed MCP URL is rejected at submission, before persistence. + + The runner persists the raw definition before normalization runs, + so the route must return 422 here -- a 202 would store the + secret-bearing URL in workflow state first. + """ + start_response = http_client.post( + "/v1/workflows/run", + json={ + "definition": { + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "one-shot-agent"}, + "spec": { + "steps": [ + { + "name": "agent", + "type": "agent", + "spawn": "none", + "output_key": "result", + "prompt": "Is pod checkout-7f9 healthy?", + "mcp_servers": [ + { + "name": "cluster", + "url": "https://admin:s3cret@example.com/mcp", + } + ], + "timeout_seconds": 120, + }, + ] + }, + }, + "provider": {"name": "openai", "model": "gpt-4o-mini"}, + }, + ) + + assert start_response.status_code == 422 + assert "workflow_id" not in start_response.json() + + def test_one_step_workflow_with_local_spawn(self, http_client: TestClient) -> None: + """One-step spawn:local workflow completes over real HTTP. + + No output_schema (see test_workflow_with_local_spawn_step): the + SubprocessExecutor has no native structured-output mode yet + (jameswnl/lightspeed-cloud-agents#235). + """ + start_response = http_client.post( + "/v1/workflows/run", + json={ + "definition": { + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "one-shot-agent-local"}, + "spec": { + "steps": [ + { + "name": "agent", + "type": "agent", + "spawn": "local", + "output_key": "result", + "prompt": ( + "Say one sentence confirming the " + "checkout-7f9 pod is healthy." + ), + "timeout_seconds": 120, + }, + ] + }, + }, + "provider": {"name": "openai", "model": "gpt-4o-mini"}, + }, + ) + + assert start_response.status_code == 202 + workflow_id = start_response.json()["workflow_id"] + assert workflow_id + + completed = wait_for_status( + http_client, + workflow_id, + lambda body: bool(body["is_terminal"]), + timeout_s=150, + ) + assert completed["status"] == "completed" + assert "result" in completed["steps"] + + @pytest.mark.ephemeral + def test_one_step_workflow_with_ephemeral_spawn( + self, http_client: TestClient + ) -> None: + """One-step spawn:ephemeral workflow completes over real HTTP.""" + skip_if_gateway_unreachable() + + start_response = http_client.post( + "/v1/workflows/run", + json={ + "definition": { + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "one-shot-agent-ephemeral"}, + "spec": { + "steps": [ + { + "name": "agent", + "type": "agent", + "spawn": "ephemeral", + "output_key": "result", + "prompt": ( + "Say one sentence confirming the " + "checkout-7f9 pod is healthy." + ), + "timeout_seconds": 120, + }, + ] + }, + }, + "provider": {"name": "openai", "model": "gpt-4o-mini"}, + }, + ) + + assert start_response.status_code == 202 + workflow_id = start_response.json()["workflow_id"] + assert workflow_id + + completed = wait_for_status( + http_client, + workflow_id, + lambda body: bool(body["is_terminal"]), + timeout_s=150, + ) + assert completed["status"] == "completed" + assert "result" in completed["steps"] diff --git a/tests/integration/cloud_agents/test_workflows_integration.py b/tests/integration/cloud_agents/test_workflows_integration.py index 5aaa7a62d..e9b565fec 100644 --- a/tests/integration/cloud_agents/test_workflows_integration.py +++ b/tests/integration/cloud_agents/test_workflows_integration.py @@ -81,7 +81,11 @@ def test_full_one_step_shape(self) -> None: assert step_input.timeout_seconds == 120 def test_bare_one_step_defaults(self) -> None: - """A bare single step defaults to name 'agent' / output_key 'result'.""" + """A bare single step defaults to agent/result naming and ephemeral spawn. + + Locks in the canonical one-step defaults: the ``agent`` / + ``result`` naming convention and the ``ephemeral`` spawn default. + """ definition = _one_step_definition({"prompt": "Inspect the cluster"}) assert validate_definition(definition) == [] @@ -90,6 +94,7 @@ def test_bare_one_step_defaults(self) -> None: assert len(agent_steps) == 1 assert agent_steps[0].name == "agent" assert agent_steps[0].output_key == "result" + assert agent_steps[0].spawn == "ephemeral" def test_one_step_matches_multi_step_normalization(self) -> None: """A one-step workflow normalizes identically to the same multi-step step. diff --git a/tests/unit/cloud_agents/test_workflows_endpoint.py b/tests/unit/cloud_agents/test_workflows_endpoint.py index ccf8b316c..ece032a45 100644 --- a/tests/unit/cloud_agents/test_workflows_endpoint.py +++ b/tests/unit/cloud_agents/test_workflows_endpoint.py @@ -108,6 +108,26 @@ def test_no_spawner_config_passes_none( mock_create_runner.assert_called_once_with(spawner=None) +def _valid_definition() -> dict[str, Any]: + """Minimal valid one-step workflow definition for handler tests.""" + return { + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "test-wf"}, + "spec": { + "steps": [ + { + "name": "agent", + "type": "agent", + "prompt": "Do it", + "output_key": "result", + "spawn": "none", + } + ] + }, + } + + class TestStartWorkflow: """Tests for start_workflow_handler.""" @@ -207,12 +227,7 @@ async def test_provider_falls_back_to_configured_defaults( mock_executor.start.return_value = "wf-abc123" body = RunWorkflowRequest( - definition={ - "apiVersion": "v1", - "kind": "AgentWorkflow", - "metadata": {"name": "test-wf"}, - "spec": {"steps": []}, - }, + definition=_valid_definition(), ) auth = ("user-1", "testuser", False, "token") request = mocker.MagicMock() @@ -237,12 +252,7 @@ async def test_explicit_provider_overrides_configured_defaults( mock_executor.start.return_value = "wf-abc123" body = RunWorkflowRequest( - definition={ - "apiVersion": "v1", - "kind": "AgentWorkflow", - "metadata": {"name": "test-wf"}, - "spec": {"steps": []}, - }, + definition=_valid_definition(), provider={"name": "anthropic", "model": "claude-sonnet-5"}, ) auth = ("user-1", "testuser", False, "token") @@ -254,6 +264,160 @@ async def test_explicit_provider_overrides_configured_defaults( assert workflow_input["provider"]["name"] == "anthropic" assert workflow_input["provider"]["model"] == "claude-sonnet-5" + @pytest.mark.asyncio + async def test_secret_bearing_mcp_definition_rejected_with_422( + self, + mocker: MockerFixture, + mock_config: Any, + mock_executor: Any, + ) -> None: + """A definition with a credentialed MCP URL is rejected before start. + + Regression test: LocalWorkflowRunner.start persists the raw + definition before build_graph/normalization runs, so without this + route-level gate the request would return 202 and store the + secret-bearing URL in workflow state first. + """ + mocker.patch("app.endpoints.workflows.check_configuration_loaded") + mock_config.spawner_configuration = None + + body = RunWorkflowRequest( + definition={ + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "test-wf"}, + "spec": { + "steps": [ + { + "name": "agent", + "type": "agent", + "prompt": "Do it", + "output_key": "result", + "mcp_servers": [ + { + "name": "cluster", + "url": "https://admin:s3cret@example.com/mcp", + } + ], + } + ] + }, + }, + provider={"name": "openai", "model": "gpt-4o-mini"}, + ) + auth = ("user-1", "testuser", False, "token") + request = mocker.MagicMock() + + from fastapi import HTTPException + + with pytest.raises(HTTPException) as exc_info: + await start_workflow_handler.__wrapped__(request, body, auth) + + assert exc_info.value.status_code == 422 + assert "validation_errors" in exc_info.value.detail + mock_executor.start.assert_not_called() + + @pytest.mark.asyncio + async def test_secret_shaped_run_credentials_rejected_with_422( + self, + mocker: MockerFixture, + mock_config: Any, + mock_executor: Any, + ) -> None: + """A secret-shaped run-level credentials_secret is rejected before start.""" + mocker.patch("app.endpoints.workflows.check_configuration_loaded") + mock_config.spawner_configuration = None + + body = RunWorkflowRequest( + definition=_valid_definition(), + provider={ + "name": "openai", + "model": "gpt-4o-mini", + "credentials_secret": "sk-test", + }, + ) + auth = ("user-1", "testuser", False, "token") + request = mocker.MagicMock() + + from fastapi import HTTPException + + with pytest.raises(HTTPException) as exc_info: + await start_workflow_handler.__wrapped__(request, body, auth) + + assert exc_info.value.status_code == 422 + mock_executor.start.assert_not_called() + + @pytest.mark.asyncio + async def test_malformed_definition_rejected_with_422( + self, + mocker: MockerFixture, + mock_config: Any, + mock_executor: Any, + ) -> None: + """A structurally invalid definition is rejected before start.""" + mocker.patch("app.endpoints.workflows.check_configuration_loaded") + mock_config.spawner_configuration = None + + body = RunWorkflowRequest( + definition={ + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "test-wf"}, + "spec": {"steps": "not-a-list"}, + }, + provider={"name": "openai", "model": "gpt-4o-mini"}, + ) + auth = ("user-1", "testuser", False, "token") + request = mocker.MagicMock() + + from fastapi import HTTPException + + with pytest.raises(HTTPException) as exc_info: + await start_workflow_handler.__wrapped__(request, body, auth) + + assert exc_info.value.status_code == 422 + mock_executor.start.assert_not_called() + + @pytest.mark.asyncio + async def test_unknown_step_type_rejected_with_422( + self, + mocker: MockerFixture, + mock_config: Any, + mock_executor: Any, + ) -> None: + """An unknown step type fails model validation before start.""" + mocker.patch("app.endpoints.workflows.check_configuration_loaded") + mock_config.spawner_configuration = None + + body = RunWorkflowRequest( + definition={ + "apiVersion": "v1", + "kind": "AgentWorkflow", + "metadata": {"name": "test-wf"}, + "spec": { + "steps": [ + { + "name": "weird", + "type": "bogus", + "prompt": "Do it", + "output_key": "result", + } + ] + }, + }, + provider={"name": "openai", "model": "gpt-4o-mini"}, + ) + auth = ("user-1", "testuser", False, "token") + request = mocker.MagicMock() + + from fastapi import HTTPException + + with pytest.raises(HTTPException) as exc_info: + await start_workflow_handler.__wrapped__(request, body, auth) + + assert exc_info.value.status_code == 422 + mock_executor.start.assert_not_called() + @pytest.mark.asyncio async def test_credentials_secret_added_for_known_provider( self, @@ -272,12 +436,7 @@ async def test_credentials_secret_added_for_known_provider( mock_executor.start.return_value = "wf-abc123" body = RunWorkflowRequest( - definition={ - "apiVersion": "v1", - "kind": "AgentWorkflow", - "metadata": {"name": "test-wf"}, - "spec": {"steps": []}, - }, + definition=_valid_definition(), provider={"name": "anthropic", "model": "claude-sonnet-5"}, ) auth = ("user-1", "testuser", False, "token") @@ -301,12 +460,7 @@ async def test_credentials_secret_omitted_for_unknown_provider( mock_executor.start.return_value = "wf-abc123" body = RunWorkflowRequest( - definition={ - "apiVersion": "v1", - "kind": "AgentWorkflow", - "metadata": {"name": "test-wf"}, - "spec": {"steps": []}, - }, + definition=_valid_definition(), provider={"name": "bedrock", "model": "some-model"}, ) auth = ("user-1", "testuser", False, "token") @@ -330,12 +484,7 @@ async def test_caller_supplied_credentials_secret_not_overridden( mock_executor.start.return_value = "wf-abc123" body = RunWorkflowRequest( - definition={ - "apiVersion": "v1", - "kind": "AgentWorkflow", - "metadata": {"name": "test-wf"}, - "spec": {"steps": []}, - }, + definition=_valid_definition(), provider={ "name": "openai", "model": "gpt-4o", @@ -370,12 +519,7 @@ async def test_sandbox_image_falls_back_to_spawner_config( mock_executor.start.return_value = "wf-abc123" body = RunWorkflowRequest( - definition={ - "apiVersion": "v1", - "kind": "AgentWorkflow", - "metadata": {"name": "test-wf"}, - "spec": {"steps": []}, - }, + definition=_valid_definition(), ) auth = ("user-1", "testuser", False, "token") request = mocker.MagicMock() @@ -400,12 +544,7 @@ async def test_sandbox_image_request_override_wins( mock_executor.start.return_value = "wf-abc123" body = RunWorkflowRequest( - definition={ - "apiVersion": "v1", - "kind": "AgentWorkflow", - "metadata": {"name": "test-wf"}, - "spec": {"steps": []}, - }, + definition=_valid_definition(), sandbox_image="custom-sandbox:v9", ) auth = ("user-1", "testuser", False, "token") @@ -429,12 +568,7 @@ async def test_no_spawner_config_falls_back_to_default_image( mock_executor.start.return_value = "wf-abc123" body = RunWorkflowRequest( - definition={ - "apiVersion": "v1", - "kind": "AgentWorkflow", - "metadata": {"name": "test-wf"}, - "spec": {"steps": []}, - }, + definition=_valid_definition(), ) auth = ("user-1", "testuser", False, "token") request = mocker.MagicMock() From 080a8edd25ca0805c098c44c8b80df356ff8de0a Mon Sep 17 00:00:00 2001 From: James Wong <2421248+jameswnl@users.noreply.github.com> Date: Tue, 29 Sep 2026 20:33:12 -0400 Subject: [PATCH 3/3] Validate run provider at submission; extract gate helper (#55) Address follow-up review feedback on PR #56. - The submission gate now also validates the run-level provider name/model against the executor's approved provider contract (inference_spec_from_provider_config), so unapproved providers 422 instead of persisting a run that fails asynchronously. Regression test: unapproved name rejected, executor never started. - Extract the gate into _validate_workflow_submission, fixing the R0914 too-many-locals on start_workflow_handler; pylint on the file is back to only the pre-existing E1101. --- src/app/endpoints/workflows.py | 84 +++++++++++++------ .../cloud_agents/test_workflows_endpoint.py | 26 ++++++ 2 files changed, 84 insertions(+), 26 deletions(-) diff --git a/src/app/endpoints/workflows.py b/src/app/endpoints/workflows.py index 9a2cefdae..e8dd27c4d 100644 --- a/src/app/endpoints/workflows.py +++ b/src/app/endpoints/workflows.py @@ -5,7 +5,10 @@ from typing import Annotated, Any from cloud_agents.workflow.core.definition import WorkflowDefinition -from cloud_agents.workflow.core.execution import validate_credential_reference +from cloud_agents.workflow.core.execution import ( + inference_spec_from_provider_config, + validate_credential_reference, +) from cloud_agents.workflow.core.validation import validate_definition from fastapi import APIRouter, Depends, HTTPException, Request, status from pydantic import ValidationError @@ -55,6 +58,56 @@ def _get_executor() -> Any: return _executor +def _validate_workflow_submission( + definition: dict[str, Any], provider: dict[str, Any] +) -> None: + """Reject invalid workflow input with 422 before the run is persisted. + + Mirrors cloud-agents' own local/api.py submission gate: the definition + must pass cloud-agents validation and model validation, and the + run-level provider must be an approved executor-known provider whose + credentials_secret (when present) is a reference, never a value. + + Parameters: + definition: Raw workflow definition from the request body. + provider: Merged run-level provider (request value or inference + defaults, with credentials_secret injected). + + Raises: + HTTPException: 422 with a ``validation_errors`` detail list when + any check fails. + """ + definition_errors = validate_definition(definition) + if definition_errors: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + detail={"validation_errors": definition_errors}, + ) + try: + WorkflowDefinition.model_validate(definition) + except ValidationError as exc: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + detail={"validation_errors": [str(exc)]}, + ) from exc + try: + inference_spec_from_provider_config(provider) + except ValueError as exc: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + detail={"validation_errors": [f"provider: {exc}"]}, + ) from exc + run_credentials_secret = provider.get("credentials_secret") + if run_credentials_secret is not None: + try: + validate_credential_reference(run_credentials_secret) + except ValueError as exc: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + detail={"validation_errors": [f"provider: {exc}"]}, + ) from exc + + @router.post( "/workflows/run", status_code=status.HTTP_202_ACCEPTED, @@ -93,31 +146,10 @@ async def start_workflow_handler( # Submission-time gate (issue #55): LocalWorkflowRunner.start persists # the raw definition and provider BEFORE build_graph/normalization - # runs, so secret-bearing or malformed input must be rejected here -- - # otherwise it returns 202 and lands in workflow state first. Mirrors - # cloud-agents' own local/api.py submission gate. - definition_errors = validate_definition(body.definition) - if definition_errors: - raise HTTPException( - status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, - detail={"validation_errors": definition_errors}, - ) - try: - WorkflowDefinition.model_validate(body.definition) - except ValidationError as exc: - raise HTTPException( - status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, - detail={"validation_errors": [str(exc)]}, - ) from exc - run_credentials_secret = provider.get("credentials_secret") - if run_credentials_secret is not None: - try: - validate_credential_reference(run_credentials_secret) - except ValueError as exc: - raise HTTPException( - status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, - detail={"validation_errors": [f"provider: {exc}"]}, - ) from exc + # runs, so secret-bearing, malformed, or unapproved-provider input + # must be rejected here -- otherwise it returns 202 and lands in + # workflow state first. + _validate_workflow_submission(body.definition, provider) spawner_config = configuration.spawner_configuration default_sandbox_image = ( diff --git a/tests/unit/cloud_agents/test_workflows_endpoint.py b/tests/unit/cloud_agents/test_workflows_endpoint.py index ece032a45..e65cebdf0 100644 --- a/tests/unit/cloud_agents/test_workflows_endpoint.py +++ b/tests/unit/cloud_agents/test_workflows_endpoint.py @@ -378,6 +378,32 @@ async def test_malformed_definition_rejected_with_422( assert exc_info.value.status_code == 422 mock_executor.start.assert_not_called() + @pytest.mark.asyncio + async def test_unapproved_run_provider_rejected_with_422( + self, + mocker: MockerFixture, + mock_config: Any, + mock_executor: Any, + ) -> None: + """A run provider outside the approved contract is rejected before start.""" + mocker.patch("app.endpoints.workflows.check_configuration_loaded") + mock_config.spawner_configuration = None + + body = RunWorkflowRequest( + definition=_valid_definition(), + provider={"name": "bogus", "model": "x"}, + ) + auth = ("user-1", "testuser", False, "token") + request = mocker.MagicMock() + + from fastapi import HTTPException + + with pytest.raises(HTTPException) as exc_info: + await start_workflow_handler.__wrapped__(request, body, auth) + + assert exc_info.value.status_code == 422 + mock_executor.start.assert_not_called() + @pytest.mark.asyncio async def test_unknown_step_type_rejected_with_422( self,