Workhorse: Add the Duo Workflow server-side execution endpoint
What does this MR do and why?
Adds the Workhorse side of Duo Workflow server-side execution:
POST /api/v4/ai/duo_workflows/workflows/:workflow_id/executeDepends on !254410 (merged) (this MR targets its branch) and pairs with !254443 (merged), which adds the Rails pre-authorization endpoint.
The problem
GET /api/v4/ai/duo_workflows/ws upgrades to a WebSocket and proxies between a
client and the Duo Workflow Service (DWS). In that model the client is the
executor: DWS sends it Action messages, it runs them and answers with an
ActionResponse. Workhorse already intercepts a few action types and executes
them itself (RunHTTPRequest, RunMCPTool, TrackLlmCallForSelfHosted).
Some callers want a flow run for them and cannot execute anything. A chat turn arriving from a Slack integration, handled in a Sidekiq job, has no filesystem, no shell and no way to answer an action. Today the only way to serve them is to start a CI job whose sole purpose is to hold the gRPC stream and answer the actions Workhorse already knows how to answer.
The endpoint
The request body is a protojson StartWorkflowRequest; Workhorse starts the
workflow itself instead of waiting for a client to send one. The response is
200 with Content-Type: application/x-ndjson: one protojson Action per
line, flushed as it arrives, over a single chunked response. Empty lines are
keepalives and are skipped by ndjson readers.
Only NewCheckpoint reaches the caller, which is enough to follow the run.
Every other action is answered by Workhorse on the caller's behalf with an error
ActionResponse. That part is not optional: DWS waits for a response to every
action it emits, so an action that is neither executed nor answered stalls the
graph until DWS times out.
The line format is deliberately unchanged from what WebSocket clients already parse, so there is no second wire format to maintain.
Notable decisions
-
The workflow ID comes from Rails, through the new
api.DuoWorkflow.WorkflowID, never from the request body, because Workhorse authorizes nothing itself. A body naming a different workflow is rejected rather than silently overridden. -
The caller's
clientCapabilities,mcpToolsandpreapprovedToolsare dropped. Capabilities describe an executor and Workhorse is not one here, so it advertises none and lets DWS fall back to its baseline behavior. The tools available to the flow are the ones Rails configured.runnerstill appends the server capabilities Rails reported. -
The response header is committed lazily, on the first action or keepalive. Both failures a caller most needs to tell apart, a held lock and a rejection by DWS, are raised before the first action, so deferring the commit lets them be reported as status codes:
Status Trigger 400Body cannot be decoded, names another workflow, or DWS rejects it as invalid 403Usage quota exhausted 409Workflow lock held by another run 500Pre-authorization response missing the service config or workflow ID 502gRPC stream to DWS cannot be opened This is a real gain over the WebSocket path, where the same conditions are close codes (
1013,1008,4400) the client has to interpret. -
After the first line, a failure only ends the stream. The caller tells a completed run from an interrupted one through the workflow's status, which DWS keeps up to date. There is no in-band terminal record, so the line format stays exactly the
Actionstream clients already parse. Trailers were considered and rejected: Ruby'sNet::HTTPdoes not expose them reliably. -
Keepalive writes an empty line every 20 s, matching the WebSocket ping interval. A flow can spend a minute inside a single model call, and proxies close a response that produces no bytes (nginx
proxy_read_timeoutdefaults to 60 s). -
No
StopWorkflowRequestwhen the caller hangs up. The gRPC stream is derived from the request context, so by then it is already gone. DWS sees the canceled stream and treats it as a disconnect, exactly as when a WebSocket client vanishes, leaving the workflow resumable from its last checkpoint. The stop request is still sent on the paths where the stream is intact, most importantly a Workhorse graceful shutdown while the caller is still connected. -
No maximum run duration, matching the WebSocket endpoint.
-
The body is capped at
MaxMessageSize(4 MB), the gRPC send limit, since a larger start request could not be forwarded anyway. Route body limits apply only to proxied requests, not to a handler reading the body itself, so this is an explicithttp.MaxBytesReader. -
connections_totalandconnection_errors_totalgain atransportlabel (websocket,http), so the new endpoint does not silently inflate the existing WebSocket numbers. No new metrics.
One fix outside the new endpoint
Only wait for the DWS reader when a stop was requested (b9abd8be) makes
runner.Close wait on agentDone only when a stop was actually requested.
The wait exists so a Recv in flight can observe the stop acknowledgment DWS
sends as an Unavailable status. When no stop was requested there is no
acknowledgment coming, and the wait lasts until the stream breaks on its own.
A workflow that never started because the lock was held elsewhere lands exactly
there: handleClientEvent fails on the lock before sending the start request, so
DWS is parked on its first Recv and Workhorse on ours. It is a pre-existing
issue that delays the 1013 close frame on the WebSocket path; on this endpoint
it would hold the 409 open indefinitely, so it is fixed here rather than
worked around. It is a separate commit for independent review.
How to set up and validate locally
Requires !254443 (merged) for the Rails endpoint. With a workflow that belongs to you:
curl -N -X POST \
-H "Authorization: Bearer $AI_WORKFLOWS_TOKEN" \
-H 'Content-Type: application/json' \
-d '{"goal": "summarize this project"}' \
"http://gdk.test:3000/api/v4/ai/duo_workflows/workflows/$WORKFLOW_ID/execute"Each line is a checkpoint action; blank lines are keepalives. Running it twice
concurrently returns 409 for the second call.
Test coverage
http_transport_test.go covers line framing, keepalives, action rejection,
write serialization under concurrency and the lazy header commit.
http_handler_test.go runs the handler end to end against the package's
in-process DWS gRPC server: streamed checkpoints, an action rejected back to DWS,
start request overrides, an oversized body, each pre-stream failure status, a
concurrent run, and caller hangup tearing down the stream.
MR acceptance checklist
This checklist encourages us to confirm any changes have been analyzed to reduce risks in quality, performance, reliability, security, and maintainability.
- I have evaluated the MR acceptance checklist for this MR.