realtime

package
v1.55.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Mar 20, 2026 License: MIT Imports: 24 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func HandleSessions

func HandleSessions(server *RealtimeServer) http.HandlerFunc

HandleSessions handles HTTP requests for session management

Types

type Adapter

type Adapter interface {
	SendAudio(base64Audio string) error
	SendText(text string) error
	HandleControl(action string, data map[string]interface{}) error
	ReceiveEvents() <-chan UnifiedEvent
	Close() error
	InputSampleRate() int // Expected input audio sample rate (e.g., 16000 or 24000)
}

Adapter interface for different realtime providers

type AudioDeltaData

type AudioDeltaData struct {
	Format     string `json:"format"`      // "pcm16"
	SampleRate int    `json:"sample_rate"` // 24000
	Chunk      string `json:"chunk"`       // Base64-encoded audio
}

AudioDeltaData contains audio chunk data

type ElevenLabsRealtimeSTT added in v1.41.4

type ElevenLabsRealtimeSTT struct {
	// contains filtered or unexported fields
}

ElevenLabsRealtimeSTT implements StreamingSTTProvider using ElevenLabs WebSocket API. Audio chunks are streamed in real-time and ElevenLabs handles VAD server-side, returning committed_transcript when it detects the user has finished speaking.

func NewElevenLabsRealtimeSTT added in v1.41.4

func NewElevenLabsRealtimeSTT(cfg *config.STTConfig) (*ElevenLabsRealtimeSTT, error)

NewElevenLabsRealtimeSTT creates a streaming STT provider

func (*ElevenLabsRealtimeSTT) Close added in v1.41.4

func (e *ElevenLabsRealtimeSTT) Close() error

Close shuts down the WebSocket connection

func (*ElevenLabsRealtimeSTT) PartialTranscripts added in v1.41.4

func (e *ElevenLabsRealtimeSTT) PartialTranscripts() <-chan string

PartialTranscripts returns the channel of partial/interim transcripts

func (*ElevenLabsRealtimeSTT) SendAudio added in v1.41.4

func (e *ElevenLabsRealtimeSTT) SendAudio(pcm16Data []byte) error

SendAudio sends a PCM16 audio chunk to ElevenLabs via WebSocket

func (*ElevenLabsRealtimeSTT) Transcripts added in v1.41.4

func (e *ElevenLabsRealtimeSTT) Transcripts() <-chan string

Transcripts returns the channel of committed transcripts

type ElevenLabsRealtimeTTS added in v1.41.4

type ElevenLabsRealtimeTTS struct {
	// contains filtered or unexported fields
}

ElevenLabsRealtimeTTS implements StreamingTTSProvider using ElevenLabs WebSocket API. Text chunks are sent as they arrive from the LLM, and audio is streamed back progressively for ultra-low latency text-to-speech.

func NewElevenLabsRealtimeTTS added in v1.41.4

func NewElevenLabsRealtimeTTS(cfg *config.TTSConfig) (*ElevenLabsRealtimeTTS, error)

NewElevenLabsRealtimeTTS creates a streaming TTS provider

func (*ElevenLabsRealtimeTTS) AudioChunks added in v1.41.4

func (e *ElevenLabsRealtimeTTS) AudioChunks() <-chan []byte

AudioChunks returns the channel of PCM16 audio chunks

func (*ElevenLabsRealtimeTTS) Close added in v1.41.4

func (e *ElevenLabsRealtimeTTS) Close() error

Close shuts down the WebSocket connection

func (*ElevenLabsRealtimeTTS) Finish added in v1.41.4

func (e *ElevenLabsRealtimeTTS) Finish() error

Finish signals end of text and waits for all audio to be received

func (*ElevenLabsRealtimeTTS) Flush added in v1.41.4

func (e *ElevenLabsRealtimeTTS) Flush() error

Flush forces generation of any buffered text

func (*ElevenLabsRealtimeTTS) GetSampleRate added in v1.41.4

func (e *ElevenLabsRealtimeTTS) GetSampleRate() int

