Compare commits

...
2 Commits
Author SHA1 Message Date
hermes-agent c345dab9c1 Merge pull request 'adapter: support OpenAI streaming (SSE) in /v1/chat/completions' (#1) from feat/adapter-streaming into main 2026-09-12 12:51:34 +02:00
hermes-agent 2d46eda619 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.
2026-09-12 10:50:47 +00:00
2 changed files with 180 additions and 19 deletions
+93 -18
View File
@@ -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.
+87 -1
View File
@@ -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: <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
conn = conn(:get, "/admin") |> Router.call(Router.init([]))