diff --git a/.gitignore b/.gitignore index c6bba59..fc21710 100644 --- a/.gitignore +++ b/.gitignore @@ -128,3 +128,6 @@ dist .yarn/build-state.yml .yarn/install-state.gz .pnp.* + +# Git worktrees +.worktrees/ diff --git a/README.md b/README.md index 77968b2..56b8ab4 100644 --- a/README.md +++ b/README.md @@ -19,17 +19,21 @@ Note that [`@comapeo/core`](https://github.com/digidem/comapeo-core) is a peer d npm install @comapeo/ipc @comapeo/core ``` +> **Release order for v10**: this version depends on `rpc-reflector` `^4.5.0`, which is not yet published — 4.5.0 must be released from its branch HEAD (which includes a fix landed after the version-bump commit) and this package's lockfile regenerated before v10 can be released. Until then `npm ci` fails; develop against a local checkout via `npm install --no-save`. + ## API -### `createComapeoCoreServer(manager: MapeoManager, messagePort: MessagePortLike): { close: () => void }` +### `createComapeoCoreServer(manager: MapeoManager, messagePort: MessagePortLike, opts?: { logger?, onRequestHook? }): { close: () => void }` Creates the IPC server instance. `manager` is a `@comapeo/core` `MapeoManager` instance and `messagePort` is an interface that resembles a [`MessagePort`](https://developer.mozilla.org/en-US/docs/Web/API/MessagePort). +`opts` is passed through to each underlying [`rpc-reflector` server](https://github.com/digidem/rpc-reflector) (manager, project routing, and per-project): `opts.logger` enables logging (`false`, the default, disables it; or pass a pino-compatible logger / the global `console`), and `opts.onRequestHook` observes each request and its response. Note that on the manager server a consumer-supplied `onRequestHook` is wrapped by the interim `leaveProject` hook (see [Lifecycle](#lifecycle)): the consumer hook runs as supplied, and completed `leaveProject` calls additionally trigger the stale-instance cleanup. + Returns an object with a `close()` method, which removes relevant event listeners from the `messagePort`. Does not close or destroy the `messagePort`. -### `createComapeoCoreClient(messagePort: MessagePortLike, opts?: { timeout?: number }): ClientApi` +### `createComapeoCoreClient(messagePort: MessagePortLike, opts?: { timeout?: number, logger? }): ClientApi` -Creates the IPC client instance. `messagePort` is an interface that resembles a [`MessagePort`](https://developer.mozilla.org/en-US/docs/Web/API/MessagePort). `opts.timeout` is an optional timeout used for sending and receiving messages over the channel. +Creates the IPC client instance. `messagePort` is an interface that resembles a [`MessagePort`](https://developer.mozilla.org/en-US/docs/Web/API/MessagePort). `opts` is passed through to each underlying [`rpc-reflector` client](https://github.com/digidem/rpc-reflector): `opts.timeout` is an optional per-call timeout for messages over the channel; `opts.logger` as for the server. Returns a client instance that reflects the interface of the `manager` provided to [`createComapeoCoreServer`](#createcomapeocoreservermanager-mapeomanager-messageport-messageportlike--close---void). Refer to the [`rpc-reflector` docs](https://github.com/digidem/rpc-reflector#const-clientapi--createclientchannel) for additional information about how to use this. @@ -51,6 +55,14 @@ Creates the services client, reflecting the `services` object passed to [`create Closes the services client. Does not close or destroy the `messagePort`. +### `notifyTransportReset(client): void` + +Tell a client (core or services) that its transport to the server has dropped: every in-flight call rejects immediately with [`RpcChannelClosedError`](#errors) instead of waiting out its timeout. The client remains fully usable. See [Transport reset](#transport-reset). + +### `resubscribe(client): void` + +Re-send a client's (core or services) event subscriptions to the server, once the transport to a restarted server is connected again. See [Transport reset](#transport-reset). + ## Behaviour These are the guarantees the wrappers add on top of [`rpc-reflector`](https://github.com/digidem/rpc-reflector); they are exercised by the test suite. @@ -76,17 +88,31 @@ The wrappers never close or destroy the `messagePort` itself — that is the cal `client.getProject(id)` resolves with a client that reflects the `MapeoProject` API, including nested namespaces such as `project.observation.*`. -- **Deduplicated.** Concurrent or repeated `getProject(id)` calls for an open project resolve to the same reference and open the project only once on the server. +- **Deduplicated.** Concurrent or repeated `getProject(id)` calls resolve to the same reference — project references are permanent for the lifetime of the client. The id is validated with the server (and the project eagerly opened) only on **first acquisition**; after one success the cached wrapper is returned with no wire round trip. - **Missing projects.** If the project does not exist, `getProject(id)` rejects with `NotFoundError` (from `@comapeo/core`). A failed lookup is not cached, so a later call for an id that does exist still succeeds. -- **Isolation.** Closing one project does not affect other open projects. +- **Left projects.** If this device has left the project (`manager.leaveProject`), every method call on a held reference rejects with [`ProjectLeftError`](#errors) until the project is re-joined via an invite, after which the same reference works again. `getProject(id)` rejects the same way on first acquisition; for a project acquired before leaving it still resolves with the cached wrapper (whose calls then reject). ### Lifecycle -- `project.close()` closes the project on the server and tears down its channel. It is idempotent — repeated calls resolve like the first. -- After a project is closed — via `project.close()` **or** by the server closing it — every method on that reference rejects with [`ProjectClosedError`](#errors). -- A project can be re-opened: after closing, `getProject(id)` opens a fresh instance and returns a new reference. Calls on the old, closed reference never reach the re-opened project — they keep rejecting. -- `closeComapeoCoreClient(client)` tears down the manager, the project-routing channel, and every open project reference. After this, all calls — including `getProject(id)` — reject with [`ClientClosedError`](#errors). (The services client is independent; close it separately with [`closeComapeoServicesClient`](#closecomapeoservicesclientservicesclient-clientapicomapeoservicesapi-void).) -- Calls already in flight when a close happens reject with [`RpcChannelClosedError`](#errors); they are not re-routed. +Project instance lifecycle is owned entirely by the server. The client cannot close a project (the reflected surface has no `project.close()`), and a project reference never goes stale. + +The server delegates the lifecycle mechanics to rpc-reflector's late-bound handlers: each project channel has one long-lived rpc-reflector server whose handler — the live `MapeoProject` instance — is bound lazily by a factory when the first call or subscription arrives, and detached when the instance closes. rpc-reflector keeps event subscriptions in a registry that outlives the instance and re-attaches them to each fresh instance before any waiting call is dispatched. Concretely: + +- The server may close a project instance at any time (resource management, `addProject` re-joining a previously-left project, a server restart). The next call on that project's channel transparently re-opens it — callers never observe the cycle. +- Event subscriptions survive server-side close/re-open, and subscribing to a project whose instance is closed re-opens it — a consumer that only listens still receives events. +- The one exception is a left project, which is never re-opened — see [`ProjectLeftError`](#errors). +- A `leaveProject` call routed through this server also closes the stale instance `@comapeo/core` leaves cached after leaving (core only cleans that up itself inside `addProject`). This hook and the server's left-project guard are an interim pair, removed together once core ships a typed PROJECT_LEFT error ([digidem/comapeo-core#1313](https://github.com/digidem/comapeo-core/issues/1313)). +- `closeComapeoCoreClient(client)` tears down the manager, the project-routing channel, and every project reference. After this, all calls — including `getProject(id)` — reject with [`ClientClosedError`](#errors). (The services client is independent; close it separately with [`closeComapeoServicesClient`](#closecomapeoservicesclientservicesclient-clientapicomapeoservicesapi-void).) +- Calls already in flight when the client closes reject with [`RpcChannelClosedError`](#errors); they are not re-routed. + +### Transport reset + +When the process hosting the server dies and restarts while the client stays alive (e.g. Android's foreground service being killed), the transport owner should drive a two-phase recovery: + +- At drop time, call `notifyTransportReset(client)` (on the core client and, if used, the services client): every in-flight call rejects immediately with [`RpcChannelClosedError`](#errors) (`code: 'RPC_CHANNEL_CLOSED'`) instead of waiting out its timeout. Reads are safe to retry once the transport reconnects; whether to replay a mutation is the caller's judgement — nothing is replayed automatically. +- Once the transport is connected to the restarted server, call `resubscribe(client)` (on the same clients): every event subscription — manager and per-project — is re-sent, since the fresh server has no subscription state. Resubscription is deliberately not done at drop time: ON frames written into a down transport can keep nudging it into reconnect attempts while the server stays down. + +Project references need no recovery: their channels are keyed by project id, which a restarted server serves identically — the next call (or a replayed subscription) transparently re-opens the project. Both functions are safe to call repeatedly and are no-ops after the client is closed. ### Events @@ -98,21 +124,19 @@ Error classes are available from the `@comapeo/ipc/errors.js` entrypoint: ```ts import { - ProjectClosedError, + ProjectLeftError, ClientClosedError, RpcChannelClosedError, RpcTimeoutError, } from '@comapeo/ipc/errors.js' ``` -After a reference is closed, calls made on it reject with a descriptive error: - -- **`ProjectClosedError`** (`code: 'PROJECT_CLOSED'`) — a method (including nested namespaces such as `project.observation.*`) was called on a project reference after that project was closed, either via `await project.close()` or by the server closing the project. A re-opened reference from a fresh `client.getProject(id)` is unaffected. +- **`ProjectLeftError`** (`code: 'PROJECT_LEFT'`) — a method (including nested namespaces such as `project.observation.*`) or `getProject(id)` was called for a project this device has left. Left projects are never transparently re-opened; re-joining via an invite makes the same reference usable again. - **`ClientClosedError`** (`code: 'CLIENT_CLOSED'`) — a method was called on the CoMapeo client, or on any project reference, after the whole client was torn down with [`closeComapeoCoreClient`](#closecomapeocoreclientclient-clientapimapeomanager-promisevoid). This includes `getProject(id)`, which after close rejects with `ClientClosedError` rather than returning a reference — whether or not that project was fetched earlier. -RPC methods return a rejected `Promise` carrying the error, so failures surface through normal `await`/`.catch()` handling. The exception is the event-emitter methods (`on`, `once`, `off`, `removeListener`, `emit`, etc.), which return synchronously rather than a promise — after close these **throw** the same error synchronously instead, so it surfaces at the call site rather than as an unhandled rejection. +RPC methods return a rejected `Promise` carrying the error, so failures surface through normal `await`/`.catch()` handling. The exception is the event-emitter methods, which return synchronously rather than a promise — after the client is closed, subscribe methods (`on`, `once`, `addListener`, and `emit`/introspection) **throw** `ClientClosedError` synchronously so the failure surfaces at the call site rather than as an unhandled rejection, while unsubscribe methods (`off`, `removeListener`, `removeAllListeners`) are safe no-ops — removing a listener from a dead client is correct teardown. -Calls that were already in flight when the close happened are not re-routed: they reject with **`RpcChannelClosedError`** as the underlying channel tears down. `RpcTimeoutError` is thrown when a call exceeds the `opts.timeout` passed to [`createComapeoCoreClient`](#createcomapeocoreclientmessageport-messageportlike-opts--timeout-number--clientapimapeomanager). +Calls that were already in flight when the close happened are not re-routed: they reject with **`RpcChannelClosedError`** (`code: 'RPC_CHANNEL_CLOSED'`) as the underlying channel tears down. The same error rejects in-flight calls when [`notifyTransportReset`](#transport-reset) is called. `RpcTimeoutError` is thrown when a call exceeds the `opts.timeout` passed to [`createComapeoCoreClient`](#createcomapeocoreclientmessageport-messageportlike-opts--timeout-number--clientapimapeomanager). ## Usage diff --git a/package-lock.json b/package-lock.json index 4e38f99..ae3eab5 100644 --- a/package-lock.json +++ b/package-lock.json @@ -11,7 +11,7 @@ "dependencies": { "custom-error-creator": "^1.4.0", "p-defer": "^4.0.1", - "rpc-reflector": "^4.2.0" + "rpc-reflector": "github:digidem/rpc-reflector#fec43e018e553e5ce1652c6495ef84e38a250750" }, "devDependencies": { "@comapeo/core": "^7.2.0", @@ -7049,9 +7049,9 @@ "license": "MIT" }, "node_modules/rpc-reflector": { - "version": "4.2.0", - "resolved": "https://registry.npmjs.org/rpc-reflector/-/rpc-reflector-4.2.0.tgz", - "integrity": "sha512-n1OQjszyYFiGJ5QXjgL2jc3kfISm5bNXxFajQY663frS0lqvpV2L8gvbCLY7slsz1tN3wn9eFByvpvLII1yLdg==", + "version": "4.4.0", + "resolved": "git+ssh://git@github.com/digidem/rpc-reflector.git#fec43e018e553e5ce1652c6495ef84e38a250750", + "integrity": "sha512-NBPABe0F7RwEyYKqwjqeWNqpqdBeyiGGuIUYr7lRQgNVwssGAEGEufqGuz+5D+JKFAvlAEGXeBuND8mu9vYK7w==", "license": "ISC", "dependencies": { "abstract-logging": "^2.0.1", diff --git a/package.json b/package.json index 924f091..4b0917d 100644 --- a/package.json +++ b/package.json @@ -58,7 +58,7 @@ "dependencies": { "custom-error-creator": "^1.4.0", "p-defer": "^4.0.1", - "rpc-reflector": "^4.2.0" + "rpc-reflector": "github:digidem/rpc-reflector#fec43e018e553e5ce1652c6495ef84e38a250750" }, "peerDependencies": { "@comapeo/core": "^7.0.1" diff --git a/src/client.js b/src/client.js index 5cb3524..82bd87c 100644 --- a/src/client.js +++ b/src/client.js @@ -2,11 +2,12 @@ import { createClient } from 'rpc-reflector/client.js' import { MANAGER_CHANNEL_ID, + PROJECT_CHANNEL_PREFIX, PROJECT_ROUTING_ID, SERVICES_ID, SubChannel, } from './lib/sub-channel.js' -import { ClientClosedError, ProjectClosedError } from './errors.js' +import { ClientClosedError, RpcChannelClosedError } from './errors.js' /** @import { ClientApi, MessagePortLike } from 'rpc-reflector' */ /** @import { MapeoProject, MapeoManager } from '@comapeo/core' */ @@ -16,13 +17,13 @@ import { ClientClosedError, ProjectClosedError } from './errors.js' // synchronously (they return the client/an array/a number, never a promise). // Mirrors the method set rpc-reflector treats specially (`prop in // EventEmitter.prototype`). -const EMITTER_METHODS = new Set([ - 'addListener', - 'on', - 'once', +const SUBSCRIBE_METHODS = new Set(['addListener', 'on', 'once']) +const UNSUBSCRIBE_METHODS = new Set([ 'removeListener', 'off', 'removeAllListeners', +]) +const OTHER_EMITTER_METHODS = new Set([ 'emit', 'eventNames', 'listeners', @@ -30,22 +31,34 @@ const EMITTER_METHODS = new Set([ ]) /** - * Build the Proxy returned for a closed client/project reference. Method calls - * (including nested namespaces such as `project.observation.*`) reject with - * `makeError()`, keeping the `Promise`-returning contract callers expect. - * EventEmitter methods are the exception: callers don't await them, so a - * rejected promise would surface as an unhandled rejection — they throw - * synchronously instead, at the call site. + * Build the Proxy returned for a reference after the whole client is closed. + * Method calls (including nested namespaces such as `project.observation.*`) + * reject with `makeError()`, keeping the `Promise`-returning contract callers + * expect. EventEmitter methods are the exception: callers don't await them, + * so a rejected promise would surface as an unhandled rejection. Subscribe + * methods throw synchronously at the call site; unsubscribe methods are + * chainable no-ops — removing a listener from a dead client is correct + * teardown (React effect cleanup runs against stale references). * * @param {() => Error} makeError */ function createClosedProxy(makeError) { + /** @type {any} */ + let proxy /** @type {ProxyHandler} */ const handler = { get(_target, prop) { - if (typeof prop === 'string' && EMITTER_METHODS.has(prop)) { - return () => { - throw makeError() + // A truthy `then` would make every nested node a thenable whose + // callbacks are never invoked — awaiting a namespace would hang. + if (prop === 'then') return undefined + if (typeof prop === 'string') { + if (SUBSCRIBE_METHODS.has(prop) || OTHER_EMITTER_METHODS.has(prop)) { + return () => { + throw makeError() + } + } + if (UNSUBSCRIBE_METHODS.has(prop)) { + return () => proxy } } return new Proxy(function () {}, handler) @@ -57,11 +70,12 @@ function createClosedProxy(makeError) { return Promise.reject(makeError()) }, } - return new Proxy({}, handler) + proxy = new Proxy({}, handler) + return proxy } /** - * @typedef {ClientApi} ComapeoProjectClientApi + * @typedef {Omit, 'close'>} ComapeoProjectClientApi */ /** @@ -75,8 +89,20 @@ function createClosedProxy(makeError) { * >} ComapeoCoreClientApi */ const CLOSE = Symbol('close') +const TRANSPORT_RESET = Symbol('transportReset') +const RESUBSCRIBE = Symbol('resubscribe') /** + * Create the client side of `createComapeoCoreServer`. + * + * Project references returned by `getProject` are permanent: each is bound to + * a channel keyed by the project's public id, which the server keeps valid + * across instance close/re-open cycles (and across server restarts). There is + * no client-visible project lifecycle — no `close()`, and no reference ever + * goes stale. Calls to a project this device has left reject with + * `ProjectLeftError` (server-side); calls to an unknown project reject with + * core's `NotFoundError`. + * * @param {MessagePortLike} messagePort * @param {Parameters[1]} [opts] * @@ -84,27 +110,22 @@ const CLOSE = Symbol('close') */ export function createComapeoCoreClient(messagePort, opts = {}) { /** - * projectPublicId → wrapper bound to a specific instance id. Only returned - * after the server confirms that instance is still current — the server - * can close a project without the client asking (e.g. `leaveProject`). - * @type {Map, - * }>} + * projectPublicId → permanent wrapper. Never evicted: the channel id is + * stable, so the wrapper stays valid for the lifetime of this client. + * @type {Map} */ - const currentProjectClients = new Map() + const projectClients = new Map() /** * projectPublicId → in-flight `getProject`. Dedupes concurrent calls; * entries are removed on settle so later calls re-validate. - * @type {Map>>} + * @type {Map>} */ const pendingProjectClients = new Map() /** - * The rpc-reflector client + SubChannel pair for every currently-open - * project. Entries are removed when the project's wrapped `close()` - * settles; `closeComapeoCoreClient` sweeps whatever is left. + * The rpc-reflector client + SubChannel pair for every project wrapper, + * swept by `closeComapeoCoreClient`. * @type {Set<{ * client: ClientApi, * channel: SubChannel, @@ -123,16 +144,61 @@ export function createComapeoCoreClient(messagePort, opts = {}) { projectRoutingChannel.start() managerChannel.start() - // Set once `closeComapeoCoreClient` has torn the whole client down. Read by the - // manager proxy and the per-project wrappers so that calls after close - // surface `ManagerClosedError` instead of rpc-reflector's `ChannelClosed`. + // Set once `closeComapeoCoreClient` has torn the whole client down. Read by + // the manager proxy and the per-project wrappers so that calls after close + // surface `ClientClosedError` instead of rpc-reflector's `ChannelClosed`. let clientClosed = false - const managerClosedProxy = createClosedProxy(() => new ClientClosedError()) + // Set at the START of the close routine: a getProject whose validation + // round trip resolves during the close's await window must not create a + // wrapper after the sweep (leaking a port listener and an unclosed rpc + // client) — it rejects instead. + let clientClosing = false + const clientClosedProxy = createClosedProxy(() => new ClientClosedError()) + + function handleTransportReset() { + if (clientClosed) return + // Fail in-flight calls fast with the channel-closed error instead of + // leaving them to hit the per-call timeout. Resubscription is + // deliberately NOT done here: at drop time the transport is down, and + // each ON frame written into it can nudge the native transport into + // retrying forever while the server stays down — the consumer calls + // `resubscribe` once the transport is back up. Project references stay + // valid: their channels are keyed by project id, which a restarted + // server serves identically. + createClient.rejectPending(managerClient, new RpcChannelClosedError()) + createClient.rejectPending( + projectRoutingClient, + new RpcChannelClosedError(), + ) + for (const entry of openProjectClients) { + createClient.rejectPending(entry.client, new RpcChannelClosedError()) + } + } + + function handleResubscribe() { + if (clientClosed) return + // Safe to call repeatedly: the server ignores duplicate ON messages. A + // replayed project subscription also re-opens that project server-side — + // an active listener is an expression of interest. + createClient.resubscribe(managerClient) + for (const entry of openProjectClients) { + createClient.resubscribe(entry.client) + } + } const client = new Proxy(managerClient, { get(target, prop, receiver) { + // Resolved ahead of the closed check: both are safe no-ops after close. + if (prop === TRANSPORT_RESET) { + return handleTransportReset + } + if (prop === RESUBSCRIBE) { + return handleResubscribe + } + if (prop === CLOSE) { return async () => { + clientClosing = true managerChannel.close() createClient.close(managerClient) @@ -147,6 +213,7 @@ export function createComapeoCoreClient(messagePort, opts = {}) { entry.channel.close() } openProjectClients.clear() + projectClients.clear() // Closed last so in-flight `assertProjectExists` calls awaited // above can complete rather than reject. @@ -158,13 +225,13 @@ export function createComapeoCoreClient(messagePort, opts = {}) { } if (prop === 'getProject') { - return createProjectClient + return getProject } // `then` must stay falsy so awaiting the client (a thenable check) does // not route into the throwing proxy. if (clientClosed && prop !== 'then') { - return Reflect.get(managerClosedProxy, prop) + return Reflect.get(clientClosedProxy, prop) } return Reflect.get(target, prop, receiver) @@ -178,11 +245,18 @@ export function createComapeoCoreClient(messagePort, opts = {}) { * @param {string} projectPublicId * @returns {Promise} */ - async function createProjectClient(projectPublicId) { + async function getProject(projectPublicId) { // Checked before the cache lookup so `getProject` rejects uniformly after // close — whether or not this id was fetched (and cached) earlier. if (clientClosed) throw new ClientClosedError() + // The wire round trip (existence check + eager server-side open) happens + // only on first acquisition; after one success the permanent wrapper is + // returned directly. A project left after acquisition surfaces on method + // calls (which reject with ProjectLeftError over the wire), not here. + const existing = projectClients.get(projectPublicId) + if (existing) return existing + const pending = pendingProjectClients.get(projectPublicId) if (pending) return pending @@ -196,96 +270,56 @@ export function createComapeoCoreClient(messagePort, opts = {}) { } /** - * Return the cached wrapper only if the server confirms its instance id is - * still current; otherwise build a fresh one. The server evicts its routing - * entry synchronously on close, so this is correct even before the close - * event reaches this client. - * * @param {string} projectPublicId - * @returns {Promise>} + * @returns {Promise} */ async function resolveProjectClient(projectPublicId) { - const instanceId = - await projectRoutingClient.assertProjectExists(projectPublicId) - - const current = currentProjectClients.get(projectPublicId) - if (current && current.instanceId === instanceId) return current.wrapper - - const wrapper = createProjectClientWrapper(projectPublicId, instanceId) - currentProjectClients.set(projectPublicId, { instanceId, wrapper }) + // One round trip on first acquisition, so a bad id rejects here (with + // `NotFoundError` / `ProjectLeftError`) rather than on the first method + // call, and so the server opens the project eagerly — attaching + // subscriptions before any project-channel frame. A failure is not + // cached; the next `getProject` retries. + await projectRoutingClient.assertProjectExists(projectPublicId) + + // The close routine may have started while the round trip above was in + // flight; its sweep only covers wrappers that already exist, so don't + // create one it can never clean up. + if (clientClosing) throw new ClientClosedError() + + const wrapper = createProjectClientWrapper(projectPublicId) + projectClients.set(projectPublicId, wrapper) return wrapper } /** * @param {string} projectPublicId - * @param {string} instanceId - * @returns {ClientApi} + * @returns {ComapeoProjectClientApi} */ - function createProjectClientWrapper(projectPublicId, instanceId) { - // Per-project messages are scoped to the current open instance, not the - // project's public id. If this project is closed and re-opened later, - // `assertProjectExists` returns a different instance id, so the new - // wrapper uses a fresh SubChannel that can't collide with the old one. - const projectChannel = new SubChannel(messagePort, instanceId) + function createProjectClientWrapper(projectPublicId) { + const projectChannel = new SubChannel( + messagePort, + `${PROJECT_CHANNEL_PREFIX}${projectPublicId}`, + ) /** @type {ClientApi} */ const projectClient = createClient(projectChannel, opts) projectChannel.start() - const registryEntry = { client: projectClient, channel: projectChannel } - openProjectClients.add(registryEntry) - - // Wrap projectClient to intercept `close`: after the wire close settles, - // tear down the local client + channel — rejecting any in-flight calls. - // Cache eviction is in the 'close' listener below, which also covers - // manager-initiated closes. - // Further method calls on this wrapper reject with `ProjectClosedError`. - // The close promise is cached so repeated `close()` calls return the - // same result instead of failing on the already-closed channel. All - // other property accesses delegate to the inner client unchanged. - /** @type {Promise | null} */ - let closePromise = null - let closed = false - // After this reference is closed, any method (including nested namespaces) - // throws a descriptive error rather than rpc-reflector's `ChannelClosed`: - // `ProjectClosedError` when this project was closed, `ManagerClosedError` - // when the whole client was torn down. In-flight calls at close time are - // left to reject with `ChannelClosed` — they were already on the wire. - const closedProxy = createClosedProxy(() => - closed ? new ProjectClosedError() : new ClientClosedError(), - ) + openProjectClients.add({ client: projectClient, channel: projectChannel }) + const wrappedProjectClient = new Proxy(projectClient, { get(target, prop, receiver) { - if (prop === 'close') { - return () => { - closePromise ??= (async () => { - try { - await target.close() - } finally { - createClient.close(projectClient) - projectChannel.close() - } - })() - return closePromise - } - } - if ((closed || clientClosed) && prop !== 'then') { - return Reflect.get(closedProxy, prop) + // Project lifecycle is server-owned: the reflected surface must not + // expose `MapeoProject.close`, which would close the server-side + // instance out from under every other consumer. + if (prop === 'close') return undefined + if (clientClosed && prop !== 'then') { + return Reflect.get(clientClosedProxy, prop) } return Reflect.get(target, prop, receiver) }, }) - wrappedProjectClient.once('close', () => { - closed = true - // A late close event must not evict a newer wrapper cached for the - // re-opened instance. - const current = currentProjectClients.get(projectPublicId) - if (current?.wrapper === wrappedProjectClient) { - currentProjectClients.delete(projectPublicId) - } - openProjectClients.delete(registryEntry) - }) - return wrappedProjectClient + return /** @type {any} */ (wrappedProjectClient) } } @@ -298,6 +332,55 @@ export async function closeComapeoCoreClient(client) { return client[CLOSE]() } +/** + * Notify a client that its transport to the server has dropped (e.g. the + * process hosting the server died): every in-flight call rejects immediately + * with rpc-reflector's `ChannelClosedError` (`code: 'RPC_CHANNEL_CLOSED'`, + * re-exported as `RpcChannelClosedError`) instead of waiting out its + * timeout. Accepts a client from `createComapeoCoreClient` (rejecting the + * manager, project-routing, and every project reference's in-flight calls) + * or from `createComapeoServicesClient`. The client remains fully usable; + * project references stay valid and serve the restarted server once the + * transport reconnects. Safe to call repeatedly and on a closed client + * (no-op). + * + * Two-phase rule: call this at drop time; deliberately do NOT replay event + * subscriptions until the transport is connected to the restarted server — + * then call {@link resubscribe}. ON frames written into a down transport can + * keep nudging it into reconnect attempts while the server stays down. + * + * @param {ComapeoCoreClientApi | ComapeoServicesClientApi} client + */ +export function notifyTransportReset(client) { + const reset = /** @type {any} */ (client)[TRANSPORT_RESET] + if (typeof reset === 'function') { + reset() + return + } + // A services client is a bare rpc-reflector client. + createClient.rejectPending(client, new RpcChannelClosedError()) +} + +/** + * Re-send every event subscription to the server — for a core client the + * manager's and every project reference's (a replayed project subscription + * also transparently re-opens that project server-side). Call once the + * transport to a restarted server is connected again — the fresh server has + * no subscription state until then; see the two-phase rule on + * {@link notifyTransportReset}. Safe to call repeatedly (the server ignores + * duplicate subscriptions) and on a closed client (no-op). + * + * @param {ComapeoCoreClientApi | ComapeoServicesClientApi} client + */ +export function resubscribe(client) { + const replay = /** @type {any} */ (client)[RESUBSCRIBE] + if (typeof replay === 'function') { + replay() + return + } + createClient.resubscribe(client) +} + /** * @typedef {ClientApi} ComapeoServicesClientApi */ diff --git a/src/errors.js b/src/errors.js index 68a0afd..d3c6132 100644 --- a/src/errors.js +++ b/src/errors.js @@ -6,13 +6,14 @@ export { } from 'rpc-reflector/errors.js' /** - * Thrown server-side when a stale call reaches a project instance that has - * already been closed. Rides the standard rpc-reflector error response back - * to the client. + * Thrown server-side when a call arrives for a project this device has left + * (`manager.leaveProject`). Left projects are never re-opened by the server; + * re-joining via an invite (`manager.addProject`) makes the project usable + * again. Rides the standard rpc-reflector error response back to the client. */ -export const ProjectClosedError = createErrorClass({ - code: 'PROJECT_CLOSED', - message: 'Project is closed', +export const ProjectLeftError = createErrorClass({ + code: 'PROJECT_LEFT', + message: 'This device has left the project', status: 410, }) diff --git a/src/index.js b/src/index.js index 1125e16..a014da9 100644 --- a/src/index.js +++ b/src/index.js @@ -3,6 +3,8 @@ export { closeComapeoCoreClient, createComapeoServicesClient, closeComapeoServicesClient, + notifyTransportReset, + resubscribe, } from './client.js' export { createComapeoCoreServer, diff --git a/src/lib/sub-channel.js b/src/lib/sub-channel.js index 1e061fb..8122b89 100644 --- a/src/lib/sub-channel.js +++ b/src/lib/sub-channel.js @@ -9,9 +9,10 @@ export const MANAGER_CHANNEL_ID = '@@comapeo/manager' export const PROJECT_ROUTING_ID = '@@comapeo/project-routing' export const SERVICES_ID = '@@comapeo/services' -// Prefix for per-project instance channel ids; the rest of the id is the -// project's public id plus a per-open counter (see `openProjectInstance`). -export const PROJECT_INSTANCE_PREFIX = '@@comapeo/project/' +// Prefix for per-project channel ids; the rest of the id is the project's +// public id. The id is stable across close/re-open cycles and server +// restarts — the server owns instance lifecycle behind it. +export const PROJECT_CHANNEL_PREFIX = '@@comapeo/project/' /** @import {MessagePortLike, MessageEvent} from 'rpc-reflector' */ diff --git a/src/server.js b/src/server.js index a0dac51..ed85ae7 100644 --- a/src/server.js +++ b/src/server.js @@ -2,127 +2,135 @@ import { createServer } from 'rpc-reflector/server.js' import { COMAPEO_PREFIX, MANAGER_CHANNEL_ID, - PROJECT_INSTANCE_PREFIX, + PROJECT_CHANNEL_PREFIX, PROJECT_ROUTING_ID, SERVICES_ID, SubChannel, } from './lib/sub-channel.js' import { isRelevantEventData } from './lib/utils.js' -import { ProjectClosedError } from './errors.js' +import { ProjectLeftError } from './errors.js' /** @import { MessagePortLike } from 'rpc-reflector' */ +/** @import { MapeoManager, MapeoProject } from '@comapeo/core' */ + +function noop() {} /** - * @param {import('@comapeo/core').MapeoManager} manager + * Serve a `MapeoManager` (and its projects) over the shared message port. + * + * Project instance lifecycle is fully owned by this server: each project has + * one channel, keyed by its public id, that is stable across close/re-open + * cycles. Behind that channel sits a single long-lived rpc-reflector server + * created with a handler *factory* (see {@link createProjectHost}): + * rpc-reflector binds a live `MapeoProject` instance lazily, keeps client + * subscriptions across instance changes, and re-attaches them to each fresh + * instance before serving any call against it. Clients never see instance + * identity and cannot close projects. + * + * The one thing the server will not transparently re-open is a project this + * device has left: calls to a left project reject with `ProjectLeftError` + * until the project is re-joined (`addProject` on re-invite). A `leaveProject` + * call routed through this server also closes the gutted instance that core + * leaves cached (core only cleans that up itself inside `addProject`). + * + * @param {MapeoManager} manager * @param {MessagePortLike} messagePort * @param {Parameters[2]} [opts] */ export function createComapeoCoreServer(manager, messagePort, opts) { - // Per-project subchannels are keyed by an *instance id* — a string that is - // unique to one open lifetime of one project. Every time a project is - // opened (or re-opened after close), a new instance id is minted and - // returned to the client by `assertProjectExists`. The client uses it as - // the SubChannel identifier for that project's per-project messages. - // - // This means stale post-close calls from a client wrapper that captured - // the old instance id cannot collide with a freshly-opened project: they - // arrive on a different SubChannel id and route to the closed-instance - // tombstone branch in `handleMessage` instead of the new project's server. - - /** @type {Map void }>} */ - const existingInstanceServers = new Map() - - /** @type {Map} */ - const existingInstanceChannels = new Map() + /** @type {Map} */ + const projectHosts = new Map() /** - * projectId → in-flight or resolved promise for the current open instance - * id. Storing the promise (rather than the resolved string) dedupes - * concurrent `assertProjectExists` calls for the same project so they all - * resolve to the same instance id, instead of racing into separate - * `manager.getProject` calls that mint duplicate SubChannels. - * @type {Map>} + * Channel ids we've already logged an error for. Reaching the drop branch + * is a "shouldn't happen" case — a prefixed id that matches no reserved + * channel and no project route; we log once per id so a repeated stray + * message can't flood logs while a genuine routing bug stays visible. + * @type {Set} */ - const currentInstanceForProject = new Map() + const droppedChannelIds = new Set() /** - * Tombstone of instance ids that have been closed. The id string itself - * is cheap (~30 bytes); however the *first* stale message that arrives - * on a tombstoned id materialises a stub SubChannel + rpc-reflector - * server (in `existingInstanceChannels` / `existingInstanceServers`) - * that lives until the top-level server close. So a project that's - * closed but never receives a stale call costs ~30 bytes; one that does - * costs the size of a SubChannel + stub server. Bounded by the number - * of distinct closed instance ids that ever receive a stale message. - * @type {Set} + * @param {string} projectPublicId + * @returns {ProjectHost} */ - const closedInstanceIds = new Set() + function getOrCreateHost(projectPublicId) { + let host = projectHosts.get(projectPublicId) + if (!host) { + host = createProjectHost({ manager, messagePort, projectPublicId, opts }) + projectHosts.set(projectPublicId, host) + } + return host + } /** - * Instance ids we've already logged an error for. Reaching the drop branch - * is a "shouldn't happen" case — a prefixed id we minted but lost track of - * (foreign traffic is dropped earlier, see `handleMessage`); we log once - * per id so a repeated stray message can't flood logs while a genuine - * routing bug stays visible. - * @type {Set} + * Close the stale instance core leaves cached after `leaveProject` (core + * opens the project to leave it, guts it, and keeps it in its cache; only + * `addProject` on re-invite cleans it up). The instance's `close` event + * also detaches the project host's handler, so the next call hits the + * left-project guard instead of the gutted instance. + * + * Interim, paired with the left-project guard in {@link createProjectHost}: + * both go away together once core ships a typed PROJECT_LEFT error + * (digidem/comapeo-core#1313). + * + * @param {string} projectPublicId */ - const droppedInstanceIds = new Set() - - let instanceCounter = 0 - - const projectRoutingApi = new ProjectRoutingApi({ - getProjectInstance(projectId) { - const existing = currentInstanceForProject.get(projectId) - if (existing) return existing - - const promise = openProjectInstance(projectId) - currentInstanceForProject.set(projectId, promise) - // If the open fails, evict so a subsequent retry can attempt again - // instead of getting back the same rejected promise. (The close - // listener handles eviction on the success path.) - promise.catch(() => { - if (currentInstanceForProject.get(projectId) === promise) { - currentInstanceForProject.delete(projectId) - } - }) - return promise - }, - }) + async function closeLeftProjectInstance(projectPublicId) { + try { + const project = await manager.getProject(projectPublicId) + await project.close() + } catch { + // Never opened, or already gone — nothing to close. + } + } /** - * @param {string} projectId - * @returns {Promise} + * Wrap the consumer's request hook (if any) so `leaveProject` completions + * trigger the stale-instance cleanup above, without touching the manager + * object itself (binding or proxying the manager breaks its private-field + * methods). + * + * @type {NonNullable[2]>['onRequestHook']} */ - async function openProjectInstance(projectId) { - // Throws if the project doesn't exist; the rejection propagates back - // to the client through rpc-reflector's standard error response. - const project = await manager.getProject(projectId) - - const instanceId = `${PROJECT_INSTANCE_PREFIX}${projectId}:${++instanceCounter}` - const projectChannel = new SubChannel(messagePort, instanceId) - existingInstanceChannels.set(instanceId, projectChannel) - - project.once('close', () => { - closedInstanceIds.add(instanceId) - currentInstanceForProject.delete(projectId) - existingInstanceServers.get(instanceId)?.close() - existingInstanceServers.delete(instanceId) - projectChannel.close() - existingInstanceChannels.delete(instanceId) - }) - - const { close } = createServer(project, projectChannel, opts) - existingInstanceServers.set(instanceId, { close }) - - projectChannel.start() - - return instanceId + const managerRequestHook = (request, next) => { + /** @type {typeof next} */ + const instrumentedNext = (req) => { + const result = next(req) + if (req.method.length === 1 && req.method[0] === 'leaveProject') { + const projectPublicId = req.args[0] + if (typeof projectPublicId === 'string') { + Promise.resolve(result).then( + () => closeLeftProjectInstance(projectPublicId), + // Leave can fail after opening (and possibly gutting) the + // instance; closing is safe either way — a healthy project + // re-opens on the next call. + () => closeLeftProjectInstance(projectPublicId), + ) + } + } + return result + } + const consumerHook = opts?.onRequestHook + if (consumerHook) { + consumerHook(request, instrumentedNext) + } else { + instrumentedNext(request) + } } + const projectRoutingApi = new ProjectRoutingApi({ + ensureProject: (projectPublicId) => + getOrCreateHost(projectPublicId).ensureHandler(), + }) + const managerChannel = new SubChannel(messagePort, MANAGER_CHANNEL_ID) const projectRoutingChannel = new SubChannel(messagePort, PROJECT_ROUTING_ID) - const managerServer = createServer(manager, managerChannel, opts) + const managerServer = createServer(manager, managerChannel, { + ...opts, + onRequestHook: managerRequestHook, + }) const projectRoutingServer = createServer( projectRoutingApi, projectRoutingChannel, @@ -138,19 +146,12 @@ export function createComapeoCoreServer(manager, messagePort, opts) { close() { messagePort.removeEventListener('message', handleMessage) - for (const [id, server] of existingInstanceServers.entries()) { - server.close() - const channel = existingInstanceChannels.get(id) - if (channel) { - channel.close() - existingInstanceChannels.delete(id) - } - existingInstanceServers.delete(id) + for (const host of projectHosts.values()) { + host.close() } + projectHosts.clear() + droppedChannelIds.clear() - currentInstanceForProject.clear() - closedInstanceIds.clear() - droppedInstanceIds.clear() managerServer.close() managerChannel.close() projectRoutingServer.close() @@ -161,7 +162,7 @@ export function createComapeoCoreServer(manager, messagePort, opts) { /** * @param {{ data: unknown }} payload */ - async function handleMessage({ data }) { + function handleMessage({ data }) { if (!isRelevantEventData(data)) return const { id } = data @@ -170,8 +171,7 @@ export function createComapeoCoreServer(manager, messagePort, opts) { // drop it silently (no warning) so unrelated traffic can't flood logs. if (!id.startsWith(COMAPEO_PREFIX)) return - // Reserved channels and currently-open project instances are routed by - // their own SubChannel listeners; nothing to do here. + // Reserved channels are routed by their own SubChannel listeners. if ( id === MANAGER_CHANNEL_ID || id === PROJECT_ROUTING_ID || @@ -180,39 +180,23 @@ export function createComapeoCoreServer(manager, messagePort, opts) { return } - if (existingInstanceChannels.has(id)) return - - if (closedInstanceIds.has(id)) { - // Stale message for a closed project instance. Build a stub - // rpc-server bound to a Proxy that throws "Project is closed" for - // any apply. The error rides the standard serializeError → RESPONSE - // path. The stub holds no reference to the (already-released) - // project. The stub stays alive on `existingInstanceChannels`/ - // `existingInstanceServers` for the rest of the session, so further - // stale messages on this instance id route through the SubChannel's - // own listener directly without re-entering this branch. - const stubChannel = new SubChannel(messagePort, id) - existingInstanceChannels.set(id, stubChannel) - const stubHandler = createClosedProjectStub() - const { close: closeStubServer } = createServer( - stubHandler, - stubChannel, - opts, - ) - existingInstanceServers.set(id, { close: closeStubServer }) - stubChannel.start() - stubChannel.dispatchEvent({ data: data.message }) - return + if (id.startsWith(PROJECT_CHANNEL_PREFIX)) { + const projectPublicId = id.slice(PROJECT_CHANNEL_PREFIX.length) + if (projectPublicId.length > 0) { + if (projectHosts.has(projectPublicId)) return + // First traffic for this project: the host's channel listener was + // not yet registered when this event was dispatched at the port + // level, so hand the frame to the channel directly. + const host = getOrCreateHost(projectPublicId) + host.channel.dispatchEvent({ data: data.message }) + return + } } - // Carries our prefix but matches no known channel. With the - // manager/project-routing/services channels and every open and closed - // project instance accounted for above, reaching here means we minted - // this id and lost track of it (or a paired client desynced) — a genuine - // routing bug, not foreign traffic. Logged once per id (see - // `droppedInstanceIds`). - if (!droppedInstanceIds.has(id)) { - droppedInstanceIds.add(id) + // Carries our prefix but matches no known channel shape — we lost track + // of an id we minted, or a paired client desynced. Logged once per id. + if (!droppedChannelIds.has(id)) { + droppedChannelIds.add(id) console.error( `comapeo-ipc: dropping message for unrecognised channel id "${id}"`, ) @@ -221,52 +205,135 @@ export function createComapeoCoreServer(manager, messagePort, opts) { } /** - * Build a Proxy bound as a stub rpc-reflector handler for a closed project - * instance: it answers property/`has` checks at any depth and throws - * `ProjectClosedError` when a method is applied (rpc-reflector catches and - * serializes it back to the client). The outer target is a plain object so the - * proxy passes rpc-reflector's `typeof handler === 'object'` invariant; nested - * accesses return a function-target proxy so `applyNestedMethod` finds - * `typeof === 'function'` and triggers the apply trap. + * @typedef {object} ProjectHost + * @property {SubChannel} channel + * @property {() => Promise} ensureHandler + * @property {() => void} close */ -function createClosedProjectStub() { - /** @type {ProxyHandler} */ - const handler = { - get() { - return new Proxy(function () {}, handler) - }, - has() { - return true - }, - apply() { - throw new ProjectClosedError() + +/** + * One per project public id, for the lifetime of the top-level server. Owns + * the project's stable SubChannel and one long-lived rpc-reflector server + * bound to it via a late-bound handler factory. rpc-reflector owns the + * instance lifecycle mechanics: it invokes the factory (single-flight) when + * the first call or subscription needing a handler arrives, keeps the + * client's event subscriptions in a registry that survives the instance, and + * re-attaches them to each fresh instance before any awaited frame is + * dispatched — so client listeners survive server-side close/re-open cycles + * they never hear about. A factory rejection (`NotFoundError`, + * `ProjectLeftError`) is answered per request with the error, and is not + * cached, so a later call retries. + * + * This host adds only what is comapeo-specific: the left-project guard, the + * close-in-flight retry around `manager.getProject`, and detaching the + * handler when the instance closes so the next call re-opens. + * + * @param {object} options + * @param {MapeoManager} options.manager + * @param {MessagePortLike} options.messagePort + * @param {string} options.projectPublicId + * @param {Parameters[2]} [options.opts] + * @returns {ProjectHost} + */ +function createProjectHost({ manager, messagePort, projectPublicId, opts }) { + const channel = new SubChannel( + messagePort, + `${PROJECT_CHANNEL_PREFIX}${projectPublicId}`, + ) + + /** + * Interim left-project guard, paired with the `leaveProject` request hook + * in `createComapeoCoreServer`; both are removed together once core ships + * a typed PROJECT_LEFT error (digidem/comapeo-core#1313). Left projects + * re-open as live-but-gutted instances (core deliberately allows this so + * an interrupted leave can finish), so leftness must be checked via + * `listProjects`, not inferred from `getProject`. + */ + async function assertNotLeft() { + const projects = await manager.listProjects({ includeLeft: true }) + const entry = projects.find((p) => p.projectId === projectPublicId) + if (entry && entry.status === 'left') { + throw new ProjectLeftError() + } + } + + /** @returns {Promise} */ + async function openProject() { + await assertNotLeft() + const project = await getOpenableProject() + // Re-checked after the open resolves: a leave can land while the open + // was in flight (leave waits for sync), and binding then would serve + // calls from the gutted instance instead of rejecting them. + await assertNotLeft() + // `once`, not `on`: core's MapeoProject emits `close` twice (once from + // `_close`, once from ready-resource). + project.once('close', () => server.detachHandler()) + return project + } + + /** + * `manager.getProject` returns the dying instance for the whole duration + * of an in-flight close (its cache evicts only on the `close` event), so + * wait out a close-in-progress and retry rather than binding to a corpse. + * + * @returns {Promise} + */ + async function getOpenableProject() { + for (let attempt = 0; attempt < 5; attempt++) { + const project = await manager.getProject(projectPublicId) + if (project.closed) continue + if (project.closing) { + await Promise.resolve(project.closing).catch(noop) + continue + } + return project + } + throw new Error(`Project ${projectPublicId} kept closing while opening`) + } + + // Created before `start()`, so the first frame — which the top-level + // router hands to the channel directly — reaches the rpc server. + const server = createServer(openProject, channel, opts) + channel.start() + + return { + channel, + /** + * Bind a live instance now (opening it if necessary), so an eager open + * attaches subscriptions before any project-channel frame. Rejects with + * `ProjectLeftError` for left projects or whatever `manager.getProject` + * throws (e.g. `NotFoundError`). Safe to call concurrently. + */ + ensureHandler: () => server.ensureHandler(), + close() { + server.close() + channel.close() }, } - return new Proxy({}, handler) } export class ProjectRoutingApi { - #getProjectInstance + #ensureProject /** - * @param {{ getProjectInstance: (projectId: string) => Promise }} opts + * @param {{ ensureProject: (projectPublicId: string) => Promise }} opts */ - constructor({ getProjectInstance }) { - this.#getProjectInstance = getProjectInstance + constructor({ ensureProject }) { + this.#ensureProject = ensureProject } /** - * Verify the project exists, opening it (or re-opening it after close) - * if necessary, and return the per-instance subchannel id the client - * should use for per-project messages. The returned id is unique to the - * current open lifetime of the project — closing and re-opening yields - * a different id. + * Verify the project exists and is usable, opening it (or re-opening it + * after a server-side close) if necessary. Rejects with `NotFoundError` + * for unknown projects and `ProjectLeftError` for projects this device has + * left. * - * @param {string} projectId - * @returns {Promise} instance id + * @param {string} projectPublicId + * @returns {Promise} */ - async assertProjectExists(projectId) { - return this.#getProjectInstance(projectId) + async assertProjectExists(projectPublicId) { + await this.#ensureProject(projectPublicId) + return true } } diff --git a/tests/basic.js b/tests/basic.js index 234fd88..215a068 100644 --- a/tests/basic.js +++ b/tests/basic.js @@ -175,6 +175,43 @@ test('Client calls fail after server closes', async (t) => { } }) +test('A getProject in flight when the client closes rejects with ClientClosedError', async (t) => { + const { client } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + + // In flight when the close starts: its validation round trip resolves + // during the close's await window, after the project-client sweep can no + // longer include a freshly minted wrapper. It must reject rather than + // create a wrapper (and its port listener) that nothing will clean up. + const inFlight = client.getProject(projectId) + const closing = closeComapeoCoreClient(client) + + await assert.rejects(() => inFlight, { code: ClientClosedError.code }) + await closing +}) + +test('Awaiting a nested namespace after close resolves instead of hanging', async (t) => { + const { client } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + await closeComapeoCoreClient(client) + + // `await` probes `then`; a truthy `then` on the closed proxy would make + // this a thenable whose callbacks never fire — a permanent hang. + const namespace = await project.observation + assert.ok(namespace) + await assert.rejects( + () => + namespace.create({ + schemaName: 'observation', + attachments: [], + tags: {}, + }), + { code: ClientClosedError.code }, + ) +}) + test('In-flight calls reject with RpcChannelClosedError when the client closes', async (t) => { const { client } = setup(t) diff --git a/tests/events.js b/tests/events.js index 169eadc..651615d 100644 --- a/tests/events.js +++ b/tests/events.js @@ -2,10 +2,25 @@ import test from 'node:test' import assert from 'node:assert/strict' import pDefer from 'p-defer' -import { ClientClosedError, ProjectClosedError } from '../src/errors.js' +import { ClientClosedError } from '../src/errors.js' import { closeComapeoCoreClient } from '../src/client.js' import { setup } from './helpers.js' + +/** + * Poll until `predicate` holds, for assertions about work the server does on + * its own initiative (with no call to await). + * + * @param {() => boolean} predicate + * @param {string} description + */ +async function waitFor(predicate, description) { + for (let i = 0; i < 200; i++) { + if (predicate()) return + await new Promise((resolve) => setTimeout(resolve, 5)) + } + assert.fail(`Timed out waiting for ${description}`) +} import { FakeManager } from './fake-manager.js' test('Server events are forwarded to client listeners', async (t) => { @@ -52,29 +67,193 @@ test('Client listeners stop receiving events after removeListener', async (t) => assert.equal(count, 1, 'no further events after removeListener') }) -test('EventEmitter methods throw synchronously after the project is closed', async (t) => { - const { client } = setup(t) +test('Project events are forwarded to client listeners', async (t) => { + const { client, serverManager } = setup(t) const projectId = await client.createProject({ name: 'mapeo' }) const project = await client.getProject(projectId) - await project.close() + /** @type {import('p-defer').DeferredPromise} */ + const deferred = pDefer() + project.on('some-event', (value) => deferred.resolve(value)) + await project.$getProjectSettings() + + const serverProject = await serverManager.getProject(projectId) + serverProject.emit('some-event', 'hello') - // Emitter methods are not awaited by callers, so a rejected promise would - // surface as an unhandled rejection — they throw at the call site instead. - assert.throws(() => project.on('some-event', () => {}), { - code: ProjectClosedError.code, - }) - assert.throws(() => project.removeListener('some-event', () => {}), { - code: ProjectClosedError.code, - }) + assert.equal(await deferred.promise, 'hello') +}) + +// The load-bearing test for server-owned lifecycle: the server-side +// subscriptions must outlive the instance they were attached to and be +// replayed into the fresh instance when it re-opens — client listeners +// survive a close/re-open they never hear about. +test('Project event subscriptions survive a server-side close and re-open', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + /** @type {unknown[]} */ + const received = [] + project.on('some-event', (value) => received.push(value)) + await project.$getProjectSettings() + + const firstInstance = await serverManager.getProject(projectId) + firstInstance.emit('some-event', 'before-close') + await project.$getProjectSettings() + assert.deepEqual(received, ['before-close']) + + await firstInstance.close() + + // Re-open via any call; the subscriptions must be re-attached before the + // buffered call is dispatched. + await project.$getProjectSettings() + + const secondInstance = await serverManager.getProject(projectId) + assert.notEqual(secondInstance, firstInstance) + secondInstance.emit('some-event', 'after-reopen') + await project.$getProjectSettings() + + assert.deepEqual(received, ['before-close', 'after-reopen']) }) -test('EventEmitter methods throw synchronously after the client is closed', async (t) => { +// Subscriptions are held per prop path, so a nested namespace that is itself +// an EventEmitter (core's `project.$sync`) must be re-attached to the matching +// namespace of the fresh instance — not just the project root. +test('Nested-namespace subscriptions survive a server-side close and re-open', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + /** @type {unknown[]} */ + const received = [] + project.$sync.on('sync-state', (value) => received.push(value)) + await project.$sync.getState() + + const firstInstance = await serverManager.getProject(projectId) + firstInstance.$sync.emit('sync-state', 'before-close') + await project.$sync.getState() + assert.deepEqual(received, ['before-close']) + + await firstInstance.close() + await project.$sync.getState() + + const secondInstance = await serverManager.getProject(projectId) + assert.notEqual(secondInstance, firstInstance) + secondInstance.$sync.emit('sync-state', 'after-reopen') + await project.$sync.getState() + + assert.deepEqual(received, ['before-close', 'after-reopen']) + + // Root and nested subscriptions are independent: the root emitter must not + // have picked up the nested namespace's listener. + secondInstance.emit('sync-state', 'from-root') + await project.$sync.getState() + assert.deepEqual(received, ['before-close', 'after-reopen']) +}) + +test('Unsubscribed project events are not replayed on re-open', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + let count = 0 + const listener = () => { + count++ + } + project.on('some-event', listener) + await project.$getProjectSettings() + project.removeListener('some-event', listener) + await project.$getProjectSettings() + + const firstInstance = await serverManager.getProject(projectId) + await firstInstance.close() + + await project.$getProjectSettings() + const secondInstance = await serverManager.getProject(projectId) + secondInstance.emit('some-event') + await project.$getProjectSettings() + + assert.equal(count, 0, 'unsubscribed event must not be re-subscribed') +}) + +test('A subscription made while the project is closed server-side still takes effect', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + await project.$getProjectSettings() + + const firstInstance = await serverManager.getProject(projectId) + await firstInstance.close() + + // Subscribing while no instance is live: the ON frame itself must wake the + // project up (subscribing expresses interest) and land on the fresh + // instance. + /** @type {import('p-defer').DeferredPromise} */ + const deferred = pDefer() + project.on('some-event', (value) => deferred.resolve(value)) + await project.$getProjectSettings() + + const secondInstance = await serverManager.getProject(projectId) + secondInstance.emit('some-event', 'woken') + assert.equal(await deferred.promise, 'woken') +}) + +// Subscribing is itself a reason to open a project: a consumer that only +// listens (a component mounted on `$sync` events, or #89's post-restart +// resubscribe) would otherwise never receive anything. Every other test here +// makes a method call after subscribing, which masks this. +test('Subscribing alone re-opens a dormant project, with no further calls', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + await project.$getProjectSettings() + + const firstInstance = await serverManager.getProject(projectId) + const opensBefore = serverManager.getProjectCallCount.get(projectId) ?? 0 + await firstInstance.close() + + /** @type {import('p-defer').DeferredPromise} */ + const deferred = pDefer() + project.on('some-event', (value) => deferred.resolve(value)) + + await waitFor( + () => (serverManager.getProjectCallCount.get(projectId) ?? 0) > opensBefore, + 'the server to re-open the project for the subscription alone', + ) + + const secondInstance = await serverManager.getProject(projectId) + assert.notEqual(secondInstance, firstInstance) + secondInstance.emit('some-event', 'woken') + assert.equal(await deferred.promise, 'woken') +}) + +test('EventEmitter subscribe methods throw synchronously after the client is closed', async (t) => { const { client } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) await closeComapeoCoreClient(client) assert.throws(() => client.on('local-peers', () => {}), { code: ClientClosedError.code, }) + assert.throws(() => project.on('some-event', () => {}), { + code: ClientClosedError.code, + }) +}) + +test('EventEmitter unsubscribe methods are no-ops after the client is closed', async (t) => { + const { client } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + const listener = () => {} + project.on('some-event', listener) + + await closeComapeoCoreClient(client) + + // Removing a listener from a dead client is correct teardown (React effect + // cleanup runs against stale references) — it must not throw. + assert.doesNotThrow(() => project.removeListener('some-event', listener)) + assert.doesNotThrow(() => project.off('some-event', listener)) + assert.doesNotThrow(() => client.removeAllListeners('local-peers')) }) diff --git a/tests/fake-manager.js b/tests/fake-manager.js index e09f899..c39a61b 100644 --- a/tests/fake-manager.js +++ b/tests/fake-manager.js @@ -20,10 +20,35 @@ import { NotFoundError } from '@comapeo/core/errors.js' * @property {number} obsCounter */ +/** + * A nested namespace that is itself an EventEmitter, mirroring core's + * `project.$sync`: the IPC server must resolve method calls and event + * subscriptions at nested paths, not just at the project root. + */ +class FakeSync extends EventEmitter { + async getState() { + return { state: 'idle' } + } +} + class FakeProject extends EventEmitter { /** @type {ProjectStore} */ #store #closed = false + /** @type {Promise | null} */ + #closeHold = null + + $sync = new FakeSync() + + // Mirror ready-resource's surface, which the IPC server reads to avoid + // binding to an instance whose close is in flight: `closing` is the close + // promise from the moment a close starts, `closed` flips once it is done. + /** @type {Promise | null} */ + closing = null + + get closed() { + return this.#closed + } /** @param {ProjectStore} store */ constructor(store) { @@ -55,9 +80,29 @@ class FakeProject extends EventEmitter { return { ...this.#store.settings } } - async close() { - if (this.#closed) return + /** + * Test knob: make an in-flight close wait for `promise` before completing, + * to hold open the window where `closing` is set but `closed` is not. + * + * @param {Promise} promise + */ + holdClose(promise) { + this.#closeHold = promise + } + + /** @returns {Promise} */ + close() { + if (this.closing) return this.closing + this.closing = this.#doClose() + return this.closing + } + + async #doClose() { + await this.#closeHold this.#closed = true + // Mirror core: `close` is emitted TWICE — once manually from + // MapeoProject's `_close` and once by ready-resource itself. + this.emit('close') this.emit('close') } } @@ -67,6 +112,8 @@ export class FakeManager extends EventEmitter { #stores = new Map() /** @type {Map} */ #liveProjects = new Map() + /** @type {Set} */ + #leftProjects = new Set() #projectCounter = 0 /** @@ -136,11 +183,51 @@ export class FakeManager extends EventEmitter { return project } - async listProjects() { - return [...this.#stores.entries()].map(([projectId, store]) => ({ - projectId, - name: store.settings.name, - })) + /** + * Mirror core: left projects are omitted unless `includeLeft`, in which + * case they appear with `status: 'left'`. + * + * @param {{ includeLeft?: boolean }} [opts] + */ + async listProjects({ includeLeft = false } = {}) { + return [...this.#stores.entries()] + .filter( + ([projectId]) => includeLeft || !this.#leftProjects.has(projectId), + ) + .map(([projectId, store]) => ({ + projectId, + name: store.settings.name, + status: this.#leftProjects.has(projectId) ? 'left' : 'joined', + })) + } + + /** + * Mirror core's `leaveProject`: marks the project left and clears its data, + * but does NOT close the live instance — core leaves the gutted instance + * cached (only `addProject` cleans it up), so the IPC server must close it. + * + * @param {string} projectId + */ + async leaveProject(projectId) { + const store = this.#stores.get(projectId) + if (!store) throw new NotFoundError(`Project ${projectId} does not exist`) + // Mirror core: leaving opens the project if it wasn't open. + await this.getProject(projectId) + this.#leftProjects.add(projectId) + store.observations.clear() + } + + /** + * Simplified stand-in for core's `addProject` on re-invite: closes any + * stale cached instance (as core does) and clears the left flag. + * + * @param {string} projectId + */ + async addProject(projectId) { + const store = this.#stores.get(projectId) + if (!store) throw new NotFoundError(`Project ${projectId} does not exist`) + await this.#liveProjects.get(projectId)?.close() + this.#leftProjects.delete(projectId) } async getIsArchiveDevice() { diff --git a/tests/integration.js b/tests/integration.js index d89dbe8..4a8eda2 100644 --- a/tests/integration.js +++ b/tests/integration.js @@ -13,7 +13,6 @@ import { closeComapeoCoreClient, } from '../src/client.js' import { createComapeoCoreServer } from '../src/server.js' -import { ProjectClosedError } from '../src/errors.js' const require = createRequire(import.meta.url) @@ -60,8 +59,6 @@ test('end-to-end against a real MapeoManager', async (t) => { }) const readBack = await project.observation.getByDocId(obs.docId) assert.equal(readBack.docId, obs.docId) - - await project.close() }) /** @@ -69,23 +66,30 @@ test('end-to-end against a real MapeoManager', async (t) => { */ test('handle manager initiating the close', async (t) => { const { server, client, manager } = setup(t) - // Manager methods round-trip. const projectId = await client.createProject({ name: 'mapeo' }) assert.ok(projectId) const clientProject = await client.getProject(projectId) + const obs = await clientProject.observation.create({ + schemaName: 'observation', + attachments: [], + tags: {}, + }) const rawProject = await manager.getProject(projectId) - // This simulates the project being closed through other means like leaveProject + // A server-side close the client never hears about (resource policy, or + // core's addProject closing a stale instance on re-invite). The client's + // reference must keep working against the transparently re-opened + // instance — with the real MapeoProject this proves a fresh instance can + // serve the same channel and reads round-trip. await rawProject.close() - await assert.rejects(() => clientProject.$getProjectSettings(), { - code: ProjectClosedError.code, - }) + const readBack = await clientProject.observation.getByDocId(obs.docId) + assert.equal(readBack.docId, obs.docId) const reOpened = await client.getProject(projectId) - + assert.equal(reOpened, clientProject, 'project references are permanent') await reOpened.$getProjectSettings() }) diff --git a/tests/project-close.js b/tests/project-close.js deleted file mode 100644 index 38df13c..0000000 --- a/tests/project-close.js +++ /dev/null @@ -1,311 +0,0 @@ -import test from 'node:test' -import assert from 'node:assert/strict' -import { NotFoundError } from '@comapeo/core/errors.js' - -import { setup } from './helpers.js' -import { ProjectClosedError } from '../src/errors.js' - -test('After close, methods on the closed reference reject', async (t) => { - const { client } = setup(t) - const projectId = await client.createProject({ name: 'mapeo' }) - const project = await client.getProject(projectId) - - // Sanity: methods work pre-close. - await project.$getProjectSettings() - - await project.close() - - await assert.rejects(() => project.$getProjectSettings(), { - code: ProjectClosedError.code, - }) -}) - -test('close() is idempotent — repeated calls resolve like the first', async (t) => { - const { client } = setup(t) - const projectId = await client.createProject({ name: 'mapeo' }) - const project = await client.getProject(projectId) - - await Promise.all([project.close(), project.close()]) - await project.close() - - // And the project can still be re-opened afterwards. - const reopened = await client.getProject(projectId) - await reopened.$getProjectSettings() -}) - -test('After close, nested-namespace methods on the closed reference reject', async (t) => { - const { client } = setup(t) - const projectId = await client.createProject({ name: 'mapeo' }) - const project = await client.getProject(projectId) - - // Pre-close: nested namespace works. - await project.observation.create({ - schemaName: 'observation', - attachments: [], - tags: {}, - }) - - await project.close() - - await assert.rejects( - () => - project.observation.create({ - schemaName: 'observation', - attachments: [], - tags: {}, - }), - { code: ProjectClosedError.code }, - ) -}) - -test('After close, observations created earlier are still readable via a re-opened reference', async (t) => { - const { client } = setup(t) - const projectId = await client.createProject({ name: 'mapeo' }) - const project = await client.getProject(projectId) - - const obs = await project.observation.create({ - schemaName: 'observation', - attachments: [], - tags: {}, - }) - - await project.close() - - const reopened = await client.getProject(projectId) - assert.notEqual( - reopened, - project, - 're-opened reference should be a fresh wrapper', - ) - - const fetched = await reopened.observation.getByDocId(obs.docId) - assert.equal(fetched.docId, obs.docId) -}) - -// The failure mode the bug report flagged: after close + re-open, a stale -// call through the OLD reference must NOT silently land on the freshly -// re-opened project — it must reject. The local teardown rejects it before -// it reaches the wire; per-instance subchannel ids guarantee that even a -// message that does reach the wire (e.g. posted while close is in flight) -// cannot route to the new instance. -test('After close + re-open, a stale call on the old reference still rejects', async (t) => { - const { client } = setup(t) - const projectId = await client.createProject({ name: 'mapeo' }) - - const oldProject = await client.getProject(projectId) - await oldProject.$getProjectSettings() - await oldProject.close() - - const newProject = await client.getProject(projectId) - await newProject.$getProjectSettings() - - await assert.rejects(() => oldProject.$getProjectSettings(), { - code: ProjectClosedError.code, - }) -}) - -test('Two parallel getProject(id) calls return one wrapper and both work', async (t) => { - const { client } = setup(t) - const projectId = await client.createProject({ name: 'mapeo' }) - - // Drive parallel `getProject(id)` calls from a freshly-cleared cache by - // closing the project first, so both calls go all the way to the server. - await (await client.getProject(projectId)).close() - - const [a, b] = await Promise.all([ - client.getProject(projectId), - client.getProject(projectId), - ]) - - assert.equal(a, b, 'both callers should resolve to the same wrapper') - await a.$getProjectSettings() - await b.$getProjectSettings() -}) - -test('Closing one project does not affect another open project', async (t) => { - const { client } = setup(t) - const projectIdA = await client.createProject({ name: 'mapeo-a' }) - const projectIdB = await client.createProject({ name: 'mapeo-b' }) - - const projectA = await client.getProject(projectIdA) - const projectB = await client.getProject(projectIdB) - - await projectA.$getProjectSettings() - await projectB.$getProjectSettings() - - await projectA.close() - - // A is closed; B is unaffected. - await assert.rejects(() => projectA.$getProjectSettings(), { - code: ProjectClosedError.code, - }) - const settingsB = await projectB.$getProjectSettings() - assert.equal(settingsB.name, 'mapeo-b') -}) - -test('When the server closes the project, client calls on the wrapper reject', async (t) => { - const { client, serverManager } = setup(t) - const projectId = await client.createProject({ name: 'mapeo' }) - const project = await client.getProject(projectId) - await project.$getProjectSettings() - - // Close the project from the server side, bypassing the client. The - // wrapper has no idea this happened until it tries a method call. - const serverProject = await serverManager.getProject(projectId) - await serverProject.close() - - await assert.rejects(() => project.$getProjectSettings(), { - code: ProjectClosedError.code, - }) - - // Even after the project is re-opened on the server, calls from the old - // wrapper still reject: they route by the old (tombstoned) instance id, - // so they cannot reach the fresh instance. - const reopenedServerProject = await serverManager.getProject(projectId) - await assert.rejects(() => project.$getProjectSettings(), { - code: ProjectClosedError.code, - }) - await reopenedServerProject.close() -}) - -// The next two tests pin recovery from a manager-initiated close — what -// `MapeoManager.addProject` does to a previously-left project when a -// re-invite is accepted (digidem/comapeo-mobile#2042). The cached client -// wrapper for the closed instance must not keep being handed out by -// `getProject` once the project can be re-opened. - -test('After a manager-initiated close is observed, getProject returns a fresh working instance', async (t) => { - const { client, serverManager } = setup(t) - const projectId = await client.createProject({ name: 'mapeo' }) - const project = await client.getProject(projectId) - - // Deliberately not node:events `once()`: its cleanup calls - // `removeListener` on the wrapper after the close event, which a closed - // wrapper may reject synchronously. - const closeObserved = new Promise((resolve) => { - project.once('close', resolve) - }) - // Round-trip so the 'close' subscription is registered server-side - // (FIFO channel) before the close below emits. - await project.$getProjectSettings() - - const serverProject = await serverManager.getProject(projectId) - await serverProject.close() - await closeObserved - - const reOpened = await client.getProject(projectId) - const settings = await reOpened.$getProjectSettings() - assert.equal(settings.name, 'mapeo') - assert.notEqual( - reOpened, - project, - 'getProject after an observed close should return a fresh wrapper', - ) -}) - -test('getProject returns a working instance immediately after a manager-initiated close', async (t) => { - const { client, serverManager } = setup(t) - const projectId = await client.createProject({ name: 'mapeo' }) - await client.getProject(projectId) - - const serverProject = await serverManager.getProject(projectId) - await serverProject.close() - - // The close notification has not reached the client yet, so its cached - // wrapper still looks open. The server already knows the project is - // closed (it updated its routing state during close() above), so to - // return a working instance here getProject must ask the server rather - // than trust its cache. - const reOpened = await client.getProject(projectId) - const settings = await reOpened.$getProjectSettings() - assert.equal(settings.name, 'mapeo') -}) - -test('A method call posted before close completes still resolves', async (t) => { - const { client } = setup(t) - const projectId = await client.createProject({ name: 'mapeo' }) - const project = await client.getProject(projectId) - - // Fire the method without awaiting it, then close. The method's request - // is posted to the server before the close request, so it should be - // processed against the still-open project and resolve normally. - const inFlight = project.$getProjectSettings() - await project.close() - - const settings = await inFlight - assert.equal(settings.name, 'mapeo') -}) - -// Repeatedly opening and closing the same project must not accumulate -// closed instances in IPC's per-project bookkeeping. An earlier draft of -// this PR risked exactly that — every close would have left a stub -// rpc-server (or a wrapper Proxy) bound to the closed project, growing -// linearly with the number of cycles. -// -// The fake manager releases its project instance on close (it holds no -// reference to a closed project), so any instance still reachable after a -// cycle is being retained by the IPC layer itself. We can therefore assert -// the strong invariant: after N cycles, zero closed instances survive. -// -// Runs only when `global.gc` is available (npm test passes --expose-gc). -test('Repeatedly opening and closing the same project does not retain prior instances', async (t) => { - if (typeof global.gc !== 'function') { - t.skip('Run with --expose-gc to verify cycle retention') - return - } - - const { client, serverManager } = setup(t) - const projectId = await client.createProject({ name: 'mapeo' }) - - const N = 5 - /** @type {Array>} */ - const refs = [] - for (let i = 0; i < N; i++) { - refs.push(await cycleAndCaptureWeakRef()) - } - - /** @returns {Promise>} */ - async function cycleAndCaptureWeakRef() { - const project = await client.getProject(projectId) - await project.$getProjectSettings() - const serverProject = await serverManager.getProject(projectId) - const ref = new WeakRef(serverProject) - await project.close() - return ref - } - - for (let i = 0; i < 20; i++) { - global.gc() - await new Promise((resolve) => setImmediate(resolve)) - } - - const survivors = refs.filter((ref) => ref.deref() !== undefined).length - assert.equal( - survivors, - 0, - `${survivors} of ${N} closed instances retained — expected 0. A non-zero count means IPC is accumulating closed instances cycle over cycle.`, - ) -}) - -test('After a failed getProject, a subsequent getProject for a real project succeeds', async (t) => { - const { client } = setup(t) - - // Different ids: first fails (project does not exist), then a real - // project is created and getProject(realId) must succeed. Without the - // cache-poisoning fix, a rejected entry could linger and break unrelated - // ids; with it, only the failed id's entry is evicted. - await assert.rejects(() => client.getProject('does-not-exist'), { - code: NotFoundError.code, - }) - - const realId = await client.createProject({ name: 'mapeo' }) - const project = await client.getProject(realId) - await project.$getProjectSettings() - - // And a retry of the original failing id still rejects (project still - // doesn't exist) — proving the failure path itself is also retried, not - // returned from a poisoned cache. - await assert.rejects(() => client.getProject('does-not-exist'), { - code: NotFoundError.code, - }) -}) diff --git a/tests/project-lifecycle.js b/tests/project-lifecycle.js new file mode 100644 index 0000000..fe22e1e --- /dev/null +++ b/tests/project-lifecycle.js @@ -0,0 +1,674 @@ +import test from 'node:test' +import assert from 'node:assert/strict' +import pDefer from 'p-defer' +import { NotFoundError } from '@comapeo/core/errors.js' + +import { setup } from './helpers.js' +import { FakeManager } from './fake-manager.js' +import { ProjectLeftError } from '../src/errors.js' +import { + createComapeoCoreClient, + closeComapeoCoreClient, +} from '../src/client.js' +import { createComapeoCoreServer } from '../src/server.js' + +/** + * Poll until `predicate` holds, for assertions about work the server does on + * its own initiative (with no call to await). + * + * @param {() => boolean} predicate + * @param {string} description + */ +async function waitFor(predicate, description) { + for (let i = 0; i < 200; i++) { + if (predicate()) return + await new Promise((resolve) => setTimeout(resolve, 5)) + } + assert.fail(`Timed out waiting for ${description}`) +} + +// Project instance lifecycle is owned by the server: project references are +// permanent, and the server transparently re-opens a project whose instance +// was closed server-side (resource policy, addProject on re-invite, server +// restart). The one deliberate exception is a left project, which rejects +// with ProjectLeftError until re-joined. These tests pin that contract. + +test('Calls work transparently after a server-side close', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + const obs = await project.observation.create({ + schemaName: 'observation', + attachments: [], + tags: {}, + }) + + // Close the project from the server side, bypassing the client. The + // client's reference must keep working: the next call re-opens the + // project on the same channel. + const serverProject = await serverManager.getProject(projectId) + await serverProject.close() + + const fetched = await project.observation.getByDocId(obs.docId) + assert.equal(fetched.docId, obs.docId) + + // The server really did cycle the instance: IPC open + this test's own + // serverManager.getProject above + IPC re-open. + assert.equal(serverManager.getProjectCallCount.get(projectId), 3) +}) + +test('Calls work immediately after a server-side close, with no event wait', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + await project.$getProjectSettings() + + const serverProject = await serverManager.getProject(projectId) + await serverProject.close() + + // No round-trip, no close notification consumed — the very next call must + // succeed against a fresh instance. + const settings = await project.$getProjectSettings() + assert.equal(settings.name, 'mapeo') +}) + +test('getProject after a server-side close returns the same working reference', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + const serverProject = await serverManager.getProject(projectId) + await serverProject.close() + + const again = await client.getProject(projectId) + assert.equal(again, project, 'project references are permanent') + await again.$getProjectSettings() +}) + +test('Nested-namespace calls also survive a server-side close', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + await project.observation.create({ + schemaName: 'observation', + attachments: [], + tags: {}, + }) + + const serverProject = await serverManager.getProject(projectId) + await serverProject.close() + + const obs = await project.observation.create({ + schemaName: 'observation', + attachments: [], + tags: {}, + }) + assert.ok(obs.docId) +}) + +test('A call in flight when the server closes the project still settles', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + await project.$getProjectSettings() + + // Fire without awaiting, then close server-side. Whatever the outcome + // (served by the dying instance or failed), the call must settle rather + // than hang. + const inFlight = project.$getProjectSettings() + const serverProject = await serverManager.getProject(projectId) + await serverProject.close() + + const result = await Promise.race([ + inFlight.then( + () => 'settled', + () => 'settled', + ), + new Promise((resolve) => setTimeout(() => resolve('hung'), 2000)), + ]) + assert.equal(result, 'settled') +}) + +test('Closing one project server-side does not affect another', async (t) => { + const { client, serverManager } = setup(t) + const projectIdA = await client.createProject({ name: 'mapeo-a' }) + const projectIdB = await client.createProject({ name: 'mapeo-b' }) + + const projectA = await client.getProject(projectIdA) + const projectB = await client.getProject(projectIdB) + await projectA.$getProjectSettings() + await projectB.$getProjectSettings() + + const serverProjectA = await serverManager.getProject(projectIdA) + await serverProjectA.close() + + const settingsB = await projectB.$getProjectSettings() + assert.equal(settingsB.name, 'mapeo-b') + assert.equal( + serverManager.getProjectCallCount.get(projectIdB), + 1, + 'project B was never cycled', + ) + + const settingsA = await projectA.$getProjectSettings() + assert.equal(settingsA.name, 'mapeo-a') +}) + +test('Two parallel getProject(id) calls return one wrapper and both work', async (t) => { + const { client } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + + const [a, b] = await Promise.all([ + client.getProject(projectId), + client.getProject(projectId), + ]) + + assert.equal(a, b, 'both callers should resolve to the same wrapper') + await a.$getProjectSettings() + await b.$getProjectSettings() +}) + +test('getProject for an already-acquired project makes no wire round trip', async (t) => { + const { client, server } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + // With the server gone, no round trip can be answered — only the cached + // wrapper from the first acquisition can resolve this. + server.close() + const again = await client.getProject(projectId) + assert.equal(again, project, 'cached wrapper returned without validation') +}) + +test('leaveProject: calls reject with ProjectLeftError and the gutted instance is closed', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + await project.$getProjectSettings() + + // Core's leaveProject opens the project, guts it, and leaves the corpse + // cached — observe the instance so we can assert the server closed it. + const liveInstance = await serverManager.getProject(projectId) + const closeObserved = pDefer() + liveInstance.once('close', () => closeObserved.resolve(undefined)) + + await client.leaveProject(projectId) + await closeObserved.promise + + // Calls on the held reference reject: left projects are never + // transparently re-opened. `getProject` itself still resolves — the + // wrapper was cached at first acquisition, and re-validating it would + // cost a round trip — but every call on it rejects the same way. + await assert.rejects(() => project.$getProjectSettings(), { + code: ProjectLeftError.code, + }) + const again = await client.getProject(projectId) + assert.equal(again, project, 'project references are permanent') + await assert.rejects(() => again.$getProjectSettings(), { + code: ProjectLeftError.code, + }) +}) + +test('Re-joining after leave makes the same reference work again', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + await project.$getProjectSettings() + + const liveInstance = await serverManager.getProject(projectId) + const closeObserved = pDefer() + liveInstance.once('close', () => closeObserved.resolve(undefined)) + await client.leaveProject(projectId) + await closeObserved.promise + + await assert.rejects(() => project.$getProjectSettings(), { + code: ProjectLeftError.code, + }) + + // Re-invite: core's addProject clears the left state (closing any stale + // instance itself). The already-held reference simply works again. + await serverManager.addProject(projectId) + + const settings = await project.$getProjectSettings() + assert.equal(settings.name, 'mapeo') +}) + +test('Left project without a prior reference: getProject rejects with ProjectLeftError', async (t) => { + const manager = new FakeManager() + const { client } = setup(t, manager) + const projectId = await client.createProject({ name: 'mapeo' }) + + await client.leaveProject(projectId) + + await assert.rejects(() => client.getProject(projectId), { + code: ProjectLeftError.code, + }) +}) + +// Cycling a project open/closed must not accumulate instances in the IPC +// layer's bookkeeping. The fake manager releases its instance on close, so +// any instance still reachable after a cycle is being retained by the IPC +// layer itself (the host detaches rpc-reflector's handler when the instance +// closes; the subscription registry holds no instance references). +// +// Runs only when `global.gc` is available (npm test passes --expose-gc). +test('Repeated server-side close/re-open cycles do not retain prior instances', async (t) => { + if (typeof global.gc !== 'function') { + t.skip('Run with --expose-gc to verify cycle retention') + return + } + + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + // Keep a live subscription across every cycle so the tape path is + // exercised while the WeakRefs are collected. + project.on('some-event', () => {}) + + const N = 5 + /** @type {Array>} */ + const refs = [] + for (let i = 0; i < N; i++) { + refs.push(await cycleAndCaptureWeakRef()) + } + + // In a helper so the instance local goes out of scope before GC runs — + // the loop frame would otherwise pin the final cycle's instance. + /** @returns {Promise>} */ + async function cycleAndCaptureWeakRef() { + await project.$getProjectSettings() + const serverProject = await serverManager.getProject(projectId) + const ref = new WeakRef(serverProject) + await serverProject.close() + return ref + } + + for (let i = 0; i < 20; i++) { + global.gc() + await new Promise((resolve) => setImmediate(resolve)) + } + + const survivors = refs.filter((ref) => ref.deref() !== undefined).length + assert.equal( + survivors, + 0, + `${survivors} of ${N} closed instances retained — expected 0. A non-zero count means IPC is accumulating closed instances cycle over cycle.`, + ) +}) + +test('After a failed getProject, a subsequent getProject for a real project succeeds', async (t) => { + const { client } = setup(t) + + await assert.rejects(() => client.getProject('does-not-exist'), { + code: NotFoundError.code, + }) + + const realId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(realId) + await project.$getProjectSettings() + + // And a retry of the original failing id still rejects (project still + // doesn't exist) — proving the failure path itself is retried, not + // returned from a poisoned cache. + await assert.rejects(() => client.getProject('does-not-exist'), { + code: NotFoundError.code, + }) +}) + +test('Method calls on a never-validated reference to an unknown project reject', async (t) => { + const { client, port2 } = setup(t) + + // Bypass getProject's existence check by writing a request frame straight + // onto an unknown project's channel — this is what a desynced or misbehaving + // client would produce. The server must answer with a per-request error + // response (the handler factory's rejection), not leave the call to time + // out. + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + await project.$getProjectSettings() + + const channelId = '@@comapeo/project/no-such-project' + /** @type {any[]} */ + const responses = [] + /** @param {any} event */ + const captureResponse = (event) => { + if (event.data?.id === channelId) responses.push(event.data.message) + } + port2.addEventListener('message', captureResponse) + t.after(() => port2.removeEventListener('message', captureResponse)) + + // REQUEST frame: [msgType.REQUEST = 0, msgId, propArray, args] + port2.postMessage({ + id: channelId, + message: [0, 99, ['$getProjectSettings'], []], + }) + + await waitFor( + () => responses.length > 0, + 'an error response on the unknown project channel', + ) + const [response] = responses + assert.equal(response[0], 1, 'RESPONSE frame (msgType.RESPONSE)') + assert.equal(response[1], 99, 'answers the request msgId') + assert.equal( + response[2]?.code, + NotFoundError.code, + 'carries the factory rejection, code preserved', + ) + + // And the server stays healthy afterwards. + const settings = await project.$getProjectSettings() + assert.equal(settings.name, 'mapeo') +}) + +// The project channel's rpc handler is bound late to whichever instance is +// live, so a bad method path is resolved against the instance at call time. +// It must fail the same way it did when the instance was bound statically. +test('Calling a method that does not exist rejects with a ReferenceError', async (t) => { + const { client } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + await assert.rejects( + // @ts-expect-error deliberately absent from the API + () => project.notAMethod(), + { name: 'ReferenceError', message: /notAMethod is not defined/ }, + ) + await assert.rejects( + // @ts-expect-error deliberately absent from the API + () => project.observation.notAMethod(), + { name: 'ReferenceError', message: /notAMethod is not defined/ }, + ) + await assert.rejects( + // @ts-expect-error deliberately absent from the API + () => project.noSuchNamespace.create(), + { name: 'ReferenceError', message: /noSuchNamespace is not defined/ }, + ) + + // The project is still healthy afterwards. + const settings = await project.$getProjectSettings() + assert.equal(settings.name, 'mapeo') +}) + +test('Concurrent calls to a dormant project open the instance exactly once', async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + const serverProject = await serverManager.getProject(projectId) + await serverProject.close() + + const before = serverManager.getProjectCallCount.get(projectId) + await Promise.all([ + project.$getProjectSettings(), + project.$getProjectSettings(), + project.observation.create({ + schemaName: 'observation', + attachments: [], + tags: {}, + }), + ]) + + assert.equal( + serverManager.getProjectCallCount.get(projectId), + (before ?? 0) + 1, + 'concurrent calls share a single open', + ) +}) + +test('project.close is not exposed on the client surface', async (t) => { + const { client } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + assert.equal( + // @ts-expect-error close is deliberately absent from the typed surface + project.close, + undefined, + 'lifecycle is server-owned; the client cannot close a project', + ) +}) + +test('A leave that races an in-flight open still rejects with ProjectLeftError', async (t) => { + const manager = new FakeManager() + const { client } = setup(t, manager) + const projectId = await client.createProject({ name: 'mapeo' }) + + // Hold the factory's `manager.getProject` open so a leave can land between + // the factory's first left check and the instance resolving — the window + // core's leaveProject keeps open for up to its sync wait. + const gate = pDefer() + const factoryBlocked = pDefer() + const originalGetProject = manager.getProject.bind(manager) + let intercepted = false + manager.getProject = async (id) => { + if (!intercepted) { + intercepted = true + factoryBlocked.resolve(undefined) + await gate.promise + } + return originalGetProject(id) + } + + const acquiring = client.getProject(projectId) + await factoryBlocked.promise + await manager.leaveProject(projectId) + gate.resolve(undefined) + + // Without the post-open re-check the factory would bind the gutted + // instance and the call would resolve with garbage. + await assert.rejects(() => acquiring, { code: ProjectLeftError.code }) +}) + +test("Core's double close emit is harmless to the host", async (t) => { + const { client, serverManager } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + /** @type {unknown[]} */ + const received = [] + project.on('some-event', (value) => received.push(value)) + await project.$getProjectSettings() + + const instance = await serverManager.getProject(projectId) + let closeEmits = 0 + instance.on('close', () => closeEmits++) + await instance.close() + assert.equal(closeEmits, 2, 'fake models core: close is emitted twice') + + // A stray extra emit after close must also be harmless (detach is + // idempotent and the host's once() is already consumed). + instance.emit('close') + + const settings = await project.$getProjectSettings() + assert.equal(settings.name, 'mapeo') + + const fresh = await serverManager.getProject(projectId) + fresh.emit('some-event', 'after-reopen') + await project.$getProjectSettings() + assert.deepEqual(received, ['after-reopen'], 'subscriptions survived') +}) + +test('Opening while an instance close is in flight waits it out and binds the fresh instance', async (t) => { + const manager = new FakeManager() + const { client } = setup(t, manager) + const projectId = await client.createProject({ name: 'mapeo' }) + + // Open server-side only, then start a close held open by a gate: `closing` + // is set, `closed` is not, and the manager cache still returns the dying + // instance (core evicts only on the `close` event). + const dying = await manager.getProject(projectId) + const gate = pDefer() + dying.holdClose(gate.promise) + const closing = dying.close() + + // First acquisition arrives mid-close: the host must wait the close out + // and retry, not bind the corpse. + const acquiring = client.getProject(projectId) + await waitFor( + () => (manager.getProjectCallCount.get(projectId) ?? 0) >= 2, + 'the factory to observe the dying instance', + ) + gate.resolve(undefined) + await closing + + const project = await acquiring + const settings = await project.$getProjectSettings() + assert.equal(settings.name, 'mapeo') + + const fresh = await manager.getProject(projectId) + assert.notEqual(fresh, dying, 'bound instance is the post-close one') + assert.equal(fresh.closed, false) + assert.equal( + manager.getProjectCallCount.get(projectId), + 4, + 'test open + dying observation + retry + this getProject', + ) +}) + +test('Gives up after the project keeps closing while opening', async (t) => { + const manager = new FakeManager() + const { client } = setup(t, manager) + const projectId = await client.createProject({ name: 'mapeo' }) + + // Every open observes an instance whose close is already in flight — a + // pathological resource policy closing projects as fast as they open. + const originalGetProject = manager.getProject.bind(manager) + manager.getProject = async (id) => { + const project = await originalGetProject(id) + project.close() + return project + } + + await assert.rejects(() => client.getProject(projectId), { + message: /kept closing while opening/, + }) +}) + +test('Re-invite after leaving a never-acquired project: getProject succeeds', async (t) => { + const manager = new FakeManager() + const { client } = setup(t, manager) + const projectId = await client.createProject({ name: 'mapeo' }) + + await client.leaveProject(projectId) + await assert.rejects(() => client.getProject(projectId), { + code: ProjectLeftError.code, + }) + + // Re-invite: nothing was cached for this id (the failed acquisition is not + // cached either), so the next getProject validates afresh and succeeds. + await manager.addProject(projectId) + const project = await client.getProject(projectId) + const settings = await project.$getProjectSettings() + assert.equal(settings.name, 'mapeo') +}) + +test('Parallel first getProject calls make exactly one validation round trip', async (t) => { + const { port1, port2 } = new MessageChannel() + const manager = new FakeManager() + let assertCalls = 0 + const server = createComapeoCoreServer(/** @type {any} */ (manager), port1, { + onRequestHook: (request, next) => { + if (request.method.join('.') === 'assertProjectExists') assertCalls++ + next(request) + }, + }) + const client = createComapeoCoreClient(port2) + port1.start() + port2.start() + t.after(async () => { + server.close() + await closeComapeoCoreClient(client) + port1.close() + port2.close() + }) + + const projectId = await client.createProject({ name: 'mapeo' }) + const [a, b, c] = await Promise.all([ + client.getProject(projectId), + client.getProject(projectId), + client.getProject(projectId), + ]) + assert.equal(a, b) + assert.equal(b, c) + assert.equal(assertCalls, 1, 'concurrent first calls share one round trip') + + await client.getProject(projectId) + assert.equal(assertCalls, 1, 'cached wrapper: no further round trips') +}) + +test('A consumer onRequestHook composes with the interim leave hook', async (t) => { + const { port1, port2 } = new MessageChannel() + const manager = new FakeManager() + /** @type {string[]} */ + const hookedMethods = [] + const server = createComapeoCoreServer(/** @type {any} */ (manager), port1, { + onRequestHook: (request, next) => { + hookedMethods.push(request.method.join('.')) + next(request) + }, + }) + const client = createComapeoCoreClient(port2) + port1.start() + port2.start() + t.after(async () => { + server.close() + await closeComapeoCoreClient(client) + port1.close() + port2.close() + }) + + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + await project.$getProjectSettings() + + const liveInstance = await manager.getProject(projectId) + const closeObserved = pDefer() + liveInstance.once('close', () => closeObserved.resolve(undefined)) + + await client.leaveProject(projectId) + // The interim leave hook still ran under the consumer hook: the gutted + // instance gets closed. + await closeObserved.promise + + assert.ok( + hookedMethods.includes('createProject'), + 'consumer hook saw manager calls', + ) + assert.ok( + hookedMethods.includes('leaveProject'), + 'consumer hook saw the leave', + ) + await assert.rejects(() => project.$getProjectSettings(), { + code: ProjectLeftError.code, + }) +}) + +test('A failed leaveProject still closes the opened instance', async (t) => { + const manager = new FakeManager() + const { client } = setup(t, manager) + const projectId = await client.createProject({ name: 'mapeo' }) + + // Leave fails after opening the instance (mirrors core: leave can fail + // mid-way, after `getProject`); the cleanup close must run regardless. + manager.leaveProject = async (id) => { + await manager.getProject(id) + throw new Error('leave failed mid-way') + } + + const instance = await manager.getProject(projectId) + const closeObserved = pDefer() + instance.once('close', () => closeObserved.resolve(undefined)) + + await assert.rejects(() => client.leaveProject(projectId), { + message: /leave failed mid-way/, + }) + await closeObserved.promise + + // The project was never actually left, so it simply re-opens and works. + const project = await client.getProject(projectId) + const settings = await project.$getProjectSettings() + assert.equal(settings.name, 'mapeo') +}) diff --git a/tests/transport-reset.js b/tests/transport-reset.js new file mode 100644 index 0000000..59d3b24 --- /dev/null +++ b/tests/transport-reset.js @@ -0,0 +1,257 @@ +import test from 'node:test' +import assert from 'node:assert/strict' +import pDefer from 'p-defer' + +import { + createComapeoCoreClient, + closeComapeoCoreClient, + notifyTransportReset, + resubscribe, + createComapeoServicesClient, + closeComapeoServicesClient, +} from '../src/client.js' +import { + createComapeoCoreServer, + createComapeoServicesServer, +} from '../src/server.js' +import { RpcChannelClosedError } from '../src/errors.js' + +import { setup } from './helpers.js' +import { FakeManager } from './fake-manager.js' + +/** + * Simulate the server process dying and restarting while the client (and the + * message port it holds) stays alive — the Android foreground-service restart + * case: close the old server and serve a fresh manager over the same port. + * + * @param {import('node:test').TestContext} t + * @param {ReturnType['server']} oldServer + * @param {MessagePort} serverPort + */ +function restartServer(t, oldServer, serverPort) { + oldServer.close() + const newManager = new FakeManager() + const newServer = createComapeoCoreServer( + /** @type {any} */ (newManager), + serverPort, + ) + t.after(() => newServer.close()) + return { newManager, newServer } +} + +test('Reset rejects in-flight manager and getProject calls with RpcChannelClosedError', async (t) => { + const { client, server } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + + // With the server gone, these calls can never be answered. getProject is a + // first acquisition, so it has a routing round trip in flight. + server.close() + const inFlightManagerCall = client.listProjects() + const inFlightGetProject = client.getProject(projectId) + + notifyTransportReset(client) + + await assert.rejects(() => inFlightManagerCall, { + code: RpcChannelClosedError.code, + }) + await assert.rejects(() => inFlightGetProject, { + code: RpcChannelClosedError.code, + }) +}) + +test('Reset rejects in-flight project method calls with RpcChannelClosedError', async (t) => { + const { client, server } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + server.close() + const inFlightProjectCall = project.$getProjectSettings() + + notifyTransportReset(client) + + await assert.rejects(() => inFlightProjectCall, { + code: RpcChannelClosedError.code, + }) + + // The reference itself is unaffected — only the in-flight call died. + // (There is no server right now, so just check the next call is a fresh + // pending promise, not an instant rejection.) + const nextCall = project.$getProjectSettings() + let settled = false + nextCall.then( + () => (settled = true), + () => (settled = true), + ) + await new Promise((resolve) => setImmediate(resolve)) + assert.equal(settled, false, 'post-reset call is pending, not dead') + notifyTransportReset(client) + await assert.rejects(() => nextCall, { code: RpcChannelClosedError.code }) +}) + +test('Project references survive a server restart', async (t) => { + const { client, server, port1 } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + await project.$getProjectSettings() + + const { newManager } = restartServer(t, server, port1) + notifyTransportReset(client) + + // The fresh manager mints the same project id ('project-1'). The held + // reference must keep working against the restarted server — its channel + // is keyed by project id, which the new server serves identically. No + // re-getProject required. + const newProjectId = await client.createProject({ name: 'mapeo-after' }) + assert.equal(newProjectId, projectId, 'test setup: same project id reminted') + + const settings = await project.$getProjectSettings() + assert.equal(settings.name, 'mapeo-after') + assert.equal(newManager.getProjectCallCount.get(projectId), 1) + + // getProject still hands back the same permanent reference. + const again = await client.getProject(projectId) + assert.equal(again, project) +}) + +test('Reset does not resubscribe; resubscribe replays subscriptions and is idempotent', async (t) => { + const { client, server, port1 } = setup(t) + + /** @type {unknown[]} */ + const received = [] + client.on('local-peers', (peers) => received.push(peers)) + await client.listProjects() + + const { newManager } = restartServer(t, server, port1) + notifyTransportReset(client) + // Round-trip barrier: had the reset replayed the subscription, the ON + // message would have been processed by now. + await client.listProjects() + + const peers = [{ deviceId: 'peer-a' }] + newManager.emit('local-peers', peers) + // Barrier to let any (unexpected) forwarded event flush through the port. + await client.listProjects() + assert.deepEqual(received, [], 'reset alone does not replay subscriptions') + + // Repeated calls are safe: the server ignores duplicate subscriptions, so + // events are not double-delivered. + resubscribe(client) + resubscribe(client) + await client.listProjects() + + newManager.emit('local-peers', peers) + await client.listProjects() + assert.deepEqual( + received, + [peers], + 'event is delivered exactly once after resubscribing', + ) +}) + +test('Project event subscriptions are replayed across a server restart', async (t) => { + const { client, server, port1 } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + /** @type {import('p-defer').DeferredPromise} */ + const deferred = pDefer() + project.on('some-event', (value) => deferred.resolve(value)) + await project.$getProjectSettings() + + const { newManager } = restartServer(t, server, port1) + notifyTransportReset(client) + await client.createProject({ name: 'mapeo' }) // remints project-1 + + // Replaying the project subscription also re-opens the project on the + // restarted server: the ON frame lands on a host whose handler factory has + // never run, and binding for it opens a fresh instance — a consumer that + // only listens starts receiving events again with no method call. + resubscribe(client) + await project.$getProjectSettings() + + const instance = await newManager.getProject(projectId) + instance.emit('some-event', 'after-restart') + assert.equal(await deferred.promise, 'after-restart') +}) + +test('A client that only listens gets events again after restart + resubscribe', async (t) => { + const { client, server, port1 } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + /** @type {import('p-defer').DeferredPromise} */ + const deferred = pDefer() + project.on('some-event', (value) => deferred.resolve(value)) + await project.$getProjectSettings() + + const { newManager } = restartServer(t, server, port1) + notifyTransportReset(client) + await client.createProject({ name: 'mapeo' }) // remints project-1 + + // No method call after resubscribing: the replayed ON alone must re-open + // the project (single-flight via the late-bound factory) and re-attach the + // subscription. + resubscribe(client) + await waitFor( + () => (newManager.getProjectCallCount.get(projectId) ?? 0) > 0, + 'the restarted server to open the project for the replayed subscription', + ) + + const instance = await newManager.getProject(projectId) + instance.emit('some-event', 'listen-only') + assert.equal(await deferred.promise, 'listen-only') +}) + +test('Reset and resubscribe are no-ops after the client is closed', async (t) => { + const { client } = setup(t) + await closeComapeoCoreClient(client) + + assert.doesNotThrow(() => notifyTransportReset(client)) + assert.doesNotThrow(() => resubscribe(client)) +}) + +test('Services client: reset rejects in-flight calls; resubscribe is safe', async (t) => { + const { port1, port2 } = new MessageChannel() + const services = { + mapServer: { + getBaseUrl: () => new Promise(() => {}), // never answers + }, + } + const servicesServer = createComapeoServicesServer(services, port1) + const servicesClient = createComapeoServicesClient(port2) + port1.start() + port2.start() + t.after(() => { + servicesServer.close() + closeComapeoServicesClient(servicesClient) + port1.close() + port2.close() + }) + + const inFlight = servicesClient.mapServer.getBaseUrl() + notifyTransportReset(servicesClient) + + await assert.rejects(() => inFlight, { code: RpcChannelClosedError.code }) + + assert.doesNotThrow(() => resubscribe(servicesClient)) + + // Both are safe after the services client is closed, too. + closeComapeoServicesClient(servicesClient) + assert.doesNotThrow(() => notifyTransportReset(servicesClient)) + assert.doesNotThrow(() => resubscribe(servicesClient)) +}) + +/** + * Poll until `predicate` holds, for assertions about work the server does on + * its own initiative (with no call to await). + * + * @param {() => boolean} predicate + * @param {string} description + */ +async function waitFor(predicate, description) { + for (let i = 0; i < 200; i++) { + if (predicate()) return + await new Promise((resolve) => setTimeout(resolve, 5)) + } + assert.fail(`Timed out waiting for ${description}`) +} diff --git a/tests/transport.js b/tests/transport.js index c05e395..783f43c 100644 --- a/tests/transport.js +++ b/tests/transport.js @@ -39,12 +39,18 @@ test('Malformed and unroutable messages are ignored and the client keeps working // Well-formed envelope without our channel prefix: a foreign sender sharing // the port. Dropped silently — not our traffic, so no warning. port2.postMessage({ id: 'someone-elses-channel', message: { value: 'x' } }) - // Well-formed envelope carrying our prefix but for an instance id the server - // never opened — a genuine routing miss. Posted twice to exercise the - // warn-once dedupe. - const unknownId = '@@comapeo/project/project-999:42' + // Well-formed envelope carrying our prefix but matching no channel shape — + // a genuine routing miss. Posted twice to exercise the warn-once dedupe. + const unknownId = '@@comapeo/bogus-channel' port2.postMessage({ id: unknownId, message: { value: 'whatever' } }) port2.postMessage({ id: unknownId, message: { value: 'whatever' } }) + // A project-shaped id for a project that doesn't exist, carrying a frame + // that isn't a request: rpc-reflector rejects it as an invalid message and + // no project open is ever attempted — nothing to respond to, no log. + port2.postMessage({ + id: '@@comapeo/project/project-999', + message: { value: 'whatever' }, + }) // Let the messages flush through the event loop. await new Promise((resolve) => setImmediate(resolve)) @@ -56,12 +62,14 @@ test('Malformed and unroutable messages are ignored and the client keeps working assert.equal(settings.name, 'mapeo') // The prefixed-but-unroutable id is logged exactly once (deduped). - const unrecognised = logs.filter((w) => w.includes('project-999:42')) + const unrecognised = logs.filter((w) => w.includes('bogus-channel')) assert.equal(unrecognised.length, 1) - // The foreign (unprefixed) envelope and the structurally-invalid messages - // are dropped without any log. + // The foreign (unprefixed) envelope, the structurally-invalid messages, and + // the non-request frame for an unknown project are dropped without any log. const foreign = logs.filter((w) => w.includes('someone-elses-channel')) assert.equal(foreign.length, 0) + const unknownProject = logs.filter((w) => w.includes('project-999')) + assert.equal(unknownProject.length, 0) }) test('createComapeoCoreServer().close() is idempotent', async (t) => {