diff --git a/README.md b/README.md index ae31cd7..993eb93 100644 --- a/README.md +++ b/README.md @@ -33,7 +33,7 @@ import { init } from "@deepgram/agents-widget"; init({ tokenFactory: () => fetch('/api/deepgram-token').then(r => r.text()), - agent: { think: { provider: { type: 'open_ai' }, model: 'gpt-4o-mini' } }, + agent: { think: { provider: { type: 'open_ai', model: 'gpt-4o-mini' } } }, }); ``` @@ -51,7 +51,7 @@ function App() { fetch('/api/deepgram-token').then(r => r.text()) }, - agent: { think: { provider: { type: 'open_ai' }, model: 'gpt-4o-mini' } }, + agent: { think: { provider: { type: 'open_ai', model: 'gpt-4o-mini' } } }, }} > @@ -69,7 +69,7 @@ import { AgentSession, AgentMicrophone, AgentPlayer } from "@deepgram/agents"; const session = new AgentSession({ auth: { tokenFactory: () => fetch('/api/deepgram-token').then(r => r.text()) }, - agent: { think: { provider: { type: 'open_ai' }, model: 'gpt-4o-mini' } }, + agent: { think: { provider: { type: 'open_ai', model: 'gpt-4o-mini' } } }, }); const player = new AgentPlayer(); diff --git a/bun.lock b/bun.lock index 19ff40f..a9efac5 100644 --- a/bun.lock +++ b/bun.lock @@ -18,7 +18,7 @@ "@playwright/test": "1.59.1", "@tailwindcss/vite": "4.2.4", "tailwindcss": "4.2.4", - "typescript": "6.0.2", + "typescript": "5.9.3", "vite": "8.0.10", }, }, @@ -26,21 +26,21 @@ "name": "@deepgram/agents", "version": "0.1.1", "dependencies": { - "@deepgram/sdk": "5.1.0", + "@deepgram/sdk": "5.9.0", "eventemitter3": "5.0.4", }, "devDependencies": { "@types/node": "25.6.0", "happy-dom": "20.9.0", "terser": "5.46.2", - "typescript": "6.0.2", + "typescript": "5.9.3", "vite": "8.0.10", "vite-plugin-dts": "4.5.4", }, }, "packages/widget": { "name": "@deepgram/agents-widget", - "version": "0.1.6", + "version": "0.1.7", "dependencies": { "@deepgram/agents": "^0.1.1", "@deepgram/react": "^0.1.0", @@ -57,7 +57,7 @@ "happy-dom": "20.9.0", "tailwindcss": "4.2.4", "terser": "5.46.2", - "typescript": "6.0.2", + "typescript": "5.9.3", "vite": "8.0.10", "vite-plugin-css-injected-by-js": "4.0.1", "vite-plugin-dts": "4.5.4", @@ -88,9 +88,9 @@ "@deepgram/react": ["@deepgram/react@0.1.0", "", { "dependencies": { "@deepgram/agents": "^0.1.1" }, "peerDependencies": { "react": ">=18.0.0", "react-dom": ">=18.0.0" } }, "sha512-l7UMFnb/JAOpw8h2jacqVK2z/9ka9iRH7xJhXFBIbH5JGyhZ9mc+a3BZoxRtz7yfQrALnCegSR0J2BGG63PsnA=="], - "@deepgram/sdk": ["@deepgram/sdk@5.1.0", "", { "dependencies": { "ws": "^8.20.0" } }, "sha512-pw4JsFLglvIAjk1C4rALHYBW5dHt30d3j25jsJhvoexTziLSTBK9LJlpeml/UDwHMILu4ZbBs3o0ZsF9DaHuZQ=="], + "@deepgram/sdk": ["@deepgram/sdk@5.9.0", "", { "dependencies": { "ws": "^8.20.0" } }, "sha512-MZge3ClQnUbFb+3cz93WI+aljPo8xPVOH5aI/oQGZNuN7JbNJBkmtXnv0Y/3ocFyqmAJltrF6U8t0gpdXzNXUw=="], - "@deepgram/ui": ["@deepgram/ui@0.1.0", "", { "dependencies": { "@deepgram/agents": "^0.1.1", "@deepgram/react": "^0.1.0", "@radix-ui/react-scroll-area": "1.2.10", "@radix-ui/react-select": "2.2.6", "@radix-ui/react-slot": "1.2.4", "@radix-ui/react-toggle": "1.1.10", "class-variance-authority": "0.7.1", "clsx": "2.1.1", "lucide-react": "1.11.0", "tailwind-merge": "3.5.0" }, "peerDependencies": { "react": ">=18.0.0", "react-dom": ">=18.0.0" } }, "sha512-vBFu5GXZnUuUonaUC9F2zmeX1er26kDouHXgGXPJ3qCMOEbFXHGb3gPe4hlljaZeQN3oWUG4ldBptbpcwY3B9w=="], + "@deepgram/ui": ["@deepgram/ui@0.1.4", "", { "dependencies": { "@deepgram/agents": "^0.1.1", "@deepgram/react": "^0.1.0", "@radix-ui/react-scroll-area": "1.2.10", "@radix-ui/react-select": "2.2.6", "@radix-ui/react-slot": "1.2.4", "@radix-ui/react-toggle": "1.1.10", "class-variance-authority": "0.7.1", "clsx": "2.1.1", "lucide-react": "1.11.0", "tailwind-merge": "3.5.0", "tailwindcss-scoped-preflight": "^4.0.6" }, "peerDependencies": { "react": ">=18.0.0", "react-dom": ">=18.0.0" } }, "sha512-v2ViBAmvb9f3+803n5b3tL1nkS3pzKvBQwYFIfpPuBfek5Iq1l3z5ybQ45wWWkmhtnrpmVM7KwoXvhCM/94VAw=="], "@emnapi/core": ["@emnapi/core@1.10.0", "", { "dependencies": { "@emnapi/wasi-threads": "1.2.1", "tslib": "^2.4.0" } }, "sha512-yq6OkJ4p82CAfPl0u9mQebQHKPJkY7WrIuk205cTYnYe+k2Z8YBh11FrbRG/H6ihirqcacOgl2BIO8oyMQLeXw=="], @@ -630,7 +630,7 @@ "tslib": ["tslib@2.8.1", "", {}, "sha512-oJFu94HQb+KVduSUQL7wnpmqnfmLsOA/nAh6b6EH0wCEoK0/mPeXU6c3wKDV83MkOuHPRHtSXKKU99IBazS/2w=="], - "typescript": ["typescript@6.0.2", "", { "bin": { "tsc": "bin/tsc", "tsserver": "bin/tsserver" } }, "sha512-bGdAIrZ0wiGDo5l8c++HWtbaNCWTS4UTv7RaTH/ThVIgjkveJt83m74bBHMJkuCbslY8ixgLBVZJIOiQlQTjfQ=="], + "typescript": ["typescript@5.9.3", "", { "bin": { "tsc": "bin/tsc", "tsserver": "bin/tsserver" } }, "sha512-jl1vZzPDinLr9eUt3J/t7V6FgNEw9QjvBPdysz9KfQDD41fQrC2Y4vKQdiaUpFT4bXlb1RHhLpp8wtm6M5TgSw=="], "ufo": ["ufo@1.6.4", "", {}, "sha512-JFNbkD1Svwe0KvGi8GOeLcP4kAWQ609twvCdcHxq1oSL8svv39ZuSvajcD8B+5D0eL4+s1Is2D/O6KN3qcTeRA=="], @@ -660,10 +660,6 @@ "ws": ["ws@8.20.0", "", { "peerDependencies": { "bufferutil": "^4.0.1", "utf-8-validate": ">=5.0.2" }, "optionalPeers": ["bufferutil", "utf-8-validate"] }, "sha512-sAt8BhgNbzCtgGbt2OxmpuryO63ZoDk/sqaB/znQm94T4fCEsy/yV+7CdC1kJhOU9lboAEU7R3kquuycDoibVA=="], - "@deepgram/agents-widget/@deepgram/ui": ["@deepgram/ui@0.1.4", "", { "dependencies": { "@deepgram/agents": "^0.1.1", "@deepgram/react": "^0.1.0", "@radix-ui/react-scroll-area": "1.2.10", "@radix-ui/react-select": "2.2.6", "@radix-ui/react-slot": "1.2.4", "@radix-ui/react-toggle": "1.1.10", "class-variance-authority": "0.7.1", "clsx": "2.1.1", "lucide-react": "1.11.0", "tailwind-merge": "3.5.0", "tailwindcss-scoped-preflight": "^4.0.6" }, "peerDependencies": { "react": ">=18.0.0", "react-dom": ">=18.0.0" } }, "sha512-v2ViBAmvb9f3+803n5b3tL1nkS3pzKvBQwYFIfpPuBfek5Iq1l3z5ybQ45wWWkmhtnrpmVM7KwoXvhCM/94VAw=="], - - "@microsoft/api-extractor/typescript": ["typescript@5.9.3", "", { "bin": { "tsc": "bin/tsc", "tsserver": "bin/tsserver" } }, "sha512-jl1vZzPDinLr9eUt3J/t7V6FgNEw9QjvBPdysz9KfQDD41fQrC2Y4vKQdiaUpFT4bXlb1RHhLpp8wtm6M5TgSw=="], - "@radix-ui/react-collection/@radix-ui/react-slot": ["@radix-ui/react-slot@1.2.3", "", { "dependencies": { "@radix-ui/react-compose-refs": "1.1.2" }, "peerDependencies": { "@types/react": "*", "react": "^16.8 || ^17.0 || ^18.0 || ^19.0 || ^19.0.0-rc" }, "optionalPeers": ["@types/react"] }, "sha512-aeNmHnBxbi2St0au6VBVC7JXFlhLlOnvIIlePNniyUNAClzmtAUEY8/pBiK3iHjufOlwA+c20/8jngo7xcrg8A=="], "@radix-ui/react-primitive/@radix-ui/react-slot": ["@radix-ui/react-slot@1.2.3", "", { "dependencies": { "@radix-ui/react-compose-refs": "1.1.2" }, "peerDependencies": { "@types/react": "*", "react": "^16.8 || ^17.0 || ^18.0 || ^19.0 || ^19.0.0-rc" }, "optionalPeers": ["@types/react"] }, "sha512-aeNmHnBxbi2St0au6VBVC7JXFlhLlOnvIIlePNniyUNAClzmtAUEY8/pBiK3iHjufOlwA+c20/8jngo7xcrg8A=="], diff --git a/examples/package.json b/examples/package.json index 384b8a6..03d8dd6 100644 --- a/examples/package.json +++ b/examples/package.json @@ -19,7 +19,7 @@ "@playwright/test": "1.59.1", "@tailwindcss/vite": "4.2.4", "tailwindcss": "4.2.4", - "typescript": "6.0.2", + "typescript": "5.9.3", "vite": "8.0.10" } } diff --git a/examples/vite.config.ts b/examples/vite.config.ts index 9c400b0..d4323de 100644 --- a/examples/vite.config.ts +++ b/examples/vite.config.ts @@ -139,7 +139,6 @@ export default defineConfig(({ mode }) => { path.resolve(".."), // agent monorepo root path.resolve("../../react"), // deepgram/react repo path.resolve("../../ui"), // deepgram/ui repo - path.resolve("../node_modules/.bun"), // bun cache (WASM/ONNX for VAD) ], }, }, diff --git a/packages/sdk/README.md b/packages/sdk/README.md index b461df6..86d66d1 100644 --- a/packages/sdk/README.md +++ b/packages/sdk/README.md @@ -15,7 +15,12 @@ import { AgentSession, AgentMicrophone, AgentPlayer } from "@deepgram/agents"; const session = new AgentSession({ auth: { tokenFactory: () => fetch('/api/deepgram-token').then(r => r.text()) }, - agent: { think: { provider: { type: 'open_ai' }, model: 'gpt-4o-mini' } }, + agent: { + listen: { provider: { type: 'deepgram', version: 'v1', model: 'nova-3' } }, + think: { provider: { type: 'open_ai', model: 'gpt-4o-mini' } }, + speak: { provider: { type: 'deepgram', model: 'aura-2-thalia-en' } }, + }, + audio: { output: { encoding: 'linear16', sampleRate: 24_000 } }, }); const player = new AgentPlayer(); @@ -59,6 +64,9 @@ Core WebSocket session. Wraps `@deepgram/sdk`'s agent connection with: - Exponential-backoff reconnect with jitter - Automatic KeepAlive pings - Audio buffering until `SettingsApplied` +- Ordered normalization of SDK binary `Blob` messages to `ArrayBuffer` + +`@deepgram/agents` depends directly on `@deepgram/sdk` 5.9.0. It disables the underlying socket's retries and owns application-level reconnect so each attempt can refresh authentication and restore Settings. Accumulated conversation context is restored for inline agent configurations only; the protocol cannot attach inline context to an agent UUID. Consumers continue to receive the documented `ArrayBuffer` from the `audio` event even though SDK 5.9 delivers binary messages to its callback as `Blob`. ### Constructor @@ -74,7 +82,7 @@ interface AgentSessionConfig { agent: AgentSettingsObject | string; // inline config or pre-built agent UUID audio?: { input?: { encoding?: AudioEncoding; sampleRate?: number }; // default: linear16 @ 16kHz - output?: { encoding?: OutputEncoding; sampleRate?: number }; // default: 24kHz + output?: { encoding?: OutputEncoding; sampleRate?: number }; // default: linear16 @ 24kHz }; keepAliveInterval?: number; // default: 10_000ms reconnect?: ReconnectConfig; @@ -89,14 +97,16 @@ interface AgentSessionConfig { |--------|-------------| | `connect()` | Open WebSocket connection | | `disconnect()` | Close connection (no reconnect) | +| `clearConversationHistory()` | Clear context accumulated for future inline-agent reconnects | | `sendAudio(data: ArrayBuffer)` | Send PCM audio frame (queued until SettingsApplied) | | `injectUserMessage(content)` | Send a text message as the user | -| `injectAgentMessage(message)` | Inject text as the agent | +| `injectAgentMessage(message, behavior?)` | Inject text as the agent with optional `default`, `queue`, or `interrupt` behavior | +| `updateListen(listen)` | Update listening settings mid-session | | `updateSpeak(speak)` | Update TTS settings mid-session | | `updateThink(think)` | Update LLM settings mid-session | | `updatePrompt(prompt)` | Update system prompt mid-session | | `sendFunctionCallResponse(id, name, content)` | Respond to a function call request | -| `getId()` | Returns session ID (available after Welcome) | +| `getId()` | Returns the request ID from Welcome | ### Events @@ -113,6 +123,9 @@ session.on("function-call-response", (msg) => {}); session.on("prompt-updated", (msg) => {}); session.on("speak-updated", (msg) => {}); session.on("think-updated", (msg) => {}); +session.on("listen-updated", (msg) => {}); +session.on("latency-report", (msg) => {}); +session.on("history", (msg) => {}); session.on("injection-refused", (msg) => {}); session.on("error", (msg) => {}); session.on("warning", (msg) => {}); @@ -136,7 +149,7 @@ session.state; // "idle" | "connecting" | "connected" | "reconnecting" | "discon ### Reconnect -Auto-reconnect is enabled by default with exponential backoff + jitter. Configure via `reconnect`: +Auto-reconnect is enabled by default with exponential backoff + jitter. After `SettingsApplied`, the latest value of each runtime setting (prompt, listen, speak, think) is replayed, ordered by when each setting was last updated, before audio queued during reconnect is flushed. Accumulated text and completed function-call context is restored for inline agent configurations only. Audio sent after an unexpected close or error stays queued across attempts and resumes only after the replacement connection receives `SettingsApplied`. The consecutive-attempt counter also resets at `SettingsApplied`, not merely when the transport opens. Configure via `reconnect`: ```ts { @@ -187,7 +200,7 @@ mic.on("error", (err: Error) => {}); ## AgentPlayer -Decodes and plays PCM Int16 audio from the agent. Provides volume/frequency analysis for visualizations and supports barge-in via `interrupt()`. +Decodes and plays raw PCM Int16 (`linear16`) audio from the agent. Provides volume/frequency analysis for visualizations and supports barge-in via `interrupt()`. `AgentPlayer` does not decode compressed output; use `linear16` output or provide a separate decoder/player. ### Usage @@ -230,12 +243,15 @@ export type { AgentSessionEvents, MicrophoneOptions, PlayerOptions, - AgentSettingsObject, ThinkSettings, SpeakSettings, + AgentSettingsObject, ListenSettings, ThinkSettings, SpeakSettings, + AgentMessageBehavior, // Server messages WelcomeMessage, SettingsAppliedMessage, ConversationTextMessage, UserStartedSpeakingMessage, AgentThinkingMessage, - FunctionCallRequestMessage, FunctionCallItem, + FunctionCallRequestMessage, FunctionCallResponseMessage, FunctionCallItem, AgentStartedSpeakingMessage, AgentAudioDoneMessage, + PromptUpdatedMessage, SpeakUpdatedMessage, ThinkUpdatedMessage, + ListenUpdatedMessage, LatencyReportMessage, HistoryMessage, AgentErrorMessage, AgentWarningMessage, InjectionRefusedMessage, ServerMessage, }; diff --git a/packages/sdk/package.json b/packages/sdk/package.json index f57bc90..f3b91d9 100644 --- a/packages/sdk/package.json +++ b/packages/sdk/package.json @@ -1,7 +1,7 @@ { "name": "@deepgram/agents", "version": "0.1.1", - "description": "Deepgram Voice Agent SDK — browser/Node WebSocket client with microphone, audio playback, and VAD", + "description": "Deepgram Voice Agent SDK — browser/Node WebSocket client with microphone and audio playback", "type": "module", "main": "dist/index.cjs", "module": "dist/index.js", @@ -33,14 +33,14 @@ "test:watch": "bun test --watch" }, "dependencies": { - "@deepgram/sdk": "5.1.0", + "@deepgram/sdk": "5.9.0", "eventemitter3": "5.0.4" }, "devDependencies": { "@types/node": "25.6.0", "happy-dom": "20.9.0", "terser": "5.46.2", - "typescript": "6.0.2", + "typescript": "5.9.3", "vite": "8.0.10", "vite-plugin-dts": "4.5.4" }, diff --git a/packages/sdk/src/__tests__/agent-session.test.ts b/packages/sdk/src/__tests__/agent-session.test.ts index 14b5511..c02bea9 100644 --- a/packages/sdk/src/__tests__/agent-session.test.ts +++ b/packages/sdk/src/__tests__/agent-session.test.ts @@ -19,11 +19,15 @@ const { AgentSession } = await import("../agent-session.js"); function createSession(overrides = {}) { return new AgentSession({ auth: { apiKey: "test-key" }, - agent: { think: { type: "open_ai", model: "gpt-4o-mini" } }, + agent: { think: { provider: { type: "open_ai", model: "gpt-4o-mini" } } }, ...overrides, }); } +async function flushMicrotasks() { + for (let i = 0; i < 8; i++) await Promise.resolve(); +} + describe("AgentSession", () => { beforeEach(() => { jest.useFakeTimers(); @@ -96,7 +100,7 @@ describe("AgentSession", () => { await session.connect(); // Simulate Welcome message - const welcomeMsg = { type: "Welcome" as const, session_id: "test-123" }; + const welcomeMsg = { type: "Welcome" as const, request_id: "test-123" }; // Get the message handler const messageHandler = mockSocket.on.mock.calls.find( (c) => c[0] === "message", @@ -110,13 +114,29 @@ describe("AgentSession", () => { expect(settings.audio.input.sample_rate).toBe(16_000); }); + it("handles Welcome emitted immediately after the socket opens", async () => { + const session = createSession(); + mockSocket.connect.mockImplementation(() => { + mockSocket._emit("open"); + mockSocket._emit("message", { + type: "Welcome", + request_id: "immediate-welcome", + }); + }); + + await session.connect(); + + expect(mockSocket.sendSettings).toHaveBeenCalledTimes(1); + expect(session.getId()).toBe("immediate-welcome"); + }); + it("emits welcome event", async () => { const session = createSession(); const welcomeHandler = jest.fn(); session.on("welcome", welcomeHandler); await session.connect(); - const welcomeMsg = { type: "Welcome" as const, session_id: "test-123" }; + const welcomeMsg = { type: "Welcome" as const, request_id: "test-123" }; const messageHandler = mockSocket.on.mock.calls.find( (c) => c[0] === "message", )![1] as (msg: unknown) => void; @@ -125,6 +145,20 @@ describe("AgentSession", () => { expect(welcomeHandler).toHaveBeenCalledWith(welcomeMsg); }); + it("tracks the current request ID and clears it on disconnect", async () => { + const session = createSession(); + await session.connect(); + const messageHandler = mockSocket.on.mock.calls.find( + (c) => c[0] === "message", + )![1] as (msg: unknown) => void; + + messageHandler({ type: "Welcome", request_id: "request-123" }); + expect(session.getId()).toBe("request-123"); + + session.disconnect(); + expect(session.getId()).toBeNull(); + }); + it("includes output config when specified", async () => { const session = createSession({ audio: { @@ -186,6 +220,40 @@ describe("AgentSession", () => { expect(mockSocket.sendMedia).toHaveBeenCalledWith(frame2); }); + it("drains accepted messages before handling an audio flush failure", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + const events: string[] = []; + session.on("audio", () => events.push("audio")); + session.on("warning", () => events.push("warning")); + session.on("sdk-error", () => events.push("sdk-error")); + session.on("reconnecting", () => events.push("reconnecting")); + await session.connect(); + session.sendAudio(new ArrayBuffer(320)); + + let resolveConversion!: (data: ArrayBuffer) => void; + const incomingAudio = new Blob(); + Object.defineProperty(incomingAudio, "arrayBuffer", { + value: () => new Promise((resolve) => { resolveConversion = resolve; }), + }); + mockSocket._emit("message", incomingAudio); + mockSocket._emit("message", { type: "SettingsApplied" }); + mockSocket._emit("message", { + type: "Warning", + code: "TEST_WARNING", + description: "queued after settings", + }); + mockSocket.sendMedia.mockImplementation(() => { + throw new Error("flush failed"); + }); + + resolveConversion(new ArrayBuffer(160)); + await flushMicrotasks(); + + expect(events).toEqual(["audio", "warning", "sdk-error", "reconnecting"]); + }); + it("sends audio immediately after SettingsApplied", async () => { const session = createSession(); await session.connect(); @@ -215,6 +283,111 @@ describe("AgentSession", () => { expect(handler).toHaveBeenCalledWith(msg); }); + + it("normalizes Blob audio to ArrayBuffer in frame order", async () => { + const session = createSession(); + const handler = jest.fn(); + session.on("audio", handler); + await session.connect(); + + let resolveFirst!: (data: ArrayBuffer) => void; + const firstData = new ArrayBuffer(160); + const secondData = new ArrayBuffer(320); + const first = new Blob(); + const second = new Blob(); + const firstArrayBuffer = jest.fn( + () => new Promise((resolve) => { resolveFirst = resolve; }), + ); + const secondArrayBuffer = jest.fn(async () => secondData); + Object.defineProperty(first, "arrayBuffer", { value: firstArrayBuffer }); + Object.defineProperty(second, "arrayBuffer", { value: secondArrayBuffer }); + + const messageHandler = mockSocket.on.mock.calls.find( + (c) => c[0] === "message", + )![1] as (msg: unknown) => void; + messageHandler(first); + messageHandler(second); + + expect(firstArrayBuffer).toHaveBeenCalledTimes(1); + expect(secondArrayBuffer).not.toHaveBeenCalled(); + + resolveFirst(firstData); + await flushMicrotasks(); + + expect(handler.mock.calls.map((call) => call[0])).toEqual([ + firstData, + secondData, + ]); + }); + + it("emits delayed Blob audio before following AgentAudioDone", async () => { + const session = createSession(); + const events: string[] = []; + session.on("audio", () => events.push("audio")); + session.on("agent-audio-done", () => events.push("done")); + await session.connect(); + + let resolveConversion!: (data: ArrayBuffer) => void; + const audio = new Blob(); + Object.defineProperty(audio, "arrayBuffer", { + value: () => new Promise((resolve) => { resolveConversion = resolve; }), + }); + const messageHandler = mockSocket.on.mock.calls.find( + (c) => c[0] === "message", + )![1] as (msg: unknown) => void; + messageHandler(audio); + messageHandler({ type: "AgentAudioDone" }); + + expect(events).toEqual([]); + + resolveConversion(new ArrayBuffer(160)); + await flushMicrotasks(); + + expect(events).toEqual(["audio", "done"]); + }); + + it("emits sdk-error when Blob conversion fails", async () => { + const session = createSession(); + const handler = jest.fn(); + session.on("sdk-error", handler); + await session.connect(); + + const conversionError = new Error("blob conversion failed"); + const audio = new Blob(); + Object.defineProperty(audio, "arrayBuffer", { + value: () => Promise.reject(conversionError), + }); + const messageHandler = mockSocket.on.mock.calls.find( + (c) => c[0] === "message", + )![1] as (msg: unknown) => void; + messageHandler(audio); + await flushMicrotasks(); + + expect(handler).toHaveBeenCalledWith(conversionError); + }); + + it("does not emit a late Blob conversion after disconnect", async () => { + const session = createSession(); + const handler = jest.fn(); + session.on("audio", handler); + await session.connect(); + + let resolveConversion!: (data: ArrayBuffer) => void; + const audio = new Blob(); + Object.defineProperty(audio, "arrayBuffer", { + value: () => new Promise((resolve) => { resolveConversion = resolve; }), + }); + const messageHandler = mockSocket.on.mock.calls.find( + (c) => c[0] === "message", + )![1] as (msg: unknown) => void; + messageHandler(audio); + + session.disconnect(); + resolveConversion(new ArrayBuffer(160)); + await flushMicrotasks(); + + expect(handler).not.toHaveBeenCalled(); + }); }); // --------------------------------------------------------------------------- @@ -241,6 +414,16 @@ describe("AgentSession", () => { expect(handler).toHaveBeenCalledWith(msg); }); + it("clears accumulated conversation history", async () => { + const { session, dispatch } = await setupConnectedSession(); + dispatch({ type: "ConversationText", role: "user", content: "forget this" }); + expect(session.conversationHistory).toHaveLength(1); + + session.clearConversationHistory(); + + expect(session.conversationHistory).toEqual([]); + }); + it("dispatches UserStartedSpeaking event", async () => { const { session, dispatch } = await setupConnectedSession(); const handler = jest.fn(); @@ -271,6 +454,107 @@ describe("AgentSession", () => { expect(handler).toHaveBeenCalledWith(msg); }); + it("records completed client-side function calls without duplicate server records", async () => { + const { session, dispatch } = await setupConnectedSession(); + const functionCall = { + id: "call-1", + name: "get_weather", + arguments: '{"city":"Austin"}', + client_side: true, + thought_signature: "signature", + }; + const completedCall = { + id: "call-1", + name: "get_weather", + arguments: '{"city":"Austin"}', + client_side: true, + response: '{"temp":72}', + thought_signature: "signature", + }; + + dispatch({ type: "FunctionCallRequest", functions: [functionCall] }); + session.sendFunctionCallResponse("call-1", "get_weather", '{"temp":72}'); + dispatch({ + type: "FunctionCallResponse", + id: "call-1", + name: "get_weather", + content: '{"temp":72}', + }); + dispatch({ type: "History", function_calls: [completedCall] }); + + expect(session.conversationHistory).toEqual([ + { type: "History", function_calls: [completedCall] }, + ]); + }); + + it("records completed server-side function calls", async () => { + const { session, dispatch } = await setupConnectedSession(); + + dispatch({ + type: "FunctionCallRequest", + functions: [{ + id: "server-call", + name: "lookup_account", + arguments: '{"id":42}', + client_side: false, + }], + }); + dispatch({ + type: "FunctionCallResponse", + id: "server-call", + name: "lookup_account", + content: '{"status":"active"}', + }); + + expect(session.conversationHistory).toEqual([{ + type: "History", + function_calls: [{ + id: "server-call", + name: "lookup_account", + arguments: '{"id":42}', + client_side: false, + response: '{"status":"active"}', + }], + }]); + }); + + it("records unseen History function calls once", async () => { + const { session, dispatch } = await setupConnectedSession(); + const history = { + type: "History" as const, + function_calls: [{ + id: "history-call", + name: "search", + arguments: '{"query":"voice"}', + client_side: false, + response: "result", + }], + }; + + dispatch(history); + dispatch(history); + + expect(session.conversationHistory).toEqual([history]); + }); + + it("clearConversationHistory clears pending function-call tracking", async () => { + const { session, dispatch } = await setupConnectedSession(); + dispatch({ + type: "FunctionCallRequest", + functions: [{ + id: "forgotten-call", + name: "search", + arguments: "{}", + client_side: true, + }], + }); + + session.clearConversationHistory(); + session.sendFunctionCallResponse("forgotten-call", "search", "result"); + + expect(session.conversationHistory).toEqual([]); + }); + it("dispatches AgentStartedSpeaking event", async () => { const { session, dispatch } = await setupConnectedSession(); const handler = jest.fn(); @@ -320,6 +604,20 @@ describe("AgentSession", () => { dispatch(msg); expect(handler).toHaveBeenCalledWith(msg); }); + + it.each([ + ["ListenUpdated", "listen-updated"], + ["LatencyReport", "latency-report"], + ["History", "history"], + ] as const)("dispatches %s event", async (type, event) => { + const { session, dispatch } = await setupConnectedSession(); + const handler = jest.fn(); + session.on(event, handler); + + const msg = { type }; + dispatch(msg); + expect(handler).toHaveBeenCalledWith(msg); + }); }); // --------------------------------------------------------------------------- @@ -327,11 +625,29 @@ describe("AgentSession", () => { // --------------------------------------------------------------------------- describe("send methods", () => { + it("updateListen delegates to socket", async () => { + const session = createSession(); + await session.connect(); + mockSocket._emit("message", { type: "SettingsApplied" }); + + const listen = { + provider: { type: "deepgram" as const, model: "flux-general-en" }, + }; + session.updateListen(listen); + expect(mockSocket.sendUpdateListen).toHaveBeenCalledWith({ + type: "UpdateListen", + listen, + }); + }); + it("updateSpeak delegates to socket", async () => { const session = createSession(); await session.connect(); + mockSocket._emit("message", { type: "SettingsApplied" }); - const speak = { type: "open_ai" as const, model: "tts-1" }; + const speak = { + provider: { type: "open_ai" as const, model: "tts-1", voice: "alloy" }, + }; session.updateSpeak(speak); expect(mockSocket.sendUpdateSpeak).toHaveBeenCalled(); }); @@ -339,8 +655,9 @@ describe("AgentSession", () => { it("updateThink delegates to socket", async () => { const session = createSession(); await session.connect(); + mockSocket._emit("message", { type: "SettingsApplied" }); - const think = { type: "open_ai" as const, model: "gpt-4o" }; + const think = { provider: { type: "open_ai" as const, model: "gpt-4o" } }; session.updateThink(think); expect(mockSocket.sendUpdateThink).toHaveBeenCalled(); }); @@ -348,6 +665,7 @@ describe("AgentSession", () => { it("updatePrompt delegates to socket", async () => { const session = createSession(); await session.connect(); + mockSocket._emit("message", { type: "SettingsApplied" }); session.updatePrompt("You are a helpful assistant"); expect(mockSocket.sendUpdatePrompt).toHaveBeenCalledWith({ @@ -378,6 +696,21 @@ describe("AgentSession", () => { }); }); + it.each(["default", "queue", "interrupt"] as const)( + "injectAgentMessage forwards %s behavior", + async (behavior) => { + const session = createSession(); + await session.connect(); + + session.injectAgentMessage("Urgent update", behavior); + expect(mockSocket.sendInjectAgentMessage).toHaveBeenCalledWith({ + type: "InjectAgentMessage", + message: "Urgent update", + behavior, + }); + }, + ); + it("sendFunctionCallResponse delegates to socket", async () => { const session = createSession(); await session.connect(); @@ -390,6 +723,23 @@ describe("AgentSession", () => { content: '{"temp":72}', }); }); + + it("routes runtime update send failures through reconnect handling", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 1, baseDelay: 100, jitter: false }, + }); + const error = new Error("update failed"); + const errorHandler = jest.fn(); + session.on("sdk-error", errorHandler); + await session.connect(); + mockSocket._emit("message", { type: "SettingsApplied" }); + mockSocket.sendUpdatePrompt.mockImplementationOnce(() => { throw error; }); + + expect(() => session.updatePrompt("new prompt")).not.toThrow(); + + expect(errorHandler).toHaveBeenCalledWith(error); + expect(session.state).toBe("reconnecting"); + }); }); // --------------------------------------------------------------------------- @@ -397,6 +747,16 @@ describe("AgentSession", () => { // --------------------------------------------------------------------------- describe("reconnection", () => { + it("disables @deepgram/sdk reconnect attempts", async () => { + const session = createSession(); + await session.connect(); + + expect(mockClient.agent.v1.connect).toHaveBeenCalledWith({ + Authorization: "Token test-key", + reconnectAttempts: 0, + }); + }); + it("schedules reconnect on socket close with exponential backoff", async () => { const session = createSession({ reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, @@ -444,6 +804,540 @@ describe("AgentSession", () => { expect(session.state).toBe("disconnected"); expect(disconnected).toHaveBeenCalled(); }); + + it("resets settings and buffers audio until reconnect settings are applied", async () => { + const session = createSession({ + agent: { + think: { provider: { type: "open_ai", model: "gpt-4o-mini" } }, + context: { + messages: [{ type: "History", role: "assistant", content: "configured context" }], + }, + }, + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + await session.connect(); + + const firstSocket = mockSocket; + const firstMessageHandler = firstSocket.on.mock.calls.find( + (c) => c[0] === "message", + )![1] as (msg: unknown) => void; + firstMessageHandler({ + type: "ConversationText", + role: "user", + content: "remember this", + }); + firstMessageHandler({ type: "SettingsApplied" }); + + const nextSocket = createMockSocket(); + createMockDeepgramClient(nextSocket); + mockClient.agent.v1.connect.mockResolvedValue(nextSocket); + + firstSocket._emit("close", { code: 1006, reason: "abnormal" }); + const queuedFrame = new ArrayBuffer(320); + session.sendAudio(queuedFrame); + expect(firstSocket.sendMedia).not.toHaveBeenCalled(); + + jest.advanceTimersByTime(100); + await flushMicrotasks(); + expect(session.state).toBe("connected"); + + const nextMessageHandler = nextSocket.on.mock.calls.find( + (c) => c[0] === "message", + )![1] as (msg: unknown) => void; + nextMessageHandler({ type: "Welcome", request_id: "reconnected" }); + expect(nextSocket.sendSettings.mock.calls[0][0].agent.context.messages).toEqual([ + { type: "History", role: "assistant", content: "configured context" }, + { type: "History", role: "user", content: "remember this" }, + ]); + expect(nextSocket.sendMedia).not.toHaveBeenCalled(); + + nextMessageHandler({ type: "SettingsApplied" }); + expect(nextSocket.sendMedia).toHaveBeenCalledWith(queuedFrame); + }); + + it("replays runtime updates in wire order before queued audio", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + await session.connect(); + const firstSocket = mockSocket; + firstSocket._emit("message", { type: "SettingsApplied" }); + + const listen = { + provider: { type: "deepgram" as const, model: "flux-general-en" }, + }; + const speak = { + provider: { type: "open_ai" as const, model: "tts-1", voice: "alloy" }, + }; + const think = { + provider: { type: "open_ai" as const, model: "gpt-4o" }, + }; + session.updatePrompt("updated prompt"); + session.updateListen(listen); + session.updateSpeak(speak); + session.updateThink(think); + + const nextSocket = createMockSocket(); + createMockDeepgramClient(nextSocket); + mockClient.agent.v1.connect.mockResolvedValue(nextSocket); + const wireOrder: string[] = []; + nextSocket.sendUpdatePrompt.mockImplementation(() => wireOrder.push("prompt")); + nextSocket.sendUpdateListen.mockImplementation(() => wireOrder.push("listen")); + nextSocket.sendUpdateSpeak.mockImplementation(() => wireOrder.push("speak")); + nextSocket.sendUpdateThink.mockImplementation(() => wireOrder.push("think")); + nextSocket.sendMedia.mockImplementation(() => wireOrder.push("audio")); + + firstSocket._emit("close", { code: 1006 }); + session.sendAudio(new ArrayBuffer(320)); + jest.advanceTimersByTime(100); + await flushMicrotasks(); + nextSocket._emit("message", { type: "SettingsApplied" }); + + expect(wireOrder).toEqual(["prompt", "listen", "speak", "think", "audio"]); + }); + + it("replays only the latest value for each runtime update type", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + await session.connect(); + const firstSocket = mockSocket; + firstSocket._emit("message", { type: "SettingsApplied" }); + + session.updatePrompt("old prompt"); + session.updateListen({ + provider: { type: "deepgram", model: "nova-3" }, + }); + session.updatePrompt("latest prompt"); + session.updateSpeak({ + provider: { type: "deepgram", model: "aura-2-thalia-en" }, + }); + session.updateListen({ + provider: { type: "deepgram", model: "flux-general-en" }, + }); + + const nextSocket = createMockSocket(); + createMockDeepgramClient(nextSocket); + mockClient.agent.v1.connect.mockResolvedValue(nextSocket); + const wireOrder: string[] = []; + nextSocket.sendUpdatePrompt.mockImplementation(() => wireOrder.push("prompt")); + nextSocket.sendUpdateSpeak.mockImplementation(() => wireOrder.push("speak")); + nextSocket.sendUpdateListen.mockImplementation(() => wireOrder.push("listen")); + + firstSocket._emit("close", { code: 1006 }); + jest.advanceTimersByTime(100); + await flushMicrotasks(); + nextSocket._emit("message", { type: "SettingsApplied" }); + + expect(wireOrder).toEqual(["prompt", "speak", "listen"]); + expect(nextSocket.sendUpdatePrompt).toHaveBeenCalledTimes(1); + expect(nextSocket.sendUpdatePrompt).toHaveBeenCalledWith({ + type: "UpdatePrompt", + prompt: "latest prompt", + }); + expect(nextSocket.sendUpdateListen).toHaveBeenCalledTimes(1); + expect(nextSocket.sendUpdateListen).toHaveBeenCalledWith({ + type: "UpdateListen", + listen: { + provider: { type: "deepgram", model: "flux-general-en" }, + }, + }); + }); + + it("snapshots runtime updates before sending and replaying them", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + await session.connect(); + const firstSocket = mockSocket; + firstSocket._emit("message", { type: "SettingsApplied" }); + + const listen = { + provider: { type: "deepgram" as const, model: "flux-general-en" }, + }; + const speak = { + provider: { type: "open_ai" as const, model: "tts-1", voice: "alloy" }, + }; + const think = { + provider: { type: "open_ai" as const, model: "gpt-4o" }, + }; + let prompt = "original prompt"; + session.updateListen(listen); + session.updateSpeak(speak); + session.updateThink(think); + session.updatePrompt(prompt); + + listen.provider.model = "mutated-listen"; + speak.provider.model = "mutated-speak"; + speak.provider.voice = "mutated-voice"; + think.provider.model = "mutated-think"; + prompt = "mutated prompt"; + + expect(firstSocket.sendUpdateListen).toHaveBeenCalledWith({ + type: "UpdateListen", + listen: { provider: { type: "deepgram", model: "flux-general-en" } }, + }); + expect(firstSocket.sendUpdatePrompt).toHaveBeenCalledWith({ + type: "UpdatePrompt", + prompt: "original prompt", + }); + + const nextSocket = createMockSocket(); + createMockDeepgramClient(nextSocket); + mockClient.agent.v1.connect.mockResolvedValue(nextSocket); + firstSocket._emit("close", { code: 1006 }); + jest.advanceTimersByTime(100); + await flushMicrotasks(); + nextSocket._emit("message", { type: "SettingsApplied" }); + + expect(nextSocket.sendUpdateListen).toHaveBeenCalledWith({ + type: "UpdateListen", + listen: { provider: { type: "deepgram", model: "flux-general-en" } }, + }); + expect(nextSocket.sendUpdateSpeak).toHaveBeenCalledWith({ + type: "UpdateSpeak", + speak: { + provider: { type: "open_ai", model: "tts-1", voice: "alloy" }, + }, + }); + expect(nextSocket.sendUpdateThink).toHaveBeenCalledWith({ + type: "UpdateThink", + think: { provider: { type: "open_ai", model: "gpt-4o" } }, + }); + expect(nextSocket.sendUpdatePrompt).toHaveBeenCalledWith({ + type: "UpdatePrompt", + prompt: "original prompt", + }); + expect(prompt).toBe("mutated prompt"); + }); + + it("replays updates requested while reconnecting", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + await session.connect(); + const firstSocket = mockSocket; + firstSocket._emit("message", { type: "SettingsApplied" }); + + const nextSocket = createMockSocket(); + createMockDeepgramClient(nextSocket); + mockClient.agent.v1.connect.mockResolvedValue(nextSocket); + firstSocket._emit("close", { code: 1006 }); + + session.updatePrompt("reconnect prompt"); + session.updateListen({ + provider: { type: "deepgram", model: "flux-general-en" }, + }); + expect(firstSocket.sendUpdatePrompt).not.toHaveBeenCalled(); + expect(firstSocket.sendUpdateListen).not.toHaveBeenCalled(); + + jest.advanceTimersByTime(100); + await flushMicrotasks(); + nextSocket._emit("message", { type: "SettingsApplied" }); + + expect(nextSocket.sendUpdatePrompt).toHaveBeenCalledWith({ + type: "UpdatePrompt", + prompt: "reconnect prompt", + }); + expect(nextSocket.sendUpdateListen).toHaveBeenCalledWith({ + type: "UpdateListen", + listen: { + provider: { type: "deepgram", model: "flux-general-en" }, + }, + }); + }); + + it("clears runtime updates on an explicit fresh connect", async () => { + const session = createSession(); + await session.connect(); + const firstSocket = mockSocket; + firstSocket._emit("message", { type: "SettingsApplied" }); + session.updatePrompt("old session prompt"); + + const nextSocket = createMockSocket(); + createMockDeepgramClient(nextSocket); + mockClient.agent.v1.connect.mockResolvedValueOnce(nextSocket); + await session.connect(); + nextSocket._emit("message", { type: "SettingsApplied" }); + + expect(nextSocket.sendUpdatePrompt).not.toHaveBeenCalled(); + }); + + it("restores completed function calls for inline agent configurations", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + await session.connect(); + const firstSocket = mockSocket; + firstSocket._emit("message", { + type: "FunctionCallRequest", + functions: [{ + id: "call-1", + name: "get_weather", + arguments: '{"city":"Austin"}', + client_side: true, + }], + }); + session.sendFunctionCallResponse("call-1", "get_weather", '{"temp":72}'); + + const nextSocket = createMockSocket(); + createMockDeepgramClient(nextSocket); + mockClient.agent.v1.connect.mockResolvedValue(nextSocket); + firstSocket._emit("close", { code: 1006 }); + jest.advanceTimersByTime(100); + await flushMicrotasks(); + nextSocket._emit("message", { type: "Welcome", request_id: "reconnected" }); + + expect(nextSocket.sendSettings.mock.calls[0][0].agent.context.messages).toEqual([{ + type: "History", + function_calls: [{ + id: "call-1", + name: "get_weather", + arguments: '{"city":"Austin"}', + client_side: true, + response: '{"temp":72}', + }], + }]); + }); + + it("does not attach accumulated context to an agent UUID", async () => { + const session = createSession({ + agent: "agent-id", + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + await session.connect(); + const firstSocket = mockSocket; + firstSocket._emit("message", { + type: "ConversationText", + role: "user", + content: "remember this", + }); + + const nextSocket = createMockSocket(); + createMockDeepgramClient(nextSocket); + mockClient.agent.v1.connect.mockResolvedValue(nextSocket); + firstSocket._emit("close", { code: 1006 }); + jest.advanceTimersByTime(100); + await flushMicrotasks(); + nextSocket._emit("message", { type: "Welcome", request_id: "reconnected" }); + + expect(nextSocket.sendSettings.mock.calls[0][0].agent).toBe("agent-id"); + }); + + it("preserves queued audio across failed reconnect attempts", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + await session.connect(); + + const firstSocket = mockSocket; + const failedSocket = createMockSocket(); + createMockDeepgramClient(failedSocket, false); + const recoveredSocket = createMockSocket(); + createMockDeepgramClient(recoveredSocket); + mockClient.agent.v1.connect + .mockResolvedValueOnce(failedSocket) + .mockResolvedValueOnce(recoveredSocket); + + firstSocket._emit("close", { code: 1006 }); + const queuedFrame = new ArrayBuffer(320); + session.sendAudio(queuedFrame); + + jest.advanceTimersByTime(100); + await flushMicrotasks(); + failedSocket._emit("error", new Error("retry failed")); + await flushMicrotasks(); + + jest.advanceTimersByTime(200); + await flushMicrotasks(); + const recoveredMessageHandler = recoveredSocket.on.mock.calls.find( + (c) => c[0] === "message", + )![1] as (msg: unknown) => void; + recoveredMessageHandler({ type: "SettingsApplied" }); + + expect(recoveredSocket.sendMedia).toHaveBeenCalledWith(queuedFrame); + }); + + it("schedules only one reconnect for error and close from the same socket", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + const reconnecting = jest.fn(); + session.on("reconnecting", reconnecting); + await session.connect(); + + mockSocket._emit("error", new Error("transport failed")); + mockSocket._emit("close", { code: 1006 }); + + expect(reconnecting).toHaveBeenCalledTimes(1); + }); + + it("drains accepted messages before processing a socket close", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + const events: string[] = []; + session.on("audio", () => events.push("audio")); + session.on("error", () => events.push("protocol-error")); + session.on("reconnecting", () => events.push("reconnecting")); + session.on("disconnected", () => events.push("disconnected")); + await session.connect(); + + let resolveConversion!: (data: ArrayBuffer) => void; + const audio = new Blob(); + Object.defineProperty(audio, "arrayBuffer", { + value: () => new Promise((resolve) => { resolveConversion = resolve; }), + }); + mockSocket._emit("message", audio); + mockSocket._emit("message", { type: "Error", message: "protocol failed" }); + mockSocket._emit("close", { code: 1006 }); + + expect(events).toEqual([]); + + resolveConversion(new ArrayBuffer(160)); + await flushMicrotasks(); + + expect(events).toEqual(["audio", "protocol-error", "reconnecting"]); + }); + + it("queues audio for reconnect while a socket close waits behind inbound audio", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + const receivedAudio = jest.fn(); + session.on("audio", receivedAudio); + await session.connect(); + + const firstSocket = mockSocket; + firstSocket._emit("message", { type: "Welcome", request_id: "original" }); + firstSocket._emit("message", { type: "SettingsApplied" }); + + const nextSocket = createMockSocket(); + createMockDeepgramClient(nextSocket); + mockClient.agent.v1.connect.mockResolvedValue(nextSocket); + + let resolveConversion!: (data: ArrayBuffer) => void; + const incomingAudio = new Blob(); + const incomingFrame = new ArrayBuffer(160); + Object.defineProperty(incomingAudio, "arrayBuffer", { + value: () => new Promise((resolve) => { resolveConversion = resolve; }), + }); + firstSocket._emit("message", incomingAudio); + firstSocket._emit("close", { code: 1006 }); + + const outgoingFrame = new ArrayBuffer(320); + session.sendAudio(outgoingFrame); + expect(session.getId()).toBeNull(); + expect(firstSocket.sendMedia).not.toHaveBeenCalled(); + + resolveConversion(incomingFrame); + await flushMicrotasks(); + expect(receivedAudio).toHaveBeenCalledWith(incomingFrame); + + jest.advanceTimersByTime(100); + await flushMicrotasks(); + expect(nextSocket.sendMedia).not.toHaveBeenCalled(); + + nextSocket._emit("message", { type: "SettingsApplied" }); + expect(nextSocket.sendMedia).toHaveBeenCalledTimes(1); + expect(nextSocket.sendMedia).toHaveBeenCalledWith(outgoingFrame); + }); + + it("drains inbound audio before handling a synchronous send failure", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + const events: string[] = []; + session.on("audio", () => events.push("audio")); + session.on("sdk-error", () => events.push("sdk-error")); + session.on("reconnecting", () => events.push("reconnecting")); + await session.connect(); + + mockSocket._emit("message", { type: "SettingsApplied" }); + let resolveConversion!: (data: ArrayBuffer) => void; + const incomingAudio = new Blob(); + Object.defineProperty(incomingAudio, "arrayBuffer", { + value: () => new Promise((resolve) => { resolveConversion = resolve; }), + }); + mockSocket._emit("message", incomingAudio); + mockSocket.sendMedia.mockImplementation(() => { + throw new Error("socket write failed"); + }); + + session.sendAudio(new ArrayBuffer(320)); + expect(events).toEqual([]); + + resolveConversion(new ArrayBuffer(160)); + await flushMicrotasks(); + + expect(events).toEqual(["audio", "sdk-error", "reconnecting"]); + }); + + it("does not write messages to a socket with a queued failure", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 3, baseDelay: 100, jitter: false }, + }); + await session.connect(); + + let resolveConversion!: (data: ArrayBuffer) => void; + const incomingAudio = new Blob(); + Object.defineProperty(incomingAudio, "arrayBuffer", { + value: () => new Promise((resolve) => { resolveConversion = resolve; }), + }); + mockSocket._emit("message", incomingAudio); + mockSocket._emit("close", { code: 1006 }); + mockSocket.sendInjectUserMessage.mockImplementation(() => { throw new Error("closed"); }); + mockSocket.sendInjectAgentMessage.mockImplementation(() => { throw new Error("closed"); }); + mockSocket.sendFunctionCallResponse.mockImplementation(() => { throw new Error("closed"); }); + + expect(() => session.injectUserMessage("Hello")).not.toThrow(); + expect(() => session.injectAgentMessage("Hi there")).not.toThrow(); + expect(() => { + session.sendFunctionCallResponse("call-1", "get_weather", '{"temp":72}'); + }).not.toThrow(); + expect(mockSocket.sendInjectUserMessage).not.toHaveBeenCalled(); + expect(mockSocket.sendInjectAgentMessage).not.toHaveBeenCalled(); + expect(mockSocket.sendFunctionCallResponse).not.toHaveBeenCalled(); + + resolveConversion(new ArrayBuffer(160)); + await flushMicrotasks(); + }); + + it("counts sockets that open but fail before SettingsApplied toward maxAttempts", async () => { + const session = createSession({ + reconnect: { enabled: true, maxAttempts: 2, baseDelay: 100, jitter: false }, + }); + const reconnecting = jest.fn(); + const disconnected = jest.fn(); + session.on("reconnecting", reconnecting); + session.on("disconnected", disconnected); + await session.connect(); + + const initialSocket = mockSocket; + const firstRetrySocket = createMockSocket(); + createMockDeepgramClient(firstRetrySocket); + const secondRetrySocket = createMockSocket(); + createMockDeepgramClient(secondRetrySocket); + mockClient.agent.v1.connect + .mockResolvedValueOnce(firstRetrySocket) + .mockResolvedValueOnce(secondRetrySocket); + + initialSocket._emit("close", { code: 1006 }); + jest.advanceTimersByTime(100); + await flushMicrotasks(); + firstRetrySocket._emit("close", { code: 1006 }); + + jest.advanceTimersByTime(200); + await flushMicrotasks(); + secondRetrySocket._emit("close", { code: 1006 }); + + expect(reconnecting.mock.calls).toEqual([ + [1, 100], + [2, 200], + ]); + expect(disconnected).toHaveBeenCalledTimes(1); + expect(session.state).toBe("disconnected"); + }); }); // --------------------------------------------------------------------------- @@ -451,6 +1345,20 @@ describe("AgentSession", () => { // --------------------------------------------------------------------------- describe("cleanup", () => { + it("closes the previous socket before opening a fresh connection", async () => { + const session = createSession(); + await session.connect(); + const firstSocket = mockSocket; + const nextSocket = createMockSocket(); + createMockDeepgramClient(nextSocket); + mockClient.agent.v1.connect.mockResolvedValueOnce(nextSocket); + + await session.connect(); + + expect(firstSocket.close).toHaveBeenCalledTimes(1); + expect(nextSocket.connect).toHaveBeenCalledTimes(1); + }); + it("disconnect closes socket and stops keepalive", async () => { const session = createSession(); await session.connect(); diff --git a/packages/sdk/src/__tests__/helpers/sdk-mocks.ts b/packages/sdk/src/__tests__/helpers/sdk-mocks.ts index 2fc0185..6717528 100644 --- a/packages/sdk/src/__tests__/helpers/sdk-mocks.ts +++ b/packages/sdk/src/__tests__/helpers/sdk-mocks.ts @@ -32,6 +32,7 @@ export function createMockSocket() { sendSettings: jest.fn(), sendMedia: jest.fn(), sendKeepAlive: jest.fn(), + sendUpdateListen: jest.fn(), sendUpdateSpeak: jest.fn(), sendUpdateThink: jest.fn(), sendUpdatePrompt: jest.fn(), diff --git a/packages/sdk/src/__tests__/microphone.test.ts b/packages/sdk/src/__tests__/microphone.test.ts index b3fb59f..7a6af4c 100644 --- a/packages/sdk/src/__tests__/microphone.test.ts +++ b/packages/sdk/src/__tests__/microphone.test.ts @@ -179,21 +179,4 @@ describe("AgentMicrophone", () => { }); }); - describe("VAD error path", () => { - it("throws a descriptive error when @ricky0123/vad-web is not available", async () => { - const mic = new AgentMicrophone(onAudioFrame, { vad: true }); - - // The dynamic import of @ricky0123/vad-web will fail because the mock - // environment doesn't provide it at the expected path. The actual error - // depends on whether the optional peer dep resolves — in test env without - // proper ONNX runtime, MicVAD.new() will throw during initialization. - // We can verify the mic at least attempts to start. - try { - await mic.start(); - mic.stop(); - } catch (err) { - expect(err).toBeInstanceOf(Error); - } - }); - }); }); diff --git a/packages/sdk/src/__tests__/sdk-runtime-contract.test.ts b/packages/sdk/src/__tests__/sdk-runtime-contract.test.ts new file mode 100644 index 0000000..e9a3aec --- /dev/null +++ b/packages/sdk/src/__tests__/sdk-runtime-contract.test.ts @@ -0,0 +1,60 @@ +import { describe, expect, it } from "bun:test"; + +describe("@deepgram/sdk runtime contract", () => { + it("delivers agent binary messages to the public callback as Blob", async () => { + const source = String.raw` + import { DeepgramClient } from "@deepgram/sdk"; + + let emitMessage; + const transport = { + send() {}, + onOpen() {}, + onMessage(listener) { emitMessage = listener; }, + onError() {}, + onClose() {}, + isOpen() { return true; }, + close() {}, + }; + const client = new DeepgramClient({ + apiKey: "test-key", + reconnect: false, + transportFactory: () => transport, + }); + const socket = await client.agent.v1.connect({ reconnectAttempts: 0 }); + const opened = new Promise((resolve, reject) => { + const timer = setTimeout(() => reject(new Error("open timeout")), 1000); + socket.on("open", () => { + clearTimeout(timer); + resolve(); + }); + }); + const received = new Promise((resolve) => socket.on("message", resolve)); + + socket.connect(); + await opened; + emitMessage(new Uint8Array([1, 2, 3]).buffer); + const message = await received; + + if (!(message instanceof Blob)) { + throw new Error("expected the SDK callback to receive Blob"); + } + const bytes = new Uint8Array(await message.arrayBuffer()); + if (bytes.join(",") !== "1,2,3") { + throw new Error("SDK Blob did not preserve binary payload"); + } + socket.close(); + `; + const subprocess = Bun.spawn([process.execPath, "--eval", source], { + cwd: import.meta.dir, + stderr: "pipe", + stdout: "pipe", + }); + const [exitCode, stderr] = await Promise.all([ + subprocess.exited, + new Response(subprocess.stderr).text(), + ]); + + expect(stderr).toBe(""); + expect(exitCode).toBe(0); + }); +}); diff --git a/packages/sdk/src/agent-session.ts b/packages/sdk/src/agent-session.ts index b6d4bcf..3875bdb 100644 --- a/packages/sdk/src/agent-session.ts +++ b/packages/sdk/src/agent-session.ts @@ -5,9 +5,11 @@ import { DeepgramClient } from "@deepgram/sdk"; import type { AgentSessionConfig, ReconnectConfig } from "./types/config.js"; import type { AgentSessionEvents } from "./types/events.js"; import type { + AgentMessageBehavior, AgentContextMessage, - AgentSettingsObject, AgentV1SettingsPayload, + FunctionCallItem, + ListenSettings, ServerMessage, SpeakSettings, ThinkSettings, @@ -18,6 +20,19 @@ import { KeepAliveTimer } from "./connection/keepalive.js"; // Runtime type returned by client.agent.v1.connect() — actually WrappedAgentV1Socket // but TypeScript sees the base V1Socket interface which has all the methods we need. type V1Socket = Awaited["agent"]["v1"]["connect"]>>; +type RuntimeUpdate = + | Parameters[0] + | Parameters[0] + | Parameters[0] + | Parameters[0]; +type IncomingSocketItem = { + socket: V1Socket; + generation: number; + lifecycle: number; +} & ( + | { type: "message"; data: unknown } + | { type: "failure"; reason: string; error?: Error } +); const DEFAULT_KEEPALIVE_MS = 10_000; const OPEN_TIMEOUT_MS = 10_000; @@ -49,9 +64,9 @@ export type AgentState = * Key SDK insight: `client.agent.v1.connect()` returns a WrappedAgentV1Socket * with `startClosed: true` — `socket.connect()` must be called explicitly to * start the WebSocket. The wrapper also calls setupBinaryHandling() so - * `socket.on("message", cb)` receives both parsed JSON and raw ArrayBuffers. + * `socket.on("message", cb)` receives parsed JSON and binary audio Blobs. * - * Browser audio I/O (microphone, playback, VAD) is provided separately by + * Browser audio I/O (microphone and playback) is provided separately by * AgentMicrophone and AgentPlayer. */ export class AgentSession extends EventEmitter { @@ -61,13 +76,21 @@ export class AgentSession extends EventEmitter { private reconnectAttempts = 0; private reconnectTimer: ReturnType | null = null; private intentionalClose = false; + private connectionGeneration = 0; + private lifecycleGeneration = 0; /** Audio frames queued before SettingsApplied; flushed once the agent is ready */ private audioQueue: ArrayBuffer[] = []; + private incomingMessageQueue: IncomingSocketItem[] = []; + private processingIncomingMessages = false; + private incomingMessageProcessingGeneration = 0; + private socketFailureQueued: V1Socket | null = null; + private runtimeUpdates: RuntimeUpdate[] = []; private settingsApplied = false; private sessionId: string | null = null; - /** Conversation history — accumulated internally so reconnects can pass context */ + /** Conversation history used to restore inline agent configurations on reconnect. */ conversationHistory: AgentContextMessage[] = []; - + private pendingFunctionCalls = new Map(); + private completedFunctionCallIds = new Set(); private _state: AgentState = "idle"; @@ -82,7 +105,14 @@ export class AgentSession extends EventEmitter { ); this.keepAlive = new KeepAliveTimer( config.keepAliveInterval ?? DEFAULT_KEEPALIVE_MS, - () => this.socket?.sendKeepAlive({ type: "KeepAlive" }), + () => { + const socket = this.socket; + if (this._canWriteToSocket(socket)) { + this._writeToSocket(socket, () => { + socket.sendKeepAlive({ type: "KeepAlive" }); + }); + } + }, ); } @@ -91,9 +121,12 @@ export class AgentSession extends EventEmitter { // --------------------------------------------------------------------------- async connect(): Promise { + this.intentionalClose = true; + this._cleanup(); this.intentionalClose = false; this.reconnectAttempts = 0; - this.conversationHistory = []; + this.runtimeUpdates.length = 0; + this.clearConversationHistory(); await this._openConnection(); } @@ -104,47 +137,78 @@ export class AgentSession extends EventEmitter { this.emit("disconnected", "user requested disconnect"); } + /** Clears context accumulated for future inline-agent reconnects. */ + clearConversationHistory(): void { + this.conversationHistory.length = 0; + this.pendingFunctionCalls.clear(); + this.completedFunctionCallIds.clear(); + } + // --------------------------------------------------------------------------- // Public send helpers // --------------------------------------------------------------------------- sendAudio(data: ArrayBuffer): void { - if (!this.settingsApplied) { + const socket = this.socket; + if (!this.settingsApplied || !this._canWriteToSocket(socket)) { this.audioQueue.push(data); return; } - this.socket?.sendMedia(data); + if (!this._writeToSocket(socket, () => { + socket.sendMedia(data); + })) { + this.audioQueue.push(data); + } + } + + updateListen(listen: ListenSettings): void { + this._recordRuntimeUpdate({ type: "UpdateListen", listen }); } updateSpeak(speak: SpeakSettings | SpeakSettings[]): void { - // eslint-disable-next-line @typescript-eslint/no-explicit-any - this.socket?.sendUpdateSpeak({ type: "UpdateSpeak", speak } as any); + this._recordRuntimeUpdate({ type: "UpdateSpeak", speak }); } updateThink(think: ThinkSettings | ThinkSettings[]): void { - // eslint-disable-next-line @typescript-eslint/no-explicit-any - this.socket?.sendUpdateThink({ type: "UpdateThink", think } as any); + this._recordRuntimeUpdate({ type: "UpdateThink", think }); } updatePrompt(prompt: string): void { - this.socket?.sendUpdatePrompt({ type: "UpdatePrompt", prompt }); + this._recordRuntimeUpdate({ type: "UpdatePrompt", prompt }); } injectUserMessage(content: string): void { - this.socket?.sendInjectUserMessage({ type: "InjectUserMessage", content }); + const socket = this.socket; + if (!this._canWriteToSocket(socket)) return; + this._writeToSocket(socket, () => { + socket.sendInjectUserMessage({ type: "InjectUserMessage", content }); + }); } - injectAgentMessage(message: string): void { - this.socket?.sendInjectAgentMessage({ type: "InjectAgentMessage", message }); + injectAgentMessage(message: string, behavior?: AgentMessageBehavior): void { + const socket = this.socket; + if (!this._canWriteToSocket(socket)) return; + this._writeToSocket(socket, () => { + socket.sendInjectAgentMessage({ + type: "InjectAgentMessage", + message, + ...(behavior === undefined ? {} : { behavior }), + }); + }); } sendFunctionCallResponse(id: string | undefined, name: string, content: string): void { - // eslint-disable-next-line @typescript-eslint/no-explicit-any - this.socket?.sendFunctionCallResponse({ type: "FunctionCallResponse", id, name, content } as any); + this._recordFunctionCallResponse(id, name, content); + const socket = this.socket; + if (!this._canWriteToSocket(socket)) return; + this._writeToSocket(socket, () => { + // eslint-disable-next-line @typescript-eslint/no-explicit-any + socket.sendFunctionCallResponse({ type: "FunctionCallResponse", id, name, content } as any); + }); } /** - * Returns the session ID assigned by the server (available after Welcome). + * Returns the request ID assigned by the server (available after Welcome). * Returns null if not yet connected. */ getId(): string | null { @@ -156,11 +220,17 @@ export class AgentSession extends EventEmitter { // --------------------------------------------------------------------------- private async _openConnection(): Promise { + const generation = ++this.connectionGeneration; + let socket: V1Socket | null = null; + this.socketFailureQueued = null; + this.settingsApplied = false; + this.keepAlive.stop(); this._setState("connecting"); try { this.tokenFactory.invalidate(); const token = await this.tokenFactory.get(); + if (generation !== this.connectionGeneration || this.intentionalClose) return; // Build client with the right auth scheme: // - Custom URL (proxy like dx-api): always Bearer @@ -175,49 +245,33 @@ export class AgentSession extends EventEmitter { const authorization = isBearer ? `Bearer ${token}` : `Token ${token}`; // Returns a WrappedAgentV1Socket with startClosed:true — NOT yet connected. - // reconnectAttempts:1 so the SDK makes one attempt; we manage retries above. - const socket = await client.agent.v1.connect({ + // Disable transport retries; AgentSession owns reconnect so it can refresh + // auth and restore inline Settings/context before queued audio resumes. + socket = await client.agent.v1.connect({ Authorization: authorization, - reconnectAttempts: 1, + reconnectAttempts: 0, }); + if (generation !== this.connectionGeneration || this.intentionalClose) { + socket.close(); + return; + } + const connectingSocket = socket; - // Ensure binary audio frames arrive as ArrayBuffer, not Blob. - // ReconnectingWebSocket defaults to binaryType:"blob"; setting arraybuffer - // here before connect() means _ws.binaryType is set correctly on open. - socket.socket.binaryType = "arraybuffer"; - - // Wait for open before binding our full event handlers. - // socket.connect() starts the WebSocket (required because startClosed:true). - // Race against close/error/timeout so we never hang indefinitely. - await new Promise((resolve, reject) => { - const timer = setTimeout( - () => reject(new Error(`open timeout after ${OPEN_TIMEOUT_MS}ms`)), - OPEN_TIMEOUT_MS, - ); - socket.on("open", () => { - clearTimeout(timer); - resolve(); - }); - socket.on("close", (e: { code: number; reason?: string }) => { - clearTimeout(timer); - reject(new Error(`socket closed before open: code ${e.code} ${e.reason ?? ""}`)); - }); - socket.on("error", (err) => { - clearTimeout(timer); - reject(err); - }); - // Actually start the WebSocket connection - socket.connect(); - }); - + this.socket = connectingSocket; + const opened = this._bindSocketEvents(connectingSocket, generation); + connectingSocket.connect(); + await opened; - this.socket = socket; - this.reconnectAttempts = 0; + if (generation !== this.connectionGeneration || connectingSocket !== this.socket) return; this._setState("connected"); - this._bindSocketEvents(socket); } catch (err) { + if (socket && this.socket === socket) this.socket = null; + if (socket) { + try { socket.close(); } catch { /* ignore */ } + } + if (generation !== this.connectionGeneration || this.intentionalClose) return; this._onConnectionError(err instanceof Error ? err : new Error(String(err))); } } @@ -240,13 +294,19 @@ export class AgentSession extends EventEmitter { // eslint-disable-next-line @typescript-eslint/no-explicit-any agent: (() => { const agent = this.config.agent as any; - if (this.conversationHistory.length === 0) return agent; + if (typeof agent === "string" || this.conversationHistory.length === 0) return agent; // There is conversation history — pass it as context and strip greeting // so the server has context but doesn't replay the opening message. return { ...agent, greeting: undefined, - context: { messages: this.conversationHistory }, + context: { + ...agent.context, + messages: [ + ...(agent.context?.messages ?? []), + ...this.conversationHistory, + ], + }, }; })(), }; @@ -261,50 +321,301 @@ export class AgentSession extends EventEmitter { return payload; } - private _bindSocketEvents(socket: V1Socket): void { + private _bindSocketEvents(socket: V1Socket, generation: number): Promise { + let opened = false; + let timer: ReturnType; + + const openPromise = new Promise((resolve, reject) => { + timer = setTimeout( + () => reject(new Error(`open timeout after ${OPEN_TIMEOUT_MS}ms`)), + OPEN_TIMEOUT_MS, + ); + socket.on("open", () => { + opened = true; + clearTimeout(timer); + resolve(); + }); + + socket.on("close", (event: { code: number; reason?: string }) => { + const reason = `socket closed: ${event.code} ${event.reason ?? ""}`; + if (!opened) { + clearTimeout(timer); + reject(new Error(reason)); + return; + } + this._queueSocketFailure(socket, generation, reason); + }); + + socket.on("error", (err) => { + const error = err instanceof Error ? err : new Error(String(err)); + if (!opened) { + clearTimeout(timer); + reject(error); + return; + } + this._queueSocketFailure(socket, generation, error.message, error); + }); + }); + // WrappedAgentV1Socket.setupBinaryHandling() already replaces the base - // handleMessage with a binary-aware handler that passes ArrayBuffers - // through to eventHandlers.message. Using socket.on("message") therefore - // receives BOTH parsed JSON objects and raw ArrayBuffers — no raw socket - // manipulation needed. + // handleMessage with a binary-aware handler. SDK 5.9 normalizes binary + // payloads to Blob at runtime even though V1Socket.Response omits Blob. socket.on("message", (msg) => { - if (msg instanceof ArrayBuffer) { - this.emit("audio", msg); - } else { - // eslint-disable-next-line @typescript-eslint/no-explicit-any - const serverMsg = msg as unknown as ServerMessage; - this._dispatchMessage(serverMsg, socket); - } + if (socket !== this.socket || generation !== this.connectionGeneration) return; + this._queueIncomingMessage(msg, socket, generation); }); - socket.on("close", (event: { code: number; reason?: string }) => { - this.keepAlive.stop(); - if (!this.intentionalClose) { - this._scheduleReconnect(`socket closed: ${event.code} ${event.reason ?? ""}`); + return openPromise; + } + + private _queueIncomingMessage( + data: unknown, + socket: V1Socket, + generation: number, + ): void { + this.incomingMessageQueue.push({ + type: "message", + data, + socket, + generation, + lifecycle: this.lifecycleGeneration, + }); + void this._drainIncomingMessages(); + } + + private _queueSocketFailure( + socket: V1Socket, + generation: number, + reason: string, + error?: Error, + ): void { + if ( + this.intentionalClose || + socket !== this.socket || + generation !== this.connectionGeneration || + this.socketFailureQueued === socket + ) { + return; + } + + this.socketFailureQueued = socket; + this.settingsApplied = false; + this.sessionId = null; + this.keepAlive.stop(); + this.incomingMessageQueue.push({ + type: "failure", + reason, + error, + socket, + generation, + lifecycle: this.lifecycleGeneration, + }); + void this._drainIncomingMessages(); + } + + private async _drainIncomingMessages(): Promise { + if (this.processingIncomingMessages) return; + this.processingIncomingMessages = true; + const processingGeneration = this.incomingMessageProcessingGeneration; + + try { + while ( + processingGeneration === this.incomingMessageProcessingGeneration && + this.incomingMessageQueue.length > 0 + ) { + const item = this.incomingMessageQueue.shift()!; + const { socket, generation, lifecycle } = item; + if ( + generation !== this.connectionGeneration || + socket !== this.socket || + this.intentionalClose + ) { + continue; + } + + if (item.type === "failure") { + if (this.socketFailureQueued === socket) { + this.socketFailureQueued = null; + } + this._handleSocketFailure(socket, item.reason, item.error); + return; + } + + const { data } = item; + if ( + data instanceof ArrayBuffer || + (typeof Blob !== "undefined" && data instanceof Blob) + ) { + try { + const chunk = data instanceof ArrayBuffer ? data : await data.arrayBuffer(); + if ( + generation === this.connectionGeneration && + socket === this.socket && + !this.intentionalClose + ) { + this.emit("audio", chunk); + } + } catch (err) { + if (lifecycle === this.lifecycleGeneration && !this.intentionalClose) { + this.emit("sdk-error", err instanceof Error ? err : new Error(String(err))); + } + } + } else { + if ( + generation === this.connectionGeneration && + socket === this.socket && + !this.intentionalClose + ) { + this._dispatchMessage(data as ServerMessage, socket); + } + } + } + } finally { + if (processingGeneration === this.incomingMessageProcessingGeneration) { + this.processingIncomingMessages = false; + } + } + } + + private _recordRuntimeUpdate(update: RuntimeUpdate): void { + const snapshot = JSON.parse(JSON.stringify(update)) as RuntimeUpdate; + const existing = this.runtimeUpdates.findIndex( + (recorded) => recorded.type === snapshot.type, + ); + if (existing !== -1) this.runtimeUpdates.splice(existing, 1); + this.runtimeUpdates.push(snapshot); + const socket = this.socket; + if (this.settingsApplied && this._canWriteToSocket(socket)) { + this._sendRuntimeUpdate(socket, snapshot); + } + } + + private _canWriteToSocket(socket: V1Socket | null): socket is V1Socket { + return socket !== null && socket === this.socket && socket !== this.socketFailureQueued; + } + + private _writeToSocket(socket: V1Socket, write: () => void): boolean { + try { + write(); + return true; + } catch (err) { + const error = err instanceof Error ? err : new Error(String(err)); + this._queueSocketFailure( + socket, + this.connectionGeneration, + error.message, + error, + ); + return false; + } + } + + private _sendRuntimeUpdate(socket: V1Socket, update: RuntimeUpdate): boolean { + return this._writeToSocket(socket, () => { + switch (update.type) { + case "UpdateListen": + socket.sendUpdateListen(update); + break; + case "UpdateThink": + socket.sendUpdateThink(update); + break; + case "UpdateSpeak": + socket.sendUpdateSpeak(update); + break; + case "UpdatePrompt": + socket.sendUpdatePrompt(update); + break; } }); + } - socket.on("error", (err) => { - this.emit("sdk-error", err); + private _replayRuntimeUpdates(socket: V1Socket): boolean { + for (const update of this.runtimeUpdates) { + if (!this._sendRuntimeUpdate(socket, update)) return false; + } + return true; + } + + private _recordFunctionCallResponse( + id: string | undefined, + name: string, + response: string, + ): void { + const request = id + ? this.pendingFunctionCalls.get(id) + : this._findPendingFunctionCall(name); + if (!request || this.completedFunctionCallIds.has(request.id)) return; + + this.completedFunctionCallIds.add(request.id); + this.pendingFunctionCalls.delete(request.id); + this.conversationHistory.push({ + type: "History", + function_calls: [{ + id: request.id, + name: request.name, + client_side: request.client_side, + arguments: request.arguments, + response, + ...(request.thought_signature === undefined + ? {} + : { thought_signature: request.thought_signature }), + }], }); } + private _findPendingFunctionCall(name: string): FunctionCallItem | undefined { + const matches = [...this.pendingFunctionCalls.values()].filter( + (request) => request.name === name, + ); + return matches.length === 1 ? matches[0] : undefined; + } + + private _recordHistoryFunctionCalls(msg: ServerMessage): void { + if (msg.type !== "History" || !("function_calls" in msg)) return; + + const functionCalls = msg.function_calls.filter((functionCall) => { + if (this.completedFunctionCallIds.has(functionCall.id)) return false; + this.completedFunctionCallIds.add(functionCall.id); + this.pendingFunctionCalls.delete(functionCall.id); + return true; + }); + if (functionCalls.length > 0) { + this.conversationHistory.push({ type: "History", function_calls: functionCalls }); + } + } + private _dispatchMessage(msg: ServerMessage, socket: V1Socket): void { switch (msg.type) { case "Welcome": { - // eslint-disable-next-line @typescript-eslint/no-explicit-any - this.sessionId = (msg as any).session_id ?? null; + if (!this._canWriteToSocket(socket)) { + this.emit("welcome", msg); + break; + } + this.sessionId = msg.request_id ?? null; const settings = this._buildSettingsPayload(); - // eslint-disable-next-line @typescript-eslint/no-explicit-any - socket.sendSettings(settings as any); + if (!this._writeToSocket(socket, () => { + // eslint-disable-next-line @typescript-eslint/no-explicit-any + socket.sendSettings(settings as any); + })) return; this.emit("welcome", msg); break; } case "SettingsApplied": + if (!this._canWriteToSocket(socket)) { + this.emit("settings-applied", msg); + break; + } this.settingsApplied = true; + this.reconnectAttempts = 0; this.keepAlive.start(); - for (const frame of this.audioQueue) { - socket.sendMedia(frame); + if (!this._replayRuntimeUpdates(socket)) return; + for (let index = 0; index < this.audioQueue.length; index++) { + if (!this._writeToSocket(socket, () => { + socket.sendMedia(this.audioQueue[index]); + })) { + this.audioQueue = this.audioQueue.slice(index); + return; + } } this.audioQueue = []; this.emit("settings-applied", msg); @@ -324,6 +635,11 @@ export class AgentSession extends EventEmitter { this.emit("agent-thinking", msg); break; case "FunctionCallRequest": + for (const functionCall of msg.functions) { + if (!this.completedFunctionCallIds.has(functionCall.id)) { + this.pendingFunctionCalls.set(functionCall.id, functionCall); + } + } this.emit("function-call-request", msg); break; case "AgentStartedSpeaking": @@ -341,6 +657,16 @@ export class AgentSession extends EventEmitter { case "ThinkUpdated": this.emit("think-updated", msg); break; + case "ListenUpdated": + this.emit("listen-updated", msg); + break; + case "LatencyReport": + this.emit("latency-report", msg); + break; + case "History": + this._recordHistoryFunctionCalls(msg); + this.emit("history", msg); + break; case "InjectionRefused": this.emit("injection-refused", msg); break; @@ -351,17 +677,38 @@ export class AgentSession extends EventEmitter { this.emit("warning", msg); break; case "FunctionCallResponse": + this._recordFunctionCallResponse(msg.id, msg.name, msg.content); this.emit("function-call-response", msg); break; } } private _onConnectionError(err: Error): void { + this.settingsApplied = false; + this.sessionId = null; + this.keepAlive.stop(); this.emit("sdk-error", err); this._scheduleReconnect(err.message); } + private _handleSocketFailure(socket: V1Socket, reason: string, err?: Error): void { + if (this.intentionalClose || socket !== this.socket) return; + + this.settingsApplied = false; + this.sessionId = null; + this.keepAlive.stop(); + this.socket = null; + this.incomingMessageProcessingGeneration++; + this.incomingMessageQueue = []; + this.processingIncomingMessages = false; + if (err) this.emit("sdk-error", err); + try { socket.close(); } catch { /* ignore */ } + this._scheduleReconnect(reason); + } + private _scheduleReconnect(reason: string): void { + if (this.intentionalClose || this.reconnectTimer !== null) return; + const cfg = { ...DEFAULT_RECONNECT, ...this.config.reconnect }; if (!cfg.enabled || this.reconnectAttempts >= cfg.maxAttempts) { this._cleanup(); @@ -381,21 +728,30 @@ export class AgentSession extends EventEmitter { this.emit("reconnecting", this.reconnectAttempts, Math.round(delay)); this.reconnectTimer = setTimeout(async () => { + this.reconnectTimer = null; await this._openConnection(); }, delay); } private _cleanup(): void { + this.connectionGeneration++; + this.lifecycleGeneration++; this.settingsApplied = false; + this.sessionId = null; this.audioQueue = []; + this.incomingMessageQueue = []; + this.socketFailureQueued = null; + this.incomingMessageProcessingGeneration++; + this.processingIncomingMessages = false; this.keepAlive.stop(); if (this.reconnectTimer !== null) { clearTimeout(this.reconnectTimer); this.reconnectTimer = null; } - if (this.socket) { - try { this.socket.close(); } catch { /* ignore */ } - this.socket = null; + const socket = this.socket; + this.socket = null; + if (socket) { + try { socket.close(); } catch { /* ignore */ } } } diff --git a/packages/sdk/src/index.ts b/packages/sdk/src/index.ts index 0a8d604..c9e1366 100644 --- a/packages/sdk/src/index.ts +++ b/packages/sdk/src/index.ts @@ -15,6 +15,8 @@ export type { AgentSettingsObject, AgentContextMessage, AgentV1SettingsPayload, + ListenSettings, + AgentMessageBehavior, ThinkSettings, ThinkProvider, SpeakSettings, @@ -28,9 +30,16 @@ export type { UserStartedSpeakingMessage, AgentThinkingMessage, FunctionCallRequestMessage, + FunctionCallResponseMessage, FunctionCallItem, AgentStartedSpeakingMessage, AgentAudioDoneMessage, + PromptUpdatedMessage, + SpeakUpdatedMessage, + ThinkUpdatedMessage, + ListenUpdatedMessage, + LatencyReportMessage, + HistoryMessage, AgentErrorMessage, AgentWarningMessage, InjectionRefusedMessage, diff --git a/packages/sdk/src/types/config.ts b/packages/sdk/src/types/config.ts index 81e5919..fcf0ff5 100644 --- a/packages/sdk/src/types/config.ts +++ b/packages/sdk/src/types/config.ts @@ -45,7 +45,7 @@ export interface AudioOutputConfig { export interface ReconnectConfig { /** Whether to auto-reconnect on unexpected disconnections. Default: true */ enabled?: boolean; - /** Max number of attempts before giving up. Default: 8 */ + /** Max consecutive failed attempts before giving up. The counter resets once a connection reaches SettingsApplied. Default: 8 */ maxAttempts?: number; /** Initial backoff delay in ms. Default: 500 */ baseDelay?: number; @@ -65,6 +65,8 @@ export interface AgentSessionConfig { /** * Agent configuration sent in the Settings message after connection. * Pass either a full settings object or a pre-built agent UUID. + * Conversation context restoration on reconnect is available only for the + * inline object form; the protocol cannot attach context to an agent UUID. */ agent: AgentSettingsObject | string; diff --git a/packages/sdk/src/types/events.ts b/packages/sdk/src/types/events.ts index 0a4e199..98dc2a6 100644 --- a/packages/sdk/src/types/events.ts +++ b/packages/sdk/src/types/events.ts @@ -10,6 +10,9 @@ import type { PromptUpdatedMessage, SpeakUpdatedMessage, ThinkUpdatedMessage, + ListenUpdatedMessage, + LatencyReportMessage, + HistoryMessage, InjectionRefusedMessage, AgentErrorMessage, AgentWarningMessage, @@ -33,6 +36,9 @@ export interface AgentSessionEvents { "prompt-updated": [msg: PromptUpdatedMessage]; "speak-updated": [msg: SpeakUpdatedMessage]; "think-updated": [msg: ThinkUpdatedMessage]; + "listen-updated": [msg: ListenUpdatedMessage]; + "latency-report": [msg: LatencyReportMessage]; + history: [msg: HistoryMessage]; "injection-refused": [msg: InjectionRefusedMessage]; error: [msg: AgentErrorMessage]; warning: [msg: AgentWarningMessage]; diff --git a/packages/sdk/src/types/messages.ts b/packages/sdk/src/types/messages.ts index 7f9655e..a46d73e 100644 --- a/packages/sdk/src/types/messages.ts +++ b/packages/sdk/src/types/messages.ts @@ -12,6 +12,10 @@ type V1Socket = Awaited["agent"][ export type AgentV1SettingsPayload = Parameters[0]; export type AgentSettingsObject = AgentV1SettingsPayload["agent"]; +export type ListenSettings = Parameters[0]["listen"]; +export type AgentMessageBehavior = NonNullable< + Parameters[0]["behavior"] +>; // Derived from the settings payload — no hand-written types type AgentObject = Extract; @@ -41,6 +45,9 @@ export type AgentAudioDoneMessage = agent.AgentV1AgentAudioDone; export type PromptUpdatedMessage = agent.AgentV1PromptUpdated; export type SpeakUpdatedMessage = agent.AgentV1SpeakUpdated; export type ThinkUpdatedMessage = agent.AgentV1ThinkUpdated; +export type ListenUpdatedMessage = agent.AgentV1ListenUpdated; +export type LatencyReportMessage = agent.AgentV1LatencyReport; +export type HistoryMessage = agent.AgentV1History; export type InjectionRefusedMessage = agent.AgentV1InjectionRefused; export type AgentErrorMessage = agent.AgentV1Error; export type AgentWarningMessage = agent.AgentV1Warning; @@ -59,6 +66,9 @@ export type ServerMessage = | PromptUpdatedMessage | SpeakUpdatedMessage | ThinkUpdatedMessage + | ListenUpdatedMessage + | LatencyReportMessage + | HistoryMessage | InjectionRefusedMessage | AgentErrorMessage | AgentWarningMessage diff --git a/packages/sdk/vite.config.ts b/packages/sdk/vite.config.ts index 2eac61f..563612d 100644 --- a/packages/sdk/vite.config.ts +++ b/packages/sdk/vite.config.ts @@ -12,8 +12,6 @@ export default defineConfig({ external: [ "@deepgram/sdk", "eventemitter3", - "@ricky0123/vad-web", - "onnxruntime-web", ], }, minify: "terser", diff --git a/packages/widget/README.md b/packages/widget/README.md index 42c8fa5..21cb6df 100644 --- a/packages/widget/README.md +++ b/packages/widget/README.md @@ -17,7 +17,7 @@ import { init } from "@deepgram/agents-widget"; const destroy = init({ tokenFactory: () => fetch('/api/deepgram-token').then(r => r.text()), - agent: { think: { provider: { type: 'open_ai' }, model: 'gpt-4o-mini' } }, + agent: { think: { provider: { type: 'open_ai', model: 'gpt-4o-mini' } } }, }); // Later: destroy() to unmount @@ -30,7 +30,7 @@ A UMD bundle ships at `dist/widget.umd.js` for ` ``` diff --git a/packages/widget/package.json b/packages/widget/package.json index 8a2000d..640448e 100644 --- a/packages/widget/package.json +++ b/packages/widget/package.json @@ -59,7 +59,7 @@ "happy-dom": "20.9.0", "tailwindcss": "4.2.4", "terser": "5.46.2", - "typescript": "6.0.2", + "typescript": "5.9.3", "vite": "8.0.10", "vite-plugin-css-injected-by-js": "4.0.1", "vite-plugin-dts": "4.5.4" diff --git a/packages/widget/src/index.ts b/packages/widget/src/index.ts index 8bd4a30..8b2993c 100644 --- a/packages/widget/src/index.ts +++ b/packages/widget/src/index.ts @@ -23,7 +23,7 @@ export type { * * init({ * tokenFactory: () => fetch('/api/deepgram-token').then(r => r.text()), - * agent: { think: { type: 'open_ai', model: 'gpt-4o-mini' } }, + * agent: { think: { provider: { type: 'open_ai', model: 'gpt-4o-mini' } } }, * }); * ``` *