diff --git a/CLAUDE.md b/CLAUDE.md index 6851dfc47..84cc8fd1b 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -73,17 +73,17 @@ This is a single-package SDK with **no monorepo**. The public surface is everyth - **`client.ts` — `StreamChat` facade.** ~5k-line class. Prefer `StreamChat.getInstance(key, options?)` — the constructor exists for advanced uses but `getInstance` is what `connectUser` warnings and most docs assume. Owns: the axios instance, WS connection lifecycle, `TokenManager`, and a registry of subsystem managers (`threads`, `polls`, `notifications`, `reminders`, `moderation`, `uploadManager`, `messageDeliveryReporter`, plus an optional `offlineDb` injected via `setOfflineDBApi`). New REST endpoints are added here as methods that call `axiosInstance` and return a type from `types.ts`. - **`channel.ts` (~2.5k lines) + `channel_state.ts` (~1.1k) + `channel_batch_updater.ts`** — per-channel object and its in-memory state. **Messages are NOT stored on `channel.state`.** The message list, thread replies, and pinned messages each live in a paginator — `channel.messagePaginator`, `thread.messagePaginator`, and `channel.pinnedMessagesPaginator` — which are the single source of truth (interval storage + a canonical `ItemIndex`). Read them via `channel.messagePaginator.state.items` / `.getItem(id)` / `.headmostItem` (newest loaded item), and mutate via the paginator (`ingestItem` / `removeItem`), never a legacy `channel.state.addMessageSorted()` / `state.messages` (removed). `channel.state.last_message_at` was **removed**; the channel's latest-message timestamp lives on `channel.messagePaginator.lastMessageAt` (its `aggregateState` store — seeded from `ChannelResponse.last_message_at`, then advanced monotonically as messages are ingested). See `docs/breaking-changes-v14-v15.md`. - **`ChannelManager.ts`** — channel _lists_. Holds one or more `ChannelPaginator`s (`state.paginators`), keeps them in sync with WS events through an `EventHandlerPipeline` per event type, and arbitrates ownership when a channel matches several lists (`ownershipResolver` / `createPriorityOwnershipResolver`). Replaced the old `channel_manager.ts` (single hand-sorted `state.channels` list with named handler overrides) in v10 — see `v9-to-v10-migration-guide-methods.md`. Filtering and ordering are the paginator's job: `matchesFilter()` runs the filter compiler over `Channel` field resolvers and ordering comes from a comparator compiled from `sort`. The manager is instantiated by the `StreamChat` constructor and lives as long as the client (`client.channelManager`) — it is not configurable through the client options; register lists with `insertPaginator({ paginator, index? })`, detach them with `removePaginator(paginatorOrId)` and set cross-list ownership with `setOwnershipResolver(resolverOrPriorityIds?)`. `setPaginators(paginators)` is the primitive the other two build on — use it (or `clearPaginators()`) for batches, since it publishes one state update instead of one per paginator, and skips the update entirely when the set is unchanged. Registration and loaded data have different owners: `disconnectUser` calls `resetPaginatorStates()`, which discards each list's channels (they belong to the user going away) while leaving the lists themselves registered, since which lists exist is the integrator's configuration. Event handling stays customizable: `ChannelManagerOptions.eventHandlers` replaces the default map wholesale at construction (start from `getDefaultHandlers()` to enrich it instead), and `addEventHandler` / `setEventHandlers` / `removeEventHandlers` adjust the pipelines afterwards — which is the only route for `client.channelManager`, since the client constructs it without options. The exported `ignoreEventsForUnknownChannels` handler, inserted at `index: 0`, is how a list opts out of pulling in channels it has not loaded. -- **`connection/wsConnection/StableWSConnection.ts` (`StableWSConnection`)** — the realtime socket, and the only transport: the long-poll fallback (`connection_fallback.ts`, `enableWSFallback`, `transport.changed`) was removed in v10. It connects to `/api/v2/connect`, which authenticates off the **first frame the client sends** (`client._buildWSAuthMessage()`) rather than the query string, and answers with a `connection.ok` hello event instead of v1's `health.check`. It runs a 25s ping and a 35s connection check, and reconnects on close/error, publishing every status transition into `client.wsConnection.state`. `connection.ok` is not in the OpenAPI spec yet, so its type is hand-written in `types.ts` and decoded through a shim in `StableWSConnection.ts` — both are marked for deletion once the backend publishes the event. - **It no longer registers `window` listeners.** The offline/online edge now arrives from `client.networkConnection` (see below), which inverts the old dependency: the socket is a _consumer_ of network status rather than the thing that detects it. Its own settings — connect timeout, ping interval, connection-check grace period, the UI's offline-notification delay, `webSocketImpl`, `urlParams` and an injectable socket — live in `client.wsConnection.config`, not in `StreamChatOptions`; `client.defaultWSTimeout` is gone. Every instance is created by `WSConnection.connect()`, which also disconnects the one it replaces, and read live from that config so an `updateConfig` reaches an already-open socket. It holds its `WSConnection` parent rather than the client, and reaches the client through it. +- **`connection/wsConnection/StableWSConnection.ts` (`StableWSConnection`)** — the realtime socket. With `client.wsConnection.config.enableWSFallback` (a `wsConnection` config key; a `StreamChatOptions` entry in v9), `client.wsConnection.connect()` switches to v9's `WSConnectionFallback` (`src/connection/wsConnection/WSConnectionFallback.ts`, long-poll on `/api/v2/longpoll`, held on `client.wsConnection.fallback`) when the socket fails with a network error — which it reports after the full `connectTimeoutMs`, where v9 cut it to 6s — dispatches `connection.fallback_activated`, and stays on it, across `disconnectUser()` too. That class is v9's, adapted only where v10 changed what it relied on: it writes `client.wsConnection.state`, drives `ConnectionIdManager`, takes network status from `WSConnection`'s one subscription through `_applyNetworkStatus` (ignored while disconnected, where v9 stopped listening for good at its first disconnect), exposes `isConnecting`, and sets `connection_id` on its own polls and close. It calls the generated `client.longPoll()` (whose `close` param is hand-added until the spec has it), with per-request `timeout`s through `StreamRequestOptions` and an `AbortController` for cancellation; `ApiClient.resolveUrl` maps a local API's `:3030` to `:8900` for that path. The socket connects to `/api/v2/connect`, which authenticates off the **first frame the client sends** (`client._buildWSAuthMessage()`) rather than the query string, and answers with a `connection.ok` hello event instead of v1's `health.check`. It runs a 25s ping and a 35s connection check, and reconnects on close/error, publishing every status transition into `client.wsConnection.state`. `connection.ok` is not in the OpenAPI spec yet, so its type is hand-written in `types.ts` and decoded through a shim in `StableWSConnection.ts` — both are marked for deletion once the backend publishes the event. + **It no longer registers `window` listeners.** The offline/online edge now arrives from `client.networkConnection` (see below), which inverts the old dependency: the socket is a _consumer_ of network status rather than the thing that detects it. Its own settings — connect timeout, `enableWSFallback`, ping interval, connection-check grace period, the UI's offline-notification delay, `webSocketImpl`, `urlParams` and an injectable socket — live in `client.wsConnection.config`, not in `StreamChatOptions`; `client.defaultWSTimeout` is gone. Every instance is created by `WSConnection.connect()`, which also disconnects the one it replaces, and read live from that config so an `updateConfig` reaches an already-open socket. It holds its `WSConnection` parent rather than the client, and reaches the client through it. - **`networkConnection/` — `NetworkConnectionObserver` (`client.networkConnection`).** The device's own network status, as a `StateStore`, written by a `NetworkStatusReporter` — a function that installs a platform listener and returns an unsubscribe, supplied as `statusReporter` config or through `setStatusReporter()`. `isOnline` is `boolean | undefined`, and `undefined` means _unknown_ rather than offline, so guards must test `=== false`. It is an accelerator, never a precondition — nothing in the SDK requires it. - **The default reporter is a compromise worth knowing.** A browser gets the real one. Every other host gets a stand-in that **mirrors the WebSocket**, so an integration that forgets the setup has a coarse signal rather than none; it reports nothing until the socket has been up once. Under it the two facts cannot disagree, which is exactly what the separation exists to express — so `isOnline === false` there also means "the socket died for its own reasons". React Native should wrap NetInfo and install it. See `docs/network-connection.md`. -- **`connection/ConnectionIdManager.ts` (`client.connectionIdManager`).** The WebSocket connection id, and the only place it lives — `client._getConnectionID()` / `_hasConnectionID()` read through it and keep no copy. The server keys watches and presence by it and answers `200` while registering nothing when it is missing, so `ApiClient._doRequest` holds any request `requiresConnectionId` recognises until one exists — the `watch` / `presence` flags, falling back to a declared `connection_id` param for `stopWatchingChannel` and `longPoll`, which carry no flag; `test/unit/codegen/connectionIdEndpoints.test.ts` pins that generated endpoint set. `StableWSConnection` drives the lifecycle: `arm()` before a socket opens, `resolveConnectionId()` on the hello frame, `invalidate()` from `_applyHealth(false)` (every status transition, so the error paths are covered) and again from `disconnect()`, and a reject when a reconnect gives up. A deliberate close invalidates rather than rejects: `closeConnection()` exists for mobile backgrounding, so a request already waiting keeps waiting for the socket `openConnection()` will build. `reset()` is the harder variant that fails those waiters instead; nothing in the SDK calls it. A caller's abort signal reaches the wait, so an abandoned request does not hold it. -- **`wsConnection/` — `WSConnection` (`client.wsConnection`).** A stable wrapper created with the client and never replaced, owning the socket's reactive status (`state`: `isHealthy`, `lastHealthyAt`, `lastUnhealthyAt`; the connection id lives on `client.connectionIdManager`), its configuration, and the one network-status subscription. The live `StableWSConnection` hangs off `.connection` and _is_ replaced per connect, which is why the store cannot live there — it would be unreadable before `connectUser` and would strand subscribers on every `closeConnection()` → `openConnection()` cycle. The store is written on **every** transition, including `disconnect()` and the error paths. It is the only description of the socket's status: there is no companion event. + **The default reporter is a compromise worth knowing.** A browser gets the real one. Every other host gets a stand-in that **mirrors the WebSocket**, so an integration that forgets the setup has a coarse signal rather than none; it reports nothing until the socket has been up once. Under it the two facts cannot disagree, which is exactly what the separation exists to express — so `isOnline === false` there also means "the socket died for its own reasons". React Native should wrap NetInfo and install it. See `docs/network-connection.md`. `usesDefaultWSConnectionNetworkStatusReporter` (`@internal`) says whether that socket-mirroring default is installed; the `enableWSFallback` switch reads it so that the default's "offline" is not mistaken for a device that is really offline. +- **`connection/ConnectionIdManager.ts` (`client.connectionIdManager`).** The connection id — the WebSocket's, or the long-poll's after an `enableWSFallback` switch — and the only place the client reads it from: `client._getConnectionID()` / `_hasConnectionID()` read through it and keep no copy. The long-poll keeps no copy either: `WSConnection.disconnect()` reads the id before the socket's `disconnect()` invalidates it, and passes it to `WSConnectionFallback.disconnect()` for the close. The server keys watches and presence by it and answers `200` while registering nothing when it is missing, so `ApiClient._doRequest` holds any request `requiresConnectionId` recognises until one exists — the `watch` / `presence` flags, and nothing else. The flagless operations that declare `connection_id` set it themselves: `StreamChat.stopWatchingChannel` (overridden to always wait for the connection id as the gate would — through a handshake or reconnect, honouring the abort signal, and rejecting when there is no connection and none is being established — then send it) and the long-poll fallback; `test/unit/codegen/connectionIdEndpoints.test.ts` pins that generated endpoint set. `StableWSConnection` drives the lifecycle: `arm()` before a socket opens, `resolveConnectionId()` on the hello frame, `invalidate()` from `_applyHealth(false)` (every status transition, so the error paths are covered) and again from `disconnect()`, and a reject when a reconnect gives up. After an `enableWSFallback` switch `WSConnectionFallback` drives it instead, in the same shape: `arm()` in `connect()`, `resolveConnectionId()` before going healthy, `invalidate()` on going closed or disconnected, and a reject when a connect fails (unless cancelled) or a poll hits a non-retryable error. A deliberate close invalidates rather than rejects: `closeConnection()` exists for mobile backgrounding, so a request already waiting keeps waiting for the socket `openConnection()` will build. `reset()` is the harder variant that fails those waiters instead; nothing in the SDK calls it. A caller's abort signal reaches the wait, so an abandoned request does not hold it. +- **`wsConnection/` — `WSConnection` (`client.wsConnection`).** A stable wrapper created with the client and never replaced, owning the socket's reactive status (`state`: `isHealthy`, `lastHealthyAt`, `lastUnhealthyAt`; the connection id lives on `client.connectionIdManager`), its configuration, and the one network-status subscription. The live `StableWSConnection` hangs off `.connection` and _is_ replaced per connect, which is why the store cannot live there — it would be unreadable before `connectUser` and would strand subscribers on every `closeConnection()` → `openConnection()` cycle. The store is written on **every** transition, including `disconnect()` and the error paths. It is the only description of the socket's status: there is no companion event. It also performs the `enableWSFallback` switch and holds the long-poll on `.fallback`, which every later `connect()` reuses and `disconnect()` closes alongside the socket. After the switch, `WSConnectionFallback` writes the store through `_setStatus()` instead — so `isHealthy`, `ConnectionRecoveryManager` and the socket-mirroring network reporter all follow the long-poll — `isConnecting` and the network-status routing follow it too, and `.connection` keeps the disconnected socket. - **No `store.ts`** — `StateStore` comes from `@stream-io/state-store`, imported directly wherever it is used and **not** re-exported from the root barrel (see "State and subscription patterns" below). - **`signing.ts` — one function, `UserFromToken`.** Decodes a JWT payload with the global `atob` and returns `user_id`. Everything else this module used to hold was server-side (JWT minting via `jsonwebtoken`, webhook/SQS/SNS verification via `crypto` + `zlib`) and was removed along with those deps — see `v9-to-v10-migration-guide-server-side.md`. Do not reintroduce secret-holding or HMAC code here; that surface lives in `@stream-io/node-sdk`. - **`middleware.ts`** — `MiddlewareExecutor` (see "Middleware pipelines" below). Used by composer pipelines, not by client request lifecycle. - **`token_manager.ts`** — handles static tokens and async token providers. Tracks a `loadTokenPromise` so concurrent calls await the same fetch. The constructor takes no arguments: there is no `secret` and no local JWT signing — every token comes from the caller (a string or a `TokenProvider`). Anonymous users may have no token at all; anyone else without one now fails at `getToken()` rather than at `setTokenOrProvider()`. -- **Event types.** There is no `events.ts` / `EVENT_MAP` any more (removed in v10). Wire events come from the generated `WSEvent` union (`src/gen/models`) so adding one means regenerating rather than hand-editing. There is no runtime decoder layer any more: `--opt response_dates_as_number` types every server-sent date as the unix-nanosecond number the wire already carries, so `src/gen/model-decoders/` (including `event-decoder-mapping.ts`) is no longer emitted at all, and `src/connection/wsConnection/StableWSConnection.ts` decodes frames with a plain cast. The unit invariant and its traps live in `src/utils/time.ts`; the consumer-facing delta is `v9-to-v10-migration-guide-dates.md`. `src/types.ts` overlays two non-generated members onto the public `Event` union: `LocalEvent` — `channels.queried`, `connection.recovered`, `capabilities.changed`, `message.read_locally`, `offline_reactions.queried`, `live_location_sharing.*`, all dispatched client-side only and never received over the wire — and `ConnectedEvent` (`connection.ok`), which _is_ a wire event but is not published in the OpenAPI spec yet. +- **Event types.** There is no `events.ts` / `EVENT_MAP` any more (removed in v10). Wire events come from the generated `WSEvent` union (`src/gen/models`) so adding one means regenerating rather than hand-editing. There is no runtime decoder layer any more: `--opt response_dates_as_number` types every server-sent date as the unix-nanosecond number the wire already carries, so `src/gen/model-decoders/` (including `event-decoder-mapping.ts`) is no longer emitted at all, and `src/connection/wsConnection/StableWSConnection.ts` decodes frames with a plain cast. The unit invariant and its traps live in `src/utils/time.ts`; the consumer-facing delta is `v9-to-v10-migration-guide-dates.md`. `src/types.ts` overlays two non-generated members onto the public `Event` union: `LocalEvent` — `channels.queried`, `connection.recovered`, `capabilities.changed`, `message.read_locally`, `offline_reactions.queried`, `live_location_sharing.*`, `connection.fallback_activated` (the `enableWSFallback` switch to long-poll), all dispatched client-side only and never received over the wire — and `ConnectedEvent` (`connection.ok`), which _is_ a wire event but is not published in the OpenAPI spec yet. **There is no `connection.changed`.** Connectivity is published as two stores — `client.wsConnection.state` and `client.networkConnection.state` — and nothing else. Publishing it twice meant two descriptions that disagreed: the event was silent on `closeConnection()` and two error paths, and held a drop for five seconds. `connection.recovered` survives, reports that a recovery pass finished rather than a status, and carries no payload. A UI that renders a "connection lost" banner holds a drop for `wsConnection.offlineNotificationDisplayDelayMs` itself, which the event used to do. - **`uploadManager.ts` / `LiveLocationManager.ts` / `CooldownTimer.ts`** — feature controllers, each owns its own `StateStore` slice. - **Domain subsystems** (each a folder with its own `index.ts` barrel): diff --git a/docs/network-connection.md b/docs/network-connection.md index 55dfd26f4..d4a049dda 100644 --- a/docs/network-connection.md +++ b/docs/network-connection.md @@ -71,7 +71,8 @@ it. ### Everywhere else, install one — and read this if you do not A host that cannot answer the question gets a stand-in that **mirrors this client's WebSocket**, so an -integration that forgets the setup has a signal rather than nothing at all. +integration that forgets the setup has a signal rather than nothing at all. (It reads the WebSocket's +status store, so after an [`enableWSFallback`](#long-poll-fallback) switch it mirrors the long-poll.) It is a safety net, not a measurement, and it has one consequence you must know about. Under it, the device's status is the socket's status, so `isOnline === false` whenever the socket dies **for its own @@ -176,7 +177,7 @@ client.wsConnection.state.getLatestValue(); // { isHealthy: boolean, lastHealthyAt: Date | null, lastUnhealthyAt: Date | null } ``` -Two things about this store are worth knowing. +A few things about this store are worth knowing. **It records every transition**, including `disconnect()` — what `client.closeConnection()` calls, the documented mobile backgrounding path — and the socket's internal error paths. There is nothing it @@ -193,6 +194,22 @@ awaits for anything that watches a channel or subscribes to presence — so thos reconnect instead of going out keyed to a connection the server has closed. You rarely need to read it; `client.connectionIdManager.connectionId` is there if you do. +### Long-poll fallback + +With `enableWSFallback` on, a WebSocket that fails to connect with a network error — reported once +`connectTimeoutMs` runs out — makes the client switch to HTTP long-polling against `/api/v2/longpoll`. +It dispatches `connection.fallback_activated` with `mode: 'longpoll'` and stays on long-poll for the +rest of the client's life, `disconnectUser()` included. The switch is skipped when a network reporter +says the device is offline; the stand-in's "offline" does not count, since it only means the socket +is down. + +After the switch, **this same store describes the long-poll**: `client.wsConnection.isHealthy` and +`state` report whether it is up, `connection.recovered` follows its reconnects, and +`client.connectionIdManager` holds its connection id. The long-poll follows `client.networkConnection` +too — it closes when the device goes offline and reconnects when it comes back — except while +`closeConnection()` or `disconnectUser()` has closed it, until it is reconnected. v9 stopped listening +for good at its first disconnect. + ### Timing and transport settings ```ts @@ -200,6 +217,7 @@ client.config.set({ client: { wsConnection: { connectTimeoutMs: 15000, // how long connect() waits for the server's hello + enableWSFallback: false, // long-poll when the WebSocket cannot connect — see above pingIntervalMs: 25000, // how often a health-check ping goes out — 25s is also the maximum healthCheckGracePeriodMs: 10000, // extra room before the socket is declared dead offlineNotificationDisplayDelayMs: 5000, // how long a UI holds a drop before reporting it @@ -227,21 +245,27 @@ client.on('connection.recovered', () => { }); ``` +With `enableWSFallback` on there is also `connection.fallback_activated`, dispatched once, when the +client switches to long-polling. It reports the transport, not its status. + There is no `connection.changed`. Connectivity used to be published twice, as the stores and as that event, and the two disagreed: the event was silent on `closeConnection()` and two error paths, and held a drop for five seconds. Subscribe to whichever store you mean instead. ## Migration -| Before | Now | -| ----------------------------------------------------------- | ----------------------------------------------------------------------------- | -| `client.wsConnection.onlineStatusChanged(fakeDomEvent)` | `client.networkConnection.setStatus(isOnline)` | -| `client.on('connection.changed', …)` | `client.wsConnection.state` or `client.networkConnection.state`, subscribed | -| `connection.recovered`'s `connection` field | gone; the event reports the socket and carries no payload | -| `NetworkStatusListenerRegistrar`, `statusListenerRegistrar` | `NetworkStatusReporter`, `statusReporter` | -| `client.threads.state.lastConnectionDropAt` | `client.wsConnection.state.lastUnhealthyAt` | -| `client.defaultWSTimeout = 5000` | `client.config.set({ client: { wsConnection: { connectTimeoutMs: 5000 } } })` | -| `new StreamChat(key, { WebSocketImpl, wsUrlParams })` | `wsConnection` config: `webSocketImpl`, `urlParams` | +| Before | Now | +| ----------------------------------------------------------- | -------------------------------------------------------------------------------- | +| `client.wsConnection.onlineStatusChanged(fakeDomEvent)` | `client.networkConnection.setStatus(isOnline)` | +| `client.on('connection.changed', …)` | `client.wsConnection.state` or `client.networkConnection.state`, subscribed | +| `connection.recovered`'s `connection` field | gone; the event reports the socket and carries no payload | +| `NetworkStatusListenerRegistrar`, `statusListenerRegistrar` | `NetworkStatusReporter`, `statusReporter` | +| `client.threads.state.lastConnectionDropAt` | `client.wsConnection.state.lastUnhealthyAt` | +| `client.defaultWSTimeout = 5000` | `client.config.set({ client: { wsConnection: { connectTimeoutMs: 5000 } } })` | +| `new StreamChat(key, { WebSocketImpl, wsUrlParams })` | `wsConnection` config: `webSocketImpl`, `urlParams` | +| `new StreamChat(key, { enableWSFallback: true })` | `wsConnection` config: `enableWSFallback` | +| `client.on('transport.changed', …)` | `client.on('connection.fallback_activated', …)`; same `mode: 'longpoll'` payload | +| `client.defaultWSTimeoutWithFallback` | gone; the switch waits `connectTimeoutMs` | `utils.ts` also had an internal `isOnline()` helper and an `addConnectionEventListeners()` pair. None was exported from the package, so there is nothing to migrate — they are gone, and diff --git a/src/api-client.ts b/src/api-client.ts index 5c5e03d60..25bc355ef 100644 --- a/src/api-client.ts +++ b/src/api-client.ts @@ -13,6 +13,9 @@ const logger = chatLoggerSystem.getLogger('api-client'); const MULTIPART_CONTENT_TYPE = 'multipart/form-data'; +/** The WebSocket fallback's endpoint, which a local API serves on its own port. */ +const LONG_POLL_PATH = '/api/v2/longpoll'; + /** * Upload requests must not inherit the axios instance timeout (3s by default) or the size * caps - either would abort a large or slow upload. @@ -62,8 +65,8 @@ export class ApiClient { params: queryParams, headers: { 'Content-Type': requestContentType }, ...(isMultipart ? UPLOAD_REQUEST_DEFAULTS : {}), - // Keep this last so a caller-supplied signal wins - and keep it returning only the keys - // it owns, so it can never clobber the upload defaults above. + // Keep this last so a caller-supplied signal and timeout win - and keep it returning only + // the keys it owns, so it can never clobber the upload defaults above. ...toAxiosRequestConfig(options), }); } @@ -107,7 +110,12 @@ export class ApiClient { } } if (resolved.startsWith('/')) { - resolved = this.client.baseURL + resolved; + const baseURL = + resolved === LONG_POLL_PATH + ? // replace port if present for testing with local API + this.client.baseURL?.replace(':3030', ':8900') + : this.client.baseURL; + resolved = baseURL + resolved; } return resolved; } @@ -185,6 +193,7 @@ export class ApiClient { }; } + await this.client.tokenManager.tokenReady(); const initialRequestConfig = this.populateRequestConfigWithDefaults(additionalConfig); const clientRequestId = initialRequestConfig.headers?.[ @@ -265,45 +274,33 @@ export class ApiClient { /** * Whether a request registers a server-side subscription, and so must not be sent before the - * handshake has produced a connection id. + * handshake has produced a connection id: one with a `watch` or `presence` flag set. * * The server keys watches and presence by that id and answers `200` while registering nothing when it * is missing, so a request that races the handshake yields a channel that never receives an event. + * + * The two generated operations that declare `connection_id` without a flag set it themselves: + * `stopWatchingChannel` through `StreamChat`'s override, and `longPoll`'s endpoint through the + * long-poll fallback. `test/unit/codegen/connectionIdEndpoints.test.ts` pins that set, so a new one + * surfaces there. */ export const requiresConnectionId = ( params: Record | undefined, body: unknown, ) => { const payload = params?.payload as Record | undefined; - // Guarded rather than `body ?? undefined`: the `in` checks below throw on a string body. const requestBody = (typeof body === 'object' && body !== null ? body : undefined) as | Record | undefined; - if ( + return Boolean( params?.watch || params?.presence || payload?.watch || payload?.presence || requestBody?.watch || - requestBody?.presence - ) { - return true; - } - - // A flag that is present and false is a deliberate "do not subscribe", so it must not fall through - // to the parameter check below — that is what made an explicit `watch: false` wait for a socket. - if ( - [params, payload, requestBody].some( - (source) => source && ('watch' in source || 'presence' in source), - ) - ) { - return false; - } - - // Some operations subscribe without carrying a flag — stop-watching and long polling — and are - // recognised by the generated `connection_id` parameter instead. - return Boolean(params && 'connection_id' in params); + requestBody?.presence, + ); }; /** @@ -323,11 +320,15 @@ const isUsableAbortSignal = (signal: unknown): signal is AbortSignal => const toAxiosRequestConfig = ({ signal, onUploadProgress, + timeout, }: StreamRequestOptions = {}): AxiosRequestConfig => ({ signal: isUsableAbortSignal(signal) ? signal : undefined, // Same reasoning as `isUsableAbortSignal`: an options object revived from a persisted // offline-db task payload has lost its functions. onUploadProgress: typeof onUploadProgress === 'function' ? onUploadProgress : undefined, + // Only when set: an explicit `undefined` would clobber the upload defaults' `timeout: 0`, and + // axios would fall back to the instance timeout. + ...(typeof timeout === 'number' ? { timeout } : {}), }); const errorIsApiError = (error: unknown): error is AxiosError => { diff --git a/src/client.ts b/src/client.ts index 6b73fe71c..01dce73ef 100644 --- a/src/client.ts +++ b/src/client.ts @@ -185,7 +185,8 @@ export class StreamChat extends ChatApi { */ networkConnection: NetworkConnectionObserver; /** - * The WebSocket connection id, and the one place it lives. + * The connection id — the WebSocket's, or the long-poll's after an `enableWSFallback` + * switch — and the one place the rest of the client reads it from. * * The server keys channel watches and presence subscriptions by it, so a request carrying either * waits here for the handshake rather than racing it. See `requiresConnectionId` in @@ -638,8 +639,12 @@ export class StreamChat extends ChatApi { * So when your app goes to background, you can call `client.closeConnection`. * And when app comes back to foreground, call `client.openConnection`. * + * After an `enableWSFallback` switch it closes the long-poll as well, telling the server to + * close its connection id. + * * @param timeout - Max number of milliseconds to wait for the WebSocket close event before forcefully assuming * successful disconnection. See https://developer.mozilla.org/en-US/docs/Web/API/CloseEvent (optional). + * The long-poll's close request uses it as its timeout (2s when omitted). */ closeConnection = async (timeout?: number) => { this._resetAIStateOnActiveChannels(); @@ -668,7 +673,8 @@ export class StreamChat extends ChatApi { }; /** - * Creates a new WebSocket connection with the current user. + * Creates a new WebSocket connection with the current user. After an `enableWSFallback` switch it + * reconnects the long-poll instead: the client never goes back to the WebSocket. * * @returns The WebSocket connect promise, or an empty resolved promise if a connection is already active. */ @@ -806,9 +812,11 @@ export class StreamChat extends ChatApi { this._rejectPendingWsPromise(teardownReason); this.wsPromise = null; - this.connectionIdManager.rejectConnectionId(teardownReason); const closePromise = this.closeConnection(timeout); + // After starting the close, which reads the connection id synchronously: after an + // `enableWSFallback` switch the long-poll's close request has to carry it. + this.connectionIdManager.rejectConnectionId(teardownReason); for (const channel of Object.values(this.activeChannels)) { channel._disconnect(); @@ -832,10 +840,11 @@ export class StreamChat extends ChatApi { this.unsubscribeClientConfiguration?.(); this.unsubscribeClientConfiguration = undefined; - // Since we wipe all user data already, we should reset token manager as well + // Deferred so the long-poll close can still authenticate. By the time it settles a + // connectUser() / connectAnonymousUser() may have set the next user and their token. closePromise .finally(() => { - this.tokenManager.reset(); + if (!this.userId) this.tokenManager.reset(); }) .catch((err) => logger @@ -1302,10 +1311,13 @@ export class StreamChat extends ChatApi { * Only `Watching` is demoted: a channel the consumer stopped on purpose, or one that was torn * down, stays `NotWatching` and must not be resurrected by a reconnect. * - * Invoked from two places, because neither covers the other: `StableWSConnection._setHealth(false)` - * for an abnormal close/error, and `closeConnection()` for a deliberate shutdown (e.g. mobile - * backgrounding), whose `disconnect()` writes the status through `_applyHealth` and so never - * reaches `_setHealth`. + * Invoked from two places on the WebSocket, because neither covers the other: + * `StableWSConnection._setHealth(false)` for an abnormal close/error, and + * `closeConnection()` for a deliberate shutdown (e.g. mobile backgrounding), whose + * `disconnect()` writes the status through `_applyHealth` and so never reaches + * `_setHealth`. After an `enableWSFallback` switch the long-poll calls it too, from + * `WSConnectionFallback._setState()` whenever going closed or disconnected takes the + * status offline. */ _markActiveChannelsWatchInterrupted() { for (const cid in this.activeChannels) { @@ -1365,10 +1377,11 @@ export class StreamChat extends ChatApi { * Requests no longer settle against these: waiting for a connection id is the * {@link ConnectionIdManager}'s job, applied centrally in `ApiClient`. * - * Called by `StableWSConnection._reconnect()`. Recovery itself is owned by - * {@link ConnectionRecoveryManager}, which subscribes to the connection lifecycle and so covers - * every reconnect path — including `closeConnection()` → `openConnection()` (mobile backgrounding), - * which never reaches `_reconnect()` at all. + * Called by `StableWSConnection._reconnect()`, and by `WSConnectionFallback.connect(true)` + * for the long-poll's own reconnects. Recovery itself is owned by + * {@link ConnectionRecoveryManager}, which subscribes to the connection lifecycle and so + * covers every reconnect path — including `closeConnection()` → `openConnection()` (mobile + * backgrounding), which reaches neither. * * @internal */ @@ -1404,17 +1417,10 @@ export class StreamChat extends ChatApi { throw Error('Property clientId is not set'); } - try { - // `wsConnection` builds and owns the socket; the reconnection logic and the connect timeout - // (`config.connectTimeoutMs`) live in there. - return await this.wsConnection.connect(); - } catch (error) { - // A failure the socket does not retry leaves nothing else to settle the pending connection id. - if (!isWSFailure(error as APIError)) { - this.connectionIdManager.rejectConnectionId(error); - } - throw error; - } + // `wsConnection` builds and owns the socket — and, with `enableWSFallback`, the long-poll it + // switches to; the reconnection logic and the connect timeout (`config.connectTimeoutMs`) live + // in there. + return await this.wsConnection.connect(); } /** @@ -1434,6 +1440,35 @@ export class StreamChat extends ChatApi { return data; } + /** + * Stops watching a channel on this client's connection. + * + * It carries no `watch` flag, so the request layer does not hold it, but it is connection-scoped + * by definition: it tells the server which connection should stop watching. So it waits for this + * client's connection id exactly as a watching request does — through a handshake or a reconnect, + * until the caller's abort signal fires — and throws when there is no connection and none is being + * established. + * + * @param ...args - `[request, requestOptions]`. `request.connection_id` is replaced by this + * client's connection id. + * @returns The server response. + */ + override async stopWatchingChannel( + ...args: Parameters + ) { + const [request, requestOptions] = args; + const signal = requestOptions?.signal; + const connectionId = await this.connectionIdManager.getConnectionId( + // Check if signal is still usable - if it is read from offline DB, it has lost `addEventListener`. + typeof signal?.addEventListener === 'function' ? signal : undefined, + ); + + return super.stopWatchingChannel( + { ...request, connection_id: connectionId }, + requestOptions, + ); + } + /** * Queries channels and returns the full API response including top-level metadata such as * `predefined_filter`. @@ -2315,15 +2350,26 @@ export class StreamChat extends ChatApi { * @returns The JSON-encoded auth message. */ _buildWSAuthMessage = () => - JSON.stringify({ - // The server requires a non-empty token even for anonymous connections, but - // skips JWT parsing for any string that is not shaped like one. Anonymous users - // have no token, so send a placeholder the server accepts and ignores. - token: this.tokenManager.getToken() || 'anonymous', - // `connect()` rejects before reaching this when `_user` is unset. - user_details: this._user as ConnectUserDetailsRequest, - products: ['chat'], - } satisfies WSAuthMessage); + JSON.stringify(this._buildWSAuthPayload(this.tokenManager.getToken())); + + /** + * The auth message itself — what {@link _buildWSAuthMessage} encodes for the WebSocket, and what + * the long-poll fallback sends as `longPoll()`'s `json` query param. + * + * @private + * + * @param token - The token the message carries. The long-poll passes none: it authenticates + * through the request's `Authorization` header instead. + */ + _buildWSAuthPayload = (token?: string): WSAuthMessage => ({ + // The server requires a non-empty token, but skips JWT parsing for any string that is not + // shaped like one. Anonymous users have no token, and the long-poll does not send one, so both + // send a placeholder the server accepts and ignores. + token: token || 'anonymous', + // `connect()` rejects before reaching this when `_user` is unset. + user_details: this._user as ConnectUserDetailsRequest, + products: ['chat'], + }); /** * Queries poll answers. diff --git a/src/configuration/shape.ts b/src/configuration/shape.ts index 65be06a1d..4df952b93 100644 --- a/src/configuration/shape.ts +++ b/src/configuration/shape.ts @@ -507,6 +507,12 @@ const WS_CONNECTION_FIELDS: Record = { kind: 'value', type: 'number', }, + enableWSFallback: { + description: + 'Falls back to HTTP long-polling (`/api/v2/longpoll`) when the WebSocket cannot connect, for networks that block WebSockets. The WebSocket connects within `connectTimeoutMs`, as without the flag; lower it to switch sooner. Dispatches `connection.fallback_activated` when it switches, and stays on long-poll from then on. Defaults to `false`. Was `StreamChatOptions.enableWSFallback`.', + kind: 'value', + type: 'boolean', + }, offlineNotificationDisplayDelayMs: { description: 'How long a drop must last before a UI reports it. Most drops resolve in under a second, so a banner showing all of them makes a working application look broken. Nothing in this package waits on it: it is here so the UI SDKs share one value. Zero holds nothing back, landing on the next task.', @@ -574,7 +580,7 @@ const CLIENT_FIELDS: Record = { }, wsConnection: { description: - "This client's WebSocket: how long to wait for it, how often to ping it, how long it may go quiet before being declared dead, and how long a drop must last before a UI reports it.", + "This client's WebSocket: how long to wait for it, how often to ping it, how long it may go quiet before being declared dead, how long a drop must last before a UI reports it, and whether to fall back to long-polling when it cannot connect.", fields: WS_CONNECTION_FIELDS, kind: 'group', }, diff --git a/src/connection/ConnectionIdManager.ts b/src/connection/ConnectionIdManager.ts index 9b71c4371..e72f0bac7 100644 --- a/src/connection/ConnectionIdManager.ts +++ b/src/connection/ConnectionIdManager.ts @@ -5,7 +5,8 @@ import { chatLoggerSystem } from '../logger'; const logger = chatLoggerSystem.getLogger('connection'); /** - * Holds the WebSocket connection id and lets callers await one that is still being negotiated. + * Holds the connection id — the WebSocket's, or the long-poll's after an `enableWSFallback` + * switch — and lets callers await one that is still being negotiated. * * The server keys channel watches and presence subscriptions by connection id, and answers `200` * while registering nothing when a request that needs one arrives without it. Requests carrying such @@ -14,11 +15,14 @@ const logger = chatLoggerSystem.getLogger('connection'); * * The id lives here and nowhere else. A copy kept on the socket outlives the socket it belongs to, * and a copy in `client.wsConnection.state` cannot be invalidated by a socket that has already been - * replaced; either way requests go out keyed to a connection the server has torn down. + * replaced; either way requests go out keyed to a connection the server has torn down. The long-poll + * keeps no copy either: `WSConnection.disconnect()` reads the id before the socket's `disconnect()` + * drops it, and hands it to the long-poll's close. * - * {@link StableWSConnection} drives the whole lifecycle: {@link arm} before a socket opens, - * {@link resolveConnectionId} on the hello frame, {@link invalidate} when it drops, {@link reset} on - * a deliberate close. + * {@link StableWSConnection} drives the lifecycle: {@link arm} before a socket opens, + * {@link resolveConnectionId} on the hello frame, {@link invalidate} when it drops or is closed + * deliberately, {@link rejectConnectionId} when a reconnect gives up. After an `enableWSFallback` + * switch, `WSConnectionFallback` drives it the same way instead. */ export class ConnectionIdManager { connectionId?: string; diff --git a/src/connection/ConnectionRecoveryManager.ts b/src/connection/ConnectionRecoveryManager.ts index 663a66921..089e126b0 100644 --- a/src/connection/ConnectionRecoveryManager.ts +++ b/src/connection/ConnectionRecoveryManager.ts @@ -69,8 +69,12 @@ export const DEFAULT_CONNECTION_RECOVERY_MANAGER_CONFIG: ConnectionRecoveryManag * transition, is guaranteed to land after the replay and the sync rather than alongside them. * * Recovery deliberately keeps no "did we drop?" flag of its own: `ChannelWatchStatus.WasWatching` - * already records exactly that, written from both truthful hooks - * (`StableWSConnection._setHealth(false)` and `closeConnection()`). + * already records exactly that, written from the truthful hooks + * (`StableWSConnection._setHealth(false)` and `closeConnection()`, plus + * `WSConnectionFallback._setState()` after an `enableWSFallback` switch). + * + * The status store it follows is the long-poll's once `enableWSFallback` has switched to it, so a + * long-poll reconnect recovers the same way a WebSocket reconnect does. */ export class ConnectionRecoveryManager extends WithSubscriptions { client: StreamChat; @@ -349,10 +353,11 @@ export class ConnectionRecoveryManager extends WithSubscriptions { } // `connection.recovered` means "recovery finished". Dispatched from here so it fires on EVERY - // reconnect path — the removed `recoverState()` was only ever called by - // `StableWSConnection._reconnect()`, so a `closeConnection()` → `openConnection()` cycle (mobile - // backgrounding) never produced it. Consumers keying post-recovery work off this event — the UI - // SDKs' mark-read-on-catch-up among them — need it after the reload above, not before. + // reconnect path — the removed `recoverState()` was only ever called on an automatic + // reconnect (`StableWSConnection._reconnect()`, or v9's long-poll `connect(true)`), so a + // `closeConnection()` → `openConnection()` cycle (mobile backgrounding) never produced it. + // Consumers keying post-recovery work off this event — the UI SDKs' mark-read-on-catch-up + // among them — need it after the reload above, not before. this.client.dispatchEvent({ type: 'connection.recovered' }); }; diff --git a/src/connection/networkConnection/NetworkConnectionObserver.ts b/src/connection/networkConnection/NetworkConnectionObserver.ts index 3deeb55b5..076c210a1 100644 --- a/src/connection/networkConnection/NetworkConnectionObserver.ts +++ b/src/connection/networkConnection/NetworkConnectionObserver.ts @@ -3,7 +3,10 @@ import { StateStore } from '@stream-io/state-store'; import { ConfigController } from '../../configuration/ConfigController'; import { deepFreezeConfig } from '../../configuration/utils/deepFreezeConfig'; import { chatLoggerSystem } from '../../logger'; -import { getDefaultNetworkStatusReporter } from './reporters'; +import { + browserNetworkStatusReporter, + getDefaultNetworkStatusReporter, +} from './reporters'; import type { NetworkConnectionObserverConfig, NetworkConnectionState, @@ -176,6 +179,21 @@ export class NetworkConnectionObserver extends WithSubscriptions { return this.state.getLatestValue().isOnline; } + /** + * Whether the installed reporter is the host default that mirrors this client's WebSocket + * (`createWSConnectionNetworkStatusReporter`). Under it `isOnline === false` means "the socket is + * down", not "the device is offline". + * + * @internal + */ + public get usesDefaultWSConnectionNetworkStatusReporter(): boolean { + return ( + this.installedReporter !== undefined && + this.installedReporter === this.defaultReporter && + this.defaultReporter !== browserNetworkStatusReporter + ); + } + /** * Installs the platform listener that reports device network status, unsubscribing whatever was * there. Pass `null` to clear it, which leaves the last known status alone rather than reverting to diff --git a/src/connection/networkConnection/reporters.ts b/src/connection/networkConnection/reporters.ts index bd5f0e487..87d187a66 100644 --- a/src/connection/networkConnection/reporters.ts +++ b/src/connection/networkConnection/reporters.ts @@ -47,10 +47,10 @@ export const browserNetworkStatusReporter: NetworkStatusReporter = (onStatusChan * and right in the common case where the device really did lose its network and took the socket with * it. * - * Nothing is reported until the socket has been up once. `isOnline` on the WebSocket store is `false` - * from construction, and forwarding that would claim the device is offline before anything had been - * attempted — a fabricated reading, which is the one thing this module refuses to produce. Until - * then the device's status stays `undefined`, meaning unknown. + * Nothing is reported until the socket has been up once. `isHealthy` on the WebSocket store is + * `false` from construction, and forwarding that would claim the device is offline before + * anything had been attempted — a fabricated reading, which is the one thing this module refuses + * to produce. Until then the device's status stays `undefined`, meaning unknown. * * Install a real reporter wherever one exists. This one can only ever repeat what the socket already * said, so it cannot tell you that the network came back before the socket noticed, which is the diff --git a/src/connection/wsConnection/WSConnection.ts b/src/connection/wsConnection/WSConnection.ts index 3f365d3d7..b47616e33 100644 --- a/src/connection/wsConnection/WSConnection.ts +++ b/src/connection/wsConnection/WSConnection.ts @@ -1,6 +1,8 @@ import { StateStore } from '@stream-io/state-store'; import { WithSubscriptions } from '../../utils/WithSubscriptions'; import { StableWSConnection } from '../../connection'; +import { WSConnectionFallback } from './WSConnectionFallback'; +import { isWSFailure } from '../../errors'; import { ConfigController } from '../../configuration/ConfigController'; import { chatLoggerSystem } from '../../logger'; import { @@ -10,6 +12,7 @@ import { } from './config'; import type { WSConnectionConfig, WSConnectionState } from './types'; import type { StreamChat } from '../../client'; +import type { APIError } from '../../errors'; import type { ConnectAPIResponse } from '../../types'; import type { Unsubscribe } from '@stream-io/state-store'; @@ -35,6 +38,9 @@ const logger = chatLoggerSystem.getLogger('client'); * socket it would leak — `client.connect()` overwrites `connection` without disconnecting the * previous one, leaving a dead socket still listening and still able to call `_reconnect()` on * itself. + * + * **With `enableWSFallback`, the store can describe the long-poll instead.** Once {@link connect} + * has switched to `WSConnectionFallback` ({@link fallback}). */ export class WSConnection extends WithSubscriptions { state: StateStore; @@ -45,6 +51,14 @@ export class WSConnection extends WithSubscriptions { * @internal */ connection: StableWSConnection | null = null; + /** + * The long-poll, once `enableWSFallback` has switched to it; `undefined` until then. Built by + * {@link connect} and never cleared: as in v9, the client does not go back to the WebSocket, not + * even across `disconnectUser()`. + * + * @internal + */ + fallback?: WSConnectionFallback; /** * Readable by the socket this object owns, which reaches the client through its parent rather than * holding one of its own. Not part of the public surface — `client.wsConnection.client` is a @@ -139,24 +153,26 @@ export class WSConnection extends WithSubscriptions { } /** - * Is this WebSocket up. Read from {@link state} rather than from the current socket, so it - * survives the socket being replaced and answers `false` rather than throwing before the first - * connect. + * Is this WebSocket up — or, after an `enableWSFallback` switch, the long-poll. Read from + * {@link state} rather than from the current socket, so it survives the socket being replaced and + * answers `false` rather than throwing before the first connect. */ get isHealthy(): boolean { return this.state.getLatestValue().isHealthy; } - /** Whether a connection attempt is in flight. */ + /** Whether a connection attempt is in flight — the long-poll's after an `enableWSFallback` switch. */ get isConnecting(): boolean { - return this.connection?.isConnecting ?? false; + return (this.fallback ?? this.connection)?.isConnecting ?? false; } /** - * Records this WebSocket's status. Returns whether it changed. + * Records this connection's status. Returns whether it changed. * - * {@link StableWSConnection} is the only caller, and it routes **every** status transition through - * here — including `disconnect()`, which `closeConnection()` uses, and the two error paths. + * {@link StableWSConnection} routes **every** status transition through here — including + * `disconnect()`, which `closeConnection()` uses, and the two error paths. The one other + * caller is {@link fallback}, which takes over once `enableWSFallback` has switched to + * long-polling. * * @internal */ @@ -180,6 +196,9 @@ export class WSConnection extends WithSubscriptions { * * Read from the store rather than from an event, so there is one description of the device's * network rather than two that can disagree. + * + * After an `enableWSFallback` switch it routes to {@link fallback} instead, leaving the + * disconnected socket out of it. */ public registerSubscriptions = (): Unsubscribe => { if (!this.hasSubscriptions) { @@ -190,7 +209,7 @@ export class WSConnection extends WithSubscriptions { // `undefined` is *unknown*, not offline: no reporter has said anything yet, and acting on // it would tear down a healthy socket on every host that cannot answer the question. if (typeof isOnline !== 'boolean') return; - this.connection?._applyNetworkStatus(isOnline); + (this.fallback ?? this.connection)?._applyNetworkStatus(isOnline); }, ), ); @@ -211,8 +230,16 @@ export class WSConnection extends WithSubscriptions { * Creation lives here rather than in the constructor because a client that never calls * `connectUser` should never build a socket, and here rather than in `client.connect()` because a * field belongs to the object that owns it. + * + * With `enableWSFallback`, a socket that fails with a network error is replaced by the long-poll + * ({@link fallback}), which every later call connects instead of building a socket. */ - connect(timeout?: number): ConnectAPIResponse | undefined { + async connect(timeout?: number): ConnectAPIResponse { + // if fallback is used before, continue using it instead of waiting for WS to fail + if (this.fallback) { + return await this.fallback.connect(); + } + const next = this.buildConnection(); const previous = this.connection; @@ -239,8 +266,42 @@ export class WSConnection extends WithSubscriptions { } this.connection = next; - // Left to the socket to default from `config.connectTimeoutMs`, so the value lives in one place. - return next.connect(timeout); + + try { + // Left to the socket to default from `config.connectTimeoutMs`, so the value lives in one place. + return await next.connect(timeout); + } catch (error) { + // run fallback only if it's WS/Network error and not a normal API error + // make sure the device is online before even trying the longpoll. The default reporter on + // hosts without a network API mirrors the WebSocket, so its "offline" only means the socket + // is down, and is ignored. + const { networkConnection } = this.client; + const isDeviceOffline = + networkConnection.isOnline === false && + !networkConnection.usesDefaultWSConnectionNetworkStatusReporter; + if ( + this.config.enableWSFallback && + isWSFailure(error as APIError) && + !isDeviceOffline + ) { + logger.withExtraTags('connect').info('WS failed, fallback to longpoll'); + this.client.dispatchEvent({ + type: 'connection.fallback_activated', + mode: 'longpoll', + }); + + next._destroyCurrentWSConnection(); + void next.disconnect(); // close WS so no retry + this.fallback = new WSConnectionFallback({ client: this.client }); + return await this.fallback.connect(); + } + + // A failure the socket does not retry leaves nothing else to settle the pending connection id. + if (!isWSFailure(error as APIError)) { + this.client.connectionIdManager.rejectConnectionId(error); + } + throw error; + } } /** @@ -258,8 +319,19 @@ export class WSConnection extends WithSubscriptions { return new StableWSConnection({ wsConnection: this }); } - disconnect(timeout?: number): Promise | undefined { - return this.connection?.disconnect(timeout); + /** + * Closes the socket and, after an `enableWSFallback` switch, the long-poll too, telling the + * server to close its connection id. The long-poll's close request uses `timeout` as its timeout + * (2s when omitted). + */ + async disconnect(timeout?: number): Promise { + // Read before the socket's `disconnect()` drops it: after a switch it is the long-poll's id, + // which the long-poll's close request has to carry. + const connectionId = this.client.connectionIdManager.connectionId; + await Promise.all([ + this.connection?.disconnect(timeout), + this.fallback?.disconnect(timeout, connectionId), + ]); } /** diff --git a/src/connection/wsConnection/WSConnectionFallback.ts b/src/connection/wsConnection/WSConnectionFallback.ts new file mode 100644 index 000000000..7bd2fa7c6 --- /dev/null +++ b/src/connection/wsConnection/WSConnectionFallback.ts @@ -0,0 +1,259 @@ +import axios, { CanceledError } from 'axios'; +import type { StreamChat } from '../../client'; +import { retryInterval, sleep } from '../../utils'; +import { isAPIError, isConnectionIDError, isErrorRetryable } from '../../errors'; +import { chatLoggerSystem } from '../../logger'; +import type { LogLevel } from '../../logger'; +import type { ConnectionOpen, Event, StreamRequestOptions } from '../../types'; + +const logger = chatLoggerSystem.getLogger('connection'); + +export enum WSFallbackConnectionState { + Closed = 'CLOSED', + Connected = 'CONNECTED', + Connecting = 'CONNECTING', + Disconnected = 'DISCONNECTED', + Init = 'INIT', +} + +export class WSConnectionFallback { + client: StreamChat; + state: WSFallbackConnectionState; + consecutiveFailures: number; + abortController?: AbortController; + + constructor({ client }: { client: StreamChat }) { + this.client = client; + this.state = WSFallbackConnectionState.Init; + this.consecutiveFailures = 0; + } + + /** + * Applies a change in the device's network status. Stands in for v9's `window` online/offline + * listeners: `WSConnection` routes the network-status store here once it has switched to this + * long-poll. + * + * @internal + */ + _applyNetworkStatus = (online: boolean) => { + // Closed on purpose by `disconnect()`: going offline would move it to `Closed`, and the next + // online would reconnect it. + if (this.state === WSFallbackConnectionState.Disconnected) return; + + this._log(`_applyNetworkStatus() - ${online ? 'online' : 'offline'}`); + + if (!online) { + this._setState(WSFallbackConnectionState.Closed); + this.abortController?.abort(); + this.abortController = undefined; + return; + } + + if (this.state === WSFallbackConnectionState.Closed) { + this.connect(true); + } + }; + + /** + * connect try to open a longpoll request + * + * @param reconnect - should be false for first call and true for subsequent calls to keep the connection alive and settle the connect promises + * + * @internal + */ + connect = async (reconnect = false) => { + if (this.state === WSFallbackConnectionState.Connecting) { + this._log('connect() - connecting already in progress', { reconnect }, 'warn'); + return; + } + if (this.state === WSFallbackConnectionState.Connected) { + this._log('connect() - already connected and polling', { reconnect }, 'warn'); + return; + } + + this._setState(WSFallbackConnectionState.Connecting); + // Before anything awaits, so a watch issued now waits for this connection's id. + this.client.connectionIdManager.arm(); + try { + const { event } = await this._req<{ event: ConnectionOpen }>( + // Authenticated by the request's `Authorization` header, so the message carries only a + // placeholder token. + { json: this.client._buildWSAuthPayload() }, + { timeout: 8000 }, // 8s + reconnect, + ); + + // The id is published before the status says the connection is up, as the WebSocket does. + this.client.connectionIdManager.resolveConnectionId(event.connection_id); + this._setState(WSFallbackConnectionState.Connected); + this.client.dispatchEvent(event); + this._poll(); + if (reconnect) { + this.client._settleConnectPromises(); + } + return event; + } catch (err) { + // `disconnect()` has already set the state of an attempt it cancelled. Overwriting it with + // Closed would let the next online edge in `_applyNetworkStatus` reconnect it. + if (this.state !== WSFallbackConnectionState.Disconnected) + this._setState(WSFallbackConnectionState.Closed); + // Nothing retries a failed connect, so fail whatever is waiting for a connection id. A cancel + // comes from `disconnect()`, which leaves them waiting for the next connection instead. + if (!axios.isCancel(err)) this.client.connectionIdManager.rejectConnectionId(err); + throw err; + } + }; + + /** + * Whether a connect request is in flight. + * + * @internal + */ + get isConnecting() { + return this.state === WSFallbackConnectionState.Connecting; + } + + /** + * Stops polling and tells the server to close `connectionId`. + * + * @param timeout - The close request's timeout, in milliseconds. + * @param connectionId - The id to close. Passed in rather than read from + * `client.connectionIdManager`, because `WSConnection.disconnect()` disconnects the old socket + * first, and that drops the manager's id. + * + * @internal + */ + disconnect = async (timeout = 2000, connectionId?: string) => { + this._setState(WSFallbackConnectionState.Disconnected); + this.abortController?.abort(); + this.abortController = undefined; + + try { + await this._req({ close: true, connection_id: connectionId }, { timeout }, false); + this._log(`disconnect() - Closed connectionID`); + } catch (err) { + this._log(`disconnect() - Failed`, { err }, 'error'); + } + }; + + private _log( + msg: string, + extra: Record = {}, + level: LogLevel = 'info', + ) { + const log = logger.withExtraTags('connection_fallback'); + log[level]('WSConnectionFallback:' + msg, extra); + } + + private _setState(state: WSFallbackConnectionState) { + this._log(`_setState() - ${state}`); + + const previous = this.state; + this.state = state; + + // transition from connecting => connected + if ( + previous === WSFallbackConnectionState.Connecting && + this.state === WSFallbackConnectionState.Connected + ) { + this.client.wsConnection._setStatus({ isHealthy: true }); + } + + if ( + this.state === WSFallbackConnectionState.Closed || + this.state === WSFallbackConnectionState.Disconnected + ) { + // The server keyed watches by this connection id, so no request may carry it any more. + this.client.connectionIdManager.invalidate(); + if (this.client.wsConnection._setStatus({ isHealthy: false })) { + this.client._markActiveChannelsWatchInterrupted(); + } + } + } + + private _req = async >( + params: NonNullable[0]>, + config: Pick, + retry: boolean, + ): Promise => { + if (!this.abortController && !params.close) { + this.abortController = new AbortController(); + } + + try { + const res = await this.client.longPoll(params, { + ...config, + signal: this.abortController?.signal, + }); + + this.consecutiveFailures = 0; // always reset in case of no error + // The spec declares no response body, so the generated response type is `{}`. + return res as T; + } catch (error: any) { + this.consecutiveFailures += 1; + + if (retry && isErrorRetryable(error)) { + this._log(`_req() - Retryable error, retrying request`); + await sleep(retryInterval(this.consecutiveFailures)); + // A reconnect, `disconnect()` or going offline during the sleep dropped the connection id + // this request carries. Sent anyway, it would come back as a ConnectionIDNotFoundError, and + // `_poll` would tear down the connection that replaced it. + if ( + params.connection_id && + params.connection_id !== this.client.connectionIdManager.connectionId + ) { + this._log(`_req() - Connection id changed, dropping the retry`); + throw new CanceledError( + 'The connection id changed while the retry was waiting', + ); + } + return this._req(params, config, retry); + } + + throw error; + } + }; + + private _poll = async () => { + while (this.state === WSFallbackConnectionState.Connected) { + try { + const data = await this._req<{ + events: Event[]; + }>( + { connection_id: this.client.connectionIdManager.connectionId }, + { timeout: 30000 }, // 30s => API responds in 20s if there is no event + true, + ); + + if (data.events?.length) { + for (let i = 0; i < data.events.length; i++) { + this.client.dispatchEvent(data.events[i]); + } + } + } catch (error: any) { + if (axios.isCancel(error)) { + this._log(`_poll() - axios canceled request`); + return; + } + + /** client.longPoll's request layer will take care of TOKEN_EXPIRED error */ + + if (isConnectionIDError(error)) { + this._log(`_poll() - ConnectionID error, connecting without ID...`); + this._setState(WSFallbackConnectionState.Disconnected); + this.connect(true); + return; + } + + if (isAPIError(error) && !isErrorRetryable(error)) { + this._setState(WSFallbackConnectionState.Closed); + // Nothing reconnects from here, so fail whatever is waiting for a connection id. + this.client.connectionIdManager.rejectConnectionId(error); + return; + } + + await sleep(retryInterval(this.consecutiveFailures)); + } + } + }; +} diff --git a/src/connection/wsConnection/config.ts b/src/connection/wsConnection/config.ts index 8a5981ec6..7ac2b3fe9 100644 --- a/src/connection/wsConnection/config.ts +++ b/src/connection/wsConnection/config.ts @@ -10,6 +10,7 @@ import type { WSConnectionConfig } from './types'; */ export const DEFAULT_WS_CONNECTION_CONFIG: WSConnectionConfig = deepFreezeConfig({ connectTimeoutMs: 15 * 1000, + enableWSFallback: false, pingIntervalMs: 25 * 1000, healthCheckGracePeriodMs: 10 * 1000, offlineNotificationDisplayDelayMs: 5 * 1000, diff --git a/src/connection/wsConnection/types.ts b/src/connection/wsConnection/types.ts index bb0e44f29..b2a42a6e4 100644 --- a/src/connection/wsConnection/types.ts +++ b/src/connection/wsConnection/types.ts @@ -2,7 +2,8 @@ import type { StableWSConnection } from '../../connection'; export type WSConnectionState = { /** - * Is this client's WebSocket up. + * Is this client's WebSocket up. After an `enableWSFallback` switch, whether its long-poll is up, + * since the long-poll writes this store from then on. * * Deliberately the same field name as `client.networkConnection.state`'s — both answer the same * question about a different connection, which is why the two objects are named as parallels. The @@ -15,7 +16,9 @@ export type WSConnectionState = { }; /** - * The WebSocket's timing knobs. + * The WebSocket's settings: its timing knobs, what it is built from (`webSocketImpl`, + * `urlParams`, `connection`), and {@link enableWSFallback}, the switch to long-polling when it + * cannot connect. * * Declared, validated, and durable across reconnects. {@link offlineNotificationDisplayDelayMs} is * the odd one out: the only field here this package does not act on itself. @@ -23,8 +26,8 @@ export type WSConnectionState = { * The network-recovery retry is deliberately **not** configurable, and stays a constant in * `config.ts`: nothing could reach it, so exposing it would add surface rather than preserve it. * - * All in **milliseconds**, and named with the unit, because a bare `pingInterval` reads equally well - * as seconds. + * The timings are all in **milliseconds**, and named with the unit, because a bare `pingInterval` + * reads equally well as seconds. */ export type WSConnectionConfig = { /** @@ -32,6 +35,16 @@ export type WSConnectionConfig = { * of 15s allows between two and three attempts of the underlying retry. */ connectTimeoutMs: number; + /** + * Whether to enable the WebSocket fallback mechanism. Only enable this feature if you expect clients to be in environments where WebSocket connections might be blocked. Most integrators shouldn't need to turn on this flag. + * + * Falls back to HTTP long-polling when the WebSocket cannot connect, for networks that block + * WebSockets. Defaults to `false`. The WebSocket connects within {@link connectTimeoutMs}, as + * without the flag; lower it to switch sooner. The client dispatches + * `connection.fallback_activated` with `mode: 'longpoll'` when it switches, and stays on long-poll + * from then on. + */ + enableWSFallback: boolean; /** * How often a health-check ping goes out while the socket is up. * diff --git a/src/errors.ts b/src/errors.ts index 59bec2230..6ca75988e 100644 --- a/src/errors.ts +++ b/src/errors.ts @@ -64,6 +64,10 @@ export function isErrorRetryable(error: APIError) { return err.retryable; } +export function isConnectionIDError(error: APIError) { + return error.code === 46; // ConnectionIDNotFoundError +} + /** * Whether an error is EPHEMERAL — a transient failure worth queueing/retrying rather than a * definitive rejection. True when the server never responded (connection/network/offline error - no diff --git a/src/gen/chat/ChatApi.ts b/src/gen/chat/ChatApi.ts index 81cef6433..5f0242fc5 100644 --- a/src/gen/chat/ChatApi.ts +++ b/src/gen/chat/ChatApi.ts @@ -2067,10 +2067,11 @@ export class ChatApi { } async longPoll( - request?: { connection_id?: string; json?: WSAuthMessage }, + request?: { close?: boolean; connection_id?: string; json?: WSAuthMessage }, requestOptions?: StreamRequestOptions, ): Promise> { const queryParams = { + close: request?.close, connection_id: request?.connection_id, json: request?.json, }; diff --git a/src/types.ts b/src/types.ts index 53ac8ffc6..172d44c0c 100644 --- a/src/types.ts +++ b/src/types.ts @@ -209,6 +209,8 @@ type LocalEvent = ( }; }) | { type: 'connection.recovered' } + // `enableWSFallback` switched from the WebSocket to long-polling. + | ({ type: 'connection.fallback_activated' } & { mode: string }) | ({ type: 'offline_reactions.queried' } & { offlineReactions: ReactionResponse[]; }) @@ -526,6 +528,10 @@ export type StreamRequestOptions = { signal?: AbortSignal; /** Only meaningful for upload (multipart) requests; ignored everywhere else. */ onUploadProgress?: (event: StreamProgressEvent) => void; + /** + * Milliseconds before this request is aborted. + */ + timeout?: number; }; export * from './gen/models'; diff --git a/test/unit/ConnectionRecoveryManager.test.ts b/test/unit/ConnectionRecoveryManager.test.ts index 0cce1951b..4b211e6d3 100644 --- a/test/unit/ConnectionRecoveryManager.test.ts +++ b/test/unit/ConnectionRecoveryManager.test.ts @@ -253,7 +253,8 @@ describe('ConnectionRecoveryManager', () => { describe('connection.recovered', () => { it('is dispatched after the reload, on a path `recoverState()` never runs on', async () => { - // `recoverState()` is only ever called by `StableWSConnection._reconnect()`, so a + // `recoverState()` was only ever called on an automatic reconnect (the socket's + // `_reconnect()`, or v9's long-poll `connect(true)`), so a // `closeConnection()` → `openConnection()` cycle (mobile backgrounding) used to produce no // `connection.recovered` at all. Consumers that key post-recovery work off it — marking a // caught-up channel read, for one — silently did nothing there. diff --git a/test/unit/api-client.test.ts b/test/unit/api-client.test.ts index 902bc5ce3..d78d91c84 100644 --- a/test/unit/api-client.test.ts +++ b/test/unit/api-client.test.ts @@ -5,7 +5,7 @@ import { beforeEach, describe, expect, it, vi } from 'vitest'; import { requiresConnectionId } from '../../src/api-client'; import { getClientWithUser } from './test-utils/getClient'; -import type { StreamChat } from '../../src/client'; +import { StreamChat } from '../../src/client'; import type { StreamRequestOptions } from '../../src/types'; describe('ApiClient request options', () => { @@ -125,6 +125,19 @@ describe('ApiClient request options', () => { expect((firstConfig() as { onUploadProgress?: unknown }).onUploadProgress).to.be .undefined; }); + + it('forwards the timeout from the request options to axios', async () => { + await sendRequest({ timeout: 30000 } as StreamRequestOptions); + + expect((firstConfig() as AxiosRequestConfig).timeout).to.equal(30000); + }); + + // An explicit `undefined` makes axios fall back to the instance timeout. + it('sends no timeout key when none is given', async () => { + await sendRequest({ signal: new AbortController().signal }); + + expect(firstConfig()).not.toHaveProperty('timeout'); + }); }); describe('ApiClient rate limit metadata', () => { @@ -390,6 +403,59 @@ describe('ApiClient multipart encoding', () => { expect(firstConfig().signal).to.equal(controller.signal); expect(firstConfig().timeout).to.equal(0); }); + + it('lets a caller timeout override the upload default', async () => { + await client.api.sendRequest( + 'POST', + '/api/v2/uploads/file', + undefined, + undefined, + { file: new File(['x'], 'a.jpg', { type: 'image/jpeg' }) }, + 'multipart/form-data', + { timeout: 1000 }, + ); + + expect(firstConfig().timeout).to.equal(1000); + expect(firstConfig().maxContentLength).to.equal(Infinity); + expect(firstConfig().maxBodyLength).to.equal(Infinity); + }); +}); + +describe('ApiClient long-poll URL', () => { + let client: StreamChat; + let requestSpy: ReturnType; + + const urlOfCall = (index: number) => + (requestSpy.mock.calls[index][0] as AxiosRequestConfig).url; + + beforeEach(() => { + client = getClientWithUser(); + requestSpy = vi + .spyOn(client.axiosInstance, 'request') + .mockResolvedValue({ data: {}, status: 200, headers: {} }); + }); + + it('sends the long-poll to its own port on a local API', async () => { + client.setBaseURL('http://localhost:3030'); + + await client.longPoll({ close: true, connection_id: 'id' }); + await client.api.sendRequest('GET', '/api/v2/chat/channels'); + + expect(urlOfCall(0)).to.equal('http://localhost:8900/api/v2/longpoll'); + expect((requestSpy.mock.calls[0][0] as AxiosRequestConfig).params).toMatchObject({ + close: true, + connection_id: 'id', + }); + expect(urlOfCall(1)).to.equal('http://localhost:3030/api/v2/chat/channels'); + }); + + it('leaves a production base URL alone', async () => { + client.setBaseURL('https://chat.stream-io-api.com'); + + await client.longPoll({ connection_id: 'id' }); + + expect(urlOfCall(0)).to.equal('https://chat.stream-io-api.com/api/v2/longpoll'); + }); }); describe('upload methods', () => { @@ -499,9 +565,10 @@ describe('ApiClient connection id gate', () => { const flush = () => new Promise((resolve) => setTimeout(resolve, 0)); describe('requiresConnectionId', () => { - it('gates on a declared connection_id query param, even without a watch flag', () => { - // the only signal stopWatchingChannel and longPoll give - expect(requiresConnectionId({ connection_id: undefined }, undefined)).to.be.true; + it('does not gate on a declared connection_id query param alone', () => { + // stopWatchingChannel and the long-poll set theirs themselves + expect(requiresConnectionId({ connection_id: undefined }, undefined)).to.be.false; + expect(requiresConnectionId({ connection_id: 'id' }, undefined)).to.be.false; }); it('gates on a watch or presence flag in the request body', () => { @@ -529,8 +596,9 @@ describe('ApiClient connection id gate', () => { .false; }); - // The regression the `connection_id` fallback used to cause: the generator emits the key for - // every operation that *can* watch, so gating on its presence gated `watch: false` too. + // The regression gating on a declared `connection_id` once caused, before the gate stopped + // reading it: the generator emits the key for every operation that *can* watch, so gating + // on its presence gated `watch: false` too. it('does not gate a declared connection_id when the request opted out of watching', () => { expect(requiresConnectionId({ connection_id: undefined }, { watch: false })).to.be .false; @@ -571,7 +639,7 @@ describe('ApiClient connection id gate', () => { expect(sentParams().connection_id).to.equal('late-id'); }); - it('gates queryThreads, getThread, sync and stopWatching, which had no gate before', async () => { + it('gates queryThreads, getThread and sync, which had no gate before', async () => { const gated = [ () => client.queryThreads({ watch: true }), () => client.getThread({ message_id: 'mid', watch: true }), @@ -581,8 +649,6 @@ describe('ApiClient connection id gate', () => { last_sync_at: new Date(), watch: true, }), - // stopWatching carries no flag at all - the declared `connection_id` is its only signal - () => client.channel('messaging', 'id').stopWatching(), ]; for (const call of gated) { @@ -600,6 +666,51 @@ describe('ApiClient connection id gate', () => { } }); + it('sends stopWatching with the current connection id', async () => { + client.connectionIdManager.resolveConnectionId('current-id'); + + await client.channel('messaging', 'id').stopWatching(); + + expect(sentParams().connection_id).to.equal('current-id'); + }); + + it('holds stopWatching while a connection is being established, then sends its id', async () => { + client.connectionIdManager.reset(); + client.connectionIdManager.arm(); + + const inFlight = client.channel('messaging', 'id').stopWatching(); + await flush(); + expect(requestSpy).not.toHaveBeenCalled(); + + client.connectionIdManager.resolveConnectionId('late-id'); + await inFlight; + expect(requestSpy).toHaveBeenCalledTimes(1); + expect(sentParams().connection_id).to.equal('late-id'); + }); + + it('rejects stopWatching when there is no id and nothing in flight', async () => { + client.connectionIdManager.reset(); + + await expect(client.channel('messaging', 'id').stopWatching()).rejects.toThrow( + 'No connection id is available', + ); + expect(requestSpy).not.toHaveBeenCalled(); + }); + + it('abandons the stopWatching wait when the caller aborts', async () => { + client.connectionIdManager.reset(); + client.connectionIdManager.arm(); + const controller = new AbortController(); + + const inFlight = client + .channel('messaging', 'id') + .stopWatching({}, { signal: controller.signal }); + controller.abort(); + + await expect(inFlight).rejects.toThrow(); + expect(requestSpy).not.toHaveBeenCalled(); + }); + it('lets a request that needs no id through while the handshake is still in flight', async () => { client.connectionIdManager.reset(); client.connectionIdManager.arm(); @@ -661,3 +772,50 @@ describe('ApiClient connection id gate', () => { expect(requestSpy).not.toHaveBeenCalled(); }); }); + +describe('ApiClient token loading', () => { + let client: StreamChat; + let requestSpy: ReturnType; + + const sentAuthorization = () => + (requestSpy.mock.calls[0][0] as AxiosRequestConfig).headers?.Authorization; + const sendRequest = () => client.api.sendRequest('GET', '/api/v2/chat/channels'); + + beforeEach(() => { + client = new StreamChat('key'); + requestSpy = vi + .spyOn(client.axiosInstance, 'request') + .mockResolvedValue({ data: {}, status: 200, headers: {} }); + }); + + it('waits for a token provider that is still loading', async () => { + let provideToken: (token: string) => void = () => undefined; + client._setToken( + { id: 'amin' }, + () => new Promise((resolve) => (provideToken = resolve)), + ); + + const request = sendRequest(); + await Promise.resolve(); + expect(requestSpy).not.toHaveBeenCalled(); + + provideToken('provided-token'); + await request; + + expect(sentAuthorization()).toBe('provided-token'); + }); + + it("rejects with the provider's error when the provider fails", async () => { + client + ._setToken({ id: 'amin' }, () => Promise.reject(new Error('provider down'))) + .catch(() => undefined); + + await expect(sendRequest()).rejects.toThrow(/Call to tokenProvider failed/); + expect(requestSpy).not.toHaveBeenCalled(); + }); + + it('still rejects when no token was ever set', async () => { + await expect(sendRequest()).rejects.toThrow(/User token is not set/); + expect(requestSpy).not.toHaveBeenCalled(); + }); +}); diff --git a/test/unit/client.test.js b/test/unit/client.test.js index 52f770833..3350cf722 100644 --- a/test/unit/client.test.js +++ b/test/unit/client.test.js @@ -1,4 +1,5 @@ import sinon from 'sinon'; +import { CanceledError } from 'axios'; import { generateMsg } from './test-utils/generateMessage'; import { getClientWithUser } from './test-utils/getClient'; @@ -7,6 +8,7 @@ import { StreamChat } from '../../src/client'; import { ChatApi } from '../../src/gen-imports'; import { chatLoggerSystem } from '../../src/logger'; import { StableWSConnection } from '../../src/connection'; +import { WSFallbackConnectionState } from '../../src/connection/wsConnection/WSConnectionFallback'; import { mockChannelQueryResponse } from './test-utils/mockChannelQueryResponse'; import { generateThreadResponse } from './test-utils/generateThreadResponse'; import { @@ -991,6 +993,58 @@ describe('Client disconnectUser', () => { expect(client.tokenManager.reset.called).to.be.true; }); + describe('when a user connects before the close settles', () => { + let client; + let resolveClose; + + beforeEach(() => { + client = new StreamChat('key', { allowServerSideConnect: true }); + // the real connectUser, without a socket + client.openConnection = () => Promise.resolve(); + const { resolve, promise } = Promise.withResolvers(); + resolveClose = resolve; + client.wsConnection = { disconnect: () => promise }; + }); + + it('keeps the next user token', async () => { + await client.connectUser({ id: 'a' }, async () => 'token-a'); + + const disconnectPromise = client.disconnectUser(); + await client.connectUser({ id: 'b' }, async () => 'token-b'); + resolveClose(); + await disconnectPromise; + + expect(client.tokenManager.token).to.equal('token-b'); + expect(client.tokenManager.user.id).to.equal('b'); + }); + + it('keeps the token of the same user reconnecting', async () => { + const user = { id: 'a' }; + const tokenProvider = async () => 'token-a'; + await client.connectUser(user, tokenProvider); + + const disconnectPromise = client.disconnectUser(); + await client.connectUser(user, tokenProvider); + resolveClose(); + await disconnectPromise; + + expect(client.tokenManager.token).to.equal('token-a'); + expect(client.tokenManager.user).to.equal(user); + }); + + it('keeps an anonymous user anonymous', async () => { + await client.connectUser({ id: 'a' }, async () => 'token-a'); + + const disconnectPromise = client.disconnectUser(); + await client.connectAnonymousUser(); + resolveClose(); + await disconnectPromise; + + expect(client.anonymous).to.be.true; + expect(client.getAuthType()).to.equal('anonymous'); + }); + }); + it('should clear upload manager records', async () => { const client = new StreamChat('', ''); client.uploadManager.state.next(() => ({ @@ -2446,3 +2500,234 @@ describe('_normalizeExpiration', () => { ); }); }); + +describe('Client WSFallback', () => { + const userToken = + 'eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJ1c2VyX2lkIjoiYW1pbiJ9.1R88K_f1CC2yrR6j1_OzMEbasfS_dxRSNbundEDBlJI'; + const wsFailure = () => + Object.assign(new Error('initial WS connection could not be established'), { + isWSFailure: true, + }); + let client; + let socketConnect; + + /** + * Answers the long-poll at the axios layer, so the generated `longPoll()` path runs: connects + * with `connectionId`, holds polls until aborted. + */ + const fakeLongPoll = (connectionId = 'new_id') => { + const calls = []; + const respond = (data) => Promise.resolve({ data, status: 200, headers: {} }); + vi.spyOn(client.axiosInstance, 'request').mockImplementation((config) => { + calls.push({ + url: config.url, + params: config.params, + timeout: config.timeout, + authorization: config.headers?.Authorization, + }); + if (config.params?.json) { + return respond({ + event: { type: 'connection.ok', connection_id: connectionId }, + }); + } + if (config.params?.close) return respond({}); + return new Promise((_, reject) => + config.signal?.addEventListener('abort', () => reject(new CanceledError())), + ); + }); + return calls; + }; + + beforeEach(() => { + client = new StreamChat('key', { allowServerSideConnect: true }); + client.config.set({ client: { wsConnection: { enableWSFallback: true } } }); + // As the real socket does: arms the connection id before failing. + socketConnect = vi + .spyOn(StableWSConnection.prototype, 'connect') + .mockImplementation(async () => { + client.connectionIdManager.arm(); + throw wsFailure(); + }); + }); + + afterEach(async () => { + await client.disconnectUser(); + vi.unstubAllGlobals(); + }); + + it('should try wsFallback if WebSocket fails', async () => { + const calls = fakeLongPoll(); + + const health = await client.connectUser({ id: 'amin' }, userToken); + + expect(calls[0]).toMatchObject({ + url: expect.stringMatching(/\/api\/v2\/longpoll$/), + // authenticated by the Authorization header, so the message carries only a placeholder + params: { + json: expect.objectContaining({ token: 'anonymous', products: ['chat'] }), + }, + timeout: 8000, + }); + expect(calls[1]).toMatchObject({ + params: { connection_id: 'new_id' }, + timeout: 30000, + }); + + expect(health).toMatchObject({ type: 'connection.ok', connection_id: 'new_id' }); + // the socket's own connect timeout, not a shortened one + expect(socketConnect).toHaveBeenCalledWith(undefined); + expect(client.wsConnection.fallback.state).toBe(WSFallbackConnectionState.Connected); + expect(client.wsConnection.isHealthy).toBe(true); + expect(client.connectionIdManager.connectionId).toBe('new_id'); + + await client.disconnectUser(); + expect(client.wsConnection.fallback.state).toBe( + WSFallbackConnectionState.Disconnected, + ); + expect(client.wsConnection.isHealthy).toBe(false); + // the close still carries the id, although the socket's disconnect dropped it first + expect(calls.at(-1)).toMatchObject({ + params: { close: true, connection_id: 'new_id' }, + }); + }); + + it('should fire connection.fallback_activated and connection.ok events', async () => { + fakeLongPoll(); + const dispatchEvent = vi.spyOn(client, 'dispatchEvent'); + + await client.connectUser({ id: 'amin' }, userToken); + + expect(dispatchEvent).toHaveBeenCalledWith( + expect.objectContaining({ + type: 'connection.fallback_activated', + mode: 'longpoll', + }), + ); + expect(dispatchEvent).toHaveBeenCalledWith( + expect.objectContaining({ type: 'connection.ok', connection_id: 'new_id' }), + ); + }); + + it('should make a watch issued during the WebSocket attempt wait for the long-poll id', async () => { + fakeLongPoll(); + + const connecting = client.connectUser({ id: 'amin' }, userToken); + const waiter = client.connectionIdManager.getConnectionId(); + await connecting; + + await expect(waiter).resolves.toBe('new_id'); + }); + + it('should close the long-poll connection id on closeConnection', async () => { + const calls = fakeLongPoll(); + await client.connectUser({ id: 'amin' }, userToken); + + await client.closeConnection(); + + expect(calls.at(-1)).toMatchObject({ + url: expect.stringMatching(/\/api\/v2\/longpoll$/), + params: { close: true, connection_id: 'new_id' }, + }); + expect(client.wsConnection.isHealthy).toBe(false); + }); + + it('should keep using the fallback after closeConnection -> openConnection', async () => { + fakeLongPoll(); + await client.connectUser({ id: 'amin' }, userToken); + const fallback = client.wsConnection.fallback; + + await client.closeConnection(); + await client.openConnection(); + + expect(client.wsConnection.fallback).toBe(fallback); + expect(socketConnect).toHaveBeenCalledTimes(1); + expect(client.wsConnection.fallback.state).toBe(WSFallbackConnectionState.Connected); + }); + + it('should connect a token provider on the kept long-poll after disconnectUser', async () => { + const calls = fakeLongPoll(); + await client.connectUser({ id: 'amin' }, userToken); + await client.disconnectUser(); + + const health = await client.connectUser({ id: 'amin' }, async () => userToken); + + expect(health).toMatchObject({ type: 'connection.ok', connection_id: 'new_id' }); + expect(socketConnect).toHaveBeenCalledTimes(1); + expect(client.wsConnection.fallback.state).toBe(WSFallbackConnectionState.Connected); + expect(calls.filter((call) => call.params?.json).at(-1)).toMatchObject({ + authorization: userToken, + }); + }); + + it('should route network status to the long-poll', async () => { + fakeLongPoll(); + await client.connectUser({ id: 'amin' }, userToken); + const socketApply = vi.spyOn(client.wsConnection.connection, '_applyNetworkStatus'); + client.networkConnection.setStatusReporter(null); + + client.networkConnection.setStatus(false); + expect(client.wsConnection.fallback.state).toBe(WSFallbackConnectionState.Closed); + expect(client.wsConnection.isHealthy).toBe(false); + + client.networkConnection.setStatus(true); + await vi.waitFor(() => + expect(client.wsConnection.fallback.state).toBe( + WSFallbackConnectionState.Connected, + ), + ); + expect(client.wsConnection.isHealthy).toBe(true); + expect(socketApply).not.toHaveBeenCalled(); + }); + + it('should report the long-poll connect as in flight', async () => { + fakeLongPoll(); + await client.connectUser({ id: 'amin' }, userToken); + await client.closeConnection(); + + const reopening = client.openConnection(); + expect(client.wsConnection.isConnecting).toBe(true); + // a second call hands back the attempt in flight instead of starting another + expect(client.openConnection()).toBe(reopening); + + await reopening; + expect(client.wsConnection.isConnecting).toBe(false); + }); + + it('should ignore fallback if flag is false', async () => { + fakeLongPoll(); + client.wsConnection.updateConfig({ enableWSFallback: false }); + + await expect(client.connectUser({ id: 'amin' }, userToken)).rejects.toThrow( + /initial WS connection could not be established/, + ); + + expect(socketConnect).toHaveBeenCalledTimes(1); + expect(client.wsConnection.fallback).toBeUndefined(); + expect(client.axiosInstance.request).not.toHaveBeenCalled(); + }); + + it('should ignore fallback if a network reporter says the device is offline', async () => { + fakeLongPoll(); + client.networkConnection.setStatusReporter(null); + client.networkConnection.setStatus(false); + + await expect(client.connectUser({ id: 'amin' }, userToken)).rejects.toThrow( + /initial WS connection could not be established/, + ); + + expect(client.wsConnection.fallback).toBeUndefined(); + }); + + it('should fall back anyway when the offline reading only mirrors the WebSocket', async () => { + fakeLongPoll(); + // the host default without a network API: its "offline" only means the socket is down + expect(client.networkConnection.usesDefaultWSConnectionNetworkStatusReporter).toBe( + true, + ); + client.networkConnection.state.partialNext({ isOnline: false }); + + await client.connectUser({ id: 'amin' }, userToken); + + expect(client.wsConnection.fallback.state).toBe(WSFallbackConnectionState.Connected); + }); +}); diff --git a/test/unit/codegen/connectionIdEndpoints.test.ts b/test/unit/codegen/connectionIdEndpoints.test.ts index 7a5a911c9..5247ff1ba 100644 --- a/test/unit/codegen/connectionIdEndpoints.test.ts +++ b/test/unit/codegen/connectionIdEndpoints.test.ts @@ -3,17 +3,14 @@ import { join } from 'node:path'; import { describe, expect, it } from 'vitest'; /** - * `requiresConnectionId` (src/api-client.ts) gates on the `watch` / `presence` flags, and falls - * back to the generated `queryParams` carrying a `connection_id` key for the operations that - * declare no flag at all: `stopWatchingChannel` and `longPoll`. The generator emits that key for - * every operation the client-side OpenAPI spec declares one on, even when the value is `undefined`. + * `requiresConnectionId` (src/api-client.ts) gates only on the `watch` / `presence` flags. The two + * generated operations that declare `connection_id` but carry no flag set it themselves: + * `stopWatchingChannel` from `StreamChat`'s override, which always waits for the connection id as the + * gate does and sends it, and `longPoll`'s endpoint from the long-poll fallback, which addresses its + * own polls. * - * That coupling is invisible at runtime: should the generator start omitting keys holding - * `undefined`, the fallback would silently stop firing for those two, and they would start racing - * the handshake again with nothing failing. This test pins the set instead. - * - * The other seven entries are not load-bearing for the fallback - they all declare a flag, which is - * read first - but they are pinned so that a regenerated spec adding a connection-scoped operation + * That makes every new flagless, connection-scoped operation a place the id has to be set by hand, + * with nothing failing if it is not. This test pins the set, so a regenerated spec adding one * surfaces here rather than passing silently. */ const GEN_ROOT = join(__dirname, '../../../src/gen'); diff --git a/test/unit/networkConnection/socketWiring.test.ts b/test/unit/networkConnection/socketWiring.test.ts index 563fd1556..e06994aa3 100644 --- a/test/unit/networkConnection/socketWiring.test.ts +++ b/test/unit/networkConnection/socketWiring.test.ts @@ -31,6 +31,29 @@ const clientWithSocket = () => { }; describe('socket ↔ network wiring', () => { + describe('usesDefaultWSConnectionNetworkStatusReporter', () => { + it('is true under the host default that mirrors the socket', () => { + expect( + new StreamChat('api-key').networkConnection + .usesDefaultWSConnectionNetworkStatusReporter, + ).toBe(true); + }); + + it('is false once another reporter is installed, or the reporter is cleared', () => { + const client = new StreamChat('api-key'); + + client.networkConnection.setStatusReporter(fakeReporter().reporter); + expect(client.networkConnection.usesDefaultWSConnectionNetworkStatusReporter).toBe( + false, + ); + + client.networkConnection.setStatusReporter(null); + expect(client.networkConnection.usesDefaultWSConnectionNetworkStatusReporter).toBe( + false, + ); + }); + }); + describe('no window listeners', () => { it('the socket registers none — the browser reporter is the only place that touches window', () => { const addEventListener = vi.fn(); diff --git a/test/unit/wsConnection/WSConnection.config.test.ts b/test/unit/wsConnection/WSConnection.config.test.ts index edf520cd1..0b272537e 100644 --- a/test/unit/wsConnection/WSConnection.config.test.ts +++ b/test/unit/wsConnection/WSConnection.config.test.ts @@ -35,6 +35,7 @@ describe('client.wsConnection configuration', () => { // `undefined` property as absent, so omitting them here would silently stop covering them. expect(client.wsConnection.config).toEqual({ connectTimeoutMs: 15_000, + enableWSFallback: false, pingIntervalMs: 25_000, healthCheckGracePeriodMs: 10_000, offlineNotificationDisplayDelayMs: 5_000, @@ -42,7 +43,7 @@ describe('client.wsConnection configuration', () => { urlParams: undefined, connection: undefined, }); - expect(Object.keys(client.wsConnection.config)).toHaveLength(7); + expect(Object.keys(client.wsConnection.config)).toHaveLength(8); }); it('keeps the 25s ping / 35s connection check pair the socket documents', () => { @@ -104,11 +105,13 @@ describe('client.wsConnection configuration', () => { // already reachable — the change is that they are declared, typed and validated in one place. // // `offlineNotificationDisplayDelayMs` is the deliberate exception: nothing here acts on it, - // and it is declared so the UI SDKs share one value. Still excluded: the network-recovery + // and it is declared so the UI SDKs share one value. `enableWSFallback` is read by + // `client.connect()`, which decides whether to fall back. Still excluded: the network-recovery // retry, a bare literal no caller could reach. expect(Object.keys(client.wsConnection.config).sort()).toEqual([ 'connectTimeoutMs', 'connection', + 'enableWSFallback', 'healthCheckGracePeriodMs', 'offlineNotificationDisplayDelayMs', 'pingIntervalMs', diff --git a/test/unit/wsConnection/WSConnection.test.ts b/test/unit/wsConnection/WSConnection.test.ts index 6358304d3..1b85dc672 100644 --- a/test/unit/wsConnection/WSConnection.test.ts +++ b/test/unit/wsConnection/WSConnection.test.ts @@ -1,6 +1,10 @@ import { beforeEach, describe, expect, it, vi } from 'vitest'; import { StreamChat } from '../../../src'; import { StableWSConnection } from '../../../src/connection'; +import { + WSFallbackConnectionState, + WSConnectionFallback, +} from '../../../src/connection/wsConnection/WSConnectionFallback'; describe('client.wsConnection', () => { let client: StreamChat; @@ -101,6 +105,21 @@ describe('client.wsConnection', () => { expect(staleApply).not.toHaveBeenCalled(); }); + it('routes network status to the long-poll once switched to it', () => { + const socket = new StableWSConnection({ wsConnection: client.wsConnection }); + client.wsConnection.connection = socket; + const fallback = new WSConnectionFallback({ client }); + client.wsConnection.fallback = fallback; + + const socketApply = vi.spyOn(socket, '_applyNetworkStatus'); + const fallbackApply = vi.spyOn(fallback, '_applyNetworkStatus'); + + client.networkConnection.setStatus(false); + + expect(fallbackApply).toHaveBeenCalledWith(false); + expect(socketApply).not.toHaveBeenCalled(); + }); + it('registers exactly one listener however many sockets come and go', () => { const before = client.listeners.get('connection.changed')?.size ?? 0; @@ -314,6 +333,51 @@ describe('client.wsConnection', () => { expect(disconnect).toHaveBeenCalledWith(0); }); + it('forwards disconnect to the long-poll too', async () => { + const connection = client.wsConnection.connection; + if (!connection) throw new Error('socket missing'); + const disconnect = vi.spyOn(connection, 'disconnect').mockResolvedValue(undefined); + const fallback = new WSConnectionFallback({ client }); + client.wsConnection.fallback = fallback; + const fallbackDisconnect = vi + .spyOn(fallback, 'disconnect') + .mockResolvedValue(undefined); + + await client.wsConnection.disconnect(0); + + expect(disconnect).toHaveBeenCalledWith(0); + expect(fallbackDisconnect).toHaveBeenCalledWith(0, undefined); + }); + + it("closes the long-poll with the connection id the socket's disconnect drops", async () => { + const fallback = new WSConnectionFallback({ client }); + client.wsConnection.fallback = fallback; + const longPoll = vi.spyOn(client, 'longPoll').mockResolvedValue({}); + client.connectionIdManager.resolveConnectionId('long-poll-id'); + + await client.wsConnection.disconnect(0); + + // the socket's disconnect ran first and dropped it from the manager + expect(client.connectionIdManager.connectionId).toBeUndefined(); + expect(longPoll).toHaveBeenCalledWith( + { close: true, connection_id: 'long-poll-id' }, + expect.objectContaining({ timeout: 0 }), + ); + }); + + it('forwards isConnecting from the long-poll once switched to it', () => { + const connection = client.wsConnection.connection; + if (!connection) throw new Error('socket missing'); + connection.isConnecting = true; + const fallback = new WSConnectionFallback({ client }); + client.wsConnection.fallback = fallback; + + expect(client.wsConnection.isConnecting).toBe(false); + + fallback.state = WSFallbackConnectionState.Connecting; + expect(client.wsConnection.isConnecting).toBe(true); + }); + it('forwards onlineStatusChanged, which React Native still calls directly', () => { // Deprecated, and must keep working for one major: `stream-chat-react-native` calls // `client.wsConnection.onlineStatusChanged({ type: … })` with a synthesized DOM event. diff --git a/test/unit/wsConnection/WSConnectionFallback.test.js b/test/unit/wsConnection/WSConnectionFallback.test.js new file mode 100644 index 000000000..7a69d863e --- /dev/null +++ b/test/unit/wsConnection/WSConnectionFallback.test.js @@ -0,0 +1,654 @@ +import sinon from 'sinon'; +import axios, { CanceledError } from 'axios'; + +import * as utils from '../../../src/utils'; +import * as errors from '../../../src/errors'; +import { ConnectionIdManager } from '../../../src/connection/ConnectionIdManager'; +import { + WSFallbackConnectionState, + WSConnectionFallback, +} from '../../../src/connection/wsConnection/WSConnectionFallback'; + +import { describe, it, expect, afterEach, vi, beforeAll, beforeEach } from 'vitest'; + +describe('WSConnectionFallback', () => { + const newClient = (overrides) => ({ + longPoll: sinon.spy(), + _buildWSAuthPayload: sinon.stub().returns('payload'), + dispatchEvent: sinon.spy(), + _settleConnectPromises: sinon.spy(), + _markActiveChannelsWatchInterrupted: sinon.spy(), + connectionIdManager: new ConnectionIdManager(), + wsConnection: { + isHealthy: false, + // Like the real one: reports whether the status changed. + _setStatus: sinon.spy(function ({ isHealthy }) { + if (this.isHealthy === isHealthy) return false; + this.isHealthy = isHealthy; + return true; + }), + }, + ...overrides, + }); + + afterEach(() => { + vi.restoreAllMocks(); + }); + + afterEach(() => { + sinon.restore(); + }); + + describe('constructor', () => { + it('should set the state correctly', () => { + const client = newClient(); + const c = new WSConnectionFallback({ client }); + + expect(c.client).to.be.eql(client); + expect(c.state).to.be.eql(WSFallbackConnectionState.Init); + expect(c.consecutiveFailures).to.be.eql(0); + }); + }); + + describe('_setState', () => { + it('should update state correctly', function () { + const c = new WSConnectionFallback({ client: newClient() }); + + expect(c.state).to.be.eql(WSFallbackConnectionState.Init); + + c._setState(WSFallbackConnectionState.Closed); + expect(c.state).to.be.eql(WSFallbackConnectionState.Closed); + + c._setState(WSFallbackConnectionState.Connected); + expect(c.state).to.be.eql(WSFallbackConnectionState.Connected); + + c._setState(WSFallbackConnectionState.Connecting); + expect(c.state).to.be.eql(WSFallbackConnectionState.Connecting); + + c._setState(WSFallbackConnectionState.Disconnected); + expect(c.state).to.be.eql(WSFallbackConnectionState.Disconnected); + }); + + it('should report online status to wsConnection', function () { + const client = newClient(); + const c = new WSConnectionFallback({ client }); + + c._setState(WSFallbackConnectionState.Connecting); + expect(client.wsConnection._setStatus.called).to.be.false; + + c._setState(WSFallbackConnectionState.Connected); + expect(client.wsConnection._setStatus.calledOnceWithExactly({ isHealthy: true })).to + .be.true; + }); + + it('should report offline status, drop the connection id and mark watches', function () { + const client = newClient(); + const c = new WSConnectionFallback({ client }); + c._setState(WSFallbackConnectionState.Connecting); + client.connectionIdManager.resolveConnectionId('id'); + c._setState(WSFallbackConnectionState.Connected); + + c._setState(WSFallbackConnectionState.Closed); + expect(client.wsConnection._setStatus.lastCall.args).to.be.eql([ + { isHealthy: false }, + ]); + expect(client.connectionIdManager.connectionId).to.be.undefined; + expect(client._markActiveChannelsWatchInterrupted.calledOnce).to.be.true; + + // no Connecting => Connected transition, so no online status + c._setState(WSFallbackConnectionState.Connected); + expect(client.wsConnection.isHealthy).to.be.false; + + // already offline: nothing new to mark + c._setState(WSFallbackConnectionState.Disconnected); + expect(client._markActiveChannelsWatchInterrupted.calledOnce).to.be.true; + }); + + it('should already be in the new state when it reports the status', function () { + const client = newClient(); + const c = new WSConnectionFallback({ client }); + const seen = []; + client.wsConnection._setStatus = sinon.spy(() => seen.push(c.state)); + + c._setState(WSFallbackConnectionState.Connecting); + c._setState(WSFallbackConnectionState.Connected); + c._setState(WSFallbackConnectionState.Closed); + c._setState(WSFallbackConnectionState.Disconnected); + expect(seen).to.be.eql([ + WSFallbackConnectionState.Connected, + WSFallbackConnectionState.Closed, + WSFallbackConnectionState.Disconnected, + ]); + }); + + it('should not be closed again by the socket-mirroring reporter on disconnect()', async function () { + const client = newClient(); + const c = new WSConnectionFallback({ client }); + c._req = sinon.stub().resolves(); + c.state = WSFallbackConnectionState.Connected; + client.wsConnection.isHealthy = true; + // the reporter forwards the socket's status as the network's, which is routed back here + const setStatus = client.wsConnection._setStatus; + client.wsConnection._setStatus = sinon.spy(function (status) { + const changed = setStatus.call(this, status); + if (changed) c._applyNetworkStatus(status.isHealthy); + return changed; + }); + + await c.disconnect(); + expect(c.state).to.be.eql(WSFallbackConnectionState.Disconnected); + // a nested `_setState(Closed)` would have reported the status a second time + expect(client.wsConnection._setStatus.calledOnce).to.be.true; + }); + + it('should keep a transition a status subscriber makes', async function () { + const client = newClient(); + const c = new WSConnectionFallback({ client }); + c._req = sinon.spy(async (params) => { + if (params.json) return { event: { connection_id: 'id' } }; + if (params.close) return; + // a poll: stop the loop, so a regression fails instead of hanging + c.state = WSFallbackConnectionState.Closed; + return {}; + }); + // an integrator closing the connection as soon as it comes up + const setStatus = client.wsConnection._setStatus; + client.wsConnection._setStatus = sinon.spy(function (status) { + const changed = setStatus.call(this, status); + if (changed && status.isHealthy) c.disconnect(); + return changed; + }); + + await c.connect(); + expect(c.state).to.be.eql(WSFallbackConnectionState.Disconnected); + // only the connect and the close: no poll on a disconnected connection + expect(c._req.callCount).to.be.eql(2); + expect(c._req.secondCall.args[0].close).to.be.true; + }); + }); + + describe('_applyNetworkStatus', () => { + it('should call connect for online status in Closed state', () => { + const c = new WSConnectionFallback({ client: newClient() }); + c.connect = sinon.spy(); + c._applyNetworkStatus(true); + expect(c.connect.called).to.be.false; + + c.state = WSFallbackConnectionState.Closed; + c._applyNetworkStatus(true); + expect(c.connect.calledOnceWithExactly(true)).to.be.true; + }); + + it('should go to Close state on offline status', () => { + const c = new WSConnectionFallback({ client: newClient() }); + const spy = sinon.spy(); + c.abortController = { abort: spy }; + c._applyNetworkStatus(false); + expect(c.state).to.be.eql(WSFallbackConnectionState.Closed); + expect(spy.calledOnce).to.be.true; + expect(c.abortController).to.be.undefined; + + c._applyNetworkStatus(false); + expect(c.state).to.be.eql(WSFallbackConnectionState.Closed); + expect(c.abortController).to.be.undefined; + }); + }); + + describe('isConnecting', () => { + it('is true only in Connecting state', () => { + const c = new WSConnectionFallback({ client: newClient() }); + expect(c.isConnecting).to.be.false; + + for (const state of Object.values(WSFallbackConnectionState)) { + c.state = state; + expect(c.isConnecting).to.be.eql(state === WSFallbackConnectionState.Connecting); + } + }); + + it('is true while the connect request is pending', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + let respond; + c._req = () => new Promise((resolve) => (respond = resolve)); + c._poll = sinon.spy(); + + const connecting = c.connect(); + expect(c.isConnecting).to.be.true; + + respond({ event: { connection_id: 'id' } }); + await connecting; + expect(c.isConnecting).to.be.false; + }); + }); + + describe('disconnect', () => { + it('should ignore network status until it connects again', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + c._req = () => null; + await c.disconnect(); + c.connect = sinon.spy(); + + // going offline would otherwise move it to Closed, and back online reconnect it + c._applyNetworkStatus(false); + c._applyNetworkStatus(true); + expect(c.state).to.be.eql(WSFallbackConnectionState.Disconnected); + expect(c.connect.called).to.be.false; + + c.state = WSFallbackConnectionState.Connected; + c._applyNetworkStatus(false); + expect(c.state).to.be.eql(WSFallbackConnectionState.Closed); + }); + + it('should cancel requests and set the state correctly', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + c._req = sinon.spy(); + const connection_id = 'id'; + const abort = sinon.spy(); + c.abortController = { abort }; + const timeout = 500; + await c.disconnect(timeout, connection_id); + + expect(c.state).to.be.eql(WSFallbackConnectionState.Disconnected); + expect(c.abortController).to.be.undefined; + expect(abort.calledOnce).to.be.true; + expect( + c._req.calledOnceWithExactly({ close: true, connection_id }, { timeout }, false), + ).to.be.true; + }); + + it('should ingore request errors', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + c._req = () => Promise.reject('error'); + await c.disconnect(); + }); + }); + + describe('_req', () => { + it('should set abort controller', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + expect(c.abortController).to.be.undefined; + await c._req({}, {}); + expect(c.abortController).to.be.instanceOf(AbortController); + + c.abortController = undefined; + await c._req({ close: true }, {}); + expect(c.abortController).to.be.undefined; + }); + + it('should send the request correctly', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + + const params = { json: 'hi' }; + const config = { timeout: 100 }; + await c._req(params, config); + expect( + c.client.longPoll.calledOnceWithExactly(params, { + ...config, + signal: c.abortController.signal, + }), + ).to.be.true; + }); + + it('should abort the request in flight when going offline', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + await c._req({}, {}); + const [, { signal }] = c.client.longPoll.lastCall.args; + + c._applyNetworkStatus(false); + expect(signal.aborted).to.be.true; + }); + + it('should keep track of consecutive failures', async () => { + // ok-err-err-ok-ok... + const longPoll = vi + .fn() + .mockResolvedValueOnce() + .mockRejectedValueOnce() + .mockRejectedValueOnce() + .mockResolvedValue(); + const c = new WSConnectionFallback({ + client: newClient({ longPoll }), + }); + + expect(c.consecutiveFailures).toBe(0); + await c._req({}); + expect(c.consecutiveFailures).toBe(0); + await expect(c._req({})).rejects.toThrow(); + expect(c.consecutiveFailures).toBe(1); + await expect(c._req({})).rejects.toThrow(); + expect(c.consecutiveFailures).toBe(2); + await c._req({}); + expect(c.consecutiveFailures).toBe(0); + await c._req({}); + expect(c.consecutiveFailures).toBe(0); + }); + + it('should not retry for non-retryable errors', async () => { + const longPoll = sinon.stub().rejects(); + sinon.stub(errors, 'isErrorRetryable').returns(false); + const c = new WSConnectionFallback({ + client: newClient({ longPoll }), + }); + sinon.spy(c); + + expect(c.consecutiveFailures).to.be.eql(0); + await expect(c._req({}, {}, true)).rejects.toThrow(); + expect(c.consecutiveFailures).to.be.eql(1); + expect(c._req.calledOnce).to.be.true; + }); + + it('should not retry when retry flag is false', async () => { + const longPoll = sinon.stub().rejects(); + sinon.stub(errors, 'isErrorRetryable').returns(true); + const c = new WSConnectionFallback({ + client: newClient({ longPoll }), + }); + sinon.spy(c); + + expect(c.consecutiveFailures).to.be.eql(0); + await expect(c._req({}, {}, false)).rejects.toThrow(); + expect(c.consecutiveFailures).to.be.eql(1); + expect(c._req.calledOnce).to.be.true; + }); + + it('should retry errors if it is retryable', async () => { + const longPoll = sinon.stub().rejects(); + + vi.spyOn(errors, 'isErrorRetryable') + .mockReturnValueOnce(true) + .mockReturnValueOnce(true) + .mockReturnValueOnce(false); + + vi.spyOn(utils, 'sleep').mockResolvedValue(); + + const c = new WSConnectionFallback({ + client: newClient({ longPoll }), + }); + sinon.spy(c, '_req'); + + expect(c.consecutiveFailures).to.be.eql(0); + await expect(c._req({}, {}, true)).rejects.toThrow(); + expect(c.consecutiveFailures).to.be.eql(3); + expect(c._req.calledThrice).to.be.true; + }); + + it('should drop a retry whose connection id changed while it waited', async () => { + const longPoll = sinon.stub().rejects(); + // once, so a retry that is not dropped fails instead of looping + vi.spyOn(errors, 'isErrorRetryable').mockReturnValueOnce(true); + const c = new WSConnectionFallback({ client: newClient({ longPoll }) }); + c.client.connectionIdManager.resolveConnectionId('old'); + // a reconnect lands while the retry sleeps + vi.spyOn(utils, 'sleep').mockImplementation(async () => { + c.client.connectionIdManager.invalidate(); + c.client.connectionIdManager.resolveConnectionId('new'); + }); + + const error = await c._req({ connection_id: 'old' }, {}, true).catch((e) => e); + expect(axios.isCancel(error)).to.be.true; + expect(longPoll.calledOnce).to.be.true; + }); + + it('should drop a retry whose connection id disconnect() dropped while it waited', async () => { + const longPoll = sinon.stub().rejects(); + vi.spyOn(errors, 'isErrorRetryable').mockReturnValueOnce(true); + const c = new WSConnectionFallback({ client: newClient({ longPoll }) }); + c.client.connectionIdManager.resolveConnectionId('old'); + vi.spyOn(utils, 'sleep').mockImplementation(async () => { + c.client.connectionIdManager.invalidate(); + }); + + const error = await c._req({ connection_id: 'old' }, {}, true).catch((e) => e); + expect(axios.isCancel(error)).to.be.true; + expect(longPoll.calledOnce).to.be.true; + }); + + it('should drop a retry whose connection went offline while it waited', async () => { + const longPoll = sinon.stub().rejects(); + vi.spyOn(errors, 'isErrorRetryable').mockReturnValueOnce(true); + const c = new WSConnectionFallback({ client: newClient({ longPoll }) }); + c.client.connectionIdManager.resolveConnectionId('old'); + c.state = WSFallbackConnectionState.Connected; + vi.spyOn(utils, 'sleep').mockImplementation(async () => { + c._applyNetworkStatus(false); + }); + + const error = await c._req({ connection_id: 'old' }, {}, true).catch((e) => e); + expect(c.state).to.be.eql(WSFallbackConnectionState.Closed); + expect(axios.isCancel(error)).to.be.true; + expect(longPoll.calledOnce).to.be.true; + }); + + it('should retry with the same connection id while it is still current', async () => { + const longPoll = sinon.stub(); + longPoll.onFirstCall().rejects(); + longPoll.resolves({ events: [] }); + vi.spyOn(errors, 'isErrorRetryable').mockReturnValue(true); + vi.spyOn(utils, 'sleep').mockResolvedValue(); + const c = new WSConnectionFallback({ client: newClient({ longPoll }) }); + c.client.connectionIdManager.resolveConnectionId('id'); + + await c._req({ connection_id: 'id' }, {}, true); + expect(longPoll.calledTwice).to.be.true; + expect(longPoll.secondCall.args[0]).to.be.eql({ connection_id: 'id' }); + }); + }); + + describe('connect', () => { + const health = { connection_id: 'connectionID' }; + it('should skip connect if already connecting or connected', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + sinon.spy(c); + c.state = WSFallbackConnectionState.Connecting; + expect(await c.connect()).to.be.undefined; + c.state = WSFallbackConnectionState.Connected; + expect(await c.connect()).to.be.undefined; + expect(c._setState.called).to.be.false; + expect(c._req.called).to.be.false; + }); + + it('should send request in correct format', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + c._req = sinon.stub().resolves({ event: health }); + c._poll = sinon.spy(); + + expect(await c.connect()).to.be.eql(health); + // authenticated by the Authorization header, so the message carries no token + expect(c.client._buildWSAuthPayload.calledOnceWithExactly()).to.be.true; + expect(c._poll.calledOnce).to.be.true; + expect(c._req.calledOnceWithExactly({ json: 'payload' }, { timeout: 8000 }, false)) + .to.be.true; + + c.state = WSFallbackConnectionState.Init; + c._req = sinon.stub().resolves({ event: health }); + expect(await c.connect(true)).to.be.eql(health); + expect(c._req.calledOnceWithExactly({ json: 'payload' }, { timeout: 8000 }, true)) + .to.be.true; + }); + + it('should update state and the connection id', async () => { + let c = new WSConnectionFallback({ client: newClient() }); + c._req = sinon.stub().resolves({ event: health }); + c._poll = sinon.spy(); + expect(await c.connect()).to.be.eql(health); + expect(c.state).to.be.eql(WSFallbackConnectionState.Connected); + expect(c.client.connectionIdManager.connectionId).to.be.eql(health.connection_id); + + c = new WSConnectionFallback({ client: newClient() }); + c._req = sinon.stub().rejects(); + c._poll = sinon.spy(); + await expect(c.connect()).rejects.toThrow(); + expect(c._poll.called).to.be.false; + expect(c.state).to.be.eql(WSFallbackConnectionState.Closed); + expect(c.client.connectionIdManager.connectionId).to.be.undefined; + }); + + it('should only start polling after connect', async () => { + let c = new WSConnectionFallback({ client: newClient() }); + c._req = sinon.stub().resolves({ event: health }); + c._poll = sinon.spy(); + expect(await c.connect()).to.be.eql(health); + expect(c._poll.calledOnce).to.be.true; + + c = new WSConnectionFallback({ client: newClient() }); + c._req = sinon.stub().rejects(); + c._poll = sinon.spy(); + await expect(c.connect()).rejects.toThrow(); + expect(c._poll.called).to.be.false; + }); + + it('should settle the connect promises on reconnect only', async () => { + let c = new WSConnectionFallback({ client: newClient() }); + c._req = sinon.stub().resolves({ event: health }); + c._poll = sinon.spy(); + await c.connect(); + expect(c.client._settleConnectPromises.called).to.be.false; + + c = new WSConnectionFallback({ client: newClient() }); + c._req = sinon.stub().rejects(); + c._poll = sinon.spy(); + await expect(c.connect()).rejects.toThrow(); + expect(c.client._settleConnectPromises.called).to.be.false; + + c = new WSConnectionFallback({ client: newClient() }); + c._req = sinon.stub().resolves({ event: health }); + c._poll = sinon.spy(); + await c.connect(true); + expect(c.client._settleConnectPromises.called).to.be.true; + }); + + it('should make a watch issued meanwhile wait for its connection id', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + let respond; + c._req = () => new Promise((resolve) => (respond = resolve)); + c._poll = sinon.spy(); + + const connecting = c.connect(); + const waiter = c.client.connectionIdManager.getConnectionId(); + respond({ event: health }); + await connecting; + await expect(waiter).resolves.toBe(health.connection_id); + }); + + it('should fail whatever waits for an id when connect fails, but not when cancelled', async () => { + let c = new WSConnectionFallback({ client: newClient() }); + c._req = sinon.stub().rejects(new Error('boom')); + let connecting = c.connect(); + let waiter = c.client.connectionIdManager.getConnectionId(); + await expect(connecting).rejects.toThrow('boom'); + await expect(waiter).rejects.toThrow('boom'); + + c = new WSConnectionFallback({ client: newClient() }); + c._req = sinon.stub().rejects(new CanceledError()); + connecting = c.connect(); + waiter = c.client.connectionIdManager.getConnectionId(); + await expect(connecting).rejects.toThrow(); + expect(c.client.connectionIdManager.loadConnectionIdPromise).to.be.equal(waiter); + }); + + it('should stay Disconnected when disconnect() cancels it', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + let cancel; + c._req = (params) => + params.close + ? Promise.resolve() + : new Promise((_, reject) => (cancel = () => reject(new CanceledError()))); + + const connecting = c.connect(); + await c.disconnect(); + cancel(); + await expect(connecting).rejects.toThrow(); + expect(c.state).to.be.eql(WSFallbackConnectionState.Disconnected); + + // so the next online edge does not reconnect it + c.connect = sinon.spy(); + c._applyNetworkStatus(false); + c._applyNetworkStatus(true); + expect(c.connect.called).to.be.false; + }); + }); + + describe('_poll', () => { + it('should do nothing if not in connect state', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + c._req = sinon.spy(); + expect(await c._poll()).to.be.undefined; + expect(c._req.called).to.be.false; + + c.state = WSFallbackConnectionState.Connecting; + expect(await c._poll()).to.be.undefined; + expect(c._req.called).to.be.false; + }); + + it('should send request in correct format', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + c.state = WSFallbackConnectionState.Connected; + c.client.connectionIdManager.resolveConnectionId('id'); + c._req = async () => { + c.state = WSFallbackConnectionState.Closed; + return {}; + }; + sinon.spy(c, '_req'); + await c._poll(); + expect( + c._req.calledOnceWithExactly({ connection_id: 'id' }, { timeout: 30000 }, true), + ).to.be.true; + }); + + it('should dispatch incoming events', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + c.state = WSFallbackConnectionState.Connected; + + c._req = async () => { + c.state = WSFallbackConnectionState.Closed; + return { events: ['1', '2'] }; + }; + await c._poll(); + expect(c.client.dispatchEvent.calledTwice).to.be.true; + expect(c.client.dispatchEvent.getCall(0).args).to.be.eql(['1']); + expect(c.client.dispatchEvent.getCall(1).args).to.be.eql(['2']); + }); + + it('should reconnect if ConnectionID err', async () => { + const c = new WSConnectionFallback({ client: newClient() }); + c.state = WSFallbackConnectionState.Connected; + c.connect = sinon.spy(); + c._req = async () => { + const err = new Error(); + err.code = 46; + throw err; + }; + + await c._poll(); + expect(c.state).to.be.eql(WSFallbackConnectionState.Disconnected); + expect(c.connect.calledOnceWithExactly(true)).to.be.true; + }); + + it('should stop for non-retryable errors', async () => { + vi.spyOn(errors, 'isErrorRetryable').mockReturnValue(false); + vi.spyOn(errors, 'isAPIError').mockReturnValue(true); + const c = new WSConnectionFallback({ client: newClient() }); + c.state = WSFallbackConnectionState.Connected; + c.client.connectionIdManager.resolveConnectionId('id'); + c._req = sinon.stub().rejects(new Error('stop')); + + await c._poll(); + expect(c.state).to.be.eql(WSFallbackConnectionState.Closed); + // nothing reconnects from here, so nothing may wait for an id + expect(() => c.client.connectionIdManager.getConnectionId()).toThrow(); + }); + + it('should continue retrying for random errors', async () => { + let counter = 0; + vi.spyOn(utils, 'sleep').mockImplementation(() => { + if (++counter > 2) c.state = WSFallbackConnectionState.Disconnected; + }); + const c = new WSConnectionFallback({ client: newClient() }); + c.state = WSFallbackConnectionState.Connected; + c._req = sinon.stub().rejects(); + + await c._poll(); + expect(c._req.calledThrice).to.be.true; + expect(utils.sleep).toHaveBeenCalledTimes(3); + }); + }); +}); diff --git a/v9-to-v10-migration-guide-client-construction.md b/v9-to-v10-migration-guide-client-construction.md index b55ef8383..94b0fa6d3 100644 --- a/v9-to-v10-migration-guide-client-construction.md +++ b/v9-to-v10-migration-guide-client-construction.md @@ -155,14 +155,14 @@ If you relied on a custom serializer, file an issue — there is no supported wa Six options were dropped in v10. All are silent no-ops if left in place — TypeScript will flag them, but a plain-JS call site will not error, so remove them explicitly. -| Option | Replacement | -| ------------------ | ------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `device` | Call [`client.createDevice({ id, push_provider, push_provider_name? })`](./v9-to-v10-migration-guide-methods.md#clientsetlocaldevice) after connecting. | -| `disableCache` | None — the client always caches. See [`disableCache`](#disablecache). | -| `enableInsights` | None — the telemetry it enabled was removed. See [`enableInsights`](#enableinsights). | -| `enableWSFallback` | None — see [long-poll fallback removed](./v9-to-v10-migration-guide-other.md#long-poll-fallback-removed). | -| `warmUp` | None — see [`warmUp`](#warmup). | -| `wsConnection` | `WebSocketImpl`, which works in both browser and node. See [`wsConnection`](#wsconnection). | +| Option | Replacement | +| ------------------ | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `device` | Call [`client.createDevice({ id, push_provider, push_provider_name? })`](./v9-to-v10-migration-guide-methods.md#clientsetlocaldevice) after connecting. | +| `disableCache` | None — the client always caches. See [`disableCache`](#disablecache). | +| `enableInsights` | None — the telemetry it enabled was removed. See [`enableInsights`](#enableinsights). | +| `enableWSFallback` | `client.config.set({ client: { wsConnection: { enableWSFallback: true } } })`. See [long-poll fallback](./v9-to-v10-migration-guide-other.md#long-poll-fallback-enablewsfallback). | +| `warmUp` | None — see [`warmUp`](#warmup). | +| `wsConnection` | `WebSocketImpl`, which works in both browser and node. See [`wsConnection`](#wsconnection). | ```diff const client = new StreamChat(API_KEY, { @@ -174,6 +174,7 @@ them, but a plain-JS call site will not error, so remove them explicitly. - wsConnection: myFakeConnection, }); + await client.createDevice({ id: pushToken, push_provider: 'firebase' }); ++ client.config.set({ client: { wsConnection: { enableWSFallback: true } } }); ``` #### `device` diff --git a/v9-to-v10-migration-guide-logging.md b/v9-to-v10-migration-guide-logging.md index ee303070f..2af998ef0 100644 --- a/v9-to-v10-migration-guide-logging.md +++ b/v9-to-v10-migration-guide-logging.md @@ -9,7 +9,7 @@ - The v9 `Logger` and `LogLevel` types (from `stream-chat`) are **removed**. Import their replacements from `stream-chat`'s new logger surface (re-exported from `./logger`): `LogLevel`, `Sink`, `ConfigureLoggersOptions`, `LogLevelEnum`, `ScopedLogger`, `ChatLoggerScope`, `chatLoggerSystem`. - Log-level enum expanded from **3** values (`'info' | 'warn' | 'error'`) to **5** (`'trace' | 'debug' | 'info' | 'warn' | 'error'`). - Configure logging by calling `chatLoggerSystem.configureLoggers({...})` **before** constructing the client (there is no constructor option for it — `logLevel` / `logOptions` fields do not exist on `StreamChatOptions`). -- Internal `_log()` methods (notably on `StableWSConnection`) are gone. If you subclassed or spied on them, switch to the scoped loggers. +- Internal `_log()` methods (notably on `StableWSConnection`) are gone. If you subclassed or spied on them, switch to the scoped loggers. The one survivor is `WSConnectionFallback._log()`, kept with the v9 long-poll; it now writes to the `connection` scope with the `connection_fallback` tag. ## What ships in v10 @@ -40,7 +40,7 @@ Every internal module attaches to one of these scopes via `chatLoggerSystem.getL | `channel` | `src/channel.ts` | | `channel-manager` | `src/channel_manager.ts` | | `client` | `src/client.ts` — connection lifecycle, event dispatch | -| `connection` | `src/StableWSConnection.ts` — primary WS transport | +| `connection` | `StableWSConnection.ts` (WS) and `WSConnectionFallback.ts` (long-poll, tag `connection_fallback`) | | `message-composer` | `src/messageComposer/messageComposer.ts` | | `offline-db` | `src/offline-support/*` **and** offline-DB paths in `client.ts` / `channel.ts` / `messageComposer.ts` | | `state-store` | reserved — declared in the scope union, not yet emitted | @@ -63,6 +63,7 @@ Unknown scope names fall through to `'default'`. The `ChatLoggerScope` union nar | `type LogLevel = 'info' \| 'error' \| 'warn'` | **REPLACED** — same name, now `'trace' \| 'debug' \| 'info' \| 'warn' \| 'error'` | | `isFunction(inputOptions.logger)` guard in constructor | gone | | `StableWSConnection._log(msg, extra?, level?)` | **REMOVED** — use `chatLoggerSystem.getLogger('connection')` | +| `WSConnectionFallback._log(msg, extra?, level?)` | **KEPT** (internal) — `connection` scope, tagged `connection_fallback` | | `extraData.tags: string[]` convention (`{ tags: ['channel', 'offlineDb'], error }`) | replaced by scope + `.withExtraTags(...)` (see below) | | Structured extra as second positional arg (`(level, msg, { tags, error, event })`) | passed as rest args after the message (`.error('msg', { error })`) | diff --git a/v9-to-v10-migration-guide-methods.md b/v9-to-v10-migration-guide-methods.md index 3656d16fc..4784f51a2 100644 --- a/v9-to-v10-migration-guide-methods.md +++ b/v9-to-v10-migration-guide-methods.md @@ -22,7 +22,7 @@ Before applying any per-method entry below, apply these repo-wide renames — th The `secret` parameter, `client.secret`, `client._isUsingServerAuth()`, and all server-only methods are gone. Where a v9 method took a `user_id?` / `userID?` / `currentUserID?` override, that argument has been dropped in v10 (the connected user is always used). -**Every request-issuing method takes an optional trailing `requestOptions`.** The per-method signatures below omit it for readability, but every method generated from the OpenAPI spec ends with `requestOptions?: StreamRequestOptions` — currently `{ signal?: AbortSignal }` — and so does every hand-written wrapper around one, on both `StreamChat` and `Channel`. It is always **last**, immediately after the method's own arguments, and is never serialized into the request: +**Every request-issuing method takes an optional trailing `requestOptions`.** The per-method signatures below omit it for readability, but every method generated from the OpenAPI spec ends with `requestOptions?: StreamRequestOptions` — currently `{ signal?: AbortSignal; onUploadProgress?; timeout?: number }`, where `timeout` overrides `axiosRequestConfig.timeout` for that request — and so does every hand-written wrapper around one, on both `StreamChat` and `Channel`. It is always **last**, immediately after the method's own arguments, and is never serialized into the request: ```ts const controller = new AbortController(); @@ -459,7 +459,7 @@ A bare URI `string` is still accepted and normalized into `{ uri }`, with the na **No node input.** v10 dropped the `form-data` dependency for the platform's global `FormData`, and with it every node-only input: `Buffer` and readable streams are no longer accepted. Backend uploads move to `@stream-io/node-sdk` — see [`v9-to-v10-migration-guide-server-side.md`](./v9-to-v10-migration-guide-server-side.md#uploads-from-node-are-gone). -**`axiosRequestConfig` → `requestOptions`.** The last argument is now `StreamRequestOptions`, carrying `onUploadProgress` and `signal`. `timeout: 0` and unbounded `maxContentLength` / `maxBodyLength` are applied automatically for multipart and are no longer caller-overridable. +**`axiosRequestConfig` → `requestOptions`.** The last argument is now `StreamRequestOptions`, carrying `onUploadProgress` and `signal`. `timeout: 0` and unbounded `maxContentLength` / `maxBodyLength` are applied automatically for multipart; of the three, only the timeout is caller-overridable, through `requestOptions.timeout`. ```ts // v9 diff --git a/v9-to-v10-migration-guide-other.md b/v9-to-v10-migration-guide-other.md index 8e1e228bc..d1cec0a39 100644 --- a/v9-to-v10-migration-guide-other.md +++ b/v9-to-v10-migration-guide-other.md @@ -18,10 +18,11 @@ ## TL;DR -- **`client.defaultWSTimeout` is gone**, and so are the `WebSocketImpl`, `wsUrlParams` and - `wsConnection` client options. Everything the WebSocket reads now lives in one configuration slice — - `client.config.set({ client: { wsConnection: { … } } })` — as `connectTimeoutMs`, `pingIntervalMs`, - `healthCheckGracePeriodMs`, `webSocketImpl`, `urlParams` and `connection`. Unlike the fields and +- **`client.defaultWSTimeout` is gone**, and so are the `WebSocketImpl`, `wsUrlParams`, + `wsConnection` and `enableWSFallback` client options. Everything the WebSocket reads now lives in one + configuration slice — `client.config.set({ client: { wsConnection: { … } } })` — as + `connectTimeoutMs`, `enableWSFallback`, `pingIntervalMs`, `healthCheckGracePeriodMs`, + `webSocketImpl`, `urlParams` and `connection`. Unlike the fields and options they replace, these survive a reconnect. - **Server-sent dates are unix-nanosecond `number`s** on every response and event type — not `Date` objects and not ISO strings, while outgoing **request** date fields are still `Date`. `new Date(ns)` @@ -37,9 +38,9 @@ means if you deploy the WebSocket client on Node 18 or 20. - **Server-side is gone.** If you construct with a `secret` or call server-only admin endpoints, switch to `@stream-io/node-sdk`. The construction guide has the full list — every feature module below that was server-only is dropped for the same reason. - Two barrels removed from the package root, one added: **`./events` and `./base64` are gone; `./logger` is new.** `./signing` survives with exactly one export left, `UserFromToken`. The `./campaign`, `./channel_batch_updater`, and `./segment` barrels are still exported but the modules are emptied (they contain only a comment pointing at the server SDK) — importing anything by name from them will fail. -- **`connection.changed` is removed.** Connectivity is published as two reactive stores, `client.wsConnection.state` for this client's socket and `client.networkConnection.state` for the device's network. A handler for the event simply stops firing, with no compile error in plain JavaScript, and a "connection lost" banner has to hold a drop itself where the event used to. The socket's own `isHealthy` is unchanged; what moved is where you read it. See below. +- **`connection.changed` is removed.** Connectivity is published as two reactive stores, `client.wsConnection.state` for this client's socket (or its long-poll, once `enableWSFallback` has switched to it) and `client.networkConnection.state` for the device's network. A handler for the event simply stops firing, with no compile error in plain JavaScript, and a "connection lost" banner has to hold a drop itself where the event used to. The socket's own `isHealthy` is unchanged; what moved is where you read it. See below. - **Watching waits instead of degrading.** A request that watches a channel or subscribes to presence is held until the WebSocket handshake produces the connection id the server keys that subscription by, rather than being sent without one and silently registering nothing. It throws only when no socket is open and none is being opened. An explicit `watch: false` is never held. See below. -- **The WebSocket connect endpoint moved to `/api/v2/connect`.** The hello event is now `connection.ok` rather than `health.check`, and the long-poll fallback (`enableWSFallback`, `transport.changed`) is gone. +- **The WebSocket connect endpoint moved to `/api/v2/connect`.** The hello event is now `connection.ok` rather than `health.check`. The long-poll fallback works as in v9, against `/api/v2/longpoll`, with `enableWSFallback` moved into the `wsConnection` configuration and `client.defaultWSTimeoutWithFallback` gone. - `Event` (type name) is kept, but its shape widened: `Event = WSEvent | ConnectedEvent | LocalEvent | keyof CustomEventTypes`. `EventPayload<''>` narrows to a specific event. - `EventTypes` (plural) renamed to `EventType` (singular). `CustomEventTypes` interface is unchanged — augment it to add custom event-type keys, same as v9. - Filter payloads now carry **per-endpoint operator constraints** (inline `Filters<{ … }>` on each request type) — previously-permissive filter objects may stop type-checking. Only one operator per field is allowed, and `null` is no longer a valid `$in` element. `QueryPollsFilters`, `QueryVotesFilters`, and `ReminderFilters` were the last hand-written holdouts and now derive from their request types too. @@ -296,9 +297,11 @@ client.networkConnection.state.getLatestValue(); // { isOnline, lastOnlineAt, lastOfflineAt } ``` -`client.wsConnection` is this client's WebSocket. `client.networkConnection` is the **device's** -network status, reported by a platform reporter you install — a separate fact that routinely disagrees -with the socket in both directions. Neither is derived from the other. +`client.wsConnection` is this client's WebSocket — or, once `enableWSFallback` has switched to +long-polling, the long-poll, which writes the same store, just as v9's long-poll dispatched +`connection.changed`. `client.networkConnection` is the **device's** network status, reported by a +platform reporter you install — a separate fact that routinely disagrees with the socket in both +directions. Neither is derived from the other. ### What breaks @@ -436,8 +439,9 @@ An explicit `watch: false` is still honoured, and a request that asks for neithe is never held at all, which is what makes it usable offline. An abort signal reaches the wait as well as the request, so an abandoned query does not sit on it. -This covers requests the old guard never did — stop-watching and long polling carry no `watch` flag -and are recognised by their declared `connection_id` parameter instead. +`channel.stopWatching()` carries no flag but is held the same way, since it tells the server which +connection should stop watching: it always waits for the connection id — rejecting, like a watching +request, when there is no connection and none is being established — and honours the abort signal. ### `connection.recovered` is withheld when the socket drops mid-recovery @@ -555,26 +559,43 @@ visible in two places: Anonymous connections are unaffected: `connectAnonymousUser()` works the same way. -### Long-poll fallback removed +### Long-poll fallback (`enableWSFallback`) -The HTTP long-poll transport that v9 fell back to when the WebSocket failed is gone. Removed -with it: +Works as in v9, but the flag moved off `StreamChatOptions` into the `wsConnection` configuration, +with the socket's other settings. Left in the constructor options, it is silently ignored: -| Removed | Notes | -| ------------------------------------------------------------------------ | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `enableWSFallback` client option | No replacement. | -| `WSConnectionFallback`, `ConnectionState` (`src/connection_fallback.ts`) | Never exported from the package root; internal. | -| `transport.changed` event | Only ever dispatched when switching to long-poll. Remove listeners. | -| `client.defaultWSTimeoutWithFallback` | `client.defaultWSTimeout` (15s) is now the only connect timeout. In v9, enabling the fallback shortened the WS timeout to 6s so the fallback could take over sooner; without it, connects get the full 15s. | +```ts +// v9 +const client = new StreamChat(API_KEY, { enableWSFallback: true }); -```diff -- const client = new StreamChat(API_KEY, { enableWSFallback: true }); -- client.on('transport.changed', ({ mode }) => reportTransport(mode)); -+ const client = new StreamChat(API_KEY); +// v10 +const client = new StreamChat(API_KEY); +client.config.set({ client: { wsConnection: { enableWSFallback: true } } }); ``` -The WebSocket's own reconnect and health-check loop is unchanged and still handles transient -network failures. If you need to react to connectivity, subscribe to `client.wsConnection.state`. +With it on, the WebSocket connects as it would without the flag, within `connectTimeoutMs` (15s by +default). v9 gave it 6s while the flag was on, through `client.defaultWSTimeoutWithFallback`, which is +gone: lower `connectTimeoutMs` to switch sooner. If the WebSocket fails with a network error, the +client switches to HTTP long-polling, dispatches `connection.fallback_activated` with +`mode: 'longpoll'`, and stays on long-poll for the rest of its lifetime. What else changed around it: + +- It polls `/api/v2/longpoll` instead of `/longpoll`. +- The switch event is `connection.fallback_activated`, renamed from `transport.changed`. Its payload + is unchanged: `mode: 'longpoll'`. +- Its status is `client.wsConnection.state`, as for the WebSocket, since `connection.changed` is gone. +- A reconnect recovers state through the usual connection recovery; `client.recoverState()` is gone. +- It takes online/offline changes from `client.networkConnection` instead of `window` events, so a + network status reporter you install, such as one wrapping NetInfo, reaches it too. It ignores them + while closed by `closeConnection()` / `disconnectUser()` and follows them again once reconnected; + v9 stopped listening for good at its first disconnect. +- `client.wsFallback` moved to `client.wsConnection.fallback`, beside the socket. + `client.wsConnection.isConnecting` reports the long-poll's connect attempts once switched. +- State is also recovered after `closeConnection()` → `openConnection()`, which v9's long-poll never + did: recovery follows the status store rather than being called from `connect()`. +- The switch is skipped when a network status reporter says the device is offline: the browser's, + or one you installed. v9 checked only `navigator.onLine`. On hosts without a network API, the + default reporter mirrors the WebSocket, so its "offline" only means the socket is down, and the + switch goes ahead. ### Watching a channel now requires a connected user @@ -656,7 +677,8 @@ affected — it carries no flag and is connection-scoped by definition, since it which connection should stop watching. Every other request is unaffected, and — unlike v9 — no longer carries a `connection_id` query -param at all. It is now attached only to the requests listed above. +param at all. It is now attached only to the requests listed above, and to the long-poll +fallback's own polls and close. #### `closeConnection()` drops the id too @@ -670,6 +692,11 @@ kept reporting an id for a socket that was gone. Read `client._getConnectionID() `client.connectionIdManager.connectionId`) instead — both are dropped the moment the connection stops being healthy. +`client.wsConnection.fallback.connectionID` (v9's `client.wsFallback.connectionID`) is **removed** +for the same reason: it was not cleared when the long-poll went down, only when it reconnected or +was disconnected. The long-poll reads its id from `client.connectionIdManager` like everything else, +so read `client._getConnectionID()` for the long-poll's id too. + The consequence for mobile apps: `closeConnection()` (the documented background/foreground seam) now makes the gated calls above throw until `openConnection()` has been called, even though the user is still set. Reopen the connection before issuing them, or pass `watch: false` for loads that @@ -1295,7 +1322,7 @@ For each source file that touches the SDK: 17. **Fix upload call sites.** `channel.sendFile` / `sendImage` are now `channel.uploadFile` / `uploadImage`, and take a request object: `{ file }`, where `file` is a `File`, a `Blob`, or a React-Native `{ uri, name, type }` descriptor — no `Buffer`, no readable streams. The MIME type still has to be explicit on the React-Native path, it just lives on the descriptor rather than in a separate `contentType` argument. `axiosRequestConfig` becomes `requestOptions` (`{ onUploadProgress, signal }`), and the routes moved to `/api/v2/…`. 18. **Delete bundler shims** added for `stream-chat`'s Node-only deps (`crypto`, `https`, `zlib`, `jsonwebtoken`, `ws`) — `package.json#browser` is gone because nothing imports them anymore. 19. **Handle the new connect hello event.** Anything keyed on the _first_ `health.check` (seeding `client.user`, unread counts, "connected" UI state) should listen for `connection.ok` instead; periodic `health.check` events are unchanged. Narrow on `event.type` before reading fields off the resolved `ConnectionOpen`. -20. **Drop long-poll fallback code.** Remove `enableWSFallback` from client options, delete `transport.changed` listeners, and delete reads of `client.defaultWSTimeoutWithFallback`. +20. **Move `enableWSFallback`** from the client options to `client.config.set({ client: { wsConnection: { enableWSFallback: true } } })`. Rename `transport.changed` listeners to `connection.fallback_activated`; the payload (`mode: 'longpoll'`) is unchanged. Delete reads of `client.defaultWSTimeoutWithFallback` (lower `connectTimeoutMs` if you relied on the faster switch), and read the long-poll's status from `client.wsConnection.state` rather than `connection.changed`. 21. **Replace `client.setLocalDevice(device)` / the `device` client option** with an explicit `await client.createDevice({ id, push_provider, push_provider_name? })` after connecting. 22. **Polyfill `atob`** if your React Native / Hermes target lacks it (`typeof atob === 'undefined'`); `UserFromToken` depends on it during `connectUser`. 23. **Call `liveLocationManager.dispose()`** when you are finished with a manager you constructed, alongside whatever `unregisterSubscriptions()` you already call. Nothing will fail to compile: `dispose()` is the _configuration_ teardown, and until it runs the client's configuration registry holds a handle to the manager — a long-lived client and many short-lived managers will accumulate them. `unregisterSubscriptions()` is unchanged and stays ref-counted, so it deliberately no longer releases configuration; it never should have, since with two callers sharing a manager the first to leave stopped a still-live instance from tracking `client.config`. `SearchController` already worked this way.