Merge pull request 'adapter: support OpenAI streaming (SSE) in /v1/chat/completions' (#1) from feat/adapter-streaming into main
This commit was merged in pull request #1.
This commit is contained in:
@@ -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
@@ -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([]))
|
||||
|
||||
|
||||
Reference in New Issue
Block a user