diff --git a/examples/nextjs-realtime/README.md b/examples/nextjs-realtime/README.md index ee8a60b..a3b622e 100644 --- a/examples/nextjs-realtime/README.md +++ b/examples/nextjs-realtime/README.md @@ -28,18 +28,39 @@ pnpm dev ## Features - Real-time webcam video transformation +- A fresh client token for every connect and reconnect (`apiKeyProvider`), minted on the server - Dynamic style prompt updates -- Connection state management +- Connection state display, including reconnects - Error handling ## How it works -1. The frontend requests a short-lived client token from `/api/realtime-token` -2. The backend uses `client.tokens.create()` to generate the token -3. The frontend uses the client token to connect to Decart's Realtime API -4. Webcam feed is captured and sent to Decart's realtime API -5. Transformed video is displayed side-by-side with the original -6. You can change the style prompt in real-time +Client tokens expire **60 seconds** after minting by default, so a token fetched on page load is +usually dead by the time the user has granted the camera and pressed Start. Instead of holding a +token, the page hands the SDK a function that fetches one: + +```ts +// components/video-stream.tsx +const client = createDecartClient({ + apiKeyProvider: async () => { + const response = await fetch("/api/realtime-token", { method: "POST" }); + const { apiKey } = await response.json(); + return apiKey; + }, +}); +``` + +1. Pressing **Start** captures the webcam and calls `client.realtime.connect(...)`. +2. The SDK calls `apiKeyProvider`, which `POST`s to `/api/realtime-token`. +3. The route mints a token with `client.tokens.create({ expiresIn: 60 })` using the permanent + `DECART_API_KEY`, which never leaves the server. +4. The SDK dials with that token. If the session drops, the SDK reconnects and calls + `apiKeyProvider` again, so a reconnect never reuses the token the session started with. +5. The transformed video is shown next to the original; the prompt can be changed live. + +If you prefer a static `apiKey`, mint it right before `connect()`. The SDK reads the token's +`exp` before dialling and rejects an expired one with `TOKEN_EXPIRED` instead of opening a +socket. ## Models diff --git a/examples/nextjs-realtime/app/api/realtime-token/route.ts b/examples/nextjs-realtime/app/api/realtime-token/route.ts index c36fd5a..5971f18 100644 --- a/examples/nextjs-realtime/app/api/realtime-token/route.ts +++ b/examples/nextjs-realtime/app/api/realtime-token/route.ts @@ -3,16 +3,22 @@ import { NextResponse } from "next/server"; const DECART_API_KEY = process.env.DECART_API_KEY; +/** + * Mints a short-lived client token with the permanent API key, which never leaves the server. + * The browser's `apiKeyProvider` calls this right before every realtime connect and reconnect. + */ export async function POST() { try { if (!DECART_API_KEY) { return NextResponse.json({ error: "DECART_API_KEY is not set" }, { status: 500 }); } - const client = createDecartClient({ - apiKey: DECART_API_KEY, + const client = createDecartClient({ apiKey: DECART_API_KEY }); + const token = await client.tokens.create({ + // Seconds until the token expires (1-3600, default 60). The SDK asks for a token right before + // each dial, so a short TTL is enough; raise it only if you mint ahead of connecting. + expiresIn: 60, }); - const token = await client.tokens.create(); return NextResponse.json(token); } catch (error) { diff --git a/examples/nextjs-realtime/components/video-stream.tsx b/examples/nextjs-realtime/components/video-stream.tsx index 2c84f5d..1a8f361 100644 --- a/examples/nextjs-realtime/components/video-stream.tsx +++ b/examples/nextjs-realtime/components/video-stream.tsx @@ -3,6 +3,29 @@ import { createDecartClient, type DecartSDKError, models, type RealTimeClient } from "@decartai/sdk"; import { useEffect, useRef, useState } from "react"; +const model = models.realtime("lucy-restyle-2"); + +/** + * Asks our backend for a fresh client token. The SDK calls this right before every connect and + * reconnect, so the token is minted for that dial: it cannot have expired while the page sat open + * or the user was deciding about camera permissions. (Client tokens live 60 s by default.) + */ +async function fetchClientToken(): Promise { + const response = await fetch("/api/realtime-token", { method: "POST" }); + if (!response.ok) throw new Error(`Token endpoint answered ${response.status}`); + const { apiKey } = await response.json(); + return apiKey; +} + +// One client for the page's lifetime. It holds no credential of its own; the provider supplies one per dial. +const client = createDecartClient({ apiKeyProvider: fetchClientToken }); + +function describeError(error: unknown): string { + if (error instanceof Error) return error.message; + const sdkError = error as Partial; + return sdkError?.code ? `${sdkError.code}: ${sdkError.message}` : String(error); +} + interface VideoStreamProps { prompt: string; } @@ -11,81 +34,71 @@ export function VideoStream({ prompt }: VideoStreamProps) { const inputRef = useRef(null); const outputRef = useRef(null); const realtimeClientRef = useRef(null); + const cameraRef = useRef(null); + // Bumped by release(), so a start() still waiting on the camera prompt or on connect() notices + // that Stop (or unmount) happened and lets go of what it was about to keep. + const attemptRef = useRef(0); const [status, setStatus] = useState("idle"); - - useEffect(() => { - let mounted = true; - - async function start() { - try { - const model = models.realtime("lucy-restyle-2"); - - setStatus("requesting camera..."); - const stream = await navigator.mediaDevices.getUserMedia({ - video: { - frameRate: model.fps, - width: model.width, - height: model.height, - }, - }); - - if (!mounted) return; - - if (inputRef.current) { - inputRef.current.srcObject = stream; - } - - // Fetch client token from our backend API - const tokenResponse = await fetch("/api/realtime-token", { - method: "POST", - }); - if (!tokenResponse.ok) { - throw new Error("Failed to get client token"); - } - const { apiKey } = await tokenResponse.json(); - - if (!mounted) return; - - setStatus("connecting..."); - - const client = createDecartClient({ apiKey }); - - const realtimeClient = await client.realtime.connect(stream, { - model, - onRemoteStream: (transformedStream: MediaStream) => { - if (outputRef.current) { - outputRef.current.srcObject = transformedStream; - } - }, - initialState: { - prompt: { text: prompt, enhance: true }, - }, - }); - - realtimeClientRef.current = realtimeClient; - - // Subscribe to events - realtimeClient.on("connectionChange", (state) => { - setStatus(state); - }); - - realtimeClient.on("error", (error: DecartSDKError) => { - setStatus(`error: ${error.message}`); - }); - } catch (error) { - setStatus(`error: ${error}`); + const [running, setRunning] = useState(false); + + function release() { + attemptRef.current++; + realtimeClientRef.current?.disconnect(); + realtimeClientRef.current = null; + for (const track of cameraRef.current?.getTracks() ?? []) track.stop(); + cameraRef.current = null; + } + + function stop() { + release(); + setRunning(false); + setStatus("idle"); + } + + async function start() { + const attempt = ++attemptRef.current; + const cancelled = () => attemptRef.current !== attempt; + setRunning(true); + try { + setStatus("requesting camera..."); + const camera = await navigator.mediaDevices.getUserMedia({ + video: { frameRate: model.fps, width: model.width, height: model.height }, + }); + if (cancelled()) { + for (const track of camera.getTracks()) track.stop(); + return; } + cameraRef.current = camera; + if (inputRef.current) inputRef.current.srcObject = camera; + + setStatus("connecting..."); + const realtimeClient = await client.realtime.connect(camera, { + model, + onRemoteStream: (transformedStream) => { + if (outputRef.current) outputRef.current.srcObject = transformedStream; + }, + onConnectionChange: setStatus, + initialState: { prompt: { text: prompt, enhance: true } }, + }); + if (cancelled()) { + realtimeClient.disconnect(); + return; + } + realtimeClientRef.current = realtimeClient; + + realtimeClient.on("error", (error) => setStatus(`error: ${describeError(error)}`)); + } catch (error) { + if (cancelled()) return; + release(); + setRunning(false); + setStatus(`error: ${describeError(error)}`); } + } - start(); - - return () => { - mounted = false; - realtimeClientRef.current?.disconnect(); - }; - }, []); + // Release the camera and the session when the component unmounts. + useEffect(() => release, []); - // Update prompt when it changes + // Update the prompt on the running session when it changes. useEffect(() => { if (realtimeClientRef.current?.isConnected()) { realtimeClientRef.current.setPrompt(prompt, { enhance: true }); @@ -94,7 +107,12 @@ export function VideoStream({ prompt }: VideoStreamProps) { return (
-

