diff --git a/lib/n8n_openai_adapter/router.ex b/lib/n8n_openai_adapter/router.ex index dffd7471fd..2bfdf5060f 100644 --- a/lib/n8n_openai_adapter/router.ex +++ b/lib/n8n_openai_adapter/router.ex @@ -235,6 +235,7 @@ defmodule N8nOpenaiAdapter.Router do model = Map.get(data, "model") messages = Map.get(data, "messages", []) + stream? = Map.get(data, "stream", false) thread_id = Map.get(data, "thread_id") session_id = if is_binary(thread_id) and thread_id != "", do: thread_id, else: "default" @@ -258,32 +259,106 @@ defmodule N8nOpenaiAdapter.Router do webhook -> case call_n8n(webhook, session_id, chat_input) do {:ok, output} -> - resp = %{ - id: - "chatcmpl-#{Base.encode16(:crypto.strong_rand_bytes(12), case: :lower)}", - object: "chat.completion", - created: System.system_time(:second), - model: model, - choices: [ - %{ - index: 0, - message: %{role: "assistant", content: output}, - finish_reason: "stop" - } - ], - usage: %{prompt_tokens: 0, completion_tokens: 0, total_tokens: 0} - } - - json(conn, 200, resp) + if stream? do + stream_completion(conn, model, output) + else + json(conn, 200, completion_body(model, output)) + end {:error, msg} -> - error(conn, 502, msg) + if stream? do + stream_error(conn, msg) + else + error(conn, 502, msg) + end end end end end end + # Build a single, non-streaming chat.completion body. + defp completion_body(model, output) do + %{ + id: "chatcmpl-#{Base.encode16(:crypto.strong_rand_bytes(12), case: :lower)}", + object: "chat.completion", + created: System.system_time(:second), + model: model, + choices: [ + %{ + index: 0, + message: %{role: "assistant", content: output}, + finish_reason: "stop" + } + ], + usage: %{prompt_tokens: 0, completion_tokens: 0, total_tokens: 0} + } + end + + # Emit an OpenAI-compatible SSE stream. The n8n reply arrives whole, so we + # deliver it as a single delta chunk (role chunk first, then content) and + # terminate with `data: [DONE]`. OpenAI SDKs that request `stream: true` + # consume this shape; without it they render an empty response. + defp stream_completion(conn, model, output) do + id = "chatcmpl-#{Base.encode16(:crypto.strong_rand_bytes(12), case: :lower)}" + created = System.system_time(:second) + + conn = + conn + |> put_resp_content_type("text/event-stream") + |> put_resp_header("cache-control", "no-cache") + |> send_chunked(200) + + chunks = [ + %{ + id: id, + object: "chat.completion.chunk", + created: created, + model: model, + choices: [%{index: 0, delta: %{role: "assistant"}, finish_reason: nil}] + }, + %{ + id: id, + object: "chat.completion.chunk", + created: created, + model: model, + choices: [%{index: 0, delta: %{content: output}, finish_reason: nil}] + }, + %{ + id: id, + object: "chat.completion.chunk", + created: created, + model: model, + choices: [%{index: 0, delta: %{}, finish_reason: "stop"}] + } + ] + + {:ok, conn} = + Enum.reduce_while(chunks, {:ok, conn}, fn chunk, {:ok, conn} -> + case chunk(conn, "data: #{Jason.encode!(chunk)}\n\n") do + {:ok, conn} -> {:cont, {:ok, conn}} + {:error, _} -> {:halt, {:ok, conn}} + end + end) + + {:ok, conn} = chunk(conn, "data: [DONE]\n\n") + conn + end + + # Emit an SSE-formatted error so a streaming client sees the failure. + defp stream_error(conn, msg) do + conn = + conn + |> put_resp_content_type("text/event-stream") + |> put_resp_header("cache-control", "no-cache") + |> send_chunked(502) + + payload = %{error: %{message: msg}} + {:ok, conn} = chunk(conn, "data: #{Jason.encode!(payload)}\n\n") + {:ok, conn} = chunk(conn, "data: [DONE]\n\n") + conn + end + # Self-contained admin page: lists agents and lets you add/remove them via # the admin API. Uses the ADMIN_API_KEY (entered in the page) as the Bearer # token for the /admin/agents calls. diff --git a/test/router_test.exs b/test/router_test.exs index 82bc7ed765..5ef30dcebe 100644 --- a/test/router_test.exs +++ b/test/router_test.exs @@ -6,11 +6,46 @@ defmodule N8nOpenaiAdapter.RouterTest do alias N8nOpenaiAdapter.AgentRegistry alias N8nOpenaiAdapter.Router - setup do + # A minimal local HTTP server that answers n8n-style chat webhooks with a + # fixed {output: ...} body, so we can exercise the adapter's n8n call path + # (including the SSE streaming route) without a real n8n instance. + defmodule EchoServer do + use Plug.Router + + plug(:match) + plug(:dispatch) + + post "/chat" do + conn + |> put_resp_content_type("application/json") + |> send_resp(200, Jason.encode!(%{"output" => "Echo reply"})) + end + + match _ do + send_resp(conn, 404, "not found") + end + end + + setup_all do + port = get_free_port() + {:ok, _} = Bandit.start_link(plug: EchoServer, port: port, scheme: :http) + {:ok, port: port} + end + + defp get_free_port do + {:ok, socket} = :gen_tcp.listen(0, [:binary, packet: :raw, active: false, reuseaddr: true]) + {:ok, port} = :inet.port(socket) + :gen_tcp.close(socket) + port + end + + setup %{port: port} do # The app starts with an empty store (no AGENTS env seeding). Seed the two # test agents via the registry so the base tests have something to list. AgentRegistry.put("scholar-agent", "https://n8n.bueso.eu/webhook/scholar-id/chat") AgentRegistry.put("media-agent", "https://n8n.bueso.eu/webhook/media-id/chat") + # A third agent pointing at the local echo server, for the streaming path. + AgentRegistry.put("echo-agent", "http://127.0.0.1:#{port}/chat") :ok end @@ -72,6 +107,57 @@ defmodule N8nOpenaiAdapter.RouterTest do assert conn.status == 400 end + test "POST /v1/chat/completions with stream=true emits SSE JSON chunks + [DONE]" do + conn = + conn( + :post, + "/v1/chat/completions", + Jason.encode!(%{ + "model" => "echo-agent", + "stream" => true, + "messages" => [%{"role" => "user", "content" => "hi"}] + }) + ) + |> put_req_header("authorization", "Bearer test-key") + |> put_req_header("content-type", "application/json") + |> Router.call(Router.init([])) + + assert conn.status == 200, "expected 200, got #{conn.status}: #{conn.resp_body}" + body = conn.resp_body + assert {"content-type", ct} = List.keyfind(conn.resp_headers, "content-type", 0) + assert String.starts_with?(ct, "text/event-stream") + # Each SSE event is `data: \n\n`; split and check the JSON is valid + events = String.split(body, "\n\n") |> Enum.reject(&(&1 == "")) + + data_events = Enum.filter(events, &String.starts_with?(&1, "data: ")) + assert Enum.any?(data_events, &String.ends_with?(&1, "[DONE]")) + + # Parse a non-DONE chunk and require the OpenAI chunk shape + chunk_event = Enum.find(data_events, &(not String.ends_with?(&1, "[DONE]"))) + chunk = Jason.decode!(String.replace_prefix(chunk_event, "data: ", "")) + assert chunk["object"] == "chat.completion.chunk" + assert is_list(chunk["choices"]) + # The content delta chunk carries the reply text + content_chunks = + Enum.filter(data_events, fn ev -> + decoded = + case Jason.decode(String.replace_prefix(ev, "data: ", "")) do + {:ok, j} -> j + _ -> :error + end + + case decoded do + %{"choices" => [choice]} -> + match?(%{"content" => c} when is_binary(c), Map.get(choice, "delta")) + + _ -> + false + end + end) + + assert length(content_chunks) == 1, "content_chunks=#{length(content_chunks)}" + end + test "GET /admin serves the web admin page (no auth needed to load the form)" do conn = conn(:get, "/admin") |> Router.call(Router.init([]))