Skip to content

Repository files navigation

Pipeflow

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.

Highlights

  • 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.

Status

🚧 Early development

The API is evolving and should be considered experimental.

Current vs. designed for

Some capabilities ship today; others are architectural boundaries the API is shaped around but not yet implemented.

AreaCurrentDesigned for
ProvidersDeepSeek, OpenRouter (LLM), Deepgram Flux (STT), Kokoro (TTS)additional providers behind the same interfaces
Persistencein-memory, SQLitePostgres, Redis, … behind the same contract
Transportin-memory (tests, development)WebSocket, WebRTC, …
Conversation addressingagent names and aliasesfloor management, multi-participant turn-taking
Realtime boundsstep-bounded coordination graphhierarchical latency budgets

Philosophy

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

Installation

bun add @moureau/pipeflow

Published 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.

Basic voice agent

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.conststt=newDeepgramSTT({apiKey: process.env.DEEPGRAM_API_KEY});consttts=newKokoroTTS();constpipeflow=newPipeflow({llm: newDeepSeekLLM({apiKey: process.env.DEEPSEEK_API_KEY}),
stt,
tts,});constjarvis=pipeflow.agent({name: "Jarvis",context: ` You are Jarvis, a helpful voice assistant. Keep your responses concise and conversational. `,});

Create a conversation:

constconversation=awaitpipeflow.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.

awaitconversation.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:

awaitconversation.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

Conversations

A conversation is the persistent entity representing a realtime interaction.

See src/conversations/README.md for the conversation domain.

constconversation=awaitpipeflow.conversations.create({agents: [jarvis],});console.log(conversation.id);awaitconversation.start();

Creation and realtime execution are deliberately separate:

create()// creates the persistent conversationstart()// moves the conversation into the started stateparticipate()// adds participantslisten()// sends audiosend()// injects a finalized text turn (no STT)stop()// finalizes the realtime session

start() 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

Participants can be added individually:

awaitconversation.participate({userId: "alice",});

Or in batches:

awaitconversation.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.

Interruption

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 generation

The 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.

Conversation events

The conversation emits a typed event stream the application can subscribe to:

  • audio-in — raw audio fed in via listen()
  • text-in — a finalized text turn via send()
  • partial-transcript — live STT partials (captions)
  • turn — a finalized participant turn
  • transcript — a transcript entry
  • audio — generated audio to play
  • generation — an agent generation
  • tool-call / tool-call-result — every tool round trip (auto-executed by default; observable even when the framework runs the tool)
  • interrupt — an interruption occurred
  • error — a provider failure
  • start / stop / state — lifecycle
conversation.on("audio",({ audio })=>voice.play(audio));conversation.on("partial-transcript",({ text })=>captions.update(text));

Multi-participant conversations

Batching participants, conversational state, and the roadmap for floor management and addressing

Conversations can contain multiple participants and agents.

constconversation=awaitpipeflow.conversations.create({agents: [jarvis],});awaitconversation.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.

Agents

An agent defines an AI persona and its capabilities.

See src/agents/README.md for the agent and tool abstractions.

constjarvis=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:

constresult=awaitjarvis.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.

Conversational agents

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).

constreceptionist=pipeflow.agent({name: "Receptionist",context: ` You are the front desk. Greet people and route requests. `,});constspecialist=pipeflow.agent({name: "Technical Specialist",aliases: ["tech"],context: ` You are the technical support specialist. `,});constconversation=awaitpipeflow.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.

Collaborative orchestration

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 clarify action — 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

Tools expose capabilities from your application to an agent.

