Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion packages/ocap-kernel/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down
46 changes: 45 additions & 1 deletion packages/ocap-kernel/src/Kernel.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<unknown>) | undefined;

// Like the real run loop, this settles only if the kernel dies.
run = vi.fn(
async () =>
async (deliver: (item: unknown) => Promise<unknown>) =>
new Promise<never>((_resolve, reject) => {
this.deliver = deliver;
this.#rejectRunLoop = reject;
}),
);
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -249,6 +255,44 @@ describe('Kernel', () => {
});
});

describe('the run loop callback', () => {
const makeRunningKernel = async (): Promise<void> => {
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(
Expand Down
22 changes: 21 additions & 1 deletion packages/ocap-kernel/src/Kernel.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 };
}
Comment thread
cursor[bot] marked this conversation as resolved.
})
.catch((error) => this.#handleRunLoopFailure(error));

// Launch new system subclusters (requires queue to be running)
Expand Down
77 changes: 68 additions & 9 deletions packages/ocap-kernel/src/KernelQueue.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'));
(
Expand All @@ -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',
]);
});
});
Expand Down
40 changes: 22 additions & 18 deletions packages/ocap-kernel/src/KernelQueue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import type {
RemoteId,
RunLoopStatus,
RunQueueItem,
RunQueueItemPeerIncarnation,
RunQueueItemRemoteInbound,
RunQueueItemNotify,
RunQueueItemSend,
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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();
Comment thread
cursor[bot] marked this conversation as resolved.
}

/**
Expand Down
10 changes: 9 additions & 1 deletion packages/ocap-kernel/src/KernelRouter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import type {
RunQueueItem,
RunQueueItemSend,
RemoteEndpointHandle,
RunQueueItemPeerIncarnation,
RunQueueItemBringOutYourDead,
RunQueueItemRemoteInbound,
RunQueueItemNotify,
Expand All @@ -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<RunQueueItem, RunQueueItemPeerIncarnation>;

type MessageRoute = {
endpointId?: EndpointId | 'kernel';
target: KRef;
Expand Down Expand Up @@ -93,7 +101,7 @@ export class KernelRouter {
* @param item - The message/notification to deliver.
* @returns The crank outcome.
*/
async deliver(item: RunQueueItem): Promise<CrankResult | undefined> {
async deliver(item: RoutedRunQueueItem): Promise<CrankResult | undefined> {
switch (item.type) {
case 'send':
return await this.#deliverSend(item);
Expand Down
14 changes: 6 additions & 8 deletions packages/ocap-kernel/src/remotes/kernel/RemoteHandle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Comment thread
cursor[bot] marked this conversation as resolved.
const pendingCount = this.#getPendingCount();
if (this.#hasPendingMessages()) {
this.#logger.log(
Expand Down
Loading
Loading