diff --git a/apps/desktop/src/renderer/hooks/useRealtimeConversation.ts b/apps/desktop/src/renderer/hooks/useRealtimeConversation.ts index 4c17483..ff9a29f 100644 --- a/apps/desktop/src/renderer/hooks/useRealtimeConversation.ts +++ b/apps/desktop/src/renderer/hooks/useRealtimeConversation.ts @@ -3,27 +3,28 @@ // Supabase realtime-token(IPC 경유)으로 ephemeral key를 받아 // 렌더러가 OpenAI와 직접 WebRTC 연결한다 (마이크 캡처 + 오디오 재생 자동). // -// 이벤트 매핑 (data channel 'oai-events'): -// response.created → thinking -// response.(output_)audio_transcript.delta → speaking + 스트리밍 텍스트 -// conversation.item.input_audio_transcription.completed → 유저 메시지 -// response.done → 어시스턴트 메시지 확정 + listening 복귀 +// 이 훅은 얇은 React 어댑터다: +// - 연결 수립/취소/재연결 정책 → services/realtimeSessionController.ts +// - 서버 이벤트 → 대화 상태 매핑 → services/realtimeConversationModel.ts +// - WebRTC·IPC·fetch IO → services/realtimeSession.ts (RealtimeTransport 포트) -import { useState, useRef, useCallback, useEffect } from 'react' +import { useEffect, useState, useSyncExternalStore } from 'react' 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' -const DATA_CHANNEL_NAME = 'oai-events' +export type { RealtimeConnectionState, RealtimeStartResult } -export type RealtimeConnectionState = 'idle' | 'connecting' | 'live' | 'error' - -interface RealtimeServerEvent { - type: string - delta?: string - transcript?: string - response?: { id?: string } - item_id?: string - error?: { message?: string } +export interface UseRealtimeConversationOptions { + /** 전송 계층 주입 (기본값: 브라우저 WebRTC) */ + createTransport?: RealtimeTransportFactory } export interface UseRealtimeConversationResult { @@ -34,273 +35,41 @@ export interface UseRealtimeConversationResult { streamingText: string streamingMsgId: string | null error: string | null - /** 연결 시작. 실패 시 false (로컬 파이프라인 fallback 유도) */ - start: () => Promise + /** + * 연결 시작. 이벤트 채널이 열린 뒤 'live'로 resolve한다. + * 'failed'면 로컬 파이프라인 fallback을 유도하고, 'cancelled'(시작 중 stop/언마운트)면 아무것도 하지 않는다. + */ + start: () => Promise stop: () => void + /** 연결 중이면 큐에 쌓았다가 채널이 열리면 보낸다 */ sendText: (text: string) => void cancelResponse: () => void clearMessages: () => void } -let messageSeq = 0 -function nextMessageId(prefix: string): string { - messageSeq += 1 - return `rt-${prefix}-${Date.now()}-${messageSeq}` -} - -export function useRealtimeConversation(): UseRealtimeConversationResult { - const [connectionState, setConnectionState] = useState('idle') - const [conversationState, setConversationState] = useState('idle') - const [messages, setMessages] = useState([]) - const [streamingText, setStreamingText] = useState('') - const [streamingMsgId, setStreamingMsgId] = useState(null) - const [error, setError] = useState(null) - - const pcRef = useRef(null) - const dcRef = useRef(null) - const micStreamRef = useRef(null) - const audioElRef = useRef(null) - const assistantBufferRef = useRef('') - const assistantMsgIdRef = useRef(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], +export function useRealtimeConversation( + options: UseRealtimeConversationOptions = {}, +): UseRealtimeConversationResult { + const [controller] = useState( + () => new RealtimeSessionController(options.createTransport ?? createBrowserRealtimeTransport), ) + const snapshot = useSyncExternalStore(controller.subscribe, controller.getSnapshot) - const start = useCallback(async (): Promise => { - if (pcRef.current) return true - - 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([]) - }, []) + // 언마운트 시 진행 중인 시작까지 취소하고 연결 정리 (컨트롤러는 재사용 가능 → StrictMode 재마운트 안전) + useEffect(() => controller.stop, [controller]) return { - connectionState, - conversationState, - messages, - isActive: connectionState === 'connecting' || connectionState === 'live', - streamingText, - streamingMsgId, - error, - start, - stop, - sendText, - cancelResponse, - clearMessages, + connectionState: snapshot.connectionState, + conversationState: snapshot.conversationState, + messages: snapshot.messages, + isActive: snapshot.connectionState === 'connecting' || snapshot.connectionState === 'live', + streamingText: snapshot.streamingText, + streamingMsgId: snapshot.streamingMsgId, + error: snapshot.error, + start: controller.start, + stop: controller.stop, + sendText: controller.sendText, + cancelResponse: controller.cancelResponse, + clearMessages: controller.clearMessages, } } diff --git a/apps/desktop/src/renderer/pages/VoiceConversationPage.tsx b/apps/desktop/src/renderer/pages/VoiceConversationPage.tsx index 23422dc..05d5791 100644 --- a/apps/desktop/src/renderer/pages/VoiceConversationPage.tsx +++ b/apps/desktop/src/renderer/pages/VoiceConversationPage.tsx @@ -124,8 +124,9 @@ export function VoiceConversationPage(): React.ReactElement { isStartingSessionRef.current = true try { if (isRealtime) { - const ok = await realtime.start() - if (!ok) { + // 'cancelled'(시작 중 정지)는 fallback 대상이 아니다 + const result = await realtime.start() + if (result === 'failed') { setBackend('local') setErrorBanner({ phase: 'llm', message: t('conversation.realtime.fallback') }) await window.electronAPI.voiceConversation.startSession() @@ -159,15 +160,21 @@ export function VoiceConversationPage(): React.ReactElement { try { if (isRealtime) { if (!realtime.isActive) { - const ok = await realtime.start() - if (!ok) { + const result = await realtime.start() + if (result === 'failed') { setBackend('local') setErrorBanner({ phase: 'llm', message: t('conversation.realtime.fallback') }) await window.electronAPI.voiceConversation.startSession() await window.electronAPI.voiceConversation.sendMessage({ text }) return } + if (result === 'cancelled') { + // 시작 중 사용자가 멈췄다 — 보내지 못한 입력을 되돌려 준다 + setTextInput((current) => current || text) + return + } } + // start()는 이벤트 채널 open 후 resolve하고, 연결 중이면 컨트롤러가 큐에 쌓는다 realtime.sendText(text) return } diff --git a/apps/desktop/src/renderer/services/realtimeConversationModel.ts b/apps/desktop/src/renderer/services/realtimeConversationModel.ts new file mode 100644 index 0000000..82e6018 --- /dev/null +++ b/apps/desktop/src/renderer/services/realtimeConversationModel.ts @@ -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 + } +} diff --git a/apps/desktop/src/renderer/services/realtimeSession.ts b/apps/desktop/src/renderer/services/realtimeSession.ts new file mode 100644 index 0000000..fee9746 --- /dev/null +++ b/apps/desktop/src/renderer/services/realtimeSession.ts @@ -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,