Realtime voice infrastructure for TypeScript.
Pipeflow is an open-source backend SDK for building voice agents, conversational applications, meeting transcribers, Discord bots, recruiters, assistants, and other realtime audio experiences.
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.
- Small — ~140 kB packed, one runtime dependency (zod).
- 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, and Kokoro today).
- Your backend stays yours — tools and the audio transport are owned by your application; Pipeflow never executes your code.
🚧 Early development
The API is evolving and should be considered experimental.
Some capabilities ship today; others are architectural boundaries the API is shaped around but not yet implemented.
| Area | Current | Designed for |
|---|---|---|
| Providers | DeepSeek, OpenRouter (LLM), Deepgram Flux (STT), Kokoro (TTS) | additional providers behind the same interfaces |
| Persistence | in-memory, SQLite | Postgres, Redis, … behind the same contract |
| Transport | in-memory (tests, development) | WebSocket, WebRTC, … |
| Conversation addressing | agent names and aliases | floor management, multi-participant turn-taking |
| Realtime bounds | step-bounded coordination graph | hierarchical latency budgets |
Pipeflow has four core concepts:
- Agent — intelligence, context, and tools.
- Conversation — a persistent realtime conversation and its participants.
- Tool — application capabilities executed by your backend.
- Provider — an implementation of STT, LLM, or TTS.
The goal is to keep these concerns independent.
An agent does not inherently own a conversation. A conversation does not require an agent. A provider should not leak into your application logic.
Pipeflow
┌──────────────┐
│ Agent │
│ context │
│ tools │
└──────┬───────┘
│
▼
┌──────────────┐
│ Conversation │
│ │
│ participants │
│ turns │
│ audio │
│ interruption │
└──────┬───────┘
│
▼
┌─────────────────────────┐
│ Orchestrator │
└──────┬──────┬──────┬────┘
│ │ │
STT LLM TTS
│ │ │
▼ ▼ ▼
Deepgram DeepSeek Kokoro
OpenRouter
bun add @moureau/pipeflowPublished on npm as @moureau/pipeflow. Alternatively, install directly from
the repository (bun add git@github.com:moureau-dev/pipeflow.git) or build
from source — see Development.
Create an agent:
import { Pipeflow } from "@moureau/pipeflow";
import { DeepSeekLLM, DeepgramSTT, KokoroTTS } from "@moureau/pipeflow/providers";
// Providers are configured explicitly with their own credentials —
// Pipeflow itself does not hold an API key.
const stt = new DeepgramSTT({ apiKey: process.env.DEEPGRAM_API_KEY });
const tts = new KokoroTTS();
const pipeflow = new Pipeflow({
llm: new DeepSeekLLM({ apiKey: process.env.DEEPSEEK_API_KEY }),
stt,
tts,
});
const jarvis = pipeflow.agent({
name: "Jarvis",
context: `
You are Jarvis, a helpful voice assistant.
Keep your responses concise and conversational.
`,
});Create a conversation:
const conversation = await pipeflow.conversations.create({
agents: [jarvis],
});Starting a conversation starts its realtime machinery: start() attaches the
orchestrator, which runs the STT → turns → LLM → TTS pipeline for the
conversation.
await conversation.start();Add a participant:
Feed it audio as it arrives:
voice.onAudio((audio) => {
conversation.listen({
userId: "alice",
audio,
});
});Listen for generated audio:
conversation.on("audio", ({ audio }) => {
voice.play(audio);
});When finished:
await conversation.stop();The application owns the audio transport. Pipeflow handles the realtime voice pipeline.
Your application
│
│ audio chunks
▼
Pipeflow
│
├── STT
├── conversation orchestration
├── LLM
└── TTS
│
│ audio chunks
▼
Your application
A conversation is the persistent entity representing a realtime interaction.
See src/conversations/README.md for the conversation domain.
const conversation =
await pipeflow.conversations.create({
agents: [jarvis],
});
console.log(conversation.id);
await conversation.start();Creation and realtime execution are deliberately separate:
create() // creates the persistent conversation
start() // moves the conversation into the started state
participate() // adds participants
listen() // sends audio
send() // injects a finalized text turn (no STT)
stop() // finalizes the realtime sessionstart() attaches the orchestrator, which subscribes to audio-in events,
runs the STT/LLM/TTS pipeline, and pushes generated audio, turns, transcripts,
and tool calls back through conversation events.
listen() is intentionally synchronous — it means "send this audio packet",
not "wait for this utterance to finish":
conversation.listen({ userId, audio });This makes it suitable for high-frequency realtime audio streams.
send({ userId, text }) does the same for a finalized text turn, bypassing
STT — text-first integrations (chat, Discord text channels) route through the
same pipeline: routing, coordination, clarification, and generation.
Participants can be added individually:
await conversation.participate({
userId: "alice",
});Or in batches:
await conversation.participate([
{ userId: "alice" },
{
userId: "bob",
aliases: ["robert", "rob"],
},
{
userId: "charlie",
aliases: ["charles"],
},
]);Participant information can be used for speaker attribution, conversation history, addressing, and multi-participant floor management.
Voice conversations need to feel immediate.
If an agent is speaking and a participant starts talking, Pipeflow can interrupt the current output immediately rather than waiting for the entire utterance to be transcribed.
Interruptions can be triggered by the application — or automatically: the orchestrator detects when a participant starts speaking while the agent is responding (barge-in).
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.
conversation.interrupt(); // stops TTS, cancels the current generationThe interrupting speech keeps flowing into STT and becomes the next turn.
Conceptually:
Agent is speaking
│
▼
Participant starts speaking
│
├── stop TTS
├── cancel current generation
└── continue receiving speech
│
▼
STT
│
▼
new conversation turn
This keeps interactions responsive even when the user interrupts the agent halfway through a sentence.
The conversation emits a typed event stream the application can subscribe to:
audio-in— raw audio fed in vialisten()text-in— a finalized text turn viasend()partial-transcript— live STT partials (captions)turn— a finalized participant turntranscript— a transcript entryaudio— generated audio to playgeneration— an agent generationtool-call/tool-call-result— tool execution round tripsinterrupt— an interruption occurrederror— a provider failurestart/stop/state— lifecycle
conversation.on("audio", ({ audio }) => voice.play(audio));
conversation.on("partial-transcript", ({ text }) => captions.update(text));Batching participants, conversational state, and the roadmap for floor management and addressing
Conversations can contain multiple participants and agents.
const conversation =
await pipeflow.conversations.create({
agents: [jarvis],
});
await conversation.participate([
{ userId: "alice" },
{ userId: "bob" },
{ userId: "charlie" },
]);In a one-to-one conversation, speech is treated as an interaction: every finalized participant turn produces an agent generation.
Addressing works at a basic level: in a multi-agent conversation, a turn is
routed to the agent whose name or alias appears in the speech, and unaddressed
turns go through the built-in understand coordination, which may delegate to
agents, ask the user, or answer directly. Floor management, multi-participant
turn-taking rules, and richer addressing heuristics are on the roadmap. A wake
word is not intended to be a fundamental requirement.
An agent defines an AI persona and its capabilities.
See src/agents/README.md for the agent and tool abstractions.
const jarvis = pipeflow.agent({
name: "Jarvis",
context: `
You are Jarvis.
You are concise, helpful, and conversational.
`,
tools: [
getWeather,
searchCalendar,
],
});Agents can also be run independently of conversations:
const result = await jarvis.run({
prompt: "Explain how a neural network works.",
});
console.log(result.text);This is useful for ordinary LLM workloads where realtime audio and conversation state are not required.
When an agent is attached to a conversation, it can take conversational turns as part of the realtime orchestration.
A conversation can coordinate multiple agents: create({ agents }) accepts
a roster. 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. In a single-agent conversation, every turn goes to that agent. Each
agent keeps its own context and its own LLM; the conversation owns the shared
runtime (STT, TTS, history, interruptions).
const receptionist = pipeflow.agent({
name: "Receptionist",
context: `
You are the front desk. Greet people and route requests.
`,
});
const specialist = pipeflow.agent({
name: "Technical Specialist",
aliases: ["tech"],
context: `
You are the technical support specialist.
`,
});
const conversation =
await pipeflow.conversations.create({
agents: [receptionist, specialist],
});"Ask the technical specialist about X" is routed to the specialist; an
unaddressed turn goes through the built-in understand coordination.
understand is a hardcoded coordination: a reasoning unit that decides
what should happen next rather than doing the work itself. It can:
- delegate to one or more agents (in parallel), each with a self-contained prompt;
- pass the work to another coordination;
- ask the user for missing details via a structured, batched
clarifyaction — every missing detail goes into one question, and question rounds are capped per run (default 2), after which reasonable assumptions are stated; - complete with a direct answer.
User: "Book a flight and check whether my calendar conflicts."
│
▼
┌────────────┐
│ understand │ ← built-in coordination
└─────┬──────┘
│
┌──────┴───────┐
▼ ▼
Travel Agent Calendar Agent
│ │
└──────┬───────┘
▼
understand
│
▼
User
Agent delegation runs as sub-generations: the target agent executes on its own LLM, context, and tools (which still run in your backend), and every delegated prompt is stamped with the current time so time-sensitive tasks reason about the right "now". Delegated agents are text-only — the coordination narrates while they work and speaks the merged answer.
Clarification is a first-class operation: an ambiguous request parks the coordination and resumes on the next turn instead of starting fresh. While waiting, participant speech is treated as the answer, not as a barge-in.
See src/conversations/orchestration/coordination/README.md for the coordination model and src/conversations/orchestration/orchestrator/README.md for how it is wired into the realtime pipeline.
Tools expose capabilities from your application to an agent.
const getWeather = new PipeflowTool({
name: "get_weather",
description: "Get the current weather for a city.",
execute: async ({ city }) => {
return weatherService.getCurrent(city);
},
});Then:
const jarvis = pipeflow.agent({
name: "Jarvis",
context: "You are a helpful assistant.",
tools: [getWeather],
});Tools execute in your backend.
Pipeflow does not execute arbitrary application code.
Pipeflow
│
LLM requests tool
│
▼
┌─────────────────┐
│ Your application │
│ │
│ execute() │
└────────┬────────┘
│
▼
tool result
│
▼
LLM
This allows tools to access your database, APIs, Discord bot, business logic, filesystem, or anything else your application controls.
In a conversation, tool calls never execute inside Pipeflow. The orchestrator
emits a tool-call event with the requested tool and its arguments; your
backend executes it and reports back:
conversation.on("tool-call", async ({ call }) => {
const result = await myBackend.execute(call.name, call.arguments);
conversation.resolveToolCall({
id: call.id,
result,
});
});The agent's narration continues while the tool runs, and the generation resumes with the tool result once it is resolved.
Tool calls the model requests together run concurrently: run() executes
them in parallel, and in a conversation your backend resolves them in
parallel.
A conversation does not require an agent.
This makes Pipeflow useful as a realtime transcription primitive.
const conversation =
await pipeflow.conversations.create();
// Without an agent, `start()` attaches the orchestrator in
// transcription-only mode: audio in, turns and transcripts out.
await conversation.start();
await conversation.participate([
{ userId: "alice" },
{ userId: "bob", aliases: ["robert"] },
]);
discordVoice.onAudio((userId, audio) => {
conversation.listen({
userId,
audio,
});
});
await conversation.stop();Retrieve the transcript afterward:
const transcript =
await pipeflow.conversations.transcript(
conversation.id,
);Transcript retrieval is separate from stop() so that ending a conversation does not require loading an arbitrarily large transcript into memory.
Pagination can be supported for long conversations.
Summarizing a meeting with a notetaker agent — as a plain LLM task or through a transcript tool
A meeting summary can simply be another agent task.
const notetaker = pipeflow.agent({
name: "Meeting Notetaker",
context: `
You are a meeting notetaker.
Produce concise notes containing:
- summary
- decisions
- action items
- unresolved questions
`,
});Retrieve the transcript:
const transcript =
await pipeflow.conversations.transcript(
conversation.id,
);Then run the agent:
const result = await notetaker.run({
prompt: `
The meeting transcription is:
${transcript.join("\n")}
Produce the meeting notes.
`,
});Or the notetaker can retrieve the transcript itself through a tool:
const getTranscript = new PipeflowTool({
name: "get_transcript",
description: "Retrieve the meeting transcript.",
execute: async () => {
return pipeflow.conversations.transcript(
conversation.id,
);
},
});The agent does not receive special access to conversations.
If it needs conversation data, a tool provides that capability.
Vendor-independent interfaces and the current adapters: Deepgram, DeepSeek, OpenRouter, Kokoro
Pipeflow separates provider interfaces (LLM, STT, TTS) from their implementations. The orchestrator works against the interfaces rather than directly against vendor APIs, so providers can be replaced without changing the conversation layer.
The project currently ships adapters for:
- STT: Deepgram
- LLM: DeepSeek, OpenRouter
- TTS: Kokoro
Providers are configured with their own credentials — Pipeflow itself does not hold an API key.
Provider availability and configuration are evolving during early development.
See src/providers/README.md for the interfaces and adapter contracts.
How the layers fit together: agents, conversations, persistence, providers, transport
Pipeflow is designed around a small number of independent layers.
src/
├── agents/
├── conversations/
│ ├── conversation/
│ ├── orchestration/
│ └── transcription/
├── persistence/
│ └── adapters/
├── providers/
│ ├── llm/
│ ├── stt/
│ └── tts/
└── transport/
Public realtime conversation API and lifecycle. See src/conversations/conversation/README.md.
See src/conversations/orchestration/README.md. The state machine coordinating:
- speech
- transcription
- turns
- floor state
- agent generation
- interruptions
- tools
- TTS
Conversation transcription and transcript state. See src/conversations/transcription/README.md.
Vendor-independent interfaces and provider adapters. See src/providers/README.md.
Persistence abstractions with adapters such as SQLite and in-memory storage. See src/persistence/README.md.
Realtime communication between Pipeflow and the application. See src/transport/README.md.
Storage adapters: in-memory for development, SQLite for lightweight persistence
Pipeflow separates persistence from the conversation domain. The in-memory adapter is useful for tests and development; SQLite provides a lightweight persistent backend suitable for local applications and early deployments. The persistence interface is intentionally provider-independent so other storage implementations can be added later.
import { SQLitePersistence } from "@moureau/pipeflow/persistence";
const pipeflow = new Pipeflow({
persistence: new SQLitePersistence({ filename: "./pipeflow.db" }),
});See src/persistence/README.md for the storage contract and adapters.
How audio flows through the pipeline as a continuous stream
A typical voice interaction looks like:
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 does it wait for a complete LLM response before beginning TTS.
The intended flow is:
audio
↓
partial STT
↓
turn detection
↓
LLM streaming
↓
TTS streaming
↓
audio
This allows the system to begin producing speech as early as possible.
Why the project is open and how the provider layer stays modular
Pipeflow is open source.
The project is designed to make realtime voice infrastructure accessible without requiring applications to implement their own orchestration layer.
The provider layer is intentionally modular so applications can choose between hosted and self-hosted services.
Clone the repository and install dependencies:
git clone git@github.com:moureau-dev/pipeflow.git
cd pipeflow
bun installRun tests:
bun testEnd-to-end tests hit the real LLM API and are skipped when no key is
available. They prefer OPENROUTER_API_KEY (default model
google/gemini-2.5-flash-lite, override with LLM_MODEL) and fall back to
DEEPSEEK_API_KEY. Add either to .env (loaded automatically) and run:
bun run test:e2eThe latency benchmark runs the pipeline repeatedly against the real model and
reports p50/p95 per hop: first token, first speechable text, TTS request, TTS
first audio, first audio delivered, and completion. STT and TTS are faked by
default; point KOKORO_URL at a Kokoro endpoint to measure the real synthesis
path, including inter-chunk audio gaps:
bun run benchmark # 10 runs; BENCH_RUNS=5 to change
KOKORO_URL=http://localhost:8880 bun run benchmark # local kokoro-fastapi
KOKORO_URL=https://api.together.ai \
KOKORO_API_KEY=... KOKORO_MODEL=hexgrad/Kokoro-82M \
bun run benchmark # Together AIBuild and type-check:
bun run build # transpile to dist/esm + dist/cjs and emit dist/types
bun run typecheckThe project uses Bun and TypeScript.
The rules the API is built around: realtime first, provider agnostic, app-owned tools
Audio is streamed continuously rather than processed as completed recordings.
STT, LLM, and TTS providers are adapters, not application-level concepts.
Your application executes your tools.
A realtime Conversation instance is a runtime handle to a persistent conversation.
An agent can participate in a conversation or simply be invoked with run().
Pipeflow owns orchestration.
Your application owns application logic.
Providers own their respective AI services.
The core API should remain centered around:
Pipeflow
Agent
Conversation
Tool
Everything else should remain replaceable implementation detail for as long as possible.
See LICENSE.