Build a Real-Time Voice Agent
End-to-end bidirectional streaming with STT and TTS
Building a conversational AI requires connecting our Speech-to-Text (STT) and Text-to-Speech (TTS) WebSockets into a single, seamless pipeline. This tutorial provides production-ready architectures that achieve ultra-low latency by streaming raw PCM audio bidirectionally.
- Zero-Latency Playback: Uses pcm_24000 streaming directly into the hardware soundcard.
- Smart Buffering: Prevents audio crackling and handles network jitter automatically.
- Echo Cancellation: Automatically pauses microphone transmission while the AI is speaking.
The Production-Ready Agent
Select your preferred language below. Simply replace the call_llm function with your actual LLM logic (e.g., OpenAI, Claude, or a local model).
import asyncio
import json
import logging
import os
import sys
import time
import pyaudio
import websockets
logging.basicConfig(level=logging.INFO, format="%(asctime)s - [%(levelname)s] - %(message)s")
logger = logging.getLogger("VoiceAgent")
API_KEY = os.getenv("MOKNAH_API_KEY", "YOUR_API_KEY")
STT_URL = "wss://api.moknah.io/api/v1/stt/ws/?language=ar-SA"
TTS_URL = "wss://api.moknah.io/api/v1/tts/ws/"
FORMAT = pyaudio.paInt16
CHANNELS = 1
RATE = 24000
CHUNK = 1200
class MoknahVoiceAgent:
def __init__(self):
self.headers = {"Authorization": f"Bearer {API_KEY}"}
self.is_ai_speaking = False
self.message_queue = asyncio.Queue()
self.audio_queue = asyncio.Queue()
self.audio = pyaudio.PyAudio()
self.mic_stream = self.audio.open(format=FORMAT, channels=CHANNELS, rate=RATE, input=True, frames_per_buffer=CHUNK)
self.speaker_stream = self.audio.open(format=FORMAT, channels=CHANNELS, rate=RATE, output=True, frames_per_buffer=4096)
async def call_llm(self, user_text):
# Replace with OpenAI/Claude call
await asyncio.sleep(0.8)
return f"أهلاً بك! لقد سمعتك تقول: {user_text}"
async def listen_loop(self, ws):
loop = asyncio.get_running_loop()
while True:
try:
data = await loop.run_in_executor(None, self.mic_stream.read, CHUNK, True)
if not self.is_ai_speaking:
await ws.send(data)
except IOError: continue
await asyncio.sleep(0.001)
async def stt_receiver(self, ws):
async for message in ws:
data = json.loads(message)
if data.get('type') == 'final' and data['text'].strip():
print(f"✅ [User]: {data['text']}")
await self.message_queue.put(data['text'])
async def audio_player_task(self):
loop = asyncio.get_running_loop()
pcm_buffer = bytearray()
play_chunk = CHUNK * 4
while True:
item = await self.audio_queue.get()
if isinstance(item, asyncio.Event):
if pcm_buffer: await loop.run_in_executor(None, self.speaker_stream.write, bytes(pcm_buffer))
pcm_buffer.clear()
item.set()
continue
pcm_buffer.extend(item)
while len(pcm_buffer) >= play_chunk:
chunk = pcm_buffer[:play_chunk]
del pcm_buffer[:play_chunk]
if len(chunk) % 2 != 0: chunk = chunk[:-1]
await loop.run_in_executor(None, self.speaker_stream.write, bytes(chunk))
async def speak(self, text):
logger.info(f"🔊 [AI]: {text}")
finished_event = asyncio.Event()
async with websockets.connect(TTS_URL, additional_headers=self.headers) as tts_ws:
await tts_ws.send(json.dumps({"type": "init", "voice_id": 1710, "output_format": "pcm_24000"}))
while True:
if json.loads(await tts_ws.recv()).get("status") == "ready": break
await tts_ws.send(json.dumps({"type": "text", "data": text}))
await tts_ws.send(json.dumps({"type": "stop"}))
async for msg in tts_ws:
if isinstance(msg, bytes): await self.audio_queue.put(msg)
elif json.loads(msg).get("type") == "end_of_stream": break
await self.audio_queue.put(finished_event)
await finished_event.wait()
async def agent_logic(self):
while True:
user_text = await self.message_queue.get()
self.is_ai_speaking = True
try:
ai_response = await self.call_llm(user_text)
if ai_response: await self.speak(ai_response)
finally:
self.is_ai_speaking = False
async def start(self):
async with websockets.connect(STT_URL, additional_headers=self.headers) as stt_ws:
logger.info(🎤 [System]: Agent is LIVE.)
await asyncio.gather(self.listen_loop(stt_ws), self.stt_receiver(stt_ws), self.audio_player_task(), self.agent_logic())
if __name__ == "__main__":
asyncio.run(MoknahVoiceAgent().start())
npm install ws node-record-lpcm16 speakerconst WebSocket = require('ws');
const record = require('node-record-lpcm16');
const Speaker = require('speaker');
const API_KEY = process.env.MOKNAH_API_KEY || "YOUR_API_KEY";
let isAiSpeaking = false;
// Hardware Speaker (24kHz PCM)
const speaker = new Speaker({
channels: 1,
bitDepth: 16,
sampleRate: 24000
});
const sttWs = new WebSocket('wss://api.moknah.io/api/v1/stt/ws/?language=ar-SA', {
headers: { "Authorization": `Bearer ${API_KEY}` }
});
// Dummy LLM Call
async function callLLM(text) {
return new Promise(resolve => setTimeout(() => resolve("أهلاً بك، لقد سمعتك تقول: " + text), 800));
}
async function speak(text) {
isAiSpeaking = true;
return new Promise((resolve) => {
const ttsWs = new WebSocket('wss://api.moknah.io/api/v1/tts/ws/', {
headers: { "Authorization": `Bearer ${API_KEY}` }
});
ttsWs.on('open', () => {
ttsWs.send(JSON.stringify({ type: "init", voice_id: 1710, output_format: "pcm_24000" }));
});
ttsWs.on('message', (data, isBinary) => {
if (isBinary) {
speaker.write(data); // Stream directly to soundcard
} else {
const msg = JSON.parse(data.toString());
if (msg.status === "ready") {
ttsWs.send(JSON.stringify({ type: "text", data: text }));
ttsWs.send(JSON.stringify({ type: "stop" }));
} else if (msg.type === "end_of_stream") {
ttsWs.close();
setTimeout(() => {
isAiSpeaking = false;
resolve();
}, 500); // Wait for buffer to drain
}
}
});
});
}
sttWs.on('open', () => {
console.log("🎤 Agent is LIVE.");
// Start Microphone (16kHz)
const mic = record.record({ sampleRate: 16000, channels: 1, threshold: 0 });
mic.stream().on('data', (data) => {
if (!isAiSpeaking && sttWs.readyState === WebSocket.OPEN) {
sttWs.send(data);
}
});
});
sttWs.on('message', async (data) => {
const msg = JSON.parse(data.toString());
if (msg.type === "final" && msg.text.trim()) {
console.log(`✅ [User]: ${msg.text}`);
const response = await callLLM(msg.text);
await speak(response);
}
});
org.java-websocket:Java-WebSocket:1.5.3 and native javax.sound.sampledimport org.java_websocket.client.WebSocketClient;
import org.java_websocket.handshake.ServerHandshake;
import javax.sound.sampled.*;
import java.net.URI;
import java.util.Map;
public class VoiceAgent {
private static final String API_KEY = "YOUR_API_KEY";
private static volatile boolean isAiSpeaking = false;
private static SourceDataLine speaker;
public static void main(String[] args) throws Exception {
// Setup Output Audio (24kHz PCM for TTS)
AudioFormat outFormat = new AudioFormat(24000, 16, 1, true, false);
speaker = AudioSystem.getSourceDataLine(outFormat);
speaker.open(outFormat, 4096);
speaker.start();
// Connect STT WebSocket
WebSocketClient sttClient = new WebSocketClient(new URI("wss://api.moknah.io/api/v1/stt/ws/?language=ar-SA"), Map.of("Authorization", "Bearer " + API_KEY)) {
@Override public void onMessage(String message) {
if (message.contains("\"type\":\"final\"") && !message.contains("\"text\":\"\"")) {
System.out.println("✅ User finished speaking.");
speak("أهلاً بك!"); // Example LLM response
}
}
@Override public void onOpen(ServerHandshake h) { System.out.println("🎤 Agent LIVE."); }
@Override public void onClose(int c, String r, boolean rmt) {}
@Override public void onError(Exception e) {}
};
sttClient.connectBlocking();
// Setup Input Audio (16kHz PCM for STT)
AudioFormat inFormat = new AudioFormat(16000, 16, 1, true, false);
TargetDataLine mic = AudioSystem.getTargetDataLine(inFormat);
mic.open(inFormat);
mic.start();
byte[] buffer = new byte[1200];
while (true) {
int bytesRead = mic.read(buffer, 0, buffer.length);
if (!isAiSpeaking && sttClient.isOpen()) {
sttClient.send(buffer); // Stream to STT
}
}
}
public static void speak(String text) {
isAiSpeaking = true;
try {
WebSocketClient ttsClient = new WebSocketClient(new URI("wss://api.moknah.io/api/v1/tts/ws/"), Map.of("Authorization", "Bearer " + API_KEY)) {
@Override public void onOpen(ServerHandshake handshakedata) {
send("{\"type\":\"init\",\"voice_id\":1710,\"output_format\":\"pcm_24000\"}");
}
@Override public void onMessage(String message) {
if (message.contains("\"status\":\"ready\"")) {
send("{\"type\":\"text\",\"data\":\"" + text + "\"}");
send("{\"type\":\"stop\"}");
} else if (message.contains("end_of_stream")) {
isAiSpeaking = false;
close();
}
}
@Override public void onMessage(java.nio.ByteBuffer bytes) {
byte[] audio = bytes.array();
speaker.write(audio, 0, audio.length); // Play to Soundcard
}
@Override public void onClose(int c, String r, boolean rmt) {}
@Override public void onError(Exception ex) {}
};
ttsClient.connectBlocking();
} catch (Exception e) { e.printStackTrace(); }
}
}
composer require ratchet/pawl react/event-loop<?php
require __DIR__ . '/vendor/autoload.php';
use React\EventLoop\Loop;
use Ratchet\Client\WebSocket;
$apiKey = getenv('MOKNAH_API_KEY') ?: 'YOUR_API_KEY';
$headers = ['Authorization' => 'Bearer ' . $apiKey];
// Simulate WebRTC Audio received from Frontend
$frontendAudioStream = fopen('php://stdin', 'r');
$isAiSpeaking = false;
\Ratchet\Client\connect('wss://api.moknah.io/api/v1/stt/ws/?language=ar-SA', [], $headers)->then(function(WebSocket $sttConn) use (&$isAiSpeaking, $headers) {
echo "🎤 Backend Bridge Connected to STT.\n";
// Listen to STT events
$sttConn->on('message', function($msg) use (&$isAiSpeaking, $headers) {
$data = json_decode($msg, true);
if (isset($data['type']) && $data['type'] === 'final' && trim($data['text'])) {
echo "✅ User: " . $data['text'] . "\n";
// Trigger LLM & TTS Pipeline
$isAiSpeaking = true;
$aiResponse = "أهلاً بك!";
\Ratchet\Client\connect('wss://api.moknah.io/api/v1/tts/ws/', [], $headers)->then(function(WebSocket $ttsConn) use ($aiResponse, &$isAiSpeaking) {
$ttsConn->send(json_encode(["type" => "init", "voice_id" => 1710, "output_format" => "pcm_24000"]));
$ttsConn->on('message', function($message) use ($ttsConn, $aiResponse, &$isAiSpeaking) {
if (is_string($message)) {
$json = json_decode($message, true);
if ($json['status'] ?? false === 'ready') {
$ttsConn->send(json_encode(["type" => "text", "data" => $aiResponse]));
$ttsConn->send(json_encode(["type" => "stop"]));
} elseif ($json['type'] ?? false === 'end_of_stream') {
$isAiSpeaking = false;
$ttsConn->close();
}
} else {
// Binary PCM data received! Stream it back to Frontend via WebRTC
echo "Streaming binary bytes to frontend...\n";
}
});
});
}
});
// Send incoming frontend audio to Moknah
Loop::addPeriodicTimer(0.05, function () use ($sttConn, &$isAiSpeaking) {
$audioChunk = fread(STDIN, 1600); // Read from frontend socket/pipe
if ($audioChunk && !$isAiSpeaking) {
$sttConn->send($audioChunk);
}
});
}, function($e) {
echo "Could not connect: {$e->getMessage()}\n";
});
Loop::run();
?>
Dive Deeper & API Reference
This tutorial demonstrates the high-level orchestration of a bidirectional voice agent. For a complete list of parameters, available languages, voice customization options, and error codes, please explore the dedicated API references:
- Speech-to-Text (STT) WebSocket API – Learn about supported languages, partial vs final transcripts, and configuring silence timeouts.
- Text-to-Speech (TTS) WebSocket API – Learn about voice IDs, fast_mode, MP3 vs PCM formats, and emotion customization.
For API-related questions or issues, contact us at api@moknah.io.
pip install websockets pyaudio