From 32a732351e7a8061c3bd96c8552eb5d717610812 Mon Sep 17 00:00:00 2001 From: Adir Amsalem Date: Wed, 7 Oct 2026 07:43:22 +0000 Subject: [PATCH 1/3] feat(realtime): refuse expired client tokens before dialling; add apiKeyProvider for fresh tokens on connect and reconnect Client tokens expire 60 s after minting by default. A token minted on page load is often already expired by the time the user grants the camera and presses start, and a reconnect after a drop reuses the same token; the server then refuses the connect after a round trip, and the SDK surfaced that as a generic connection error. The SDK now decodes the token's `exp` before every dial (connect, connect retry, reconnect) and fails with TOKEN_EXPIRED client-side, with a 5 s clock-skew tolerance and a message that says how late the token is and what to do. `createDecartClient({ apiKeyProvider })` accepts an async function returning a fresh credential; it is called before every connect and reconnect (and before `subscribe`). The static `apiKey` path is unchanged, including its timing: the first socket still opens in the same tick. README documents the 60 s TTL, `expiresIn`, minting right before connect, and the provider. --- packages/sdk/AGENTS.md | 1 + packages/sdk/README.md | 34 +++ packages/sdk/src/create-client.ts | 32 ++- packages/sdk/src/public-api.ts | 1 + packages/sdk/src/realtime/browser/index.ts | 2 + packages/sdk/src/realtime/client.ts | 44 +++- packages/sdk/src/realtime/config-realtime.ts | 7 + packages/sdk/src/realtime/credential.ts | 76 ++++++ packages/sdk/src/realtime/factory.ts | 2 + .../sdk/src/realtime/react-native/index.ts | 2 + .../sdk/src/realtime/signaling-channel.ts | 13 +- packages/sdk/src/realtime/stream-session.ts | 38 ++- packages/sdk/src/realtime/subscribe-client.ts | 26 +- packages/sdk/src/utils/errors.ts | 64 ++++- packages/sdk/tests/credential.unit.test.ts | 123 +++++++++ packages/sdk/tests/realtime.unit.test.ts | 237 +++++++++++++++++- packages/sdk/tests/unit.test.ts | 12 + 17 files changed, 676 insertions(+), 38 deletions(-) create mode 100644 packages/sdk/src/realtime/credential.ts create mode 100644 packages/sdk/tests/credential.unit.test.ts 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..0944355 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,9 @@ 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", { message: "apiKeyProvider must be a function" }) + .optional(), baseUrl: z.url().optional(), proxy: proxySchema.optional(), integration: z.string().optional(), @@ -39,6 +43,7 @@ export type DecartClientOptions = | { proxy: string; apiKey?: never; + apiKeyProvider?: ApiKeyProvider; baseUrl?: string; realtimeBaseUrl?: string; integration?: string; @@ -48,6 +53,7 @@ export type DecartClientOptions = | { proxy?: never; apiKey?: string; + apiKeyProvider?: ApiKeyProvider; baseUrl?: string; realtimeBaseUrl?: string; integration?: string; @@ -61,6 +67,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 +86,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 +104,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 +131,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 +158,7 @@ export const createDecartClientForPlatform = (createRealtime: CreateRealtime, op publishBaseUrl: wsBaseUrl, subscribeBaseUrl, apiKey: apiKey || "", + apiKeyProvider, integration, logger, telemetryEnabled, @@ -265,9 +288,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..93e6c20 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, @@ -191,6 +195,10 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { let preparedConnection: PreparedConnection | undefined; try { + // Before any setup: an expired client token is refused here, not by the server after a round + // trip, and the provider (when set) is asked for a token minted for this dial. + const credential = await nextCredential(); + const initialImageRef = isFileRefId(initialState?.image) ? initialState.image : undefined; const initialImage = initialImageRef === undefined && initialState?.image ? await imageToBase64(initialState.image) : undefined; @@ -209,7 +217,7 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { logger, observability: { telemetryEnabled: opts.telemetryEnabled, - apiKey, + apiKey: credential, model: options.model.name, integration, logger, @@ -234,17 +242,26 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { ); } - const queryParams = new URLSearchParams({ - ...(preparedConnection.queryParams ?? {}), - ...(options.queryParams ?? {}), - api_key: apiKey, - model: options.model.name, - ...(resolution ? { resolution } : {}), - ...(speed ? { speed } : {}), - }); + const preparedQueryParams = preparedConnection.queryParams ?? {}; + const signalingUrl = (apiKey: string) => { + const queryParams = new URLSearchParams({ + ...preparedQueryParams, + ...(options.queryParams ?? {}), + api_key: apiKey, + model: options.model.name, + ...(resolution ? { resolution } : {}), + ...(speed ? { speed } : {}), + }); + return `${url}?${queryParams.toString()}`; + }; session = new StreamSession({ - url: `${url}?${queryParams.toString()}`, + url: signalingUrl(credential), + // Each retry and reconnect dials with a credential checked (or minted) for that dial. + redialUrl: () => { + const next = nextCredential(); + return typeof next === "string" ? signalingUrl(next) : next.then(signalingUrl); + }, integration, observability, frameTiming: preparedConnection.frameTiming, @@ -330,7 +347,8 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { observability?.stop(); session?.disconnect(); preparedConnection?.dispose(); - throw error; + // SDK-judged failures (expired client token, unusable provider result) reject as plain SDK errors. + 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..4646e3e 100644 --- a/packages/sdk/src/realtime/config-realtime.ts +++ b/packages/sdk/src/realtime/config-realtime.ts @@ -46,6 +46,13 @@ 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, 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..8efe17e --- /dev/null +++ b/packages/sdk/src/realtime/credential.ts @@ -0,0 +1,76 @@ +import { decodeClientToken, isClientTokenJwt } from "../tokens/claims"; +import { + createApiKeyProviderResultError, + createClientTokenExpiredError, + DecartSDKException, + isDecartSDKError, +} 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; + +/** Resolves and validates the credential for one dial. Synchronous for a static key, so the first socket still opens in the same tick. */ +export type CredentialSource = () => 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 `DecartSDKException` wrapping a TOKEN_EXPIRED error + */ +export function assertClientTokenNotExpired(credential: string, now = Date.now()): void { + if (!isClientTokenJwt(credential)) return; + let expiresAt: string; + try { + expiresAt = decodeClientToken(credential).expiresAt; + } catch { + return; + } + const expiredSecondsAgo = (now - Date.parse(expiresAt)) / 1000; + if (!(expiredSecondsAgo > REALTIME_CONFIG.session.clientTokenExpiryToleranceSeconds)) return; + throw new DecartSDKException(createClientTokenExpiredError(Math.round(expiredSecondsAgo), expiresAt)); +} + +function describeThrown(value: unknown): string { + if (typeof value === "object" && value !== null && "message" in value) return String(value.message); + return String(value); +} + +/** + * Build the per-dial credential source: the provider's fresh token when one is set, else the static + * key, checked for expiry either way. Every thrown value is an `Error`, which the session's retry + * loop requires; SDK errors travel as {@link DecartSDKException} and are never retried. + */ +export function createCredentialSource(options: { apiKey: string; apiKeyProvider?: ApiKeyProvider }): CredentialSource { + const { apiKey, apiKeyProvider } = options; + if (!apiKeyProvider) { + return () => { + assertClientTokenNotExpired(apiKey); + return apiKey; + }; + } + return async () => { + let credential: unknown; + try { + credential = await apiKeyProvider(); + } catch (error) { + if (isDecartSDKError(error)) throw new DecartSDKException(error); + if (error instanceof Error) throw error; + throw new Error(`apiKeyProvider rejected: ${describeThrown(error)}`); + } + if (typeof credential !== "string" || credential.length === 0) { + throw new DecartSDKException(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/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..00a8b49 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 } from "../utils/errors"; import { createConsoleLogger, type Logger } from "../utils/logger"; import { REALTIME_CONFIG } from "./config-realtime"; import type { MediaChannel, MediaChannelFactory, VideoCodec } from "./media-channel"; @@ -93,7 +94,15 @@ 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`. A + * `DecartSDKException` thrown here (expired token, unusable provider result) is final: the dial + * is not retried. + */ + redialUrl?: () => string | Promise; integration?: string; observability?: RealtimeObservability; frameTiming?: boolean; @@ -118,6 +127,7 @@ export class StreamSession { private queue: QueuePosition | null = null; private disposed = false; + private dialed = false; private currentAttempt = 0; private teardownGeneration = 0; @@ -234,6 +244,15 @@ export class StreamSession { this.logger.error("realtime connect: session refused, not retrying", { reason: terminal }); return false; } + // The SDK judged this dial itself (expired client token, unusable provider result); another + // dial would get the same answer from the server. + if (error instanceof DecartSDKException) { + this.logger.error("realtime connect: credential rejected, not retrying", { + code: error.sdkError.code, + error: error.message, + }); + return false; + } const msg = error.message.toLowerCase(); const permanent = REALTIME_CONFIG.session.permanentErrorSubstrings.some((err) => msg.includes(err)); if (permanent) { @@ -444,9 +463,26 @@ export class StreamSession { }); } + /** + * The first dial uses the connect-time URL; every later one asks `redialUrl`, so a reconnect + * never reuses a token that expired while the session ran. + */ + 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; + return url.then((resolved) => { + // Disconnected while the credential was being fetched: do not dial at all. + if (this.disposed) 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..c10ed68 100644 --- a/packages/sdk/src/realtime/subscribe-client.ts +++ b/packages/sdk/src/realtime/subscribe-client.ts @@ -1,8 +1,9 @@ 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, unwrapSDKError } 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 +27,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 +71,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 +119,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); @@ -146,6 +141,8 @@ export const createRealTimeSubscribeClient = (opts: RealTimeSubscribeClientOptio }; try { + // An expired client token is refused before anything is loaded or fetched. + const apiKey = await nextCredential(); const { Room: LiveKitRoom, RoomEvent } = await loadLiveKitClient(); observability = new RealtimeObservability({ telemetryEnabled: false, @@ -246,9 +243,10 @@ export const createRealTimeSubscribeClient = (opts: RealTimeSubscribeClientOptio room.disconnect().catch(() => {}); } frameMetadataWorker?.terminate(); - if (isDecartSDKError(error)) { - logger.error("Realtime subscribe error", { error: error.message }); - throw error; + const sdkError = unwrapSDKError(error); + if (sdkError) { + logger.error("Realtime subscribe error", { error: sdkError.message }); + throw sdkError; } const err = error instanceof Error ? error : new Error(String(error)); logger.error("Realtime subscribe error", { error: err.message }); diff --git a/packages/sdk/src/utils/errors.ts b/packages/sdk/src/utils/errors.ts index dc1d785..cc78474 100644 --- a/packages/sdk/src/utils/errors.ts +++ b/packages/sdk/src/utils/errors.ts @@ -46,10 +46,69 @@ 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}, for code paths that must throw real errors: p-retry + * rejects anything else, and the session's `error` channel is typed `Error`. `classifyWebrtcError` + * and the public entry points 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; + } +} + +/** The plain SDK error when `error` is one or wraps one, else `undefined`. */ +export function unwrapSDKError(error: unknown): DecartSDKError | undefined { + if (error instanceof DecartSDKException) return error.sdkError; + return isDecartSDKError(error) ? error : undefined; +} + 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" + : received === null + ? "null" + : Array.isArray(received) + ? "an array" + : 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}.`, + { received: got }, ); } @@ -106,6 +165,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..b175e64 --- /dev/null +++ b/packages/sdk/tests/credential.unit.test.ts @@ -0,0 +1,123 @@ +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 { DecartSDKException, ERROR_CODES } from "../src/utils/errors.js"; + +const base64url = (text: string) => btoa(text).replace(/\+/g, "-").replace(/\//g, "_").replace(/=+$/, ""); +/** An unsigned JWT in the platform's client-token shape; only `exp` matters here. */ +const clientTokenJwt = (claims: Record) => + `eyJhbGciOiJFZERTQSJ9.${base64url(JSON.stringify({ sub: "user_1", ...claims }))}.sig`; + +const NOW = Date.UTC(2026, 9, 7, 12, 0, 0); +const expIn = (seconds: number) => Math.floor(NOW / 1000) + 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)); + let thrown: unknown; + try { + assertClientTokenNotExpired(clientTokenJwt({ exp }), NOW); + } catch (error) { + thrown = error; + } + + expect(thrown).toBeInstanceOf(DecartSDKException); + const { sdkError } = thrown as DecartSDKException; + expect(sdkError.code).toBe(ERROR_CODES.TOKEN_EXPIRED); + expect(sdkError.message).toContain(`Client token expired ${tolerance + 1} s ago`); + expect(sdkError.message).toContain("apiKeyProvider"); + expect(sdkError.message).toContain("60 s"); + expect(sdkError.data).toEqual({ + claim: "exp", + expiresAt: new Date(exp * 1000).toISOString(), + expiredSecondsAgo: tolerance + 1, + }); + }); + + it("reports whole seconds for tokens hours late", () => { + const exp = expIn(-7200); + expect(() => assertClientTokenNotExpired(clientTokenJwt({ exp }), NOW + 400)).toThrow( + expect.objectContaining({ + sdkError: 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: Math.floor(Date.now() / 1000) - 120 }); + expect(() => createCredentialSource({ apiKey: expired })()).toThrow( + expect.objectContaining({ sdkError: expect.objectContaining({ code: ERROR_CODES.TOKEN_EXPIRED }) }), + ); + await expect(createCredentialSource({ apiKey: "", apiKeyProvider: async () => expired })()).rejects.toMatchObject({ + sdkError: { 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({ sdkError: { code: ERROR_CODES.INVALID_API_KEY } }); + } + const next = createCredentialSource({ + apiKey: "", + apiKeyProvider: async () => ({ apiKey: "x" }) as unknown as string, + }); + await expect(next()).rejects.toThrow("got object"); + }); + + it("rethrows a provider Error as-is and wraps anything else so the retry loop sees an Error", async () => { + const failure = new TypeError("Failed to fetch"); + await expect(createCredentialSource({ apiKey: "", apiKeyProvider: () => Promise.reject(failure) })()).rejects.toBe( + failure, + ); + + await expect( + createCredentialSource({ apiKey: "", apiKeyProvider: () => Promise.reject("endpoint down") })(), + ).rejects.toThrow("apiKeyProvider rejected: endpoint down"); + + const sdkError = { code: ERROR_CODES.TOKEN_CREATE_ERROR, message: "Failed to create token: 503" }; + const thrown = await createCredentialSource({ apiKey: "", apiKeyProvider: () => Promise.reject(sdkError) })().catch( + (e: unknown) => e, + ); + expect(thrown).toBeInstanceOf(DecartSDKException); + expect((thrown as DecartSDKException).sdkError).toBe(sdkError); + }); +}); diff --git a/packages/sdk/tests/realtime.unit.test.ts b/packages/sdk/tests/realtime.unit.test.ts index 665222b..243d303 100644 --- a/packages/sdk/tests/realtime.unit.test.ts +++ b/packages/sdk/tests/realtime.unit.test.ts @@ -1,9 +1,10 @@ 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"; const liveKitMock = vi.hoisted(() => { const roomInstances: MockRoom[] = []; @@ -109,6 +110,11 @@ const flushMicrotasks = async () => { await Promise.resolve(); }; +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). */ +const clientTokenJwt = (claims: Record) => + `eyJhbGciOiJFZERTQSJ9.${base64url(JSON.stringify({ sub: "user_1", ...claims }))}.sig`; + type FakeWebSocketMessageEvent = { data: string; }; @@ -403,6 +409,42 @@ 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 subscriber = createRealTimeSubscribeClient({ + baseUrl: "https://api.example.test", + apiKey: clientTokenJwt({ exp: Math.floor(Date.now() / 1000) - 600 }), + logger, + }); + + await expect( + subscriber.subscribe({ token: btoa(JSON.stringify({ room_name: "room-1" })), onRemoteStream: () => {} }), + ).rejects.toMatchObject({ code: ERROR_CODES.TOKEN_EXPIRED, data: { expiredSecondsAgo: 600 } }); + expect(fetch).not.toHaveBeenCalled(); + expect(liveKitMock.roomInstances).toHaveLength(0); + }); + 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 +582,8 @@ describe("realtime.connect options", () => { afterEach(() => { vi.unstubAllGlobals(); + vi.restoreAllMocks(); + vi.useRealTimers(); }); it("adds resolution to the realtime URL when provided", async () => { @@ -758,6 +802,197 @@ 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) => Math.floor(NOW / 1000) + seconds; + const apiKeyOf = (ws: FakeWebSocket) => new URL(ws.url).searchParams.get("api_key"); + + const createClient = async ( + credentials: { apiKey: string; apiKeyProvider?: () => string | Promise }, + logger: { debug(): void; info(): void; warn(): void; error(): void } = { + debug() {}, + info() {}, + warn() {}, + error() {}, + }, + ) => { + 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", () => { 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", From 25d62960ef85540f4fa863b8281018cfa7a94c67 Mon Sep 17 00:00:00 2001 From: Adir Amsalem Date: Wed, 7 Oct 2026 08:39:26 +0000 Subject: [PATCH 2/3] docs(examples): nextjs-realtime mints a fresh client token per dial via apiKeyProvider The Next.js example now creates one client with `apiKeyProvider` and lets the SDK fetch a token from /api/realtime-token right before every connect and reconnect, instead of fetching one itself and connecting with it. A Start button makes the point: the token is minted when the user acts, not when the page loads. The route sets `expiresIn` explicitly and the README walks through the flow. --- examples/nextjs-realtime/README.md | 35 ++++- .../app/api/realtime-token/route.ts | 12 +- .../components/video-stream.tsx | 146 ++++++++++-------- 3 files changed, 115 insertions(+), 78 deletions(-) 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..dd39bc2 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,63 @@ export function VideoStream({ prompt }: VideoStreamProps) { const inputRef = useRef(null); const outputRef = useRef(null); const realtimeClientRef = useRef(null); + const cameraRef = useRef(null); 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 stop() { + realtimeClientRef.current?.disconnect(); + realtimeClientRef.current = null; + for (const track of cameraRef.current?.getTracks() ?? []) track.stop(); + cameraRef.current = null; + setRunning(false); + setStatus("idle"); + } + + async function start() { + setRunning(true); + try { + // However long the page has been open, the token is minted only now, inside connect(). + setStatus("requesting camera..."); + const camera = await navigator.mediaDevices.getUserMedia({ + video: { frameRate: model.fps, width: model.width, height: model.height }, + }); + 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; + }, + // "connected" | "generating" | "reconnecting" | "disconnected". A reconnect asks + // fetchClientToken again, so it never dials with the token the session started with. + onConnectionChange: setStatus, + initialState: { prompt: { text: prompt, enhance: true } }, + }); + realtimeClientRef.current = realtimeClient; + + realtimeClient.on("error", (error) => setStatus(`error: ${describeError(error)}`)); + } catch (error) { + setStatus(`error: ${describeError(error)}`); + realtimeClientRef.current?.disconnect(); + realtimeClientRef.current = null; + for (const track of cameraRef.current?.getTracks() ?? []) track.stop(); + cameraRef.current = null; + setRunning(false); } + } - start(); - + // Release the camera and the session when the component unmounts. + useEffect(() => { return () => { - mounted = false; realtimeClientRef.current?.disconnect(); + for (const track of cameraRef.current?.getTracks() ?? []) track.stop(); }; }, []); - // 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 +99,12 @@ export function VideoStream({ prompt }: VideoStreamProps) { return (
-

Status: {status}

+

+ + Status: {status} +

Input

From f94641713f0f7344d8a607fb091dbe59dccb8c94 Mon Sep 17 00:00:00 2001 From: Adir Amsalem Date: Wed, 7 Oct 2026 08:54:52 +0000 Subject: [PATCH 3/3] refactor(realtime): plain SDK errors from the credential path; wrap once at the retry boundary - credential.ts throws plain DecartSDKErrors like the rest of the SDK and lets provider rejections through unchanged; StreamSession wraps them into DecartSDKException once, where p-retry needs an Error. - Retry permanence is a config list of SDK error codes (REALTIME_CONFIG.session.permanentErrorCodes) next to the existing substring list, so a provider's transient SDK error is retried as the docs say; only TOKEN_EXPIRED and INVALID_API_KEY end the attempts. - subscribe() resolves the credential before its try block, so a provider rejection reaches the caller unchanged, as from connect(). - The per-dial URL builder captures only what it needs instead of the whole connect options; the provider round trip overlaps image encoding; telemetry reports with the credential of the dial that started the session; the dial resolver's stale guard also checks the attempt. - Example: one release() for the teardown copies, and an attempt counter so a Stop or unmount during the camera prompt or connect() releases the camera and the session instead of leaking them. - Tests: shared client-token JWT helper; two session-level tests for the retry boundary; subscribe provider-rejection test; pinned clock in the subscribe expiry test. --- .../components/video-stream.tsx | 38 ++++--- packages/sdk/src/create-client.ts | 4 +- packages/sdk/src/realtime/client.ts | 33 +++--- packages/sdk/src/realtime/config-realtime.ts | 8 ++ packages/sdk/src/realtime/credential.ts | 48 +++------ .../observability/realtime-observability.ts | 5 +- packages/sdk/src/realtime/stream-session.ts | 44 ++++---- packages/sdk/src/realtime/subscribe-client.ts | 21 ++-- packages/sdk/src/utils/errors.ts | 22 +--- packages/sdk/tests/credential.unit.test.ts | 76 +++++-------- packages/sdk/tests/helpers/client-token.ts | 11 ++ packages/sdk/tests/realtime.unit.test.ts | 100 ++++++++++++++---- 12 files changed, 226 insertions(+), 184 deletions(-) create mode 100644 packages/sdk/tests/helpers/client-token.ts diff --git a/examples/nextjs-realtime/components/video-stream.tsx b/examples/nextjs-realtime/components/video-stream.tsx index dd39bc2..1a8f361 100644 --- a/examples/nextjs-realtime/components/video-stream.tsx +++ b/examples/nextjs-realtime/components/video-stream.tsx @@ -35,26 +35,39 @@ export function VideoStream({ prompt }: VideoStreamProps) { 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"); const [running, setRunning] = useState(false); - function stop() { + 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 { - // However long the page has been open, the token is minted only now, inside connect(). 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; @@ -64,31 +77,26 @@ export function VideoStream({ prompt }: VideoStreamProps) { onRemoteStream: (transformedStream) => { if (outputRef.current) outputRef.current.srcObject = transformedStream; }, - // "connected" | "generating" | "reconnecting" | "disconnected". A reconnect asks - // fetchClientToken again, so it never dials with the token the session started with. 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) { - setStatus(`error: ${describeError(error)}`); - realtimeClientRef.current?.disconnect(); - realtimeClientRef.current = null; - for (const track of cameraRef.current?.getTracks() ?? []) track.stop(); - cameraRef.current = null; + if (cancelled()) return; + release(); setRunning(false); + setStatus(`error: ${describeError(error)}`); } } // Release the camera and the session when the component unmounts. - useEffect(() => { - return () => { - realtimeClientRef.current?.disconnect(); - for (const track of cameraRef.current?.getTracks() ?? []) track.stop(); - }; - }, []); + useEffect(() => release, []); // Update the prompt on the running session when it changes. useEffect(() => { diff --git a/packages/sdk/src/create-client.ts b/packages/sdk/src/create-client.ts index 0944355..a8f35fa 100644 --- a/packages/sdk/src/create-client.ts +++ b/packages/sdk/src/create-client.ts @@ -16,9 +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", { message: "apiKeyProvider must be a function" }) - .optional(), + apiKeyProvider: z.custom((val) => typeof val === "function").optional(), baseUrl: z.url().optional(), proxy: proxySchema.optional(), integration: z.string().optional(), diff --git a/packages/sdk/src/realtime/client.ts b/packages/sdk/src/realtime/client.ts index 93e6c20..bef2e30 100644 --- a/packages/sdk/src/realtime/client.ts +++ b/packages/sdk/src/realtime/client.ts @@ -183,6 +183,7 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { onConnectionQuality, onQueuePosition, initialState, + queryParams: extraQueryParams, resolution, speed, retries, @@ -195,13 +196,13 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { let preparedConnection: PreparedConnection | undefined; try { - // Before any setup: an expired client token is refused here, not by the server after a round - // trip, and the provider (when set) is asked for a token minted for this dial. - const credential = await nextCredential(); - 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; @@ -242,13 +243,17 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { ); } - const preparedQueryParams = preparedConnection.queryParams ?? {}; - const signalingUrl = (apiKey: string) => { + // 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, - ...(options.queryParams ?? {}), + ...extraQueryParams, api_key: apiKey, - model: options.model.name, + model: modelName, ...(resolution ? { resolution } : {}), ...(speed ? { speed } : {}), }); @@ -256,11 +261,10 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { }; session = new StreamSession({ - url: signalingUrl(credential), - // Each retry and reconnect dials with a credential checked (or minted) for that dial. + url: dialUrl(credential), redialUrl: () => { const next = nextCredential(); - return typeof next === "string" ? signalingUrl(next) : next.then(signalingUrl); + return typeof next === "string" ? dialUrl(next) : next.then(dialUrl); }, integration, observability, @@ -294,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)); @@ -347,7 +351,6 @@ export const createRealTimeClient = (opts: RealTimeClientOptions) => { observability?.stop(); session?.disconnect(); preparedConnection?.dispose(); - // SDK-judged failures (expired client token, unusable provider result) reject as plain SDK errors. 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 4646e3e..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). */ @@ -53,6 +55,12 @@ export const REALTIME_CONFIG = { * 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 index 8efe17e..5a1dec8 100644 --- a/packages/sdk/src/realtime/credential.ts +++ b/packages/sdk/src/realtime/credential.ts @@ -1,10 +1,5 @@ -import { decodeClientToken, isClientTokenJwt } from "../tokens/claims"; -import { - createApiKeyProviderResultError, - createClientTokenExpiredError, - DecartSDKException, - isDecartSDKError, -} from "../utils/errors"; +import { decodeClientToken } from "../tokens/claims"; +import { createApiKeyProviderResultError, createClientTokenExpiredError } from "../utils/errors"; import { REALTIME_CONFIG } from "./config-realtime"; /** @@ -17,40 +12,32 @@ import { REALTIME_CONFIG } from "./config-realtime"; */ export type ApiKeyProvider = () => string | Promise; -/** Resolves and validates the credential for one dial. Synchronous for a static key, so the first socket still opens in the same tick. */ -export type CredentialSource = () => 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 `DecartSDKException` wrapping a TOKEN_EXPIRED error + * @throws `TOKEN_EXPIRED` */ export function assertClientTokenNotExpired(credential: string, now = Date.now()): void { - if (!isClientTokenJwt(credential)) return; let expiresAt: string; try { - expiresAt = decodeClientToken(credential).expiresAt; + ({ expiresAt } = decodeClientToken(credential)); } catch { return; } const expiredSecondsAgo = (now - Date.parse(expiresAt)) / 1000; - if (!(expiredSecondsAgo > REALTIME_CONFIG.session.clientTokenExpiryToleranceSeconds)) return; - throw new DecartSDKException(createClientTokenExpiredError(Math.round(expiredSecondsAgo), expiresAt)); -} - -function describeThrown(value: unknown): string { - if (typeof value === "object" && value !== null && "message" in value) return String(value.message); - return String(value); + if (expiredSecondsAgo > REALTIME_CONFIG.session.clientTokenExpiryToleranceSeconds) { + throw createClientTokenExpiredError(Math.round(expiredSecondsAgo), expiresAt); + } } /** - * Build the per-dial credential source: the provider's fresh token when one is set, else the static - * key, checked for expiry either way. Every thrown value is an `Error`, which the session's retry - * loop requires; SDK errors travel as {@link DecartSDKException} and are never retried. + * 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 }): CredentialSource { +export function createCredentialSource(options: { apiKey: string; apiKeyProvider?: ApiKeyProvider }): ApiKeyProvider { const { apiKey, apiKeyProvider } = options; if (!apiKeyProvider) { return () => { @@ -59,17 +46,8 @@ export function createCredentialSource(options: { apiKey: string; apiKeyProvider }; } return async () => { - let credential: unknown; - try { - credential = await apiKeyProvider(); - } catch (error) { - if (isDecartSDKError(error)) throw new DecartSDKException(error); - if (error instanceof Error) throw error; - throw new Error(`apiKeyProvider rejected: ${describeThrown(error)}`); - } - if (typeof credential !== "string" || credential.length === 0) { - throw new DecartSDKException(createApiKeyProviderResultError(credential)); - } + 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/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/stream-session.ts b/packages/sdk/src/realtime/stream-session.ts index 00a8b49..05ca691 100644 --- a/packages/sdk/src/realtime/stream-session.ts +++ b/packages/sdk/src/realtime/stream-session.ts @@ -1,7 +1,7 @@ import mitt, { type Emitter } from "mitt"; import pRetry, { AbortError } from "p-retry"; -import { DecartSDKException } from "../utils/errors"; +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"; @@ -69,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 } : {}) })); } @@ -98,9 +108,8 @@ interface StreamSessionConfig { 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`. A - * `DecartSDKException` thrown here (expired token, unusable provider result) is final: the dial - * is not retried. + * 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; @@ -244,11 +253,10 @@ export class StreamSession { this.logger.error("realtime connect: session refused, not retrying", { reason: terminal }); return false; } - // The SDK judged this dial itself (expired client token, unusable provider result); another - // dial would get the same answer from the server. - if (error instanceof DecartSDKException) { - this.logger.error("realtime connect: credential rejected, not retrying", { - code: error.sdkError.code, + 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; @@ -313,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; } } @@ -463,19 +469,17 @@ export class StreamSession { }); } - /** - * The first dial uses the connect-time URL; every later one asks `redialUrl`, so a reconnect - * never reuses a token that expired while the session ran. - */ + /** 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 while the credential was being fetched: do not dial at all. - if (this.disposed) throw new AbortError("Stale connect attempt"); + // 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; }); } diff --git a/packages/sdk/src/realtime/subscribe-client.ts b/packages/sdk/src/realtime/subscribe-client.ts index c10ed68..911fa78 100644 --- a/packages/sdk/src/realtime/subscribe-client.ts +++ b/packages/sdk/src/realtime/subscribe-client.ts @@ -1,6 +1,12 @@ import type { RemoteParticipant, RemoteTrack, Room } from "livekit-client"; -import { classifyWebrtcError, createSDKError, type DecartSDKError, ERROR_CODES, unwrapSDKError } 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"; @@ -140,9 +146,11 @@ 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 { - // An expired client token is refused before anything is loaded or fetched. - const apiKey = await nextCredential(); const { Room: LiveKitRoom, RoomEvent } = await loadLiveKitClient(); observability = new RealtimeObservability({ telemetryEnabled: false, @@ -243,10 +251,9 @@ export const createRealTimeSubscribeClient = (opts: RealTimeSubscribeClientOptio room.disconnect().catch(() => {}); } frameMetadataWorker?.terminate(); - const sdkError = unwrapSDKError(error); - if (sdkError) { - logger.error("Realtime subscribe error", { error: sdkError.message }); - throw sdkError; + if (isDecartSDKError(error)) { + logger.error("Realtime subscribe error", { error: error.message }); + throw error; } const err = error instanceof Error ? error : new Error(String(error)); logger.error("Realtime subscribe error", { error: err.message }); diff --git a/packages/sdk/src/utils/errors.ts b/packages/sdk/src/utils/errors.ts index cc78474..c4725c6 100644 --- a/packages/sdk/src/utils/errors.ts +++ b/packages/sdk/src/utils/errors.ts @@ -56,9 +56,9 @@ export function isDecartSDKError(error: unknown): error is DecartSDKError { } /** - * An `Error` carrying a {@link DecartSDKError}, for code paths that must throw real errors: p-retry - * rejects anything else, and the session's `error` channel is typed `Error`. `classifyWebrtcError` - * and the public entry points unwrap it, so callers still receive the plain SDK error. + * 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; @@ -71,12 +71,6 @@ export class DecartSDKException extends Error { } } -/** The plain SDK error when `error` is one or wraps one, else `undefined`. */ -export function unwrapSDKError(error: unknown): DecartSDKError | undefined { - if (error instanceof DecartSDKException) return error.sdkError; - return isDecartSDKError(error) ? error : undefined; -} - export function createInvalidApiKeyError(): DecartSDKError { return createSDKError( ERROR_CODES.INVALID_API_KEY, @@ -97,18 +91,10 @@ export function createClientTokenExpiredError(expiredSecondsAgo: number, expires } export function createApiKeyProviderResultError(received: unknown): DecartSDKError { - const got = - received === "" - ? "an empty string" - : received === null - ? "null" - : Array.isArray(received) - ? "an array" - : typeof received; + 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}.`, - { received: got }, ); } diff --git a/packages/sdk/tests/credential.unit.test.ts b/packages/sdk/tests/credential.unit.test.ts index b175e64..6c7be48 100644 --- a/packages/sdk/tests/credential.unit.test.ts +++ b/packages/sdk/tests/credential.unit.test.ts @@ -2,15 +2,11 @@ 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 { DecartSDKException, ERROR_CODES } from "../src/utils/errors.js"; - -const base64url = (text: string) => btoa(text).replace(/\+/g, "-").replace(/\//g, "_").replace(/=+$/, ""); -/** An unsigned JWT in the platform's client-token shape; only `exp` matters here. */ -const clientTokenJwt = (claims: Record) => - `eyJhbGciOiJFZERTQSJ9.${base64url(JSON.stringify({ sub: "user_1", ...claims }))}.sig`; +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) => Math.floor(NOW / 1000) + seconds; +const expIn = (seconds: number) => expAt(NOW, seconds); const tolerance = REALTIME_CONFIG.session.clientTokenExpiryToleranceSeconds; describe("assertClientTokenNotExpired", () => { @@ -30,32 +26,19 @@ describe("assertClientTokenNotExpired", () => { it("rejects a token past the tolerance with TOKEN_EXPIRED and says how late it is", () => { const exp = expIn(-(tolerance + 1)); - let thrown: unknown; - try { - assertClientTokenNotExpired(clientTokenJwt({ exp }), NOW); - } catch (error) { - thrown = error; - } - expect(thrown).toBeInstanceOf(DecartSDKException); - const { sdkError } = thrown as DecartSDKException; - expect(sdkError.code).toBe(ERROR_CODES.TOKEN_EXPIRED); - expect(sdkError.message).toContain(`Client token expired ${tolerance + 1} s ago`); - expect(sdkError.message).toContain("apiKeyProvider"); - expect(sdkError.message).toContain("60 s"); - expect(sdkError.data).toEqual({ - claim: "exp", - expiresAt: new Date(exp * 1000).toISOString(), - expiredSecondsAgo: 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", () => { - const exp = expIn(-7200); - expect(() => assertClientTokenNotExpired(clientTokenJwt({ exp }), NOW + 400)).toThrow( - expect.objectContaining({ - sdkError: expect.objectContaining({ data: expect.objectContaining({ expiredSecondsAgo: 7200 }) }), - }), + expect(() => assertClientTokenNotExpired(clientTokenJwt({ exp: expIn(-7200) }), NOW + 400)).toThrow( + expect.objectContaining({ data: expect.objectContaining({ expiredSecondsAgo: 7200 }) }), ); }); }); @@ -82,42 +65,35 @@ describe("createCredentialSource", () => { }); it("checks expiry on provider tokens and static tokens alike", async () => { - const expired = clientTokenJwt({ exp: Math.floor(Date.now() / 1000) - 120 }); + const expired = clientTokenJwt({ exp: expAt(Date.now(), -120) }); expect(() => createCredentialSource({ apiKey: expired })()).toThrow( - expect.objectContaining({ sdkError: expect.objectContaining({ code: ERROR_CODES.TOKEN_EXPIRED }) }), + expect.objectContaining({ code: ERROR_CODES.TOKEN_EXPIRED }), ); await expect(createCredentialSource({ apiKey: "", apiKeyProvider: async () => expired })()).rejects.toMatchObject({ - sdkError: { code: ERROR_CODES.TOKEN_EXPIRED }, + 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({ sdkError: { code: ERROR_CODES.INVALID_API_KEY } }); + 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.toThrow("got object"); + await expect(next()).rejects.toMatchObject({ message: expect.stringContaining("got object") }); }); - it("rethrows a provider Error as-is and wraps anything else so the retry loop sees an Error", async () => { - const failure = new TypeError("Failed to fetch"); - await expect(createCredentialSource({ apiKey: "", apiKeyProvider: () => Promise.reject(failure) })()).rejects.toBe( - failure, - ); - - await expect( - createCredentialSource({ apiKey: "", apiKeyProvider: () => Promise.reject("endpoint down") })(), - ).rejects.toThrow("apiKeyProvider rejected: endpoint down"); - - const sdkError = { code: ERROR_CODES.TOKEN_CREATE_ERROR, message: "Failed to create token: 503" }; - const thrown = await createCredentialSource({ apiKey: "", apiKeyProvider: () => Promise.reject(sdkError) })().catch( - (e: unknown) => e, - ); - expect(thrown).toBeInstanceOf(DecartSDKException); - expect((thrown as DecartSDKException).sdkError).toBe(sdkError); + 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 243d303..aff6c44 100644 --- a/packages/sdk/tests/realtime.unit.test.ts +++ b/packages/sdk/tests/realtime.unit.test.ts @@ -5,6 +5,7 @@ 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[] = []; @@ -110,11 +111,6 @@ const flushMicrotasks = async () => { await Promise.resolve(); }; -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). */ -const clientTokenJwt = (claims: Record) => - `eyJhbGciOiJFZERTQSJ9.${base64url(JSON.stringify({ sub: "user_1", ...claims }))}.sig`; - type FakeWebSocketMessageEvent = { data: string; }; @@ -432,17 +428,39 @@ describe("realtime.subscribe", () => { 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: clientTokenJwt({ exp: Math.floor(Date.now() / 1000) - 600 }), + apiKey: "", + apiKeyProvider: () => Promise.reject(failure), logger, }); await expect( subscriber.subscribe({ token: btoa(JSON.stringify({ room_name: "room-1" })), onRemoteStream: () => {} }), - ).rejects.toMatchObject({ code: ERROR_CODES.TOKEN_EXPIRED, data: { expiredSecondsAgo: 600 } }); + ).rejects.toBe(failure); expect(fetch).not.toHaveBeenCalled(); - expect(liveKitMock.roomInstances).toHaveLength(0); }); it("terminates the frame-metadata worker when subscribe connect fails", async () => { @@ -805,18 +823,10 @@ describe("realtime.connect options", () => { describe("client token expiry preflight and apiKeyProvider", () => { const NOW = Date.UTC(2026, 9, 7, 12, 0, 0); - const expIn = (seconds: number) => Math.floor(NOW / 1000) + seconds; + 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 }, - logger: { debug(): void; info(): void; warn(): void; error(): void } = { - debug() {}, - info() {}, - warn() {}, - error() {}, - }, - ) => { + const createClient = async (credentials: { apiKey: string; apiKeyProvider?: () => string | Promise }) => { const { createRealTimeClient } = await import("../src/realtime/client.js"); return createRealTimeClient({ baseUrl: "wss://api3.decart.ai", @@ -1808,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", @@ -2068,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 });