adapter: support OpenAI streaming (SSE) in /v1/chat/completions
OpenCode's @ai-sdk/openai-compatible sends stream:true and would render an empty response because the adapter always returned a single non-streaming chat.completion JSON body. Now when stream:true, emit OpenAI-compatible SSE chat.completion.chunk events (role, content, [DONE]) so streaming clients render text. Non-streaming path unchanged. Adds a local EchoServer test that exercises the streaming route end-to-end.
This commit is contained in:
@@ -235,6 +235,7 @@ defmodule N8nOpenaiAdapter.Router do
|
|||||||
|
|
||||||
model = Map.get(data, "model")
|
model = Map.get(data, "model")
|
||||||
messages = Map.get(data, "messages", [])
|
messages = Map.get(data, "messages", [])
|
||||||
|
stream? = Map.get(data, "stream", false)
|
||||||
thread_id = Map.get(data, "thread_id")
|
thread_id = Map.get(data, "thread_id")
|
||||||
session_id = if is_binary(thread_id) and thread_id != "", do: thread_id, else: "default"
|
session_id = if is_binary(thread_id) and thread_id != "", do: thread_id, else: "default"
|
||||||
|
|
||||||
@@ -258,32 +259,106 @@ defmodule N8nOpenaiAdapter.Router do
|
|||||||
webhook ->
|
webhook ->
|
||||||
case call_n8n(webhook, session_id, chat_input) do
|
case call_n8n(webhook, session_id, chat_input) do
|
||||||
{:ok, output} ->
|
{:ok, output} ->
|
||||||
resp = %{
|
if stream? do
|
||||||
id:
|
stream_completion(conn, model, output)
|
||||||
"chatcmpl-#{Base.encode16(:crypto.strong_rand_bytes(12), case: :lower)}",
|
else
|
||||||
object: "chat.completion",
|
json(conn, 200, completion_body(model, output))
|
||||||
created: System.system_time(:second),
|
end
|
||||||
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)
|
|
||||||
|
|
||||||
{:error, msg} ->
|
{:error, msg} ->
|
||||||
error(conn, 502, msg)
|
if stream? do
|
||||||
|
stream_error(conn, msg)
|
||||||
|
else
|
||||||
|
error(conn, 502, msg)
|
||||||
|
end
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
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
|
# 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
|
# the admin API. Uses the ADMIN_API_KEY (entered in the page) as the Bearer
|
||||||
# token for the /admin/agents calls.
|
# token for the /admin/agents calls.
|
||||||
|
|||||||
+87
-1
@@ -6,11 +6,46 @@ defmodule N8nOpenaiAdapter.RouterTest do
|
|||||||
alias N8nOpenaiAdapter.AgentRegistry
|
alias N8nOpenaiAdapter.AgentRegistry
|
||||||
alias N8nOpenaiAdapter.Router
|
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
|
# 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.
|
# 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("scholar-agent", "https://n8n.bueso.eu/webhook/scholar-id/chat")
|
||||||
AgentRegistry.put("media-agent", "https://n8n.bueso.eu/webhook/media-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
|
:ok
|
||||||
end
|
end
|
||||||
|
|
||||||
@@ -72,6 +107,57 @@ defmodule N8nOpenaiAdapter.RouterTest do
|
|||||||
assert conn.status == 400
|
assert conn.status == 400
|
||||||
end
|
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: <json>\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
|
test "GET /admin serves the web admin page (no auth needed to load the form)" do
|
||||||
conn = conn(:get, "/admin") |> Router.call(Router.init([]))
|
conn = conn(:get, "/admin") |> Router.call(Router.init([]))
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user