Repository navigation
Expand file tree
/
Copy pathllms.txt
More file actions
640 lines (478 loc) · 22.3 KB
/
Copy pathllms.txt
File metadata and controls
640 lines (478 loc) · 22.3 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
# Pipeflow — AI Reference Guide
Pipeflow is an open-source TypeScript backend SDK for building voice agents, conversational applications, meeting transcribers, Discord bots, recruiters, assistants, and other realtime audio experiences. Published as `@moureau/pipeflow`.
It handles the plumbing between audio, speech-to-text, LLMs, text-to-speech, conversations, tools, and persistence — while keeping your application in control.
**Pipeflow is the pipe. You build what flows through it.**
**Realtime by default** — audio, transcripts, and speech stream continuously, with built-in interruption and barge-in handling.
**Provider-agnostic** — STT, LLM, and TTS are swappable adapters (Deepgram, DeepSeek, OpenRouter, Kokoro, OpenAI, Claude).
**Your backend stays yours** — tools and the audio transport are owned by your application; Pipeflow never executes your code.
---
## Core Concepts
Three independent concepts:
* **Pipeflow** — the orchestrator that wires STT, LLM, and TTS together. Created with providers.
* **Agent** — intelligence, context, and tools. Can run standalone (`agent.run()`) or participate in a conversation.
* **Conversation** — a persistent realtime conversation with participants, turns, audio, and interruption handling.
* **Tool** — application capabilities executed by your backend. `PipeflowTool` with Zod `schema.in` for model contract and `execute()` callback.
* **Provider** — an implementation of STT, LLM, or TTS (Deepgram, DeepSeek, OpenRouter, OpenAI, Claude, Kokoro).
* **Logger** — pluggable logging. Pass `logger` or `verbose: true` to `Pipeflow`. `ConsoleLogger` and `SilentLogger` exported from the main entry.
An agent does not inherently own a conversation. A conversation does not require an agent. A provider should not leak into your application logic.
---
## Installation
```bash
bun add @moureau/pipeflow
```
Published on npm as `@moureau/pipeflow`.
---
## Basic Voice Agent
```ts
import { Pipeflow } from "@moureau/pipeflow";
import { DeepSeekLLM, DeepgramSTT, KokoroTTS } from "@moureau/pipeflow/providers";
const stt = new DeepgramSTT({ apiKey: process.env.DEEPGRAM_API_KEY });
const tts = new KokoroTTS();
const llm = new DeepSeekLLM({ apiKey: process.env.DEEPSEEK_API_KEY });
const pipeflow = new Pipeflow({ llm, stt, tts });
const jarvis = pipeflow.agent({
name: "Jarvis",
context: `
You are Jarvis, a helpful voice assistant.
Keep your responses concise and conversational.
`,
});
const conversation = await pipeflow.conversations.create({
agents: [jarvis],
});
await conversation.start();
// Feed audio as it arrives:
voice.onAudio((audio) => {
conversation.listen({ userId: "alice", audio });
});
// Listen for generated audio:
conversation.on("audio", ({ audio }) => {
voice.play(audio);
});
```
---
## Imports
```ts
// Core
import { Pipeflow, PipeflowTool, ConsoleLogger } from "@moureau/pipeflow";
// Browser/Node client
import { PipeflowClient, WebSocketProtocol } from "@moureau/pipeflow/client";
// Providers
import { DeepSeekLLM, OpenRouterLLM, OpenAILLM, ClaudeLLM } from "@moureau/pipeflow/providers";
import { DeepgramSTT } from "@moureau/pipeflow/providers";
import { KokoroTTS } from "@moureau/pipeflow/providers";
// Persistence
import { SQLitePersistence } from "@moureau/pipeflow/persistence";
// Transport (server adapters)
import { BunServerAdapter, ConversationWebSocketServer } from "@moureau/pipeflow/transport";
// Bring your own: implement ServerAdapter for socket.io, node-ws, etc.
// Conversations (types, subpaths)
import type { Conversation, ConversationEvent } from "@moureau/pipeflow/conversations";
```
---
## Logging
Pipeflow has a pluggable logger interface injected at construction:
```ts
import { Pipeflow, ConsoleLogger } from "@moureau/pipeflow";
// Shortcut — prints [pipeflow] prefix messages
const pipeflow = new Pipeflow({ llm, verbose: true });
// Or with a custom logger:
import type { Logger } from "@moureau/pipeflow";
const pipeflow = new Pipeflow({
llm,
logger: new ConsoleLogger("myapp"),
});
```
The `Logger` interface has `info`, `warn`, `error`, `debug` methods. The default is a no-op `SilentLogger`.
---
## Pipeline
### Conversation Lifecycle
```ts
create() // creates the persistent conversation
## Conversations API
Conversations are created through `pipeflow.conversations`:
```ts
// Create (options: agents, createdBy, autoExecuteTools,
// maxConcurrentTtsRequests, audioReorderMs, coordinations,
// retryDirectAnswer)
const conversation = await pipeflow.conversations.create({ agents: [jarvis] });
// Retrieve by id. Agent instances and per-conversation tuning are not
// persisted, so pass them back to restore them (options: agents,
// coordinations, retryDirectAnswer, maxConcurrentTtsRequests,
// audioReorderMs). stt, tts, logger, and the autoExecuteTools default are
// inherited from the Conversations instance.
const conv = await pipeflow.conversations.get(id, { agents: [jarvis] });
// List with filters, ordering, and pagination
const { conversations, total } = await pipeflow.conversations.list({
userId: "alice",
status: "active", // "active" | "ended"
archived: false, // exclude archived (default)
orderBy: "createdAt", // "createdAt" | "endedAt"
orderDir: "desc",
page: 1,
pageSize: 100,
});
// Soft-delete (hides from list, still retrievable via get())
await pipeflow.conversations.archive(id);
// Permanently delete
await pipeflow.conversations.delete(id);
// Retrieve transcript
const transcript = await pipeflow.conversations.transcript(id);
```
### Conversation Lifecycle
```ts
start() // moves the conversation into the started state (attaches orchestrator)
participate() // adds participants
listen() // sends audio (synchronous — feed audio packets)
send() // injects a finalized text turn (no STT)
stop() // finalizes the realtime session
interrupt() // stops TTS, cancels current generation (barge-in)
```
`listen()` is intentionally synchronous — it means "send this audio packet", not "wait for this utterance to finish":
```ts
conversation.listen({ userId, audio });
```
Pass the sender's capture-order `sequence` to enable netcode-style reordering: chunks arriving ahead of a gap are held for `audioReorderMs` (default 100, set on `create()`) so a packet still in flight can fill it, then released in order before STT. In-order audio (the common case) is emitted immediately. Late and duplicate sequences are dropped.
```ts
conversation.listen({ userId, audio, sequence: wire.sequence });
```
`send({ userId, text })` does the same for a finalized text turn, bypassing STT.
### Participants
```ts
await conversation.participate([
{ userId: "alice" },
{ userId: "bob", aliases: ["robert", "rob"] },
{ userId: "charlie", aliases: ["charles"] },
]);
```
### Interruption
`interrupt()` is a semantic guarantee: it stops the active generation and prevents any further generated audio from reaching the conversation (tested to stop audio within ~100ms). The next turn starts fresh.
```ts
conversation.interrupt();
```
The orchestrator also detects automatic barge-in when a participant starts speaking while the agent is responding.
---
## Conversation Events
```ts
conversation.on("audio-in", ({ userId, audio }) => {});
conversation.on("text-in", ({ userId, text }) => {});
conversation.on("partial-transcript", ({ userId, text }) => {}); // live STT partials (captions)
conversation.on("turn", ({ userId, text }) => {}); // finalized participant turn
conversation.on("transcript", ({ userId, text }) => {}); // transcript entry
conversation.on("audio", ({ userId, audio }) => {}); // generated audio to play
conversation.on("generation", ({ userId, text }) => {}); // agent generation
conversation.on("generation-complete", ({ generation }) => {}); // top-level generation finished
conversation.on("text-delta", ({ text, agentName }) => {}); // streaming LLM deltas
conversation.on("agent-delta", ({ agent, text }) => {}); // delegated specialist progress
conversation.on("tool-call", ({ call }) => {}); // every tool round-trip
conversation.on("tool-call-result", ({ callId, result }) => {}); // tool execution result
conversation.on("plan", ({ steps }) => {}); // coordinator produced a plan
conversation.on("plan-step", ({ step, status, text }) => {}); // step started/completed/failed
conversation.on("interrupt", () => {});
conversation.on("error", ({ error }) => {});
conversation.on("start", ({}) => {});
conversation.on("stop", ({}) => {});
conversation.on("state", ({ state }) => {});
```
---
## Agents
An agent defines an AI persona and its capabilities.
```ts
const jarvis = pipeflow.agent({
name: "Jarvis",
context: `
You are Jarvis.
You are concise, helpful, and conversational.
`,
tools: [getWeather, searchCalendar],
});
```
`context` accepts either a static string or a `ContextFn` — a function called
each time the agent generates a response. The function receives the triggering
prompt and any available conversation context and may return a string or a
promise:
```ts
pipeflow.agent({
name: "DocBot",
context: ({ prompt, conversationId, participants, annotations }) =>
`You are helping in conversation ${conversationId ?? "standalone"}. ` +
`The user said: "${prompt}".`,
});
// Async variant:
pipeflow.agent({
name: "AsyncBot",
context: async ({ prompt, annotations }) => {
const config = await loadConfig();
return `Context: ${config.mode}. Prompt: ${prompt}`;
},
});
```
Dynamic app state can be injected mid-conversation via
`conversation.annotations` (an `Annotations` instance). Values are rendered
as system messages before every generation:
```ts
// When the user navigates to a file:
conversation.annotations.set("currentFile", "src/app.ts");
conversation.annotations.set({ userRole: "admin", mode: "review" });
// Remove when no longer relevant:
conversation.annotations.delete("currentFile");
conversation.annotations.clear();
// In ContextFn:
context: ({ prompt, annotations }) => {
const file = annotations.get("currentFile") ?? "unknown";
return `You are reviewing ${file}. User: "${prompt}"`;
}
```
The `ContextFn` type and `ContextParams` interface are exported from
`@moureau/pipeflow`:
Agents can also be run independently of conversations:
```ts
const result = await jarvis.run({
prompt: "Explain how a neural network works.",
});
console.log(result.text);
```
### Multi-agent conversations
A conversation can coordinate **multiple agents**. A turn that explicitly addresses an agent by name or alias goes straight to that agent; unaddressed turns go through the built-in `understand` coordination, which decides whether to delegate, ask the user, or answer directly.
```ts
const conversation = await pipeflow.conversations.create({
agents: [receptionist, specialist],
});
```
### Collaborative orchestration (plan-based)
`understand` is a hardcoded **coordination**: a reasoning unit that outputs
a structured **plan** instead of re-delegating reactively. It produces one
`delegate` call with all agent steps, then the framework executes deterministically.
The coordinator can:
- **plan** — output a structured plan with one step per agent. Each step has
an `id`, optional `agent` (omit for a pure-LLM step), `prompt`, and
`dependsOn` (string ids referencing other steps). Independent steps run in
parallel; dependent steps receive prior outputs. The model speaks the final
answer in its narration before emitting the plan, so there is no extra
composition round-trip.
- **clarify** — ask the user for missing details via a structured, batched
`clarify` action (capped at 2 question rounds, then assumptions are stated).
- **complete** — answer directly for trivial requests that need no agents.
- **agents** — simpler syntax for models that can't handle the full plan
schema. Converted to a plan internally and executed the same way.
Agent delegation runs as **sub-generations** on their own LLM, context, and
tools (text-only, 3-tool-round max, 512-token output cap). Clarification
resumes on the next turn — participant speech during a clarify wait is
treated as the answer, not as barge-in. Plans can be observed via `plan` and
`plan-step` events; specialist progress streams via `agent-delta`.
Available as `retryDirectAnswer` (on `pipeflow.conversations.create()`,
`Conversation`, or `Pipeflow`) to retry once when multi-agent answers directly
without planning. Extra coordinations are registered on the same call via the
`coordinations` option.
---
## Tools
```ts
import { z } from "zod";
const getWeather = new PipeflowTool({
name: "get_weather",
description: "Get the current weather for a city.",
schema: {
in: z.object({ city: z.string().describe("The city to look up.") }),
},
execute: async ({ city }) => {
return weatherService.getCurrent(city);
},
});
```
`schema.in` is what the model sends; `schema.out` (optional, defaults to `in`) is what `execute` receives after validation.
### Auto-execute (default)
```ts
const jarvis = pipeflow.agent({
name: "Jarvis",
context: "You are a helpful assistant.",
tools: [getWeather],
});
const conversation = await pipeflow.conversations.create({ agents: [jarvis] });
// tool calls auto-execute — no handler needed
```
The `tool-call` and `tool-call-result` events still fire for observation. Thrown errors are returned to the model as `{ "error": ... }`. Unknown tool names report `Unknown tool "..."`. Concurrent tool calls run in parallel. Hung tools are bounded by `toolTimeoutMs` (default 30s).
### Opt out: manual tool execution
```ts
const conversation = await pipeflow.conversations.create({
agents: [jarvis],
autoExecuteTools: false,
});
conversation.on("tool-call", async ({ call }) => {
const result = await myBackend.executeTool(call.name, call.arguments);
conversation.resolveToolCall({ id: call.id, result });
});
```
---
## Server Transport (Bun-agnostic, pluggable)
Pipeflow bridges a `Conversation` to remote clients via a runtime-independent `ServerAdapter` interface. This works with Bun, Node `ws`, socket.io, or any other WebSocket library.
```ts
import { ConversationWebSocketServer, BunServerAdapter } from "@moureau/pipeflow/transport";
const adapter = new BunServerAdapter({ port: 3000 });
const server = new ConversationWebSocketServer({ conversation, adapter });
await server.start();
```
You can bring your own adapter — implement the `ServerAdapter` interface with `start()`, `stop()`, and `onConnection()`:
```ts
import type { ServerAdapter, ServerClient } from "@moureau/pipeflow/transport";
class MySocketIOAdapter implements ServerAdapter {
// implement start(), stop(), onConnection()
}
```
---
## Meeting Transcription
A conversation does not require an agent — Pipeflow works as a realtime transcription primitive.
```ts
const conversation = await pipeflow.conversations.create();
await conversation.start();
await conversation.participate([
{ userId: "alice" },
{ userId: "bob", aliases: ["robert"] },
]);
discordVoice.onAudio((userId, audio) => {
conversation.listen({ userId, audio });
});
await conversation.stop();
const transcript = await pipeflow.conversations.transcript(conversation.id);
```
---
## Meeting Summaries
```ts
const notetaker = pipeflow.agent({
name: "Meeting Notetaker",
context: `
You are a meeting notetaker.
Produce concise notes containing:
- summary
- decisions
- action items
- unresolved questions
`,
});
const transcript = await pipeflow.conversations.transcript(conversation.id);
const result = await notetaker.run({
prompt: `
The meeting transcription is:
${transcript.join("\n")}
Produce the meeting notes.
`,
});
```
---
## Browser Client
The `@moureau/pipeflow/client` subpath ships a typed browser/Node SDK for talking to a Pipeflow server: `PipeflowClient` mirrors the `Conversation` event surface over a swappable `Protocol` (a `WebSocketProtocol` is included).
```ts
import { PipeflowClient, WebSocketProtocol } from "@moureau/pipeflow/client";
const client = new PipeflowClient({
protocol: new WebSocketProtocol({
url: "wss://api.example.com/conversations/abc",
token: () => localStorage.getItem("jwt"),
}),
audio: { input: true },
reconnect: { maxAttempts: 5 },
userId: "alice",
});
client.on("turn", (e) => console.log(e.turn.text));
client.on("audio", (e) => client.playAudio(e.audio));
await client.connect();
await client.sendText("hello");
```
`Protocol` is a small contract: `connect`, `close`, `send`, `onMessage`, `onStatus`, plus a `status` getter — bring your own transport (SSE, WebRTC, in-memory tests) without touching application code.
---
## Providers
| Interface | Provider | Import path |
|-----------|----------|-------------|
| LLM | DeepSeek | `@moureau/pipeflow/providers` |
| LLM | OpenRouter | `@moureau/pipeflow/providers` |
| LLM | OpenAI | `@moureau/pipeflow/providers` |
| LLM | Claude (Anthropic) | `@moureau/pipeflow/providers` |
| STT | Deepgram | `@moureau/pipeflow/providers` |
| STT | OpenRouter | `@moureau/pipeflow/providers` |
| TTS | Kokoro | `@moureau/pipeflow/providers` |
| TTS | OpenRouter | `@moureau/pipeflow/providers` |
Providers are configured with their own credentials — Pipeflow itself does not hold an API key.
---
## Persistence
```ts
import { SQLitePersistence } from "@moureau/pipeflow/persistence";
const pipeflow = new Pipeflow({
persistence: new SQLitePersistence({ filename: "./pipeflow.db" }),
});
```
In-memory is the default (no persistence adapter needed). Separate persistence from the conversation domain.
---
## Architecture
```
src/
├── agents/ — Agent and tool abstractions
├── conversations/
│ ├── conversation/ — Public realtime conversation API and lifecycle
│ ├── orchestration/ — State machine: speech, transcription, turns, floor, generation, interruptions, tools, TTS
│ └── transcription/ — Conversation transcription and transcript state
├── logger/ — Pluggable logging interface and built-in adapters
├── persistence/ — Storage abstractions with adapters (in-memory, SQLite)
├── providers/ — Vendor-independent interfaces and adapters (STT, LLM, TTS)
└── transport/ — Realtime communication adapters (in-memory, Bun, pluggable ServerAdapter)
```
### Realtime pipeline
```
Audio input → Speech detection → STT stream → partial transcript → Conversation state
→ turn completed → LLM stream → token stream → TTS stream → audio chunks → Application
```
Everything happens as a stream — Pipeflow does not wait for a complete recording before beginning transcription, nor for a complete LLM response before beginning TTS.
---
## Exports
| Import path | Exports |
|---|---|
| `@moureau/pipeflow` | Core: `Pipeflow`, `PipeflowTool`, `Agent`, `ConsoleLogger`, `SilentLogger`, `ContextFn`, `ContextParams`, `Annotations`, Logger types |
| `@moureau/pipeflow/client` | Browser/Node client: `PipeflowClient`, `WebSocketProtocol`, `Protocol` |
| `@moureau/pipeflow/providers` | All provider adapters |
| `@moureau/pipeflow/persistence` | Persistence adapters (`SQLitePersistence`, `MemoryPersistence`) |
| `@moureau/pipeflow/transport` | Transport types, `MemoryTransport`, `ServerAdapter`/`ServerClient` types |
| `@moureau/pipeflow/conversations` | Conversation types and utilities |
---
## Server transport
Pipeflow provides `ServerAdapter`/`ServerClient` interfaces so you can bridge
conversations to remote clients with your own WebSocket server
(Bun, Node `ws`, socket.io, etc.). You own the routing and event subscription:
```ts
import type { ServerAdapter, ServerClient } from "@moureau/pipeflow/transport";
// Subscribe to exactly the events you want to forward
conversation.on("turn", ({ turn }) => client.send(JSON.stringify({ type: "turn", turn })));
conversation.on("generation-complete", ({ generation }) => client.send(...));
// skip text-delta if you don't need streaming text
```
---
## Design Principles
- **Realtime first:** Audio is streamed continuously, not processed as completed recordings.
- **Provider agnostic:** STT, LLM, and TTS providers are adapters, not application-level concepts.
- **Application-owned tools:** Your application executes your tools — Pipeflow only invokes callbacks.
- **Conversations are persistent entities:** A `Conversation` is a runtime handle to a persistent conversation.
- **Agents are independent:** An agent can participate in a conversation or simply be invoked with `run()`.
- **Explicit boundaries:** Pipeflow owns orchestration. Your application owns logic. Providers own their services.
- **Runtime independent:** Transport adapters are pluggable — Bun, Node `ws`, socket.io, or custom.
- **Small public API:** `Pipeflow`, `Agent`, `Conversation`, `Tool` — everything else is replaceable implementation detail.
---
## Development
```bash
git clone git@github.com:moureau-dev/pipeflow.git
cd pipeflow
bun install
bun test # 562+ unit tests (44 files, text-only)
bun run test:e2e # end-to-end (requires API key in .env)
bun run build # transpile to dist/esm + dist/cjs + dist/types
bun run typecheck
```
Default to Bun (not Node): `bun <file>`, `bun test`, `bun run build`, `bunx tsc`.
---
## Common Mistakes
| Wrong | Right |
|---|---|
| Awaiting before calling `conversation.listen()` in a hot audio stream | `listen()` is synchronous — call it each time a buffer arrives |
| Creating a new `Pipeflow` per conversation | One `Pipeflow` instance, many `conversations.create()` calls |
| Expecting `conversation.on("transcript")` to have full text before turn ends | Transcripts are streamed — `turn` event has the finalized text |
| Adding tools via `conversation.on("tool-call")` when auto-execute is on | Tools auto-execute by default — `tool-call` events are for observation |
| Assuming `stop()` also retrieves the transcript | `stop()` and `transcript()` are separate — retrieve after stop |
| Using a provider key without configuring it on the provider instance | Each provider takes its own credentials; Pipeflow does not hold keys |
| Waiting for `start()` to finish before adding participants | `participate()` works before and after `start()` |
| Hardcoding Bun in server transport | Use `ServerAdapter` interface — swap Bun for node-ws, socket.io, etc. |
| Ignoring logger | Add `verbose: true` or pass a custom `Logger` to diagnose issues |