Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
6aeec5d
feat(ocap-kernel): carry out a vat restart on the run loop
sirtimid Sep 15, 2026
3a100bc
fix(ocap-kernel): let every restart request queue its own item
sirtimid Sep 15, 2026
b9256e6
test(ocap-kernel): pass restartVat to the router in the result-promis…
sirtimid Oct 1, 2026
e717fef
fix(ocap-kernel): abort the crank of a restart whose relaunch fails
sirtimid Oct 1, 2026
7614c6d
docs(ocap-kernel): cut the restart comments down to what the code can…
sirtimid Oct 1, 2026
cda91df
fix(ocap-kernel): answer a restart's callers once its crank ends
sirtimid Oct 2, 2026
dfa78c2
fix(ocap-kernel): hold a restart request made mid-crank until the cra…
sirtimid Oct 2, 2026
8935a9e
fix(ocap-kernel): answer a restart with the vat it restarted
sirtimid Oct 2, 2026
8506a4e
fix(ocap-kernel): answer a restart's callers only once its crank commits
sirtimid Oct 6, 2026
b8761d0
fix(ocap-kernel): reject a queued restart when the kernel stops
sirtimid Oct 6, 2026
fb16122
fix(ocap-kernel): relaunch a restarted vat whose old worker will not …
sirtimid Oct 6, 2026
c287cca
fix(ocap-kernel): give up on a restarted vat's worker that does not s…
sirtimid Oct 6, 2026
bfd5ec6
docs(ocap-kernel): list the restart changes as breaking
sirtimid Oct 6, 2026
079def1
fix(ocap-kernel): refuse a relaunch timeout setTimeout cannot honour
sirtimid Oct 6, 2026
975495b
Merge remote-tracking branch 'origin/main' into sirtimid/restart-vat-…
sirtimid Oct 6, 2026
3fab908
fix(ocap-kernel): close the channel of a relaunch that timed out
sirtimid Oct 6, 2026
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
39 changes: 39 additions & 0 deletions packages/kernel-test/src/vat-lifecycle.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import type { VatId } from '@metamask/ocap-kernel';
import { describe, expect, it, beforeEach } from 'vitest';