GetSampleRate returns the configured output sample rate

func (*ElevenLabsRealtimeTTS) SendText added in v1.41.4

func (e *ElevenLabsRealtimeTTS) SendText(text string) error

SendText sends a text chunk to the TTS WebSocket

type ElevenLabsSTT added in v1.41.4

type ElevenLabsSTT struct {
	// contains filtered or unexported fields
}

ElevenLabsSTT implements STTProvider using ElevenLabs Scribe API

func NewElevenLabsSTT added in v1.41.4

func NewElevenLabsSTT(cfg *config.STTConfig) (*ElevenLabsSTT, error)

NewElevenLabsSTT creates a new ElevenLabs STT provider

func (*ElevenLabsSTT) Close added in v1.41.4

func (e *ElevenLabsSTT) Close() error

Close cleans up resources

func (*ElevenLabsSTT) Transcribe added in v1.41.4

func (e *ElevenLabsSTT) Transcribe(audioData []byte, sampleRate int) (string, error)

Transcribe sends PCM16 audio to ElevenLabs Scribe API and returns transcript

type ElevenLabsTTS added in v1.41.4

type ElevenLabsTTS struct {
	// contains filtered or unexported fields
}

ElevenLabsTTS implements TTSProvider using ElevenLabs streaming API

func NewElevenLabsTTS added in v1.41.4

func NewElevenLabsTTS(cfg *config.TTSConfig) (*ElevenLabsTTS, error)

NewElevenLabsTTS creates a new ElevenLabs TTS provider

func (*ElevenLabsTTS) Close added in v1.41.4

func (e *ElevenLabsTTS) Close() error

Close cleans up resources

func (*ElevenLabsTTS) GetSampleRate added in v1.41.4

func (e *ElevenLabsTTS) GetSampleRate() int

GetSampleRate returns the configured output sample rate

func (*ElevenLabsTTS) Synthesize added in v1.41.4

func (e *ElevenLabsTTS) Synthesize(text string) (<-chan []byte, error)

Synthesize streams PCM16 audio chunks from ElevenLabs

func (*ElevenLabsTTS) SynthesizeFull added in v1.41.4

func (e *ElevenLabsTTS) SynthesizeFull(text string) ([]byte, error)

SynthesizeFull returns all audio as a single buffer

type ErrorData

type ErrorData struct {
	Code    string `json:"code"`
	Message string `json:"message"`
}

ErrorData contains error information

type EventType

type EventType string

EventType represents the type of unified event

const (
	EventTypeAudioDelta     EventType = "audio_delta"
	EventTypeAudioComplete  EventType = "audio_complete"
	EventTypeAudioInterrupt EventType = "audio_interrupt"
	EventTypeTranscript     EventType = "transcript"
	EventTypeToolCall       EventType = "tool_call"
	EventTypeToolResult     EventType = "tool_result"
	EventTypeTurnStart      EventType = "turn_start"
	EventTypeTurnEnd        EventType = "turn_end"
	EventTypeError          EventType = "error"
	EventTypeSessionCreated EventType = "session_created"
)

type GeminiRealtimeAdapter

type GeminiRealtimeAdapter struct {
	// contains filtered or unexported fields
}

GeminiRealtimeAdapter implements the Adapter interface for Gemini Live API

func NewGeminiRealtimeAdapter

func NewGeminiRealtimeAdapter(session *Session, messageSaver threads.MessageSaver,
	eventBus *events.EventBus) (*GeminiRealtimeAdapter, error)

NewGeminiRealtimeAdapter creates a new Gemini Live adapter

func (*GeminiRealtimeAdapter) Close

func (a *GeminiRealtimeAdapter) Close() error

Close closes the adapter

func (*GeminiRealtimeAdapter) GetResumptionToken

func (a *GeminiRealtimeAdapter) GetResumptionToken() string

GetResumptionToken returns the session resumption token (Gemini-specific)

func (*GeminiRealtimeAdapter) HandleControl

