Packages

A spec-compliant Open Responses server for Elixir and Phoenix, with pluggable provider adapters and first-class streaming.

Current section

Files

Jump to
open_responses lib open_responses_web controllers response_controller.ex
Raw

lib/open_responses_web/controllers/response_controller.ex

defmodule OpenResponsesWeb.ResponseController do
@moduledoc """
Handles POST /v1/responses.
Supports both streaming (SSE) and non-streaming modes.
## Required fields
| Field | Type | Description |
|---|---|---|
| `model` | string | Model identifier — determines provider routing |
| `api_key` | string | Caller's provider API key — passed directly to the adapter |
| `input` | list | Conversation input items |
Requests missing `api_key` are rejected with `400 invalid_request` before
any provider call is made. The server holds no shared API keys for client
requests.
"""
use OpenResponsesWeb, :controller
alias OpenResponses.Responses
alias OpenResponses.Responses.Response
alias OpenResponses.LoopSupervisor
alias OpenResponses.Loop
alias Phoenix.PubSub
@spec create(Plug.Conn.t(), map()) :: Plug.Conn.t()
def create(conn, params) do
:telemetry.execute(
[:open_responses, :request, :start],
%{system_time: System.system_time()},
%{model: params["model"]}
)
with :ok <- validate_provider_key(params),
{:ok, response} <- create_response(params) do
if params["stream"] == true do
stream_response(conn, response, params)
else
sync_response(conn, response, params)
end
else
{:error, %{type: type, message: message}} ->
conn
|> put_status(400)
|> json(%{error: %{type: type, message: message}})
{:error, %Ash.Error.Invalid{} = error} ->
conn
|> put_status(400)
|> json(%{error: %{type: "invalid_request", message: inspect(error)}})
{:error, reason} ->
conn
|> put_status(500)
|> json(%{error: %{type: "server_error", message: inspect(reason)}})
end
end
defp validate_provider_key(%{"api_key" => key}) when is_binary(key) and key != "", do: :ok
defp validate_provider_key(_params) do
{:error, %{type: "invalid_request", message: "api_key is required"}}
end
defp create_response(%{"input" => _} = params) do
Ash.create(Response, %{
model: params["model"],
input: params["input"] || [],
tools: params["tools"] || [],
tool_choice: params["tool_choice"],
temperature: params["temperature"],
top_p: params["top_p"],
max_output_tokens: params["max_output_tokens"],
previous_response_id: params["previous_response_id"],
metadata: params["metadata"] || %{}
}, domain: Responses)
end
defp create_response(params) do
missing = Enum.filter(["model", "input"], &(not Map.has_key?(params, &1)))
message = "Missing required fields: #{Enum.join(missing, ", ")}"
{:error, %Ash.Error.Invalid{errors: [%Ash.Error.Changes.InvalidAttribute{field: :input, message: message}]}}
end
defp stream_response(conn, response, params) do
topic = Loop.topic(response.id)
PubSub.subscribe(OpenResponses.PubSub, topic)
{:ok, _pid} = LoopSupervisor.start_loop(response: response, input: params["input"] || [], provider: %{"api_key" => params["api_key"]})
conn =
conn
|> put_resp_content_type("text/event-stream")
|> put_resp_header("cache-control", "no-cache")
|> put_resp_header("connection", "keep-alive")
|> send_chunked(200)
conn = send_sse_event(conn, "response.created", encode_response(response))
drain_loop_events(conn, response)
end
defp drain_loop_events(conn, response) do
receive do
{:loop_event, %{"type" => "response.completed"} = event} ->
conn = send_sse_event(conn, event["type"], event)
final = load_response(response.id)
conn = send_sse_event(conn, "response.completed", encode_response(final))
send_sse_done(conn)
{:loop_event, %{"type" => type} = event} when type in ["response.failed", "response.incomplete"] ->
conn = send_sse_event(conn, event["type"], event)
send_sse_done(conn)
{:loop_event, event} ->
conn = send_sse_event(conn, event["type"] || "event", event)
drain_loop_events(conn, response)
after
30_000 ->
send_sse_done(conn)
end
end
defp sync_response(conn, response, params) do
topic = Loop.topic(response.id)
PubSub.subscribe(OpenResponses.PubSub, topic)
{:ok, _pid} = LoopSupervisor.start_loop(response: response, input: params["input"] || [], provider: %{"api_key" => params["api_key"]})
result = await_completion(response)
json(conn, result)
end
defp await_completion(response) do
receive do
{:loop_event, %{"type" => "response.completed"}} ->
response.id |> load_response() |> encode_response()
{:loop_event, %{"type" => "response.failed"}} ->
response.id |> load_response() |> encode_response()
{:loop_event, %{"type" => "response.incomplete"}} ->
response.id |> load_response() |> encode_response()
{:loop_event, _} ->
await_completion(response)
after
30_000 ->
encode_response(response)
end
end
defp load_response(id) do
{:ok, resp} = Ash.get(Response, id, domain: Responses)
resp
end
defp send_sse_event(conn, type, data) do
payload = "event: #{type}\ndata: #{Jason.encode!(data)}\n\n"
{:ok, conn} = Plug.Conn.chunk(conn, payload)
conn
end
defp send_sse_done(conn) do
{:ok, conn} = Plug.Conn.chunk(conn, "data: [DONE]\n\n")
conn
end
defp encode_response(%Response{} = r) do
%{
id: r.id,
object: r.object,
model: r.model,
status: r.status,
output: r.output,
usage: r.usage,
created_at: r.created_at,
metadata: r.metadata
}
end
end