import {
extractTestLogs,
getBundleSpec,
makeKernel,
makeTestLogger,
Expand Down Expand Up @@ -147,6 +148,44 @@ describe('Vat Lifecycle', { timeout: 30_000 }, () => {
expect(kernelStore.getRootObject(deadVatId)).toBeUndefined();
});

it('restarts one vat twice through the run loop', async () => {
const kernelDatabase = await makeSQLKernelDatabase({
dbFilename: ':memory:',
});
const { logger: restartLogger, entries } = makeTestLogger();
const kernel = await makeKernel(kernelDatabase, true, restartLogger);

await kernel.launchSubcluster({
bootstrap: 'alice',
forceReset: true,
vats: {
alice: {
bundleSpec: getBundleSpec('resume-vat'),
parameters: { name: 'Alice' },
},
// `resume-vat`'s bootstrap introduces its peers to each other, so the
// subcluster needs all three.
bob: {
bundleSpec: getBundleSpec('resume-vat'),
parameters: { name: 'Bob' },
},
carol: {
bundleSpec: getBundleSpec('resume-vat'),
parameters: { name: 'Carol' },
},
},
});
await waitUntilQuiescent();

// `start count` lives in the vat's baggage, so a restart that answered its
// caller without replacing the worker would not move it.
await kernel.restartVat('v1');
await kernel.restartVat('v1');
await waitUntilQuiescent(1000);

expect(extractTestLogs(entries, 'v1')).toContain('Alice: start count: 3');
});

it('leaves no record of a terminated vat for the next boot to restore', async () => {
const kernelDatabase = await makeSQLKernelDatabase({
dbFilename: ':memory:',
Expand Down
3 changes: 3 additions & 0 deletions packages/ocap-kernel/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Added

- Add a `vatRelaunchTimeoutMs` option to `Kernel.make`: how long a vat restart waits for the new worker before giving up, 30 seconds by default ([#1096](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1096))
- A `CrankResult` may carry `afterCommit`, work the run loop runs once the crank has committed and skips if it aborts, for what a rollback could not undo anyway: in-memory state, and sending a message the crank has already written down. While it runs it may not write the kernel store, since by then there is no transaction to write into ([#1101](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1101))
- Add `orphanKernelObject` to the kernel store, which gives up the kernel's record of who owns an object ([#1091](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1091))
- Add `IOListener`, an endpoint peers connect to that yields one `IOChannel` per connection via `accept()`, replacing the previous one-client-at-a-time channel. Each accepted connection is a distinct object, so holding one conveys no way to reach another, and `direction` is enforced per connection. `accept()` resolves `null` once the listener is closed so an accept loop can terminate rather than hang ([#1007](https://github.com/MetaMask/ocap-kernel/pull/1007))
Expand Down Expand Up @@ -65,6 +66,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- Migrate to `DataView`, whose `setFloat*` methods lockdown repairs to write only canonical `NaN`s. This applies to bundled dependencies as much as to vat code: a library reaching for `Float64Array` at module scope will now throw on vat startup
- **BREAKING:** Remove `TextEncoder` and `TextDecoder` from `AllowedGlobalName`, since SES 2 permits them in every compartment and the kernel no longer endows them. A `VatConfig.globals` still naming either now fails validation, rejecting the whole launch with `invalid cluster config`, so drop them from cluster configs ([#1112](https://github.com/MetaMask/ocap-kernel/pull/1112))
- Vats keep access to both, and `allowedGlobalNames` can no longer withhold them
- **BREAKING:** `restartVat` is carried out by the run loop in a crank of its own, so a crank can no longer observe a vat between workers as dead. It now waits behind the run queue, and rejects if the run loop dies or the kernel is stopped, reset or has its storage cleared first. Concurrent restarts of one vat are carried out once ([#1096](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1096))
- **BREAKING:** A vat whose relaunch fails, or whose new worker does not start within `vatRelaunchTimeoutMs`, is terminated rather than left persisted with no worker ([#1096](https://github.com/Consensys-Incorporated/ocap-kernel/pull/1096))

### Fixed

Expand Down
173 changes: 102 additions & 71 deletions packages/ocap-kernel/src/Kernel.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,8 @@ import type {
OnRunLoopFailure,
PlatformServices,
ClusterConfig,
CrankResult,
RunQueueItem,
} from './types.ts';
import { VatHandle } from './vats/VatHandle.ts';
import { makeMapKernelDatabase } from '../test/storage.ts';
Expand All @@ -31,13 +33,19 @@ const mocks = vi.hoisted(() => {

#rejectRunLoop: ((error: Error) => void) | undefined;

/**
* The router's dispatch, so queued work can be carried out as a crank
* would carry it out.
*/
#deliver: ((item: RunQueueItem) => Promise<unknown>) | undefined;

// Like the real run loop, this settles only if the kernel dies.
run = vi.fn(
async () =>
new Promise<never>((_resolve, reject) => {
this.#rejectRunLoop = reject;
}),
);
run = vi.fn(async (deliver: (item: RunQueueItem) => Promise<unknown>) => {
this.#deliver = deliver;
return new Promise<never>((_resolve, reject) => {
this.#rejectRunLoop = reject;
});
});

/**
* Fail the run loop, in the order the real `KernelQueue.run` does: the
Expand All @@ -49,6 +57,10 @@ const mocks = vi.hoisted(() => {
killRunLoop(error: Error): void {
this.#runLoopFailure = error;
this.#rejectRunLoop?.(error);
for (const reject of this.#runLoopDeathWaiters) {
reject(error);
}
this.#runLoopDeathWaiters.clear();
}

getRunLoopStatus = vi.fn(() =>
Expand All @@ -61,6 +73,45 @@ const mocks = vi.hoisted(() => {
: { state: 'running' },
);

// A stand-in for the run loop: the item is delivered through the router,
// on a later turn, exactly as a crank would deliver it.
enqueueRestartVat = vi.fn((vatId: string) => {
this.#deliverLater({ type: 'restartVat', vatId } as RunQueueItem);
});

/**
* @param item - The queued item to hand to the router on a later turn.
*/
#deliverLater(item: RunQueueItem): void {
queueMicrotask(() => {
this.#runCrank(item).catch((error: unknown) => {
throw error;
});
});
}

/**
* Deliver an item and run what follows its commit, as a crank does. No
* test here aborts one.
*
* @param item - The item to deliver.
*/
async #runCrank(item: RunQueueItem): Promise<void> {
const crankResult = (await this.#deliver?.(item)) as
| CrankResult
| undefined;
await crankResult?.afterCommit?.();
}

onRunLoopDeath = vi.fn((reject: (error: Error) => void) => {
this.#runLoopDeathWaiters.add(reject);
return () => {
this.#runLoopDeathWaiters.delete(reject);
};
});

readonly #runLoopDeathWaiters = new Set<(error: Error) => void>();

assertRunLoopAlive = vi.fn((what: string) => {
if (this.#runLoopFailure) {
throw new Error(`Kernel run loop died; cannot ${what}`, {
Expand Down Expand Up @@ -179,7 +230,7 @@ describe('Kernel', () => {
vatId,
config: vatConfig,
init: vi.fn(),
terminate: vi.fn(),
terminate: vi.fn().mockResolvedValue(undefined),
handleMessage: vi.fn(),
deliverMessage: vi.fn(),
deliverNotify: vi.fn(),
Expand Down Expand Up @@ -790,94 +841,74 @@ describe('Kernel', () => {
});

describe('restartVat()', () => {
it('preserves vat state across multiple restarts', async () => {
const kernel = await Kernel.make(
mockPlatformServices,
mockKernelDatabase,
);
await kernel.launchSubcluster(makeSingleVatClusterConfig());
await kernel.restartVat('v1');
expect(kernel.getVatIds()).toStrictEqual(['v1']);
await kernel.restartVat('v1');
expect(kernel.getVatIds()).toStrictEqual(['v1']);
expect(vatHandles).toHaveLength(3); // Three instances created
expect(vatHandles[0]?.terminate).toHaveBeenCalledTimes(1);
expect(vatHandles[1]?.terminate).toHaveBeenCalledTimes(1);
expect(vatHandles[2]?.terminate).not.toHaveBeenCalled();
expect(launchWorkerMock).toHaveBeenCalledTimes(3); // initial + 2 restarts
expect(launchWorkerMock).toHaveBeenLastCalledWith(
'v1',
makeMockVatConfig(),
);
});

it('restarts a vat', async () => {
const kernel = await Kernel.make(
mockPlatformServices,
mockKernelDatabase,
);
await kernel.launchSubcluster(makeSingleVatClusterConfig());
expect(kernel.getVatIds()).toStrictEqual(['v1']);
await kernel.restartVat('v1');

const returnedHandle = await kernel.restartVat('v1');

expect(
mocks.KernelQueue.lastInstance.enqueueRestartVat,
).toHaveBeenCalledWith('v1');
expect(vatHandles[0]?.terminate).toHaveBeenCalledOnce();
expect(terminateWorkerMock).toHaveBeenCalledOnce();
expect(launchWorkerMock).toHaveBeenCalledTimes(2);
expect(launchWorkerMock).toHaveBeenLastCalledWith(
'v1',
makeMockVatConfig(),
);
expect(kernel.getVatIds()).toStrictEqual(['v1']);
expect(makeVatHandleMock).toHaveBeenCalledTimes(2);
});

it('throws error when restarting non-existent vat', async () => {
const kernel = await Kernel.make(
mockPlatformServices,
mockKernelDatabase,
);
await expect(kernel.restartVat('v999')).rejects.toThrow(VatNotFoundError);
expect(vatHandles).toHaveLength(0);
expect(launchWorkerMock).not.toHaveBeenCalled();
expect(returnedHandle).toBe(vatHandles[1]);
});

it('handles restart failure during termination', async () => {
it('rejects its caller when the run loop dies first', async () => {
const kernel = await Kernel.make(
mockPlatformServices,
mockKernelDatabase,
);
await kernel.launchSubcluster(makeSingleVatClusterConfig());
vatHandles[0]?.terminate.mockRejectedValueOnce(
new Error('Termination failed'),
);
await expect(kernel.restartVat('v1')).rejects.toThrow(
'Termination failed',
);
expect(launchWorkerMock).toHaveBeenCalledTimes(1);
});
const queue = mocks.KernelQueue.lastInstance;
queue.enqueueRestartVat.mockImplementationOnce(() => undefined);

it('handles restart failure during launch', async () => {
const kernel = await Kernel.make(
mockPlatformServices,
mockKernelDatabase,
);
await kernel.launchSubcluster(makeSingleVatClusterConfig());
launchWorkerMock.mockRejectedValueOnce(new Error('Launch failed'));
await expect(kernel.restartVat('v1')).rejects.toThrow('Launch failed');
expect(vatHandles[0]?.terminate).toHaveBeenCalledOnce();
expect(kernel.getVatIds()).toStrictEqual([]);
const restarting = kernel.restartVat('v1');
queue.killRunLoop(new Error('run loop boom'));

await expect(restarting).rejects.toThrow('run loop boom');
});

it('returns the new vat handle', async () => {
it.each([
{ method: 'clearStorage', message: 'Kernel storage was cleared' },
{ method: 'reset', message: 'Kernel was reset' },
{ method: 'stop', message: 'Kernel was stopped' },
] as const)(
'rejects a queued restart on $method',
async ({ method, message }) => {
const kernel = await Kernel.make(
mockPlatformServices,
mockKernelDatabase,
);
await kernel.launchSubcluster(makeSingleVatClusterConfig());
mocks.KernelQueue.lastInstance.enqueueRestartVat.mockImplementationOnce(
() => undefined,
);

const restarting = kernel.restartVat('v1');
await kernel[method]();

await expect(restarting).rejects.toThrow(message);
},
);

it('throws error when restarting non-existent vat', async () => {
const kernel = await Kernel.make(
mockPlatformServices,
mockKernelDatabase,
);
await kernel.launchSubcluster(makeSingleVatClusterConfig());
const originalHandle = vatHandles[0];
const returnedHandle = await kernel.restartVat('v1');
expect(returnedHandle).not.toBe(originalHandle);
expect(returnedHandle).toBe(vatHandles[1]);
expect(returnedHandle.vatId).toBe('v1');
await expect(kernel.restartVat('v999')).rejects.toThrow(VatNotFoundError);
expect(vatHandles).toHaveLength(0);
expect(launchWorkerMock).not.toHaveBeenCalled();
expect(
mocks.KernelQueue.lastInstance.enqueueRestartVat,
).not.toHaveBeenCalled();
});
});

Expand Down
19 changes: 18 additions & 1 deletion packages/ocap-kernel/src/Kernel.ts
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,7 @@ export class Kernel {
* @param options.ioListenerFactory - Optional factory for creating IO listeners.
* @param options.allowedGlobalNames - Optional list of allowed global names for vat endowments.
* @param options.onRunLoopFailure - Optional handler called if the run loop dies.
* @param options.vatRelaunchTimeoutMs - How long a vat restart waits for the new worker before terminating the vat.
* @param options.auditRefCounts - If true, verify every kref's reference
* counts against the references the kernel actually holds at the end of each
* crank, and throw on any mismatch. Intended for tests and debugging; the
Expand All @@ -126,6 +127,7 @@ export class Kernel {
ioListenerFactory?: IOListenerFactory;
allowedGlobalNames?: AllowedGlobalName[];
onRunLoopFailure?: OnRunLoopFailure;
vatRelaunchTimeoutMs?: number;
auditRefCounts?: boolean;
} = {},
) {
Expand Down Expand Up @@ -161,6 +163,7 @@ export class Kernel {
kernelQueue: this.#kernelQueue,
logger: this.#logger.subLogger({ tags: ['VatManager'] }),
allowedGlobalNames: options.allowedGlobalNames,
vatRelaunchTimeoutMs: options.vatRelaunchTimeoutMs,
});

this.#remoteManager = new RemoteManager({
Expand Down Expand Up @@ -224,6 +227,7 @@ export class Kernel {
this.#kernelServiceManager.invokeKernelService.bind(
this.#kernelServiceManager,
),
this.#vatManager.performVatRestart.bind(this.#vatManager),
this.#logger,
);

Expand Down Expand Up @@ -256,6 +260,7 @@ export class Kernel {
* @param options.systemSubclusters - Optional array of system subcluster configurations.
* @param options.allowedGlobalNames - Optional list of allowed global names for vat endowments. When set, only these names from the `VatSupervisor`'s configured endowments (see `createDefaultEndowments`) are available to vats.
* @param options.onRunLoopFailure - Optional handler called if the run loop dies. The kernel must be restarted after that, so an embedder that outlives it (e.g. a daemon) should use this to terminate or restart.
* @param options.vatRelaunchTimeoutMs - How long a vat restart waits for the new worker before terminating the vat, in milliseconds: more than 0 and at most 2^31 - 1. Defaults to 30 seconds.
* @param options.auditRefCounts - If true, verify reference counts against
* ground truth at the end of each crank and throw on any mismatch.
* @returns A promise for the new kernel instance.
Expand All @@ -272,6 +277,7 @@ export class Kernel {
systemSubclusters?: SystemSubclusterConfig[];
allowedGlobalNames?: AllowedGlobalName[];
onRunLoopFailure?: OnRunLoopFailure;
vatRelaunchTimeoutMs?: number;
auditRefCounts?: boolean;
} = {},
): Promise<Kernel> {
Expand Down Expand Up @@ -614,7 +620,9 @@ export class Kernel {
}

/**
* Restarts a vat.
* Restarts a vat. The run loop carries the restart out, so this waits behind
* the run queue and rejects if the run loop dies. A vat whose relaunch fails
* is terminated.
*
* @param vatId - The ID of the vat to restart.
* @returns A promise for the restarted vat handle.
Expand All @@ -640,6 +648,9 @@ export class Kernel {
async clearStorage(): Promise<void> {
await this.#kernelQueue.waitForCrank();
this.#kernelStore.clear();
this.#vatManager.abandonRestarts(
new Error('Kernel storage was cleared; the restart was abandoned'),
);
}

/**
Expand Down Expand Up @@ -841,6 +852,9 @@ export class Kernel {
await this.terminateAllVats();
this.#subclusterManager.clearSystemSubclusters();
this.#resetKernelState();
this.#vatManager.abandonRestarts(
new Error('Kernel was reset; the restart was abandoned'),
);
} catch (error) {
this.#logger.error('Error resetting kernel:', error);
throw error;
Expand All @@ -860,6 +874,9 @@ export class Kernel {
*/
async stop(): Promise<void> {
await this.#kernelQueue.waitForCrank();
this.#vatManager.abandonRestarts(
new Error('Kernel was stopped; the restart was abandoned'),
);
this.#kernelStore.recordLastActiveTime();
await this.#platformServices.stopRemoteComms();
this.#remoteManager.cleanup();
Expand Down
Loading
Loading