From 287a12237f1558a1e8364f9018c689f4de624e96 Mon Sep 17 00:00:00 2001 From: Dimitris Marlagkoutsos Date: Tue, 15 Sep 2026 23:48:32 +0200 Subject: [PATCH 1/4] feat(ocap-kernel): carry out a peer incarnation change on the run loop MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The change took a `peerIncarnation_*` savepoint of its own, which nested inside whichever crank was open — the same defect the inbound message path just shed. It becomes a run queue item, carried out in a crank of its own. The handshake still gets its answer immediately: whether this is a restart is a read of what the store already says, and the transport needs it to decide whether to reset the connection. Only the writes are queued. Queued behind anything that peer has already sent, which is what keeps the two incarnations apart: the messages ahead of the change belong to the one that is ending and are recorded against it, the ones behind it to the one that is starting. So the eager discard the previous branch needed goes away — it would now throw away the new incarnation's messages rather than the old one's. Rejecting the promises the restarted remote was deciding, and resetting its in-memory state, move to `afterCommit`. Co-Authored-By: Claude Opus 5 (1M context) --- packages/ocap-kernel/CHANGELOG.md | 4 +- packages/ocap-kernel/src/Kernel.ts | 9 +- packages/ocap-kernel/src/KernelQueue.test.ts | 23 ++-- packages/ocap-kernel/src/KernelQueue.ts | 40 +++--- packages/ocap-kernel/src/KernelRouter.ts | 1 - .../src/remotes/kernel/RemoteHandle.ts | 14 +-- .../src/remotes/kernel/RemoteManager.test.ts | 114 ++++++++++++++---- .../src/remotes/kernel/RemoteManager.ts | 103 +++++++++------- packages/ocap-kernel/src/types.ts | 11 ++ packages/ocap-kernel/test/remotes-mocks.ts | 2 +- 10 files changed, 215 insertions(+), 106 deletions(-) diff --git a/packages/ocap-kernel/CHANGELOG.md b/packages/ocap-kernel/CHANGELOG.md index c2aab043aa..a8f86f46b6 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 ([#1105](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1105)) + - 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.ts b/packages/ocap-kernel/src/Kernel.ts index c041d060be..9b1c0d1794 100644 --- a/packages/ocap-kernel/src/Kernel.ts +++ b/packages/ocap-kernel/src/Kernel.ts @@ -343,7 +343,14 @@ 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. + item.type === 'peerIncarnation' + ? await this.#remoteManager.applyIncarnationChange(item) + : await this.#kernelRouter.deliver(item), + ) .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..9d01fe6fde 100644 --- a/packages/ocap-kernel/src/KernelQueue.test.ts +++ b/packages/ocap-kernel/src/KernelQueue.test.ts @@ -801,23 +801,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..9f49f452d5 100644 --- a/packages/ocap-kernel/src/KernelRouter.ts +++ b/packages/ocap-kernel/src/KernelRouter.ts @@ -108,7 +108,6 @@ export class KernelRouter { case 'remoteInbound': return await this.#deliverRemoteInbound(item); default: - // @ts-expect-error Runtime does not respect "never". Fail`unsupported or unknown run queue item type ${item.type}`; } return undefined; 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..fccb07f86a 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,7 +915,7 @@ describe('RemoteManager', () => { const resolvePromisesSpy = vi.spyOn(mockKernelQueue, 'resolvePromises'); - await getOnIncarnationChange()(peerId, 'incarnation-B'); + await handshakeAndRunCrank(peerId, 'incarnation-B'); expect(resolvePromisesSpy).toHaveBeenCalledWith(remoteId, [ [ @@ -922,7 +937,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'); @@ -936,12 +951,53 @@ describe('RemoteManager', () => { kernelStore.setPeerIncarnation(peerId, 'incarnation-A'); const resolvePromisesSpy = vi.spyOn(mockKernelQueue, 'resolvePromises'); - await getOnIncarnationChange()(peerId, 'incarnation-B'); + await handshakeAndRunCrank(peerId, 'incarnation-B'); 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('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 +1008,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..36acaaa3e2 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,46 @@ 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); - // 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)', - ); - for (const kpid of promisesToReject) { - this.#kernelQueue.resolvePromises(remote.remoteId, [ - [kpid, true, failure], - ]); - } - } + if (!isRestart || !remote) { + return {}; } - - return isRestart; + return { + afterCommit: async () => { + // In-memory state changes and run-queue mutations are not reversible + // by a rollback, so they wait until the kv layer is durable. + remote.finalizePeerRestart(); + 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], + ]); + } + } + }, + }; } /** 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(), From 48d754dbcb74929bf8d41bc0641047881279e9eb Mon Sep 17 00:00:00 2001 From: Dimitris Marlagkoutsos Date: Wed, 16 Sep 2026 00:14:32 +0200 Subject: [PATCH 2/4] fix(ocap-kernel): reject a restarted peer's promises inside the crank MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review follow-up, and the same mistake the branch below it made: `afterCommit` must not write the kernel store, and `resolvePromises` writes — promise state, reference counts, and a notify row per subscriber, all in autocommit once `endCrank` has committed. It also raced the transport's own give-up handling, which rejects the same promises from a send continuation; whichever arrived second would `Fail` out of a post-commit hook and kill the kernel. The rejections move into the crank, buffered with `immediate: false` so the notifies they produce still wait for the commit, the way a vat's syscalls do. `afterCommit` keeps only the in-memory counter reset. An incarnation change that cannot be recorded no longer kills the run loop — the same containment the inbound message path already has — and the router's exhaustiveness check comes back by excluding the item type the kernel handles itself, rather than by deleting the directive that proved it. Co-Authored-By: Claude Opus 5 (1M context) --- packages/ocap-kernel/src/Kernel.test.ts | 46 ++++++++++++++++++- packages/ocap-kernel/src/Kernel.ts | 23 ++++++++-- packages/ocap-kernel/src/KernelQueue.test.ts | 17 +++++++ packages/ocap-kernel/src/KernelRouter.ts | 11 ++++- .../src/remotes/kernel/RemoteManager.test.ts | 43 ++++++++++++++--- .../src/remotes/kernel/RemoteManager.ts | 42 ++++++++++------- 6 files changed, 152 insertions(+), 30 deletions(-) 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 9b1c0d1794..77655b94f9 100644 --- a/packages/ocap-kernel/src/Kernel.ts +++ b/packages/ocap-kernel/src/Kernel.ts @@ -343,14 +343,27 @@ export class Kernel { // Start the kernel queue processing (non-blocking) // This runs for the entire lifetime of the kernel this.#kernelQueue - .run(async (item) => + .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. - item.type === 'peerIncarnation' - ? await this.#remoteManager.applyIncarnationChange(item) - : await this.#kernelRouter.deliver(item), - ) + 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 9d01fe6fde..807284e582 100644 --- a/packages/ocap-kernel/src/KernelQueue.test.ts +++ b/packages/ocap-kernel/src/KernelQueue.test.ts @@ -784,6 +784,23 @@ describe('KernelQueue', () => { }); }); + 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')); ( diff --git a/packages/ocap-kernel/src/KernelRouter.ts b/packages/ocap-kernel/src/KernelRouter.ts index 9f49f452d5..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); @@ -108,6 +116,7 @@ export class KernelRouter { case 'remoteInbound': return await this.#deliverRemoteInbound(item); default: + // @ts-expect-error Runtime does not respect "never". Fail`unsupported or unknown run queue item type ${item.type}`; } return undefined; diff --git a/packages/ocap-kernel/src/remotes/kernel/RemoteManager.test.ts b/packages/ocap-kernel/src/remotes/kernel/RemoteManager.test.ts index fccb07f86a..856da41fbd 100644 --- a/packages/ocap-kernel/src/remotes/kernel/RemoteManager.test.ts +++ b/packages/ocap-kernel/src/remotes/kernel/RemoteManager.test.ts @@ -917,15 +917,21 @@ describe('RemoteManager', () => { 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 () => { @@ -997,6 +1003,29 @@ describe('RemoteManager', () => { 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); diff --git a/packages/ocap-kernel/src/remotes/kernel/RemoteManager.ts b/packages/ocap-kernel/src/remotes/kernel/RemoteManager.ts index 36acaaa3e2..ce6d287ef2 100644 --- a/packages/ocap-kernel/src/remotes/kernel/RemoteManager.ts +++ b/packages/ocap-kernel/src/remotes/kernel/RemoteManager.ts @@ -284,23 +284,33 @@ export class RemoteManager { if (!isRestart || !remote) { return {}; } + + // 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, + ); + } + } + return { - afterCommit: async () => { - // In-memory state changes and run-queue mutations are not reversible - // by a rollback, so they wait until the kv layer is durable. - remote.finalizePeerRestart(); - 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], - ]); - } - } - }, + // 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(), }; } From 7553e26d11df0f82b5676c30817f8b7131d99e80 Mon Sep 17 00:00:00 2001 From: Dimitris Marlagkoutsos Date: Wed, 16 Sep 2026 00:15:22 +0200 Subject: [PATCH 3/4] chore: link the changelog entries to the real PR number Co-Authored-By: Claude Opus 5 (1M context) --- packages/ocap-kernel/CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/ocap-kernel/CHANGELOG.md b/packages/ocap-kernel/CHANGELOG.md index a8f86f46b6..119f41fddd 100644 --- a/packages/ocap-kernel/CHANGELOG.md +++ b/packages/ocap-kernel/CHANGELOG.md @@ -49,7 +49,7 @@ 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 ([#1105](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1105)) +- **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)) From ff8bb8531358128e341e732c4bfff368673cafd4 Mon Sep 17 00:00:00 2001 From: Dimitris Marlagkoutsos Date: Wed, 16 Sep 2026 00:22:22 +0200 Subject: [PATCH 4/4] test(ocap-kernel): pin the run loop wake, on both paths that arm it An arrival is the one kind of work that does not go through `#enqueueRun`, so its wake is its own. Deleting either call left every test green; a parked loop with work waiting is a permanent wedge. `does not reject promises when there are none` also asserted only a negative, and passed whether or not the restart it describes had happened. Co-Authored-By: Claude Opus 5 (1M context) --- packages/ocap-kernel/src/KernelQueue.test.ts | 37 +++++++++++++++++++ .../src/remotes/kernel/RemoteManager.test.ts | 5 ++- 2 files changed, 41 insertions(+), 1 deletion(-) diff --git a/packages/ocap-kernel/src/KernelQueue.test.ts b/packages/ocap-kernel/src/KernelQueue.test.ts index 807284e582..fe8b779014 100644 --- a/packages/ocap-kernel/src/KernelQueue.test.ts +++ b/packages/ocap-kernel/src/KernelQueue.test.ts @@ -784,6 +784,43 @@ 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')); ( diff --git a/packages/ocap-kernel/src/remotes/kernel/RemoteManager.test.ts b/packages/ocap-kernel/src/remotes/kernel/RemoteManager.test.ts index 856da41fbd..38103caa15 100644 --- a/packages/ocap-kernel/src/remotes/kernel/RemoteManager.test.ts +++ b/packages/ocap-kernel/src/remotes/kernel/RemoteManager.test.ts @@ -953,12 +953,15 @@ 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 handshakeAndRunCrank(peerId, 'incarnation-B'); + // The restart has to have happened, or this asserts nothing. + expect(persistSpy).toHaveBeenCalledOnce(); expect(resolvePromisesSpy).not.toHaveBeenCalled(); });