feat: add qwen realtime practice beta
This commit is contained in:
@@ -0,0 +1,327 @@
|
||||
import { apiRequest } from './api';
|
||||
|
||||
type RealtimeStatus = 'idle' | 'connecting' | 'connected' | 'listening' | 'speaking';
|
||||
|
||||
interface RealtimeSdpResponse {
|
||||
answerSdp: string;
|
||||
model: string;
|
||||
}
|
||||
|
||||
interface RealtimeEvent {
|
||||
type?: string;
|
||||
transcript?: string;
|
||||
error?: { message?: string };
|
||||
}
|
||||
|
||||
interface ProviderEventHandlers {
|
||||
onSessionCreated: () => void;
|
||||
onSessionUpdated: () => void;
|
||||
onSessionError: (error: Error) => void;
|
||||
}
|
||||
|
||||
export interface RealtimePracticeCallbacks {
|
||||
onStatus: (status: RealtimeStatus) => void;
|
||||
onTraineeTranscript: (text: string) => void;
|
||||
onCustomerTranscript: (text: string) => void;
|
||||
onAudioPlaybackBlocked: (blocked: boolean) => void;
|
||||
onTerminated: () => void;
|
||||
onError: (message: string) => void;
|
||||
}
|
||||
|
||||
export interface RealtimePracticeController {
|
||||
start: (instructions: string) => Promise<void>;
|
||||
stop: () => void;
|
||||
resumeAudio: () => Promise<boolean>;
|
||||
}
|
||||
|
||||
const isSupported = () =>
|
||||
typeof navigator !== 'undefined'
|
||||
&& Boolean(navigator.mediaDevices?.getUserMedia)
|
||||
&& typeof RTCPeerConnection !== 'undefined';
|
||||
|
||||
const waitForIceComplete = (connection: RTCPeerConnection) =>
|
||||
new Promise<void>((resolve) => {
|
||||
if (connection.iceGatheringState === 'complete') {
|
||||
resolve();
|
||||
return;
|
||||
}
|
||||
const previous = connection.onicegatheringstatechange;
|
||||
let timeout: ReturnType<typeof setTimeout>;
|
||||
const done = () => {
|
||||
clearTimeout(timeout);
|
||||
connection.onicegatheringstatechange = previous;
|
||||
resolve();
|
||||
};
|
||||
timeout = setTimeout(done, 3000);
|
||||
connection.onicegatheringstatechange = () => {
|
||||
previous?.call(connection, new Event('icegatheringstatechange'));
|
||||
if (connection.iceGatheringState === 'complete') done();
|
||||
};
|
||||
});
|
||||
|
||||
const normalizeSdp = (sdp: string) => `${sdp.replace(/\r\n|\r|\n/g, '\r\n').replace(/\r\n$/, '')}\r\n`;
|
||||
const trimTranscript = (value?: string) => (value || '').trim();
|
||||
|
||||
export const createRealtimePracticeController = (callbacks: RealtimePracticeCallbacks): RealtimePracticeController => {
|
||||
let activeAttempt = 0;
|
||||
let connection: RTCPeerConnection | null = null;
|
||||
let microphone: MediaStream | null = null;
|
||||
let speaker: HTMLAudioElement | null = null;
|
||||
let controlChannel: RTCDataChannel | null = null;
|
||||
let cancelSessionWait: (() => void) | null = null;
|
||||
let sessionWaitAttempt: number | null = null;
|
||||
|
||||
const isCurrent = (attempt: number) => activeAttempt === attempt;
|
||||
|
||||
const releaseResources = () => {
|
||||
controlChannel?.close();
|
||||
controlChannel = null;
|
||||
microphone?.getTracks().forEach((track) => track.stop());
|
||||
microphone = null;
|
||||
if (speaker) {
|
||||
speaker.pause();
|
||||
speaker.srcObject = null;
|
||||
speaker = null;
|
||||
}
|
||||
connection?.close();
|
||||
connection = null;
|
||||
};
|
||||
|
||||
const sendSessionUpdate = (channel: RTCDataChannel, instructions: string) => {
|
||||
if (channel.readyState !== 'open') return;
|
||||
channel.send(JSON.stringify({
|
||||
event_id: `event_${Date.now()}`,
|
||||
type: 'session.update',
|
||||
session: {
|
||||
modalities: ['text', 'audio'],
|
||||
voice: 'Ethan',
|
||||
input_audio_format: 'pcm',
|
||||
output_audio_format: 'pcm',
|
||||
input_audio_transcription: { model: 'qwen3-asr-flash-realtime' },
|
||||
instructions,
|
||||
turn_detection: {
|
||||
type: 'semantic_vad',
|
||||
threshold: 0.5,
|
||||
prefix_padding_ms: 500,
|
||||
silence_duration_ms: 800
|
||||
}
|
||||
}
|
||||
}));
|
||||
};
|
||||
|
||||
const handleEvent = (raw: string, attempt: number, provider: ProviderEventHandlers) => {
|
||||
if (!isCurrent(attempt)) return;
|
||||
let event: RealtimeEvent;
|
||||
try {
|
||||
event = JSON.parse(raw) as RealtimeEvent;
|
||||
} catch {
|
||||
return;
|
||||
}
|
||||
if (event.type === 'session.created') provider.onSessionCreated();
|
||||
if (event.type === 'session.updated') provider.onSessionUpdated();
|
||||
if (event.type === 'input_audio_buffer.speech_started') callbacks.onStatus('listening');
|
||||
if (event.type === 'response.created') callbacks.onStatus('speaking');
|
||||
if (event.type === 'response.done') callbacks.onStatus('connected');
|
||||
if (event.type === 'conversation.item.input_audio_transcription.completed') {
|
||||
const transcript = trimTranscript(event.transcript);
|
||||
if (transcript) callbacks.onTraineeTranscript(transcript);
|
||||
}
|
||||
if (event.type === 'response.audio_transcript.done') {
|
||||
const transcript = trimTranscript(event.transcript);
|
||||
if (transcript) callbacks.onCustomerTranscript(transcript);
|
||||
}
|
||||
if (event.type === 'error') {
|
||||
const error = new Error(event.error?.message || '实时陪练连接异常,请改用录音对练');
|
||||
callbacks.onError(error.message);
|
||||
provider.onSessionError(error);
|
||||
}
|
||||
};
|
||||
|
||||
const bindChannel = (channel: RTCDataChannel, attempt: number, provider: ProviderEventHandlers) => {
|
||||
channel.onmessage = (event) => handleEvent(String(event.data || ''), attempt, provider);
|
||||
channel.onerror = () => {
|
||||
if (!isCurrent(attempt)) return;
|
||||
const error = new Error('实时陪练数据通道异常,请改用录音对练');
|
||||
callbacks.onError(error.message);
|
||||
provider.onSessionError(error);
|
||||
};
|
||||
};
|
||||
|
||||
const stop = () => {
|
||||
activeAttempt += 1;
|
||||
cancelSessionWait?.();
|
||||
cancelSessionWait = null;
|
||||
sessionWaitAttempt = null;
|
||||
releaseResources();
|
||||
callbacks.onAudioPlaybackBlocked(false);
|
||||
callbacks.onStatus('idle');
|
||||
};
|
||||
|
||||
const terminate = (attempt: number, message: string) => {
|
||||
if (!isCurrent(attempt)) return;
|
||||
callbacks.onError(message);
|
||||
stop();
|
||||
callbacks.onTerminated();
|
||||
};
|
||||
|
||||
const resumeAudio = async () => {
|
||||
if (!speaker) return false;
|
||||
try {
|
||||
await speaker.play();
|
||||
callbacks.onAudioPlaybackBlocked(false);
|
||||
return true;
|
||||
} catch {
|
||||
callbacks.onAudioPlaybackBlocked(true);
|
||||
return false;
|
||||
}
|
||||
};
|
||||
|
||||
const start = async (instructions: string) => {
|
||||
if (!isSupported()) throw new Error('当前设备不支持实时语音,请使用录音对练');
|
||||
stop();
|
||||
const attempt = ++activeAttempt;
|
||||
callbacks.onStatus('connecting');
|
||||
|
||||
let localMicrophone: MediaStream | null = null;
|
||||
let localConnection: RTCPeerConnection | null = null;
|
||||
let localChannel: RTCDataChannel | null = null;
|
||||
let localSpeaker: HTMLAudioElement | null = null;
|
||||
let readyTimeout: ReturnType<typeof setTimeout> | undefined;
|
||||
let disconnectTimeout: ReturnType<typeof setTimeout> | undefined;
|
||||
let sessionUpdateSent = false;
|
||||
let resolveReady: (() => void) | undefined;
|
||||
let rejectReady: ((error: Error) => void) | undefined;
|
||||
let sessionReady: Promise<void> | null = null;
|
||||
const releaseLocalResources = () => {
|
||||
if (disconnectTimeout) clearTimeout(disconnectTimeout);
|
||||
localChannel?.close();
|
||||
localMicrophone?.getTracks().forEach((track) => track.stop());
|
||||
if (localSpeaker) {
|
||||
localSpeaker.pause();
|
||||
localSpeaker.srcObject = null;
|
||||
}
|
||||
localConnection?.close();
|
||||
};
|
||||
const ensureCurrent = () => {
|
||||
if (isCurrent(attempt)) return true;
|
||||
releaseLocalResources();
|
||||
return false;
|
||||
};
|
||||
|
||||
try {
|
||||
localMicrophone = await navigator.mediaDevices.getUserMedia({ audio: true });
|
||||
if (!ensureCurrent()) return;
|
||||
microphone = localMicrophone;
|
||||
localConnection = new RTCPeerConnection({ iceServers: [] });
|
||||
connection = localConnection;
|
||||
const audioSender = localConnection.addTransceiver('audio', { direction: 'sendrecv' }).sender;
|
||||
const provider: ProviderEventHandlers = {
|
||||
onSessionCreated: () => {
|
||||
if (!isCurrent(attempt) || !localChannel || sessionUpdateSent) return;
|
||||
sessionUpdateSent = true;
|
||||
sendSessionUpdate(localChannel, instructions);
|
||||
},
|
||||
onSessionUpdated: () => {
|
||||
if (!isCurrent(attempt) || !sessionUpdateSent) return;
|
||||
const track = localMicrophone?.getAudioTracks()[0] || null;
|
||||
void audioSender.replaceTrack(track)
|
||||
.then(() => {
|
||||
if (isCurrent(attempt)) resolveReady?.();
|
||||
})
|
||||
.catch((error: unknown) => rejectReady?.(error instanceof Error ? error : new Error('实时陪练麦克风启动失败')));
|
||||
},
|
||||
onSessionError: (error) => {
|
||||
rejectReady?.(error);
|
||||
terminate(attempt, error.message);
|
||||
}
|
||||
};
|
||||
|
||||
localChannel = localConnection.createDataChannel('oai-events');
|
||||
controlChannel = localChannel;
|
||||
bindChannel(localChannel, attempt, provider);
|
||||
localConnection.ondatachannel = (event) => bindChannel(event.channel, attempt, provider);
|
||||
localConnection.ontrack = (event) => {
|
||||
if (!isCurrent(attempt)) return;
|
||||
const stream = event.streams[0] || new MediaStream([event.track]);
|
||||
localSpeaker = new Audio();
|
||||
speaker = localSpeaker;
|
||||
localSpeaker.autoplay = true;
|
||||
localSpeaker.setAttribute('playsinline', 'true');
|
||||
localSpeaker.srcObject = stream;
|
||||
void localSpeaker.play()
|
||||
.then(() => {
|
||||
if (isCurrent(attempt)) callbacks.onAudioPlaybackBlocked(false);
|
||||
})
|
||||
.catch(() => {
|
||||
if (isCurrent(attempt)) callbacks.onAudioPlaybackBlocked(true);
|
||||
});
|
||||
};
|
||||
localConnection.onconnectionstatechange = () => {
|
||||
if (!isCurrent(attempt)) return;
|
||||
if (localConnection?.connectionState === 'failed' || localConnection?.connectionState === 'closed') {
|
||||
terminate(attempt, '实时陪练已断开,请改用录音对练');
|
||||
return;
|
||||
}
|
||||
if (localConnection?.connectionState === 'disconnected') {
|
||||
disconnectTimeout ??= setTimeout(() => {
|
||||
if (localConnection?.connectionState === 'disconnected') {
|
||||
terminate(attempt, '实时陪练已断开,请改用录音对练');
|
||||
}
|
||||
}, 3000);
|
||||
} else if (disconnectTimeout) {
|
||||
clearTimeout(disconnectTimeout);
|
||||
disconnectTimeout = undefined;
|
||||
}
|
||||
};
|
||||
|
||||
const offer = await localConnection.createOffer();
|
||||
if (!ensureCurrent()) return;
|
||||
await localConnection.setLocalDescription(offer);
|
||||
if (!ensureCurrent()) return;
|
||||
await waitForIceComplete(localConnection);
|
||||
if (!ensureCurrent()) return;
|
||||
const offerSdp = localConnection.localDescription?.sdp;
|
||||
if (!offerSdp) throw new Error('实时陪练连接请求创建失败');
|
||||
const response = await apiRequest<RealtimeSdpResponse>({
|
||||
url: '/api/train/practice/realtime/sdp',
|
||||
method: 'POST',
|
||||
data: { offerSdp: normalizeSdp(offerSdp) },
|
||||
timeout: 30000
|
||||
});
|
||||
if (!ensureCurrent()) return;
|
||||
sessionReady = new Promise<void>((resolve, reject) => {
|
||||
resolveReady = resolve;
|
||||
rejectReady = reject;
|
||||
readyTimeout = setTimeout(() => reject(new Error('实时陪练会话初始化超时,请使用录音对练')), 15000);
|
||||
});
|
||||
void sessionReady.catch(() => undefined);
|
||||
sessionWaitAttempt = attempt;
|
||||
cancelSessionWait = () => {
|
||||
if (readyTimeout) clearTimeout(readyTimeout);
|
||||
rejectReady?.(new Error('实时陪练已结束'));
|
||||
};
|
||||
await localConnection.setRemoteDescription({ type: 'answer', sdp: normalizeSdp(response.answerSdp) });
|
||||
if (!ensureCurrent()) return;
|
||||
const waitForSessionReady = sessionReady;
|
||||
if (!waitForSessionReady) throw new Error('实时陪练会话初始化失败');
|
||||
await waitForSessionReady;
|
||||
if (!ensureCurrent()) return;
|
||||
callbacks.onStatus('connected');
|
||||
} catch (error) {
|
||||
if (!isCurrent(attempt)) {
|
||||
releaseLocalResources();
|
||||
return;
|
||||
}
|
||||
stop();
|
||||
throw error instanceof Error ? error : new Error('实时陪练连接失败,请改用录音对练');
|
||||
} finally {
|
||||
if (readyTimeout) clearTimeout(readyTimeout);
|
||||
if (sessionWaitAttempt === attempt) {
|
||||
sessionWaitAttempt = null;
|
||||
cancelSessionWait = null;
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
return { start, stop, resumeAudio };
|
||||
};
|
||||
Reference in New Issue
Block a user