diff --git a/packages/ocap-kernel/CHANGELOG.md b/packages/ocap-kernel/CHANGELOG.md index c2aab043aa..119f41fddd 100644 --- a/packages/ocap-kernel/CHANGELOG.md +++ b/packages/ocap-kernel/CHANGELOG.md @@ -49,12 +49,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Changed +- **BREAKING:** A peer's incarnation change is carried out by the run loop, in a crank of its own, instead of in a `peerIncarnation_*` savepoint nested inside whichever crank was open. The handshake is still answered immediately, from what the store already says, because the transport cannot wait for a crank to finish; only the writes it implies are queued ([#1104](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1104)) + - Queued behind anything that peer has already sent, so a message from the incarnation that is ending is recorded against it and one from the incarnation that is starting is not. Ordering is what keeps them apart, which is why nothing has to be discarded + - Rejecting the promises the restarted remote was deciding, and resetting its in-memory state, wait for that crank to commit — neither is reversible by a rollback - **BREAKING:** A message from a remote peer is now a `remoteInbound` run queue item, taken delivery of in a crank of its own, instead of being processed on arrival inside a `receive_*` savepoint of its own. That savepoint nested inside whichever crank happened to be open, which — depending on the order — either deferred that crank's commit or committed a half-finished message. `RemoteHandle.handleRemoteMessage` is replaced by `receiveFromPeer`, which acknowledges and validates where the message arrives and queues the rest, and `deliverInbound`, which the run loop calls ([#1103](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1103)) - Duplicate detection now reads the store rather than memory, because several messages from one peer can be waiting their turn at once and the in-memory sequence number only catches up as each crank commits - A message whose sequence number is missing or not a positive integer is refused where it arrives. Queued, it would have thrown in its crank, killed the run loop, and been rolled back onto the queue to kill the next boot too — and left `highestReceivedSeq` reading back as `NaN`, which silently disables duplicate detection and acknowledgement for that peer forever - A message for a remote that has gone, or one whose payload the crank cannot make sense of — an unknown method, a reference that does not resolve, a reply to a redemption that has already timed out — is logged and dropped. Thrown, it would have killed the run loop and been rolled back onto the queue to kill the next boot too, which a peer could provoke at will - An arrival waits in memory rather than on the run queue, because it arrives while a crank is open and writing it there would put it inside that crank's transaction for an abort to swallow. Nothing is lost: the peer is not acknowledged until the crank that takes the message commits, so an arrival the kernel forgets is one the peer sends again - - Arrivals from a peer that has restarted are discarded, since recording an old incarnation's sequence number against the new one would make the new incarnation's first message look like a duplicate - **BREAKING:** `Kernel.make`'s `ioChannelFactory` option is now `ioListenerFactory`, and the exported `IOChannelFactory` type is replaced by `IOListener` and `IOListenerFactory`. A cluster config's `io` entries now create listeners; vats call `accept()` to obtain a channel instead of reading and writing the endowment directly ([#1007](https://github.com/MetaMask/ocap-kernel/pull/1007)) - Attribute a failed subcluster vat launch to the specific vat by kernel id and `ClusterConfig` name (e.g. `Failed to launch vat v3 (bob)`), preserving the original error as the `cause` ([#975](https://github.com/MetaMask/ocap-kernel/pull/975)) - **BREAKING:** Remove `VatConfig.platformConfig.fetch` — migrate to `globals: ['fetch', ...]` + `network.allowedHosts` ([#942](https://github.com/MetaMask/ocap-kernel/pull/942)) diff --git a/packages/ocap-kernel/src/Kernel.test.ts b/packages/ocap-kernel/src/Kernel.test.ts index d85cb02234..3a0c660f15 100644 --- a/packages/ocap-kernel/src/Kernel.test.ts +++ b/packages/ocap-kernel/src/Kernel.test.ts @@ -31,10 +31,14 @@ const mocks = vi.hoisted(() => { #rejectRunLoop: ((error: Error) => void) | undefined; + /** The callback the kernel gave the run loop, for tests to drive. */ + deliver: ((item: unknown) => Promise) | undefined; + // Like the real run loop, this settles only if the kernel dies. run = vi.fn( - async () => + async (deliver: (item: unknown) => Promise) => new Promise((_resolve, reject) => { + this.deliver = deliver; this.#rejectRunLoop = reject; }), ); @@ -91,6 +95,8 @@ const mocks = vi.hoisted(() => { setMessageHandler = vi.fn(); + applyIncarnationChange = vi.fn().mockResolvedValue({}); + initIdentity = vi.fn().mockResolvedValue(undefined); initRemoteComms = vi.fn().mockResolvedValue(undefined); @@ -249,6 +255,44 @@ describe('Kernel', () => { }); }); + describe('the run loop callback', () => { + const makeRunningKernel = async (): Promise => { + await Kernel.make(mockPlatformServices, mockKernelDatabase); + }; + + it('carries out a peer incarnation change itself', async () => { + await makeRunningKernel(); + const item = { + type: 'peerIncarnation', + peerId: 'peer-1', + incarnation: 'incarnation-B', + }; + + // Not a delivery to an endpoint, so the router never sees it — and it + // would throw on an item type it does not know. + await mocks.KernelQueue.lastInstance.deliver?.(item); + + expect( + mocks.RemoteManager.lastInstance.applyIncarnationChange, + ).toHaveBeenCalledWith(item); + }); + + it('does not die of an incarnation change it cannot record', async () => { + await makeRunningKernel(); + mocks.RemoteManager.lastInstance.applyIncarnationChange.mockRejectedValue( + new Error('the store is gone'), + ); + + expect( + await mocks.KernelQueue.lastInstance.deliver?.({ + type: 'peerIncarnation', + peerId: 'peer-1', + incarnation: 'incarnation-B', + }), + ).toStrictEqual({ abort: true }); + }); + }); + describe('queueMessage()', () => { it('enqueues a message and returns the result', async () => { const kernel = await Kernel.make( diff --git a/packages/ocap-kernel/src/Kernel.ts b/packages/ocap-kernel/src/Kernel.ts index c041d060be..77655b94f9 100644 --- a/packages/ocap-kernel/src/Kernel.ts +++ b/packages/ocap-kernel/src/Kernel.ts @@ -343,7 +343,27 @@ export class Kernel { // Start the kernel queue processing (non-blocking) // This runs for the entire lifetime of the kernel this.#kernelQueue - .run(this.#kernelRouter.deliver.bind(this.#kernelRouter)) + .run(async (item) => { + // Not a delivery to an endpoint, so not the router's to route: a peer + // restarting is the kernel's own bookkeeping, and for a peer with no + // live remote there is no endpoint to route it to at all. + if (item.type !== 'peerIncarnation') { + return await this.#kernelRouter.deliver(item); + } + try { + return await this.#remoteManager.applyIncarnationChange(item); + } catch (error) { + // A peer restarting is the peer's doing, and the crank has rolled + // back whatever this started. The handshake that queued it is long + // answered, and the peer will hand us its incarnation again on the + // next one. + this.#logger.error( + `Could not record the incarnation change for peer ${item.peerId}:`, + error, + ); + return { abort: true }; + } + }) .catch((error) => this.#handleRunLoopFailure(error)); // Launch new system subclusters (requires queue to be running) diff --git a/packages/ocap-kernel/src/KernelQueue.test.ts b/packages/ocap-kernel/src/KernelQueue.test.ts index a4e121dc89..fe8b779014 100644 --- a/packages/ocap-kernel/src/KernelQueue.test.ts +++ b/packages/ocap-kernel/src/KernelQueue.test.ts @@ -784,6 +784,60 @@ describe('KernelQueue', () => { }); }); + // A parked run loop with work waiting is a permanent wedge, and an arrival + // is the one kind of work that does not go through `#enqueueRun`. + it.each([ + { + what: 'a message arrives', + accept: (queue: KernelQueue) => + queue.acceptRemoteInbound('r1' as RemoteId, '{"seq":1}'), + }, + { + what: 'a peer restarts', + accept: (queue: KernelQueue) => + queue.acceptPeerIncarnation('peer-1', 'incarnation-B'), + }, + ])('wakes a parked run loop when $what', async ({ accept }) => { + (kernelStore.runQueueLength as unknown as MockInstance).mockReturnValue( + 0, + ); + let turns = 0; + (kernelStore.endCrank as unknown as MockInstance).mockImplementation( + () => { + turns += 1; + if (turns === 1) { + // The loop found nothing and has armed its wake. + accept(kernelQueue); + } else if (turns > 2) { + throw new Error(STOP_RUN_LOOP); + } + }, + ); + const deliver = vi.fn().mockResolvedValue({}); + + await expect(kernelQueue.run(deliver)).rejects.toThrow(STOP_RUN_LOOP); + + expect(mockPromiseKit.resolve).toHaveBeenCalled(); + expect(deliver).toHaveBeenCalledOnce(); + }); + + it('refuses an incarnation change once the run loop has died', async () => { + const deliver = vi.fn().mockRejectedValue(new Error('dead')); + ( + kernelStore.runQueueLength as unknown as MockInstance + ).mockReturnValueOnce(1); + (kernelStore.dequeueRun as unknown as MockInstance).mockReturnValue({ + type: 'send', + target: 'ko1', + message: {} as KernelMessage, + }); + await expect(kernelQueue.run(deliver)).rejects.toThrow('dead'); + + expect(() => + kernelQueue.acceptPeerIncarnation('peer-1', 'incarnation-B'), + ).toThrow('Kernel run loop died'); + }); + it('refuses an arrival once the run loop has died', async () => { const deliver = vi.fn().mockRejectedValue(new Error('dead')); ( @@ -801,23 +855,28 @@ describe('KernelQueue', () => { ).toThrow('Kernel run loop died'); }); - it('discards what a remote sent before its incarnation changed', async () => { + it("keeps a peer's messages and its incarnation change in the order they arrived", async () => { kernelQueue.acceptRemoteInbound('r1' as RemoteId, '{"seq":47}'); - kernelQueue.acceptRemoteInbound('r2' as RemoteId, '{"seq":3}'); - - kernelQueue.discardRemoteInbound('r1' as RemoteId); + kernelQueue.acceptPeerIncarnation('peer-1', 'incarnation-B'); + kernelQueue.acceptRemoteInbound('r1' as RemoteId, '{"seq":1}'); const delivered: RunQueueItem[] = []; const deliver = vi.fn().mockImplementation((item: RunQueueItem) => { delivered.push(item); - throw new Error(STOP_RUN_LOOP); + if (delivered.length === 3) { + throw new Error(STOP_RUN_LOOP); + } + return {}; }); await expect(kernelQueue.run(deliver)).rejects.toThrow(STOP_RUN_LOOP); - // Kept, the old incarnation's seq would be recorded against the new one - // and the new one's first message discarded as a duplicate. - expect(delivered).toStrictEqual([ - { type: 'remoteInbound', remoteId: 'r2', message: '{"seq":3}' }, + // Order is what keeps the incarnations apart: the message ahead of the + // change belongs to the incarnation that ended and is recorded against + // it, the one behind it to the incarnation that started. + expect(delivered.map((item) => item.type)).toStrictEqual([ + 'remoteInbound', + 'peerIncarnation', + 'remoteInbound', ]); }); }); diff --git a/packages/ocap-kernel/src/KernelQueue.ts b/packages/ocap-kernel/src/KernelQueue.ts index 68ae020d64..f6bb18a075 100644 --- a/packages/ocap-kernel/src/KernelQueue.ts +++ b/packages/ocap-kernel/src/KernelQueue.ts @@ -14,6 +14,7 @@ import type { RemoteId, RunLoopStatus, RunQueueItem, + RunQueueItemPeerIncarnation, RunQueueItemRemoteInbound, RunQueueItemNotify, RunQueueItemSend, @@ -86,8 +87,11 @@ export class KernelQueue { */ #runLoopState: RunLoopState = { state: 'idle' }; - /** Messages from peers, waiting for a crank to take delivery of them. */ - #arrivedFromRemotes: RunQueueItemRemoteInbound[] = []; + /** What peers have sent, waiting for a crank to take delivery of it. */ + readonly #arrivedFromRemotes: ( + | RunQueueItemRemoteInbound + | RunQueueItemPeerIncarnation + )[] = []; /** * Construct a new KernelQueue instance. @@ -589,25 +593,25 @@ export class KernelQueue { } /** - * Forget what a remote sent before an incarnation change, none of which the - * peer that sent it is still waiting on. + * Accept a peer's incarnation change, for the run loop to carry out in a + * crank of its own. * - * Left queued, a message from the old incarnation would record its sequence - * number against the new one, and the new incarnation's first message would - * then be discarded as a duplicate. + * Held with the arrivals, and behind any this peer has already sent: the + * messages ahead of it belong to the incarnation that is ending and are its + * to account for, and the ones behind it to the incarnation that is + * starting. * - * Reaches only what is still waiting. An arrival already handed to a crank - * has been shifted off this list, and a restart detected while that crank is - * suspended still records the old incarnation's sequence number — the - * handshake runs on the transport's flow, with nothing serializing it - * against an open crank. - * - * @param remoteId - The remote whose arrivals to discard. + * @param peerId - The peer that restarted. + * @param incarnation - The incarnation it now reports. */ - discardRemoteInbound(remoteId: RemoteId): void { - this.#arrivedFromRemotes = this.#arrivedFromRemotes.filter( - (item) => item.remoteId !== remoteId, - ); + acceptPeerIncarnation(peerId: string, incarnation: string): void { + this.assertRunLoopAlive('accept a peer incarnation change'); + this.#arrivedFromRemotes.push({ + type: 'peerIncarnation', + peerId, + incarnation, + }); + this.#wakeTheRunLoop(); } /** diff --git a/packages/ocap-kernel/src/KernelRouter.ts b/packages/ocap-kernel/src/KernelRouter.ts index 66279d1c42..5c00515d14 100644 --- a/packages/ocap-kernel/src/KernelRouter.ts +++ b/packages/ocap-kernel/src/KernelRouter.ts @@ -17,6 +17,7 @@ import type { RunQueueItem, RunQueueItemSend, RemoteEndpointHandle, + RunQueueItemPeerIncarnation, RunQueueItemBringOutYourDead, RunQueueItemRemoteInbound, RunQueueItemNotify, @@ -25,6 +26,13 @@ import type { } from './types.ts'; import { assert, Fail } from './utils/assert.ts'; +/** + * Every run queue item the router routes. A peer's incarnation change is not a + * delivery to an endpoint and is carried out by the kernel instead, so leaving + * it out is what keeps the exhaustiveness check below honest. + */ +type RoutedRunQueueItem = Exclude; + type MessageRoute = { endpointId?: EndpointId | 'kernel'; target: KRef; @@ -93,7 +101,7 @@ export class KernelRouter { * @param item - The message/notification to deliver. * @returns The crank outcome. */ - async deliver(item: RunQueueItem): Promise { + async deliver(item: RoutedRunQueueItem): Promise { switch (item.type) { case 'send': return await this.#deliverSend(item); diff --git a/packages/ocap-kernel/src/remotes/kernel/RemoteHandle.ts b/packages/ocap-kernel/src/remotes/kernel/RemoteHandle.ts index 648a0a287f..2debc3c62e 100644 --- a/packages/ocap-kernel/src/remotes/kernel/RemoteHandle.ts +++ b/packages/ocap-kernel/src/remotes/kernel/RemoteHandle.ts @@ -1587,16 +1587,14 @@ export class RemoteHandle implements EndpointHandle { } /** - * Apply the in-memory side of a peer restart: discard arrivals from the - * incarnation that is gone, cancel timers, reject in-flight URL redemption - * promises, and reset sequence counters. Must be called after - * {@link persistPeerRestart} and after the caller's savepoint has been - * released, with nothing awaited in between: a send whose turn comes there - * finds the queue emptied in the store but not yet retired here, and writes - * a message of the incarnation that has ended. + * Apply the in-memory side of a peer restart: cancel timers, reject + * in-flight URL redemption promises, and reset sequence counters. Must be + * called after {@link persistPeerRestart} and after the caller's savepoint + * has been released, with nothing awaited in between: a send whose turn + * comes there finds the queue emptied in the store but not yet retired + * here, and writes a message of the incarnation that has ended. */ finalizePeerRestart(): void { - this.#kernelQueue.discardRemoteInbound(this.remoteId); const pendingCount = this.#getPendingCount(); if (this.#hasPendingMessages()) { this.#logger.log( diff --git a/packages/ocap-kernel/src/remotes/kernel/RemoteManager.test.ts b/packages/ocap-kernel/src/remotes/kernel/RemoteManager.test.ts index ded6d90f30..38103caa15 100644 --- a/packages/ocap-kernel/src/remotes/kernel/RemoteManager.test.ts +++ b/packages/ocap-kernel/src/remotes/kernel/RemoteManager.test.ts @@ -832,6 +832,33 @@ describe('RemoteManager', () => { ) => Promise; } + /** + * Complete a handshake and run the crank it queues, the way the transport + * and the run loop do between them. + * + * @param peerId - The peer handshaking. + * @param incarnation - The incarnation it reports. + * @returns Whether the kernel judged this a restart. + */ + async function handshakeAndRunCrank( + peerId: string, + incarnation: string, + ): Promise { + const verdict = await getOnIncarnationChange()(peerId, incarnation); + const queued = vi + .mocked(mockKernelQueue.acceptPeerIncarnation) + .mock.calls.at(-1); + if (queued) { + const result = await remoteManager.applyIncarnationChange({ + type: 'peerIncarnation', + peerId: queued[0], + incarnation: queued[1], + }); + await result.afterCommit?.(); + } + return verdict; + } + it('triggers persist + finalize peer restart when persisted incarnation differs from observed', async () => { const peerId = 'peer-that-restarted'; const remote = remoteManager.establishRemote(peerId); @@ -840,7 +867,7 @@ describe('RemoteManager', () => { // Seed the persisted incarnation (as if a prior handshake recorded it). kernelStore.setPeerIncarnation(peerId, 'incarnation-A'); - const verdict = await getOnIncarnationChange()(peerId, 'incarnation-B'); + const verdict = await handshakeAndRunCrank(peerId, 'incarnation-B'); expect(verdict).toBe(true); expect(persistSpy).toHaveBeenCalled(); @@ -848,25 +875,13 @@ describe('RemoteManager', () => { expect(kernelStore.getPeerIncarnation(peerId)).toBe('incarnation-B'); }); - // A live restart arrives here, never through `handlePeerRestart`. - it('discards what the old incarnation sent', async () => { - const peerId = 'peer-that-restarted'; - const { remoteId } = remoteManager.establishRemote(peerId); - const discardSpy = vi.spyOn(mockKernelQueue, 'discardRemoteInbound'); - kernelStore.setPeerIncarnation(peerId, 'incarnation-A'); - - await getOnIncarnationChange()(peerId, 'incarnation-B'); - - expect(discardSpy).toHaveBeenCalledWith(remoteId); - }); - it('does not trigger restart on first observation of a peer', async () => { const peerId = 'peer-first-contact'; const remote = remoteManager.establishRemote(peerId); const persistSpy = vi.spyOn(remote, 'persistPeerRestart'); const finalizeSpy = vi.spyOn(remote, 'finalizePeerRestart'); - const verdict = await getOnIncarnationChange()(peerId, 'incarnation-A'); + const verdict = await handshakeAndRunCrank(peerId, 'incarnation-A'); expect(verdict).toBe(false); expect(persistSpy).not.toHaveBeenCalled(); @@ -881,7 +896,7 @@ describe('RemoteManager', () => { const finalizeSpy = vi.spyOn(remote, 'finalizePeerRestart'); kernelStore.setPeerIncarnation(peerId, 'incarnation-A'); - const verdict = await getOnIncarnationChange()(peerId, 'incarnation-A'); + const verdict = await handshakeAndRunCrank(peerId, 'incarnation-A'); expect(verdict).toBe(false); expect(persistSpy).not.toHaveBeenCalled(); @@ -900,17 +915,23 @@ describe('RemoteManager', () => { const resolvePromisesSpy = vi.spyOn(mockKernelQueue, 'resolvePromises'); - await getOnIncarnationChange()(peerId, 'incarnation-B'); + await handshakeAndRunCrank(peerId, 'incarnation-B'); - expect(resolvePromisesSpy).toHaveBeenCalledWith(remoteId, [ + expect(resolvePromisesSpy).toHaveBeenCalledWith( + remoteId, [ - kpid, - true, - expect.objectContaining({ - body: expect.stringContaining('[KERNEL:PEER_RESTARTED]'), - }), + [ + kpid, + true, + expect.objectContaining({ + body: expect.stringContaining('[KERNEL:PEER_RESTARTED]'), + }), + ], ], - ]); + // Buffered: resolving writes, so it belongs in the crank, and the + // notifies it produces must not go out before that crank commits. + false, + ); }); it('persists the incarnation and reports restart even when no remote handle exists', async () => { @@ -922,7 +943,7 @@ describe('RemoteManager', () => { // there's a RemoteHandle to reset. Transport callers use the verdict // to suppress stale outbound messages; that decision is correct // regardless of local handle presence. - const verdict = await getOnIncarnationChange()(peerId, 'incarnation-B'); + const verdict = await handshakeAndRunCrank(peerId, 'incarnation-B'); expect(verdict).toBe(true); expect(kernelStore.getPeerIncarnation(peerId)).toBe('incarnation-B'); @@ -932,16 +953,83 @@ describe('RemoteManager', () => { it('does not reject promises when there are none', async () => { const peerId = 'peer-without-promises'; - remoteManager.establishRemote(peerId); + const remote = remoteManager.establishRemote(peerId); + const persistSpy = vi.spyOn(remote, 'persistPeerRestart'); kernelStore.setPeerIncarnation(peerId, 'incarnation-A'); const resolvePromisesSpy = vi.spyOn(mockKernelQueue, 'resolvePromises'); - await getOnIncarnationChange()(peerId, 'incarnation-B'); + await handshakeAndRunCrank(peerId, 'incarnation-B'); + // The restart has to have happened, or this asserts nothing. + expect(persistSpy).toHaveBeenCalledOnce(); expect(resolvePromisesSpy).not.toHaveBeenCalled(); }); - it('rolls back the savepoint and preserves stored state when persistPeerRestart throws', async () => { + it('answers the handshake without waiting for a crank', async () => { + const peerId = 'peer-mid-handshake'; + const remote = remoteManager.establishRemote(peerId); + const persistSpy = vi.spyOn(remote, 'persistPeerRestart'); + kernelStore.setPeerIncarnation(peerId, 'incarnation-A'); + + const verdict = await getOnIncarnationChange()(peerId, 'incarnation-B'); + + // The transport needs the verdict to finish the handshake, and the + // store already knows the answer. + expect(verdict).toBe(true); + expect(mockKernelQueue.acceptPeerIncarnation).toHaveBeenCalledWith( + peerId, + 'incarnation-B', + ); + // None of the writes it implies have happened yet. + expect(persistSpy).not.toHaveBeenCalled(); + expect(kernelStore.getPeerIncarnation(peerId)).toBe('incarnation-A'); + }); + + it('rejects the promises the restarted remote was deciding only after the commit', async () => { + const peerId = 'peer-deferred-rejects'; + const remote = remoteManager.establishRemote(peerId); + kernelStore.setPeerIncarnation(peerId, 'incarnation-A'); + const finalizeSpy = vi.spyOn(remote, 'finalizePeerRestart'); + await getOnIncarnationChange()(peerId, 'incarnation-B'); + + const result = await remoteManager.applyIncarnationChange({ + type: 'peerIncarnation', + peerId, + incarnation: 'incarnation-B', + }); + + // Neither a rejection nor the in-memory reset is reversible by a + // rollback, so both wait for the crank to commit. + expect(finalizeSpy).not.toHaveBeenCalled(); + expect(kernelStore.getPeerIncarnation(peerId)).toBe('incarnation-B'); + await result.afterCommit?.(); + expect(finalizeSpy).toHaveBeenCalledOnce(); + }); + + it('records a change once however many handshakes queued it', async () => { + const peerId = 'peer-reconnecting'; + const remote = remoteManager.establishRemote(peerId); + const persistSpy = vi.spyOn(remote, 'persistPeerRestart'); + kernelStore.setPeerIncarnation(peerId, 'incarnation-A'); + // Two handshakes before either crank runs — a peer re-dialling while + // the first change is still waiting its turn. + await getOnIncarnationChange()(peerId, 'incarnation-B'); + await getOnIncarnationChange()(peerId, 'incarnation-B'); + + const item = { + type: 'peerIncarnation' as const, + peerId, + incarnation: 'incarnation-B', + }; + await remoteManager.applyIncarnationChange(item); + await remoteManager.applyIncarnationChange(item); + + // Tearing the c-list down twice would take the new incarnation's own + // state with it the second time. + expect(persistSpy).toHaveBeenCalledOnce(); + }); + + it('leaves the incarnation unadvanced when persistPeerRestart throws', async () => { const peerId = 'peer-handler-throws'; const remote = remoteManager.establishRemote(peerId); kernelStore.setPeerIncarnation(peerId, 'incarnation-A'); @@ -952,16 +1040,24 @@ describe('RemoteManager', () => { }); const finalizeSpy = vi.spyOn(remote, 'finalizePeerRestart'); + // The handshake is answered from what the store already says, so it + // succeeds; the writes it queued are what fail. + expect(await getOnIncarnationChange()(peerId, 'incarnation-B')).toBe( + true, + ); await expect( - getOnIncarnationChange()(peerId, 'incarnation-B'), + remoteManager.applyIncarnationChange({ + type: 'peerIncarnation', + peerId, + incarnation: 'incarnation-B', + }), ).rejects.toThrow(failure); - // Persisted incarnation must NOT have advanced — the savepoint - // rollback should have reverted the would-be setPeerIncarnation that - // runs after persistPeerRestart in the wrapped block. + // `setPeerIncarnation` runs after `persistPeerRestart`, so it never + // happened. Anything it had reached would go back with the crank, which + // this map-backed store cannot model. expect(kernelStore.getPeerIncarnation(peerId)).toBe('incarnation-A'); - // finalize must not run if the persisted phase failed: in-memory - // mutations would otherwise drift from the rolled-back kv view. + // finalize is post-commit work, and there was no commit. expect(finalizeSpy).not.toHaveBeenCalled(); }); }); diff --git a/packages/ocap-kernel/src/remotes/kernel/RemoteManager.ts b/packages/ocap-kernel/src/remotes/kernel/RemoteManager.ts index b3a5aa8151..ce6d287ef2 100644 --- a/packages/ocap-kernel/src/remotes/kernel/RemoteManager.ts +++ b/packages/ocap-kernel/src/remotes/kernel/RemoteManager.ts @@ -5,7 +5,12 @@ import { RemoteHandle } from './RemoteHandle.ts'; import type { KernelQueue } from '../../KernelQueue.ts'; import { makeKernelError } from '../../liveslots/kernel-marshal.ts'; import type { KernelStore } from '../../store/index.ts'; -import type { PlatformServices, RemoteId } from '../../types.ts'; +import type { + CrankResult, + PlatformServices, + RemoteId, + RunQueueItemPeerIncarnation, +} from '../../types.ts'; import type { RemoteIdentity, RemoteComms, @@ -225,7 +230,27 @@ export class RemoteManager { if (stored === observedIncarnation) { return false; } + // Answered from what the store already says, because the transport needs + // it to finish the handshake and cannot wait for a crank. The writes it + // implies are queued, and ordered against anything this peer has sent. + this.#kernelQueue.acceptPeerIncarnation(peerId, observedIncarnation); + return stored !== undefined; + } + /** + * Record a peer's incarnation change, in a crank of its own. + * + * @param item - The change, as accepted at handshake time. + * @returns The crank outcome, carrying the work that follows the commit. + */ + async applyIncarnationChange( + item: RunQueueItemPeerIncarnation, + ): Promise { + const { peerId, incarnation } = item; + const stored = this.#kernelStore.getPeerIncarnation(peerId); + if (stored === incarnation) { + return {}; + } const isRestart = stored !== undefined; const remote = isRestart ? this.#remotesByPeer.get(peerId) : undefined; @@ -237,52 +262,56 @@ export class RemoteManager { ? Array.from(this.#kernelStore.getPromisesByDecider(remote.remoteId)) : []; - const savepoint = `peerIncarnation_${peerId}`; - this.#kernelStore.createSavepoint(savepoint); - try { - if (isRestart) { - this.#logger?.log( - `Peer ${peerId.slice(0, 8)} restarted (incarnation ${stored.slice(0, 8)} → ${observedIncarnation.slice(0, 8)})`, + if (isRestart) { + this.#logger?.log( + `Peer ${peerId.slice(0, 8)} restarted (incarnation ${stored.slice(0, 8)} → ${incarnation.slice(0, 8)})`, + ); + if (remote) { + remote.persistPeerRestart(); + } else { + // No live RemoteHandle for the peer but a persisted incarnation + // exists — usually a transient race during kernel boot before + // initRemoteComms has finished restoring remotes. The persisted + // bookkeeping the missing handle would have cleaned up may leak. + // Surfacing as a warning so operators can correlate. + this.#logger?.warn( + `Peer ${peerId.slice(0, 8)} restart detected but no live RemoteHandle; advancing persisted incarnation without c-list cleanup`, ); - if (remote) { - remote.persistPeerRestart(); - } else { - // No live RemoteHandle for the peer but a persisted incarnation - // exists — usually a transient race during kernel boot before - // initRemoteComms has finished restoring remotes. The persisted - // bookkeeping the missing handle would have cleaned up may leak. - // Surfacing as a warning so operators can correlate. - this.#logger?.warn( - `Peer ${peerId.slice(0, 8)} restart detected but no live RemoteHandle; advancing persisted incarnation without c-list cleanup`, - ); - } } - this.#kernelStore.setPeerIncarnation(peerId, observedIncarnation); - this.#kernelStore.releaseSavepoint(savepoint); - } catch (error) { - this.#kernelStore.rollbackSavepoint(savepoint); - throw error; + } + this.#kernelStore.setPeerIncarnation(peerId, incarnation); + + if (!isRestart || !remote) { + return {}; } - // Post-commit fan-out: in-memory state changes and run-queue - // mutations are not reversible by a savepoint, so they wait until the - // kv layer is durable. - if (isRestart && remote) { - remote.finalizePeerRestart(); - if (promisesToReject.length > 0) { - const failure = makeKernelError( - 'PEER_RESTARTED', - 'Remote peer restarted (incarnation changed)', + // Rejected here, inside the crank, because resolving a promise writes: + // its state, its reference counts, and a notify for every subscriber. + // `immediate: false` buffers those notifies until the crank commits, the + // way a vat's own syscalls are buffered, so nothing is delivered on the + // strength of writes that may yet roll back. Doing it post-commit instead + // would put the writes outside every transaction, and would race the + // transport's own give-up handling, which rejects the same promises from + // a send continuation — the second rejection `Fail`s. + if (promisesToReject.length > 0) { + const failure = makeKernelError( + 'PEER_RESTARTED', + 'Remote peer restarted (incarnation changed)', + ); + for (const kpid of promisesToReject) { + this.#kernelQueue.resolvePromises( + remote.remoteId, + [[kpid, true, failure]], + false, ); - for (const kpid of promisesToReject) { - this.#kernelQueue.resolvePromises(remote.remoteId, [ - [kpid, true, failure], - ]); - } } } - return isRestart; + return { + // The in-memory counters only, which no rollback could put back and + // which must not be visible before the writes above are durable. + afterCommit: async () => remote.finalizePeerRestart(), + }; } /** diff --git a/packages/ocap-kernel/src/types.ts b/packages/ocap-kernel/src/types.ts index df1bb4cd47..629e6064de 100644 --- a/packages/ocap-kernel/src/types.ts +++ b/packages/ocap-kernel/src/types.ts @@ -386,12 +386,23 @@ export type RunQueueItemRemoteInbound = Infer< typeof RunQueueItemRemoteInboundStruct >; +const RunQueueItemPeerIncarnationStruct = object({ + type: literal('peerIncarnation'), + peerId: string(), + incarnation: string(), +}); + +export type RunQueueItemPeerIncarnation = Infer< + typeof RunQueueItemPeerIncarnationStruct +>; + export const RunQueueItemStruct = union([ RunQueueItemSendStruct, RunQueueItemNotifyStruct, RunQueueItemGCActionStruct, RunQueueItemBringOutYourDeadStruct, RunQueueItemRemoteInboundStruct, + RunQueueItemPeerIncarnationStruct, ]); export type RunQueueItem = Infer; diff --git a/packages/ocap-kernel/test/remotes-mocks.ts b/packages/ocap-kernel/test/remotes-mocks.ts index 277684ea23..c80cfbf58b 100644 --- a/packages/ocap-kernel/test/remotes-mocks.ts +++ b/packages/ocap-kernel/test/remotes-mocks.ts @@ -77,7 +77,7 @@ export class MockRemotesFactory { enqueueSend: vi.fn(), enqueueNotify: vi.fn(), acceptRemoteInbound: vi.fn(), - discardRemoteInbound: vi.fn(), + acceptPeerIncarnation: vi.fn(), resolvePromises: vi.fn(), assertRunLoopAlive: vi.fn(), waitForCrank: vi.fn(),