func (a *GeminiRealtimeAdapter) HandleControl(action string, data map[string]interface{}) error

HandleControl handles control messages

func (*GeminiRealtimeAdapter) InputSampleRate added in v1.47.0

func (a *GeminiRealtimeAdapter) InputSampleRate() int

InputSampleRate returns the expected input sample rate for Gemini (16kHz)

func (*GeminiRealtimeAdapter) ReceiveEvents

func (a *GeminiRealtimeAdapter) ReceiveEvents() <-chan UnifiedEvent

ReceiveEvents returns channel for receiving events

func (*GeminiRealtimeAdapter) SendAudio

func (a *GeminiRealtimeAdapter) SendAudio(base64Audio string) error

func (*GeminiRealtimeAdapter) SendText

func (a *GeminiRealtimeAdapter) SendText(text string) error

SendText sends text message to Gemini

type Message

type Message struct {
	Type      string      `json:"type"`       // "audio", "text", "control", event types
	SessionID string      `json:"session_id"` // Session identifier
	Data      interface{} `json:"data"`       // Payload (varies by type)
	Timestamp int64       `json:"timestamp"`  // Unix milliseconds
}

Message represents a unified message format for client-server communication

type OpenAIRealtimeAdapter

type OpenAIRealtimeAdapter struct {
	// contains filtered or unexported fields
}

OpenAIRealtimeAdapter implements the Adapter interface for OpenAI Realtime API

func NewOpenAIRealtimeAdapter

func NewOpenAIRealtimeAdapter(session *Session, messageSaver threads.MessageSaver,
	eventBus *events.EventBus) (*OpenAIRealtimeAdapter, error)

NewOpenAIRealtimeAdapter creates a new OpenAI Realtime adapter

func (*OpenAIRealtimeAdapter) Close

func (a *OpenAIRealtimeAdapter) Close() error

Close closes the adapter

func (*OpenAIRealtimeAdapter) HandleControl

func (a *OpenAIRealtimeAdapter) HandleControl(action string, data map[string]interface{}) error

HandleControl handles control messages

func (*OpenAIRealtimeAdapter) InputSampleRate added in v1.47.0

func (a *OpenAIRealtimeAdapter) InputSampleRate() int

InputSampleRate returns the expected input sample rate for OpenAI (24kHz)

func (*OpenAIRealtimeAdapter) ReceiveEvents

func (a *OpenAIRealtimeAdapter) ReceiveEvents() <-chan UnifiedEvent

ReceiveEvents returns channel for receiving events

func (*OpenAIRealtimeAdapter) SendAudio

func (a *OpenAIRealtimeAdapter) SendAudio(base64Audio string) error

SendAudio sends audio chunk to OpenAI

func (*OpenAIRealtimeAdapter) SendText

func (a *OpenAIRealtimeAdapter) SendText(text string) error

SendText sends text message to OpenAI

type RealtimeServer

type RealtimeServer struct {
	// contains filtered or unexported fields
}

RealtimeServer manages WebSocket connections for real-time voice communication

func NewServer

func NewServer(db *sql.DB, messageSaver threads.MessageSaver, eventBus *events.EventBus) *RealtimeServer

NewServer creates a new realtime server instance

func (*RealtimeServer) GetSessions

func (s *RealtimeServer) GetSessions() []*Session

GetSessions returns all active sessions

func (*RealtimeServer) HandleWebSocket

func (s *RealtimeServer) HandleWebSocket(w http.ResponseWriter, r *http.Request)

HandleWebSocket handles WebSocket connections for real-time communication Auto-detects format: JSON (browser/app) or Binary (telephony like Vonage/Twilio)

type STTProvider added in v1.41.4

type STTProvider interface {
	// Transcribe converts audio data (PCM16 little-endian) to text
	Transcribe(audioData []byte, sampleRate int) (string, error)
	Close() error
}

STTProvider transcribes audio to text (batch mode)

func NewSTTProvider added in v1.41.4

func NewSTTProvider(cfg *config.STTConfig) (STTProvider, error)

