| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| function backoffMs(attempt, baseMs = 500, capMs = 10_000) { |
| const exp = Math.min(capMs, baseMs * 2 ** attempt); |
| return Math.floor(Math.random() * exp); |
| } |
| |
| const POLICY_VIOLATION = 1008; |
| const NORMAL_CLOSURE = 1000; |
| function closeReasonFrom(code, text) { |
| const reason = (text || '').toLowerCase(); |
| if (code === POLICY_VIOLATION) { |
| if (reason.includes('resume-expired')) |
| return 'resume_expired'; |
| if (reason.includes('bad-resume-token')) |
| return 'bad_resume_token'; |
| if (reason.includes('session-not-found')) |
| return 'session_not_found'; |
| if (reason.includes('session-ended')) |
| return 'session_ended'; |
| if (reason.includes('websocket-disabled')) |
| return 'websocket_disabled'; |
| } |
| if (code === NORMAL_CLOSURE && reason.includes('max-duration')) |
| return 'max_duration'; |
| if (reason.includes('idle')) |
| return 'idle_timeout'; |
| return 'unknown'; |
| } |
| |
| export class CallSocket { |
| constructor(opts) { |
| this.ws = null; |
| this.status = 'idle'; |
| this.lastServerSeq = 0; |
| this.outboundQueue = []; |
| this.heartbeatTimer = null; |
| this.lastServerActivityMs = 0; |
| this.deadPeerTimer = null; |
| this.reconnectAttempt = 0; |
| this.reconnectTimer = null; |
| this.connectedAtMs = 0; |
| this.shuttingDown = false; |
| this.listeners = {}; |
| this.url = opts.url; |
| this.resumeWindowMs = (opts.resumeWindowSec ?? 20) * 1000; |
| this.heartbeatIntervalMs = opts.heartbeatIntervalMs ?? 15_000; |
| this.WebSocketImpl = opts.webSocketImpl ?? WebSocket; |
| this.log = opts.log ?? { |
| info: (m, e) => console.info(`[callSocket] ${m}`, e ?? ''), |
| warn: (m, e) => console.warn(`[callSocket] ${m}`, e ?? ''), |
| error: (m, e) => console.error(`[callSocket] ${m}`, e ?? ''), |
| }; |
| } |
| |
| getStatus() { return this.status; } |
| on(evt, fn) { |
| let set = this.listeners[evt]; |
| if (!set) { |
| set = new Set(); |
| this.listeners[evt] = set; |
| } |
| set.add(fn); |
| return () => set.delete(fn); |
| } |
| |
| connect() { |
| if (this.status !== 'idle') { |
| this.log.warn('connect() called in non-idle state', { status: this.status }); |
| return; |
| } |
| this.openSocket(); |
| } |
| |
| |
| sendTranscript(p) { |
| this.enqueueLine({ type: 'transcript.final', ts: Date.now(), payload: p }); |
| } |
| |
| sendUiState(p) { |
| this.enqueueLine({ type: 'ui.state', ts: Date.now(), payload: p }); |
| } |
| |
| sendControl(action) { |
| this.enqueueLine({ type: 'call.control', ts: Date.now(), payload: { action } }); |
| } |
| |
| |
| |
| |
| sendTranscriptPartial(p) { |
| this.enqueueLine({ |
| type: 'transcript.partial', |
| ts: Date.now(), |
| payload: p, |
| }); |
| } |
| |
| |
| |
| |
| sendBargeIn(turn_id) { |
| this.enqueueLine({ |
| type: 'user.barge_in', |
| ts: Date.now(), |
| payload: { turn_id }, |
| }); |
| } |
| |
| end() { |
| if (this.shuttingDown) |
| return; |
| this.shuttingDown = true; |
| this.transition('draining'); |
| this.sendControl('end'); |
| |
| |
| |
| window.setTimeout(() => this.dispose('user_ended'), 400); |
| } |
| |
| |
| dispose(reason = 'unmounted', code = NORMAL_CLOSURE) { |
| if (this.status === 'closed') |
| return; |
| this.shuttingDown = true; |
| this.clearTimers(); |
| try { |
| this.ws?.close(code, reason); |
| } |
| catch { } |
| this.ws = null; |
| this.outboundQueue = []; |
| this.transition('closed'); |
| this.emit('closed', { reason, code }); |
| } |
| |
| openSocket() { |
| this.transition(this.reconnectAttempt === 0 ? 'connecting' : 'reconnecting'); |
| let ws; |
| try { |
| ws = new this.WebSocketImpl(this.url); |
| } |
| catch (err) { |
| this.log.error('WebSocket constructor threw', { err: String(err) }); |
| this.scheduleReconnect(); |
| return; |
| } |
| this.ws = ws; |
| this.lastServerActivityMs = Date.now(); |
| ws.addEventListener('open', () => { |
| this.reconnectAttempt = 0; |
| this.connectedAtMs = Date.now(); |
| this.transition('live'); |
| this.startHeartbeat(); |
| this.flushQueue(); |
| }); |
| ws.addEventListener('message', (ev) => this.onMessage(ev)); |
| ws.addEventListener('close', (ev) => this.onClose(ev)); |
| ws.addEventListener('error', () => { |
| |
| |
| |
| this.log.warn('socket error event'); |
| }); |
| } |
| onMessage(ev) { |
| this.lastServerActivityMs = Date.now(); |
| if (typeof ev.data !== 'string') |
| return; |
| let env = null; |
| try { |
| env = JSON.parse(ev.data); |
| } |
| catch { |
| this.log.warn('non-JSON frame dropped'); |
| return; |
| } |
| if (!env || typeof env.type !== 'string') |
| return; |
| |
| |
| if (typeof env.seq === 'number' && env.seq <= this.lastServerSeq) { |
| this.log.warn('non-monotonic seq', { |
| got: env.seq, last: this.lastServerSeq, type: env.type, |
| }); |
| } |
| else if (typeof env.seq === 'number') { |
| this.lastServerSeq = env.seq; |
| } |
| this.dispatch(env); |
| } |
| dispatch(env) { |
| const raw = env.payload; |
| switch (env.type) { |
| case 'call.state': { |
| const p = raw; |
| this.emit('callState', p); |
| if (p.status === 'ended') |
| this.dispose('user_ended'); |
| return; |
| } |
| case 'transcript.final': |
| this.emit('assistantTranscript', raw); |
| return; |
| case 'assistant.partial': |
| this.emit('assistantPartial', raw); |
| return; |
| case 'assistant.turn_end': |
| this.emit('assistantTurnEnd', raw); |
| return; |
| case 'assistant.cancel': |
| this.emit('assistantCancel', raw); |
| return; |
| case 'assistant.filler': |
| this.emit('assistantFiller', raw); |
| return; |
| case 'assistant.backchannel': |
| this.emit('assistantBackchannel', raw); |
| return; |
| case 'safety.notice': |
| this.emit('safetyNotice', env.payload); |
| return; |
| case 'error': |
| this.emit('serverError', raw); |
| return; |
| case 'pong': |
| this.emit('pong', undefined); |
| return; |
| case 'ping': |
| |
| |
| this.sendControl('ping'); |
| return; |
| default: |
| |
| this.log.info(`unknown server event: ${env.type}`); |
| } |
| } |
| onClose(ev) { |
| const reason = closeReasonFrom(ev.code, ev.reason); |
| this.clearTimers(); |
| |
| const terminal = new Set([ |
| 'user_ended', |
| 'max_duration', |
| 'resume_expired', |
| 'bad_resume_token', |
| 'session_not_found', |
| 'session_ended', |
| 'websocket_disabled', |
| ]); |
| if (this.shuttingDown || terminal.has(reason)) { |
| this.ws = null; |
| this.transition('closed'); |
| this.emit('closed', { reason, code: ev.code, detail: ev.reason }); |
| return; |
| } |
| |
| this.log.warn('socket dropped; scheduling reconnect', { |
| code: ev.code, reason: ev.reason, |
| }); |
| this.ws = null; |
| this.scheduleReconnect(); |
| } |
| scheduleReconnect() { |
| if (this.shuttingDown) |
| return; |
| const sinceLiveMs = this.connectedAtMs > 0 ? Date.now() - this.connectedAtMs : 0; |
| if (sinceLiveMs > this.resumeWindowMs) { |
| this.log.warn('resume window elapsed; giving up'); |
| this.dispose('resume_expired', POLICY_VIOLATION); |
| return; |
| } |
| const delay = backoffMs(this.reconnectAttempt); |
| this.reconnectAttempt += 1; |
| this.transition('reconnecting'); |
| this.reconnectTimer = window.setTimeout(() => { |
| this.reconnectTimer = null; |
| this.openSocket(); |
| }, delay); |
| } |
| |
| enqueueLine(env) { |
| const line = JSON.stringify(env); |
| if (this.status === 'live' && this.ws?.readyState === this.WebSocketImpl.OPEN) { |
| try { |
| this.ws.send(line); |
| } |
| catch (err) { |
| this.log.warn('send failed; queueing', { err: String(err) }); |
| this.outboundQueue.push(line); |
| } |
| } |
| else { |
| this.outboundQueue.push(line); |
| } |
| } |
| flushQueue() { |
| if (!this.ws || this.ws.readyState !== this.WebSocketImpl.OPEN) |
| return; |
| const q = this.outboundQueue; |
| this.outboundQueue = []; |
| for (const line of q) { |
| try { |
| this.ws.send(line); |
| } |
| catch (err) { |
| this.log.error('flush failed; re-queueing remaining', { err: String(err) }); |
| this.outboundQueue.push(line); |
| return; |
| } |
| } |
| } |
| startHeartbeat() { |
| this.clearHeartbeat(); |
| this.heartbeatTimer = window.setInterval(() => { |
| this.sendControl('ping'); |
| }, this.heartbeatIntervalMs); |
| |
| |
| |
| this.deadPeerTimer = window.setInterval(() => { |
| if (Date.now() - this.lastServerActivityMs > this.heartbeatIntervalMs * 2) { |
| this.log.warn('peer silent; forcing reconnect'); |
| try { |
| this.ws?.close(); |
| } |
| catch { } |
| } |
| }, this.heartbeatIntervalMs); |
| } |
| clearHeartbeat() { |
| if (this.heartbeatTimer !== null) { |
| window.clearInterval(this.heartbeatTimer); |
| this.heartbeatTimer = null; |
| } |
| if (this.deadPeerTimer !== null) { |
| window.clearInterval(this.deadPeerTimer); |
| this.deadPeerTimer = null; |
| } |
| } |
| clearTimers() { |
| this.clearHeartbeat(); |
| if (this.reconnectTimer !== null) { |
| window.clearTimeout(this.reconnectTimer); |
| this.reconnectTimer = null; |
| } |
| } |
| |
| emit(evt, payload) { |
| const set = this.listeners[evt]; |
| if (!set) |
| return; |
| for (const fn of set) { |
| try { |
| fn(payload); |
| } |
| catch (err) { |
| this.log.error('listener threw', { evt, err: String(err) }); |
| } |
| } |
| } |
| transition(next) { |
| if (this.status === next) |
| return; |
| this.status = next; |
| this.emit('statusChange', next); |
| } |
| } |
|
|