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.
469 lines
15 KiB
Elixir
469 lines
15 KiB
Elixir
defmodule N8nOpenaiAdapter.Router do
|
|
@moduledoc """
|
|
OpenAI-compatible HTTP surface for n8n chat agents.
|
|
|
|
GET /v1/models -> lists configured agents (OpenAI shape)
|
|
POST /v1/chat/completions -> {model, messages} -> forwards last user msg
|
|
to the agent's n8n webhook, returns an
|
|
OpenAI chat.completion response.
|
|
|
|
Auth: clients send `Authorization: Bearer <ADAPTER_API_KEY>`. The webhook
|
|
call to n8n can carry its own basic auth via `CHAT_WEBHOOK_BASIC="user:pass"`.
|
|
|
|
Plug order matters: `match` -> JSON parser (consumes body) -> `dispatch`.
|
|
Parsed JSON lands in `conn.body_params`, so handlers read body_params, not
|
|
a manual body read.
|
|
"""
|
|
use Plug.Router
|
|
|
|
alias N8nOpenaiAdapter.AgentRegistry
|
|
|
|
plug(:match)
|
|
plug(Plug.Parsers, parsers: [:json], json_decoder: Jason)
|
|
plug(:dispatch)
|
|
|
|
# --- Helpers ---
|
|
|
|
# Ensure the request carries a valid ADAPTER_API_KEY bearer token.
|
|
defp authorize!(conn) do
|
|
expected = System.get_env("ADAPTER_API_KEY")
|
|
|
|
case get_req_header(conn, "authorization") do
|
|
[auth] ->
|
|
token =
|
|
case Regex.run(~r/^Bearer\s+(.+)$/i, auth) do
|
|
[_, t] -> t
|
|
_ -> auth
|
|
end
|
|
|
|
if expected && Plug.Crypto.secure_compare(token, expected) do
|
|
conn
|
|
else
|
|
{:halt, send_resp(conn, 401, Jason.encode!(%{error: %{message: "Invalid API key"}}))}
|
|
end
|
|
|
|
_ ->
|
|
{:halt,
|
|
send_resp(conn, 401, Jason.encode!(%{error: %{message: "Missing Authorization header"}}))}
|
|
end
|
|
end
|
|
|
|
# Admin endpoints use a SEPARATE ADMIN_API_KEY so you can grant agent
|
|
# management without exposing the OpenAI-facing key.
|
|
defp authorize_admin!(conn) do
|
|
expected = System.get_env("ADMIN_API_KEY")
|
|
|
|
case get_req_header(conn, "authorization") do
|
|
[auth] ->
|
|
token =
|
|
case Regex.run(~r/^Bearer\s+(.+)$/i, auth) do
|
|
[_, t] -> t
|
|
_ -> auth
|
|
end
|
|
|
|
if expected && Plug.Crypto.secure_compare(token, expected) do
|
|
conn
|
|
else
|
|
{:halt, send_resp(conn, 401, Jason.encode!(%{error: %{message: "Invalid admin key"}}))}
|
|
end
|
|
|
|
_ ->
|
|
{:halt,
|
|
send_resp(conn, 401, Jason.encode!(%{error: %{message: "Missing Authorization header"}}))}
|
|
end
|
|
end
|
|
|
|
defp json(conn, status, body) do
|
|
conn
|
|
|> put_resp_content_type("application/json")
|
|
|> send_resp(status, Jason.encode!(body))
|
|
end
|
|
|
|
defp error(conn, status, message) do
|
|
json(conn, status, %{error: %{message: message}})
|
|
end
|
|
|
|
# Pull the last user message content out of an OpenAI messages array.
|
|
defp last_user_text(messages) do
|
|
messages
|
|
|> Enum.reverse()
|
|
|> Enum.find_value(fn
|
|
%{"role" => "user", "content" => c} when is_binary(c) -> c
|
|
_ -> nil
|
|
end)
|
|
end
|
|
|
|
# Forward a message to the n8n chat webhook and return the assistant reply.
|
|
defp call_n8n(webhook, session_id, chat_input) do
|
|
body = %{sessionId: session_id, action: "sendMessage", chatInput: chat_input}
|
|
|
|
headers =
|
|
case System.get_env("CHAT_WEBHOOK_BASIC") do
|
|
nil ->
|
|
[{"content-type", "application/json"}]
|
|
|
|
basic ->
|
|
[
|
|
{"content-type", "application/json"},
|
|
{"authorization", "Basic " <> Base.encode64(basic)}
|
|
]
|
|
end
|
|
|
|
case Req.post(webhook, json: body, headers: headers, receive_timeout: 120_000) do
|
|
{:ok, %{status: 200, body: %{"output" => output}}} -> {:ok, output}
|
|
{:ok, %{status: s}} -> {:error, "n8n webhook returned HTTP #{s}"}
|
|
{:error, e} -> {:error, "n8n webhook error: #{Exception.message(e)}"}
|
|
end
|
|
end
|
|
|
|
# --- Routes ---
|
|
|
|
get "/v1/models" do
|
|
case authorize!(conn) do
|
|
{:halt, conn} ->
|
|
conn
|
|
|
|
conn ->
|
|
models =
|
|
Enum.map(AgentRegistry.model_names(), fn name ->
|
|
%{id: name, object: "model", created: 0, owned_by: "n8n"}
|
|
end)
|
|
|
|
json(conn, 200, %{object: "list", data: models})
|
|
end
|
|
end
|
|
|
|
post "/v1/chat/completions" do
|
|
case authorize!(conn) do
|
|
{:halt, conn} ->
|
|
conn
|
|
|
|
conn ->
|
|
handle_chat(conn)
|
|
end
|
|
end
|
|
|
|
# --- Admin API (separate ADMIN_API_KEY) ---
|
|
|
|
# The admin page itself is served WITHOUT the Bearer auth gate — it's a static
|
|
# form. The user enters their ADMIN_API_KEY in the page, and that key drives
|
|
# the auth'd /admin/agents CRUD calls. Serving the page openly is harmless (it
|
|
# exposes no data; the agent list is only fetched with a valid key).
|
|
get "/admin" do
|
|
conn
|
|
|> put_resp_content_type("text/html")
|
|
|> send_resp(200, admin_page())
|
|
end
|
|
|
|
get "/admin/agents" do
|
|
case authorize_admin!(conn) do
|
|
{:halt, conn} ->
|
|
conn
|
|
|
|
conn ->
|
|
json(conn, 200, %{agents: AgentRegistry.all()})
|
|
end
|
|
end
|
|
|
|
post "/admin/agents" do
|
|
case authorize_admin!(conn) do
|
|
{:halt, conn} ->
|
|
conn
|
|
|
|
conn ->
|
|
data = conn.body_params
|
|
model = Map.get(data, "model")
|
|
webhook = Map.get(data, "webhook")
|
|
|
|
cond do
|
|
not is_binary(model) or model == "" ->
|
|
error(conn, 400, "Bad request: missing or invalid \"model\"")
|
|
|
|
not is_binary(webhook) or webhook == "" ->
|
|
error(conn, 400, "Bad request: missing or invalid \"webhook\"")
|
|
|
|
true ->
|
|
AgentRegistry.put(model, webhook)
|
|
json(conn, 200, %{ok: true, model: model, webhook: webhook})
|
|
end
|
|
end
|
|
end
|
|
|
|
delete "/admin/agents/:model" do
|
|
case authorize_admin!(conn) do
|
|
{:halt, conn} ->
|
|
conn
|
|
|
|
conn ->
|
|
case AgentRegistry.delete(conn.params["model"]) do
|
|
:ok -> json(conn, 200, %{ok: true})
|
|
:not_found -> error(conn, 404, "Unknown model")
|
|
end
|
|
end
|
|
end
|
|
|
|
match _ do
|
|
case authorize!(conn) do
|
|
{:halt, conn} ->
|
|
conn
|
|
|
|
conn ->
|
|
json(conn, 404, %{error: %{message: "Not found"}})
|
|
end
|
|
end
|
|
|
|
# Temporary diagnostics: log the raw chat request (stream flag + headers) so
|
|
# we can see exactly what SwiftChat sends. Remove after diagnosis.
|
|
defp log_request(conn, data) do
|
|
stream = Map.get(data, "stream")
|
|
accept = get_req_header(conn, "accept") |> List.first()
|
|
content_type = get_req_header(conn, "content-type") |> List.first()
|
|
model = Map.get(data, "model")
|
|
|
|
IO.puts(
|
|
"CHAT_REQ model=#{inspect(model)} stream=#{inspect(stream)} " <>
|
|
"accept=#{inspect(accept)} content_type=#{inspect(content_type)} " <>
|
|
"body=#{Jason.encode!(data) |> String.slice(0, 500)}"
|
|
)
|
|
end
|
|
|
|
defp handle_chat(conn) do
|
|
data = conn.body_params
|
|
# Temporary diagnostics: log the raw chat request so we can see exactly what
|
|
# SwiftChat sends (stream flag, headers). Remove after diagnosis.
|
|
log_request(conn, data)
|
|
|
|
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"
|
|
|
|
cond do
|
|
not is_binary(model) or model == "" ->
|
|
error(conn, 400, "Bad request: missing or invalid \"model\"")
|
|
|
|
not is_list(messages) ->
|
|
error(conn, 400, "Bad request: expected \"messages\" array")
|
|
|
|
true ->
|
|
case last_user_text(messages) do
|
|
nil ->
|
|
error(conn, 400, "No user message in request")
|
|
|
|
chat_input ->
|
|
case AgentRegistry.webhook_for(model) do
|
|
:error ->
|
|
error(conn, 400, "Unknown model: #{model}")
|
|
|
|
webhook ->
|
|
case call_n8n(webhook, session_id, chat_input) do
|
|
{:ok, output} ->
|
|
if stream? do
|
|
stream_completion(conn, model, output)
|
|
else
|
|
json(conn, 200, completion_body(model, output))
|
|
end
|
|
|
|
{:error, 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.
|
|
defp admin_page do
|
|
"""
|
|
<!doctype html>
|
|
<html lang="en">
|
|
<head>
|
|
<meta charset="utf-8">
|
|
<meta name="viewport" content="width=device-width, initial-scale=1">
|
|
<title>n8n OpenAI Adapter — Agents</title>
|
|
<style>
|
|
body { font-family: system-ui, sans-serif; max-width: 720px; margin: 2rem auto; padding: 0 1rem; color: #1a1a1a; }
|
|
h1 { font-size: 1.4rem; }
|
|
input[type=text], input[type=password] { width: 100%; padding: .5rem; margin: .25rem 0 .75rem; box-sizing: border-box; }
|
|
button { padding: .5rem 1rem; cursor: pointer; }
|
|
table { width: 100%; border-collapse: collapse; margin-top: 1rem; }
|
|
th, td { text-align: left; padding: .5rem; border-bottom: 1px solid #ddd; }
|
|
.msg { margin-top: 1rem; padding: .5rem; border-radius: 4px; }
|
|
.ok { background: #e6f4ea; color: #1e7e34; }
|
|
.err { background: #fdecea; color: #c62828; }
|
|
.del { color: #c62828; background: none; border: none; cursor: pointer; }
|
|
</style>
|
|
</head>
|
|
<body>
|
|
<h1>n8n OpenAI Adapter — Agents</h1>
|
|
<p>Manage which n8n chat agents are exposed as OpenAI models.</p>
|
|
|
|
<label for="key">Admin API key</label>
|
|
<input type="password" id="key" placeholder="ADMIN_API_KEY" autocomplete="off">
|
|
<button onclick="loadAgents()">Load agents</button>
|
|
|
|
<h2>Add / update agent</h2>
|
|
<label for="model">Model name</label>
|
|
<input type="text" id="model" placeholder="e.g. scholar-agent">
|
|
<label for="webhook">n8n chat webhook URL</label>
|
|
<input type="text" id="webhook" placeholder="https://n8n.bueso.eu/webhook/<id>/chat">
|
|
<button onclick="saveAgent()">Save agent</button>
|
|
|
|
<div id="msg"></div>
|
|
|
|
<h2>Configured agents</h2>
|
|
<table id="agents"><thead><tr><th>Model</th><th>Webhook</th><th></th></tr></thead><tbody></tbody></table>
|
|
|
|
<script>
|
|
const key = () => document.getElementById('key').value.trim();
|
|
const msg = (text, ok) => {
|
|
const el = document.getElementById('msg');
|
|
el.className = 'msg ' + (ok ? 'ok' : 'err');
|
|
el.textContent = text;
|
|
};
|
|
const auth = () => ({ 'Authorization': 'Bearer ' + key(), 'Content-Type': 'application/json' });
|
|
|
|
async function loadAgents() {
|
|
if (!key()) return msg('Enter the admin API key first', false);
|
|
try {
|
|
const r = await fetch('/admin/agents', { headers: auth() });
|
|
if (!r.ok) return msg('Failed to load agents: HTTP ' + r.status, false);
|
|
const data = await r.json();
|
|
const tbody = document.querySelector('#agents tbody');
|
|
tbody.innerHTML = '';
|
|
for (const [model, webhook] of Object.entries(data.agents)) {
|
|
const tr = document.createElement('tr');
|
|
tr.innerHTML = '<td>' + model + '</td><td>' + webhook + '</td>' +
|
|
'<td><button class="del" onclick="deleteAgent(\\'' + model + '\\')">Delete</button></td>';
|
|
tbody.appendChild(tr);
|
|
}
|
|
msg('Loaded ' + Object.keys(data.agents).length + ' agent(s)', true);
|
|
} catch (e) { msg('Error: ' + e, false); }
|
|
}
|
|
|
|
async function saveAgent() {
|
|
const model = document.getElementById('model').value.trim();
|
|
const webhook = document.getElementById('webhook').value.trim();
|
|
if (!key()) return msg('Enter the admin API key first', false);
|
|
if (!model || !webhook) return msg('Model and webhook are required', false);
|
|
try {
|
|
const r = await fetch('/admin/agents', {
|
|
method: 'POST', headers: auth(),
|
|
body: JSON.stringify({ model, webhook })
|
|
});
|
|
if (!r.ok) return msg('Failed to save agent: HTTP ' + r.status, false);
|
|
document.getElementById('model').value = '';
|
|
document.getElementById('webhook').value = '';
|
|
msg('Saved agent "' + model + '"', true);
|
|
loadAgents();
|
|
} catch (e) { msg('Error: ' + e, false); }
|
|
}
|
|
|
|
async function deleteAgent(model) {
|
|
if (!key()) return msg('Enter the admin API key first', false);
|
|
if (!confirm('Delete agent "' + model + '"?')) return;
|
|
try {
|
|
const r = await fetch('/admin/agents/' + encodeURIComponent(model), {
|
|
method: 'DELETE', headers: auth()
|
|
});
|
|
if (!r.ok) return msg('Failed to delete: HTTP ' + r.status, false);
|
|
msg('Deleted agent "' + model + '"', true);
|
|
loadAgents();
|
|
} catch (e) { msg('Error: ' + e, false); }
|
|
}
|
|
</script>
|
|
</body>
|
|
</html>
|
|
"""
|
|
end
|
|
end
|