Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion apps/desktop/electron/main/live-voice/audio-port.ts
Original file line number Diff line number Diff line change
Expand Up @@ -123,7 +123,7 @@ export function createLivePcmBridge(options: LivePcmBridgeOptions): LivePcmBridg
expectedInputSequence += 1;
const accepted = !muted && epoch === captureEpoch;
if (accepted) options.onInput(new Uint8Array(data), captureEpoch);
post({ kind: "credit", direction: "uplink", consumedSequence: sequence, accepted });
post({ kind: "credit", callId: options.callId, direction: "uplink", consumedSequence: sequence, accepted });
return;
}
case "gate-applied": {
Expand Down
26 changes: 17 additions & 9 deletions apps/desktop/electron/main/live-voice/codex-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,12 +46,17 @@ export function createCodexAdapter(context: LiveAdapterContext, deps: {
}
if (closed || context.signal.aborted) throw Object.assign(new Error("Live call was cancelled"), { errorCode: "LIVE_STALE_CALL" });
const url = `${CODEX_LIVE_BASE}/realtime/calls?intent=quicksilver&architecture=avas`;
await deps.assertEndpoint(url);
requestController = new AbortController();
const onAbort = () => requestController?.abort(context.signal.reason);
const controller = new AbortController();
requestController = controller;
const onAbort = () => controller.abort(context.signal.reason);
context.signal.addEventListener("abort", onAbort, { once: true });
const timeout = setTimeout(() => requestController?.abort(new Error("timeout")), 15_000);
let timeout: ReturnType<typeof setTimeout> | undefined;
try {
await deps.assertEndpoint(url);
if (closed || context.signal.aborted) {
throw Object.assign(new Error("Live call was cancelled"), { errorCode: "LIVE_STALE_CALL" });
}
timeout = setTimeout(() => controller.abort(new Error("timeout")), 15_000);
const headers = {
Authorization: `Bearer ${auth.accessToken}`,
"chatgpt-account-id": auth.accountId,
Expand All @@ -65,7 +70,7 @@ export function createCodexAdapter(context: LiveAdapterContext, deps: {
method: "POST",
headers,
body: JSON.stringify(buildCodexCallBody({ sdp: offerSdp, voice: binding.voice, ...(context.workProfile ? { workProfile: context.workProfile } : {}) })),
signal: requestController.signal,
signal: controller.signal,
redirect: "manual",
});
if (response.status !== 201) {
Expand All @@ -85,20 +90,23 @@ export function createCodexAdapter(context: LiveAdapterContext, deps: {
if (closed || context.signal.aborted) throw Object.assign(new Error("Live call was cancelled"), { errorCode: "LIVE_STALE_CALL" });
return { answerSdp };
} catch (error) {
if (closed || context.signal.aborted) {
throw Object.assign(new Error("Live call was cancelled"), { errorCode: "LIVE_STALE_CALL" });
}
if (error && typeof error === "object" && "errorCode" in error) throw error;
if (requestController.signal.aborted) {
if (controller.signal.aborted) {
throw Object.assign(new Error("Codex Live call creation timed out or was cancelled"), {
errorCode: context.signal.aborted ? "LIVE_STALE_CALL" : "LIVE_TIMEOUT",
errorCode: "LIVE_TIMEOUT",
});
}
throw Object.assign(new Error("Codex Live connection failed"), {
errorCode: ErrorCodes.LIVE_NETWORK_ERROR,
cause: error,
});
} finally {
clearTimeout(timeout);
if (timeout) clearTimeout(timeout);
context.signal.removeEventListener("abort", onAbort);
requestController = null;
if (requestController === controller) requestController = null;
}
},
async close() {
Expand Down
35 changes: 25 additions & 10 deletions apps/desktop/electron/main/live-voice/gemini-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ import { geminiAudioEndMessage, geminiAudioMessage, geminiSetupMessage, geminiTo
import type { LiveAdapter, LiveAdapterContext, LiveReceiptDelivery } from "./types";
import type { LiveWorkFeedback } from "@pi-desktop/shared";
import { openLiveWebSocket } from "./websocket-transport";
import { sendJsonBounded, waitForReady, waitForSocketReady, websocketJson } from "./websocket-wire";
import { sendJsonBounded, sendJsonConfirmed, waitForReady, waitForSocketReady, websocketJson } from "./websocket-wire";
import { WebSocket } from "ws";

const GEMINI_LIVE_ENDPOINT = "wss://generativelanguage.googleapis.com/ws/google.ai.generativelanguage.v1beta.GenerativeService.BidiGenerateContent";
Expand All @@ -19,6 +19,7 @@ export function createGeminiAdapter(context: LiveAdapterContext): LiveAdapter {
let socket: Awaited<ReturnType<typeof openLiveWebSocket>> | null = null;
let closed = false;
let muted = context.signal.aborted;
let setupSent = false;
let readyResolve: (() => void) | null = null;
let readyReject: ((error: Error) => void) | null = null;
const completedFunctionCalls = new Set<string>();
Expand Down Expand Up @@ -60,7 +61,8 @@ export function createGeminiAdapter(context: LiveAdapterContext): LiveAdapter {
async (receipt) => {
const deliveryId = randomUUID();
try {
sendFunctionReceipt({ ...candidate, receipt });
if (!socket) throw new Error("Gemini Live socket is not open");
await sendJsonConfirmed(socket, geminiToolResponseMessage({ ...candidate, receipt }), context.signal);
return { status: "sent", deliveryId };
} catch {
return { status: "not-sent", deliveryId, code: "LIVE_WORK_FEEDBACK_UNDELIVERED" };
Expand All @@ -69,7 +71,12 @@ export function createGeminiAdapter(context: LiveAdapterContext): LiveAdapter {
).catch(() => context.onEvent({ kind: "error", code: "LIVE_WORK_INTENT_UNAVAILABLE" }));
}

function onSocketFailure(): void {
if (!closed && !context.signal.aborted) context.onEvent({ kind: "error", code: "LIVE_NETWORK_ERROR" });
}

function onSocketMessage(data: unknown): void {
if (closed || context.signal.aborted) return;
try {
const value = websocketJson(data, MAX_LIVE_JSON_BYTES);
const root = value as Record<string, unknown>;
Expand Down Expand Up @@ -107,16 +114,21 @@ export function createGeminiAdapter(context: LiveAdapterContext): LiveAdapter {
const url = new URL(GEMINI_LIVE_ENDPOINT);
url.searchParams.set("key", apiKey);
const liveSocket = await openLiveWebSocket({ url: url.toString(), signal: context.signal, endpointOrigin: "third-party" });
if (closed || context.signal.aborted) {
liveSocket.terminate();
throw Object.assign(new Error("Live call was cancelled"), { errorCode: "LIVE_STALE_CALL" });
}
socket = liveSocket;
liveSocket.on("message", onSocketMessage);
liveSocket.once("close", (code) => {
if (!closed) context.onEvent({ kind: "error", code: code === 1000 ? "LIVE_NETWORK_ERROR" : "LIVE_NETWORK_ERROR" });
});
liveSocket.once("error", () => {
if (!closed) context.onEvent({ kind: "error", code: "LIVE_NETWORK_ERROR" });
});
liveSocket.once("close", onSocketFailure);
liveSocket.once("error", onSocketFailure);
await waitForSocketReady(liveSocket, context.signal);
if (closed || context.signal.aborted) {
liveSocket.terminate();
throw Object.assign(new Error("Live call was cancelled"), { errorCode: "LIVE_STALE_CALL" });
}
sendJsonBounded(liveSocket, geminiSetupMessage({ modelId: binding.modelId, voice: binding.voice, ...(context.workProfile ? { workProfile: context.workProfile } : {}) }));
setupSent = true;
await waitForReady(ready, context.signal, 12_000, "Gemini Live setup timed out");
return {};
},
Expand All @@ -137,7 +149,7 @@ export function createGeminiAdapter(context: LiveAdapterContext): LiveAdapter {
return { status: "not-sent", deliveryId, code: "LIVE_WORK_FEEDBACK_UNDELIVERED" };
}
try {
sendJsonBounded(socket, geminiWorkFeedbackMessage(feedback));
await sendJsonConfirmed(socket, geminiWorkFeedbackMessage(feedback), context.signal);
return { status: "sent", deliveryId };
} catch {
return { status: "not-sent", deliveryId, code: "LIVE_WORK_FEEDBACK_UNDELIVERED" };
Expand All @@ -146,7 +158,10 @@ export function createGeminiAdapter(context: LiveAdapterContext): LiveAdapter {
async close() {
if (closed) return;
closed = true;
readyReject?.(new Error("Gemini Live call closed"));
if (setupSent) readyReject?.(new Error("Gemini Live call closed"));
socket?.off("message", onSocketMessage);
socket?.off("close", onSocketFailure);
socket?.off("error", onSocketFailure);
if (socket && socket.readyState !== WebSocket.CLOSED) socket.close(1000, "call ended");
completedFunctionCalls.clear();
socket = null;
Expand Down
19 changes: 14 additions & 5 deletions apps/desktop/electron/main/live-voice/openai-realtime-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ import { LIVE_WORK_TOOL_NAME, MAX_LIVE_AUDIO_BYTES, MAX_LIVE_JSON_BYTES, parseLi
import type { LiveAdapter, LiveAdapterContext, LivePlaybackCursor, LiveReceiptDelivery } from "./types";
import type { LiveWorkFeedback } from "@pi-desktop/shared";
import { openLiveWebSocket } from "./websocket-transport";
import { sendJsonBounded, waitForReady, waitForSocketReady, websocketJson } from "./websocket-wire";
import { sendJsonBounded, sendJsonConfirmed, waitForReady, waitForSocketReady, websocketJson } from "./websocket-wire";
import { WebSocket } from "ws";

const MAX_FRAME_BYTES = 24000 * 2 / 10;
Expand Down Expand Up @@ -88,7 +88,8 @@ export function createOpenAIRealtimeAdapter(context: LiveAdapterContext): LiveAd
async (receipt) => {
const deliveryId = randomUUID();
try {
sendToolReceipt({ providerRequestId: candidate.providerRequestId, receipt, resume: false });
if (!socket) throw new Error("Realtime socket is not open");
await sendJsonConfirmed(socket, realtimeToolReceiptMessage({ providerRequestId: candidate.providerRequestId, receipt }), context.signal);
return { status: "sent", deliveryId };
} catch {
return { status: "not-sent", deliveryId, code: "LIVE_WORK_FEEDBACK_UNDELIVERED" };
Expand All @@ -97,7 +98,12 @@ export function createOpenAIRealtimeAdapter(context: LiveAdapterContext): LiveAd
).catch(() => context.onEvent({ kind: "error", code: "LIVE_WORK_INTENT_UNAVAILABLE" }));
}

function onSocketFailure(): void {
if (!closed) context.onEvent({ kind: "error", code: "LIVE_NETWORK_ERROR" });
}

function onSocketMessage(data: unknown): void {
if (closed || context.signal.aborted) return;
try {
const value = websocketJson(data, MAX_LIVE_JSON_BYTES);
const event = value as Record<string, unknown>;
Expand Down Expand Up @@ -170,8 +176,8 @@ export function createOpenAIRealtimeAdapter(context: LiveAdapterContext): LiveAd
const liveSocket = await openLiveWebSocket({ url, headers: { Authorization: `Bearer ${apiKey}`, ...(binding.wireProfile === "realtime-compat-v1" ? { "OpenAI-Beta": "realtime=v1" } : {}) }, signal: context.signal, endpointOrigin: "user" });
socket = liveSocket;
liveSocket.on("message", onSocketMessage);
liveSocket.once("close", () => { if (!closed) context.onEvent({ kind: "error", code: "LIVE_NETWORK_ERROR" }); });
liveSocket.once("error", () => { if (!closed) context.onEvent({ kind: "error", code: "LIVE_NETWORK_ERROR" }); });
liveSocket.once("close", onSocketFailure);
liveSocket.once("error", onSocketFailure);
await waitForSocketReady(liveSocket, context.signal);
await waitForReady(created, context.signal, 12_000, "Realtime session creation timed out");
await waitForReady(updated, context.signal, 12_000, "Realtime session update timed out");
Expand Down Expand Up @@ -202,7 +208,7 @@ export function createOpenAIRealtimeAdapter(context: LiveAdapterContext): LiveAd
return { status: "not-sent", deliveryId, code: "LIVE_WORK_FEEDBACK_UNDELIVERED" };
}
try {
for (const message of realtimeWorkFeedbackMessages(feedback)) sendJsonBounded(socket, message);
for (const message of realtimeWorkFeedbackMessages(feedback)) await sendJsonConfirmed(socket, message, context.signal);
return { status: "sent", deliveryId };
} catch {
return { status: "not-sent", deliveryId, code: "LIVE_WORK_FEEDBACK_UNDELIVERED" };
Expand All @@ -215,6 +221,9 @@ export function createOpenAIRealtimeAdapter(context: LiveAdapterContext): LiveAd
updatedReject?.(new Error("Realtime call closed"));
responseTracker.clear();
settledFunctionCalls.clear();
socket?.off("message", onSocketMessage);
socket?.off("close", onSocketFailure);
socket?.off("error", onSocketFailure);
if (socket && socket.readyState !== WebSocket.CLOSED) socket.close(1000, "call ended");
socket = null;
},
Expand Down
2 changes: 1 addition & 1 deletion apps/desktop/electron/main/live-voice/owner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,7 @@ export function isTrustedRendererUrl(value: string): boolean {
const url = new URL(value);
if (url.protocol === "file:") {
if (url.search || url.hash) return false;
return resolve(fileURLToPath(url)) === resolve(__dirname, "../../renderer/index.html");
return resolve(fileURLToPath(url)) === resolve(__dirname, "../renderer/index.html");
}
if (url.protocol !== "http:") return false;
const configured = process.env.ELECTRON_RENDERER_URL;
Expand Down
21 changes: 15 additions & 6 deletions apps/desktop/electron/main/live-voice/websocket-transport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -88,10 +88,14 @@ export async function openLiveWebSocket(input: {
await endpointGuard.assertPublicUrl(input.url.replace(/^wss:/, "https:"), input.endpointOrigin);
if (input.signal.aborted) throw input.signal.reason;
const route = await session.defaultSession.resolveProxy(input.url);
if (input.signal.aborted) {
throw Object.assign(new Error("Live provider connection was cancelled"), { errorCode: "LIVE_STALE_CALL" });
}
const proxy = parseResolvedProxy(route);
const agent = proxy ? new ProxyTunnelAgent(proxy) : undefined;

try {
let cleanupAbort = () => {};
const connected = await new Promise<WebSocket>((resolve, reject) => {
const socket = new WebSocket(input.url, {
...(agent ? { agent } : {}),
Expand All @@ -103,18 +107,18 @@ export async function openLiveWebSocket(input: {
});
let settled = false;
const onAbort = () => {
const wasSettled = settled;
settled = true;
if (!wasSettled) reject(Object.assign(new Error("Live provider connection was cancelled"), { errorCode: "LIVE_STALE_CALL" }));
cleanup();
socket.terminate();
if (!settled) {
settled = true;
reject(input.signal.reason);
}
};
const cleanup = () => input.signal.removeEventListener("abort", onAbort);
cleanupAbort = cleanup;
input.signal.addEventListener("abort", onAbort, { once: true });
socket.once("open", () => {
if (settled) return;
settled = true;
cleanup();
resolve(socket);
});
socket.once("error", (error) => {
Expand Down Expand Up @@ -143,9 +147,14 @@ export async function openLiveWebSocket(input: {
}));
});
});
connected.once("close", () => agent?.destroy());
connected.once("close", () => {
cleanupAbort();
agent?.destroy();
});
return connected;
} catch (error) {
// Handshake failures normally clean up in their event handlers. This also
// covers a constructor failure before those handlers can run.
agent?.destroy();
throw error;
}
Expand Down
46 changes: 44 additions & 2 deletions apps/desktop/electron/main/live-voice/websocket-wire.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,18 +15,60 @@ export function websocketJson(data: unknown, maxBytes = MAX_LIVE_JSON_BYTES): un
return JSON.parse(text) as unknown;
}

export function sendJsonBounded(socket: WebSocket, value: unknown): void {
function boundedJsonBody(socket: WebSocket, value: unknown): string {
const body = JSON.stringify(value);
const byteLength = Buffer.byteLength(body, "utf8");
if (byteLength > MAX_LIVE_JSON_BYTES || socket.bufferedAmount + byteLength > MAX_BUFFERED_JSON_BYTES) {
throw Object.assign(new Error("Live audio transport queue is full"), { errorCode: "LIVE_AUDIO_BACKPRESSURE" });
}
socket.send(body, (error) => {
return body;
}

export function sendJsonBounded(socket: WebSocket, value: unknown): void {
socket.send(boundedJsonBody(socket, value), (error) => {
if (error) socket.terminate();
});
}

/** A receipt is sent only after the local socket write succeeds, not when it is queued. */
export async function sendJsonConfirmed(socket: WebSocket, value: unknown, signal: AbortSignal): Promise<void> {
const body = boundedJsonBody(socket, value);
if (socket.readyState !== WebSocket.OPEN || signal.aborted) {
throw Object.assign(new Error("Live provider message was not sent"), { errorCode: "LIVE_WORK_FEEDBACK_UNDELIVERED" });
}
await new Promise<void>((resolve, reject) => {
let settled = false;
const failure = () => Object.assign(new Error("Live provider message was not sent"), { errorCode: "LIVE_WORK_FEEDBACK_UNDELIVERED" });
const onFailure = () => finish(failure());
const timer = setTimeout(onFailure, 1_000);
function finish(error?: Error): void {
if (settled) return;
settled = true;
clearTimeout(timer);
socket.off("close", onFailure);
socket.off("error", onFailure);
signal.removeEventListener("abort", onFailure);
if (error) reject(error);
else resolve();
}
socket.once("close", onFailure);
socket.once("error", onFailure);
signal.addEventListener("abort", onFailure, { once: true });
try {
socket.send(body, (error) => {
finish(error ? failure() : undefined);
if (error) socket.terminate();
});
} catch {
finish(failure());
}
});
}

export async function waitForSocketReady(socket: WebSocket, signal: AbortSignal): Promise<void> {
if (signal.aborted) {
throw Object.assign(new Error("Live provider connection was cancelled"), { errorCode: "LIVE_STALE_CALL" });
}
if (socket.readyState === WebSocket.OPEN) return;
await new Promise<void>((resolve, reject) => {
const timeout = setTimeout(() => finish(Object.assign(new Error("Live provider socket timed out"), { errorCode: "LIVE_TIMEOUT" })), SOCKET_READY_TIMEOUT_MS);
Expand Down
2 changes: 1 addition & 1 deletion apps/desktop/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -158,7 +158,7 @@
"entitlementsInherit": "build/entitlements.mac.plist",
"extendInfo": {
"NSLocalNetworkUsageDescription": "PI-Desktop needs local network access so agents can reach OpenAI-compatible providers, Ollama, and other AI endpoints running on your LAN. / PI-Desktop 需要本地网络访问权限,以便助手请求同一局域网中的 OpenAI 兼容服务、Ollama 等本地 AI 端点。",
"NSMicrophoneUsageDescription": "PI-Desktop plugins may use the microphone after you grant them audio access. / 在你授权后,PI-Desktop 插件可以使用麦克风。",
"NSMicrophoneUsageDescription": "PI-Desktop uses the microphone for Live voice calls and authorized plugin audio access. / PI-Desktop 使用麦克风进行实时语音通话,以及你已授权的插件音频访问。",
"NSBonjourServices": [
"_http._tcp",
"_https._tcp"
Expand Down
Loading
Loading