const trimTranscript = (value) => (value || '').trim(); const normalizeSdp = (sdp) => `${sdp.replace(/\r\n|\r|\n/g, '\r\n').replace(/\r\n$/, '')}\r\n`; const TOOL_CALL_FALLBACK_OUTPUT = '工具暂时不可用,请直接回应员工'; export const isRealtimeBrowserSupported = () => typeof navigator !== 'undefined' && Boolean(navigator.mediaDevices?.getUserMedia) && typeof RTCPeerConnection !== 'undefined'; const waitForIceComplete = (connection) => new Promise((resolve) => { if (connection.iceGatheringState === 'complete') { resolve(); return; } const previous = connection.onicegatheringstatechange; let timeout; 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(); }; }); export const createRealtimeBrowserEngine = (callbacks, exchangeSdp) => { let activeAttempt = 0; let connection = null; let microphone = null; let speaker = null; let controlChannel = null; let cancelSessionWait = null; let sessionWaitAttempt = null; let toolCallChain = Promise.resolve(); const isCurrent = (attempt) => 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, sessionUpdate) => { if (channel.readyState !== 'open') return; channel.send(JSON.stringify({ event_id: `event_${Date.now()}`, ...sessionUpdate })); }; const handleToolCall = (attempt, event) => { const callId = trimTranscript(event.call_id); const name = trimTranscript(event.name); if (!callId || !name) return; const argsJson = typeof event.arguments === 'string' ? event.arguments : '{}'; toolCallChain = toolCallChain.then(async () => { if (!isCurrent(attempt)) return; const channel = controlChannel; if (!channel || channel.readyState !== 'open') return; callbacks.onStatus('calling-tool'); let output; try { output = await callbacks.onToolCall(callId, name, argsJson); } catch { output = TOOL_CALL_FALLBACK_OUTPUT; } if (!isCurrent(attempt) || channel.readyState !== 'open') return; channel.send(JSON.stringify({ event_id: `event_${Date.now()}`, type: 'conversation.item.create', item: { type: 'function_call_output', call_id: callId, output } })); channel.send(JSON.stringify({ event_id: `event_${Date.now()}`, type: 'response.create' })); callbacks.onStatus('connected'); }).catch(() => { // 链内任何意外抛错不得毒化后续工具调用;本次调用软失败,对话继续 }); }; const handleEvent = (raw, attempt, provider) => { if (!isCurrent(attempt)) return; let event; try { event = JSON.parse(raw); } 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 === 'response.function_call_arguments.done') { handleToolCall(attempt, event); } if (event.type === 'error') { const error = new Error(event.error?.message || '实时陪练连接异常,请改用录音对练'); callbacks.onError(error.message); provider.onSessionError(error); } }; const bindChannel = (channel, attempt, provider) => { 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; toolCallChain = Promise.resolve(); releaseResources(); callbacks.onAudioPlaybackBlocked(false); callbacks.onStatus('idle'); }; const terminate = (attempt, message) => { 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 (sessionUpdate, sessionId) => { if (!isRealtimeBrowserSupported()) throw new Error('当前设备的 WebView 不支持实时语音,请更新系统 WebView 后重试'); stop(); const attempt = ++activeAttempt; callbacks.onStatus('connecting'); let localMicrophone = null; let localConnection = null; let localChannel = null; let localSpeaker = null; let readyTimeout; let disconnectTimeout; let sessionUpdateSent = false; let resolveReady; let rejectReady; let sessionReady = 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 = { onSessionCreated: () => { if (!isCurrent(attempt) || !localChannel || sessionUpdateSent) return; sessionUpdateSent = true; sendSessionUpdate(localChannel, sessionUpdate); }, onSessionUpdated: () => { if (!isCurrent(attempt) || !sessionUpdateSent) return; const track = localMicrophone?.getAudioTracks()[0] || null; void audioSender.replaceTrack(track) .then(() => { if (isCurrent(attempt)) resolveReady?.(); }) .catch((error) => 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 answerSdp = await exchangeSdp(normalizeSdp(offerSdp), sessionId); if (!ensureCurrent()) return; sessionReady = new Promise((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(answerSdp) }); if (!ensureCurrent()) return; await sessionReady; 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 }; };