OpenAI GPT-Live全双工语音系统深度解析:语音层与推理层彻底解耦

一、引言:语音AI的"寒武纪大爆发"

2026年7月8日,OpenAI正式发布GPT-Live系列语音模型,这是其第三代语音交互系统。随后在8月3日,OpenAI工程团队发布技术博客《How we built a realtime system for responsive voice AI in six months》,详细披露了GPT-Live的底层架构设计。8月4日,更多技术细节通过媒体渠道公开。

这套系统最核心的突破在于:彻底抛弃了语音AI领域沿用多年的"轮流说话"(turn-based)架构,转向原生全双工(full-duplex)架构,并将语音交互层与深度推理层在系统层面完全解耦。

本文将从工程实现角度,深入剖析GPT-Live的架构设计,涵盖全双工音频管道、语音-推理解耦设计、并发状态管理、协议优化,并提供完整的代码实现。

来源说明:本文技术细节主要基于OpenAI官方博客《How we built a realtime system for responsive voice AI in six months》(https://openai.com/index/continuous-voice-interaction-with-gpt-live/)及《Introducing GPT‑Live》(https://openai.com/index/introducing-gpt-live/)。


二、架构演进:从级联到全双工

2.1 第一代:级联系统(Cascaded System)

OpenAI最初的语音系统(2023年ChatGPT语音版)采用经典的三段式级联架构:

用户语音 → [STT(Whisper)] → 文本 → [LLM(GPT-4)] → 文本 → [TTS] → 语音回复

问题

  • 每段之间串行执行,延迟累积(典型端到端延迟>2秒)
  • 语音中的语调、节奏、情感信息在STT转录中丢失
  • 无法处理打断、语气词等自然对话行为

2.2 第二代:语音到语音 + 轮次检测(Turn-based S2S)

2024年发布的Advanced Voice Mode将STT和TTS合并到单一模型中,直接处理音频:

用户语音 → [语音模型(多模态)] → 语音回复
                ↑
         [轮次检测器] ← 静音检测

改进:延迟降低,保留了副语言信息(语调、语速等)。

遗留问题

  • 仍然依赖轮次检测器(turn detector)判断对话边界
  • 轮次检测器面临两难:猜早了打断用户,猜晚了反应迟钝
  • 本质仍是回合制,无法实现真正的并行对话

2.3 第三代:GPT-Live 全双工架构

GPT-Live彻底重构了语音交互的底层逻辑:

                           ┌─────────────────┐
         用户音频 ────────►│  GPT-Live-1     │◄──────── 用户音频
                           │  (全双工语音模型) │────────► 用户音频
         用户音频 ────────►│                  │
                           └────────┬────────┘
                                    │ 异步RPC
                                    ▼
                           ┌─────────────────┐
                           │  GPT-5.5        │
                           │  (前沿推理模型)   │
                           │  搜索 / 推理 /    │
                           │  工具调用        │
                           └─────────────────┘

关键创新

  1. 全双工音频:模型同时处理输入和输出音频流
  2. 移除轮次检测器:模型自身控制对话节奏
  3. 异步委派:语音模型可以将复杂任务委派给GPT-5.5
  4. 专用媒体路径:音频流走独立快速通道,与业务逻辑隔离

据OpenAI工程团队介绍,GPT-Live的架构决策是在过去六个月内将系统从Python asyncio迁移到Go语言,重写了模型推理、上下文管理和媒体传输层(来源:OpenAI工程博客)。


三、全双工音频架构深度解析

3.1 核心概念:全双工 vs 半双工

特性半双工(旧系统)全双工(GPT-Live)
同时收发
打断支持需等待当前轮次结束随时可打断
语气词处理困难(需单独处理)原生支持
对话节奏控制轮次检测器模型自身
延迟特征突发式(处理完一轮才响应)连续流式

3.2 音频流管道设计

GPT-Live的音频管道由以下几个关键组件构成:

[客户端麦克风]
      │
      ▼
[WebRTC音频轨道] ───► [抖动缓冲器(Jitter Buffer)]
      │                        │
      │                        ▼
      │               [音频预处理]
      │               ├── 回声消除(AEC)
      │               ├── 噪声抑制(NS)
      │               ├── 自动增益控制(AGC)
      │               └── 语音活动检测(VAD)
      │                        │
      │                        ▼
      │               [音频编码器] ◄── Opus编码
      │                        │
      │                        ▼
      │               [Go媒体前端] ────► [状态推理引擎]
      │                                          │
      │                                          ▼
      │                               [GPT-Live-1全双工模型]
      │                                          │
      │                               ┌──────────┴──────────┐
      │                               ▼                     ▼
      │                        [音频生成器]          [异步委派器]
      │                               │                     │
      │                               ▼                     ▼
      │                        [Opus解码]           [GPT-5.5推理]
      │                               │
      │                               ▼
      └─────────────────────── [WebRTC音频输出]

关键设计原则:媒体流与应用逻辑彻底分离。音频在客户端和语音模型之间的专用快速通道上移动,工具调用等应用工作发生在异步RPC边界之后。(来源:OpenAI工程博客)

3.3 代码实现:Python全双工音频管道

以下是一个简化的全双工音频管道实现,展示了WebRTC流处理、音频缓冲区管理和VAD检测的核心逻辑:

"""
全双工音频管道实现
模拟GPT-Live中的音频流处理核心逻辑
"""

import asyncio
import collections
import struct
import time
import wave
from dataclasses import dataclass, field
from enum import Enum
from typing import Optional, Callable, Awaitable


class AudioState(Enum):
    """音频流状态"""
    SILENCE = "silence"
    LISTENING = "listening"
    SPEAKING = "speaking"
    BARGE_IN = "barge_in"  # 用户打断


@dataclass
class AudioFrame:
    """音频帧数据结构"""
    timestamp: float
    data: bytes
    sample_rate: int = 16000
    channels: int = 1
    sample_width: int = 2  # 16-bit PCM
    
    @property
    def duration_seconds(self) -> float:
        return len(self.data) / (self.sample_rate * self.channels * self.sample_width)
    
    @property
    def rms(self) -> float:
        """计算音频帧的RMS能量"""
        if not self.data:
            return 0.0
        samples = struct.unpack_from(f"<{len(self.data)//2}h", self.data)
        sum_squares = sum(s * s for s in samples)
        return (sum_squares / len(samples)) ** 0.5


class AudioRingBuffer:
    """
    音频环形缓冲区
    用于管理持续流入的音频数据,支持滑动窗口读取
    """
    
    def __init__(self, max_duration_seconds: float = 30.0, sample_rate: int = 16000):
        self.sample_rate = sample_rate
        self.max_samples = int(max_duration_seconds * sample_rate)
        self.buffer = collections.deque(maxlen=self.max_samples)
        self._lock = asyncio.Lock()
    
    async def write(self, samples: list[int]) -> None:
        """写入音频采样数据"""
        async with self._lock:
            self.buffer.extend(samples)
    
    async def read(self, num_samples: int) -> list[int]:
        """读取最近N个采样点"""
        async with self._lock:
            if len(self.buffer) < num_samples:
                return list(self.buffer)
            return list(self.buffer)[-num_samples:]
    
    async def clear(self) -> None:
        async with self._lock:
            self.buffer.clear()
    
    @property
    async def duration_seconds(self) -> float:
        async with self._lock:
            return len(self.buffer) / self.sample_rate


class VoiceActivityDetector:
    """
    语音活动检测器(VAD)
    基于能量水平和自适应阈值,模拟GPT-Live中移除独立turn detector后的
    模型内建VAD逻辑
    """
    
    def __init__(
        self,
        frame_duration_ms: int = 30,
        silence_ratio_threshold: float = 0.3,
        min_speech_frames: int = 3,
        min_silence_frames: int = 15,  # ~450ms静音视为停顿
        energy_threshold: float = 100.0,
    ):
        self.frame_duration_ms = frame_duration_ms
        self.silence_ratio_threshold = silence_ratio_threshold
        self.min_speech_frames = min_speech_frames
        self.min_silence_frames = min_silence_frames
        self.energy_threshold = energy_threshold
        
        self._speech_frames = 0
        self._silence_frames = 0
        self._is_speech = False
        self._recent_frames: list[bool] = []
    
    def is_speech_frame(self, frame: AudioFrame) -> bool:
        """判断单帧是否为语音"""
        return frame.rms > self.energy_threshold
    
    def process_frame(self, frame: AudioFrame) -> AudioState:
        """
        处理音频帧并返回当前状态
        GPT-Live的核心创新:不再依赖外部turn detector做二元判决,
        而是让模型在每一帧都做出交互决策
        """
        is_speech = self.is_speech_frame(frame)
        self._recent_frames.append(is_speech)
        
        # 保持最近N帧用于平滑判断
        window_size = 10
        if len(self._recent_frames) > window_size:
            self._recent_frames.pop(0)
        
        speech_ratio = sum(self._recent_frames) / len(self._recent_frames)
        
        if speech_ratio > self.silence_ratio_threshold:
            self._speech_frames += 1
            self._silence_frames = 0
        else:
            self._silence_frames += 1
            self._speech_frames = 0
        
        # 状态转换逻辑
        if not self._is_speech and self._speech_frames >= self.min_speech_frames:
            self._is_speech = True
            return AudioState.LISTENING
        
        if self._is_speech and self._silence_frames >= self.min_silence_frames:
            self._is_speech = False
            return AudioState.SILENCE
        
        return AudioState.LISTENING if self._is_speech else AudioState.SILENCE
    
    def reset(self) -> None:
        self._speech_frames = 0
        self._silence_frames = 0
        self._is_speech = False
        self._recent_frames.clear()


class EchoCanceller:
    """
    简化的回声消除器
    使用自适应NLMS滤波器的概念性实现
    """
    
    def __init__(self, filter_length: int = 512, mu: float = 0.01):
        self.filter_length = filter_length
        self.mu = mu
        self.weights = [0.0] * filter_length
        self.reference_buffer: list[float] = [0.0] * filter_length
    
    def process(self, mic_signal: list[float], ref_signal: list[float]) -> list[float]:
        """
        处理麦克风信号,减去参考信号中的回声分量
        mic_signal: 麦克风采集信号
        ref_signal: 扬声器播放的参考信号
        """
        output = []
        for n in range(len(mic_signal)):
            # 更新参考缓冲区
            self.reference_buffer.pop(0)
            self.reference_buffer.append(ref_signal[n] if n < len(ref_signal) else 0.0)
            
            # 估计回声
            echo_estimate = sum(
                self.weights[i] * self.reference_buffer[self.filter_length - 1 - i]
                for i in range(self.filter_length)
            )
            
            # 误差信号(消除回声后的信号)
            error = mic_signal[n] - echo_estimate
            
            # 更新滤波器权重(NLMS自适应)
            norm = sum(x * x for x in self.reference_buffer)
            if norm > 1e-10:
                for i in range(self.filter_length):
                    self.weights[i] += (
                        self.mu * error * self.reference_buffer[self.filter_length - 1 - i] / norm
                    )
            
            output.append(error)
        
        return output


class FullDuplexAudioPipeline:
    """
    全双工音频管道(完整实现)
    模拟GPT-Live的音频流处理核心
    """
    
    def __init__(
        self,
        sample_rate: int = 16000,
        frame_duration_ms: int = 20,  # 每帧20ms,与Opus编码一致
        vad: Optional[VoiceActivityDetector] = None,
        echo_canceller: Optional[EchoCanceller] = None,
    ):
        self.sample_rate = sample_rate
        self.frame_duration_ms = frame_duration_ms
        self.frame_size = int(sample_rate * frame_duration_ms / 1000)
        
        self.vad = vad or VoiceActivityDetector()
        self.echo_canceller = echo_canceller or EchoCanceller()
        
        self.input_buffer = AudioRingBuffer(max_duration_seconds=60.0)
        self.output_buffer = AudioRingBuffer(max_duration_seconds=10.0)
        
        self.state = AudioState.SILENCE
        self._running = False
        self._on_audio_frame: Optional[Callable[[AudioFrame], Awaitable[None]]] = None
        self._on_state_change: Optional[Callable[[AudioState], Awaitable[None]]] = None
    
    async def push_input_frame(self, frame: AudioFrame) -> None:
        """
        处理输入音频帧(来自麦克风)
        这是全双工架构的核心:输入处理不阻塞输出
        """
        # 1. 回声消除
        ref_samples = await self.output_buffer.read(self.frame_size)
        if ref_samples:
            mic_samples = struct.unpack_from(f"<{len(frame.data)//2}h", frame.data)
            cleaned = self.echo_canceller.process(
                list(mic_samples), ref_samples
            )
            frame.data = struct.pack(f"<{len(cleaned)}h", *[int(s) for s in cleaned])
        
        # 2. 写入输入缓冲区
        samples = struct.unpack_from(f"<{len(frame.data)//2}h", frame.data)
        await self.input_buffer.write(list(samples))
        
        # 3. VAD检测(GPT-Live中此逻辑内建于模型,此处为模拟)
        new_state = self.vad.process_frame(frame)
        if new_state != self.state:
            self.state = new_state
            if self._on_state_change:
                await self._on_state_change(new_state)
        
        # 4. 回调通知
        if self._on_audio_frame:
            await self._on_audio_frame(frame)
    
    async def push_output_frame(self, frame: AudioFrame) -> None:
        """
        处理输出音频帧(来自模型生成的语音)
        也写入输出缓冲区供回声消除使用
        """
        samples = struct.unpack_from(f"<{len(frame.data)//2}h", frame.data)
        await self.output_buffer.write(list(samples))
    
    async def start(self) -> None:
        self._running = True
        self.state = AudioState.LISTENING
    
    async def stop(self) -> None:
        self._running = False
        await self.input_buffer.clear()
        await self.output_buffer.clear()
    
    def on_audio_frame(self, callback: Callable[[AudioFrame], Awaitable[None]]) -> None:
        self._on_audio_frame = callback
    
    def on_state_change(self, callback: Callable[[AudioState], Awaitable[None]]) -> None:
        self._on_state_change = callback

3.4 WebRTC传输层优化

OpenAI团队在传输层做了大量优化工作。核心决策是用Go重写媒体前端和推理逻辑,替换了之前的Python asyncio实现。新系统的p95延迟直接追平了旧系统的p50延迟(来源:OpenAI工程博客)。

WebRTC的天然优势:

  • 支持丢包恢复(通过FEC和RTX)
  • 时钟漂移补偿(NTP时间同步)
  • 网络变化自适应(通过拥塞控制算法)
  • 音频拉伸/压缩(解决丢包和延迟抖动)

下面是用Go实现WebRTC音频轨道的核心逻辑:

// 全双工WebRTC音频会话管理
// 模拟GPT-Live中Go语言实现的媒体前端

package media

import (
	"context"
	"encoding/binary"
	"io"
	"log"
	"sync"
	"time"

	"github.com/pion/webrtc/v4"
	"github.com/pion/rtp"
)

// AudioCodec 定义音频编解码器参数
type AudioCodec struct {
	SampleRate   int
	Channels     int
	FrameSize    int // 每帧采样数
	Bitrate      int
}

// DefaultOpusCodec Opus默认配置(GPT-Live使用的编码器)
var DefaultOpusCodec = AudioCodec{
	SampleRate: 48000,
	Channels:   1,
	FrameSize:  960, // 20ms @ 48kHz
	Bitrate:    32000,
}

// AudioFrame 音频帧
type AudioFrame struct {
	Timestamp time.Time
	Sequence  uint16
	Data      []byte
	Duration  time.Duration
}

// AudioTrack 封装WebRTC音频轨道
type AudioTrack struct {
	track      *webrtc.TrackLocalStaticSample
	codec      AudioCodec
	seqNum     uint16
	mu         sync.Mutex
	onFrame    func(AudioFrame)
	writeCh    chan AudioFrame
	ctx        context.Context
	cancel     context.CancelFunc
}

// NewAudioTrack 创建新的音频轨道
func NewAudioTrack(codec AudioCodec) (*AudioTrack, error) {
	opusCodec := webrtc.RTPCodecCapability{
		MimeType:    webrtc.MimeTypeOpus,
		ClockRate:   uint32(codec.SampleRate),
		Channels:    uint16(codec.Channels),
	}

	track, err := webrtc.NewTrackLocalStaticSample(
		opusCodec, "audio", "voice-stream",
	)
	if err != nil {
		return nil, err
	}

	ctx, cancel := context.WithCancel(context.Background())
	return &AudioTrack{
		track:   track,
		codec:   codec,
		seqNum:  0,
		writeCh: make(chan AudioFrame, 256), // 缓冲256帧
		ctx:     ctx,
		cancel:  cancel,
	}, nil
}

// WriteFrame 写入音频帧(非阻塞,由后台goroutine发送)
func (at *AudioTrack) WriteFrame(frame AudioFrame) error {
	select {
	case at.writeCh <- frame:
		return nil
	default:
		// 缓冲区满,丢弃最旧的帧(保持实时性)
		select {
		case <-at.writeCh:
			// 丢弃一帧
		default:
		}
		at.writeCh <- frame
		return nil
	}
}

// Start 启动音频发送循环
func (at *AudioTrack) Start() {
	go func() {
		ticker := time.NewTicker(20 * time.Millisecond) // 每20ms发送一帧
		defer ticker.Stop()

		for {
			select {
			case <-at.ctx.Done():
				return
			case frame := <-at.writeCh:
				at.mu.Lock()
				at.seqNum++
				seq := at.seqNum
				at.mu.Unlock()

				// 构建RTP包
				rtpPacket := &rtp.Packet{
					Header: rtp.Header{
						Version:        2,
						PayloadType:    111, // dynamic payload type for Opus
						SequenceNumber: seq,
						Timestamp:      uint32(frame.Timestamp.UnixMicro() * at.codec.SampleRate / 1_000_000),
						SSRC:           0xABCDEF01,
						Marker:         false,
					},
					Payload: frame.Data,
				}

				raw, err := rtpPacket.Marshal()
				if err != nil {
					log.Printf("failed to marshal RTP packet: %v", err)
					continue
				}

				if _, err := at.track.Write(raw); err != nil {
					log.Printf("failed to write to track: %v", err)
				}

			case <-ticker.C:
				// 心跳:确保连接活跃
				// GPT-Live使用持续流确保无间隙
			}
		}
	}()
}

// Stop 停止音频轨道
func (at *AudioTrack) Stop() {
	at.cancel()
}

// MediaSession 全双工媒体会话
type MediaSession struct {
	pc         *webrtc.PeerConnection
	inputTrack *webrtc.TrackRemote
	outputTrack *AudioTrack
	clientID   string
	onFrame    func(AudioFrame)
	mu         sync.RWMutex
}

// NewMediaSession 创建新的全双工媒体会话
func NewMediaSession(clientID string, api *webrtc.API) (*MediaSession, error) {
	config := webrtc.Configuration{
		ICEServers: []webrtc.ICEServer{
			{
				URLs: []string{"stun:stun.l.google.com:19302"},
			},
		},
	}

	pc, err := api.NewPeerConnection(config)
	if err != nil {
		return nil, err
	}

	session := &MediaSession{
		pc:       pc,
		clientID: clientID,
	}

	// 设置输入轨道处理
	pc.OnTrack(func(track *webrtc.TrackRemote, receiver *webrtc.RTPReceiver) {
		session.mu.Lock()
		session.inputTrack = track
		session.mu.Unlock()

		go session.handleIncomingTrack(track)
	})

	// 创建输出轨道
	outputTrack, err := NewAudioTrack(DefaultOpusCodec)
	if err != nil {
		pc.Close()
		return nil, err
	}

	if _, err := pc.AddTrack(outputTrack.track); err != nil {
		pc.Close()
		return nil, err
	}

	session.outputTrack = outputTrack
	return session, nil
}

// handleIncomingTrack 处理输入音频流
func (ms *MediaSession) handleIncomingTrack(track *webrtc.TrackRemote) {
	// 抖动缓冲区:200ms的抖动容限
	jitterBuffer := NewJitterBuffer(200 * time.Millisecond)

	for {
		rtpPacket, _, err := track.ReadRTP()
		if err != nil {
			if err == io.EOF {
				return
			}
			log.Printf("read RTP error: %v", err)
			continue
		}

		frame := AudioFrame{
			Timestamp: time.Now(),
			Sequence:  rtpPacket.Header.SequenceNumber,
			Data:      rtpPacket.Payload,
			Duration:  20 * time.Millisecond,
		}

		// 通过抖动缓冲区处理
		jitterBuffer.Push(frame)
	}
}

// Start 启动媒体会话
func (ms *MediaSession) Start() {
	ms.outputTrack.Start()
}

// Stop 停止媒体会话
func (ms *MediaSession) Stop() {
	ms.outputTrack.Stop()
	ms.pc.Close()
}

// JitterBuffer 抖动缓冲区
// 用于平滑网络延迟波动,GPT-Live中同样使用此机制
type JitterBuffer struct {
	buffer    []AudioFrame
	maxSize   int
	cond      *sync.Cond
	mu        sync.Mutex
}

func NewJitterBuffer(maxDuration time.Duration) *JitterBuffer {
	// 假设每帧20ms
	maxSize := int(maxDuration / (20 * time.Millisecond))
	return &JitterBuffer{
		buffer:  make([]AudioFrame, 0, maxSize),
		maxSize: maxSize,
		cond:    sync.NewCond(&sync.Mutex{}),
	}
}

func (jb *JitterBuffer) Push(frame AudioFrame) {
	jb.mu.Lock()
	defer jb.mu.Unlock()

	if len(jb.buffer) >= jb.maxSize {
		// 缓冲区满,丢弃最旧的帧
		jb.buffer = jb.buffer[1:]
	}
	jb.buffer = append(jb.buffer, frame)
	jb.cond.Signal()
}

func (jb *JitterBuffer) Pop() (AudioFrame, bool) {
	jb.mu.Lock()
	defer jb.mu.Unlock()

	for len(jb.buffer) == 0 {
		jb.cond.Wait()
	}

	frame := jb.buffer[0]
	jb.buffer = jb.buffer[1:]
	return frame, true
}

四、语音层与推理层解耦架构

4.1 双模型架构设计

GPT-Live最关键的架构创新是将"说话"(conversational interaction)与"思考"(deep reasoning)在系统层面彻底分离。这通过两个独立的模型实现:

  • GPT-Live-1(前台):全双工语音模型,负责低延迟、高自然的对话交互
  • GPT-5.5(后台):前沿推理模型,负责搜索、复杂推理、工具调用

这两层之间的交互通过异步RPC边界完成,核心媒体路径永远不会被后台推理阻塞。

4.2 架构图

┌─────────────────────────────────────────────────────────────────┐
│                        GPT-Live 系统架构                           │
│                                                                   │
│  ┌──────────────────────────────────────────────────────┐        │
│  │              实时媒体路径(Live Path)                  │        │
│  │                                                      │        │
│  │  ┌────────┐    ┌──────────┐    ┌────────────────┐   │        │
│  │  │客户端  │◄──►│ Go媒体   │◄──►│ GPT-Live-1     │   │        │
│  │  │WebRTC │    │ 前端     │    │ 全双工语音模型   │   │        │
│  │  └────────┘    └──────────┘    └───────┬────────┘   │        │
│  │                                        │            │        │
│  └────────────────────────────────────────┼────────────┘        │
│                                           │                     │
│                                 异步RPC边界 │                     │
│                                           ▼                     │
│  ┌──────────────────────────────────────────────────────┐        │
│  │              异步委派路径(Delegation Path)           │        │
│  │                                                      │        │
│  │  ┌──────────────────────────────────────────────┐   │        │
│  │  │        推理调度器 (Go)                        │   │        │
│  │  │  ┌─────────┐  ┌─────────┐  ┌─────────┐     │   │        │
│  │  │  │ 会话管理 │  │ 上下文  │  │ 结果    │     │   │        │
│  │  │  │         │  │ 缓存    │  │ 合并    │     │   │        │
│  │  │  └─────────┘  └─────────┘  └─────────┘     │   │        │
│  │  └──────────────────────┬───────────────────────┘   │        │
│  │                         │                            │        │
│  │  ┌──────────────────────▼───────────────────────┐   │        │
│  │  │           GPT-5.5 推理实例                     │   │        │
│  │  │  ┌──────────┐ ┌──────────┐ ┌──────────┐     │   │        │
│  │  │  │ 搜索     │ │ 推理     │ │ 工具调用  │     │   │        │
│  │  │  │ (Web)    │ │ (思维链)  │ │ (Function)│     │   │        │
│  │  │  └──────────┘ └──────────┘ └──────────┘     │   │        │
│  │  └──────────────────────────────────────────────┘   │        │
│  └──────────────────────────────────────────────────────┘        │
└─────────────────────────────────────────────────────────────────┘

4.3 异步委派调度器实现(Go)

以下是GPT-Live中异步推理调度器的核心实现:

// 异步推理调度器
// 管理前台语音模型与后台推理模型的并发任务调度

package scheduler

import (
	"context"
	"fmt"
	"log"
	"sync"
	"time"
)

// DelegationRequest 委派请求
type DelegationRequest struct {
	ID        string
	SessionID string
	Query     string
	Context   *ConversationContext
	Tools     []ToolDefinition
	Effort    ReasoningEffort
	CreatedAt time.Time
	ResultCh  chan *DelegationResult
}

// DelegationResult 委派结果
type DelegationResult struct {
	RequestID   string
	Content     string
	ToolResults []ToolResult
	TokenUsage  TokenUsage
	Latency     time.Duration
	Error       error
}

// ReasoningEffort 推理强度
type ReasoningEffort int

const (
	EffortInstant ReasoningEffort = iota // 即时响应
	EffortMedium                         // 中等推理
	EffortHigh                           // 深度推理
)

// TokenUsage token用量统计
type TokenUsage struct {
	PromptTokens     int
	CompletionTokens int
	TotalTokens      int
}

// ToolDefinition 工具定义
type ToolDefinition struct {
	Name        string
	Description string
	Parameters  map[string]interface{}
}

// ToolResult 工具调用结果
type ToolResult struct {
	ToolName string
	Content  string
	Success  bool
}

// ConversationContext 会话上下文
type ConversationContext struct {
	Messages     []Message
	TokenCount   int
	LastActivity time.Time
}

// Message 对话消息
type Message struct {
	Role    string
	Content string
	AudioID string
	Time    time.Time
}

// SessionState 会话状态
type SessionState struct {
	ID              string
	VoiceModelID    string
	FrontierModelID string
	Context         *ConversationContext
	PendingRequests map[string]*DelegationRequest
	CreatedAt       time.Time
	mu              sync.RWMutex
}

// InferenceScheduler 推理调度器
// 核心职责:管理前台语音模型与后台推理模型之间的异步任务调度
type InferenceScheduler struct {
	mu            sync.RWMutex
	sessions      map[string]*SessionState
	frontierModel FrontierModelClient
	promptCache   *PromptCache
	workerPool    *WorkerPool
	maxRetries    int
}

// FrontierModelClient 前沿模型客户端接口
type FrontierModelClient interface {
	Infer(ctx context.Context, req *DelegationRequest) (*DelegationResult, error)
	PrefillContext(ctx context.Context, sessionID string, context *ConversationContext) error
	HealthCheck(ctx context.Context) bool
}

// NewInferenceScheduler 创建推理调度器
func NewInferenceScheduler(
	frontier FrontierModelClient,
	maxWorkers int,
	cacheSize int,
) *InferenceScheduler {
	return &InferenceScheduler{
		sessions:      make(map[string]*SessionState),
		frontierModel: frontier,
		promptCache:   NewPromptCache(cacheSize),
		workerPool:    NewWorkerPool(maxWorkers),
		maxRetries:    3,
	}
}

// RegisterSession 注册新会话
// 在语音会话启动时预创建推理会话并预填充上下文
func (s *InferenceScheduler) RegisterSession(
	sessionID string,
	initialContext *ConversationContext,
) error {
	s.mu.Lock()
	defer s.mu.Unlock()

	state := &SessionState{
		ID:              sessionID,
		Context:         initialContext,
		PendingRequests: make(map[string]*DelegationRequest),
		CreatedAt:       time.Now(),
	}

	// 预填充前沿模型的上下文(GPT-Live的关键优化)
	// 确保在第一次委派请求前,prompt已处理完毕
	ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
	defer cancel()

	if err := s.frontierModel.PrefillContext(ctx, sessionID, initialContext); err != nil {
		log.Printf("failed to prefill context for session %s: %v", sessionID, err)
		// 非致命错误,继续注册
	}

	s.sessions[sessionID] = state
	log.Printf("registered session %s with prefilled frontier context", sessionID)
	return nil
}

// DelegateTask 异步委派任务
// 这是GPT-Live的核心操作:语音模型将复杂任务委派给后台推理模型
// 返回的channels允许语音模型在等待结果时继续对话
func (s *InferenceScheduler) DelegateTask(
	sessionID string,
	query string,
	tools []ToolDefinition,
	effort ReasoningEffort,
) (<-chan *DelegationResult, error) {
	s.mu.RLock()
	session, exists := s.sessions[sessionID]
	s.mu.RUnlock()

	if !exists {
		return nil, fmt.Errorf("session %s not found", sessionID)
	}

	req := &DelegationRequest{
		ID:        generateRequestID(),
		SessionID: sessionID,
		Query:     query,
		Context:   session.Context,
		Tools:     tools,
		Effort:    effort,
		CreatedAt: time.Now(),
		ResultCh:  make(chan *DelegationResult, 1),
	}

	session.mu.Lock()
	session.PendingRequests[req.ID] = req
	session.mu.Unlock()

	// 异步提交到worker池
	s.workerPool.Submit(func() {
		result := s.executeDelegationWithRetry(req)
		
		session.mu.Lock()
		delete(session.PendingRequests, req.ID)
		session.mu.Unlock()
		
		req.ResultCh <- result
		close(req.ResultCh)
	})

	return req.ResultCh, nil
}

// executeDelegationWithRetry 带重试的委派执行
func (s *InferenceScheduler) executeDelegationWithRetry(req *DelegationRequest) *DelegationResult {
	var lastErr error
	
	for attempt := 0; attempt <= s.maxRetries; attempt++ {
		if attempt > 0 {
			// 指数退避
			backoff := time.Duration(100*attempt) * time.Millisecond
			time.Sleep(backoff)
		}

		ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
		start := time.Now()
		
		result, err := s.frontierModel.Infer(ctx, req)
		cancel()
		
		if err == nil {
			result.Latency = time.Since(start)
			return result
		}
		
		lastErr = err
		log.Printf("delegation attempt %d failed for request %s: %v",
			attempt+1, req.ID, err)
	}

	return &DelegationResult{
		RequestID: req.ID,
		Error:     fmt.Errorf("all retries failed: %v", lastErr),
	}
}

// GetSessionContext 获取会话上下文
func (s *InferenceScheduler) GetSessionContext(sessionID string) *ConversationContext {
	s.mu.RLock()
	defer s.mu.RUnlock()
	
	if session, exists := s.sessions[sessionID]; exists {
		return session.Context
	}
	return nil
}

// UpdateSessionContext 更新会话上下文
// 处理后台任务结果与用户实时修正的同步
func (s *InferenceScheduler) UpdateSessionContext(
	sessionID string,
	messages []Message,
) error {
	s.mu.Lock()
	defer s.mu.Unlock()

	session, exists := s.sessions[sessionID]
	if !exists {
		return fmt.Errorf("session %s not found", sessionID)
	}

	session.mu.Lock()
	defer session.mu.Unlock()

	session.Context.Messages = append(session.Context.Messages, messages...)
	session.Context.TokenCount = calculateTokenCount(session.Context.Messages)
	session.Context.LastActivity = time.Now()

	return nil
}

// UnregisterSession 注销会话
func (s *InferenceScheduler) UnregisterSession(sessionID string) {
	s.mu.Lock()
	defer s.mu.Unlock()

	if session, exists := s.sessions[sessionID]; exists {
		// 清理所有待处理请求
		session.mu.Lock()
		for _, req := range session.PendingRequests {
			close(req.ResultCh)
		}
		session.mu.Unlock()
		
		delete(s.sessions, sessionID)
		log.Printf("unregistered session %s", sessionID)
	}
}

// WorkerPool goroutine工作池
type WorkerPool struct {
	workers chan struct{}
	wg      sync.WaitGroup
}

func NewWorkerPool(maxWorkers int) *WorkerPool {
	return &WorkerPool{
		workers: make(chan struct{}, maxWorkers),
	}
}

func (wp *WorkerPool) Submit(task func()) {
	wp.workers <- struct{}{}
	wp.wg.Add(1)
	go func() {
		defer func() {
			<-wp.workers
			wp.wg.Done()
		}()
		task()
	}()
}

func (wp *WorkerPool) Wait() {
	wp.wg.Wait()
}

// PromptCache 提示词缓存
type PromptCache struct {
	mu    sync.RWMutex
	items map[string]*CachedPrompt
	size  int
}

type CachedPrompt struct {
	Prefix    string
	KVState   []byte // 缓存的KV cache状态
	CreatedAt time.Time
	HitCount  int
}

func NewPromptCache(size int) *PromptCache {
	return &PromptCache{
		items: make(map[string]*CachedPrompt),
		size:  size,
	}
}

func (pc *PromptCache) Get(key string) (*CachedPrompt, bool) {
	pc.mu.RLock()
	defer pc.mu.RUnlock()
	item, ok := pc.items[key]
	if ok {
		item.HitCount++
	}
	return item, ok
}

func (pc *PromptCache) Set(key string, item *CachedPrompt) {
	pc.mu.Lock()
	defer pc.mu.Unlock()
	
	if len(pc.items) >= pc.size {
		// LRU驱逐:淘汰命中次数最少的
		var minKey string
		minHits := int(^uint(0) >> 1)
		for k, v := range pc.items {
			if v.HitCount < minHits {
				minHits = v.HitCount
				minKey = k
			}
		}
		delete(pc.items, minKey)
	}
	
	pc.items[key] = item
}

// 辅助函数
func generateRequestID() string {
	return fmt.Sprintf("req_%d_%d", time.Now().UnixNano(), randInt(1000, 9999))
}

func randInt(min, max int) int {
	return min + int(time.Now().UnixNano()%int64(max-min+1))
}

func calculateTokenCount(messages []Message) int {
	total := 0
	for _, msg := range messages {
		total += len(msg.Content) / 4 // 粗略估算
	}
	return total
}

4.4 延迟优化策略

GPT-Live在委派路径上做了多层优化,确保后台推理结果能快速融入对话:

1. 预填充推理会话

当语音会话启动时,应用服务器会同时为前沿模型创建一个推理会话,预先填充初始上下文。这确保在第一次发起委派请求前,提示词已经处理完毕。推理会话在整个语音通话期间保持可用状态(来源:OpenAI工程博客)。

2. 稳定会话亲和性(Session Affinity)

连续请求绑定到同一推理实例,利用prompt caching技术,大大缩短后续请求的延迟。

3. 推理强度分级

用户可选择Instant(即时)、Medium(中等)、High(高)三种推理强度,对应GPT-5.5 Instant和GPT-5.5 Thinking两种模型配置。

// 推理强度配置
func (s *InferenceScheduler) getEffortConfig(effort ReasoningEffort) EffortConfig {
	switch effort {
	case EffortInstant:
		return EffortConfig{
			ModelName:       "gpt-5.5-instant",
			MaxTokens:       512,
			ReasoningEffort: 0,   // 无额外推理
			Temperature:     0.7,
			Timeout:         10 * time.Second,
		}
	case EffortMedium:
		return EffortConfig{
			ModelName:       "gpt-5.5-thinking",
			MaxTokens:       2048,
			ReasoningEffort: 1,   // 中等推理
			Temperature:     0.5,
			Timeout:         30 * time.Second,
		}
	case EffortHigh:
		return EffortConfig{
			ModelName:       "gpt-5.5-thinking",
			MaxTokens:       4096,
			ReasoningEffort: 2,   // 深度推理
			Temperature:     0.3,
			Timeout:         60 * time.Second,
		}
	default:
		return EffortConfig{}
	}
}

type EffortConfig struct {
	ModelName       string
	MaxTokens       int
	ReasoningEffort int
	Temperature     float64
	Timeout         time.Duration
}

五、状态管理与上下文压缩

5.1 有状态推理的挑战

语音会话可能持续很长时间(数十分钟甚至更长),上下文不断增长,模型实例也会根据需求启动或关闭。GPT-Live需要解决三个核心问题:

  1. 上下文超限:累积的上下文超出模型窗口
  2. KV缓存失效:上下文压缩后,注意力键值缓存需要重建
  3. 无缝切换:模型实例切换不能影响音频流

5.2 无缝交接机制

OpenAI团队构建了跨模型实例的无缝交接机制:

时间线:
┌──────────────────────────────────────────────────────────────┐
│ 模型实例A(运行中)                                          │
│ ├── 持续处理音频流                                           │
│ ├── 上下文逐渐增长                                           │
│ └── 触发上下文压缩条件                                       │
│                                                              │
│  ┌─ 后台操作 ───────────────────────────────────────┐       │
│  │ ① 使用当前上下文创建压缩后的新上下文               │       │
│  │ ② 预热模型实例B,预填充压缩后的上下文              │       │
│  │ ③ 实例A和实例B并行运行推理                         │       │
│  │ ④ 实例B完全就绪后,音频流无缝切换到实例B           │       │
│  └──────────────────────────────────────────────────┘       │
│                                                              │
│ 模型实例B(接管)                                            │
│ ├── 继续处理音频流(无中断)                                 │
│ └── 上下文已压缩,可继续增长                                │
└──────────────────────────────────────────────────────────────┘

5.3 上下文管理器实现

// 上下文管理器
// 负责会话上下文的压缩、交接和状态同步

package context

import (
	"context"
	"log"
	"sync"
	"time"
)

// ContextState 上下文状态
type ContextState int

const (
	StateActive    ContextState = iota // 活跃状态
	StateCompacting                    // 正在压缩
	StateReady                         // 已准备好新实例
	StateSwitching                     // 正在切换
)

// SessionContext 会话上下文
type SessionContext struct {
	ID              string
	Messages        []Message
	TokenCount      int
	KVState         []byte          // KV缓存状态
	LastActivity    time.Time
	State           ContextState
	StateChangedAt  time.Time
	
	mu              sync.RWMutex
}

// Message 对话消息
type Message struct {
	Role    string `json:"role"`
	Content string `json:"content"`
	Time    time.Time `json:"time"`
}

// ContextManager 上下文管理器
type ContextManager struct {
	mu               sync.RWMutex
	sessions         map[string]*SessionContext
	maxContextTokens int           // 最大上下文token数
	compactionRatio  float64       // 压缩比例(如0.5表示压缩到一半)
	compactionFunc   CompactionFunc
}

// CompactionFunc 上下文压缩函数
type CompactionFunc func(messages []Message, maxTokens int) ([]Message, int)

// NewContextManager 创建上下文管理器
func NewContextManager(
	maxContextTokens int,
	compactionRatio float64,
	compactionFunc CompactionFunc,
) *ContextManager {
	if compactionFunc == nil {
		compactionFunc = defaultCompactionFunc
	}
	
	return &ContextManager{
		sessions:         make(map[string]*SessionContext),
		maxContextTokens: maxContextTokens,
		compactionRatio:  compactionRatio,
		compactionFunc:   compactionFunc,
	}
}

// defaultCompactionFunc 默认压缩函数
// 策略:保留系统提示和最近的消息,对中间历史进行摘要
func defaultCompactionFunc(messages []Message, maxTokens int) ([]Message, int) {
	if len(messages) == 0 {
		return messages, 0
	}
	
	// 保留系统提示(第一条)
	systemMsg := messages[0]
	
	// 保留最近N条消息
	keepRecent := 10
	if keepRecent >= len(messages) {
		keepRecent = len(messages) - 1
	}
	
	compacted := make([]Message, 0, keepRecent+2)
	compacted = append(compacted, systemMsg)
	
	// 中间消息摘要
	midMessages := messages[1 : len(messages)-keepRecent]
	if len(midMessages) > 0 {
		summary := summarizeMessages(midMessages)
		compacted = append(compacted, Message{
			Role:    "system",
			Content: summary,
			Time:    time.Now(),
		})
	}
	
	// 追加最近消息
	compacted = append(compacted, messages[len(messages)-keepRecent:]...)
	
	totalTokens := 0
	for _, msg := range compacted {
		totalTokens += len(msg.Content) / 4
	}
	
	return compacted, totalTokens
}

// summarizeMessages 消息摘要(简化实现)
func summarizeMessages(messages []Message) string {
	if len(messages) == 0 {
		return ""
	}
	
	summary := "[Compressed conversation summary: "
	summary += "The conversation covers "
	if len(messages) <= 5 {
		summary += "a brief exchange"
	} else {
		summary += "an extended discussion"
	}
	summary += " with " + messages[0].Role + " starting the conversation"
	summary += " ending with " + messages[len(messages)-1].Role + "'s last message at "
	summary += messages[len(messages)-1].Time.Format("15:04:05")
	summary += "]"
	
	return summary
}

// RegisterSession 注册新会话
func (cm *ContextManager) RegisterSession(id string, initialMessages []Message) {
	cm.mu.Lock()
	defer cm.mu.Unlock()
	
	tokenCount := 0
	for _, msg := range initialMessages {
		tokenCount += len(msg.Content) / 4
	}
	
	cm.sessions[id] = &SessionContext{
		ID:             id,
		Messages:       initialMessages,
		TokenCount:     tokenCount,
		LastActivity:   time.Now(),
		State:          StateActive,
		StateChangedAt: time.Now(),
	}
	
	log.Printf("registered session context: %s (%d tokens)", id, tokenCount)
}

// AppendMessage 追加消息
func (cm *ContextManager) AppendMessage(sessionID string, msg Message) error {
	cm.mu.Lock()
	session, exists := cm.sessions[sessionID]
	cm.mu.Unlock()
	
	if !exists {
		return nil
	}
	
	session.mu.Lock()
	defer session.mu.Unlock()
	
	session.Messages = append(session.Messages, msg)
	session.TokenCount += len(msg.Content) / 4
	session.LastActivity = time.Now()
	
	// 检查是否需要触发压缩
	if session.TokenCount > cm.maxContextTokens {
		go cm.CompactContext(sessionID)
	}
	
	return nil
}

// CompactContext 异步压缩上下文
// 核心逻辑:在后台准备新实例的同时保持原实例继续服务
func (cm *ContextManager) CompactContext(sessionID string) {
	cm.mu.RLock()
	session, exists := cm.sessions[sessionID]
	cm.mu.RUnlock()
	
	if !exists {
		return
	}
	
	session.mu.Lock()
	if session.State == StateCompacting || session.State == StateSwitching {
		session.mu.Unlock()
		return // 已经在压缩中
	}
	session.State = StateCompacting
	session.StateChangedAt = time.Now()
	
	// 复制当前消息用于压缩
	currentMessages := make([]Message, len(session.Messages))
	copy(currentMessages, session.Messages)
	session.mu.Unlock()
	
	// 后台执行压缩(不阻塞媒体路径)
	targetTokens := int(float64(cm.maxContextTokens) * cm.compactionRatio)
	compactedMessages, newTokenCount := cm.compactionFunc(currentMessages, targetTokens)
	
	log.Printf("compacted session %s: %d -> %d tokens",
		sessionID, session.TokenCount, newTokenCount)
	
	// 更新上下文
	session.mu.Lock()
	session.Messages = compactedMessages
	session.TokenCount = newTokenCount
	session.KVState = nil // KV缓存失效,需要重建
	session.State = StateReady
	session.StateChangedAt = time.Now()
	session.mu.Unlock()
	
	// 通知调度器准备新实例
	// 实际系统中,此处会触发新模型实例的预热和切换
	log.Printf("context compaction complete for session %s, ready for handoff", sessionID)
}

// GetSessionContext 获取会话上下文快照
func (cm *ContextManager) GetSessionContext(sessionID string) *SessionContext {
	cm.mu.RLock()
	session, exists := cm.sessions[sessionID]
	cm.mu.RUnlock()
	
	if !exists {
		return nil
	}
	
	session.mu.RLock()
	defer session.mu.RUnlock()
	
	// 返回副本
	ctxCopy := &SessionContext{
		ID:            session.ID,
		Messages:      make([]Message, len(session.Messages)),
		TokenCount:    session.TokenCount,
		LastActivity:  session.LastActivity,
		State:         session.State,
	}
	copy(ctxCopy.Messages, session.Messages)
	
	return ctxCopy
}

// UnregisterSession 注销会话
func (cm *ContextManager) UnregisterSession(sessionID string) {
	cm.mu.Lock()
	defer cm.mu.Unlock()
	delete(cm.sessions, sessionID)
	log.Printf("unregistered session context: %s", sessionID)
}

六、连续语音中的离散消息提取

6.1 问题本质

虽然GPT-Live的语音模型处理的是连续语音流,但ChatGPT的UI界面、安全分析系统、日志系统等外围系统仍然需要离散的用户和助手消息。这带来了一个核心矛盾:连续输入中如何提取离散的对话轮次?

6.2 推测性消息队列

OpenAI的解决方案是维护一个推测性消息队列(speculative message queue):

时间线 ────────────────────────────────────────────────────────────►

用户语音:   "嗯...我想知道...那个...GPT-5.5的性能..."
              ↑         ↑           ↑           ↑
              │         │           │           │
消息状态:   [暂定]    [暂定]      [暂定]     [确认]
            归属:用户  归属:用户   归属:用户   归属:用户
            文本可变   文本可变    文本可变    文本已锁定

核心逻辑:

  1. 最新消息始终是暂定状态(provisional)
  2. 随着更多语音到达,其文本、时间戳和发言人归属都可以改变
  3. 只有当一个人发言足够久,归属变得可靠时,服务器才会最终确认该消息
  4. 系统维护两种视图:推测视图(UI实时更新)和权威视图(分析日志用)

6.3 消息提取器实现

"""
从连续语音流中提取离散消息
模拟GPT-Live的推测性消息队列设计
"""

import asyncio
import time
from dataclasses import dataclass, field
from enum import Enum
from typing import Optional


class Speaker(Enum):
    USER = "user"
    ASSISTANT = "assistant"
    UNKNOWN = "unknown"


class MessageStatus(Enum):
    PROVISIONAL = "provisional"  # 暂定状态
    FINALIZED = "finalized"      # 已确认


@dataclass
class TranscriptSegment:
    """转录片段"""
    text: str
    speaker: Speaker
    start_time: float
    end_time: float
    confidence: float
    is_acknowledgment: bool = False  # 是否为语气词(如"嗯""好的")


@dataclass
class DiscreteMessage:
    """离散消息"""
    id: str
    speaker: Speaker
    text: str
    start_time: float
    end_time: float
    status: MessageStatus
    segments: list[TranscriptSegment] = field(default_factory=list)
    is_acknowledgment: bool = False


class SpeculativeMessageQueue:
    """
    推测性消息队列
    从连续语音流中提取离散消息
    维护两种视图:推测视图(UI)和权威视图(日志)
    """
    
    def __init__(
        self,
        finalize_after_silence_ms: float = 800.0,  # 800ms静音后确认
        min_utterance_ms: float = 300.0,           # 最短有效话语
        acknowledgment_threshold_ms: float = 500.0, # 语气词判定阈值
    ):
        self.finalize_after_silence_ms = finalize_after_silence_ms
        self.min_utterance_ms = min_utterance_ms
        self.acknowledgment_threshold_ms = acknowledgment_threshold_ms
        
        self._queue: list[DiscreteMessage] = []
        self._current_message: Optional[DiscreteMessage] = None
        self._last_segment_time: float = 0.0
        self._message_counter = 0
        
        # 两种视图
        self._speculative_view: list[DiscreteMessage] = []
        self._authoritative_view: list[DiscreteMessage] = []
    
    def process_segment(self, segment: TranscriptSegment) -> list[DiscreteMessage]:
        """
        处理新的转录片段
        返回状态变更的消息列表
        """
        now = segment.start_time
        changes = []
        
        # 判断是否为新轮次
        time_since_last = (now - self._last_segment_time) * 1000  # 转为ms
        
        if time_since_last > self.finalize_after_silence_ms:
            # 静音超过阈值,确认当前消息
            if self._current_message and self._current_message.status == MessageStatus.PROVISIONAL:
                self._finalize_current_message()
                changes.append(self._current_message)
            
            # 开始新消息
            self._start_new_message(segment)
            changes.append(self._current_message)
        
        elif self._current_message is None:
            # 没有活跃消息,开始新消息
            self._start_new_message(segment)
            changes.append(self._current_message)
        
        elif segment.speaker != self._current_message.speaker:
            # 发言人变更
            if self._current_message.status == MessageStatus.PROVISIONAL:
                if self._is_acknowledgment(segment):
                    # 是语气词,不打断,合并到当前消息
                    segment.is_acknowledgment = True
                    self._current_message.segments.append(segment)
                    self._current_message.is_acknowledgment = True
                else:
                    # 真正的打断,确认当前消息,开始新消息
                    self._finalize_current_message()
                    changes.append(self._current_message)
                    self._start_new_message(segment)
                    changes.append(self._current_message)
        
        else:
            # 同一发言人继续,追加到当前消息
            self._current_message.segments.append(segment)
            self._current_message.text = self._build_text(self._current_message.segments)
            self._current_message.end_time = segment.end_time
            
            # 更新推测视图
            self._update_speculative_view()
        
        self._last_segment_time = segment.end_time
        return changes
    
    def _start_new_message(self, segment: TranscriptSegment) -> None:
        """开始新消息"""
        self._message_counter += 1
        self._current_message = DiscreteMessage(
            id=f"msg_{self._message_counter}",
            speaker=segment.speaker,
            text=segment.text,
            start_time=segment.start_time,
            end_time=segment.end_time,
            status=MessageStatus.PROVISIONAL,
            segments=[segment],
            is_acknowledgment=segment.is_acknowledgment,
        )
        self._queue.append(self._current_message)
        self._update_speculative_view()
    
    def _finalize_current_message(self) -> None:
        """确认当前消息"""
        if self._current_message is None:
            return
        
        duration = (self._current_message.end_time - self._current_message.start_time) * 1000
        
        if duration < self.min_utterance_ms and not self._current_message.is_acknowledgment:
            # 太短且不是语气词,丢弃
            self._queue.remove(self._current_message)
            self._update_speculative_view()
            self._current_message = None
            return
        
        self._current_message.status = MessageStatus.FINALIZED
        
        # 更新权威视图
        self._authoritative_view.append(self._current_message)
        
        # 从推测视图移除(已被确认)
        # 实际系统中,推测视图会保留但标记为已确认
        self._current_message = None
    
    def _is_acknowledgment(self, segment: TranscriptSegment) -> bool:
        """判断是否为语气词"""
        acknowledgment_words = {
            "嗯", "嗯哼", "哦", "好的", "对的", "是的", "没错",
            "mhmm", "uh-huh", "yeah", "okay", "got it", "right",
            "i see", "aha", "mm", "hmm",
        }
        return segment.text.strip().lower() in acknowledgment_words
    
    def _build_text(self, segments: list[TranscriptSegment]) -> str:
        """从片段构建完整文本"""
        return " ".join(seg.text for seg in segments)
    
    def _update_speculative_view(self) -> None:
        """更新推测视图"""
        self._speculative_view = [
            msg for msg in self._queue
            if msg.status == MessageStatus.PROVISIONAL
        ]
    
    def get_speculative_view(self) -> list[DiscreteMessage]:
        """获取推测视图(供UI使用)"""
        return list(self._speculative_view)
    
    def get_authoritative_view(self) -> list[DiscreteMessage]:
        """获取权威视图(供日志和分析使用)"""
        return list(self._authoritative_view)
    
    def get_all_messages(self) -> list[DiscreteMessage]:
        """获取所有消息"""
        return list(self._queue)


# 使用示例
async def demo_speculative_queue():
    """演示推测性消息队列的工作流程"""
    queue = SpeculativeMessageQueue()
    
    segments = [
        # 用户开始说话
        TranscriptSegment("嗯", Speaker.USER, 0.0, 0.3, 0.95),
        TranscriptSegment("我想知道", Speaker.USER, 0.3, 0.8, 0.92),
        TranscriptSegment("GPT-5.5的性能怎么样", Speaker.USER, 0.8, 1.8, 0.88),
        # 静音800ms后确认消息
        
        # 助手开始回答
        TranscriptSegment("好的", Speaker.ASSISTANT, 2.8, 3.0, 0.97),
        TranscriptSegment("GPT-5.5相比前代", Speaker.ASSISTANT, 3.0, 3.8, 0.95),
        # 用户打断
        TranscriptSegment("具体数字", Speaker.USER, 3.8, 4.2, 0.90),
    ]
    
    for seg in segments:
        changes = queue.process_segment(seg)
        if changes:
            for msg in changes:
                print(f"[{msg.status.value}] {msg.speaker.value}: {msg.text}")
    
    print(f"\n推测视图: {len(queue.get_speculative_view())} 条")
    print(f"权威视图: {len(queue.get_authoritative_view())} 条")
    
    for msg in queue.get_authoritative_view():
        print(f"  [{msg.speaker.value}] {msg.text}")


asyncio.run(demo_speculative_queue())

七、会话启动优化:WARP与Instant Connect

7.1 标准的WebRTC握手开销

标准的WebRTC连接建立需要6次网络往返:

客户端                          服务器
  │                               │
  │── ICE Binding Request ───────►│  RTT 1
  │◄── ICE Binding Response ─────│
  │                               │
  │── DTLS ClientHello ──────────►│  RTT 2
  │◄── DTLS ServerHello ─────────│
  │◄── DTLS Certificate ─────────│  RTT 3
  │── DTLS Finished ────────────►│
  │                               │
  │── SCTP INIT ────────────────►│  RTT 4
  │◄── SCTP INIT_ACK ────────────│
  │── SCTP COOKIE_ECHO ─────────►│  RTT 5
  │◄── SCTP COOKIE_ACK ──────────│
  │                               │
  │── DCEP ---- data channel ----►│  RTT 6
  │                               │
  6 RTTs = 300-600ms 延迟

7.2 WARP协议优化

OpenAI与WebRTC社区合作设计了WARP(WebRTC Abridged Roundtrip Protocol),将6次网络往返减少到1次(来源:OpenAI工程博客)。

WARP的核心优化:

  1. SPED:将DTLS握手搭载在ICE之上
  2. DTLS 1.3:使用更快的DTLS 1.3握手(1-RTT vs 2-RTT)
  3. SNAP:预先协商SCTP握手
  4. 预协商数据通道:绕过DCEP

WARP已被提交为IETF草案(https://datatracker.ietf.org/doc/draft-uberti-tsvwg-warp/),并已集成到libwebrtc和Pion中。

7.3 Instant Connect

除了WARP,OpenAI还开发了Instant Connect技术,在连接前预先协商SDP参数:

// Instant Connect 实现
// 预先协商SDP参数,实现单UDP包启动会话

package transport

import (
	"crypto/rand"
	"encoding/hex"
	"encoding/json"
	"log"
	"net"
	"time"
)

// PreNegotiatedSession 预协商会话
type PreNegotiatedSession struct {
	SessionID     string
	Ufrag         string
	Pwd           string
	Fingerprint   string
	AudioCodec    string
	CreatedAt     time.Time
	ExpiresAt     time.Time
	IsUsed        bool
}

// InstantConnectManager 即时连接管理器
type InstantConnectManager struct {
	preSessions    map[string]*PreNegotiatedSession
	sessionTTL     time.Duration
	audioCodec     string
	serverUfrag    string
	serverPwd      string
	serverFingerprint string
}

// NewInstantConnectManager 创建即时连接管理器
func NewInstantConnectManager(
	sessionTTL time.Duration,
	audioCodec string,
) *InstantConnectManager {
	return &InstantConnectManager{
		preSessions:    make(map[string]*PreNegotiatedSession),
		sessionTTL:     sessionTTL,
		audioCodec:     audioCodec,
		serverUfrag:    generateICEUfrag(),
		serverPwd:      generateICEPwd(),
		serverFingerprint: generateFingerprint(),
	}
}

// GeneratePreSession 生成预协商会话
// 在用户点击按钮前就已经完成
func (m *InstantConnectManager) GeneratePreSession() *PreNegotiatedSession {
	sessionID := generateSessionID()
	
	session := &PreNegotiatedSession{
		SessionID:   sessionID,
		Ufrag:       generateICEUfrag(),
		Pwd:         generateICEPwd(),
		Fingerprint: m.serverFingerprint,
		AudioCodec:  m.audioCodec,
		CreatedAt:   time.Now(),
		ExpiresAt:   time.Now().Add(m.sessionTTL),
		IsUsed:      false,
	}
	
	m.preSessions[sessionID] = session
	log.Printf("generated pre-negotiated session: %s (expires %v)",
		sessionID, session.ExpiresAt)
	
	return session
}

// HandleFirstPacket 处理第一个UDP包
// 当客户端发送第一个媒体包时,服务器立即响应
func (m *InstantConnectManager) HandleFirstPacket(
	conn *net.UDPConn,
	addr *net.UDPAddr,
	packet []byte,
) (*PreNegotiatedSession, error) {
	// 解析包中的SessionID
	sessionID := extractSessionID(packet)
	
	session, exists := m.preSessions[sessionID]
	if !exists {
		return nil, nil // 走标准信令流程
	}
	
	if time.Now().After(session.ExpiresAt) {
		delete(m.preSessions, sessionID)
		return nil, nil // 会话已过期
	}
	
	if session.IsUsed {
		return nil, nil // 会话已被使用
	}
	
	session.IsUsed = true
	
	// 立即响应
	response := m.buildImmediateResponse(session)
	conn.WriteTo(response, addr)
	
	log.Printf("instant connect for session %s from %s", sessionID, addr)
	return session, nil
}

// buildImmediateResponse 构建即时响应
func (m *InstantConnectManager) buildImmediateResponse(
	session *PreNegotiatedSession,
) []byte {
	response := map[string]interface{}{
		"type":        "instant_connect",
		"session_id":  session.SessionID,
		"ufrag":       m.serverUfrag,
		"pwd":         m.serverPwd,
		"fingerprint": m.serverFingerprint,
		"codec":       session.AudioCodec,
		"timestamp":   time.Now().UnixMilli(),
	}
	
	data, _ := json.Marshal(response)
	return data
}

// 辅助函数
func generateSessionID() string {
	b := make([]byte, 16)
	rand.Read(b)
	return hex.EncodeToString(b)
}

func generateICEUfrag() string {
	b := make([]byte, 8)
	rand.Read(b)
	return hex.EncodeToString(b)
}

func generateICEPwd() string {
	b := make([]byte, 22)
	rand.Read(b)
	return hex.EncodeToString(b)
}

func generateFingerprint() string {
	return "sha-256 " + generateSessionID()
}

func extractSessionID(packet []byte) string {
	if len(packet) < 4 {
		return ""
	}
	// 简化实现:假设packet前4字节是sessionID长度
	// 实际实现更复杂
	return string(packet[4:])
}

WARP和Instant Connect的结合使得客户端只需发送一个UDP数据包就能启动会话,服务器可以立即响应,让系统马上开始监听。这对比传统的6次握手,启动延迟降低了约80-90%(来源:OpenAI工程博客)。


八、生产环境的测试与部署

8.1 影子测试(Shadow Testing)

OpenAI在实际部署前进行了数月的影子测试——将生产流量同时路由到旧版Advanced Voice Mode和新版GPT-Live,新系统以只读模式运行,不影响用户实际体验。

影子测试中发现的几个关键教训:

  1. 容量不等于GPU吞吐量:语音会话持续打开并连续发送帧,CPU侧的流处理器、队列和网络路径必须与推理同步扩展。实际负载下,支撑组件比负载测试估计的更早饱和。

  2. 地理路由是第一优先级:将会话路由到远端容量会在启动和流式传输的多个环节增加延迟。

  3. 长时间会话暴露深层问题:长时间运行暴露了内存和持久化压力、重连场景下的压缩和状态恢复、客户端断开时的竞态条件。

8.2 可观测性

OpenAI团队增加了更细粒度的遥测、针对已知良好配置的验证、分阶段灰度以及快速隔离或禁用单个路径的能力。


九、总结与展望

GPT-Live的全双工架构代表了语音AI从"回合制对话"到"连续流对话"的根本性转变。其核心突破可以概括为三个层面:

  1. 架构层面:全双工音频 + 媒体/逻辑分离 + 异步委派,实现了语音层与推理层的彻底解耦
  2. 工程层面:Go重写媒体前端(p95=p50旧系统)、WARP协议(6 RTT→1 RTT)、Instant Connect(单UDP包启动)
  3. 产品层面:移除轮次检测器、推测性消息队列、无缝上下文压缩

OpenAI官方表示,这套架构将成为实时交互的底层平台,未来将扩展到更多设备和应用场景,让每一次语音交流都如同面对面般真实(来源:OpenAI工程博客)。


参考文献

  1. OpenAI, “Introducing GPT‑Live”, July 8, 2026. https://openai.com/index/introducing-gpt-live/
  2. Justin Uberti and Zahan Malkani, “How we built a realtime system for responsive voice AI in six months”, August 3, 2026. https://openai.com/index/continuous-voice-interaction-with-gpt-live/
  3. OpenAI, “Delivering low-latency voice AI at scale”, 2026. https://openai.com/index/delivering-low-latency-voice-ai-at-scale/
  4. IETF, “WARP: WebRTC Abridged Roundtrip Protocol”, https://datatracker.ietf.org/doc/draft-uberti-tsvwg-warp/
  5. IETF, “SPED: Speedy DTLS Encapsulation”, https://datatracker.ietf.org/doc/draft-hancke-webrtc-sped/
  6. IETF, “SNAP: SCTP Negotiation Acceleration Protocol”, https://datatracker.ietf.org/doc/draft-hancke-tsvwg-snap/
  7. 腾讯新闻, “OpenAI彻底干掉语音延迟!6个月重写底层架构,GPT Live全双工技术揭秘”, August 4, 2026.
  8. AI Times, “턴 디텍터 없애고 추론 분리… 오픈AI가 밝힌 ‘GPT-라이브’ 초저지연 비밀”, August 4, 2026.