@aiconnect/flue-openai-adapter
v0.2.0
Published
OpenAI-compatible /v1/chat/completions route for Flue apps (Hono)
Maintainers
Readme
@aiconnect/flue-openai-adapter
OpenAI-compatible POST /v1/chat/completions route for Flue apps, as a mountable Hono sub-app. Supports streaming (SSE chunks) and non-streaming responses, with usage mapping and honest error signaling.
Install
npm install @aiconnect/flue-openai-adapterhono >= 4 is a peer dependency. @flue/runtime is not imported — drivers talk to Flue over its HTTP/SSE surface, so one package works across Flue majors.
Usage (Flue 0.9.x workflow app)
import { flue } from '@flue/runtime/routing';
import { Hono } from 'hono';
import { openaiCompat, workflowV0 } from '@aiconnect/flue-openai-adapter';
const app = new Hono();
app.route('/v1', openaiCompat({
driver: workflowV0({
workflow: 'my-workflow',
app: () => app,
}),
}));
app.route('/', flue());
export default app;The driver self-dispatches POST /workflows/<name> through the app (no network hop) and relays the run's text_delta frames as OpenAI chunks.
Workflow contract (streaming: 'deltas', default)
The target workflow must:
- accept
stream: truein its payload and synthesize plain text in that mode (structured output would stream raw JSON into the deltas); - return
{ answer?: string, usage?: { input, output, totalTokens } }in its result, used by shortcut paths that produce no deltas (cache hits, static responses).
Workflows with structured output (streaming: 'result-only')
workflowV0({ workflow: 'my-workflow', app: () => app, streaming: 'result-only' })No stream flag is sent, deltas are ignored, and the answer comes exclusively from the final run result — streaming clients receive it as a single chunk.
Usage (Flue 2.x agent app)
import { openaiCompat, agentV2 } from '@aiconnect/flue-openai-adapter';
app.route('/v1', openaiCompat({
driver: agentV2({ agent: 'host', app: () => app }),
}));Flue 2 prompts are fire-and-forget: the driver POSTs { kind: 'user', body } to /agents/<name>/<id> (202 + stream coordinates), reads the conversation stream (GET ?view=updates&live=sse), relays message-delta items with kind: "text" for the admitted submission (reasoning deltas are dropped), and aborts the stream once submission-settled arrives. A non-completed outcome (failed, aborted) is reported as an error.
By default each request runs in a fresh conversation (random UUID). Pass conversation to keep server-side history:
agentV2({ agent: 'host', app: () => app, conversation: (_q, req) => String(req.user ?? 'default') })Multi-agent mounts (drivers)
One /v1 endpoint can front several agents: pass a drivers map and request.model
selects the driver by key. GET /v1/models lists the keys (OpenAI list format), so
model-picker UIs (Open WebUI, connect chat) discover the agents automatically.
app.route('/v1', openaiCompat({
drivers: {
sarah: agentV2({ agent: 'sarah', app: () => app }),
hello: agentV2({ agent: 'hello', app: () => app }),
},
defaultModel: 'hello', // lookup key when the request omits `model`
}));An unknown model gets a 404 with code: 'model_not_found'. When the request omits
model, defaultModel is used as the lookup key; if it is not a key and a single
driver is also set, that driver handles the request (0.1.0 fallback). Single-driver
mounts (driver) behave exactly as in 0.1.0, and GET /models then lists just
defaultModel.
Options
openaiCompat(options)
| Option | Default | Purpose |
|---|---|---|
| driver | — | Single transport adapter; required unless drivers is set |
| drivers | — | model → driver map; keys become GET /models entries |
| resolveAnswer(result) | result.answer ?? '' | Assistant content for runs without deltas |
| extractQuestion(messages) | last user message with string content | Question sent to the driver |
| defaultModel | 'flue' | Model echoed when the request omits model; with drivers, also its lookup key |
workflowV0(options)
| Option | Default | Purpose |
|---|---|---|
| workflow | — | Workflow name (/workflows/<name>) |
| app | — | Lazy ref to the Hono app for self-dispatch |
| streaming | 'deltas' | 'deltas' | 'result-only' |
| buildPayload(question, request) | { question, stream? } | Custom workflow payload |
| baseUrl | http://flue.internal | Internal dispatch base URL |
agentV2(options)
| Option | Default | Purpose |
|---|---|---|
| agent | — | Agent route mount name (/agents/<name>) |
| app | global fetch | Lazy ref to the Hono app; omit to reach an external server via baseUrl |
| baseUrl | http://flue.internal | Base URL for admission/stream requests |
| conversation(question, request) | random UUID | Conversation instance id (stable id keeps history) |
Error semantics
- Driver invocation failure → HTTP 502 with
{ error: { type: 'upstream_error' } }. - Errored run, non-streaming → HTTP 502, even if partial text accumulated.
- Errored run, streaming → a top-level
data: { "error": ... }SSE payload (the OpenAI wire convention;finish_reasonis never"stop"on an errored run), thendata: [DONE]. - A stream that closes without its terminal frame (
run_end/submission-settled) is treated as truncated and reported as an error, never as a clean completion. - Client disconnect on a streaming response aborts the underlying Flue run (
cancel()→ driverabort()).
Usage accounting
Non-streaming responses carry usage mapped from the driver result. Streaming responses include a final usage chunk (empty choices) only when the client sends stream_options: { include_usage: true }, matching OpenAI behavior.
Custom drivers
Implement the Driver interface to adapt other transports:
interface Driver {
run(question: string, request: ChatCompletionRequest): Promise<AsyncIterable<DriverEvent>>;
}
// DriverEvent = { type: 'delta', text } | { type: 'end', result, errored }Publishing
npm version patch # ou minor / major — cria o commit e a tag
npm run release # prepublishOnly (types + testes + build) → publish → push da tagRoadmap
workflowV0is transitional: it will be removed in 1.0 once consumers migrate to Flue 2.usagemapping foragentV2(the Flue 2 conversation stream does not expose token usage; reported as zeros).
Limitations
- Only the last
usermessage with text is used (override withextractQuestion); conversation history is ignored. Array-of-parts content is supported — itstextparts are joined; image/audio parts are ignored. - One choice per response (
nis not supported).
