diff --git a/src/client.js b/src/client.js index 5cb3524..9503ba9 100644 --- a/src/client.js +++ b/src/client.js @@ -6,7 +6,11 @@ import { SERVICES_ID, SubChannel, } from './lib/sub-channel.js' -import { ClientClosedError, ProjectClosedError } from './errors.js' +import { + ClientClosedError, + ProjectClosedError, + TransportClosedError, +} from './errors.js' /** @import { ClientApi, MessagePortLike } from 'rpc-reflector' */ /** @import { MapeoProject, MapeoManager } from '@comapeo/core' */ @@ -29,13 +33,24 @@ const EMITTER_METHODS = new Set([ 'listenerCount', ]) +// Removing a listener from a client that is already dead is correct teardown +// behaviour (e.g. React effect cleanup running against a stale reference), so +// unlike the other emitter methods these must not throw on a closed proxy. +const EMITTER_UNSUBSCRIBE_METHODS = new Set([ + 'removeListener', + 'off', + 'removeAllListeners', +]) + /** * 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. + * synchronously instead, at the call site. Unsubscribe methods are a further + * exception: they are no-ops that return the proxy for chaining, because + * removing a listener from a dead client is valid teardown, not a bug. * * @param {() => Error} makeError */ @@ -44,6 +59,9 @@ function createClosedProxy(makeError) { const handler = { get(_target, prop) { if (typeof prop === 'string' && EMITTER_METHODS.has(prop)) { + if (EMITTER_UNSUBSCRIBE_METHODS.has(prop)) { + return () => proxy + } return () => { throw makeError() } @@ -57,7 +75,8 @@ function createClosedProxy(makeError) { return Promise.reject(makeError()) }, } - return new Proxy({}, handler) + const proxy = new Proxy({}, handler) + return proxy } /** @@ -75,6 +94,8 @@ function createClosedProxy(makeError) { * >} ComapeoCoreClientApi */ const CLOSE = Symbol('close') +const TRANSPORT_RESET = Symbol('transportReset') +const RESUBSCRIBE = Symbol('resubscribe') /** * @param {MessagePortLike} messagePort @@ -108,6 +129,7 @@ export function createComapeoCoreClient(messagePort, opts = {}) { * @type {Set<{ * client: ClientApi, * channel: SubChannel, + * hardClose: (error: Error) => void, * }>} */ const openProjectClients = new Set() @@ -129,6 +151,46 @@ export function createComapeoCoreClient(messagePort, opts = {}) { let clientClosed = false const managerClosedProxy = createClosedProxy(() => new ClientClosedError()) + // Bumped on every transport reset so a `getProject` whose routing response + // arrived just before the reset cannot cache a wrapper bound to the dead + // server (see `resolveProjectClient`). + let resetGeneration = 0 + + function handleTransportReset() { + if (clientClosed) return + resetGeneration++ + + // Fail in-flight calls fast with a distinguishable, retryable 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 + // `resubscribeCoreClient` once the transport is back up. + createClient.rejectPending(managerClient, new TransportClosedError()) + createClient.rejectPending(projectRoutingClient, new TransportClosedError()) + + // Project instance ids are minted by a counter that restarts with the + // server, so a restarted server can mint an id equal to the one a cached + // wrapper is bound to — the instance-id currency check in + // `resolveProjectClient` could then falsely pass. Hard-close every + // wrapper and drop the cache so `getProject` always builds a fresh + // wrapper against the new server. + for (const entry of openProjectClients) { + entry.hardClose(new TransportClosedError()) + } + openProjectClients.clear() + currentProjectClients.clear() + } + + function handleResubscribe() { + // rpc-reflector's resubscribe is a no-op on a closed client, and the + // server ignores duplicate ON messages, so this is safe to call + // repeatedly and after close. + if (clientClosed) return + createClient.resubscribe(managerClient) + createClient.resubscribe(projectRoutingClient) + } + const client = new Proxy(managerClient, { get(target, prop, receiver) { if (prop === CLOSE) { @@ -157,6 +219,14 @@ export function createComapeoCoreClient(messagePort, opts = {}) { } } + if (prop === TRANSPORT_RESET) { + return handleTransportReset + } + + if (prop === RESUBSCRIBE) { + return handleResubscribe + } + if (prop === 'getProject') { return createProjectClient } @@ -205,9 +275,15 @@ export function createComapeoCoreClient(messagePort, opts = {}) { * @returns {Promise>} */ async function resolveProjectClient(projectPublicId) { + const generation = resetGeneration const instanceId = await projectRoutingClient.assertProjectExists(projectPublicId) + // A reset can land between the routing response arriving and this + // continuation running; a wrapper minted now would be bound to the dead + // server, so reject like any other call in flight during the reset. + if (generation !== resetGeneration) throw new TransportClosedError() + const current = currentProjectClients.get(projectPublicId) if (current && current.instanceId === instanceId) return current.wrapper @@ -232,9 +308,6 @@ export function createComapeoCoreClient(messagePort, opts = {}) { 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 @@ -254,6 +327,31 @@ export function createComapeoCoreClient(messagePort, opts = {}) { const closedProxy = createClosedProxy(() => closed ? new ProjectClosedError() : new ClientClosedError(), ) + + const registryEntry = { + client: projectClient, + channel: projectChannel, + // Local-only teardown for a transport reset: the server this instance + // belonged to is gone, so there is no wire close to await. Stale + // references then behave like a closed project (`ProjectClosedError`), + // and `close()` on them resolves like an already-closed project. + hardClose: (/** @type {Error} */ error) => { + createClient.rejectPending(projectClient, error) + // Fire 'close' listeners (the app's teardown listeners and the + // cache-eviction listener below) before closing the client, matching + // what a server-initiated close delivers. Emitted before close so a + // `.off` called from inside a close listener hits a still-open + // client (a harmless OFF frame into a dead socket) rather than + // throwing. + createClient.emitLocal(projectClient, 'close') + createClient.close(projectClient) + projectChannel.close() + closed = true + closePromise = Promise.resolve() + }, + } + openProjectClients.add(registryEntry) + const wrappedProjectClient = new Proxy(projectClient, { get(target, prop, receiver) { if (prop === 'close') { @@ -298,6 +396,53 @@ export async function closeComapeoCoreClient(client) { return client[CLOSE]() } +/** + * Notify the core client that the underlying transport has dropped (e.g. + * Android killed the foreground service hosting the server). Call at drop + * time — the transport does not need to be back up. This: + * + * - rejects every call that was in flight with `TransportClosedError` + * (`code: 'RPC_TRANSPORT_CLOSED'`), so callers can fail fast and retry + * instead of waiting for the per-call timeout; + * - hard-closes every open project client (each fires its `'close'` event + * locally, so app-held `once('close')` teardown listeners run) and drops + * the project cache, so a later `getProject` builds a fresh wrapper + * against the new server. Stale project references held by the app behave + * like closed projects (`ProjectClosedError`), except that removing + * listeners from them is a harmless no-op. A `getProject` in flight during + * the reset rejects with `TransportClosedError`. + * + * This deliberately does NOT replay event subscriptions: writing into a + * still-down transport can keep nudging it into a hot retry loop. Once the + * transport has reconnected to the restarted server, call + * {@link resubscribeCoreClient} to restore subscriptions. + * + * No-op after `closeComapeoCoreClient`. + * + * @param {ComapeoCoreClientApi} client client created with `createComapeoCoreClient` + * @returns {void} + */ +export function notifyCoreClientTransportReset(client) { + // @ts-expect-error + return client[TRANSPORT_RESET]() +} + +/** + * Re-send the core client's event subscriptions (manager events and project + * routing) after the transport has reconnected to a restarted server, which + * lost all subscription state. Call once the transport is back up, after + * having called `notifyCoreClientTransportReset` at drop time. Safe to call + * repeatedly (the server ignores duplicate subscriptions); no-op after + * `closeComapeoCoreClient`. + * + * @param {ComapeoCoreClientApi} client client created with `createComapeoCoreClient` + * @returns {void} + */ +export function resubscribeCoreClient(client) { + // @ts-expect-error + return client[RESUBSCRIBE]() +} + /** * @typedef {ClientApi} ComapeoServicesClientApi */ @@ -329,3 +474,33 @@ export function createComapeoServicesClient(messagePort, opts = {}) { export function closeComapeoServicesClient(servicesClient) { createClient.close(servicesClient) } + +/** + * Notify the services client that the underlying transport has dropped: + * rejects every call that was in flight with `TransportClosedError` + * (`code: 'RPC_TRANSPORT_CLOSED'`). Call at drop time. Like + * {@link notifyCoreClientTransportReset} this does not replay event + * subscriptions — call {@link resubscribeServicesClient} once the transport + * has reconnected. No-op after `closeComapeoServicesClient`. + * + * @param {ComapeoServicesClientApi} servicesClient client created with `createComapeoServicesClient` + * @returns {void} + */ +export function notifyServicesClientTransportReset(servicesClient) { + createClient.rejectPending(servicesClient, new TransportClosedError()) +} + +/** + * Re-send the services client's event subscriptions after the transport has + * reconnected to a restarted server, which lost all subscription state. Call + * once the transport is back up, after having called + * `notifyServicesClientTransportReset` at drop time. Safe to call repeatedly + * (the server ignores duplicate subscriptions); no-op after + * `closeComapeoServicesClient`. + * + * @param {ComapeoServicesClientApi} servicesClient client created with `createComapeoServicesClient` + * @returns {void} + */ +export function resubscribeServicesClient(servicesClient) { + createClient.resubscribe(servicesClient) +} diff --git a/src/errors.js b/src/errors.js index 68a0afd..6c95b6c 100644 --- a/src/errors.js +++ b/src/errors.js @@ -16,6 +16,20 @@ export const ProjectClosedError = createErrorClass({ status: 410, }) +/** + * Rejected client-side into calls that were in flight when the underlying + * transport dropped (e.g. the process hosting the server was killed and + * restarted). Distinguishable from a real failure or a timeout: the call + * never completed on the server's side of the connection, so a read is safe + * to retry once the transport has reconnected. + */ +export const TransportClosedError = createErrorClass({ + code: 'RPC_TRANSPORT_CLOSED', + message: + 'Transport closed: the connection to the server dropped while the call was in flight', + status: 503, +}) + /** * Thrown client-side when a method is called after the CoMapeo core client * (the whole IPC client) has been closed via `closeComapeoCoreClient`. diff --git a/src/index.js b/src/index.js index 1125e16..ad5ad77 100644 --- a/src/index.js +++ b/src/index.js @@ -1,8 +1,12 @@ export { createComapeoCoreClient, closeComapeoCoreClient, + notifyCoreClientTransportReset, + resubscribeCoreClient, createComapeoServicesClient, closeComapeoServicesClient, + notifyServicesClientTransportReset, + resubscribeServicesClient, } from './client.js' export { createComapeoCoreServer, diff --git a/tests/events.js b/tests/events.js index 169eadc..877a9f4 100644 --- a/tests/events.js +++ b/tests/events.js @@ -52,7 +52,7 @@ 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) => { +test('EventEmitter subscribe methods throw synchronously after the project is closed; unsubscribe methods are no-ops', async (t) => { const { client } = setup(t) const projectId = await client.createProject({ name: 'mapeo' }) const project = await client.getProject(projectId) @@ -60,13 +60,14 @@ test('EventEmitter methods throw synchronously after the project is closed', asy await project.close() // Emitter methods are not awaited by callers, so a rejected promise would - // surface as an unhandled rejection — they throw at the call site instead. + // surface as an unhandled rejection — subscribe methods throw at the call + // site instead. Unsubscribe methods are valid teardown on a closed + // reference (e.g. React effect cleanup), so they are no-ops. assert.throws(() => project.on('some-event', () => {}), { code: ProjectClosedError.code, }) - assert.throws(() => project.removeListener('some-event', () => {}), { - code: ProjectClosedError.code, - }) + assert.doesNotThrow(() => project.removeListener('some-event', () => {})) + assert.doesNotThrow(() => project.off('some-event', () => {})) }) test('EventEmitter methods throw synchronously after the client is closed', async (t) => { diff --git a/tests/transport-reset.js b/tests/transport-reset.js new file mode 100644 index 0000000..1e00c4b --- /dev/null +++ b/tests/transport-reset.js @@ -0,0 +1,334 @@ +import test from 'node:test' +import assert from 'node:assert/strict' +import { EventEmitter } from 'node:events' +import pDefer from 'p-defer' + +import { + createComapeoCoreClient, + closeComapeoCoreClient, + notifyCoreClientTransportReset, + resubscribeCoreClient, + createComapeoServicesClient, + closeComapeoServicesClient, + notifyServicesClientTransportReset, + resubscribeServicesClient, +} from '../src/client.js' +import { + createComapeoCoreServer, + createComapeoServicesServer, +} from '../src/server.js' +import { ProjectClosedError, TransportClosedError } 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 TransportClosedError', async (t) => { + const { client, server } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + + // With the server gone, these calls can never be answered. + server.close() + const inFlightManagerCall = client.listProjects() + const inFlightGetProject = client.getProject(projectId) + + notifyCoreClientTransportReset(client) + + await assert.rejects(() => inFlightManagerCall, { + code: TransportClosedError.code, + }) + await assert.rejects(() => inFlightGetProject, { + code: TransportClosedError.code, + }) +}) + +test('Reset rejects in-flight project method calls with TransportClosedError', 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() + + notifyCoreClientTransportReset(client) + + await assert.rejects(() => inFlightProjectCall, { + code: TransportClosedError.code, + }) +}) + +test('Reset does not resubscribe; resubscribeCoreClient 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) + notifyCoreClientTransportReset(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. + resubscribeCoreClient(client) + resubscribeCoreClient(client) + await client.listProjects() + + newManager.emit('local-peers', peers) + await client.listProjects() + assert.deepEqual( + received, + [peers], + 'event is delivered exactly once after resubscribing', + ) +}) + +test('Stale project wrapper is not reused after reset, even when the new server mints the same instance id', async (t) => { + const { client, server, port1 } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const staleProject = await client.getProject(projectId) + await staleProject.$getProjectSettings() + + const { newManager } = restartServer(t, server, port1) + notifyCoreClientTransportReset(client) + + // The fresh manager mints the same project id ('project-1') and the fresh + // server's instance counter restarts, so `assertProjectExists` returns an + // instance id identical to the one the stale wrapper is bound to — the + // cache must have been dropped for this to return a working wrapper. + const newProjectId = await client.createProject({ name: 'mapeo' }) + assert.equal(newProjectId, projectId, 'test setup: same project id reminted') + + const freshProject = await client.getProject(projectId) + assert.notEqual( + freshProject, + staleProject, + 'getProject returns a fresh wrapper, not the stale one', + ) + const settings = await freshProject.$getProjectSettings() + assert.equal(settings.name, 'mapeo') + + // Only the fresh manager's project should be reachable. + assert.equal(newManager.getProjectCallCount.get(projectId), 1) +}) + +test('Reset fires each project wrapper’s close event exactly once', async (t) => { + const { client, server, port1 } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + + let closeCount = 0 + const noop = () => {} + project.on('close', () => { + closeCount++ + // Teardown code commonly removes listeners from inside a close + // listener; this must not throw mid-teardown. + assert.doesNotThrow(() => project.off('own-role-change', noop)) + }) + + restartServer(t, server, port1) + notifyCoreClientTransportReset(client) + assert.equal(closeCount, 1, 'close listener ran synchronously with reset') + + // A second reset finds no open project clients; the closed wrapper's + // listeners must not fire again. + notifyCoreClientTransportReset(client) + assert.equal(closeCount, 1) +}) + +test('Stale project references behave like closed projects after reset', async (t) => { + const { client, server, port1 } = setup(t) + const projectId = await client.createProject({ name: 'mapeo' }) + const staleProject = await client.getProject(projectId) + + restartServer(t, server, port1) + notifyCoreClientTransportReset(client) + + await assert.rejects(() => staleProject.$getProjectSettings(), { + code: ProjectClosedError.code, + }) + await assert.rejects( + () => + staleProject.observation.create({ + schemaName: 'observation', + attachments: [], + tags: {}, + }), + { code: ProjectClosedError.code }, + ) + // `close()` on a stale reference resolves like an already-closed project. + await staleProject.close() + + // Removing listeners from a stale reference is valid teardown (e.g. React + // effect cleanup) and must not throw; subscribing is still an error. + const noop = () => {} + assert.doesNotThrow(() => staleProject.off('close', noop)) + assert.doesNotThrow(() => staleProject.removeListener('close', noop)) + assert.doesNotThrow(() => staleProject.removeAllListeners()) + assert.doesNotThrow( + () => staleProject.off('own-role-change', noop).off('close', noop), + 'unsubscribe no-ops keep chaining semantics', + ) + assert.throws( + () => staleProject.on('close', noop), + { code: ProjectClosedError.code }, + 'subscribe methods still throw on a stale reference', + ) + assert.throws(() => staleProject.once('close', noop), { + code: ProjectClosedError.code, + }) +}) + +test('getProject that is in flight during reset rejects, and a retry returns a working client', async (t) => { + const manager = new FakeManager() + /** @type {import('p-defer').DeferredPromise} */ + const gate = pDefer() + // Hold `getProject` open server-side so the routing call is still in + // flight when the reset lands. + const originalGetProject = manager.getProject.bind(manager) + manager.getProject = () => gate.promise + + const { client, server, port1 } = setup(t, manager) + const projectId = await client.createProject({ name: 'mapeo' }) + + const inFlightGetProject = client.getProject(projectId) + + const { newManager } = restartServer(t, server, port1) + notifyCoreClientTransportReset(client) + + await assert.rejects(() => inFlightGetProject, { + code: TransportClosedError.code, + }) + manager.getProject = originalGetProject + + // The dedupe map entry for the rejected call must not poison the retry. + await client.createProject({ name: 'mapeo' }) + const project = await client.getProject(projectId) + const settings = await project.$getProjectSettings() + assert.equal(settings.name, 'mapeo') + assert.ok(newManager.getProjectCallCount.get(projectId)) +}) + +test('Reset is a no-op after the client is closed', async (t) => { + const { port1, port2 } = new MessageChannel() + const manager = new FakeManager() + const server = createComapeoCoreServer(/** @type {any} */ (manager), port1) + const client = createComapeoCoreClient(port2) + port1.start() + port2.start() + t.after(() => { + server.close() + port1.close() + port2.close() + }) + + await closeComapeoCoreClient(client) + + assert.doesNotThrow(() => notifyCoreClientTransportReset(client)) + assert.doesNotThrow(() => resubscribeCoreClient(client)) +}) + +test('Services client reset rejects in-flight calls; resubscribeServicesClient replays subscriptions', async (t) => { + const { port1, port2 } = new MessageChannel() + t.after(() => { + port1.close() + port2.close() + }) + + const makeApi = () => + Object.assign(new EventEmitter(), { + mapServer: { + async getBaseUrl() { + return 'http://localhost:3000' + }, + }, + }) + + const oldServer = createComapeoServicesServer(makeApi(), port1) + const client = createComapeoServicesClient(port2) + port1.start() + port2.start() + t.after(() => closeComapeoServicesClient(client)) + + /** @type {unknown[]} */ + const received = [] + const eventfulClient = /** @type {any} */ (client) + eventfulClient.on('service-event', (/** @type {unknown} */ e) => + received.push(e), + ) + await client.mapServer.getBaseUrl() + + oldServer.close() + const inFlightCall = client.mapServer.getBaseUrl() + + // Restart: a fresh services server (with a fresh emitter) on the same port. + const newApi = makeApi() + const newServer = createComapeoServicesServer(newApi, port1) + t.after(() => newServer.close()) + + notifyServicesClientTransportReset(client) + + await assert.rejects(() => inFlightCall, { + code: TransportClosedError.code, + }) + + // Reset alone does not replay subscriptions. + await client.mapServer.getBaseUrl() + newApi.emit('service-event', 'dropped') + await client.mapServer.getBaseUrl() + assert.deepEqual(received, []) + + // Safe to call repeatedly: the server ignores duplicate subscriptions. + resubscribeServicesClient(client) + resubscribeServicesClient(client) + // Round-trip barrier so the replayed subscribe has been processed. + await client.mapServer.getBaseUrl() + newApi.emit('service-event', 'payload') + // Let the forwarded event flush through the port. + await client.mapServer.getBaseUrl() + + assert.deepEqual(received, ['payload']) +}) + +test('Services client reset and resubscribe are no-ops after close', async (t) => { + const { port1, port2 } = new MessageChannel() + t.after(() => { + port1.close() + port2.close() + }) + + const client = createComapeoServicesClient(port2) + port2.start() + closeComapeoServicesClient(client) + + assert.doesNotThrow(() => notifyServicesClientTransportReset(client)) + assert.doesNotThrow(() => resubscribeServicesClient(client)) +})