diff --git a/packages/keyring-sdk/CHANGELOG.md b/packages/keyring-sdk/CHANGELOG.md index 6a081453b..3bc632012 100644 --- a/packages/keyring-sdk/CHANGELOG.md +++ b/packages/keyring-sdk/CHANGELOG.md @@ -7,6 +7,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added + +- Add `AtomicKeyring` base class, `AtomicUpdater` type, and `isAtomicKeyring` helper for keyrings that serialize their own operations and commit state through controller-mediated vault writes ([#640](https://github.com/MetaMask/accounts/pull/640)) + - Atomic keyrings run long-running operations inside their own read-write lock, entirely outside the controller mutex, and persist state through the injected `update` callback. +- Add `ReadWriteLock` primitive with parallel readers, exclusive writers ([#640](https://github.com/MetaMask/accounts/pull/640)) + - The lock is write-preferring by default: a pending writer blocks new readers, so a steady stream of reads cannot starve a long write. + - Pass `priority: 'read'` to let new readers join active readers even while a writer is queued, at the cost that continuous reads can starve a queued writer. + ### Changed - **BREAKING:** Remove envelope from migration framework ([#619](https://github.com/MetaMask/accounts/pull/619)) diff --git a/packages/keyring-sdk/src/atomic-keyring.test.ts b/packages/keyring-sdk/src/atomic-keyring.test.ts new file mode 100644 index 000000000..526749546 --- /dev/null +++ b/packages/keyring-sdk/src/atomic-keyring.test.ts @@ -0,0 +1,396 @@ +import type { AtomicUpdater } from './atomic-keyring'; +import { AtomicKeyring, isAtomicKeyring } from './atomic-keyring'; + +/** + * The states of the test keyring's state machine: an uninitialized keyring, + * a staged (mid-transition) keyring, and a ready keyring. + */ +type TestState = 'uninitialized' | 'staged' | 'ready'; + +/** + * The remote service the test keyring talks to, standing in for the + * multi-second ceremonies of a real atomic keyring. + */ +type RemoteService = { + performCeremony: () => Promise; + activate: () => Promise; +}; + +/** + * A promise that can be resolved or rejected from the outside. + */ +type Deferred = { + promise: Promise; + resolve: (value: Value) => void; + reject: (reason: Error) => void; +}; + +/** + * Create a deferred promise. + * + * @returns A deferred promise, with external `resolve` and `reject`. + */ +function createDeferred(): Deferred { + let resolve!: (value: Value) => void; + let reject!: (reason: Error) => void; + const promise = new Promise((_resolve, _reject) => { + resolve = _resolve; + reject = _reject; + }); + return { promise, resolve, reject }; +} + +/** + * Let pending promise continuations (microtasks) run, so queued waiters and + * callbacks get a chance to start. + */ +async function flush(): Promise { + for (let i = 0; i < 5; i++) { + await Promise.resolve(); + } +} + +/** + * Create a mock remote service whose operations resolve immediately. + * + * @returns The mock remote service. + */ +function createRemote(): jest.Mocked { + return { + performCeremony: jest.fn().mockResolvedValue(undefined), + activate: jest.fn().mockResolvedValue(undefined), + }; +} + +/** + * A minimal atomic keyring used to exercise the base class: a checkpointed + * creation flow (`uninitialized` → `staged` → `ready`) around a remote + * activation, plus a recovery path that advances a staged state without + * re-running the ceremony. + */ +class TestAtomicKeyring extends AtomicKeyring { + readonly #remote: RemoteService; + + #state: TestState = 'uninitialized'; + + constructor(options: { updater: AtomicUpdater; remote: RemoteService }) { + super({ updater: options.updater }); + this.#remote = options.remote; + } + + /** + * The current state. Deliberately lock-free, like the snapshot reads the + * controller performs itself. + * + * @returns The current state. + */ + get state(): TestState { + return this.#state; + } + + /** + * Checkpointed creation: a long-running ceremony, a staged commit, a + * remote activation, then promotion of the staged state. A failure after + * the first checkpoint leaves the staged state committed, so the next + * write operation can recover instead of re-running the ceremony. + */ + async createAccount(): Promise { + await this.withWriteLock(async () => { + await this.#remote.performCeremony(); + + await this.update(() => { + this.#state = 'staged'; + }); + + await this.#remote.activate(); + + await this.update(() => { + this.#state = 'ready'; + }); + }); + } + + /** + * Recovery: advance a staged state machine to ready, without re-running + * the ceremony. + */ + async recover(): Promise { + await this.withWriteLock(async () => { + if (this.#state !== 'staged') { + return; + } + await this.#remote.activate(); + await this.update(() => { + this.#state = 'ready'; + }); + }); + } + + /** + * Expose `update` for direct testing. + * + * @param action - The state mutation to commit. + * @returns A promise that resolves when the update completes. + */ + async commit(action: () => void | Promise): Promise { + return this.update(action); + } + + /** + * Expose `withReadLock` for direct testing. + * + * @param fn - The operation to run while holding the read lock. + * @returns Whatever `fn` resolves to. + */ + async read(fn: () => Promise): Promise { + return this.withReadLock(fn); + } + + /** + * Expose `withWriteLock` for direct testing. + * + * @param fn - The operation to run while holding the write lock. + * @returns Whatever `fn` resolves to. + */ + async write(fn: () => Promise): Promise { + return this.withWriteLock(fn); + } +} + +/** + * Create a test keyring with a recording updater, which stands in for the + * controller's vault-write callback and records the keyring state after + * every committed action. + * + * @param options - Setup options. + * @param options.updater - A custom updater, replacing the recording one. + * @returns The test keyring, its remote service, and the recorded writes. + */ +function setup(options: { updater?: AtomicUpdater } = {}): { + keyring: TestAtomicKeyring; + remote: jest.Mocked; + writes: TestState[]; + updater: AtomicUpdater; +} { + const remote = createRemote(); + const writes: TestState[] = []; + const updater: AtomicUpdater = + options.updater ?? + (async (keyring, action): Promise => { + await action(); + writes.push((keyring as TestAtomicKeyring).state); + }); + const keyring = new TestAtomicKeyring({ updater, remote }); + return { keyring, remote, writes, updater }; +} + +describe('AtomicKeyring', () => { + describe('constructor', () => { + it('brands the instance as an atomic keyring', () => { + const { keyring } = setup(); + + expect(isAtomicKeyring(keyring)).toBe(true); + }); + }); + + describe('update', () => { + it('passes the keyring instance and the action to the updater', async () => { + const updater = jest.fn( + async (_keyring: AtomicKeyring, action: () => void | Promise) => { + await action(); + }, + ); + const { keyring } = setup({ updater }); + const action = jest.fn(async () => undefined); + + await keyring.commit(action); + + expect(updater).toHaveBeenCalledWith(keyring, action); + expect(action).toHaveBeenCalledTimes(1); + }); + + it('runs the action before resolving', async () => { + const { keyring } = setup(); + let actionDone = false; + + await keyring.commit(async () => { + actionDone = true; + }); + + expect(actionDone).toBe(true); + }); + + it('rethrows errors thrown by the action', async () => { + const updater = jest.fn( + async (_keyring: AtomicKeyring, action: () => void | Promise) => { + await action(); + }, + ); + const { keyring } = setup({ updater }); + + await expect( + keyring.commit(async () => { + throw new Error('keyring bug'); + }), + ).rejects.toThrow('keyring bug'); + + // The updater ran and rethrew the action error. + expect(updater).toHaveBeenCalledTimes(1); + }); + + it('propagates errors thrown by the updater', async () => { + const updater: AtomicUpdater = async () => { + throw new Error('vault write failed'); + }; + const { keyring } = setup({ updater }); + + await expect(keyring.commit(async () => undefined)).rejects.toThrow( + 'vault write failed', + ); + }); + + it('rejects nested update calls with an explicit error', async () => { + const { keyring } = setup(); + + await expect( + keyring.commit(async () => { + await keyring.commit(async () => undefined); + }), + ).rejects.toThrow('Nested update() calls are not allowed'); + }); + + it('allows update calls after a nested call was rejected', async () => { + const { keyring, writes } = setup(); + + await expect( + keyring.commit(async () => { + await keyring.commit(async () => undefined); + }), + ).rejects.toThrow('Nested update() calls are not allowed'); + + // The re-entrancy guard was released, so a fresh update goes through. + await keyring.createAccount(); + expect(writes).toStrictEqual(['staged', 'ready']); + }); + }); + + describe('withReadLock and withWriteLock', () => { + it('runs read operations in parallel', async () => { + const { keyring } = setup(); + const gate = createDeferred(); + const started: string[] = []; + + const first = keyring.read(async () => { + started.push('first'); + await gate.promise; + }); + const second = keyring.read(async () => { + started.push('second'); + await gate.promise; + }); + await flush(); + + // Both reads started, even though neither has finished. + expect(started).toStrictEqual(['first', 'second']); + + gate.resolve(); + await Promise.all([first, second]); + }); + + it('excludes a write operation from a concurrent read operation', async () => { + const { keyring } = setup(); + const gate = createDeferred(); + const order: string[] = []; + + const read = keyring.read(async () => { + order.push('read'); + await gate.promise; + order.push('read:end'); + }); + const write = keyring.write(async () => { + order.push('write'); + }); + await flush(); + + // The write operation has not started while the read is active. + expect(order).toStrictEqual(['read']); + + gate.resolve(); + await Promise.all([read, write]); + + expect(order).toStrictEqual(['read', 'read:end', 'write']); + }); + }); + + describe('isAtomicKeyring', () => { + it('returns true for an AtomicKeyring subclass instance', () => { + const { keyring } = setup(); + + expect(isAtomicKeyring(keyring)).toBe(true); + }); + + it('returns false for values that are not atomic keyrings', () => { + expect(isAtomicKeyring(null)).toBe(false); + expect(isAtomicKeyring(undefined)).toBe(false); + expect(isAtomicKeyring({})).toBe(false); + expect(isAtomicKeyring({ type: 'Simple Keyring' })).toBe(false); + expect(isAtomicKeyring(() => undefined)).toBe(false); + expect(isAtomicKeyring('keyring')).toBe(false); + }); + }); + + describe('checkpointed commit flow', () => { + it('persists each checkpoint as it is committed', async () => { + const { keyring, writes, remote } = setup(); + + await keyring.createAccount(); + + expect(keyring.state).toBe('ready'); + expect(writes).toStrictEqual(['staged', 'ready']); + expect(remote.performCeremony).toHaveBeenCalledTimes(1); + expect(remote.activate).toHaveBeenCalledTimes(1); + }); + + it('keeps committed checkpoints when a later step fails', async () => { + const { keyring, writes, remote } = setup(); + remote.activate.mockRejectedValueOnce( + new Error('remote activation failed'), + ); + + await expect(keyring.createAccount()).rejects.toThrow( + 'remote activation failed', + ); + + // The staged checkpoint survived the failure — there is no rollback. + expect(keyring.state).toBe('staged'); + expect(writes).toStrictEqual(['staged']); + }); + + it('recovers a staged state on the next write operation without re-running the ceremony', async () => { + const { keyring, writes, remote } = setup(); + remote.activate.mockRejectedValueOnce( + new Error('remote activation failed'), + ); + await expect(keyring.createAccount()).rejects.toThrow( + 'remote activation failed', + ); + + await keyring.recover(); + + expect(keyring.state).toBe('ready'); + expect(writes).toStrictEqual(['staged', 'ready']); + expect(remote.performCeremony).toHaveBeenCalledTimes(1); + expect(remote.activate).toHaveBeenCalledTimes(2); + }); + + it('leaves non-staged states untouched during recovery', async () => { + const { keyring, writes, remote } = setup(); + + await keyring.recover(); + + expect(keyring.state).toBe('uninitialized'); + expect(writes).toStrictEqual([]); + expect(remote.activate).not.toHaveBeenCalled(); + }); + }); +}); diff --git a/packages/keyring-sdk/src/atomic-keyring.ts b/packages/keyring-sdk/src/atomic-keyring.ts new file mode 100644 index 000000000..5eeeae068 --- /dev/null +++ b/packages/keyring-sdk/src/atomic-keyring.ts @@ -0,0 +1,139 @@ +import { ReadWriteLock } from './read-write-lock'; + +/** + * A vault-write callback injected into an `AtomicKeyring` at construction, + * via the keyring builder context. + * + * Stateless: the calling keyring is the first parameter, so a single + * function serves every keyring. Acquires the controller lock, verifies the + * keyring is still registered, runs `action()`, writes the vault if and only + * if the serialized state changed, then releases — a sub-millisecond commit + * window. + */ +export type AtomicUpdater = ( + keyring: AtomicKeyring, + action: () => void | Promise, +) => Promise; + +/** + * Brand used to mark `AtomicKeyring` instances. Registered on the global + * symbol registry, so the brand survives duplicate copies of this package — + * where `instanceof` would fail — and `isAtomicKeyring` keeps working. + */ +const ATOMIC_KEYRING_BRAND = Symbol.for('@metamask/keyring-sdk/atomic-keyring'); + +/** + * Base class for keyrings that serialize their own operations and call back + * into the controller only to persist state. + * + * Long-running work runs inside the lock helpers, entirely outside the + * controller mutex; state is committed through `update`. + * + * Subclasses implementing an `init` lifecycle hook must keep it cheap — + * in-memory work only. The controller invokes `init` under its lock, at + * construction and vault-restore time; an `init` that performs long work, + * or calls `update` (which needs that same, non-reentrant lock), would + * block or deadlock the controller. Long work belongs in write operations, + * which the controller dispatches without holding its lock. + * + * Snapshot reads the controller performs itself — `getAccounts` and + * `serialize` — must stay lock-free in subclasses: the controller calls + * them while holding its own mutex, and the updater calls `serialize` while + * the keyring is inside a write operation. They must return immediately + * from the current state. + */ +export abstract class AtomicKeyring { + readonly #updater: AtomicUpdater; + + readonly #rwLock: ReadWriteLock = new ReadWriteLock(); + + #updating: boolean = false; + + /** + * The updater is required at construction — no unbound state, ever. + * + * @param options - Constructor options. + * @param options.updater - The controller's vault-write callback, which + * this keyring invokes through `update` to commit state. + */ + constructor(options: { updater: AtomicUpdater }) { + this.#updater = options.updater; + (this as Record)[ATOMIC_KEYRING_BRAND] = true; + } + + /** + * Commit a state mutation on behalf of this keyring. Injects `this` into + * the updater, so the controller can verify registration. + * + * Actions must be pure assignments of pre-computed next state: compute + * and validate before calling `update`. A throwing action is a keyring + * bug, and is rethrown as such. Nested `update` calls are rejected with + * an error, never a deadlock. + * + * Call `update` only from within a write-locked operation, so commits + * cannot interleave with other operations on this keyring. + * + * @param action - The state mutation to commit. Must be a pure assignment + * of pre-computed next state. + * @throws If called while another `update` call is in progress, or if the + * action or updater throws. + */ + protected async update(action: () => void | Promise): Promise { + if (this.#updating) { + throw new Error( + 'AtomicKeyring - Nested update() calls are not allowed. The update() action must not call update() again.', + ); + } + + this.#updating = true; + try { + await this.#updater(this, action); + } finally { + this.#updating = false; + } + } + + /** + * Run a read operation — parallel with other readers, queued behind + * writers. Only for operations that can tolerate waiting behind a write + * section. + * + * @param fn - The operation to run while holding the read lock. + * @returns Whatever `fn` resolves to. + */ + protected async withReadLock( + fn: () => Promise, + ): Promise { + return this.#rwLock.withReadLock(fn); + } + + /** + * Run a write operation — exclusive with readers and other writers. + * Long-running network work belongs here, outside the controller mutex. + * + * @param fn - The operation to run while holding the write lock. + * @returns Whatever `fn` resolves to. + */ + protected async withWriteLock( + fn: () => Promise, + ): Promise { + return this.#rwLock.withWriteLock(fn); + } +} + +/** + * Whether the value is an `AtomicKeyring`. + * + * Brand-based — a global-registry symbol set by the constructor — so it + * works across duplicate copies of the package, where `instanceof` fails. + * + * @param value - The value to check. + * @returns Whether the value is an `AtomicKeyring` instance. + */ +export function isAtomicKeyring(value: unknown): value is AtomicKeyring { + return ( + typeof value === 'object' && + value !== null && + (value as Record)[ATOMIC_KEYRING_BRAND] === true + ); +} diff --git a/packages/keyring-sdk/src/index.ts b/packages/keyring-sdk/src/index.ts index 1390bd72f..b62b1f0a0 100644 --- a/packages/keyring-sdk/src/index.ts +++ b/packages/keyring-sdk/src/index.ts @@ -1,5 +1,7 @@ +export * from './atomic-keyring'; export * from './keyring-account-registry'; export * from './migration'; export * from './mnemonic'; export * from './entropy'; +export * from './read-write-lock'; export * from './eth'; diff --git a/packages/keyring-sdk/src/read-write-lock.test.ts b/packages/keyring-sdk/src/read-write-lock.test.ts new file mode 100644 index 000000000..2bf6584ac --- /dev/null +++ b/packages/keyring-sdk/src/read-write-lock.test.ts @@ -0,0 +1,371 @@ +import { ReadWriteLock } from './read-write-lock'; + +/** + * A promise that can be resolved or rejected from the outside. + */ +type Deferred = { + promise: Promise; + resolve: (value: Value) => void; + reject: (reason: Error) => void; +}; + +/** + * Create a deferred promise. + * + * @returns A deferred promise, with external `resolve` and `reject`. + */ +function createDeferred(): Deferred { + let resolve!: (value: Value) => void; + let reject!: (reason: Error) => void; + const promise = new Promise((_resolve, _reject) => { + resolve = _resolve; + reject = _reject; + }); + return { promise, resolve, reject }; +} + +/** + * Let pending promise continuations (microtasks) run, so queued waiters and + * callbacks get a chance to start. + */ +async function flush(): Promise { + for (let i = 0; i < 5; i++) { + await Promise.resolve(); + } +} + +describe('ReadWriteLock', () => { + describe('withReadLock', () => { + it('returns the result of the callback', async () => { + const lock = new ReadWriteLock(); + + expect(await lock.withReadLock(async () => 'result')).toBe('result'); + }); + + it('allows multiple readers to hold the lock simultaneously', async () => { + const lock = new ReadWriteLock(); + const gate = createDeferred(); + const started: number[] = []; + + const first = lock.withReadLock(async () => { + started.push(1); + await gate.promise; + }); + const second = lock.withReadLock(async () => { + started.push(2); + await gate.promise; + }); + const third = lock.withReadLock(async () => { + started.push(3); + await gate.promise; + }); + await flush(); + + // All three readers started, even though none has finished: readers + // hold the lock in parallel. + expect(started).toStrictEqual([1, 2, 3]); + + gate.resolve(); + await Promise.all([first, second, third]); + }); + + it('releases the lock when the callback throws synchronously', async () => { + const lock = new ReadWriteLock(); + + await expect( + lock.withReadLock(async (): Promise => { + throw new Error('fail'); + }), + ).rejects.toThrow('fail'); + + // The lock is usable again after the failure. + expect(await lock.withReadLock(async () => 'ok')).toBe('ok'); + }); + + it('releases the lock when the callback rejects', async () => { + const lock = new ReadWriteLock(); + + await expect( + lock.withReadLock(async () => { + throw new Error('fail'); + }), + ).rejects.toThrow('fail'); + + expect(await lock.withReadLock(async () => 'ok')).toBe('ok'); + }); + }); + + describe('withWriteLock', () => { + it('returns the result of the callback', async () => { + const lock = new ReadWriteLock(); + + expect(await lock.withWriteLock(async () => 'result')).toBe('result'); + }); + + it('waits for active readers to finish before granting the write lock', async () => { + const lock = new ReadWriteLock(); + const gate = createDeferred(); + const order: string[] = []; + + const read = lock.withReadLock(async () => { + order.push('read:start'); + await gate.promise; + order.push('read:end'); + }); + const write = lock.withWriteLock(async () => { + order.push('write:start'); + }); + await flush(); + + // The write operation has not started while the read is active. + expect(order).toStrictEqual(['read:start']); + + gate.resolve(); + await Promise.all([read, write]); + + expect(order).toStrictEqual(['read:start', 'read:end', 'write:start']); + }); + + it('excludes other writers while a writer holds the lock', async () => { + const lock = new ReadWriteLock(); + const gate = createDeferred(); + const order: string[] = []; + + const first = lock.withWriteLock(async () => { + order.push('first:start'); + await gate.promise; + order.push('first:end'); + }); + const second = lock.withWriteLock(async () => { + order.push('second:start'); + }); + await flush(); + + expect(order).toStrictEqual(['first:start']); + + gate.resolve(); + await Promise.all([first, second]); + + expect(order).toStrictEqual(['first:start', 'first:end', 'second:start']); + }); + + it('releases the lock when the callback throws synchronously', async () => { + const lock = new ReadWriteLock(); + + await expect( + lock.withWriteLock(async (): Promise => { + throw new Error('fail'); + }), + ).rejects.toThrow('fail'); + + expect(await lock.withWriteLock(async () => 'ok')).toBe('ok'); + }); + + it('releases the lock when the callback rejects', async () => { + const lock = new ReadWriteLock(); + + await expect( + lock.withWriteLock(async () => { + throw new Error('fail'); + }), + ).rejects.toThrow('fail'); + + expect(await lock.withWriteLock(async () => 'ok')).toBe('ok'); + }); + }); + + describe('grant ordering', () => { + it('grants waiters in first-in-first-out order', async () => { + const lock = new ReadWriteLock(); + const gate = createDeferred(); + const order: string[] = []; + + const firstWrite = lock.withWriteLock(async () => { + await gate.promise; + }); + + // Queued in arrival order: read, write, read. + const queuedRead = lock.withReadLock(async () => { + order.push('queued-read'); + }); + const queuedWrite = lock.withWriteLock(async () => { + order.push('queued-write'); + }); + const queuedRead2 = lock.withReadLock(async () => { + order.push('queued-read-2'); + }); + await flush(); + + expect(order).toStrictEqual([]); + + gate.resolve(); + await Promise.all([firstWrite, queuedRead, queuedWrite, queuedRead2]); + + expect(order).toStrictEqual([ + 'queued-read', + 'queued-write', + 'queued-read-2', + ]); + }); + + it('queues new readers behind a pending writer', async () => { + const lock = new ReadWriteLock(); + const readGate = createDeferred(); + const writeGate = createDeferred(); + const started: Record = { + read: false, + write: false, + lateRead: false, + }; + const order: string[] = []; + + const read = lock.withReadLock(async () => { + started.read = true; + order.push('read'); + await readGate.promise; + }); + + // A writer queues up while the reader is active... + const write = lock.withWriteLock(async () => { + started.write = true; + order.push('write'); + await writeGate.promise; + }); + + // ...and a later reader must queue behind the pending writer, even + // though a reader is active. + const lateRead = lock.withReadLock(async () => { + started.lateRead = true; + order.push('late-read'); + }); + await flush(); + + expect(started).toStrictEqual({ + read: true, + write: false, + lateRead: false, + }); + + readGate.resolve(); + await flush(); + + // The pending writer was granted next, not the later reader. + expect(started).toStrictEqual({ + read: true, + write: true, + lateRead: false, + }); + + writeGate.resolve(); + await Promise.all([read, write, lateRead]); + + expect(order).toStrictEqual(['read', 'write', 'late-read']); + }); + + it('grants consecutive queued readers together', async () => { + const lock = new ReadWriteLock(); + const writeGate = createDeferred(); + const readGate = createDeferred(); + const order: string[] = []; + + const write = lock.withWriteLock(async () => { + await writeGate.promise; + }); + + const firstRead = lock.withReadLock(async () => { + order.push('read-1'); + await readGate.promise; + }); + const secondRead = lock.withReadLock(async () => { + order.push('read-2'); + }); + const thirdRead = lock.withReadLock(async () => { + order.push('read-3'); + }); + await flush(); + + expect(order).toStrictEqual([]); + + writeGate.resolve(); + await flush(); + + // All queued readers were granted together: the first is active and + // holding the others open, before any of them finished. + expect(order).toStrictEqual(['read-1', 'read-2', 'read-3']); + + readGate.resolve(); + await Promise.all([write, firstRead, secondRead, thirdRead]); + }); + }); + + describe('priority option', () => { + it('lets new readers start while a writer is pending, with read priority', async () => { + const lock = new ReadWriteLock({ priority: 'read' }); + const readGate = createDeferred(); + const writeGate = createDeferred(); + const order: string[] = []; + + const read = lock.withReadLock(async () => { + order.push('read:start'); + await readGate.promise; + order.push('read:end'); + }); + + // A writer queues up while the reader is active... + const write = lock.withWriteLock(async () => { + order.push('write'); + await writeGate.promise; + }); + + // ...and a later reader starts immediately, jumping the merely + // pending writer. + const lateRead = lock.withReadLock(async () => { + order.push('late-read'); + }); + await flush(); + + expect(order).toStrictEqual(['read:start', 'late-read']); + + // The pending writer still waits for the jumping reader to finish. + readGate.resolve(); + writeGate.resolve(); + await Promise.all([read, write, lateRead]); + + expect(order).toStrictEqual([ + 'read:start', + 'late-read', + 'read:end', + 'write', + ]); + }); + + it('still queues readers behind an active writer, with read priority', async () => { + const lock = new ReadWriteLock({ priority: 'read' }); + const writeGate = createDeferred(); + const order: string[] = []; + + const write = lock.withWriteLock(async () => { + order.push('write'); + await writeGate.promise; + order.push('write:end'); + }); + await flush(); + + expect(order).toStrictEqual(['write']); + + const read = lock.withReadLock(async () => { + order.push('read'); + }); + await flush(); + + // The writer is active — read priority does not jump an active + // writer, only a queued one. + expect(order).toStrictEqual(['write']); + + writeGate.resolve(); + await Promise.all([write, read]); + + expect(order).toStrictEqual(['write', 'write:end', 'read']); + }); + }); +}); diff --git a/packages/keyring-sdk/src/read-write-lock.ts b/packages/keyring-sdk/src/read-write-lock.ts new file mode 100644 index 000000000..16490b3de --- /dev/null +++ b/packages/keyring-sdk/src/read-write-lock.ts @@ -0,0 +1,237 @@ +/** + * A mode in which the lock can be held: shared (`'read'`) or exclusive + * (`'write'`). + */ +type LockMode = 'read' | 'write'; + +/** + * A queued lock-acquisition request. + */ +type LockWaiter = { + /** The mode being acquired. */ + mode: LockMode; + /** Resolves once the lock has been granted. */ + resolve: () => void; +}; + +/** + * Which mode a {@link ReadWriteLock} prioritizes when both are contended. + */ +export type ReadWriteLockPriority = 'read' | 'write'; + +/** + * Options for constructing a {@link ReadWriteLock}. + */ +export type ReadWriteLockOptions = { + /** + * Which mode takes priority when both are contended. Defaults to + * `'write'`. + * + * - `'write'`: a pending writer blocks new readers, so a steady stream of + * reads cannot starve a long write. + * - `'read'`: new readers join active readers even while a writer is + * queued, for better read latency — at the cost that a continuous stream + * of reads can starve a queued writer. + */ + priority?: ReadWriteLockPriority; +}; + +/** + * A read-write lock: multiple readers may hold the lock simultaneously, + * while writers hold it exclusively. + * + * Grant order is first-in-first-out, and the lock is write-preferring by + * default: a pending writer blocks new readers, and consecutive queued + * readers are granted together once the lock turns over to them. Set the + * `priority` option to `'read'` to let new readers join active readers even + * while a writer is queued — better read latency, at the cost that + * continuous reads can starve a queued writer. The lock is always released + * when the callback settles, whether it resolves or throws. + */ +export class ReadWriteLock { + #readers: number = 0; + + #writer: boolean = false; + + readonly #waiters: LockWaiter[] = []; + + readonly #priority: ReadWriteLockPriority; + + /** + * @param options - Lock options. + * @param options.priority - Which mode takes priority when both are + * contended. Defaults to `'write'`. + */ + constructor(options: ReadWriteLockOptions = {}) { + this.#priority = options.priority ?? 'write'; + } + + /** + * Run a read operation — parallel with other readers, exclusive with + * writers. New readers queue behind any pending writer. + * + * @param fn - The operation to run while holding the read lock. + * @returns Whatever `fn` resolves to. + */ + async withReadLock(fn: () => Promise): Promise { + return this.#withLock('read', fn); + } + + /** + * Run a write operation — exclusive with readers and other writers. + * + * @param fn - The operation to run while holding the write lock. + * @returns Whatever `fn` resolves to. + */ + async withWriteLock(fn: () => Promise): Promise { + return this.#withLock('write', fn); + } + + /** + * Acquire the lock in the given mode, run the operation, and release the + * lock — including on a throw, so the lock can never leak. + * + * @param mode - The mode in which to hold the lock. + * @param fn - The operation to run while holding the lock. + * @returns Whatever `fn` resolves to. + */ + async #withLock( + mode: LockMode, + fn: () => Promise, + ): Promise { + await this.#acquire(mode); + try { + return await fn(); + } finally { + this.#release(mode); + } + } + + /** + * Whether any acquisition request is queued, waiting for the lock to turn + * over to it. A queued waiter always blocks newcomers, which keeps the + * grant order first-in-first-out. + * + * @returns Whether any acquisition request is queued. + */ + #hasWaiters(): boolean { + return this.#waiters.length > 0; + } + + /** + * Whether the lock is currently held for reading. Readers hold the lock + * in parallel, so this does not mean an exclusive hold. + * + * @returns Whether any reader is active. + */ + #isReading(): boolean { + return this.#readers > 0; + } + + /** + * Whether the lock is currently held for writing. A writer holds the lock + * exclusively: while writing, no other holder — reader or writer — can + * be active. + * + * @returns Whether a writer is active. + */ + #isWriting(): boolean { + return this.#writer; + } + + /** + * Acquire the lock in the given mode, queuing when it cannot be granted + * immediately. + * + * @param mode - The mode to acquire. + * @returns A promise that resolves once the lock is held. + */ + async #acquire(mode: LockMode): Promise { + return new Promise((resolve) => { + let acquire = true; + if (this.#isWriting()) { + // Writing is exclusive, so we cannot acquire the lock if a writer is active. + acquire = false; + } else if (mode === 'write' && this.#isReading()) { + // Writing is exclusive with readers, so we cannot acquire the lock if any reader is active. + acquire = false; + } else if ( + mode === 'read' && + // We might have pending write requests in the queue. + this.#hasWaiters() && + this.#priority === 'write' + ) { + // We cannot acquire the lock for reading if there are waiters, unless priority is 'read'. + acquire = false; + } + + if (acquire) { + this.#grant(mode); + resolve(); + } else { + // We cannot grant the lock immediately, so we queue the waiter. + this.#waiters.push({ mode, resolve }); + } + }); + } + + /** + * Grant the lock in the given mode. Only called when the grant is legal. + * + * @param mode - The mode to grant. + */ + #grant(mode: LockMode): void { + if (mode === 'read') { + this.#readers += 1; + } else { + this.#writer = true; + } + } + + /** + * Release the lock in the given mode, then grant queued waiters in order. + * + * @param mode - The mode being released. + */ + #release(mode: LockMode): void { + if (mode === 'read') { + this.#readers -= 1; + } else { + this.#writer = false; + } + this.#drain(); + } + + /** + * Grant queued waiters, in order, while the front of the queue can be + * granted. Consecutive queued readers are granted together; a granted + * writer blocks everything behind it. Unlike a newcomer, the front + * waiter is never blocked by queued waiters — everyone behind it + * arrived later — so only active holders block a grant here. + */ + #drain(): void { + let waiter: LockWaiter | undefined = this.#waiters.shift(); + while (waiter !== undefined) { + let acquire = true; + + if (this.#isWriting()) { + // Writing is exclusive, so the waiter cannot acquire the lock if a writer is active. + acquire = false; + } else if (waiter.mode === 'write' && this.#isReading()) { + // Writing is exclusive with readers, so the waiter cannot acquire the lock if any reader is active. + acquire = false; + } + + if (acquire) { + // Grant the front waiter and continue with the next one, so consecutive queued readers are granted together. + this.#grant(waiter.mode); + waiter.resolve(); + waiter = this.#waiters.shift(); + } else { + // The waiter is blocked by an active holder; put it back at the front and stop. + this.#waiters.unshift(waiter); + return; + } + } + } +}