NewSTTProvider creates a batch STT provider based on config

type Session

type Session struct {
	ID       string
	ThreadID string
	Provider string // "openai-realtime"

	// Tool execution state
	PendingTools map[string]*ToolExecution

	State        SessionState
	CreatedAt    time.Time
	LastActivity time.Time
	// contains filtered or unexported fields
}

Session represents an active realtime voice session

func (*Session) AddPendingTool

func (s *Session) AddPendingTool(execution *ToolExecution)

AddPendingTool adds a tool execution to pending state

func (*Session) Close

func (s *Session) Close()

Close closes the session connections

func (*Session) GetPendingTool

func (s *Session) GetPendingTool(callID string) (*ToolExecution, bool)

GetPendingTool retrieves a pending tool execution

func (*Session) SetTurnState

func (s *Session) SetTurnState(state string)

SetTurnState sets the current turn state

func (*Session) UpdateLastActivity

func (s *Session) UpdateLastActivity()

UpdateLastActivity updates the last activity timestamp

func (*Session) UpdateToolStatus

func (s *Session) UpdateToolStatus(callID, status string)

UpdateToolStatus updates the status of a tool execution

type SessionState

type SessionState struct {
	IsActive          bool
	TurnState         string // "idle", "user_speaking", "assistant_speaking", "tool_executing"
	ConversationTurns int
	CurrentResponseID string
}

SessionState tracks the current state of a session

type StandardVoiceAdapter added in v1.41.4

type StandardVoiceAdapter struct {
	// contains filtered or unexported fields
}

StandardVoiceAdapter implements the Adapter interface using STT + Core LLM + TTS. This allows any LLM provider (Anthropic, OpenAI, Gemini, Fireworks, Groq, etc.) to be used for voice conversations without requiring a dedicated realtime API.

Full streaming pipeline:

Audio → WebSocket STT (realtime) → transcript
                                     ↓
                                  LLM (SSE stream)
                                     ↓ tokens
                                  WebSocket TTS (realtime) → audio chunks → client

func NewStandardVoiceAdapter added in v1.41.4

func NewStandardVoiceAdapter(session *Session, messageSaver threads.MessageSaver,
	eventBus *events.EventBus) (*StandardVoiceAdapter, error)

NewStandardVoiceAdapter creates a new standard voice adapter

func (*StandardVoiceAdapter) Close added in v1.41.4

func (a *StandardVoiceAdapter) Close() error

Close cleans up all resources

func (*StandardVoiceAdapter) HandleControl added in v1.41.4

func (a *StandardVoiceAdapter) HandleControl(action string, data map[string]interface{}) error

HandleControl handles control messages

func (*StandardVoiceAdapter) InputSampleRate added in v1.47.0

func (a *StandardVoiceAdapter) InputSampleRate() int

InputSampleRate returns the expected input sample rate for Standard mode (16kHz for ElevenLabs STT)

func (*StandardVoiceAdapter) ReceiveEvents added in v1.41.4

func (a *StandardVoiceAdapter) ReceiveEvents() <-chan UnifiedEvent

ReceiveEvents returns the event channel

func (*StandardVoiceAdapter) SendAudio added in v1.41.4

func (a *StandardVoiceAdapter) SendAudio(base64Audio string) error

SendAudio receives base64-encoded PCM16 audio and forwards directly to streaming STT

func (*StandardVoiceAdapter) SendText added in v1.41.4

func (a *StandardVoiceAdapter) SendText(text string) error

SendText sends a text message directly to the LLM (bypasses STT)

type StreamingSTTProvider added in v1.41.4

type StreamingSTTProvider interface {
	// SendAudio sends a PCM16 audio chunk to the STT service
	SendAudio(pcm16Data []byte) error
	// Transcripts returns a channel that receives committed (final) transcripts
	Transcripts() <-chan string
	// PartialTranscripts returns a channel that receives interim/partial transcripts
	PartialTranscripts() <-chan string
	Close() error
}