import{z}from"zod";constgetWeather=newPipeflowTool({name: "get_weather",description: "Get the current weather for a city.",// One zod schema drives both sides: it derives the JSON schema the model// sees, and validates the arguments before execute() runs.schema: {in: z.object({city: z.string().describe("The city to look up.")}),},execute: async({ city })=>{returnweatherService.getCurrent(city);},});

schema.in is what the model sends; schema.out (optional, defaults to in) is what execute receives after validation, and may transform the arguments. Prefer schema — it keeps the model contract and your execute signature in sync and rejects bad arguments with a clear error instead of crashing inside the tool. If you need a hand-written JSON schema instead, pass parameters (no runtime validation).

Then:

constjarvis=pipeflow.agent({name: "Jarvis",context: "You are a helpful assistant.",tools: [getWeather],});

Tools execute in your backend: the execute function is your code, and it runs in your process — Pipeflow only invokes the callbacks you registered.

 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: tools auto-execute

Attach tools to an agent and they just work — the orchestrator runs each tool and feeds the result back into the generation (same loop as Agent.run()):

constjarvis=pipeflow.agent({name: "Jarvis",context: "You are a helpful assistant.",tools: [getWeather],});constconversation=awaitpipeflow.conversations.create({agents: [jarvis]});// no tool-call handler needed — get_weather runs automatically

The tool-call (and tool-call-result) events still fire for every call, so you can observe or log them. A thrown tool error is caught and returned to the model as { "error": ... } so the agent can recover; an unknown tool name reports Unknown tool "..." the same way. Tool calls the model issues together run concurrently, and a hung tool is bounded by toolTimeoutMs (default 30s).

Opt out: resolve tool calls yourself

When a tool needs approval, or runs in a different backend, pass autoExecuteTools: false (on Pipeflow, pipeflow.conversations.create(), or Conversation) and the orchestrator hands the call to you:

constconversation=awaitpipeflow.conversations.create({agents: [jarvis],autoExecuteTools: false,});conversation.on("tool-call",async({ call })=>{// call.arguments is already parsed JSONconstresult=awaitmyBackend.executeTool(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.

Meeting transcription

A conversation does not require an agent.

This makes Pipeflow useful as a realtime transcription primitive.

constconversation=awaitpipeflow.conversations.create();// Without an agent, `start()` attaches the orchestrator in// transcription-only mode: audio in, turns and transcripts out.awaitconversation.start();awaitconversation.participate([{userId: "alice"},{userId: "bob",aliases: ["robert"]},]);discordVoice.onAudio((userId,audio)=>{conversation.listen({
userId,
audio,});});awaitconversation.stop();

Retrieve the transcript afterward:

consttranscript=awaitpipeflow.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.

Meeting summaries

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.

constnotetaker=pipeflow.agent({name: "Meeting Notetaker",context: ` You are a meeting notetaker. Produce concise notes containing: - summary - decisions - action items - unresolved questions `,});

Retrieve the transcript:

consttranscript=awaitpipeflow.conversations.transcript(conversation.id,);

Then run the agent:

constresult=awaitnotetaker.run({prompt: ` The meeting transcription is:${transcript.join("\n")} Produce the meeting notes. `,});

Or the notetaker can retrieve the transcript itself through a tool:

constgetTranscript=newPipeflowTool({name: "get_transcript",description: "Retrieve the meeting transcript.",execute: async()=>{returnpipeflow.conversations.transcript(conversation.id,);},});

The agent does not receive special access to conversations.

If it needs conversation data, a tool provides that capability.

Providers

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, OpenAI, Claude (Anthropic)
  • 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.

Architecture

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/

Conversation

Public realtime conversation API and lifecycle. See src/conversations/conversation/README.md.

Orchestration

See src/conversations/orchestration/README.md. The state machine coordinating:

  • speech
  • transcription
  • turns
  • floor state
  • agent generation
  • interruptions
  • tools
  • TTS

Transcription

Conversation transcription and transcript state. See src/conversations/transcription/README.md.

Providers

Vendor-independent interfaces and provider adapters. See src/providers/README.md.

Persistence

Persistence abstractions with adapters such as SQLite and in-memory storage. See src/persistence/README.md.

Transport

Realtime communication between Pipeflow and the application. See src/transport/README.md.

Persistence

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";constpipeflow=newPipeflow({persistence: newSQLitePersistence({filename: "./pipeflow.db"}),});

See src/persistence/README.md for the storage contract and adapters.

Realtime architecture

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.

Open source

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.

Development

Clone the repository and install dependencies:

git clone git@github.com:moureau-dev/pipeflow.git
cd pipeflow
bun install

Run tests:

bun test

End-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:e2e

The 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 AI

Build and type-check:

bun run build # transpile to dist/esm + dist/cjs and emit dist/types
bun run typecheck

The project uses Bun and TypeScript.

Design principles

The rules the API is built around: realtime first, provider agnostic, app-owned tools

Realtime first

Audio is streamed continuously rather than 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.

Conversations are persistent entities

A realtime Conversation instance 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 application logic.

Providers own their respective AI services.

Small public API

The core API should remain centered around:

Pipeflow
Agent
Conversation
Tool

Everything else should remain replaceable implementation detail for as long as possible.

License

See LICENSE.

About

Pipeflow is an open-source realtime, hierarchical multi-agent runtime for TypeScript, with voice, tool execution, agent routing, and parallel sub-agent orchestration.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Contributors

Languages