fix(conversation): make realtime sessions cancellable and recoverable
This commit is contained in:
parent
a85b24d385
commit
f4f9653361
6 changed files with 1248 additions and 279 deletions
|
|
@ -3,27 +3,28 @@
|
||||||
// Supabase realtime-token(IPC 경유)으로 ephemeral key를 받아
|
// Supabase realtime-token(IPC 경유)으로 ephemeral key를 받아
|
||||||
// 렌더러가 OpenAI와 직접 WebRTC 연결한다 (마이크 캡처 + 오디오 재생 자동).
|
// 렌더러가 OpenAI와 직접 WebRTC 연결한다 (마이크 캡처 + 오디오 재생 자동).
|
||||||
//
|
//
|
||||||
// 이벤트 매핑 (data channel 'oai-events'):
|
// 이 훅은 얇은 React 어댑터다:
|
||||||
// response.created → thinking
|
// - 연결 수립/취소/재연결 정책 → services/realtimeSessionController.ts
|
||||||
// response.(output_)audio_transcript.delta → speaking + 스트리밍 텍스트
|
// - 서버 이벤트 → 대화 상태 매핑 → services/realtimeConversationModel.ts
|
||||||
// conversation.item.input_audio_transcription.completed → 유저 메시지
|
// - WebRTC·IPC·fetch IO → services/realtimeSession.ts (RealtimeTransport 포트)
|
||||||
// response.done → 어시스턴트 메시지 확정 + listening 복귀
|
|
||||||
|
|
||||||
import { useState, useRef, useCallback, useEffect } from 'react'
|
import { useEffect, useState, useSyncExternalStore } from 'react'
|
||||||
import type { ConversationState, ConversationMessage } from '@d3ro/core/types'
|
import type { ConversationState, ConversationMessage } from '@d3ro/core/types'
|
||||||
|
import {
|
||||||
|
createBrowserRealtimeTransport,
|
||||||
|
type RealtimeTransportFactory,
|
||||||
|
} from '../services/realtimeSession'
|
||||||
|
import {
|
||||||
|
RealtimeSessionController,
|
||||||
|
type RealtimeConnectionState,
|
||||||
|
type RealtimeStartResult,
|
||||||
|
} from '../services/realtimeSessionController'
|
||||||
|
|
||||||
const REALTIME_CALLS_URL = 'https://api.openai.com/v1/realtime/calls'
|
export type { RealtimeConnectionState, RealtimeStartResult }
|
||||||
const DATA_CHANNEL_NAME = 'oai-events'
|
|
||||||
|
|
||||||
export type RealtimeConnectionState = 'idle' | 'connecting' | 'live' | 'error'
|
export interface UseRealtimeConversationOptions {
|
||||||
|
/** 전송 계층 주입 (기본값: 브라우저 WebRTC) */
|
||||||
interface RealtimeServerEvent {
|
createTransport?: RealtimeTransportFactory
|
||||||
type: string
|
|
||||||
delta?: string
|
|
||||||
transcript?: string
|
|
||||||
response?: { id?: string }
|
|
||||||
item_id?: string
|
|
||||||
error?: { message?: string }
|
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface UseRealtimeConversationResult {
|
export interface UseRealtimeConversationResult {
|
||||||
|
|
@ -34,273 +35,41 @@ export interface UseRealtimeConversationResult {
|
||||||
streamingText: string
|
streamingText: string
|
||||||
streamingMsgId: string | null
|
streamingMsgId: string | null
|
||||||
error: string | null
|
error: string | null
|
||||||
/** 연결 시작. 실패 시 false (로컬 파이프라인 fallback 유도) */
|
/**
|
||||||
start: () => Promise<boolean>
|
* 연결 시작. 이벤트 채널이 열린 뒤 'live'로 resolve한다.
|
||||||
|
* 'failed'면 로컬 파이프라인 fallback을 유도하고, 'cancelled'(시작 중 stop/언마운트)면 아무것도 하지 않는다.
|
||||||
|
*/
|
||||||
|
start: () => Promise<RealtimeStartResult>
|
||||||
stop: () => void
|
stop: () => void
|
||||||
|
/** 연결 중이면 큐에 쌓았다가 채널이 열리면 보낸다 */
|
||||||
sendText: (text: string) => void
|
sendText: (text: string) => void
|
||||||
cancelResponse: () => void
|
cancelResponse: () => void
|
||||||
clearMessages: () => void
|
clearMessages: () => void
|
||||||
}
|
}
|
||||||
|
|
||||||
let messageSeq = 0
|
export function useRealtimeConversation(
|
||||||
function nextMessageId(prefix: string): string {
|
options: UseRealtimeConversationOptions = {},
|
||||||
messageSeq += 1
|
): UseRealtimeConversationResult {
|
||||||
return `rt-${prefix}-${Date.now()}-${messageSeq}`
|
const [controller] = useState(
|
||||||
}
|
() => new RealtimeSessionController(options.createTransport ?? createBrowserRealtimeTransport),
|
||||||
|
|
||||||
export function useRealtimeConversation(): UseRealtimeConversationResult {
|
|
||||||
const [connectionState, setConnectionState] = useState<RealtimeConnectionState>('idle')
|
|
||||||
const [conversationState, setConversationState] = useState<ConversationState>('idle')
|
|
||||||
const [messages, setMessages] = useState<ConversationMessage[]>([])
|
|
||||||
const [streamingText, setStreamingText] = useState('')
|
|
||||||
const [streamingMsgId, setStreamingMsgId] = useState<string | null>(null)
|
|
||||||
const [error, setError] = useState<string | null>(null)
|
|
||||||
|
|
||||||
const pcRef = useRef<RTCPeerConnection | null>(null)
|
|
||||||
const dcRef = useRef<RTCDataChannel | null>(null)
|
|
||||||
const micStreamRef = useRef<MediaStream | null>(null)
|
|
||||||
const audioElRef = useRef<HTMLAudioElement | null>(null)
|
|
||||||
const assistantBufferRef = useRef('')
|
|
||||||
const assistantMsgIdRef = useRef<string | null>(null)
|
|
||||||
|
|
||||||
const cleanup = useCallback(() => {
|
|
||||||
dcRef.current?.close()
|
|
||||||
dcRef.current = null
|
|
||||||
pcRef.current?.close()
|
|
||||||
pcRef.current = null
|
|
||||||
micStreamRef.current?.getTracks().forEach((track) => track.stop())
|
|
||||||
micStreamRef.current = null
|
|
||||||
if (audioElRef.current) {
|
|
||||||
audioElRef.current.srcObject = null
|
|
||||||
audioElRef.current.remove()
|
|
||||||
audioElRef.current = null
|
|
||||||
}
|
|
||||||
assistantBufferRef.current = ''
|
|
||||||
assistantMsgIdRef.current = null
|
|
||||||
}, [])
|
|
||||||
|
|
||||||
// 언마운트 시 연결 정리
|
|
||||||
useEffect(() => cleanup, [cleanup])
|
|
||||||
|
|
||||||
const flushAssistantMessage = useCallback(() => {
|
|
||||||
const content = assistantBufferRef.current.trim()
|
|
||||||
const msgId = assistantMsgIdRef.current
|
|
||||||
assistantBufferRef.current = ''
|
|
||||||
assistantMsgIdRef.current = null
|
|
||||||
setStreamingText('')
|
|
||||||
setStreamingMsgId(null)
|
|
||||||
if (content && msgId) {
|
|
||||||
setMessages((prev) => [
|
|
||||||
...prev,
|
|
||||||
{ id: msgId, role: 'assistant', content, timestamp: Date.now() },
|
|
||||||
])
|
|
||||||
}
|
|
||||||
}, [])
|
|
||||||
|
|
||||||
const handleServerEvent = useCallback(
|
|
||||||
(event: RealtimeServerEvent) => {
|
|
||||||
switch (event.type) {
|
|
||||||
case 'response.created': {
|
|
||||||
assistantBufferRef.current = ''
|
|
||||||
assistantMsgIdRef.current = nextMessageId('assistant')
|
|
||||||
setConversationState('thinking')
|
|
||||||
break
|
|
||||||
}
|
|
||||||
|
|
||||||
// GA/신규 이벤트명 모두 수용
|
|
||||||
case 'response.output_audio_transcript.delta':
|
|
||||||
case 'response.audio_transcript.delta':
|
|
||||||
case 'response.output_text.delta':
|
|
||||||
case 'response.text.delta': {
|
|
||||||
if (event.delta) {
|
|
||||||
assistantBufferRef.current += event.delta
|
|
||||||
setStreamingMsgId(assistantMsgIdRef.current)
|
|
||||||
setStreamingText(assistantBufferRef.current)
|
|
||||||
setConversationState('speaking')
|
|
||||||
}
|
|
||||||
break
|
|
||||||
}
|
|
||||||
|
|
||||||
case 'conversation.item.input_audio_transcription.completed': {
|
|
||||||
const transcript = (event.transcript ?? '').trim()
|
|
||||||
if (transcript) {
|
|
||||||
setMessages((prev) => [
|
|
||||||
...prev,
|
|
||||||
{
|
|
||||||
id: nextMessageId('user'),
|
|
||||||
role: 'user',
|
|
||||||
content: transcript,
|
|
||||||
timestamp: Date.now(),
|
|
||||||
},
|
|
||||||
])
|
|
||||||
}
|
|
||||||
break
|
|
||||||
}
|
|
||||||
|
|
||||||
case 'response.done': {
|
|
||||||
flushAssistantMessage()
|
|
||||||
setConversationState('listening')
|
|
||||||
break
|
|
||||||
}
|
|
||||||
|
|
||||||
case 'error': {
|
|
||||||
setError(event.error?.message ?? 'Realtime error')
|
|
||||||
break
|
|
||||||
}
|
|
||||||
|
|
||||||
default:
|
|
||||||
break
|
|
||||||
}
|
|
||||||
},
|
|
||||||
[flushAssistantMessage],
|
|
||||||
)
|
)
|
||||||
|
const snapshot = useSyncExternalStore(controller.subscribe, controller.getSnapshot)
|
||||||
|
|
||||||
const start = useCallback(async (): Promise<boolean> => {
|
// 언마운트 시 진행 중인 시작까지 취소하고 연결 정리 (컨트롤러는 재사용 가능 → StrictMode 재마운트 안전)
|
||||||
if (pcRef.current) return true
|
useEffect(() => controller.stop, [controller])
|
||||||
|
|
||||||
setError(null)
|
|
||||||
setConnectionState('connecting')
|
|
||||||
|
|
||||||
try {
|
|
||||||
// 1) ephemeral token (Supabase realtime-token → OpenAI client_secrets)
|
|
||||||
const tokenResult = await window.electronAPI.voiceConversation.getRealtimeToken({})
|
|
||||||
if (!tokenResult.success) {
|
|
||||||
throw new Error(tokenResult.error?.message ?? 'token request failed')
|
|
||||||
}
|
|
||||||
const { value: ephemeralKey, model } = tokenResult.data
|
|
||||||
|
|
||||||
// 2) 마이크 + peer connection
|
|
||||||
const micStream = await navigator.mediaDevices.getUserMedia({ audio: true })
|
|
||||||
micStreamRef.current = micStream
|
|
||||||
|
|
||||||
const pc = new RTCPeerConnection()
|
|
||||||
pcRef.current = pc
|
|
||||||
const micTrack = micStream.getAudioTracks()[0]
|
|
||||||
if (micTrack) pc.addTrack(micTrack, micStream)
|
|
||||||
|
|
||||||
// 3) 원격 오디오 재생
|
|
||||||
const audioEl = document.createElement('audio')
|
|
||||||
audioEl.autoplay = true
|
|
||||||
audioElRef.current = audioEl
|
|
||||||
pc.ontrack = (e) => {
|
|
||||||
audioEl.srcObject = e.streams[0] ?? null
|
|
||||||
}
|
|
||||||
|
|
||||||
pc.onconnectionstatechange = () => {
|
|
||||||
const cs = pc.connectionState
|
|
||||||
if (cs === 'failed' || cs === 'disconnected' || cs === 'closed') {
|
|
||||||
setConnectionState((prev) => (prev === 'live' ? 'error' : prev))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// 4) 이벤트 데이터 채널
|
|
||||||
const dc = pc.createDataChannel(DATA_CHANNEL_NAME)
|
|
||||||
dcRef.current = dc
|
|
||||||
dc.addEventListener('message', (e: MessageEvent) => {
|
|
||||||
try {
|
|
||||||
handleServerEvent(JSON.parse(e.data as string) as RealtimeServerEvent)
|
|
||||||
} catch {
|
|
||||||
// 파싱 불가 이벤트 무시
|
|
||||||
}
|
|
||||||
})
|
|
||||||
dc.addEventListener('open', () => {
|
|
||||||
// 유저 발화 전사 활성화 — 미지원 모델이면 error 이벤트만 오고 대화는 유지됨
|
|
||||||
dc.send(
|
|
||||||
JSON.stringify({
|
|
||||||
type: 'session.update',
|
|
||||||
session: {
|
|
||||||
type: 'realtime',
|
|
||||||
audio: {
|
|
||||||
input: { transcription: { model: 'whisper-1' } },
|
|
||||||
},
|
|
||||||
},
|
|
||||||
}),
|
|
||||||
)
|
|
||||||
})
|
|
||||||
|
|
||||||
// 5) SDP 교환
|
|
||||||
const offer = await pc.createOffer()
|
|
||||||
await pc.setLocalDescription(offer)
|
|
||||||
|
|
||||||
const sdpResponse = await fetch(`${REALTIME_CALLS_URL}?model=${encodeURIComponent(model)}`, {
|
|
||||||
method: 'POST',
|
|
||||||
body: offer.sdp,
|
|
||||||
headers: {
|
|
||||||
Authorization: `Bearer ${ephemeralKey}`,
|
|
||||||
'Content-Type': 'application/sdp',
|
|
||||||
},
|
|
||||||
})
|
|
||||||
if (!sdpResponse.ok) {
|
|
||||||
const text = await sdpResponse.text()
|
|
||||||
throw new Error(`SDP exchange failed (${sdpResponse.status}): ${text.slice(0, 200)}`)
|
|
||||||
}
|
|
||||||
|
|
||||||
const answerSdp = await sdpResponse.text()
|
|
||||||
await pc.setRemoteDescription({ type: 'answer', sdp: answerSdp })
|
|
||||||
|
|
||||||
setConnectionState('live')
|
|
||||||
setConversationState('listening')
|
|
||||||
return true
|
|
||||||
} catch (err) {
|
|
||||||
cleanup()
|
|
||||||
setConnectionState('error')
|
|
||||||
setConversationState('idle')
|
|
||||||
setError(err instanceof Error ? err.message : String(err))
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
}, [cleanup, handleServerEvent])
|
|
||||||
|
|
||||||
const stop = useCallback(() => {
|
|
||||||
flushAssistantMessage()
|
|
||||||
cleanup()
|
|
||||||
setConnectionState('idle')
|
|
||||||
setConversationState('idle')
|
|
||||||
}, [cleanup, flushAssistantMessage])
|
|
||||||
|
|
||||||
const sendText = useCallback((text: string) => {
|
|
||||||
const dc = dcRef.current
|
|
||||||
if (!dc || dc.readyState !== 'open') return
|
|
||||||
setMessages((prev) => [
|
|
||||||
...prev,
|
|
||||||
{ id: nextMessageId('user'), role: 'user', content: text, timestamp: Date.now() },
|
|
||||||
])
|
|
||||||
dc.send(
|
|
||||||
JSON.stringify({
|
|
||||||
type: 'conversation.item.create',
|
|
||||||
item: {
|
|
||||||
type: 'message',
|
|
||||||
role: 'user',
|
|
||||||
content: [{ type: 'input_text', text }],
|
|
||||||
},
|
|
||||||
}),
|
|
||||||
)
|
|
||||||
dc.send(JSON.stringify({ type: 'response.create' }))
|
|
||||||
}, [])
|
|
||||||
|
|
||||||
const cancelResponse = useCallback(() => {
|
|
||||||
const dc = dcRef.current
|
|
||||||
if (!dc || dc.readyState !== 'open') return
|
|
||||||
dc.send(JSON.stringify({ type: 'response.cancel' }))
|
|
||||||
flushAssistantMessage()
|
|
||||||
setConversationState('listening')
|
|
||||||
}, [flushAssistantMessage])
|
|
||||||
|
|
||||||
const clearMessages = useCallback(() => {
|
|
||||||
setMessages([])
|
|
||||||
}, [])
|
|
||||||
|
|
||||||
return {
|
return {
|
||||||
connectionState,
|
connectionState: snapshot.connectionState,
|
||||||
conversationState,
|
conversationState: snapshot.conversationState,
|
||||||
messages,
|
messages: snapshot.messages,
|
||||||
isActive: connectionState === 'connecting' || connectionState === 'live',
|
isActive: snapshot.connectionState === 'connecting' || snapshot.connectionState === 'live',
|
||||||
streamingText,
|
streamingText: snapshot.streamingText,
|
||||||
streamingMsgId,
|
streamingMsgId: snapshot.streamingMsgId,
|
||||||
error,
|
error: snapshot.error,
|
||||||
start,
|
start: controller.start,
|
||||||
stop,
|
stop: controller.stop,
|
||||||
sendText,
|
sendText: controller.sendText,
|
||||||
cancelResponse,
|
cancelResponse: controller.cancelResponse,
|
||||||
clearMessages,
|
clearMessages: controller.clearMessages,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -124,8 +124,9 @@ export function VoiceConversationPage(): React.ReactElement {
|
||||||
isStartingSessionRef.current = true
|
isStartingSessionRef.current = true
|
||||||
try {
|
try {
|
||||||
if (isRealtime) {
|
if (isRealtime) {
|
||||||
const ok = await realtime.start()
|
// 'cancelled'(시작 중 정지)는 fallback 대상이 아니다
|
||||||
if (!ok) {
|
const result = await realtime.start()
|
||||||
|
if (result === 'failed') {
|
||||||
setBackend('local')
|
setBackend('local')
|
||||||
setErrorBanner({ phase: 'llm', message: t('conversation.realtime.fallback') })
|
setErrorBanner({ phase: 'llm', message: t('conversation.realtime.fallback') })
|
||||||
await window.electronAPI.voiceConversation.startSession()
|
await window.electronAPI.voiceConversation.startSession()
|
||||||
|
|
@ -159,15 +160,21 @@ export function VoiceConversationPage(): React.ReactElement {
|
||||||
try {
|
try {
|
||||||
if (isRealtime) {
|
if (isRealtime) {
|
||||||
if (!realtime.isActive) {
|
if (!realtime.isActive) {
|
||||||
const ok = await realtime.start()
|
const result = await realtime.start()
|
||||||
if (!ok) {
|
if (result === 'failed') {
|
||||||
setBackend('local')
|
setBackend('local')
|
||||||
setErrorBanner({ phase: 'llm', message: t('conversation.realtime.fallback') })
|
setErrorBanner({ phase: 'llm', message: t('conversation.realtime.fallback') })
|
||||||
await window.electronAPI.voiceConversation.startSession()
|
await window.electronAPI.voiceConversation.startSession()
|
||||||
await window.electronAPI.voiceConversation.sendMessage({ text })
|
await window.electronAPI.voiceConversation.sendMessage({ text })
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
if (result === 'cancelled') {
|
||||||
|
// 시작 중 사용자가 멈췄다 — 보내지 못한 입력을 되돌려 준다
|
||||||
|
setTextInput((current) => current || text)
|
||||||
|
return
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
// start()는 이벤트 채널 open 후 resolve하고, 연결 중이면 컨트롤러가 큐에 쌓는다
|
||||||
realtime.sendText(text)
|
realtime.sendText(text)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
|
||||||
139
apps/desktop/src/renderer/services/realtimeConversationModel.ts
Normal file
139
apps/desktop/src/renderer/services/realtimeConversationModel.ts
Normal file
|
|
@ -0,0 +1,139 @@
|
||||||
|
// src/renderer/services/realtimeConversationModel.ts
|
||||||
|
// OpenAI Realtime 서버 이벤트 → 대화 UI 상태 매핑 (순수 함수).
|
||||||
|
// IO·React에 의존하지 않으므로 node 환경에서 그대로 테스트할 수 있다.
|
||||||
|
//
|
||||||
|
// 이벤트 매핑 (data channel 'oai-events'):
|
||||||
|
// response.created → thinking
|
||||||
|
// response.(output_)audio_transcript.delta → speaking + 스트리밍 텍스트
|
||||||
|
// conversation.item.input_audio_transcription.completed → 유저 메시지
|
||||||
|
// response.done → 어시스턴트 메시지 확정 + listening 복귀
|
||||||
|
// error → error 메시지
|
||||||
|
|
||||||
|
import type { ConversationState, ConversationMessage } from '@d3ro/core/types'
|
||||||
|
|
||||||
|
export interface RealtimeServerEvent {
|
||||||
|
type: string
|
||||||
|
delta?: string
|
||||||
|
transcript?: string
|
||||||
|
response?: { id?: string }
|
||||||
|
item_id?: string
|
||||||
|
error?: { message?: string }
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface RealtimeConversationModel {
|
||||||
|
conversationState: ConversationState
|
||||||
|
messages: ConversationMessage[]
|
||||||
|
streamingText: string
|
||||||
|
streamingMsgId: string | null
|
||||||
|
error: string | null
|
||||||
|
/** 진행 중인 어시스턴트 응답 누적 텍스트 */
|
||||||
|
assistantBuffer: string
|
||||||
|
/** 진행 중인 어시스턴트 응답의 메시지 id (response.created에서 발급) */
|
||||||
|
assistantMsgId: string | null
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface RealtimeModelContext {
|
||||||
|
nextId: (prefix: 'user' | 'assistant') => string
|
||||||
|
now: () => number
|
||||||
|
}
|
||||||
|
|
||||||
|
export const INITIAL_REALTIME_MODEL: RealtimeConversationModel = {
|
||||||
|
conversationState: 'idle',
|
||||||
|
messages: [],
|
||||||
|
streamingText: '',
|
||||||
|
streamingMsgId: null,
|
||||||
|
error: null,
|
||||||
|
assistantBuffer: '',
|
||||||
|
assistantMsgId: null,
|
||||||
|
}
|
||||||
|
|
||||||
|
let messageSeq = 0
|
||||||
|
export const defaultRealtimeModelContext: RealtimeModelContext = {
|
||||||
|
nextId: (prefix) => {
|
||||||
|
messageSeq += 1
|
||||||
|
return `rt-${prefix}-${Date.now()}-${messageSeq}`
|
||||||
|
},
|
||||||
|
now: () => Date.now(),
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 진행 중인 어시스턴트 응답을 메시지로 확정하고 스트리밍 상태를 비운다. */
|
||||||
|
export function flushAssistantMessage(
|
||||||
|
model: RealtimeConversationModel,
|
||||||
|
ctx: RealtimeModelContext,
|
||||||
|
): RealtimeConversationModel {
|
||||||
|
const content = model.assistantBuffer.trim()
|
||||||
|
const msgId = model.assistantMsgId
|
||||||
|
const messages =
|
||||||
|
content && msgId
|
||||||
|
? [...model.messages, { id: msgId, role: 'assistant' as const, content, timestamp: ctx.now() }]
|
||||||
|
: model.messages
|
||||||
|
return {
|
||||||
|
...model,
|
||||||
|
messages,
|
||||||
|
assistantBuffer: '',
|
||||||
|
assistantMsgId: null,
|
||||||
|
streamingText: '',
|
||||||
|
streamingMsgId: null,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export function appendUserMessage(
|
||||||
|
model: RealtimeConversationModel,
|
||||||
|
content: string,
|
||||||
|
ctx: RealtimeModelContext,
|
||||||
|
): RealtimeConversationModel {
|
||||||
|
return {
|
||||||
|
...model,
|
||||||
|
messages: [
|
||||||
|
...model.messages,
|
||||||
|
{ id: ctx.nextId('user'), role: 'user', content, timestamp: ctx.now() },
|
||||||
|
],
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 서버 이벤트 하나를 적용한다. 관심 없는 이벤트면 같은 참조를 돌려준다. */
|
||||||
|
export function reduceRealtimeServerEvent(
|
||||||
|
model: RealtimeConversationModel,
|
||||||
|
event: RealtimeServerEvent,
|
||||||
|
ctx: RealtimeModelContext,
|
||||||
|
): RealtimeConversationModel {
|
||||||
|
switch (event.type) {
|
||||||
|
case 'response.created':
|
||||||
|
return {
|
||||||
|
...model,
|
||||||
|
assistantBuffer: '',
|
||||||
|
assistantMsgId: ctx.nextId('assistant'),
|
||||||
|
conversationState: 'thinking',
|
||||||
|
}
|
||||||
|
|
||||||
|
// GA/신규 이벤트명 모두 수용
|
||||||
|
case 'response.output_audio_transcript.delta':
|
||||||
|
case 'response.audio_transcript.delta':
|
||||||
|
case 'response.output_text.delta':
|
||||||
|
case 'response.text.delta': {
|
||||||
|
if (!event.delta) return model
|
||||||
|
const assistantBuffer = model.assistantBuffer + event.delta
|
||||||
|
return {
|
||||||
|
...model,
|
||||||
|
assistantBuffer,
|
||||||
|
streamingMsgId: model.assistantMsgId,
|
||||||
|
streamingText: assistantBuffer,
|
||||||
|
conversationState: 'speaking',
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
case 'conversation.item.input_audio_transcription.completed': {
|
||||||
|
const transcript = (event.transcript ?? '').trim()
|
||||||
|
return transcript ? appendUserMessage(model, transcript, ctx) : model
|
||||||
|
}
|
||||||
|
|
||||||
|
case 'response.done':
|
||||||
|
return { ...flushAssistantMessage(model, ctx), conversationState: 'listening' }
|
||||||
|
|
||||||
|
case 'error':
|
||||||
|
return { ...model, error: event.error?.message ?? 'Realtime error' }
|
||||||
|
|
||||||
|
default:
|
||||||
|
return model
|
||||||
|
}
|
||||||
|
}
|
||||||
308
apps/desktop/src/renderer/services/realtimeSession.ts
Normal file
308
apps/desktop/src/renderer/services/realtimeSession.ts
Normal file
|
|
@ -0,0 +1,308 @@
|
||||||
|
// src/renderer/services/realtimeSession.ts
|
||||||
|
// OpenAI Realtime 세션 전송 계층.
|
||||||
|
//
|
||||||
|
// - RealtimeTransport: 컨트롤러/훅이 의존하는 포트 (연결·송신·종료·이벤트 구독).
|
||||||
|
// - WebRtcRealtimeTransport: WebRTC 구현. 토큰 → 마이크 → PeerConnection → SDP 교환 →
|
||||||
|
// data channel 'open'까지를 취소 가능한 connect(signal) 하나로 묶는다.
|
||||||
|
// 각 await 뒤에 취소 여부를 확인해, 취소되면 그때까지 얻은 마이크·PC를 즉시 해제한다.
|
||||||
|
// - IO 세부(IPC 토큰, getUserMedia, RTCPeerConnection, fetch, <audio>)는 전부 주입받는다.
|
||||||
|
// createBrowserRealtimeTransport()가 실제 브라우저/Electron 의존성을 묶는 기본 조립점이다.
|
||||||
|
|
||||||
|
import type { RealtimeServerEvent } from './realtimeConversationModel'
|
||||||
|
|
||||||
|
export const REALTIME_CALLS_URL = 'https://api.openai.com/v1/realtime/calls'
|
||||||
|
export const REALTIME_DATA_CHANNEL_NAME = 'oai-events'
|
||||||
|
|
||||||
|
/** 클라이언트 → 서버 이벤트 (JSON 직렬화 가능한 객체) */
|
||||||
|
export interface RealtimeClientEvent {
|
||||||
|
type: string
|
||||||
|
[key: string]: unknown
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 컨트롤러가 의존하는 전송 포트. 인스턴스는 1회용이다 (connect는 한 번만). */
|
||||||
|
export interface RealtimeTransport {
|
||||||
|
/**
|
||||||
|
* 연결을 수립하고 이벤트 채널이 열리면 resolve한다.
|
||||||
|
* signal이 abort되면 잡고 있던 자원을 모두 해제하고 RealtimeAbortError로 reject한다.
|
||||||
|
*/
|
||||||
|
connect(signal: AbortSignal): Promise<void>
|
||||||
|
/** 채널이 열려 있으면 전송하고 true, 아니면 false */
|
||||||
|
send(event: RealtimeClientEvent): boolean
|
||||||
|
/** 자원 해제 (멱등). 호출자가 직접 닫은 경우 onClosed는 불리지 않는다. */
|
||||||
|
close(): void
|
||||||
|
onEvent(listener: (event: RealtimeServerEvent) => void): () => void
|
||||||
|
/** 연결이 열린 뒤 원격/네트워크 사유로 끊겼을 때 한 번 호출된다. */
|
||||||
|
onClosed(listener: () => void): () => void
|
||||||
|
}
|
||||||
|
|
||||||
|
export type RealtimeTransportFactory = () => RealtimeTransport
|
||||||
|
|
||||||
|
export class RealtimeAbortError extends Error {
|
||||||
|
constructor() {
|
||||||
|
super('realtime connect aborted')
|
||||||
|
this.name = 'RealtimeAbortError'
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export function isRealtimeAbortError(err: unknown): err is RealtimeAbortError {
|
||||||
|
return err instanceof RealtimeAbortError
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── WebRTC 구현 ────────────────────────────────────────────────
|
||||||
|
|
||||||
|
export interface RealtimeToken {
|
||||||
|
value: string
|
||||||
|
model: string
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface SdpHttpResponse {
|
||||||
|
ok: boolean
|
||||||
|
status: number
|
||||||
|
text(): Promise<string>
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 원격 오디오 재생 대상 */
|
||||||
|
export interface RealtimeAudioSink {
|
||||||
|
setStream(stream: MediaStream | null): void
|
||||||
|
dispose(): void
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface WebRtcTransportDeps {
|
||||||
|
getToken(): Promise<RealtimeToken>
|
||||||
|
getUserMedia(): Promise<MediaStream>
|
||||||
|
createPeerConnection(): RTCPeerConnection
|
||||||
|
postSdp(url: string, init: { body: string; headers: Record<string, string> }): Promise<SdpHttpResponse>
|
||||||
|
createAudioSink(): RealtimeAudioSink
|
||||||
|
callsUrl?: string
|
||||||
|
}
|
||||||
|
|
||||||
|
type Phase = 'new' | 'connecting' | 'open' | 'closed'
|
||||||
|
|
||||||
|
const TERMINAL_PC_STATES: ReadonlySet<RTCPeerConnectionState> = new Set(['failed', 'closed'])
|
||||||
|
|
||||||
|
export class WebRtcRealtimeTransport implements RealtimeTransport {
|
||||||
|
private phase: Phase = 'new'
|
||||||
|
private stream: MediaStream | null = null
|
||||||
|
private pc: RTCPeerConnection | null = null
|
||||||
|
private dc: RTCDataChannel | null = null
|
||||||
|
private sink: RealtimeAudioSink | null = null
|
||||||
|
private pendingOpen: { resolve: () => void; reject: (err: Error) => void } | null = null
|
||||||
|
private readonly eventListeners = new Set<(event: RealtimeServerEvent) => void>()
|
||||||
|
private readonly closedListeners = new Set<() => void>()
|
||||||
|
|
||||||
|
constructor(private readonly deps: WebRtcTransportDeps) {}
|
||||||
|
|
||||||
|
async connect(signal: AbortSignal): Promise<void> {
|
||||||
|
if (this.phase !== 'new') throw new Error('realtime transport is single-use')
|
||||||
|
this.phase = 'connecting'
|
||||||
|
|
||||||
|
const onAbort = (): void => this.close()
|
||||||
|
signal.addEventListener('abort', onAbort, { once: true })
|
||||||
|
try {
|
||||||
|
this.ensureActive(signal)
|
||||||
|
|
||||||
|
// 1) ephemeral token
|
||||||
|
const { value: ephemeralKey, model } = await this.deps.getToken()
|
||||||
|
this.ensureActive(signal)
|
||||||
|
|
||||||
|
// 2) 마이크 — 먼저 필드에 담아야 취소 시 close()가 트랙을 멈춘다
|
||||||
|
this.stream = await this.deps.getUserMedia()
|
||||||
|
this.ensureActive(signal)
|
||||||
|
|
||||||
|
// 3) peer connection + 원격 오디오
|
||||||
|
const pc = this.deps.createPeerConnection()
|
||||||
|
this.pc = pc
|
||||||
|
const micTrack = this.stream.getAudioTracks()[0]
|
||||||
|
if (micTrack) pc.addTrack(micTrack, this.stream)
|
||||||
|
|
||||||
|
const sink = this.deps.createAudioSink()
|
||||||
|
this.sink = sink
|
||||||
|
pc.ontrack = (e: RTCTrackEvent): void => {
|
||||||
|
sink.setStream(e.streams[0] ?? null)
|
||||||
|
}
|
||||||
|
pc.onconnectionstatechange = (): void => {
|
||||||
|
if (TERMINAL_PC_STATES.has(pc.connectionState)) {
|
||||||
|
this.handleRemoteClose(`peer connection ${pc.connectionState}`)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 4) 이벤트 데이터 채널
|
||||||
|
const dc = pc.createDataChannel(REALTIME_DATA_CHANNEL_NAME)
|
||||||
|
this.dc = dc
|
||||||
|
dc.onmessage = (e: MessageEvent): void => this.handleMessage(e.data)
|
||||||
|
dc.onopen = (): void => this.handleOpen()
|
||||||
|
dc.onclose = (): void => this.handleRemoteClose('data channel closed')
|
||||||
|
|
||||||
|
// 5) SDP 교환
|
||||||
|
const offer = await pc.createOffer()
|
||||||
|
this.ensureActive(signal)
|
||||||
|
await pc.setLocalDescription(offer)
|
||||||
|
this.ensureActive(signal)
|
||||||
|
|
||||||
|
const url = `${this.deps.callsUrl ?? REALTIME_CALLS_URL}?model=${encodeURIComponent(model)}`
|
||||||
|
const sdpResponse = await this.deps.postSdp(url, {
|
||||||
|
body: offer.sdp ?? '',
|
||||||
|
headers: {
|
||||||
|
Authorization: `Bearer ${ephemeralKey}`,
|
||||||
|
'Content-Type': 'application/sdp',
|
||||||
|
},
|
||||||
|
})
|
||||||
|
this.ensureActive(signal)
|
||||||
|
const answerSdp = await sdpResponse.text()
|
||||||
|
this.ensureActive(signal)
|
||||||
|
if (!sdpResponse.ok) {
|
||||||
|
throw new Error(`SDP exchange failed (${sdpResponse.status}): ${answerSdp.slice(0, 200)}`)
|
||||||
|
}
|
||||||
|
|
||||||
|
await pc.setRemoteDescription({ type: 'answer', sdp: answerSdp })
|
||||||
|
this.ensureActive(signal)
|
||||||
|
|
||||||
|
// 6) 이벤트 채널 open까지 대기 — 이후 send가 버려지지 않도록
|
||||||
|
await this.waitForOpen()
|
||||||
|
} catch (err) {
|
||||||
|
this.close()
|
||||||
|
throw err
|
||||||
|
} finally {
|
||||||
|
signal.removeEventListener('abort', onAbort)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
send(event: RealtimeClientEvent): boolean {
|
||||||
|
const dc = this.dc
|
||||||
|
if (this.phase !== 'open' || !dc || dc.readyState !== 'open') return false
|
||||||
|
dc.send(JSON.stringify(event))
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
close(): void {
|
||||||
|
this.phase = 'closed'
|
||||||
|
this.release(new RealtimeAbortError())
|
||||||
|
}
|
||||||
|
|
||||||
|
onEvent(listener: (event: RealtimeServerEvent) => void): () => void {
|
||||||
|
this.eventListeners.add(listener)
|
||||||
|
return () => this.eventListeners.delete(listener)
|
||||||
|
}
|
||||||
|
|
||||||
|
onClosed(listener: () => void): () => void {
|
||||||
|
this.closedListeners.add(listener)
|
||||||
|
return () => this.closedListeners.delete(listener)
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── 내부 ──
|
||||||
|
|
||||||
|
private ensureActive(signal: AbortSignal): void {
|
||||||
|
if (signal.aborted || this.phase === 'closed') {
|
||||||
|
this.close()
|
||||||
|
throw new RealtimeAbortError()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private waitForOpen(): Promise<void> {
|
||||||
|
if (this.phase === 'open') return Promise.resolve()
|
||||||
|
if (this.dc?.readyState === 'open') {
|
||||||
|
this.phase = 'open'
|
||||||
|
return Promise.resolve()
|
||||||
|
}
|
||||||
|
return new Promise<void>((resolve, reject) => {
|
||||||
|
this.pendingOpen = { resolve, reject }
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
private handleOpen(): void {
|
||||||
|
if (this.phase !== 'connecting') return
|
||||||
|
this.phase = 'open'
|
||||||
|
const pending = this.pendingOpen
|
||||||
|
this.pendingOpen = null
|
||||||
|
pending?.resolve()
|
||||||
|
}
|
||||||
|
|
||||||
|
private handleMessage(data: unknown): void {
|
||||||
|
if (this.phase === 'closed' || typeof data !== 'string') return
|
||||||
|
let parsed: unknown
|
||||||
|
try {
|
||||||
|
parsed = JSON.parse(data)
|
||||||
|
} catch {
|
||||||
|
return // 파싱 불가 이벤트 무시
|
||||||
|
}
|
||||||
|
if (!parsed || typeof parsed !== 'object') return
|
||||||
|
const event = parsed as RealtimeServerEvent
|
||||||
|
if (typeof event.type !== 'string') return
|
||||||
|
for (const listener of this.eventListeners) listener(event)
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 원격/네트워크 사유 종료: 자원을 스스로 해제하고, 열린 뒤였다면 onClosed를 한 번 알린다. */
|
||||||
|
private handleRemoteClose(reason: string): void {
|
||||||
|
if (this.phase === 'closed' || this.phase === 'new') return
|
||||||
|
const wasOpen = this.phase === 'open'
|
||||||
|
this.phase = 'closed'
|
||||||
|
this.release(new Error(`Realtime connection lost: ${reason}`))
|
||||||
|
if (wasOpen) {
|
||||||
|
for (const listener of this.closedListeners) listener()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private release(pendingReason: Error): void {
|
||||||
|
const pending = this.pendingOpen
|
||||||
|
this.pendingOpen = null
|
||||||
|
|
||||||
|
const dc = this.dc
|
||||||
|
this.dc = null
|
||||||
|
if (dc) {
|
||||||
|
dc.onopen = null
|
||||||
|
dc.onclose = null
|
||||||
|
dc.onmessage = null
|
||||||
|
dc.close()
|
||||||
|
}
|
||||||
|
const pc = this.pc
|
||||||
|
this.pc = null
|
||||||
|
if (pc) {
|
||||||
|
pc.ontrack = null
|
||||||
|
pc.onconnectionstatechange = null
|
||||||
|
pc.close()
|
||||||
|
}
|
||||||
|
this.stream?.getTracks().forEach((track) => track.stop())
|
||||||
|
this.stream = null
|
||||||
|
if (this.sink) {
|
||||||
|
this.sink.setStream(null)
|
||||||
|
this.sink.dispose()
|
||||||
|
this.sink = null
|
||||||
|
}
|
||||||
|
|
||||||
|
pending?.reject(pendingReason)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── 기본 조립 (브라우저/Electron 렌더러) ─────────────────────────
|
||||||
|
|
||||||
|
function createBrowserAudioSink(): RealtimeAudioSink {
|
||||||
|
const audioEl = document.createElement('audio')
|
||||||
|
audioEl.autoplay = true
|
||||||
|
return {
|
||||||
|
setStream: (stream) => {
|
||||||
|
audioEl.srcObject = stream
|
||||||
|
},
|
||||||
|
dispose: () => {
|
||||||
|
audioEl.srcObject = null
|
||||||
|
audioEl.remove()
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export const browserWebRtcDeps: WebRtcTransportDeps = {
|
||||||
|
getToken: async () => {
|
||||||
|
// Supabase realtime-token → OpenAI client_secrets (IPC 경유)
|
||||||
|
const result = await window.electronAPI.voiceConversation.getRealtimeToken({})
|
||||||
|
if (!result.success) {
|
||||||
|
throw new Error(result.error?.message ?? 'token request failed')
|
||||||
|
}
|
||||||
|
return { value: result.data.value, model: result.data.model }
|
||||||
|
},
|
||||||
|
getUserMedia: () => navigator.mediaDevices.getUserMedia({ audio: true }),
|
||||||
|
createPeerConnection: () => new RTCPeerConnection(),
|
||||||
|
postSdp: (url, init) => fetch(url, { method: 'POST', body: init.body, headers: init.headers }),
|
||||||
|
createAudioSink: createBrowserAudioSink,
|
||||||
|
}
|
||||||
|
|
||||||
|
export const createBrowserRealtimeTransport: RealtimeTransportFactory = () =>
|
||||||
|
new WebRtcRealtimeTransport(browserWebRtcDeps)
|
||||||
217
apps/desktop/src/renderer/services/realtimeSessionController.ts
Normal file
217
apps/desktop/src/renderer/services/realtimeSessionController.ts
Normal file
|
|
@ -0,0 +1,217 @@
|
||||||
|
// src/renderer/services/realtimeSessionController.ts
|
||||||
|
// Realtime 대화 세션 수명주기 컨트롤러 (React 비의존).
|
||||||
|
//
|
||||||
|
// - start/stop/언마운트 경쟁: 시도마다 AbortController를 두고, stop이 오면 abort + close.
|
||||||
|
// 취소된 시도는 'cancelled'로 끝나며 연결·마이크가 되살아나지 않는다.
|
||||||
|
// - 동시 start는 진행 중인 시도 하나로 합친다.
|
||||||
|
// - 연결이 끊기면(transport.onClosed) 전송 계층을 버리고 'error'로 전환해 재시작을 허용한다.
|
||||||
|
// - 연결 중 sendText는 큐에 쌓았다가 채널이 열리면 순서대로 보낸다.
|
||||||
|
// - 상태는 불변 스냅샷으로 노출해 useSyncExternalStore로 구독할 수 있다.
|
||||||
|
|
||||||
|
import {
|
||||||
|
INITIAL_REALTIME_MODEL,
|
||||||
|
appendUserMessage,
|
||||||
|
defaultRealtimeModelContext,
|
||||||
|
flushAssistantMessage,
|
||||||
|
reduceRealtimeServerEvent,
|
||||||
|
type RealtimeConversationModel,
|
||||||
|
type RealtimeModelContext,
|
||||||
|
type RealtimeServerEvent,
|
||||||
|
} from './realtimeConversationModel'
|
||||||
|
import type { RealtimeClientEvent, RealtimeTransport, RealtimeTransportFactory } from './realtimeSession'
|
||||||
|
|
||||||
|
export type RealtimeConnectionState = 'idle' | 'connecting' | 'live' | 'error'
|
||||||
|
|
||||||
|
/** 'failed'일 때만 로컬 파이프라인 fallback을 유도한다. 'cancelled'는 사용자가 멈춘 경우. */
|
||||||
|
export type RealtimeStartResult = 'live' | 'failed' | 'cancelled'
|
||||||
|
|
||||||
|
export interface RealtimeSessionSnapshot extends RealtimeConversationModel {
|
||||||
|
connectionState: RealtimeConnectionState
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 채널이 열리면 보내는 세션 설정 — 유저 발화 전사 활성화 (미지원 모델이면 error 이벤트만 오고 대화는 유지됨) */
|
||||||
|
export const REALTIME_SESSION_UPDATE: RealtimeClientEvent = {
|
||||||
|
type: 'session.update',
|
||||||
|
session: {
|
||||||
|
type: 'realtime',
|
||||||
|
audio: {
|
||||||
|
input: { transcription: { model: 'whisper-1' } },
|
||||||
|
},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
function userTextEvents(text: string): RealtimeClientEvent[] {
|
||||||
|
return [
|
||||||
|
{
|
||||||
|
type: 'conversation.item.create',
|
||||||
|
item: {
|
||||||
|
type: 'message',
|
||||||
|
role: 'user',
|
||||||
|
content: [{ type: 'input_text', text }],
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{ type: 'response.create' },
|
||||||
|
]
|
||||||
|
}
|
||||||
|
|
||||||
|
export class RealtimeSessionController {
|
||||||
|
private snapshot: RealtimeSessionSnapshot = { ...INITIAL_REALTIME_MODEL, connectionState: 'idle' }
|
||||||
|
private readonly listeners = new Set<() => void>()
|
||||||
|
private transport: RealtimeTransport | null = null
|
||||||
|
private attempt: AbortController | null = null
|
||||||
|
private pendingStart: Promise<RealtimeStartResult> | null = null
|
||||||
|
private outbox: RealtimeClientEvent[] = []
|
||||||
|
|
||||||
|
constructor(
|
||||||
|
private readonly createTransport: RealtimeTransportFactory,
|
||||||
|
private readonly ctx: RealtimeModelContext = defaultRealtimeModelContext,
|
||||||
|
) {}
|
||||||
|
|
||||||
|
// ── 구독 (useSyncExternalStore 계약) ──
|
||||||
|
|
||||||
|
subscribe = (listener: () => void): (() => void) => {
|
||||||
|
this.listeners.add(listener)
|
||||||
|
return () => {
|
||||||
|
this.listeners.delete(listener)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
getSnapshot = (): RealtimeSessionSnapshot => this.snapshot
|
||||||
|
|
||||||
|
// ── 명령 ──
|
||||||
|
|
||||||
|
start = (): Promise<RealtimeStartResult> => {
|
||||||
|
if (this.snapshot.connectionState === 'live' && this.transport) return Promise.resolve('live')
|
||||||
|
if (this.pendingStart) return this.pendingStart
|
||||||
|
|
||||||
|
const attempt = new AbortController()
|
||||||
|
const transport = this.createTransport()
|
||||||
|
this.attempt = attempt
|
||||||
|
this.transport = transport
|
||||||
|
this.outbox = []
|
||||||
|
transport.onEvent((event) => {
|
||||||
|
if (this.transport === transport) this.applyServerEvent(event)
|
||||||
|
})
|
||||||
|
// 채널이 열리자마자 끊겨 connect가 아직 반환 전이면, 그 시도를 실패로 마무리한다
|
||||||
|
let lostWhileConnecting = false
|
||||||
|
transport.onClosed(() => {
|
||||||
|
if (this.transport !== transport) return
|
||||||
|
if (this.snapshot.connectionState === 'connecting') lostWhileConnecting = true
|
||||||
|
else this.handleConnectionLost()
|
||||||
|
})
|
||||||
|
this.update({ connectionState: 'connecting', error: null })
|
||||||
|
|
||||||
|
const run = async (): Promise<RealtimeStartResult> => {
|
||||||
|
const isCurrent = (): boolean => !attempt.signal.aborted && this.transport === transport
|
||||||
|
const fail = (message: string): RealtimeStartResult => {
|
||||||
|
this.transport = null
|
||||||
|
this.attempt = null
|
||||||
|
this.outbox = []
|
||||||
|
this.update({
|
||||||
|
...flushAssistantMessage(this.snapshot, this.ctx),
|
||||||
|
connectionState: 'error',
|
||||||
|
conversationState: 'idle',
|
||||||
|
error: message,
|
||||||
|
})
|
||||||
|
return 'failed'
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
await transport.connect(attempt.signal)
|
||||||
|
} catch (err) {
|
||||||
|
transport.close()
|
||||||
|
if (!isCurrent()) return 'cancelled'
|
||||||
|
return fail(err instanceof Error ? err.message : String(err))
|
||||||
|
}
|
||||||
|
if (!isCurrent()) {
|
||||||
|
transport.close()
|
||||||
|
return 'cancelled'
|
||||||
|
}
|
||||||
|
if (lostWhileConnecting) {
|
||||||
|
transport.close()
|
||||||
|
return fail('Realtime connection lost')
|
||||||
|
}
|
||||||
|
|
||||||
|
this.attempt = null
|
||||||
|
transport.send(REALTIME_SESSION_UPDATE)
|
||||||
|
const queued = this.outbox
|
||||||
|
this.outbox = []
|
||||||
|
for (const event of queued) transport.send(event)
|
||||||
|
this.update({ connectionState: 'live', conversationState: 'listening' })
|
||||||
|
return 'live'
|
||||||
|
}
|
||||||
|
|
||||||
|
const pending = run().finally(() => {
|
||||||
|
if (this.pendingStart === pending) this.pendingStart = null
|
||||||
|
})
|
||||||
|
this.pendingStart = pending
|
||||||
|
return pending
|
||||||
|
}
|
||||||
|
|
||||||
|
stop = (): void => {
|
||||||
|
this.attempt?.abort()
|
||||||
|
this.attempt = null
|
||||||
|
const transport = this.transport
|
||||||
|
this.transport = null
|
||||||
|
this.pendingStart = null
|
||||||
|
this.outbox = []
|
||||||
|
transport?.close()
|
||||||
|
this.update({
|
||||||
|
...flushAssistantMessage(this.snapshot, this.ctx),
|
||||||
|
connectionState: 'idle',
|
||||||
|
conversationState: 'idle',
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
sendText = (text: string): void => {
|
||||||
|
const { connectionState } = this.snapshot
|
||||||
|
const transport = this.transport
|
||||||
|
if (!transport) return
|
||||||
|
|
||||||
|
if (connectionState === 'connecting') {
|
||||||
|
// 채널이 열리면 start()가 순서대로 보낸다
|
||||||
|
this.outbox.push(...userTextEvents(text))
|
||||||
|
this.update(appendUserMessage(this.snapshot, text, this.ctx))
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if (connectionState !== 'live') return
|
||||||
|
|
||||||
|
const [itemCreate, responseCreate] = userTextEvents(text)
|
||||||
|
if (!transport.send(itemCreate)) return
|
||||||
|
transport.send(responseCreate)
|
||||||
|
this.update(appendUserMessage(this.snapshot, text, this.ctx))
|
||||||
|
}
|
||||||
|
|
||||||
|
cancelResponse = (): void => {
|
||||||
|
const transport = this.transport
|
||||||
|
if (!transport || this.snapshot.connectionState !== 'live') return
|
||||||
|
if (!transport.send({ type: 'response.cancel' })) return
|
||||||
|
this.update({ ...flushAssistantMessage(this.snapshot, this.ctx), conversationState: 'listening' })
|
||||||
|
}
|
||||||
|
|
||||||
|
clearMessages = (): void => {
|
||||||
|
this.update({ messages: [] })
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── 내부 ──
|
||||||
|
|
||||||
|
private applyServerEvent(event: RealtimeServerEvent): void {
|
||||||
|
const next = reduceRealtimeServerEvent(this.snapshot, event, this.ctx)
|
||||||
|
if (next !== this.snapshot) this.update(next)
|
||||||
|
}
|
||||||
|
|
||||||
|
private handleConnectionLost(): void {
|
||||||
|
this.transport = null
|
||||||
|
this.attempt = null
|
||||||
|
this.outbox = []
|
||||||
|
this.update({
|
||||||
|
...flushAssistantMessage(this.snapshot, this.ctx),
|
||||||
|
connectionState: 'error',
|
||||||
|
conversationState: 'idle',
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
private update(patch: Partial<RealtimeSessionSnapshot>): void {
|
||||||
|
this.snapshot = { ...this.snapshot, ...patch }
|
||||||
|
for (const listener of this.listeners) listener()
|
||||||
|
}
|
||||||
|
}
|
||||||
529
apps/desktop/tests/unit/realtime-session-redteam-r1-30.test.ts
Normal file
529
apps/desktop/tests/unit/realtime-session-redteam-r1-30.test.ts
Normal file
|
|
@ -0,0 +1,529 @@
|
||||||
|
// tests/unit/realtime-session-redteam-r1-30.test.ts
|
||||||
|
// Realtime 대화 세션: 시작 중 취소 누수, 끊김 후 재연결, open 전 첫 메시지 유실 회귀 테스트.
|
||||||
|
// node 환경에서 fake 전송 계층/WebRTC 의존성으로 검증한다.
|
||||||
|
|
||||||
|
import { describe, it, expect, vi } from 'vitest'
|
||||||
|
import {
|
||||||
|
INITIAL_REALTIME_MODEL,
|
||||||
|
reduceRealtimeServerEvent,
|
||||||
|
type RealtimeModelContext,
|
||||||
|
type RealtimeServerEvent,
|
||||||
|
} from '../../src/renderer/services/realtimeConversationModel'
|
||||||
|
import {
|
||||||
|
RealtimeAbortError,
|
||||||
|
WebRtcRealtimeTransport,
|
||||||
|
type RealtimeClientEvent,
|
||||||
|
type RealtimeTransport,
|
||||||
|
type WebRtcTransportDeps,
|
||||||
|
} from '../../src/renderer/services/realtimeSession'
|
||||||
|
import {
|
||||||
|
REALTIME_SESSION_UPDATE,
|
||||||
|
RealtimeSessionController,
|
||||||
|
} from '../../src/renderer/services/realtimeSessionController'
|
||||||
|
|
||||||
|
// ── helpers ──
|
||||||
|
|
||||||
|
interface Deferred<T> {
|
||||||
|
promise: Promise<T>
|
||||||
|
resolve: (value: T) => void
|
||||||
|
reject: (err: Error) => void
|
||||||
|
}
|
||||||
|
|
||||||
|
function deferred<T>(): Deferred<T> {
|
||||||
|
let resolve!: (value: T) => void
|
||||||
|
let reject!: (err: Error) => void
|
||||||
|
const promise = new Promise<T>((res, rej) => {
|
||||||
|
resolve = res
|
||||||
|
reject = rej
|
||||||
|
})
|
||||||
|
return { promise, resolve, reject }
|
||||||
|
}
|
||||||
|
|
||||||
|
const flush = async (): Promise<void> => {
|
||||||
|
for (let i = 0; i < 10; i++) await Promise.resolve()
|
||||||
|
}
|
||||||
|
|
||||||
|
function makeCtx(): RealtimeModelContext {
|
||||||
|
let seq = 0
|
||||||
|
return {
|
||||||
|
nextId: (prefix) => `${prefix}-${++seq}`,
|
||||||
|
now: () => 1000,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── 순수 reducer ──
|
||||||
|
|
||||||
|
describe('reduceRealtimeServerEvent', () => {
|
||||||
|
it('maps a full assistant turn to thinking → speaking → message + listening', () => {
|
||||||
|
const ctx = makeCtx()
|
||||||
|
let m = reduceRealtimeServerEvent(INITIAL_REALTIME_MODEL, { type: 'response.created' }, ctx)
|
||||||
|
expect(m.conversationState).toBe('thinking')
|
||||||
|
expect(m.assistantMsgId).toBe('assistant-1')
|
||||||
|
|
||||||
|
m = reduceRealtimeServerEvent(m, { type: 'response.output_audio_transcript.delta', delta: '안녕' }, ctx)
|
||||||
|
m = reduceRealtimeServerEvent(m, { type: 'response.audio_transcript.delta', delta: '하세요' }, ctx)
|
||||||
|
expect(m.conversationState).toBe('speaking')
|
||||||
|
expect(m.streamingText).toBe('안녕하세요')
|
||||||
|
expect(m.streamingMsgId).toBe('assistant-1')
|
||||||
|
|
||||||
|
m = reduceRealtimeServerEvent(m, { type: 'response.done' }, ctx)
|
||||||
|
expect(m.conversationState).toBe('listening')
|
||||||
|
expect(m.streamingText).toBe('')
|
||||||
|
expect(m.streamingMsgId).toBeNull()
|
||||||
|
expect(m.messages).toEqual([
|
||||||
|
{ id: 'assistant-1', role: 'assistant', content: '안녕하세요', timestamp: 1000 },
|
||||||
|
])
|
||||||
|
})
|
||||||
|
|
||||||
|
it('adds trimmed user transcripts and ignores empty ones', () => {
|
||||||
|
const ctx = makeCtx()
|
||||||
|
const ev = (transcript: string): RealtimeServerEvent => ({
|
||||||
|
type: 'conversation.item.input_audio_transcription.completed',
|
||||||
|
transcript,
|
||||||
|
})
|
||||||
|
let m = reduceRealtimeServerEvent(INITIAL_REALTIME_MODEL, ev(' hi '), ctx)
|
||||||
|
expect(m.messages).toEqual([{ id: 'user-1', role: 'user', content: 'hi', timestamp: 1000 }])
|
||||||
|
const same = reduceRealtimeServerEvent(m, ev(' '), ctx)
|
||||||
|
expect(same).toBe(m)
|
||||||
|
m = reduceRealtimeServerEvent(m, { type: 'error', error: { message: 'boom' } }, ctx)
|
||||||
|
expect(m.error).toBe('boom')
|
||||||
|
expect(reduceRealtimeServerEvent(m, { type: 'error' }, ctx).error).toBe('Realtime error')
|
||||||
|
})
|
||||||
|
|
||||||
|
it('returns the same reference for unknown events and empty deltas', () => {
|
||||||
|
const ctx = makeCtx()
|
||||||
|
expect(reduceRealtimeServerEvent(INITIAL_REALTIME_MODEL, { type: 'session.created' }, ctx)).toBe(
|
||||||
|
INITIAL_REALTIME_MODEL,
|
||||||
|
)
|
||||||
|
expect(
|
||||||
|
reduceRealtimeServerEvent(INITIAL_REALTIME_MODEL, { type: 'response.text.delta', delta: '' }, ctx),
|
||||||
|
).toBe(INITIAL_REALTIME_MODEL)
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
||||||
|
// ── WebRTC transport (fake 의존성) ──
|
||||||
|
|
||||||
|
class FakeTrack {
|
||||||
|
stop = vi.fn()
|
||||||
|
}
|
||||||
|
|
||||||
|
class FakeStream {
|
||||||
|
readonly track = new FakeTrack()
|
||||||
|
getTracks = (): FakeTrack[] => [this.track]
|
||||||
|
getAudioTracks = (): FakeTrack[] => [this.track]
|
||||||
|
}
|
||||||
|
|
||||||
|
class FakeDataChannel {
|
||||||
|
readyState = 'connecting'
|
||||||
|
onopen: (() => void) | null = null
|
||||||
|
onclose: (() => void) | null = null
|
||||||
|
onmessage: ((e: { data: unknown }) => void) | null = null
|
||||||
|
send = vi.fn()
|
||||||
|
close = vi.fn(() => {
|
||||||
|
this.readyState = 'closed'
|
||||||
|
})
|
||||||
|
open(): void {
|
||||||
|
this.readyState = 'open'
|
||||||
|
this.onopen?.()
|
||||||
|
}
|
||||||
|
receive(data: unknown): void {
|
||||||
|
this.onmessage?.({ data })
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
class FakePeerConnection {
|
||||||
|
connectionState = 'new'
|
||||||
|
ontrack: ((e: { streams: unknown[] }) => void) | null = null
|
||||||
|
onconnectionstatechange: (() => void) | null = null
|
||||||
|
readonly dc = new FakeDataChannel()
|
||||||
|
readonly remote = deferred<void>()
|
||||||
|
addTrack = vi.fn()
|
||||||
|
createDataChannel = vi.fn(() => this.dc)
|
||||||
|
createOffer = vi.fn(async () => ({ type: 'offer', sdp: 'offer-sdp' }))
|
||||||
|
setLocalDescription = vi.fn(async () => undefined)
|
||||||
|
setRemoteDescription = vi.fn(() => this.remote.promise)
|
||||||
|
close = vi.fn()
|
||||||
|
setState(state: string): void {
|
||||||
|
this.connectionState = state
|
||||||
|
this.onconnectionstatechange?.()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
interface TransportHarness {
|
||||||
|
transport: WebRtcRealtimeTransport
|
||||||
|
token: Deferred<{ value: string; model: string }>
|
||||||
|
media: Deferred<FakeStream>
|
||||||
|
stream: FakeStream
|
||||||
|
pcs: FakePeerConnection[]
|
||||||
|
sink: { setStream: ReturnType<typeof vi.fn>; dispose: ReturnType<typeof vi.fn> }
|
||||||
|
deps: {
|
||||||
|
getUserMedia: ReturnType<typeof vi.fn>
|
||||||
|
createPeerConnection: ReturnType<typeof vi.fn>
|
||||||
|
postSdp: ReturnType<typeof vi.fn>
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function makeTransport(sdpOk = true): TransportHarness {
|
||||||
|
const token = deferred<{ value: string; model: string }>()
|
||||||
|
const media = deferred<FakeStream>()
|
||||||
|
const stream = new FakeStream()
|
||||||
|
const pcs: FakePeerConnection[] = []
|
||||||
|
const sink = { setStream: vi.fn(), dispose: vi.fn() }
|
||||||
|
const getUserMedia = vi.fn(() => media.promise)
|
||||||
|
const createPeerConnection = vi.fn(() => {
|
||||||
|
const pc = new FakePeerConnection()
|
||||||
|
pcs.push(pc)
|
||||||
|
return pc
|
||||||
|
})
|
||||||
|
const postSdp = vi.fn(async () => ({
|
||||||
|
ok: sdpOk,
|
||||||
|
status: sdpOk ? 201 : 401,
|
||||||
|
text: async () => (sdpOk ? 'answer-sdp' : 'unauthorized'),
|
||||||
|
}))
|
||||||
|
const deps = {
|
||||||
|
getToken: () => token.promise,
|
||||||
|
getUserMedia,
|
||||||
|
createPeerConnection,
|
||||||
|
postSdp,
|
||||||
|
createAudioSink: () => sink,
|
||||||
|
callsUrl: 'https://example.test/calls',
|
||||||
|
} as unknown as WebRtcTransportDeps
|
||||||
|
return {
|
||||||
|
transport: new WebRtcRealtimeTransport(deps),
|
||||||
|
token,
|
||||||
|
media,
|
||||||
|
stream,
|
||||||
|
pcs,
|
||||||
|
sink,
|
||||||
|
deps: { getUserMedia, createPeerConnection, postSdp },
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async function driveToRemoteDescription(h: TransportHarness): Promise<FakePeerConnection> {
|
||||||
|
h.token.resolve({ value: 'ek_1', model: 'gpt-realtime' })
|
||||||
|
await flush()
|
||||||
|
h.media.resolve(h.stream)
|
||||||
|
await flush()
|
||||||
|
const pc = h.pcs[0]
|
||||||
|
expect(pc).toBeDefined()
|
||||||
|
return pc
|
||||||
|
}
|
||||||
|
|
||||||
|
describe('WebRtcRealtimeTransport', () => {
|
||||||
|
it('resolves connect only after the data channel opens', async () => {
|
||||||
|
const h = makeTransport()
|
||||||
|
const ac = new AbortController()
|
||||||
|
let settled = false
|
||||||
|
const p = h.transport.connect(ac.signal).then(() => {
|
||||||
|
settled = true
|
||||||
|
})
|
||||||
|
const pc = await driveToRemoteDescription(h)
|
||||||
|
expect(h.deps.postSdp).toHaveBeenCalledWith(
|
||||||
|
'https://example.test/calls?model=gpt-realtime',
|
||||||
|
expect.objectContaining({ body: 'offer-sdp' }),
|
||||||
|
)
|
||||||
|
pc.remote.resolve()
|
||||||
|
await flush()
|
||||||
|
expect(settled).toBe(false)
|
||||||
|
expect(h.transport.send({ type: 'x' })).toBe(false)
|
||||||
|
|
||||||
|
pc.dc.open()
|
||||||
|
await p
|
||||||
|
expect(settled).toBe(true)
|
||||||
|
expect(h.transport.send({ type: 'x' })).toBe(true)
|
||||||
|
expect(pc.dc.send).toHaveBeenCalledWith(JSON.stringify({ type: 'x' }))
|
||||||
|
})
|
||||||
|
|
||||||
|
it('abort during the token request never opens the microphone', async () => {
|
||||||
|
const h = makeTransport()
|
||||||
|
const ac = new AbortController()
|
||||||
|
const p = h.transport.connect(ac.signal)
|
||||||
|
ac.abort()
|
||||||
|
h.token.resolve({ value: 'ek', model: 'm' })
|
||||||
|
await expect(p).rejects.toBeInstanceOf(RealtimeAbortError)
|
||||||
|
expect(h.deps.getUserMedia).not.toHaveBeenCalled()
|
||||||
|
expect(h.deps.createPeerConnection).not.toHaveBeenCalled()
|
||||||
|
})
|
||||||
|
|
||||||
|
it('abort during getUserMedia stops the late microphone stream and creates no peer connection', async () => {
|
||||||
|
const h = makeTransport()
|
||||||
|
const ac = new AbortController()
|
||||||
|
const p = h.transport.connect(ac.signal)
|
||||||
|
h.token.resolve({ value: 'ek', model: 'm' })
|
||||||
|
await flush()
|
||||||
|
expect(h.deps.getUserMedia).toHaveBeenCalledTimes(1)
|
||||||
|
ac.abort()
|
||||||
|
h.media.resolve(h.stream)
|
||||||
|
await expect(p).rejects.toBeInstanceOf(RealtimeAbortError)
|
||||||
|
expect(h.stream.track.stop).toHaveBeenCalled()
|
||||||
|
expect(h.deps.createPeerConnection).not.toHaveBeenCalled()
|
||||||
|
})
|
||||||
|
|
||||||
|
it('abort while waiting for the channel releases pc, channel, mic and audio sink', async () => {
|
||||||
|
const h = makeTransport()
|
||||||
|
const ac = new AbortController()
|
||||||
|
const p = h.transport.connect(ac.signal)
|
||||||
|
const pc = await driveToRemoteDescription(h)
|
||||||
|
pc.remote.resolve()
|
||||||
|
await flush()
|
||||||
|
ac.abort()
|
||||||
|
await expect(p).rejects.toBeInstanceOf(RealtimeAbortError)
|
||||||
|
expect(pc.close).toHaveBeenCalled()
|
||||||
|
expect(pc.dc.close).toHaveBeenCalled()
|
||||||
|
expect(h.stream.track.stop).toHaveBeenCalled()
|
||||||
|
expect(h.sink.dispose).toHaveBeenCalled()
|
||||||
|
})
|
||||||
|
|
||||||
|
it('a failed SDP exchange rejects with the status and releases resources', async () => {
|
||||||
|
const h = makeTransport(false)
|
||||||
|
const p = h.transport.connect(new AbortController().signal)
|
||||||
|
const pc = await driveToRemoteDescription(h)
|
||||||
|
await expect(p).rejects.toThrow('SDP exchange failed (401): unauthorized')
|
||||||
|
expect(pc.close).toHaveBeenCalled()
|
||||||
|
expect(h.stream.track.stop).toHaveBeenCalled()
|
||||||
|
})
|
||||||
|
|
||||||
|
it('a peer connection failure before open rejects connect', async () => {
|
||||||
|
const h = makeTransport()
|
||||||
|
const p = h.transport.connect(new AbortController().signal)
|
||||||
|
const pc = await driveToRemoteDescription(h)
|
||||||
|
pc.remote.resolve()
|
||||||
|
await flush()
|
||||||
|
pc.setState('failed')
|
||||||
|
await expect(p).rejects.toThrow('Realtime connection lost')
|
||||||
|
expect(h.stream.track.stop).toHaveBeenCalled()
|
||||||
|
})
|
||||||
|
|
||||||
|
it('after open, a failed connection closes itself and notifies onClosed exactly once', async () => {
|
||||||
|
const h = makeTransport()
|
||||||
|
const onClosed = vi.fn()
|
||||||
|
h.transport.onClosed(onClosed)
|
||||||
|
const p = h.transport.connect(new AbortController().signal)
|
||||||
|
const pc = await driveToRemoteDescription(h)
|
||||||
|
pc.remote.resolve()
|
||||||
|
await flush()
|
||||||
|
pc.dc.open()
|
||||||
|
await p
|
||||||
|
|
||||||
|
pc.setState('disconnected') // 일시적 끊김은 종료로 보지 않는다
|
||||||
|
expect(onClosed).not.toHaveBeenCalled()
|
||||||
|
|
||||||
|
pc.setState('failed')
|
||||||
|
pc.setState('closed')
|
||||||
|
expect(onClosed).toHaveBeenCalledTimes(1)
|
||||||
|
expect(pc.close).toHaveBeenCalled()
|
||||||
|
expect(h.stream.track.stop).toHaveBeenCalled()
|
||||||
|
expect(h.transport.send({ type: 'x' })).toBe(false)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('explicit close does not report onClosed, and forwards only valid JSON events', async () => {
|
||||||
|
const h = makeTransport()
|
||||||
|
const onClosed = vi.fn()
|
||||||
|
const onEvent = vi.fn()
|
||||||
|
h.transport.onClosed(onClosed)
|
||||||
|
h.transport.onEvent(onEvent)
|
||||||
|
const p = h.transport.connect(new AbortController().signal)
|
||||||
|
const pc = await driveToRemoteDescription(h)
|
||||||
|
pc.remote.resolve()
|
||||||
|
await flush()
|
||||||
|
pc.dc.open()
|
||||||
|
await p
|
||||||
|
|
||||||
|
pc.dc.receive('not json')
|
||||||
|
pc.dc.receive(JSON.stringify({ nope: true }))
|
||||||
|
pc.dc.receive(JSON.stringify({ type: 'response.created' }))
|
||||||
|
expect(onEvent).toHaveBeenCalledTimes(1)
|
||||||
|
expect(onEvent).toHaveBeenCalledWith({ type: 'response.created' })
|
||||||
|
|
||||||
|
h.transport.close()
|
||||||
|
expect(onClosed).not.toHaveBeenCalled()
|
||||||
|
expect(pc.close).toHaveBeenCalled()
|
||||||
|
})
|
||||||
|
|
||||||
|
it('is single-use', async () => {
|
||||||
|
const h = makeTransport()
|
||||||
|
h.transport.close()
|
||||||
|
await expect(h.transport.connect(new AbortController().signal)).rejects.toThrow('single-use')
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
||||||
|
// ── 컨트롤러 (fake transport) ──
|
||||||
|
|
||||||
|
class FakeTransport implements RealtimeTransport {
|
||||||
|
readonly connectCall = deferred<void>()
|
||||||
|
open = false
|
||||||
|
ignoreAbort = false
|
||||||
|
sent: RealtimeClientEvent[] = []
|
||||||
|
close = vi.fn(() => {
|
||||||
|
this.open = false
|
||||||
|
})
|
||||||
|
private eventCb: ((e: RealtimeServerEvent) => void) | null = null
|
||||||
|
private closedCb: (() => void) | null = null
|
||||||
|
|
||||||
|
connect(signal: AbortSignal): Promise<void> {
|
||||||
|
signal.addEventListener('abort', () => {
|
||||||
|
if (!this.ignoreAbort) this.connectCall.reject(new RealtimeAbortError())
|
||||||
|
})
|
||||||
|
return this.connectCall.promise.then(() => {
|
||||||
|
this.open = true
|
||||||
|
})
|
||||||
|
}
|
||||||
|
send(event: RealtimeClientEvent): boolean {
|
||||||
|
if (!this.open) return false
|
||||||
|
this.sent.push(event)
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
onEvent(cb: (e: RealtimeServerEvent) => void): () => void {
|
||||||
|
this.eventCb = cb
|
||||||
|
return () => undefined
|
||||||
|
}
|
||||||
|
onClosed(cb: () => void): () => void {
|
||||||
|
this.closedCb = cb
|
||||||
|
return () => undefined
|
||||||
|
}
|
||||||
|
emit(e: RealtimeServerEvent): void {
|
||||||
|
this.eventCb?.(e)
|
||||||
|
}
|
||||||
|
drop(): void {
|
||||||
|
this.open = false
|
||||||
|
this.closedCb?.()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function makeController(): { controller: RealtimeSessionController; transports: FakeTransport[] } {
|
||||||
|
const transports: FakeTransport[] = []
|
||||||
|
const controller = new RealtimeSessionController(() => {
|
||||||
|
const t = new FakeTransport()
|
||||||
|
transports.push(t)
|
||||||
|
return t
|
||||||
|
}, makeCtx())
|
||||||
|
return { controller, transports }
|
||||||
|
}
|
||||||
|
|
||||||
|
describe('RealtimeSessionController', () => {
|
||||||
|
it('stop during connecting cancels the attempt, closes the transport and never goes live', async () => {
|
||||||
|
const { controller, transports } = makeController()
|
||||||
|
const p = controller.start()
|
||||||
|
expect(controller.getSnapshot().connectionState).toBe('connecting')
|
||||||
|
controller.stop()
|
||||||
|
expect(transports[0].close).toHaveBeenCalled()
|
||||||
|
await expect(p).resolves.toBe('cancelled')
|
||||||
|
expect(controller.getSnapshot().connectionState).toBe('idle')
|
||||||
|
})
|
||||||
|
|
||||||
|
it('a connect that completes after stop (abort ignored) is closed and never goes live', async () => {
|
||||||
|
const { controller, transports } = makeController()
|
||||||
|
const p = controller.start()
|
||||||
|
const t = transports[0]
|
||||||
|
t.ignoreAbort = true
|
||||||
|
controller.stop()
|
||||||
|
t.close.mockClear()
|
||||||
|
t.connectCall.resolve()
|
||||||
|
await expect(p).resolves.toBe('cancelled')
|
||||||
|
expect(t.close).toHaveBeenCalled()
|
||||||
|
expect(t.sent).toEqual([])
|
||||||
|
expect(controller.getSnapshot().connectionState).toBe('idle')
|
||||||
|
})
|
||||||
|
|
||||||
|
it('concurrent start calls share one connection attempt', async () => {
|
||||||
|
const { controller, transports } = makeController()
|
||||||
|
const a = controller.start()
|
||||||
|
const b = controller.start()
|
||||||
|
expect(transports).toHaveLength(1)
|
||||||
|
transports[0].connectCall.resolve()
|
||||||
|
await expect(a).resolves.toBe('live')
|
||||||
|
await expect(b).resolves.toBe('live')
|
||||||
|
await expect(controller.start()).resolves.toBe('live')
|
||||||
|
expect(transports).toHaveLength(1)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('sends session.update when live and sendText right after start is delivered', async () => {
|
||||||
|
const { controller, transports } = makeController()
|
||||||
|
const p = controller.start()
|
||||||
|
transports[0].connectCall.resolve()
|
||||||
|
await expect(p).resolves.toBe('live')
|
||||||
|
controller.sendText('hello')
|
||||||
|
const t = transports[0]
|
||||||
|
expect(t.sent[0]).toEqual(REALTIME_SESSION_UPDATE)
|
||||||
|
expect(t.sent.slice(1).map((e) => e.type)).toEqual(['conversation.item.create', 'response.create'])
|
||||||
|
expect(controller.getSnapshot().messages.map((m) => m.content)).toEqual(['hello'])
|
||||||
|
expect(controller.getSnapshot().conversationState).toBe('listening')
|
||||||
|
})
|
||||||
|
|
||||||
|
it('queues text sent while connecting and flushes it after session.update', async () => {
|
||||||
|
const { controller, transports } = makeController()
|
||||||
|
const p = controller.start()
|
||||||
|
controller.sendText('first')
|
||||||
|
expect(controller.getSnapshot().messages.map((m) => m.content)).toEqual(['first'])
|
||||||
|
expect(transports[0].sent).toEqual([])
|
||||||
|
transports[0].connectCall.resolve()
|
||||||
|
await p
|
||||||
|
expect(transports[0].sent.map((e) => e.type)).toEqual([
|
||||||
|
'session.update',
|
||||||
|
'conversation.item.create',
|
||||||
|
'response.create',
|
||||||
|
])
|
||||||
|
})
|
||||||
|
|
||||||
|
it('after the connection drops, the session reports error and start reconnects with a new transport', async () => {
|
||||||
|
const { controller, transports } = makeController()
|
||||||
|
const p = controller.start()
|
||||||
|
transports[0].connectCall.resolve()
|
||||||
|
await p
|
||||||
|
transports[0].emit({ type: 'response.created' })
|
||||||
|
transports[0].emit({ type: 'response.text.delta', delta: 'partial' })
|
||||||
|
|
||||||
|
transports[0].drop()
|
||||||
|
const snap = controller.getSnapshot()
|
||||||
|
expect(snap.connectionState).toBe('error')
|
||||||
|
expect(snap.conversationState).toBe('idle')
|
||||||
|
expect(snap.messages.map((m) => m.content)).toEqual(['partial'])
|
||||||
|
|
||||||
|
const again = controller.start()
|
||||||
|
expect(transports).toHaveLength(2)
|
||||||
|
transports[1].connectCall.resolve()
|
||||||
|
await expect(again).resolves.toBe('live')
|
||||||
|
// 이전 전송 계층의 늦은 이벤트는 무시된다
|
||||||
|
transports[0].emit({ type: 'response.created' })
|
||||||
|
expect(controller.getSnapshot().conversationState).toBe('listening')
|
||||||
|
})
|
||||||
|
|
||||||
|
it('a failed connect reports failed with the error message and allows retry', async () => {
|
||||||
|
const { controller, transports } = makeController()
|
||||||
|
const p = controller.start()
|
||||||
|
transports[0].connectCall.reject(new Error('token request failed'))
|
||||||
|
await expect(p).resolves.toBe('failed')
|
||||||
|
expect(controller.getSnapshot()).toMatchObject({
|
||||||
|
connectionState: 'error',
|
||||||
|
conversationState: 'idle',
|
||||||
|
error: 'token request failed',
|
||||||
|
})
|
||||||
|
void controller.start()
|
||||||
|
expect(transports).toHaveLength(2)
|
||||||
|
expect(controller.getSnapshot().error).toBeNull()
|
||||||
|
})
|
||||||
|
|
||||||
|
it('notifies subscribers and applies server events only from the current transport', async () => {
|
||||||
|
const { controller, transports } = makeController()
|
||||||
|
const listener = vi.fn()
|
||||||
|
const unsubscribe = controller.subscribe(listener)
|
||||||
|
const p = controller.start()
|
||||||
|
transports[0].connectCall.resolve()
|
||||||
|
await p
|
||||||
|
listener.mockClear()
|
||||||
|
transports[0].emit({ type: 'response.created' })
|
||||||
|
expect(listener).toHaveBeenCalledTimes(1)
|
||||||
|
expect(controller.getSnapshot().conversationState).toBe('thinking')
|
||||||
|
controller.cancelResponse()
|
||||||
|
expect(transports[0].sent.at(-1)).toEqual({ type: 'response.cancel' })
|
||||||
|
expect(controller.getSnapshot().conversationState).toBe('listening')
|
||||||
|
unsubscribe()
|
||||||
|
})
|
||||||
|
|
||||||
|
it('sendText is ignored when no session exists', () => {
|
||||||
|
const { controller, transports } = makeController()
|
||||||
|
controller.sendText('x')
|
||||||
|
expect(transports).toHaveLength(0)
|
||||||
|
expect(controller.getSnapshot().messages).toEqual([])
|
||||||
|
})
|
||||||
|
})
|
||||||
Loading…
Add table
Add a link
Reference in a new issue