mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-02 16:31:34 +00:00
fix(xai): own realtime connection lifecycle
This commit is contained in:
@@ -24,6 +24,10 @@ import {
|
||||
type XaiRealtimeEvent,
|
||||
} from "./realtime-voice-config.js";
|
||||
import { XaiRealtimeMalformedAudioError, XaiRealtimeVoiceEvents } from "./realtime-voice-events.js";
|
||||
import {
|
||||
XaiRealtimeVoiceLifecycle,
|
||||
type XaiRealtimeVoiceConnection,
|
||||
} from "./realtime-voice-lifecycle.js";
|
||||
import { xaiUserAgentHeaderFor } from "./src/xai-user-agent.js";
|
||||
|
||||
export class XaiRealtimeVoiceBridge extends XaiRealtimeVoiceEvents implements RealtimeVoiceBridge {
|
||||
@@ -32,11 +36,10 @@ export class XaiRealtimeVoiceBridge extends XaiRealtimeVoiceEvents implements Re
|
||||
readonly supportsToolResultContinuation = false;
|
||||
|
||||
private ws: WebSocket | null = null;
|
||||
private connected = false;
|
||||
private sessionConfigured = false;
|
||||
private intentionallyClosed = false;
|
||||
private terminalError: Error | null = null;
|
||||
private reconnectAttempts = 0;
|
||||
private readonly lifecycle = new XaiRealtimeVoiceLifecycle();
|
||||
private connection: XaiRealtimeVoiceConnection | undefined;
|
||||
private connectPromise: Promise<void> | undefined;
|
||||
private pendingAudio: Buffer[] = [];
|
||||
private pendingAudioBytes = 0;
|
||||
private pendingToolResults: Array<{
|
||||
@@ -48,25 +51,35 @@ export class XaiRealtimeVoiceBridge extends XaiRealtimeVoiceEvents implements Re
|
||||
private connectionUrl = "";
|
||||
private readonly flowId = randomUUID();
|
||||
private sessionReadyFired = false;
|
||||
private reconnectAbortController = new AbortController();
|
||||
|
||||
async connect(): Promise<void> {
|
||||
if (this.terminalError) {
|
||||
throw this.terminalError;
|
||||
}
|
||||
this.intentionallyClosed = false;
|
||||
if (this.reconnectAbortController.signal.aborted) {
|
||||
this.reconnectAbortController = new AbortController();
|
||||
if (this.lifecycle.isReady()) {
|
||||
return;
|
||||
}
|
||||
if (this.connectPromise) {
|
||||
return this.connectPromise;
|
||||
}
|
||||
const connection = this.lifecycle.connect();
|
||||
this.connection = connection;
|
||||
const connectPromise = this.doConnect(connection);
|
||||
this.connectPromise = connectPromise;
|
||||
try {
|
||||
await connectPromise;
|
||||
} finally {
|
||||
if (this.connectPromise === connectPromise) {
|
||||
this.connectPromise = undefined;
|
||||
}
|
||||
}
|
||||
this.reconnectAttempts = 0;
|
||||
await this.doConnect();
|
||||
}
|
||||
|
||||
sendAudio(audio: Buffer): void {
|
||||
if (this.intentionallyClosed) {
|
||||
if (this.lifecycle.phase() === "terminal") {
|
||||
return;
|
||||
}
|
||||
if (!this.connected || !this.sessionConfigured || this.ws?.readyState !== WebSocket.OPEN) {
|
||||
if (!this.isConnected()) {
|
||||
this.enqueuePendingAudio(audio);
|
||||
return;
|
||||
}
|
||||
@@ -81,7 +94,7 @@ export class XaiRealtimeVoiceBridge extends XaiRealtimeVoiceEvents implements Re
|
||||
}
|
||||
|
||||
sendUserMessage(text: string): void {
|
||||
if (this.intentionallyClosed) {
|
||||
if (this.lifecycle.phase() === "terminal") {
|
||||
return;
|
||||
}
|
||||
if (!this.canSubmitInput()) {
|
||||
@@ -108,7 +121,7 @@ export class XaiRealtimeVoiceBridge extends XaiRealtimeVoiceEvents implements Re
|
||||
result: unknown,
|
||||
options?: RealtimeVoiceToolResultOptions,
|
||||
): void {
|
||||
if (this.intentionallyClosed) {
|
||||
if (this.lifecycle.phase() === "terminal") {
|
||||
return;
|
||||
}
|
||||
if (!this.canSubmitToolResult()) {
|
||||
@@ -125,24 +138,258 @@ export class XaiRealtimeVoiceBridge extends XaiRealtimeVoiceEvents implements Re
|
||||
}
|
||||
|
||||
close(): void {
|
||||
this.intentionallyClosed = true;
|
||||
// The bridge owns both its active socket and reconnect delay; canceling
|
||||
// both keeps terminal close from retaining callbacks for the full backoff.
|
||||
this.reconnectAbortController.abort();
|
||||
this.connected = false;
|
||||
this.sessionConfigured = false;
|
||||
this.resetTerminalState();
|
||||
if (this.ws) {
|
||||
this.ws.close(1000, "Bridge closed");
|
||||
this.ws = null;
|
||||
const connection = this.connection;
|
||||
if (!this.lifecycle.cancel()) {
|
||||
return;
|
||||
}
|
||||
this.resetTerminalState();
|
||||
if (!connection) {
|
||||
return;
|
||||
}
|
||||
const ws = this.ws;
|
||||
this.ws = null;
|
||||
if (ws?.readyState !== WebSocket.CLOSED) {
|
||||
ws?.close(1000, "Bridge closed");
|
||||
}
|
||||
this.notifyClose(connection, "completed");
|
||||
}
|
||||
|
||||
isConnected(): boolean {
|
||||
return this.connected && this.sessionConfigured;
|
||||
return this.lifecycle.isReady() && this.ws?.readyState === WebSocket.OPEN;
|
||||
}
|
||||
|
||||
private async doConnect(): Promise<void> {
|
||||
private async doConnect(connection: XaiRealtimeVoiceConnection): Promise<void> {
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
let settled = false;
|
||||
let reachedReady = false;
|
||||
let startupFailureClosing = false;
|
||||
let activeWs: WebSocket | undefined;
|
||||
let removeAbortListener = () => {};
|
||||
const settleResolve = () => {
|
||||
if (settled) {
|
||||
return;
|
||||
}
|
||||
settled = true;
|
||||
if (connectTimeout) {
|
||||
clearTimeout(connectTimeout);
|
||||
}
|
||||
removeAbortListener();
|
||||
resolve();
|
||||
};
|
||||
const settleReject = (error: Error) => {
|
||||
if (settled) {
|
||||
return;
|
||||
}
|
||||
settled = true;
|
||||
if (connectTimeout) {
|
||||
clearTimeout(connectTimeout);
|
||||
}
|
||||
removeAbortListener();
|
||||
reject(error);
|
||||
};
|
||||
const connectTimeout = setTimeout(() => {
|
||||
if (
|
||||
this.lifecycle.isCurrent(connection) &&
|
||||
!reachedReady &&
|
||||
this.lifecycle.terminalOutcome(connection) !== "completed"
|
||||
) {
|
||||
startupFailureClosing = true;
|
||||
activeWs?.terminate();
|
||||
settleReject(new Error("xAI realtime voice connection timeout"));
|
||||
}
|
||||
}, XAI_REALTIME_CONNECT_TIMEOUT_MS);
|
||||
const onAbort = () => {
|
||||
const terminalOutcome = this.lifecycle.terminalOutcome(connection);
|
||||
if (terminalOutcome !== "error" && activeWs && activeWs.readyState !== WebSocket.CLOSED) {
|
||||
activeWs.close(1000, "connection canceled");
|
||||
}
|
||||
if (terminalOutcome === "completed") {
|
||||
settleResolve();
|
||||
return;
|
||||
}
|
||||
const reason = connection.signal.reason;
|
||||
settleReject(reason instanceof Error ? reason : new Error(String(reason)));
|
||||
};
|
||||
connection.signal.addEventListener("abort", onAbort, { once: true });
|
||||
removeAbortListener = () => connection.signal.removeEventListener("abort", onAbort);
|
||||
if (connection.signal.aborted) {
|
||||
onAbort();
|
||||
}
|
||||
|
||||
const openWebSocket = (resolvedConnection: {
|
||||
url: string;
|
||||
headers: Record<string, string>;
|
||||
}) => {
|
||||
if (settled) {
|
||||
return;
|
||||
}
|
||||
if (!this.lifecycle.isCurrent(connection) || connection.signal.aborted) {
|
||||
settleResolve();
|
||||
return;
|
||||
}
|
||||
const { url, headers } = resolvedConnection;
|
||||
this.connectionUrl = url;
|
||||
const proxyAgent = createDebugProxyWebSocketAgent(resolveDebugProxySettings());
|
||||
const ws = new WebSocket(url, {
|
||||
headers,
|
||||
maxPayload: XAI_REALTIME_WS_MAX_PAYLOAD_BYTES,
|
||||
...(proxyAgent ? { agent: proxyAgent } : {}),
|
||||
});
|
||||
activeWs = ws;
|
||||
this.ws = ws;
|
||||
|
||||
const rejectStartup = (error: Error) => {
|
||||
if (!this.lifecycle.acceptsEvents(connection) || reachedReady) {
|
||||
return;
|
||||
}
|
||||
startupFailureClosing = true;
|
||||
settleReject(error);
|
||||
if (ws.readyState !== WebSocket.CLOSED) {
|
||||
ws.close(1000, "startup failed");
|
||||
}
|
||||
};
|
||||
|
||||
ws.on("open", () => {
|
||||
if (!this.lifecycle.acceptsEvents(connection)) {
|
||||
ws.close(1000, "stale connection");
|
||||
return;
|
||||
}
|
||||
// Resumed sessions replay prior items, so preserve unresolved tool calls until
|
||||
// their outputs are accepted on the replacement socket.
|
||||
this.resetRealtimeSessionState({
|
||||
preserveToolCallState:
|
||||
this.config.sessionResumption === true && this.conversationId !== null,
|
||||
});
|
||||
captureWsEvent({
|
||||
url,
|
||||
direction: "local",
|
||||
kind: "ws-open",
|
||||
flowId: this.flowId,
|
||||
meta: { provider: "xai", capability: "realtime-voice" },
|
||||
});
|
||||
this.sendEvent(this.buildSessionUpdate());
|
||||
});
|
||||
|
||||
ws.on("message", (data: Buffer) => {
|
||||
if (!this.lifecycle.acceptsEvents(connection) || this.ws !== ws) {
|
||||
return;
|
||||
}
|
||||
if (settled && !reachedReady) {
|
||||
return;
|
||||
}
|
||||
captureWsEvent({
|
||||
url,
|
||||
direction: "inbound",
|
||||
kind: "ws-frame",
|
||||
flowId: this.flowId,
|
||||
payload: data,
|
||||
meta: { provider: "xai", capability: "realtime-voice" },
|
||||
});
|
||||
try {
|
||||
const event = JSON.parse(data.toString()) as XaiRealtimeEvent;
|
||||
if (event.type === "error" && !reachedReady) {
|
||||
rejectStartup(new Error(readXaiRealtimeErrorDetail(event.error)));
|
||||
return;
|
||||
}
|
||||
this.handleEvent(event, connection);
|
||||
if (
|
||||
event.type === "session.updated" &&
|
||||
this.lifecycle.isCurrent(connection) &&
|
||||
this.lifecycle.isReady()
|
||||
) {
|
||||
reachedReady = true;
|
||||
settleResolve();
|
||||
}
|
||||
} catch (error) {
|
||||
if (error instanceof XaiRealtimeMalformedAudioError) {
|
||||
settleReject(error);
|
||||
this.failConnection(error, ws, connection);
|
||||
return;
|
||||
}
|
||||
console.error("[xai] realtime event parse failed:", error);
|
||||
}
|
||||
});
|
||||
|
||||
ws.on("error", (error) => {
|
||||
if (!this.lifecycle.acceptsEvents(connection) || this.ws !== ws) {
|
||||
return;
|
||||
}
|
||||
captureWsEvent({
|
||||
url,
|
||||
direction: "local",
|
||||
kind: "error",
|
||||
flowId: this.flowId,
|
||||
errorText: error instanceof Error ? error.message : String(error),
|
||||
meta: { provider: "xai", capability: "realtime-voice" },
|
||||
});
|
||||
if (!reachedReady) {
|
||||
rejectStartup(error instanceof Error ? error : new Error(String(error)));
|
||||
return;
|
||||
}
|
||||
this.config.onError?.(error instanceof Error ? error : new Error(String(error)));
|
||||
});
|
||||
|
||||
ws.on("close", (code, reasonBuffer) => {
|
||||
captureWsEvent({
|
||||
url,
|
||||
direction: "local",
|
||||
kind: "ws-close",
|
||||
flowId: this.flowId,
|
||||
closeCode: typeof code === "number" ? code : undefined,
|
||||
meta: {
|
||||
provider: "xai",
|
||||
capability: "realtime-voice",
|
||||
reason:
|
||||
Buffer.isBuffer(reasonBuffer) && reasonBuffer.length > 0
|
||||
? reasonBuffer.toString("utf8")
|
||||
: undefined,
|
||||
},
|
||||
});
|
||||
if (!this.lifecycle.isCurrent(connection)) {
|
||||
return;
|
||||
}
|
||||
if (this.ws === ws) {
|
||||
this.ws = null;
|
||||
}
|
||||
if (startupFailureClosing) {
|
||||
return;
|
||||
}
|
||||
if (this.terminalError) {
|
||||
this.notifyClose(connection, "error");
|
||||
return;
|
||||
}
|
||||
if (this.lifecycle.terminalOutcome(connection) === "completed") {
|
||||
settleResolve();
|
||||
this.notifyClose(connection, "completed");
|
||||
return;
|
||||
}
|
||||
if (!reachedReady && !settled) {
|
||||
settleReject(new Error("xAI realtime voice connection closed before ready"));
|
||||
return;
|
||||
}
|
||||
void this.attemptReconnect("websocket-close", connection);
|
||||
});
|
||||
};
|
||||
|
||||
void this.resolveConnectionParams()
|
||||
.then(openWebSocket)
|
||||
.catch((error: unknown) => {
|
||||
if (
|
||||
!this.lifecycle.isCurrent(connection) ||
|
||||
this.lifecycle.terminalOutcome(connection) === "completed"
|
||||
) {
|
||||
settleResolve();
|
||||
return;
|
||||
}
|
||||
settleReject(error instanceof Error ? error : new Error(String(error)));
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
private async resolveConnectionParams(): Promise<{
|
||||
url: string;
|
||||
headers: Record<string, string>;
|
||||
}> {
|
||||
const apiKey = this.config.resolveApiKey
|
||||
? await this.config.resolveApiKey()
|
||||
: await resolveXaiRealtimeApiKey(this.config.apiKey, this.config.cfg);
|
||||
@@ -152,209 +399,19 @@ export class XaiRealtimeVoiceBridge extends XaiRealtimeVoiceEvents implements Re
|
||||
model,
|
||||
this.config.sessionResumption === true ? (this.conversationId ?? undefined) : undefined,
|
||||
);
|
||||
const headers = {
|
||||
Authorization: `Bearer ${apiKey}`,
|
||||
...xaiUserAgentHeaderFor(this.config.baseUrl),
|
||||
return {
|
||||
url,
|
||||
headers: {
|
||||
Authorization: `Bearer ${apiKey}`,
|
||||
...xaiUserAgentHeaderFor(this.config.baseUrl),
|
||||
},
|
||||
};
|
||||
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
let settled = false;
|
||||
let startupFailureClosing = false;
|
||||
let terminalFailure: Error | undefined;
|
||||
let terminalCloseNotified = false;
|
||||
const settleResolve = () => {
|
||||
if (!settled) {
|
||||
settled = true;
|
||||
clearTimeout(connectTimeout);
|
||||
resolve();
|
||||
}
|
||||
};
|
||||
const settleReject = (error: Error) => {
|
||||
if (!settled) {
|
||||
settled = true;
|
||||
clearTimeout(connectTimeout);
|
||||
reject(error);
|
||||
}
|
||||
};
|
||||
const notifyTerminalClose = (error: Error) => {
|
||||
settleReject(error);
|
||||
if (terminalCloseNotified) {
|
||||
return;
|
||||
}
|
||||
terminalCloseNotified = true;
|
||||
this.config.onClose?.("error");
|
||||
};
|
||||
const connectTimeout = setTimeout(() => {
|
||||
if (!this.sessionConfigured && !this.intentionallyClosed) {
|
||||
startupFailureClosing = true;
|
||||
this.ws?.terminate();
|
||||
settleReject(new Error("xAI realtime voice connection timeout"));
|
||||
}
|
||||
}, XAI_REALTIME_CONNECT_TIMEOUT_MS);
|
||||
|
||||
if (this.intentionallyClosed) {
|
||||
settleResolve();
|
||||
return;
|
||||
}
|
||||
|
||||
this.connectionUrl = url;
|
||||
const proxyAgent = createDebugProxyWebSocketAgent(resolveDebugProxySettings());
|
||||
const ws = new WebSocket(url, {
|
||||
headers,
|
||||
maxPayload: XAI_REALTIME_WS_MAX_PAYLOAD_BYTES,
|
||||
...(proxyAgent ? { agent: proxyAgent } : {}),
|
||||
});
|
||||
this.ws = ws;
|
||||
|
||||
const rejectStartup = (error: Error) => {
|
||||
startupFailureClosing = true;
|
||||
settleReject(error);
|
||||
if (ws.readyState !== WebSocket.CLOSED) {
|
||||
ws.close(1000, "startup failed");
|
||||
}
|
||||
};
|
||||
const failConnection = (error: Error) => {
|
||||
if (terminalFailure) {
|
||||
return;
|
||||
}
|
||||
terminalFailure = error;
|
||||
this.terminalError = error;
|
||||
this.intentionallyClosed = true;
|
||||
this.reconnectAbortController.abort();
|
||||
this.connected = false;
|
||||
this.sessionConfigured = false;
|
||||
this.resetTerminalState();
|
||||
try {
|
||||
this.config.onError?.(error);
|
||||
} finally {
|
||||
if (ws.readyState !== WebSocket.CLOSED) {
|
||||
ws.close(1002, "Malformed audio payload");
|
||||
} else {
|
||||
this.ws = null;
|
||||
notifyTerminalClose(error);
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
ws.on("open", () => {
|
||||
// Resumed sessions replay prior items, so preserve unresolved tool calls until
|
||||
// their outputs are accepted on the replacement socket.
|
||||
this.resetRealtimeSessionState({
|
||||
preserveToolCallState:
|
||||
this.config.sessionResumption === true && this.conversationId !== null,
|
||||
});
|
||||
this.connected = true;
|
||||
this.sessionConfigured = false;
|
||||
captureWsEvent({
|
||||
url,
|
||||
direction: "local",
|
||||
kind: "ws-open",
|
||||
flowId: this.flowId,
|
||||
meta: { provider: "xai", capability: "realtime-voice" },
|
||||
});
|
||||
this.sendEvent(this.buildSessionUpdate());
|
||||
});
|
||||
|
||||
ws.on("message", (data: Buffer) => {
|
||||
if (settled && !this.sessionConfigured) {
|
||||
return;
|
||||
}
|
||||
captureWsEvent({
|
||||
url,
|
||||
direction: "inbound",
|
||||
kind: "ws-frame",
|
||||
flowId: this.flowId,
|
||||
payload: data,
|
||||
meta: { provider: "xai", capability: "realtime-voice" },
|
||||
});
|
||||
try {
|
||||
const event = JSON.parse(data.toString()) as XaiRealtimeEvent;
|
||||
if (event.type === "error" && !this.sessionConfigured) {
|
||||
rejectStartup(new Error(readXaiRealtimeErrorDetail(event.error)));
|
||||
return;
|
||||
}
|
||||
this.handleEvent(event);
|
||||
if (event.type === "session.updated") {
|
||||
settleResolve();
|
||||
}
|
||||
} catch (error) {
|
||||
if (error instanceof XaiRealtimeMalformedAudioError) {
|
||||
failConnection(error);
|
||||
return;
|
||||
}
|
||||
console.error("[xai] realtime event parse failed:", error);
|
||||
}
|
||||
});
|
||||
|
||||
ws.on("error", (error) => {
|
||||
captureWsEvent({
|
||||
url,
|
||||
direction: "local",
|
||||
kind: "error",
|
||||
flowId: this.flowId,
|
||||
errorText: error instanceof Error ? error.message : String(error),
|
||||
meta: { provider: "xai", capability: "realtime-voice" },
|
||||
});
|
||||
if (!this.sessionConfigured) {
|
||||
rejectStartup(error instanceof Error ? error : new Error(String(error)));
|
||||
return;
|
||||
}
|
||||
this.config.onError?.(error instanceof Error ? error : new Error(String(error)));
|
||||
});
|
||||
|
||||
ws.on("close", (code, reasonBuffer) => {
|
||||
captureWsEvent({
|
||||
url,
|
||||
direction: "local",
|
||||
kind: "ws-close",
|
||||
flowId: this.flowId,
|
||||
closeCode: typeof code === "number" ? code : undefined,
|
||||
meta: {
|
||||
provider: "xai",
|
||||
capability: "realtime-voice",
|
||||
reason:
|
||||
Buffer.isBuffer(reasonBuffer) && reasonBuffer.length > 0
|
||||
? reasonBuffer.toString("utf8")
|
||||
: undefined,
|
||||
},
|
||||
});
|
||||
if (terminalFailure) {
|
||||
if (this.ws === ws) {
|
||||
this.ws = null;
|
||||
}
|
||||
this.connected = false;
|
||||
this.sessionConfigured = false;
|
||||
notifyTerminalClose(terminalFailure);
|
||||
return;
|
||||
}
|
||||
if (startupFailureClosing) {
|
||||
if (this.ws === ws) {
|
||||
this.connected = false;
|
||||
this.sessionConfigured = false;
|
||||
}
|
||||
return;
|
||||
}
|
||||
const wasSessionConfigured = this.sessionConfigured;
|
||||
this.connected = false;
|
||||
this.sessionConfigured = false;
|
||||
if (this.intentionallyClosed) {
|
||||
settleResolve();
|
||||
this.config.onClose?.("completed");
|
||||
return;
|
||||
}
|
||||
if (!wasSessionConfigured && !settled) {
|
||||
settleReject(new Error("xAI realtime voice connection closed before ready"));
|
||||
return;
|
||||
}
|
||||
void this.attemptReconnect("websocket-close");
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
private async attemptReconnect(reason: string): Promise<void> {
|
||||
if (this.intentionallyClosed) {
|
||||
return;
|
||||
}
|
||||
private async attemptReconnect(
|
||||
reason: string,
|
||||
connection: XaiRealtimeVoiceConnection,
|
||||
): Promise<void> {
|
||||
const blocked = this.reconnectBlockReason();
|
||||
if (blocked) {
|
||||
this.config.onEvent?.({
|
||||
@@ -362,53 +419,59 @@ export class XaiRealtimeVoiceBridge extends XaiRealtimeVoiceEvents implements Re
|
||||
type: "session.reconnect.blocked",
|
||||
detail: `reason=${reason} ${blocked}`,
|
||||
});
|
||||
this.enterTerminalState();
|
||||
this.config.onClose?.("error");
|
||||
this.enterTerminalState(connection);
|
||||
return;
|
||||
}
|
||||
if (this.reconnectAttempts >= XAI_REALTIME_MAX_RECONNECT_ATTEMPTS) {
|
||||
const retry = this.lifecycle.retry(connection, XAI_REALTIME_MAX_RECONNECT_ATTEMPTS);
|
||||
if (!retry) {
|
||||
return;
|
||||
}
|
||||
if (retry === "exhausted") {
|
||||
this.config.onEvent?.({
|
||||
direction: "client",
|
||||
type: "session.reconnect.exhausted",
|
||||
detail: `reason=${reason} attempts=${this.reconnectAttempts}`,
|
||||
detail: `reason=${reason} attempts=${XAI_REALTIME_MAX_RECONNECT_ATTEMPTS}`,
|
||||
});
|
||||
this.enterTerminalState();
|
||||
this.config.onClose?.("error");
|
||||
this.enterTerminalState(connection);
|
||||
return;
|
||||
}
|
||||
this.reconnectAttempts += 1;
|
||||
const attempt = this.reconnectAttempts;
|
||||
const attempt = retry.attempt;
|
||||
const delay = XAI_REALTIME_BASE_RECONNECT_DELAY_MS * 2 ** (attempt - 1);
|
||||
this.config.onEvent?.({
|
||||
direction: "client",
|
||||
type: "session.reconnect.scheduled",
|
||||
detail: `reason=${reason} attempt=${attempt} delayMs=${delay}`,
|
||||
});
|
||||
const reconnectSignal = this.reconnectAbortController.signal;
|
||||
try {
|
||||
await sleepWithAbort(delay, reconnectSignal);
|
||||
await sleepWithAbort(delay, retry.signal);
|
||||
} catch (error) {
|
||||
if (!reconnectSignal.aborted) {
|
||||
if (!retry.signal.aborted) {
|
||||
throw error;
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (this.intentionallyClosed) {
|
||||
const nextConnection = this.lifecycle.reconnect(connection);
|
||||
if (!nextConnection) {
|
||||
return;
|
||||
}
|
||||
this.connection = nextConnection;
|
||||
try {
|
||||
await this.doConnect();
|
||||
await this.doConnect(nextConnection);
|
||||
this.config.onEvent?.({
|
||||
direction: "client",
|
||||
type: "session.reconnect.ready",
|
||||
detail: `reason=${reason} attempt=${attempt}`,
|
||||
});
|
||||
} catch (error) {
|
||||
if (this.terminalError) {
|
||||
if (
|
||||
this.terminalError ||
|
||||
!this.lifecycle.isCurrent(nextConnection) ||
|
||||
this.lifecycle.terminalOutcome(nextConnection)
|
||||
) {
|
||||
return;
|
||||
}
|
||||
this.config.onError?.(error instanceof Error ? error : new Error(String(error)));
|
||||
await this.attemptReconnect(reason);
|
||||
await this.attemptReconnect(reason, nextConnection);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -427,9 +490,14 @@ export class XaiRealtimeVoiceBridge extends XaiRealtimeVoiceEvents implements Re
|
||||
return undefined;
|
||||
}
|
||||
|
||||
protected onSessionUpdated(): void {
|
||||
this.sessionConfigured = true;
|
||||
this.reconnectAttempts = 0;
|
||||
protected acceptsEvent(connection: XaiRealtimeVoiceConnection): boolean {
|
||||
return this.lifecycle.acceptsEvents(connection);
|
||||
}
|
||||
|
||||
protected onSessionUpdated(connection: XaiRealtimeVoiceConnection): void {
|
||||
if (!this.lifecycle.ready(connection)) {
|
||||
return;
|
||||
}
|
||||
const pendingAudio = this.pendingAudio.splice(0);
|
||||
this.pendingAudioBytes = 0;
|
||||
for (const chunk of pendingAudio) {
|
||||
@@ -470,11 +538,11 @@ export class XaiRealtimeVoiceBridge extends XaiRealtimeVoiceEvents implements Re
|
||||
}
|
||||
|
||||
private canSubmitToolResult(): boolean {
|
||||
return this.connected && this.sessionConfigured && this.ws?.readyState === WebSocket.OPEN;
|
||||
return this.isConnected();
|
||||
}
|
||||
|
||||
private canSubmitInput(): boolean {
|
||||
return this.connected && this.sessionConfigured && this.ws?.readyState === WebSocket.OPEN;
|
||||
return this.isConnected();
|
||||
}
|
||||
|
||||
private enqueuePendingAudio(audio: Buffer): void {
|
||||
@@ -489,12 +557,45 @@ export class XaiRealtimeVoiceBridge extends XaiRealtimeVoiceEvents implements Re
|
||||
this.pendingAudioBytes += queuedAudio.byteLength;
|
||||
}
|
||||
|
||||
private enterTerminalState(): void {
|
||||
this.intentionallyClosed = true;
|
||||
this.reconnectAbortController.abort();
|
||||
this.connected = false;
|
||||
this.sessionConfigured = false;
|
||||
private failConnection(
|
||||
error: XaiRealtimeMalformedAudioError,
|
||||
ws: WebSocket,
|
||||
connection: XaiRealtimeVoiceConnection,
|
||||
): void {
|
||||
if (this.terminalError) {
|
||||
return;
|
||||
}
|
||||
this.terminalError = error;
|
||||
this.lifecycle.failure(connection);
|
||||
this.resetTerminalState();
|
||||
try {
|
||||
this.config.onError?.(error);
|
||||
} finally {
|
||||
if (ws.readyState !== WebSocket.CLOSED) {
|
||||
ws.close(1002, "Malformed audio payload");
|
||||
} else {
|
||||
this.notifyClose(connection, "error");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private enterTerminalState(connection: XaiRealtimeVoiceConnection): void {
|
||||
if (this.lifecycle.failure(connection)) {
|
||||
this.resetTerminalState();
|
||||
}
|
||||
this.notifyClose(connection, "error");
|
||||
}
|
||||
|
||||
private notifyClose(
|
||||
connection: XaiRealtimeVoiceConnection,
|
||||
outcome: "completed" | "error",
|
||||
): void {
|
||||
const terminalOutcome = this.lifecycle.close(connection, outcome);
|
||||
if (!terminalOutcome) {
|
||||
return;
|
||||
}
|
||||
this.resetTerminalState();
|
||||
this.config.onClose?.(terminalOutcome);
|
||||
}
|
||||
|
||||
private resetTerminalState(): void {
|
||||
|
||||
@@ -6,6 +6,7 @@ import {
|
||||
readXaiRealtimeErrorDetail,
|
||||
type XaiRealtimeEvent,
|
||||
} from "./realtime-voice-config.js";
|
||||
import type { XaiRealtimeVoiceConnection } from "./realtime-voice-lifecycle.js";
|
||||
import { XaiRealtimeVoiceProtocol } from "./realtime-voice-protocol.js";
|
||||
|
||||
export class XaiRealtimeMalformedAudioError extends Error {}
|
||||
@@ -15,9 +16,10 @@ export abstract class XaiRealtimeVoiceEvents extends XaiRealtimeVoiceProtocol {
|
||||
private assistantTranscriptFinalized = false;
|
||||
private inputTranscriptReplacements = new Map<string, string>();
|
||||
|
||||
protected abstract onSessionUpdated(): void;
|
||||
protected abstract acceptsEvent(connection: XaiRealtimeVoiceConnection): boolean;
|
||||
protected abstract onSessionUpdated(connection: XaiRealtimeVoiceConnection): void;
|
||||
|
||||
protected handleEvent(event: XaiRealtimeEvent): void {
|
||||
protected handleEvent(event: XaiRealtimeEvent, connection: XaiRealtimeVoiceConnection): void {
|
||||
this.config.onEvent?.({
|
||||
direction: "server",
|
||||
type: event.type,
|
||||
@@ -27,6 +29,9 @@ export abstract class XaiRealtimeVoiceEvents extends XaiRealtimeVoiceProtocol {
|
||||
? { responseId: event.response_id ?? event.response?.id }
|
||||
: {}),
|
||||
});
|
||||
if (!this.acceptsEvent(connection)) {
|
||||
return;
|
||||
}
|
||||
switch (event.type) {
|
||||
case "session.created":
|
||||
return;
|
||||
@@ -58,7 +63,7 @@ export abstract class XaiRealtimeVoiceEvents extends XaiRealtimeVoiceProtocol {
|
||||
return;
|
||||
}
|
||||
case "session.updated":
|
||||
this.onSessionUpdated();
|
||||
this.onSessionUpdated(connection);
|
||||
return;
|
||||
case "response.created":
|
||||
this.responseActive = true;
|
||||
|
||||
165
extensions/xai/realtime-voice-lifecycle.ts
Normal file
165
extensions/xai/realtime-voice-lifecycle.ts
Normal file
@@ -0,0 +1,165 @@
|
||||
type XaiRealtimeVoiceLifecyclePhase = "idle" | "connecting" | "ready" | "retry-wait" | "terminal";
|
||||
|
||||
type XaiRealtimeVoiceTerminalOutcome = "completed" | "error";
|
||||
|
||||
export type XaiRealtimeVoiceConnection = Readonly<{
|
||||
id: symbol;
|
||||
signal: AbortSignal;
|
||||
}>;
|
||||
|
||||
type XaiRealtimeVoiceIdleState = {
|
||||
phase: "idle" | "terminal";
|
||||
terminalOutcome?: "completed";
|
||||
};
|
||||
|
||||
type XaiRealtimeVoiceConnectionState = {
|
||||
connection: XaiRealtimeVoiceConnection;
|
||||
controller: AbortController;
|
||||
phase: Exclude<XaiRealtimeVoiceLifecyclePhase, "idle">;
|
||||
retryAttempts: number;
|
||||
terminalOutcome?: XaiRealtimeVoiceTerminalOutcome;
|
||||
terminalNotified: boolean;
|
||||
};
|
||||
|
||||
export class XaiRealtimeVoiceLifecycle {
|
||||
private state: XaiRealtimeVoiceIdleState | XaiRealtimeVoiceConnectionState = {
|
||||
phase: "idle",
|
||||
};
|
||||
|
||||
connect(): XaiRealtimeVoiceConnection {
|
||||
if ("controller" in this.state) {
|
||||
this.state.controller.abort(new Error("xAI realtime voice connection replaced"));
|
||||
}
|
||||
const controller = new AbortController();
|
||||
const connection = this.createConnection(controller);
|
||||
this.state = {
|
||||
connection,
|
||||
controller,
|
||||
phase: "connecting",
|
||||
retryAttempts: 0,
|
||||
terminalNotified: false,
|
||||
};
|
||||
return connection;
|
||||
}
|
||||
|
||||
reconnect(connection: XaiRealtimeVoiceConnection): XaiRealtimeVoiceConnection | undefined {
|
||||
const state = this.currentState(connection);
|
||||
if (!state || state.phase !== "retry-wait" || state.terminalOutcome) {
|
||||
return undefined;
|
||||
}
|
||||
const nextConnection = this.createConnection(state.controller);
|
||||
state.connection = nextConnection;
|
||||
state.phase = "connecting";
|
||||
return nextConnection;
|
||||
}
|
||||
|
||||
ready(connection: XaiRealtimeVoiceConnection): boolean {
|
||||
const state = this.currentState(connection);
|
||||
if (!state || state.phase !== "connecting" || state.terminalOutcome) {
|
||||
return false;
|
||||
}
|
||||
state.phase = "ready";
|
||||
state.retryAttempts = 0;
|
||||
return true;
|
||||
}
|
||||
|
||||
retry(
|
||||
connection: XaiRealtimeVoiceConnection,
|
||||
maxAttempts: number,
|
||||
): { attempt: number; signal: AbortSignal } | "exhausted" | undefined {
|
||||
const state = this.currentState(connection);
|
||||
if (!state || state.phase === "retry-wait" || state.terminalOutcome) {
|
||||
return undefined;
|
||||
}
|
||||
if (state.retryAttempts >= maxAttempts) {
|
||||
return "exhausted";
|
||||
}
|
||||
state.retryAttempts += 1;
|
||||
state.phase = "retry-wait";
|
||||
return { attempt: state.retryAttempts, signal: state.controller.signal };
|
||||
}
|
||||
|
||||
cancel(): boolean {
|
||||
const state = this.state;
|
||||
if (state.phase === "terminal") {
|
||||
return false;
|
||||
}
|
||||
if (!("controller" in state)) {
|
||||
this.state = {
|
||||
phase: "terminal",
|
||||
terminalOutcome: "completed",
|
||||
};
|
||||
return true;
|
||||
}
|
||||
state.phase = "terminal";
|
||||
state.terminalOutcome = "completed";
|
||||
state.controller.abort(new Error("xAI realtime voice session canceled"));
|
||||
return true;
|
||||
}
|
||||
|
||||
failure(connection: XaiRealtimeVoiceConnection): boolean {
|
||||
const state = this.currentState(connection);
|
||||
if (!state || state.terminalOutcome) {
|
||||
return false;
|
||||
}
|
||||
state.phase = "terminal";
|
||||
state.terminalOutcome = "error";
|
||||
state.controller.abort(new Error("xAI realtime voice session failed"));
|
||||
return true;
|
||||
}
|
||||
|
||||
close(
|
||||
connection: XaiRealtimeVoiceConnection,
|
||||
outcome: XaiRealtimeVoiceTerminalOutcome,
|
||||
): XaiRealtimeVoiceTerminalOutcome | undefined {
|
||||
const state = this.currentState(connection);
|
||||
if (!state) {
|
||||
return undefined;
|
||||
}
|
||||
if (!state.terminalOutcome) {
|
||||
state.phase = "terminal";
|
||||
state.terminalOutcome = outcome;
|
||||
state.controller.abort(new Error("xAI realtime voice session closed"));
|
||||
}
|
||||
if (state.terminalNotified) {
|
||||
return undefined;
|
||||
}
|
||||
state.terminalNotified = true;
|
||||
return state.terminalOutcome;
|
||||
}
|
||||
|
||||
isCurrent(connection: XaiRealtimeVoiceConnection): boolean {
|
||||
return this.currentState(connection) !== undefined;
|
||||
}
|
||||
|
||||
acceptsEvents(connection: XaiRealtimeVoiceConnection): boolean {
|
||||
const phase = this.currentState(connection)?.phase;
|
||||
return phase === "connecting" || phase === "ready";
|
||||
}
|
||||
|
||||
isReady(): boolean {
|
||||
return this.state.phase === "ready";
|
||||
}
|
||||
|
||||
phase(): XaiRealtimeVoiceLifecyclePhase {
|
||||
return this.state.phase;
|
||||
}
|
||||
|
||||
terminalOutcome(
|
||||
connection: XaiRealtimeVoiceConnection,
|
||||
): XaiRealtimeVoiceTerminalOutcome | undefined {
|
||||
return this.currentState(connection)?.terminalOutcome;
|
||||
}
|
||||
|
||||
private createConnection(controller: AbortController): XaiRealtimeVoiceConnection {
|
||||
return { id: Symbol("xai-realtime-voice-connection"), signal: controller.signal };
|
||||
}
|
||||
|
||||
private currentState(
|
||||
connection: XaiRealtimeVoiceConnection,
|
||||
): XaiRealtimeVoiceConnectionState | undefined {
|
||||
return "connection" in this.state && this.state.connection.id === connection.id
|
||||
? this.state
|
||||
: undefined;
|
||||
}
|
||||
}
|
||||
@@ -15,6 +15,7 @@ import {
|
||||
type XaiRealtimeSessionUpdate,
|
||||
type XaiRealtimeVoiceBridgeConfig,
|
||||
} from "./realtime-voice-config.js";
|
||||
import type { XaiRealtimeVoiceConnection } from "./realtime-voice-lifecycle.js";
|
||||
|
||||
export abstract class XaiRealtimeVoiceProtocol {
|
||||
protected readonly audioFormat: RealtimeVoiceAudioFormat;
|
||||
@@ -287,5 +288,8 @@ export abstract class XaiRealtimeVoiceProtocol {
|
||||
}
|
||||
|
||||
protected abstract resetInputTranscripts(): void;
|
||||
protected abstract handleEvent(event: XaiRealtimeEvent): void;
|
||||
protected abstract handleEvent(
|
||||
event: XaiRealtimeEvent,
|
||||
connection: XaiRealtimeVoiceConnection,
|
||||
): void;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user