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/execute

Depends 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, mcpTools and preapprovedTools are 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. runner still 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
    400 Body cannot be decoded, names another workflow, or DWS rejects it as invalid
    403 Usage quota exhausted
    409 Workflow lock held by another run
    500 Pre-authorization response missing the service config or workflow ID
    502 gRPC 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 Action stream clients already parse. Trailers were considered and rejected: Ruby's Net::HTTP does 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_timeout defaults to 60 s).

  • No StopWorkflowRequest when 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 explicit http.MaxBytesReader.

  • connections_total and connection_errors_total gain a transport label (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.

Merge request reports

Loading
Loading