Status: {status}

+

+ + Status: {status} +

Input

diff --git a/packages/sdk/AGENTS.md b/packages/sdk/AGENTS.md index 99049d7..5fe132d 100644 --- a/packages/sdk/AGENTS.md +++ b/packages/sdk/AGENTS.md @@ -25,6 +25,7 @@ - `types.ts` - TypeScript types for queue operations - **src/realtime/** - LiveKit-backed real-time video streaming logic - `client.ts` - Real-time client implementation and public event surface + - `credential.ts` - Per-dial credential: `apiKeyProvider` and the client-token expiry preflight (TOKEN_EXPIRED before any dial) - `livekit-manager.ts` - LiveKit connection lifecycle and retry management - `livekit-connection.ts` - LiveKit room connection and control WebSocket handling - `methods.ts` - Realtime method implementations diff --git a/packages/sdk/README.md b/packages/sdk/README.md index 56ccfc9..2037381 100644 --- a/packages/sdk/README.md +++ b/packages/sdk/README.md @@ -347,6 +347,40 @@ clockTolerance }`. This is an offline check, not an API call. It needs WebCrypto runtimes without it (React Native) `verify` rejects with `UNSUPPORTED_PLATFORM_FEATURE`; `decodeClientToken` works everywhere. +#### Token lifetime: mint right before connecting + +Client tokens expire **60 seconds** after minting by default (`expiresIn`, 1-3600 s). A token +minted on page load or component mount is usually dead by the time the user has granted camera +access and pressed start, and a reconnect later with the same token is refused too. The SDK reads +the token's `exp` before every realtime dial (connect, connect retry and reconnect) and rejects an +expired one with `TOKEN_EXPIRED` instead of dialling: no socket is opened, and the message says how +long ago it expired. A few seconds of clock skew are tolerated. + +Two ways to keep the token fresh: + +- Mint it right before `connect()`, after `getUserMedia` and the user's click, and size + `expiresIn` to your connect window, e.g. `client.tokens.create({ expiresIn: 300 })`. +- Pass `apiKeyProvider` instead of (or alongside) `apiKey`. The SDK calls it before every connect + and reconnect and dials with what it returns, so each dial carries a token minted for it: + +```ts +const client = createDecartClient({ + apiKeyProvider: async () => { + // Your server: `client.tokens.create(...)` with your permanent API key. + const res = await fetch("/api/realtime-token", { method: "POST" }); + const { apiKey } = await res.json(); + return apiKey; + }, +}); + +const realtimeClient = await client.realtime.connect(stream, { model, onRemoteStream }); +``` + +If the provider rejects during `connect()`, `connect()` rejects with that error; during a +reconnect a provider failure is retried like any other dial failure. The provider serves the +realtime APIs (`connect` and `subscribe`); `process`, `queue`, `files` and `tokens` still use +`apiKey` or `proxy`. + ### React Native / Expo React Native realtime requires [LiveKit's React Native packages](https://github.com/livekit/client-sdk-react-native) diff --git a/packages/sdk/src/create-client.ts b/packages/sdk/src/create-client.ts index 472769e..a8f35fa 100644 --- a/packages/sdk/src/create-client.ts +++ b/packages/sdk/src/create-client.ts @@ -2,10 +2,11 @@ import { z } from "zod"; import { createFilesClient } from "./files/client"; import { createProcessClient } from "./process/client"; import { createQueueClient } from "./queue/client"; +import type { ApiKeyProvider } from "./realtime/credential"; import type { CreateRealtime } from "./realtime/factory"; import { createTokensClient } from "./tokens/client"; import { readEnv } from "./utils/env"; -import { createInvalidApiKeyError, createInvalidBaseUrlError } from "./utils/errors"; +import { createInvalidApiKeyError, createInvalidBaseUrlError, createSDKError, ERROR_CODES } from "./utils/errors"; import { createConsoleLogger, type Logger } from "./utils/logger"; // Schema with validation to ensure proxy and apiKey are mutually exclusive @@ -15,6 +16,7 @@ const proxySchema = z.union([z.string().url(), z.string().startsWith("/")]); const decartClientOptionsSchema = z .object({ apiKey: z.string().min(1).optional(), + apiKeyProvider: z.custom((val) => typeof val === "function").optional(), baseUrl: z.url().optional(), proxy: proxySchema.optional(), integration: z.string().optional(), @@ -39,6 +41,7 @@ export type DecartClientOptions = | { proxy: string; apiKey?: never; + apiKeyProvider?: ApiKeyProvider; baseUrl?: string; realtimeBaseUrl?: string; integration?: string; @@ -48,6 +51,7 @@ export type DecartClientOptions = | { proxy?: never; apiKey?: string; + apiKeyProvider?: ApiKeyProvider; baseUrl?: string; realtimeBaseUrl?: string; integration?: string; @@ -61,6 +65,9 @@ export type DecartClientOptions = * @param options - Configuration options * @param options.proxy - URL of the proxy server. When set, the client will use the proxy instead of direct API access and apiKey is not required. * @param options.apiKey - API key for authentication. + * @param options.apiKeyProvider - Called before every realtime connect and reconnect for a fresh credential (a client + * token's `apiKey`), so a token minted on page load is never dialled after it expired. Wins over `apiKey` for + * realtime; the HTTP APIs (`process`, `queue`, `files`, `tokens`) still use `apiKey` or `proxy`. * @param options.baseUrl - Override the default API base URL. * @param options.realtimeBaseUrl - Override the default WebSocket base URL for realtime connections. * @param options.integration - Optional integration identifier. @@ -77,6 +84,11 @@ export type DecartClientOptions = * * // Option 3: Using proxy (client-side, no API key needed) * const client = createDecartClient({ proxy: "https://your-server.com/api/decart" }); + * + * // Option 4: A fresh client token for every realtime connect and reconnect + * const client = createDecartClient({ + * apiKeyProvider: async () => (await (await fetch("/api/realtime-token", { method: "POST" })).json()).apiKey, + * }); * ``` */ export const createDecartClientForPlatform = (createRealtime: CreateRealtime, options: DecartClientOptions = {}) => { @@ -90,6 +102,13 @@ export const createDecartClientForPlatform = (createRealtime: CreateRealtime, op throw createInvalidApiKeyError(); } + if (issue.path.includes("apiKeyProvider")) { + throw createSDKError( + ERROR_CODES.INVALID_OPTIONS, + "apiKeyProvider must be a function returning a client token's apiKey (a string or a Promise of one).", + ); + } + if (issue.path.includes("baseUrl") || issue.path.includes("realtimeBaseUrl")) { const urlField = issue.path.includes("realtimeBaseUrl") ? "realtimeBaseUrl" : "baseUrl"; throw createInvalidBaseUrlError((options as Record)[urlField]); @@ -110,8 +129,9 @@ export const createDecartClientForPlatform = (createRealtime: CreateRealtime, op const apiKey = isProxyMode ? undefined : (("apiKey" in parsedOptions.data ? parsedOptions.data.apiKey : undefined) ?? readEnv("DECART_API_KEY")); + const { apiKeyProvider } = parsedOptions.data; - if (!isProxyMode && !apiKey) { + if (!isProxyMode && !apiKey && !apiKeyProvider) { throw createInvalidApiKeyError(); } @@ -136,6 +156,7 @@ export const createDecartClientForPlatform = (createRealtime: CreateRealtime, op publishBaseUrl: wsBaseUrl, subscribeBaseUrl, apiKey: apiKey || "", + apiKeyProvider, integration, logger, telemetryEnabled, @@ -265,9 +286,12 @@ export const createDecartClientForPlatform = (createRealtime: CreateRealtime, op * const claims = await serverClient.tokens.verify(token.token); * console.log(claims.serviceTier, claims.pool, claims.organizationId); * - * // Client-side: Use the client token + * // Client-side: Use the client token (mint it right before connecting: the default TTL is 60 s) * const client = createDecartClient({ apiKey: token.apiKey }); * const realtimeClient = await client.realtime.connect(stream, options); + * + * // Or let the SDK fetch a fresh token before every connect and reconnect + * const client = createDecartClient({ apiKeyProvider: () => fetchClientTokenFromYourServer() }); * ``` */ tokens, diff --git a/packages/sdk/src/public-api.ts b/packages/sdk/src/public-api.ts index 4505e6b..9beabb7 100644 --- a/packages/sdk/src/public-api.ts +++ b/packages/sdk/src/public-api.ts @@ -18,6 +18,7 @@ export type { RealTimeClientInitialState, RealtimeMediaStream, } from "./realtime/client"; +export type { ApiKeyProvider } from "./realtime/credential"; export type { SetInput } from "./realtime/methods"; export type { ConnectionQuality, diff --git a/packages/sdk/src/realtime/browser/index.ts b/packages/sdk/src/realtime/browser/index.ts index f7617b4..53e7598 100644 --- a/packages/sdk/src/realtime/browser/index.ts +++ b/packages/sdk/src/realtime/browser/index.ts @@ -9,6 +9,7 @@ export const createBrowserRealtime: CreateRealtime = (options) => { const publish = createRealTimeClient({ baseUrl: options.publishBaseUrl, apiKey: options.apiKey, + apiKeyProvider: options.apiKeyProvider, integration: options.integration, logger: options.logger, telemetryEnabled: options.telemetryEnabled, @@ -17,6 +18,7 @@ export const createBrowserRealtime: CreateRealtime = (options) => { const subscribe = createRealTimeSubscribeClient({ baseUrl: options.subscribeBaseUrl, apiKey: options.apiKey, + apiKeyProvider: options.apiKeyProvider, integration: options.integration, logger: options.logger, createFrameMetadataWorker, diff --git a/packages/sdk/src/realtime/client.ts b/packages/sdk/src/realtime/client.ts index 7c0bcf0..bef2e30 100644 --- a/packages/sdk/src/realtime/client.ts +++ b/packages/sdk/src/realtime/client.ts @@ -8,9 +8,10 @@ import { resolveFpsNumber, } from "../shared/model"; import { modelStateSchema } from "../shared/types"; -import { classifyWebrtcError, type DecartSDKError } from "../utils/errors"; +import { classifyWebrtcError, type DecartSDKError, DecartSDKException } from "../utils/errors"; import { createConsoleLogger, type Logger } from "../utils/logger"; import { imageToBase64 } from "../utils/media"; +import { type ApiKeyProvider, createCredentialSource } from "./credential"; import { createEventBuffer } from "./event-buffer"; import type { MediaChannelFactory, VideoCodec } from "./media-channel"; import { realtimeMethods, type SetInput } from "./methods"; @@ -31,6 +32,8 @@ import type { export type RealTimeClientOptions = { baseUrl: string; apiKey: string; + /** Called before every dial (connect, connect retry, reconnect) for a fresh credential; wins over `apiKey`. */ + apiKeyProvider?: ApiKeyProvider; integration?: string; logger: Logger; telemetryEnabled: boolean; @@ -164,8 +167,9 @@ export type RealTimeClient = { }; export const createRealTimeClient = (opts: RealTimeClientOptions) => { - const { baseUrl, apiKey, integration } = opts; + const { baseUrl, integration } = opts; const logger = opts.logger ?? createConsoleLogger("info"); + const nextCredential = createCredentialSource({ apiKey: opts.apiKey, apiKeyProvider: opts.apiKeyProvider }); const connect = async ( stream: MediaStream | null, @@ -179,6 +183,7 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { onConnectionQuality, onQueuePosition, initialState, + queryParams: extraQueryParams, resolution, speed, retries, @@ -192,8 +197,12 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { try { const initialImageRef = isFileRefId(initialState?.image) ? initialState.image : undefined; - const initialImage = - initialImageRef === undefined && initialState?.image ? await imageToBase64(initialState.image) : undefined; + // Before any setup: an expired client token is refused here, not by the server after a round + // trip. The provider round trip and the image encoding do not depend on each other. + const [credential, initialImage] = await Promise.all([ + nextCredential(), + initialImageRef === undefined && initialState?.image ? imageToBase64(initialState.image) : undefined, + ]); const initialPrompt = initialState?.prompt ? { text: initialState.prompt.text, enhance: initialState.prompt.enhance } : undefined; @@ -209,7 +218,7 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { logger, observability: { telemetryEnabled: opts.telemetryEnabled, - apiKey, + apiKey: credential, model: options.model.name, integration, logger, @@ -234,17 +243,29 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { ); } - const queryParams = new URLSearchParams({ - ...(preparedConnection.queryParams ?? {}), - ...(options.queryParams ?? {}), - api_key: apiKey, - model: options.model.name, - ...(resolution ? { resolution } : {}), - ...(speed ? { speed } : {}), - }); + // Captures only what a dial needs, not `options` (it lives as long as the session). + const preparedQueryParams = preparedConnection.queryParams; + const modelName = options.model.name; + let dialCredential = credential; + const dialUrl = (apiKey: string) => { + dialCredential = apiKey; + const queryParams = new URLSearchParams({ + ...preparedQueryParams, + ...extraQueryParams, + api_key: apiKey, + model: modelName, + ...(resolution ? { resolution } : {}), + ...(speed ? { speed } : {}), + }); + return `${url}?${queryParams.toString()}`; + }; session = new StreamSession({ - url: `${url}?${queryParams.toString()}`, + url: dialUrl(credential), + redialUrl: () => { + const next = nextCredential(); + return typeof next === "string" ? dialUrl(next) : next.then(dialUrl); + }, integration, observability, frameTiming: preparedConnection.frameTiming, @@ -277,7 +298,7 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { session.on("sessionStarted", ({ sessionId: id, subscribeToken: token }) => { sessionId = id; subscribeToken = token; - observability?.sessionStarted(id); + observability?.sessionStarted(id, dialCredential); }); session.on("generationTick", (e) => emitOrBuffer("generationTick", e)); @@ -330,7 +351,7 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { observability?.stop(); session?.disconnect(); preparedConnection?.dispose(); - throw error; + throw error instanceof DecartSDKException ? error.sdkError : error; } }; diff --git a/packages/sdk/src/realtime/config-realtime.ts b/packages/sdk/src/realtime/config-realtime.ts index 08c7b1c..4b21780 100644 --- a/packages/sdk/src/realtime/config-realtime.ts +++ b/packages/sdk/src/realtime/config-realtime.ts @@ -1,3 +1,5 @@ +import { ERROR_CODES } from "../utils/errors"; + /** Publish bitrate floor for the primary video layer (bps). */ const MIN_VIDEO_BITRATE_BPS = 1_100_000; /** Publish bitrate cap for the primary video layer (bps). */ @@ -46,6 +48,19 @@ export const REALTIME_CONFIG = { closeReason: "session limit reached", errorText: "concurrent session limit reached", }, + /** + * Client-token expiry preflight: a JWT credential whose `exp` is more than this many seconds + * in the past by the local clock is refused before the dial with TOKEN_EXPIRED, instead of by + * the server after a round trip. The allowance covers clock skew only: the server grants none, + * so a token this late would be refused anyway. + */ + clientTokenExpiryToleranceSeconds: 5, + /** + * SDK errors a dial can fail with before the socket opens that no re-dial would change: the + * credential itself was refused (see credential.ts). Any other SDK error an `apiKeyProvider` + * rejects with, such as a failed mint, is transient and retried. + */ + permanentErrorCodes: [ERROR_CODES.TOKEN_EXPIRED, ERROR_CODES.INVALID_API_KEY], permanentErrorSubstrings: [ "permission denied", "not allowed", diff --git a/packages/sdk/src/realtime/credential.ts b/packages/sdk/src/realtime/credential.ts new file mode 100644 index 0000000..5a1dec8 --- /dev/null +++ b/packages/sdk/src/realtime/credential.ts @@ -0,0 +1,54 @@ +import { decodeClientToken } from "../tokens/claims"; +import { createApiKeyProviderResultError, createClientTokenExpiredError } from "../utils/errors"; +import { REALTIME_CONFIG } from "./config-realtime"; + +/** + * Returns the credential for the next realtime dial: a client token's `apiKey` (minted on your + * server with `client.tokens.create`) or a permanent key. The SDK calls it right before every + * connect and reconnect, so each dial carries a token minted for it instead of one from page load. + * + * A rejection on `connect()` rejects the connect with the same error. During reconnects a + * rejection is retried like any other dial failure. + */ +export type ApiKeyProvider = () => string | Promise; + +/** + * Refuse a client token the server is certain to refuse: a JWT whose `exp` is already in the past, + * beyond the clock-skew tolerance. Opaque keys and JWTs that do not decode pass through for the + * server to judge. + * + * @throws `TOKEN_EXPIRED` + */ +export function assertClientTokenNotExpired(credential: string, now = Date.now()): void { + let expiresAt: string; + try { + ({ expiresAt } = decodeClientToken(credential)); + } catch { + return; + } + const expiredSecondsAgo = (now - Date.parse(expiresAt)) / 1000; + if (expiredSecondsAgo > REALTIME_CONFIG.session.clientTokenExpiryToleranceSeconds) { + throw createClientTokenExpiredError(Math.round(expiredSecondsAgo), expiresAt); + } +} + +/** + * The per-dial credential: the provider's fresh token when one is set, else the static key, checked + * for expiry either way. Synchronous for a static key, so the first socket still opens in the same + * tick as `connect()`. Provider rejections propagate as they are. + */ +export function createCredentialSource(options: { apiKey: string; apiKeyProvider?: ApiKeyProvider }): ApiKeyProvider { + const { apiKey, apiKeyProvider } = options; + if (!apiKeyProvider) { + return () => { + assertClientTokenNotExpired(apiKey); + return apiKey; + }; + } + return async () => { + const credential: unknown = await apiKeyProvider(); + if (typeof credential !== "string" || credential.length === 0) throw createApiKeyProviderResultError(credential); + assertClientTokenNotExpired(credential); + return credential; + }; +} diff --git a/packages/sdk/src/realtime/factory.ts b/packages/sdk/src/realtime/factory.ts index 0d8cc56..a6ac974 100644 --- a/packages/sdk/src/realtime/factory.ts +++ b/packages/sdk/src/realtime/factory.ts @@ -1,5 +1,6 @@ import type { Logger } from "../utils/logger"; import type { RealTimeClient, RealTimeClientConnectOptions, RealtimeMediaStream } from "./client"; +import type { ApiKeyProvider } from "./credential"; import type { CheckConnectivityOptions, ConnectivityReport } from "./preflight-types"; import type { RealTimeSubscribeClient, SubscribeOptions } from "./subscribe-client"; @@ -7,6 +8,7 @@ export type RealtimeFactoryOptions = { publishBaseUrl: string; subscribeBaseUrl: string; apiKey: string; + apiKeyProvider?: ApiKeyProvider; integration?: string; logger: Logger; telemetryEnabled: boolean; diff --git a/packages/sdk/src/realtime/observability/realtime-observability.ts b/packages/sdk/src/realtime/observability/realtime-observability.ts index e81d80d..21b8518 100644 --- a/packages/sdk/src/realtime/observability/realtime-observability.ts +++ b/packages/sdk/src/realtime/observability/realtime-observability.ts @@ -151,7 +151,8 @@ export class RealtimeObservability { ); } - sessionStarted(sessionId: string): void { + /** `apiKey` is the credential the session dialled with; a reconnect may carry a fresher one than connect did. */ + sessionStarted(sessionId: string, apiKey = this.options.apiKey): void { if (!this.options.telemetryEnabled) { return; } @@ -161,7 +162,7 @@ export class RealtimeObservability { } const reporter = new TelemetryReporter({ - apiKey: this.options.apiKey, + apiKey, sessionId, model: this.options.model, integration: this.options.integration, diff --git a/packages/sdk/src/realtime/react-native/index.ts b/packages/sdk/src/realtime/react-native/index.ts index 9b77fc1..f3245ff 100644 --- a/packages/sdk/src/realtime/react-native/index.ts +++ b/packages/sdk/src/realtime/react-native/index.ts @@ -12,6 +12,7 @@ export const createReactNativeRealtime: CreateRealtime = (options) => { const publish = createRealTimeClient({ baseUrl: options.publishBaseUrl, apiKey: options.apiKey, + apiKeyProvider: options.apiKeyProvider, integration: options.integration, logger: options.logger, telemetryEnabled: options.telemetryEnabled, @@ -20,6 +21,7 @@ export const createReactNativeRealtime: CreateRealtime = (options) => { const subscribe = createRealTimeSubscribeClient({ baseUrl: options.subscribeBaseUrl, apiKey: options.apiKey, + apiKeyProvider: options.apiKeyProvider, integration: options.integration, logger: options.logger, }); diff --git a/packages/sdk/src/realtime/signaling-channel.ts b/packages/sdk/src/realtime/signaling-channel.ts index 90ecfbf..912209d 100644 --- a/packages/sdk/src/realtime/signaling-channel.ts +++ b/packages/sdk/src/realtime/signaling-channel.ts @@ -40,7 +40,8 @@ export type SignalingChannelEvents = { }; export interface SignalingChannelConfig { - url: string; + /** The signaling URL, or a function resolving it right before each socket open (a fresh credential per dial). */ + url: string | (() => string | Promise); integration?: string; logger?: Logger; observability?: RealtimeObservability; @@ -142,8 +143,10 @@ export class SignalingChannel { const connectTimeout = opts.connectTimeout ?? REALTIME_CONFIG.signaling.connectTimeoutMs; const handshakeTimeout = opts.handshakeTimeout ?? REALTIME_CONFIG.signaling.handshakeTimeoutMs; + const dial = typeof this.config.url === "function" ? this.config.url() : this.config.url; + const url = typeof dial === "string" ? dial : await dial; this.config.observability?.startPhase("websocket-open"); - await this.openSocket(connectTimeout); + await this.openSocket(url, connectTimeout); this.config.observability?.endPhase("websocket-open", { success: true }); this.config.observability?.startPhase("room-join"); @@ -249,10 +252,10 @@ export class SignalingChannel { if (!ack.success) throw new Error(ack.error ?? "Failed to send image"); } - private async openSocket(timeout: number): Promise { + private async openSocket(url: string, timeout: number): Promise { const userAgent = encodeURIComponent(buildUserAgent(this.config.integration)); - const separator = this.config.url.includes("?") ? "&" : "?"; - const wsUrl = `${this.config.url}${separator}user_agent=${userAgent}`; + const separator = url.includes("?") ? "&" : "?"; + const wsUrl = `${url}${separator}user_agent=${userAgent}`; this.closing = false; await new Promise((resolve, reject) => { diff --git a/packages/sdk/src/realtime/stream-session.ts b/packages/sdk/src/realtime/stream-session.ts index b424c33..05ca691 100644 --- a/packages/sdk/src/realtime/stream-session.ts +++ b/packages/sdk/src/realtime/stream-session.ts @@ -1,6 +1,7 @@ import mitt, { type Emitter } from "mitt"; import pRetry, { AbortError } from "p-retry"; +import { DecartSDKException, isDecartSDKError } from "../utils/errors"; import { createConsoleLogger, type Logger } from "../utils/logger"; import { REALTIME_CONFIG } from "./config-realtime"; import type { MediaChannel, MediaChannelFactory, VideoCodec } from "./media-channel"; @@ -68,6 +69,16 @@ function terminalReasonFromError(error: unknown): SessionEndReason | null { return null; } +/** + * p-retry only accepts `Error` instances: a plain SDK error thrown inside a dial (an expired client + * token, an `apiKeyProvider` rejection) is carried across the retry loop as a `DecartSDKException`. + */ +function toDialError(thrown: unknown): Error { + if (thrown instanceof Error) return thrown; + if (isDecartSDKError(thrown)) return new DecartSDKException(thrown); + return new Error(String(thrown)); +} + export function encodeSubscribeToken(roomName: string, options: { frameTiming?: boolean } = {}): string { return btoa(JSON.stringify({ room_name: roomName, ...(options.frameTiming ? { frame_timing: true } : {}) })); } @@ -93,7 +104,14 @@ type StreamSessionEvents = { }; interface StreamSessionConfig { + /** Signaling URL of the first dial. */ url: string; + /** + * Signaling URL for every later dial (connect retries and reconnects), resolved right before the + * socket opens so it can carry a freshly minted client token. Defaults to `url`. An SDK error + * thrown here whose code is in `permanentErrorCodes` ends the dial attempts. + */ + redialUrl?: () => string | Promise; integration?: string; observability?: RealtimeObservability; frameTiming?: boolean; @@ -118,6 +136,7 @@ export class StreamSession { private queue: QueuePosition | null = null; private disposed = false; + private dialed = false; private currentAttempt = 0; private teardownGeneration = 0; @@ -234,6 +253,14 @@ export class StreamSession { this.logger.error("realtime connect: session refused, not retrying", { reason: terminal }); return false; } + const sdkCode = error instanceof DecartSDKException ? error.sdkError.code : undefined; + if (sdkCode && (REALTIME_CONFIG.session.permanentErrorCodes as readonly string[]).includes(sdkCode)) { + this.logger.error("realtime connect: credential refused, not retrying", { + code: sdkCode, + error: error.message, + }); + return false; + } const msg = error.message.toLowerCase(); const permanent = REALTIME_CONFIG.session.permanentErrorSubstrings.some((err) => msg.includes(err)); if (permanent) { @@ -294,11 +321,9 @@ export class StreamSession { sessionId: roomInfo.sessionId, subscribeToken: encodeSubscribeToken(roomInfo.roomName, { frameTiming: this.config.frameTiming }), }); - } catch (error) { - this.config.observability?.finishConnectionBreakdown({ - success: false, - error: error instanceof Error ? error.message : String(error), - }); + } catch (thrown) { + const error = toDialError(thrown); + this.config.observability?.finishConnectionBreakdown({ success: false, error: error.message }); throw error; } } @@ -444,9 +469,24 @@ export class StreamSession { }); } + /** The first dial uses the connect-time URL; every later one asks `redialUrl`. */ + private resolveDialUrl(): string | Promise { + const redial = this.dialed ? this.config.redialUrl : undefined; + this.dialed = true; + if (!redial) return this.config.url; + const url = redial(); + if (typeof url === "string") return url; + const attempt = this.currentAttempt; + return url.then((resolved) => { + // Disconnected or superseded while the credential was being fetched: do not dial at all. + if (this.disposed || this.currentAttempt !== attempt) throw new AbortError("Stale connect attempt"); + return resolved; + }); + } + private createTransport(): void { this.signaling = new SignalingChannel({ - url: this.config.url, + url: () => this.resolveDialUrl(), integration: this.config.integration, logger: this.logger, observability: this.config.observability, diff --git a/packages/sdk/src/realtime/subscribe-client.ts b/packages/sdk/src/realtime/subscribe-client.ts index 3d66700..911fa78 100644 --- a/packages/sdk/src/realtime/subscribe-client.ts +++ b/packages/sdk/src/realtime/subscribe-client.ts @@ -1,8 +1,15 @@ import type { RemoteParticipant, RemoteTrack, Room } from "livekit-client"; -import { classifyWebrtcError, createSDKError, type DecartSDKError, ERROR_CODES } from "../utils/errors"; +import { + classifyWebrtcError, + createSDKError, + type DecartSDKError, + ERROR_CODES, + isDecartSDKError, +} from "../utils/errors"; import { createConsoleLogger, type Logger } from "../utils/logger"; import { REALTIME_CONFIG } from "./config-realtime"; +import { type ApiKeyProvider, createCredentialSource } from "./credential"; import { createEventBuffer } from "./event-buffer"; import { loadLiveKitClient } from "./livekit"; import type { DiagnosticEvent } from "./observability/diagnostics"; @@ -26,15 +33,6 @@ type WatchStreamCredentialsRequest = { roomName: string; }; -function isDecartSDKError(error: unknown): error is DecartSDKError { - return ( - typeof error === "object" && - error !== null && - typeof (error as Partial).code === "string" && - typeof (error as Partial).message === "string" - ); -} - function createFrameMetadataSubscribeUnsupportedError(detail?: string): DecartSDKError { const suffix = detail ? `: ${detail}` : ""; return createSDKError( @@ -79,6 +77,8 @@ export type SubscribeOptions = { export type RealTimeSubscribeClientOptions = { baseUrl: string; apiKey: string; + /** Called before each subscribe for a fresh credential; wins over `apiKey`. */ + apiKeyProvider?: ApiKeyProvider; integration?: string; logger: Logger; createFrameMetadataWorker?: () => Worker; @@ -125,8 +125,9 @@ async function fetchWatchStreamCredentials(opts: WatchStreamCredentialsRequest): } export const createRealTimeSubscribeClient = (opts: RealTimeSubscribeClientOptions) => { - const { baseUrl, apiKey, integration } = opts; + const { baseUrl, integration } = opts; const logger = opts.logger ?? createConsoleLogger("info"); + const nextCredential = createCredentialSource({ apiKey: opts.apiKey, apiKeyProvider: opts.apiKeyProvider }); const subscribe = async (options: SubscribeOptions): Promise => { const { room_name: roomName, frame_timing: frameTiming } = decodeSubscribeToken(options.token); @@ -145,6 +146,10 @@ export const createRealTimeSubscribeClient = (opts: RealTimeSubscribeClientOptio emitOrBuffer("connectionChange", state); }; + // Outside the try: an expired client token is refused before anything is loaded or fetched, and + // a provider rejection reaches the caller unchanged, as it does from `connect()`. + const apiKey = await nextCredential(); + try { const { Room: LiveKitRoom, RoomEvent } = await loadLiveKitClient(); observability = new RealtimeObservability({ diff --git a/packages/sdk/src/utils/errors.ts b/packages/sdk/src/utils/errors.ts index dc1d785..c4725c6 100644 --- a/packages/sdk/src/utils/errors.ts +++ b/packages/sdk/src/utils/errors.ts @@ -46,10 +46,55 @@ export function createSDKError( return { code, message, data, cause }; } +export function isDecartSDKError(error: unknown): error is DecartSDKError { + return ( + typeof error === "object" && + error !== null && + typeof (error as Partial).code === "string" && + typeof (error as Partial).message === "string" + ); +} + +/** + * An `Error` carrying a {@link DecartSDKError} across the realtime session's retry loop, which + * (like p-retry) only accepts real errors. `connect()` and `classifyWebrtcError` unwrap it, so + * callers still receive the plain SDK error. + */ +export class DecartSDKException extends Error { + readonly sdkError: DecartSDKError; + + constructor(sdkError: DecartSDKError) { + super(sdkError.message); + this.name = "DecartSDKException"; + this.sdkError = sdkError; + if (sdkError.cause) this.cause = sdkError.cause; + } +} + export function createInvalidApiKeyError(): DecartSDKError { return createSDKError( ERROR_CODES.INVALID_API_KEY, - "Missing API key. Pass `apiKey` to createDecartClient() or set the DECART_API_KEY environment variable.", + "Missing API key. Pass `apiKey` or `apiKeyProvider` to createDecartClient(), or set the DECART_API_KEY environment variable.", + ); +} + +/** The platform's default client-token TTL (`tokens.create({ expiresIn })`), quoted in the expiry message. */ +const CLIENT_TOKEN_DEFAULT_TTL_SECONDS = 60; + +/** An expired client token caught by the realtime preflight, before any dial. */ +export function createClientTokenExpiredError(expiredSecondsAgo: number, expiresAt: string): DecartSDKError { + return createSDKError( + ERROR_CODES.TOKEN_EXPIRED, + `Client token expired ${expiredSecondsAgo} s ago (exp ${expiresAt}). Mint a new one right before connecting, or pass apiKeyProvider to createDecartClient so the SDK fetches a fresh token before every connect and reconnect. Client tokens expire ${CLIENT_TOKEN_DEFAULT_TTL_SECONDS} s after minting by default (tokens.create({ expiresIn })).`, + { claim: "exp", expiresAt, expiredSecondsAgo }, + ); +} + +export function createApiKeyProviderResultError(received: unknown): DecartSDKError { + const got = received === "" ? "an empty string" : typeof received; + return createSDKError( + ERROR_CODES.INVALID_API_KEY, + `apiKeyProvider must resolve to a non-empty API key string (a client token's apiKey); got ${got}.`, ); } @@ -106,6 +151,9 @@ export function createWebrtcSignalingError(error: Error): DecartSDKError { * Classify a raw WebRTC error into a specific SDK error based on its message. */ export function classifyWebrtcError(error: Error): DecartSDKError { + // Already judged by the SDK (e.g. an expired client token on reconnect): surface it unchanged. + if (error instanceof DecartSDKException) return error.sdkError; + const msg = error.message.toLowerCase(); const source = (error as Error & { source?: string }).source; diff --git a/packages/sdk/tests/credential.unit.test.ts b/packages/sdk/tests/credential.unit.test.ts new file mode 100644 index 0000000..6c7be48 --- /dev/null +++ b/packages/sdk/tests/credential.unit.test.ts @@ -0,0 +1,99 @@ +import { describe, expect, it, vi } from "vitest"; + +import { REALTIME_CONFIG } from "../src/realtime/config-realtime.js"; +import { assertClientTokenNotExpired, createCredentialSource } from "../src/realtime/credential.js"; +import { ERROR_CODES } from "../src/utils/errors.js"; +import { clientTokenJwt, expAt } from "./helpers/client-token.js"; + +const NOW = Date.UTC(2026, 9, 7, 12, 0, 0); +const expIn = (seconds: number) => expAt(NOW, seconds); +const tolerance = REALTIME_CONFIG.session.clientTokenExpiryToleranceSeconds; + +describe("assertClientTokenNotExpired", () => { + it("passes opaque keys and JWTs it cannot decode through to the server", () => { + expect(() => assertClientTokenNotExpired("ek_opaque", NOW)).not.toThrow(); + expect(() => assertClientTokenNotExpired("dct_permanent", NOW)).not.toThrow(); + expect(() => assertClientTokenNotExpired("eyJ.not-base64-json.sig", NOW)).not.toThrow(); + // A JWT with no `exp` is malformed for a client token; the server judges it. + expect(() => assertClientTokenNotExpired(clientTokenJwt({ exp: undefined }), NOW)).not.toThrow(); + }); + + it("passes a live token and one within the clock-skew tolerance", () => { + expect(() => assertClientTokenNotExpired(clientTokenJwt({ exp: expIn(300) }), NOW)).not.toThrow(); + expect(() => assertClientTokenNotExpired(clientTokenJwt({ exp: expIn(0) }), NOW)).not.toThrow(); + expect(() => assertClientTokenNotExpired(clientTokenJwt({ exp: expIn(-tolerance) }), NOW)).not.toThrow(); + }); + + it("rejects a token past the tolerance with TOKEN_EXPIRED and says how late it is", () => { + const exp = expIn(-(tolerance + 1)); + + expect(() => assertClientTokenNotExpired(clientTokenJwt({ exp }), NOW)).toThrow( + expect.objectContaining({ + code: ERROR_CODES.TOKEN_EXPIRED, + message: expect.stringMatching(/^Client token expired 6 s ago .*apiKeyProvider.*60 s/), + data: { claim: "exp", expiresAt: new Date(exp * 1000).toISOString(), expiredSecondsAgo: tolerance + 1 }, + }), + ); + }); + + it("reports whole seconds for tokens hours late", () => { + expect(() => assertClientTokenNotExpired(clientTokenJwt({ exp: expIn(-7200) }), NOW + 400)).toThrow( + expect.objectContaining({ data: expect.objectContaining({ expiredSecondsAgo: 7200 }) }), + ); + }); +}); + +describe("createCredentialSource", () => { + it("returns the static key synchronously on every call without a provider", () => { + const next = createCredentialSource({ apiKey: "ek_static" }); + expect(next()).toBe("ek_static"); + expect(next()).toBe("ek_static"); + }); + + it("asks the provider on every call and prefers it over the static key", async () => { + const provider = vi.fn().mockResolvedValueOnce("fresh-1").mockResolvedValueOnce("fresh-2"); + const next = createCredentialSource({ apiKey: "stale", apiKeyProvider: provider }); + + await expect(next()).resolves.toBe("fresh-1"); + await expect(next()).resolves.toBe("fresh-2"); + expect(provider).toHaveBeenCalledTimes(2); + }); + + it("accepts a synchronous provider", async () => { + const next = createCredentialSource({ apiKey: "", apiKeyProvider: () => "sync-token" }); + await expect(next()).resolves.toBe("sync-token"); + }); + + it("checks expiry on provider tokens and static tokens alike", async () => { + const expired = clientTokenJwt({ exp: expAt(Date.now(), -120) }); + expect(() => createCredentialSource({ apiKey: expired })()).toThrow( + expect.objectContaining({ code: ERROR_CODES.TOKEN_EXPIRED }), + ); + await expect(createCredentialSource({ apiKey: "", apiKeyProvider: async () => expired })()).rejects.toMatchObject({ + code: ERROR_CODES.TOKEN_EXPIRED, + }); + }); + + it("rejects a provider result that is not a non-empty string", async () => { + for (const result of ["", undefined, 42, { apiKey: "eyJ" }] as unknown[]) { + const next = createCredentialSource({ apiKey: "", apiKeyProvider: async () => result as string }); + await expect(next()).rejects.toMatchObject({ code: ERROR_CODES.INVALID_API_KEY }); + } + const next = createCredentialSource({ + apiKey: "", + apiKeyProvider: async () => ({ apiKey: "x" }) as unknown as string, + }); + await expect(next()).rejects.toMatchObject({ message: expect.stringContaining("got object") }); + }); + + it("lets a provider rejection through unchanged, whatever its shape", async () => { + for (const rejection of [ + new TypeError("Failed to fetch"), + "endpoint down", + { code: "TOKEN_CREATE_ERROR", message: "503" }, + ]) { + const next = createCredentialSource({ apiKey: "", apiKeyProvider: () => Promise.reject(rejection) }); + await expect(next()).rejects.toBe(rejection); + } + }); +}); diff --git a/packages/sdk/tests/helpers/client-token.ts b/packages/sdk/tests/helpers/client-token.ts new file mode 100644 index 0000000..8dee347 --- /dev/null +++ b/packages/sdk/tests/helpers/client-token.ts @@ -0,0 +1,11 @@ +const base64url = (text: string) => btoa(text).replace(/\+/g, "-").replace(/\//g, "_").replace(/=+$/, ""); + +/** + * An unsigned JWT in the platform's client-token shape (`apiKey` and `token` both carry it). Only + * the payload matters to the SDK's expiry preflight; `sub` is required by the claims mapper. + */ +export const clientTokenJwt = (claims: Record) => + `eyJhbGciOiJFZERTQSJ9.${base64url(JSON.stringify({ sub: "user_1", ...claims }))}.sig`; + +/** An `exp` claim `seconds` away from `now` (milliseconds). */ +export const expAt = (now: number, seconds: number) => Math.floor(now / 1000) + seconds; diff --git a/packages/sdk/tests/realtime.unit.test.ts b/packages/sdk/tests/realtime.unit.test.ts index 665222b..aff6c44 100644 --- a/packages/sdk/tests/realtime.unit.test.ts +++ b/packages/sdk/tests/realtime.unit.test.ts @@ -1,9 +1,11 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; -import { models } from "../src/index.js"; +import { createDecartClient, ERROR_CODES, models } from "../src/index.js"; import { prepareBrowserConnection } from "../src/realtime/browser/prepare-connection.js"; import { REALTIME_CONFIG } from "../src/realtime/config-realtime.js"; import { createLiveKitMediaChannel, type MediaChannel } from "../src/realtime/media-channel.js"; import type { ServerError } from "../src/realtime/types.js"; +import type { DecartSDKError } from "../src/utils/errors.js"; +import { clientTokenJwt, expAt } from "./helpers/client-token.js"; const liveKitMock = vi.hoisted(() => { const roomInstances: MockRoom[] = []; @@ -403,6 +405,64 @@ describe("realtime.subscribe", () => { client.disconnect(); }); + it("asks apiKeyProvider for the watch-stream credential instead of the static key", async () => { + const { createRealTimeSubscribeClient } = await import("../src/realtime/subscribe-client.js"); + const apiKeyProvider = vi.fn(async () => "fresh-viewer-token"); + const subscriber = createRealTimeSubscribeClient({ + baseUrl: "https://api.example.test", + apiKey: "stale-key", + apiKeyProvider, + logger, + }); + + const client = await subscriber.subscribe({ + token: btoa(JSON.stringify({ room_name: "room-1" })), + onRemoteStream: () => {}, + }); + + expect(apiKeyProvider).toHaveBeenCalledTimes(1); + const [, init] = (fetch as unknown as ReturnType).mock.calls[0] as [string, RequestInit]; + expect((init.headers as Record)["x-api-key"]).toBe("fresh-viewer-token"); + client.disconnect(); + }); + + it("refuses an expired client token before fetching watch-stream credentials", async () => { + const { createRealTimeSubscribeClient } = await import("../src/realtime/subscribe-client.js"); + const now = Date.UTC(2026, 9, 7, 12, 0, 0); + const nowSpy = vi.spyOn(Date, "now").mockReturnValue(now); + const subscriber = createRealTimeSubscribeClient({ + baseUrl: "https://api.example.test", + apiKey: clientTokenJwt({ exp: expAt(now, -600) }), + logger, + }); + + try { + await expect( + subscriber.subscribe({ token: btoa(JSON.stringify({ room_name: "room-1" })), onRemoteStream: () => {} }), + ).rejects.toMatchObject({ code: ERROR_CODES.TOKEN_EXPIRED, data: { expiredSecondsAgo: 600 } }); + } finally { + nowSpy.mockRestore(); + } + expect(fetch).not.toHaveBeenCalled(); + expect(liveKitMock.roomInstances).toHaveLength(0); + }); + + it("lets an apiKeyProvider rejection reach the subscriber unchanged", async () => { + const { createRealTimeSubscribeClient } = await import("../src/realtime/subscribe-client.js"); + const failure = new Error("token endpoint answered 503"); + const subscriber = createRealTimeSubscribeClient({ + baseUrl: "https://api.example.test", + apiKey: "", + apiKeyProvider: () => Promise.reject(failure), + logger, + }); + + await expect( + subscriber.subscribe({ token: btoa(JSON.stringify({ room_name: "room-1" })), onRemoteStream: () => {} }), + ).rejects.toBe(failure); + expect(fetch).not.toHaveBeenCalled(); + }); + it("terminates the frame-metadata worker when subscribe connect fails", async () => { const { createRealTimeSubscribeClient } = await import("../src/realtime/subscribe-client.js"); const worker = { terminate: vi.fn() } as unknown as Worker; @@ -540,6 +600,8 @@ describe("realtime.connect options", () => { afterEach(() => { vi.unstubAllGlobals(); + vi.restoreAllMocks(); + vi.useRealTimers(); }); it("adds resolution to the realtime URL when provided", async () => { @@ -758,6 +820,189 @@ describe("realtime.connect options", () => { ).rejects.toThrow(); expect(FakeWebSocket.instances).toHaveLength(0); }); + + describe("client token expiry preflight and apiKeyProvider", () => { + const NOW = Date.UTC(2026, 9, 7, 12, 0, 0); + const expIn = (seconds: number) => expAt(NOW, seconds); + const apiKeyOf = (ws: FakeWebSocket) => new URL(ws.url).searchParams.get("api_key"); + + const createClient = async (credentials: { apiKey: string; apiKeyProvider?: () => string | Promise }) => { + const { createRealTimeClient } = await import("../src/realtime/client.js"); + return createRealTimeClient({ + baseUrl: "wss://api3.decart.ai", + ...credentials, + logger, + telemetryEnabled: false, + prepareConnection: prepareBrowserConnection, + }); + }; + + it("rejects an expired client token with TOKEN_EXPIRED before opening a socket", async () => { + vi.spyOn(Date, "now").mockReturnValue(NOW); + const exp = expIn(-83); + const client = await createClient({ apiKey: clientTokenJwt({ exp }) }); + const onConnectionChange = vi.fn(); + + await expect( + client.connect(null, { model: models.realtime("lucy-2.5"), onRemoteStream: vi.fn(), onConnectionChange }), + ).rejects.toMatchObject({ + code: ERROR_CODES.TOKEN_EXPIRED, + message: expect.stringMatching(/^Client token expired 83 s ago \(exp 2026-10-07T11:58:37\.000Z\)\./), + data: { claim: "exp", expiresAt: new Date(exp * 1000).toISOString(), expiredSecondsAgo: 83 }, + }); + // Nothing was dialled and no state was reported: the token never left the client. + expect(FakeWebSocket.instances).toHaveLength(0); + expect(onConnectionChange).not.toHaveBeenCalled(); + }); + + it("connects with a token inside the clock-skew tolerance, a live token and an opaque key", async () => { + vi.spyOn(Date, "now").mockReturnValue(NOW); + const skewed = clientTokenJwt({ exp: expIn(-REALTIME_CONFIG.session.clientTokenExpiryToleranceSeconds) }); + const live = clientTokenJwt({ exp: expIn(60) }); + + for (const apiKey of [skewed, live, "ek_opaque_key"]) { + FakeWebSocket.instances = []; + const client = await createClient({ apiKey }); + const realtimeClient = await client.connect(null, { + model: models.realtime("lucy-2.5"), + onRemoteStream: vi.fn(), + }); + expect(FakeWebSocket.instances).toHaveLength(1); + expect(apiKeyOf(FakeWebSocket.instances[0])).toBe(apiKey); + realtimeClient.disconnect(); + } + }); + + it("asks apiKeyProvider before the first dial and again before each reconnect", async () => { + let minted = 0; + const apiKeyProvider = vi.fn(async () => `fresh-${++minted}`); + const client = await createClient({ apiKey: "", apiKeyProvider }); + + const realtimeClient = await client.connect(null, { + model: models.realtime("lucy-2.5"), + onRemoteStream: vi.fn(), + }); + expect(apiKeyProvider).toHaveBeenCalledTimes(1); + const first = FakeWebSocket.instances[0]; + expect(apiKeyOf(first)).toBe("fresh-1"); + + // A non-terminal close of a connected session reconnects on a fresh socket with a fresh token. + first.onclose?.({ code: 1006, reason: "" }); + await vi.waitFor(() => expect(FakeWebSocket.instances).toHaveLength(2)); + expect(apiKeyProvider).toHaveBeenCalledTimes(2); + expect(apiKeyOf(FakeWebSocket.instances[1])).toBe("fresh-2"); + await vi.waitFor(() => expect(realtimeClient.getConnectionState()).toBe("connected")); + + FakeWebSocket.instances[1].onclose?.({ code: 1006, reason: "" }); + await vi.waitFor(() => expect(FakeWebSocket.instances).toHaveLength(3)); + expect(apiKeyOf(FakeWebSocket.instances[2])).toBe("fresh-3"); + realtimeClient.disconnect(); + }); + + it("asks apiKeyProvider again for a connect retry after a transient refusal", async () => { + class RefuseOnceWebSocket extends FakeWebSocket { + static refused = false; + send(data: string): void { + if (JSON.parse(data).type === "livekit_join" && !RefuseOnceWebSocket.refused) { + RefuseOnceWebSocket.refused = true; + setTimeout(() => this.onclose?.({ code: 1013, reason: "Try Again Later" }), 0); + return; + } + super.send(data); + } + } + vi.stubGlobal("WebSocket", RefuseOnceWebSocket); + let minted = 0; + const apiKeyProvider = vi.fn(async () => `fresh-${++minted}`); + const client = await createClient({ apiKey: "", apiKeyProvider }); + + // The retry re-dials after the 1 s backoff floor, with a token minted for that dial. + const realtimeClient = await client.connect(null, { + model: models.realtime("lucy-2.5"), + onRemoteStream: vi.fn(), + }); + + expect(FakeWebSocket.instances).toHaveLength(2); + expect(apiKeyProvider).toHaveBeenCalledTimes(2); + expect(apiKeyOf(FakeWebSocket.instances[0])).toBe("fresh-1"); + expect(apiKeyOf(FakeWebSocket.instances[1])).toBe("fresh-2"); + realtimeClient.disconnect(); + }); + + it("stops a reconnect with TOKEN_EXPIRED instead of dialling when the static token expired mid-session", async () => { + const now = vi.spyOn(Date, "now").mockReturnValue(NOW); + const client = await createClient({ apiKey: clientTokenJwt({ exp: expIn(30) }) }); + const realtimeClient = await client.connect(null, { + model: models.realtime("lucy-2.5"), + onRemoteStream: vi.fn(), + }); + const errors: DecartSDKError[] = []; + const states: string[] = []; + realtimeClient.on("error", (error) => errors.push(error)); + realtimeClient.on("connectionChange", (state) => states.push(state)); + + now.mockReturnValue(NOW + 40_000); + FakeWebSocket.instances[0].onclose?.({ code: 1006, reason: "" }); + + await vi.waitFor(() => expect(errors).toHaveLength(1)); + expect(errors[0]).toMatchObject({ + code: ERROR_CODES.TOKEN_EXPIRED, + data: { expiredSecondsAgo: 10 }, + }); + // One socket only: the expired token was never dialled, and no retry followed. + expect(FakeWebSocket.instances).toHaveLength(1); + // The buffer replays the connect-time states to the late listener; then the failed reconnect. + expect(states).toEqual(["connecting", "connected", "reconnecting", "disconnected"]); + expect(realtimeClient.getConnectionState()).toBe("disconnected"); + }); + + it("rejects connect when apiKeyProvider returns an expired token", async () => { + vi.spyOn(Date, "now").mockReturnValue(NOW); + const client = await createClient({ apiKey: "", apiKeyProvider: () => clientTokenJwt({ exp: expIn(-100) }) }); + + await expect( + client.connect(null, { model: models.realtime("lucy-2.5"), onRemoteStream: vi.fn() }), + ).rejects.toMatchObject({ code: ERROR_CODES.TOKEN_EXPIRED, data: { expiredSecondsAgo: 100 } }); + expect(FakeWebSocket.instances).toHaveLength(0); + }); + + it("rejects connect with the provider's own error when apiKeyProvider fails", async () => { + const client = await createClient({ + apiKey: "", + apiKeyProvider: async () => { + throw new Error("token endpoint answered 503"); + }, + }); + + await expect( + client.connect(null, { model: models.realtime("lucy-2.5"), onRemoteStream: vi.fn() }), + ).rejects.toThrow("token endpoint answered 503"); + expect(FakeWebSocket.instances).toHaveLength(0); + }); + + it("rejects connect when apiKeyProvider resolves to something other than a token string", async () => { + const client = await createClient({ apiKey: "", apiKeyProvider: async () => ({ apiKey: "eyJ" }) as never }); + + await expect( + client.connect(null, { model: models.realtime("lucy-2.5"), onRemoteStream: vi.fn() }), + ).rejects.toMatchObject({ code: ERROR_CODES.INVALID_API_KEY }); + expect(FakeWebSocket.instances).toHaveLength(0); + }); + + it("flows apiKeyProvider from createDecartClient to the realtime dial", async () => { + const apiKeyProvider = vi.fn(async () => "fresh-from-app-server"); + const client = createDecartClient({ apiKeyProvider, telemetry: false }); + + const realtimeClient = await client.realtime.connect(null, { + model: models.realtime("lucy-2.5"), + onRemoteStream: vi.fn(), + }); + + expect(apiKeyProvider).toHaveBeenCalledTimes(1); + expect(apiKeyOf(FakeWebSocket.instances[0])).toBe("fresh-from-app-server"); + realtimeClient.disconnect(); + }); + }); }); describe("SignalingChannel initial handshake", () => { @@ -1573,7 +1818,9 @@ describe("StreamSession startup orchestration", () => { const CAPACITY_REFUSAL = /\b1013\b|at capacity|session limit|concurrent session|try again later/i; /** Starts a connect and opens the socket, leaving the join unanswered. */ - const openHandshake = async (opts: { initialImage?: string; connectRetries?: number } = {}) => { + const openHandshake = async ( + opts: { initialImage?: string; connectRetries?: number; redialUrl?: () => string | Promise } = {}, + ) => { const { StreamSession } = await import("../src/realtime/stream-session.js"); const session = new StreamSession({ url: "wss://example.test/realtime", @@ -1833,6 +2080,56 @@ describe("StreamSession startup orchestration", () => { session.disconnect(); }); + it("retries a redial whose credential source failed transiently, whatever it rejected with", async () => { + vi.useFakeTimers(); + const { retry } = REALTIME_CONFIG.session; + const redialUrl = vi + .fn<() => Promise>() + // A plain SDK error (e.g. a failed mint behind the provider), then a bare string. + .mockRejectedValueOnce({ code: "TOKEN_CREATE_ERROR", message: "Failed to create token: 503" }) + .mockRejectedValueOnce("token endpoint down") + .mockResolvedValue("wss://example.test/realtime?api_key=fresh"); + const { session, ws, ended, connectPromise } = await openHandshake({ redialUrl }); + + ws.onclose?.({ code: 1013, reason: "Try Again Later" }); + await vi.advanceTimersByTimeAsync(retry.minTimeout); + expect(redialUrl).toHaveBeenCalledTimes(1); + await vi.advanceTimersByTimeAsync(retry.minTimeout * retry.factor); + expect(redialUrl).toHaveBeenCalledTimes(2); + expect(FakeWebSocket.instances).toHaveLength(1); + await vi.advanceTimersByTimeAsync(retry.minTimeout * retry.factor ** 2); + expect(redialUrl).toHaveBeenCalledTimes(3); + expect(FakeWebSocket.instances).toHaveLength(2); + + const retried = FakeWebSocket.instances[1]; + expect(retried.url).toContain("api_key=fresh"); + retried.onopen?.(); + await flushMicrotasks(); + sendRoomInfo(retried); + await expect(connectPromise).resolves.toBeUndefined(); + expect(ended).toEqual([]); + session.disconnect(); + }); + + it("does not retry a redial the SDK refused outright, such as an expired client token", async () => { + vi.useFakeTimers(); + const redialUrl = vi.fn<() => Promise>().mockRejectedValue({ + code: "TOKEN_EXPIRED", + message: "Client token expired 10 s ago", + }); + const { session, ws, ended, connectPromise } = await openHandshake({ redialUrl }); + + ws.onclose?.({ code: 1013, reason: "Try Again Later" }); + await vi.advanceTimersByTimeAsync(REALTIME_CONFIG.session.retry.minTimeout); + await expect(connectPromise).rejects.toMatchObject({ sdkError: { code: "TOKEN_EXPIRED" } }); + await vi.advanceTimersByTimeAsync(REALTIME_CONFIG.session.retry.maxTimeout); + + expect(redialUrl).toHaveBeenCalledTimes(1); + expect(FakeWebSocket.instances).toHaveLength(1); + expect(ended).toEqual([]); + expect(session.getConnectionState()).toBe("disconnected"); + }); + it("connectRetries: 0 rejects a transient handshake failure after a single dial", async () => { vi.useFakeTimers(); const { session, ws, ended, connectPromise } = await openHandshake({ connectRetries: 0 }); diff --git a/packages/sdk/tests/unit.test.ts b/packages/sdk/tests/unit.test.ts index 218b7b9..3720f52 100644 --- a/packages/sdk/tests/unit.test.ts +++ b/packages/sdk/tests/unit.test.ts @@ -83,6 +83,18 @@ describe("Decart SDK", () => { expect(() => createDecartClient({ apiKey: "test", realtimeBaseUrl: "not-a-url" })).toThrow("Invalid base URL"); }); + it("creates a client with apiKeyProvider and no apiKey", () => { + expect(() => createDecartClient({ apiKeyProvider: async () => "fresh" })).not.toThrow(); + expect(() => createDecartClient({ apiKey: "parent", apiKeyProvider: () => "fresh" })).not.toThrow(); + expect(() => createDecartClient({ proxy: "/api/decart", apiKeyProvider: () => "fresh" })).not.toThrow(); + }); + + it("rejects an apiKeyProvider that is not a function", () => { + expect(() => createDecartClient({ apiKeyProvider: "ek_token" as never })).toThrow( + expect.objectContaining({ code: "INVALID_OPTIONS", message: expect.stringContaining("apiKeyProvider") }), + ); + }); + it("creates a client with custom realtimeBaseUrl", () => { const decart = createDecartClient({ apiKey: "test",