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