mirror of
https://github.com/openclaw/openclaw.git
synced 2026-07-21 10:16:44 +00:00
fix(realtime): add timeout to WebRTC offer requests (#104662)
* fix(realtime): add timeout to WebRTC offer requests * test(realtime): reject offer aborts with errors * fix(ui): abort pending realtime offer on stop * test(ui): normalize realtime abort rejection --------- Co-authored-by: Peter Steinberger <steipete@gmail.com>
This commit is contained in:
co-authored by
Peter Steinberger
parent
26c4187297
commit
41a1fd696b
@@ -174,6 +174,7 @@ function expectSpokenStatusMessage(events: SentRealtimeEvent[], message: string)
|
||||
describe("WebRtcSdpRealtimeTalkTransport", () => {
|
||||
afterEach(() => {
|
||||
vi.unstubAllGlobals();
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
beforeEach(() => {
|
||||
@@ -325,10 +326,102 @@ describe("WebRtcSdpRealtimeTalkTransport", () => {
|
||||
Authorization: "Bearer client-secret-123",
|
||||
"Content-Type": "application/sdp",
|
||||
},
|
||||
signal: expect.any(AbortSignal),
|
||||
});
|
||||
transport.stop();
|
||||
});
|
||||
|
||||
it("aborts stalled WebRTC SDP answer body reads after the offer timeout", async () => {
|
||||
vi.useFakeTimers();
|
||||
let offerSignal: AbortSignal | undefined;
|
||||
const fetchMock = vi.fn(async (_url: string | URL | Request, init?: RequestInit) => {
|
||||
offerSignal = init?.signal ?? undefined;
|
||||
return {
|
||||
ok: true,
|
||||
status: 200,
|
||||
text: () =>
|
||||
new Promise<string>((_, reject) => {
|
||||
offerSignal?.addEventListener(
|
||||
"abort",
|
||||
() => {
|
||||
const reason = offerSignal?.reason;
|
||||
reject(reason instanceof Error ? reason : new Error("offer request aborted"));
|
||||
},
|
||||
{ once: true },
|
||||
);
|
||||
}),
|
||||
} as Response;
|
||||
});
|
||||
vi.stubGlobal("fetch", fetchMock as unknown as typeof fetch);
|
||||
const transport = createOpenAiTransport();
|
||||
|
||||
const startResult = transport.start().then(
|
||||
() => undefined,
|
||||
(error: unknown) => error,
|
||||
);
|
||||
|
||||
await vi.waitFor(() => expect(fetchMock).toHaveBeenCalledTimes(1));
|
||||
expect(offerSignal?.aborted).toBe(false);
|
||||
|
||||
await vi.runAllTimersAsync();
|
||||
|
||||
await expect(startResult).resolves.toMatchObject(
|
||||
new Error("Realtime WebRTC offer request timed out after 30000ms"),
|
||||
);
|
||||
expect(offerSignal?.aborted).toBe(true);
|
||||
});
|
||||
|
||||
it("aborts a pending WebRTC SDP answer body read when stopped", async () => {
|
||||
let offerSignal: AbortSignal | undefined;
|
||||
const fetchMock = vi.fn(async (_url: string | URL | Request, init?: RequestInit) => {
|
||||
offerSignal = init?.signal ?? undefined;
|
||||
return {
|
||||
ok: true,
|
||||
status: 200,
|
||||
text: () =>
|
||||
new Promise<string>((_, reject) => {
|
||||
offerSignal?.addEventListener(
|
||||
"abort",
|
||||
() => {
|
||||
const reason = offerSignal?.reason;
|
||||
reject(reason instanceof Error ? reason : new Error("offer request aborted"));
|
||||
},
|
||||
{ once: true },
|
||||
);
|
||||
}),
|
||||
} as Response;
|
||||
});
|
||||
vi.stubGlobal("fetch", fetchMock as unknown as typeof fetch);
|
||||
const transport = createOpenAiTransport();
|
||||
|
||||
const startResult = transport.start();
|
||||
await vi.waitFor(() => expect(fetchMock).toHaveBeenCalledTimes(1));
|
||||
|
||||
transport.stop();
|
||||
|
||||
await expect(startResult).resolves.toBeUndefined();
|
||||
expect(offerSignal?.aborted).toBe(true);
|
||||
});
|
||||
|
||||
it("clears the WebRTC offer timeout after setup succeeds", async () => {
|
||||
vi.useFakeTimers();
|
||||
let offerSignal: AbortSignal | undefined;
|
||||
vi.stubGlobal(
|
||||
"fetch",
|
||||
vi.fn(async (_url: string | URL | Request, init?: RequestInit) => {
|
||||
offerSignal = init?.signal ?? undefined;
|
||||
return new Response("answer-sdp");
|
||||
}) as unknown as typeof fetch,
|
||||
);
|
||||
const transport = createOpenAiTransport();
|
||||
|
||||
await transport.start();
|
||||
await vi.runAllTimersAsync();
|
||||
|
||||
expect(offerSignal?.aborted).toBe(false);
|
||||
transport.stop();
|
||||
});
|
||||
|
||||
it("surfaces realtime provider errors from the OpenAI data channel", async () => {
|
||||
vi.stubGlobal(
|
||||
"fetch",
|
||||
|
||||
@@ -37,6 +37,12 @@ type ToolBuffer = {
|
||||
};
|
||||
|
||||
const cancelledSetup = Symbol("cancelledSetup");
|
||||
const REALTIME_WEBRTC_OFFER_TIMEOUT_MS = 30_000;
|
||||
|
||||
type PendingOfferRequest = {
|
||||
controller: AbortController;
|
||||
timeout: ReturnType<typeof globalThis.setTimeout>;
|
||||
};
|
||||
|
||||
export class WebRtcSdpRealtimeTalkTransport implements RealtimeTalkTransport {
|
||||
private peer: RTCPeerConnection | null = null;
|
||||
@@ -49,6 +55,7 @@ export class WebRtcSdpRealtimeTalkTransport implements RealtimeTalkTransport {
|
||||
private responseCreateInFlight = false;
|
||||
private responseCreatePending = false;
|
||||
private toolBuffers = new Map<string, ToolBuffer>();
|
||||
private pendingOfferRequest: PendingOfferRequest | null = null;
|
||||
private readonly consultAbortControllers = new Set<AbortController>();
|
||||
private readonly emitTalkEvent: ReturnType<typeof createRealtimeTalkEventEmitter>;
|
||||
|
||||
@@ -126,28 +133,7 @@ export class WebRtcSdpRealtimeTalkTransport implements RealtimeTalkTransport {
|
||||
if (!this.isCurrentPeer(peer)) {
|
||||
return;
|
||||
}
|
||||
const sdp = await this.awaitSetupStep(
|
||||
peer,
|
||||
fetch(this.session.offerUrl ?? "https://api.openai.com/v1/realtime/calls", {
|
||||
method: "POST",
|
||||
body: offer.sdp,
|
||||
headers: {
|
||||
...this.session.offerHeaders,
|
||||
Authorization: `Bearer ${this.session.clientSecret}`,
|
||||
"Content-Type": "application/sdp",
|
||||
},
|
||||
}),
|
||||
);
|
||||
if (sdp === cancelledSetup) {
|
||||
return;
|
||||
}
|
||||
if (!this.isCurrentPeer(peer)) {
|
||||
return;
|
||||
}
|
||||
if (!sdp.ok) {
|
||||
throw new Error(`Realtime WebRTC setup failed (${sdp.status})`);
|
||||
}
|
||||
const answerSdp = await this.awaitSetupStep(peer, sdp.text());
|
||||
const answerSdp = await this.readOfferAnswer(peer, offer);
|
||||
if (answerSdp === cancelledSetup) {
|
||||
return;
|
||||
}
|
||||
@@ -163,6 +149,83 @@ export class WebRtcSdpRealtimeTalkTransport implements RealtimeTalkTransport {
|
||||
);
|
||||
}
|
||||
|
||||
private async readOfferAnswer(
|
||||
peer: RTCPeerConnection,
|
||||
offer: RTCSessionDescriptionInit,
|
||||
): Promise<string | typeof cancelledSetup> {
|
||||
const request = this.beginOfferRequest();
|
||||
try {
|
||||
const sdp = await this.awaitSetupStep(
|
||||
peer,
|
||||
fetch(this.session.offerUrl ?? "https://api.openai.com/v1/realtime/calls", {
|
||||
method: "POST",
|
||||
body: offer.sdp,
|
||||
headers: {
|
||||
...this.session.offerHeaders,
|
||||
Authorization: `Bearer ${this.session.clientSecret}`,
|
||||
"Content-Type": "application/sdp",
|
||||
},
|
||||
signal: request.controller.signal,
|
||||
}),
|
||||
);
|
||||
if (sdp === cancelledSetup) {
|
||||
return cancelledSetup;
|
||||
}
|
||||
if (!this.isCurrentPeer(peer)) {
|
||||
return cancelledSetup;
|
||||
}
|
||||
if (!sdp.ok) {
|
||||
throw new Error(`Realtime WebRTC setup failed (${sdp.status})`);
|
||||
}
|
||||
const answerSdp = await this.awaitSetupStep(peer, sdp.text());
|
||||
if (answerSdp === cancelledSetup) {
|
||||
return cancelledSetup;
|
||||
}
|
||||
if (!this.isCurrentPeer(peer)) {
|
||||
return cancelledSetup;
|
||||
}
|
||||
return answerSdp;
|
||||
} finally {
|
||||
this.finishOfferRequest(request);
|
||||
}
|
||||
}
|
||||
|
||||
private beginOfferRequest(): PendingOfferRequest {
|
||||
this.abortOfferRequest();
|
||||
const controller = new AbortController();
|
||||
const request = {
|
||||
controller,
|
||||
timeout: globalThis.setTimeout(() => {
|
||||
controller.abort(
|
||||
new Error(
|
||||
`Realtime WebRTC offer request timed out after ${REALTIME_WEBRTC_OFFER_TIMEOUT_MS}ms`,
|
||||
),
|
||||
);
|
||||
}, REALTIME_WEBRTC_OFFER_TIMEOUT_MS),
|
||||
};
|
||||
this.pendingOfferRequest = request;
|
||||
return request;
|
||||
}
|
||||
|
||||
private finishOfferRequest(request: PendingOfferRequest): void {
|
||||
globalThis.clearTimeout(request.timeout);
|
||||
// A stopped transport may already have started a replacement request.
|
||||
// Never let the old request's finally block detach the new lifecycle owner.
|
||||
if (this.pendingOfferRequest === request) {
|
||||
this.pendingOfferRequest = null;
|
||||
}
|
||||
}
|
||||
|
||||
private abortOfferRequest(): void {
|
||||
const request = this.pendingOfferRequest;
|
||||
if (!request) {
|
||||
return;
|
||||
}
|
||||
this.pendingOfferRequest = null;
|
||||
globalThis.clearTimeout(request.timeout);
|
||||
request.controller.abort();
|
||||
}
|
||||
|
||||
private isCurrentPeer(peer: RTCPeerConnection): boolean {
|
||||
return !this.closed && this.peer === peer;
|
||||
}
|
||||
@@ -186,6 +249,7 @@ export class WebRtcSdpRealtimeTalkTransport implements RealtimeTalkTransport {
|
||||
this.emitTalkEvent({ type: "session.closed", final: true });
|
||||
}
|
||||
this.closed = true;
|
||||
this.abortOfferRequest();
|
||||
this.channel?.close();
|
||||
this.channel = null;
|
||||
this.peer?.close();
|
||||
|
||||
Reference in New Issue
Block a user