StreamingSTTProvider transcribes audio in real-time via streaming

func NewStreamingSTTProvider added in v1.41.4

func NewStreamingSTTProvider(cfg *config.STTConfig) (StreamingSTTProvider, error)

NewStreamingSTTProvider creates a streaming STT provider based on config

type StreamingTTSProvider added in v1.41.4

type StreamingTTSProvider interface {
	// SendText sends a text chunk to the TTS service (call multiple times as LLM streams)
	SendText(text string) error
	// Flush forces generation of any buffered text
	Flush() error
	// Finish signals end of text input and waits for all audio
	Finish() error
	// AudioChunks returns a channel that receives PCM16 audio chunks
	AudioChunks() <-chan []byte
	// GetSampleRate returns the output sample rate
	GetSampleRate() int
	Close() error
}

StreamingTTSProvider converts text to speech via streaming WebSocket. Text chunks can be sent progressively as LLM tokens arrive.

func NewStreamingTTSProvider added in v1.41.4

func NewStreamingTTSProvider(cfg *config.TTSConfig) (StreamingTTSProvider, error)

NewStreamingTTSProvider creates a streaming TTS provider based on config

type TTSProvider added in v1.41.4

type TTSProvider interface {
	// Synthesize converts text to PCM16 audio chunks streamed via channel.
	// The channel receives raw PCM16 little-endian byte slices.
	// Channel is closed when synthesis is complete.
	Synthesize(text string) (<-chan []byte, error)

	// SynthesizeFull converts text to a single PCM16 audio buffer (non-streaming)
	SynthesizeFull(text string) ([]byte, error)

	// GetSampleRate returns the output sample rate
	GetSampleRate() int

	Close() error
}

TTSProvider converts text to speech audio

func NewTTSProvider added in v1.41.4

func NewTTSProvider(cfg *config.TTSConfig) (TTSProvider, error)

NewTTSProvider creates a batch TTS provider based on config

type ToolCallData

type ToolCallData struct {
	ID    string                 `json:"id"`
	Name  string                 `json:"name"`
	Input map[string]interface{} `json:"input"`
}

ToolCallData contains tool invocation information

type ToolExecution

type ToolExecution struct {
	CallID    string
	Name      string
	Arguments map[string]interface{}
	Status    string // "pending", "executing", "completed", "failed"
	Result    interface{}
	Error     error
	StartTime time.Time
	EndTime   time.Time
}

ToolExecution represents a tool call in progress

type ToolResultData

type ToolResultData struct {
	CallID string      `json:"call_id"`
	Name   string      `json:"name"`
	Result interface{} `json:"result"`
	Error  error       `json:"error,omitempty"`
}

ToolResultData contains tool execution result

type TranscriptData

type TranscriptData struct {
	Role    string `json:"role"`              // "user" or "assistant"
	Content string `json:"content"`           // Transcribed text
	Partial bool   `json:"partial,omitempty"` // True for interim/partial transcripts
}

TranscriptData contains transcription information

type UnifiedEvent

type UnifiedEvent struct {
	Type      EventType   `json:"type"`
	SessionID string      `json:"session_id,omitempty"`
	Timestamp time.Time   `json:"timestamp"`
	Data      interface{} `json:"data"`
}

UnifiedEvent represents a provider-agnostic event

type WhisperSTT added in v1.41.4

type WhisperSTT struct {
	// contains filtered or unexported fields
}

WhisperSTT implements STTProvider using OpenAI Whisper API

func NewWhisperSTT added in v1.41.4

func NewWhisperSTT(cfg *config.STTConfig) (*WhisperSTT, error)

NewWhisperSTT creates a new Whisper STT provider

func (*WhisperSTT) Close added in v1.41.4

func (w *WhisperSTT) Close() error

Close cleans up resources

func (*WhisperSTT) Transcribe added in v1.41.4

func (w *WhisperSTT) Transcribe(audioData []byte, sampleRate int) (string, error)

Transcribe sends PCM16 audio to Whisper API and returns transcript

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL