User Guide API Reference

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.

🧠 What this architecture handles for you:
  • 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).

Requires: pip install websockets pyaudio
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())
Requires: npm install ws node-record-lpcm16 speaker
const 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);
    }
});
Requires: org.java-websocket:Java-WebSocket:1.5.3 and native javax.sound.sampled
import 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(); }
    }
}
Note: PHP is usually used as an Async Server Bridge (Ratchet/ReactPHP) connecting a web-frontend to Moknah, rather than accessing hardware microphones directly. Requires: 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:

API Support

For API-related questions or issues, contact us at api@moknah.io.