Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
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
48 changes: 48 additions & 0 deletions packages/kernel-test/src/cluster-launch.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { makeSQLKernelDatabase } from '@metamask/kernel-store/sqlite/nodejs';
import { waitUntilQuiescent } from '@metamask/kernel-utils';
import { Logger } from '@metamask/logger';
import type { LogEntry } from '@metamask/logger';
import type { Kernel } from '@metamask/ocap-kernel';
Expand Down Expand Up @@ -123,3 +124,50 @@ describe('cluster initialization', { timeout: 4_000 }, () => {
]);
});
});

describe('peer rejection propagation', { timeout: 10_000 }, () => {
let logger: Logger;
let entries: LogEntry[];
let kernel: Kernel;

beforeEach(async () => {
const testLogger = makeTestLogger();
logger = testLogger.logger;
entries = testLogger.entries;
const database = await makeSQLKernelDatabase({});
kernel = await makeKernel(
database,
true,
logger.subLogger({ tags: ['test'] }),
);
});

it('bootstrap observes peer rejection when a peer vat fails to launch', async () => {
await expect(
kernel.launchSubcluster({
bootstrap: 'main',
vats: {
main: {
bundleSpec: getBundleSpec('peer-rejection-bootstrap'),
parameters: {},
},
peer: {
bundleSpec: getBundleSpec('error-build-throw'),
parameters: {},
},
},
}),
).rejects.toMatchObject({
message: expect.stringMatching(/^Failed to launch vat \S+ \(peer\)$/u),
});

// Let the kernel run loop deliver the parked bootstrap message to the
// main vat, which will observe the peer's rejected root promise.
await waitUntilQuiescent(200);

const vatLogs = extractTestLogs(entries, 'console');
expect(vatLogs).toContainEqual(
expect.stringMatching(/^peer rejected:.*VAT_TERMINATED/u),
);
});
});
14 changes: 10 additions & 4 deletions packages/kernel-test/src/persistence.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -96,10 +96,15 @@ describe('persistent storage', { timeout: 20_000 }, () => {
false,
logger.logger.subLogger({ tags: ['test'] }),
);
const result1 = await runTestVats(kernel1, multiVatCluster);
expect(result1).toBe('Coordinator initialized with 2 workers');
// Capture rootKref directly: concurrent vat launch means the coordinator
// may not be assigned ko4, so we cannot use a hardcoded ref here.
const { bootstrapResult: launch1Result, rootKref: coordinatorRoot } =
await kernel1.launchSubcluster(multiVatCluster);
await waitUntilQuiescent();
const workResult1 = await runResume(kernel1, v1Root);
expect(kunser(launch1Result as CapData<string>)).toBe(
'Coordinator initialized with 2 workers',
);
const workResult1 = await runResume(kernel1, coordinatorRoot);
expect(workResult1).toBe('Work completed: Worker1(1), Worker2(1)');
await waitUntilQuiescent();
await kernel1.stop();
Expand All @@ -110,7 +115,8 @@ describe('persistent storage', { timeout: 20_000 }, () => {
logger.logger.subLogger({ tags: ['test'] }),
);
await new Promise((resolve) => setTimeout(resolve, 1000));
const workResult2 = await runResume(kernel2, v1Root);
// coordinatorRoot (ko<N>) is stable across kernel restarts.
const workResult2 = await runResume(kernel2, coordinatorRoot);
expect(workResult2).toBe('Work completed: Worker1(2), Worker2(2)');
await kernel2.stop();
});
Expand Down
2 changes: 1 addition & 1 deletion packages/kernel-test/src/rejection.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ describe('rejection', () => {
});
expect(vat).toBeDefined();
const vats = kernel.getVatIds();
expect(vats).toStrictEqual(vatIds);
expect([...vats].sort()).toStrictEqual([...vatIds].sort());

await waitUntilQuiescent();
const vatLogs = vatIds.map((vatId) => extractTestLogs(entries, vatId));
Expand Down
29 changes: 18 additions & 11 deletions packages/kernel-test/src/resume.test.ts
Original file line number Diff line number Diff line change
@@ -1,14 +1,14 @@
import type { CapData } from '@endo/marshal';
import { makeSQLKernelDatabase } from '@metamask/kernel-store/sqlite/nodejs';
import { waitUntilQuiescent } from '@metamask/kernel-utils';
import type { KRef } from '@metamask/ocap-kernel';
import { kunser } from '@metamask/ocap-kernel';
import { describe, expect, it } from 'vitest';

import {
getBundleSpec,
makeKernel,
makeTestLogger,
runResume,
runTestVats,
sortLogs,
extractTestLogs,
} from './utils.ts';
Expand Down Expand Up @@ -100,21 +100,22 @@ const reference = sortLogs([
...carolResumeReference,
]);

// Vat root objects start with ko4 due to the kernel facet and other kernel service objects being created before any vats.
const v1Root: KRef = 'ko4';
const v2Root: KRef = 'ko5';
const v3Root: KRef = 'ko6';

describe('restarting vats', async () => {
it('exercise restart vats individually', async () => {
const kernelDatabase = await makeSQLKernelDatabase({
dbFilename: ':memory:',
});
const { logger, entries } = makeTestLogger();
const kernel = await makeKernel(kernelDatabase, true, logger);
const bootstrapResult = await runTestVats(kernel, testSubcluster);
expect(bootstrapResult).toBe('bootstrap Alice');
// Use launchSubcluster directly to get vatRootKrefs: concurrent vat launch
// means ko<N> assignment order depends on worker startup speed.
const { bootstrapResult, vatRootKrefs } =
await kernel.launchSubcluster(testSubcluster);
await waitUntilQuiescent();
expect(kunser(bootstrapResult as CapData<string>)).toBe('bootstrap Alice');
const v1Root = vatRootKrefs.alice;
const v2Root = vatRootKrefs.bob;
const v3Root = vatRootKrefs.carol;
await kernel.restartVat('v1');
await kernel.restartVat('v2');
await kernel.restartVat('v3');
Expand All @@ -136,9 +137,15 @@ describe('restarting vats', async () => {
});
const { logger: logger1, entries: entries1 } = makeTestLogger();
const kernel1 = await makeKernel(kernelDatabase, true, logger1);
const bootstrapResult = await runTestVats(kernel1, testSubcluster);
expect(bootstrapResult).toBe('bootstrap Alice');
// Capture vatRootKrefs from first kernel: ko<N> refs are stable across
// kernel restarts because they are persisted in the kernel store.
const { bootstrapResult, vatRootKrefs } =
await kernel1.launchSubcluster(testSubcluster);
await waitUntilQuiescent();
expect(kunser(bootstrapResult as CapData<string>)).toBe('bootstrap Alice');
const v1Root = vatRootKrefs.alice;
const v2Root = vatRootKrefs.bob;
const v3Root = vatRootKrefs.carol;
const { logger: logger2, entries: entries2 } = makeTestLogger();
const kernel2 = await makeKernel(kernelDatabase, false, logger2);
await new Promise((resolve) => setTimeout(resolve, 1000));
Expand Down
28 changes: 28 additions & 0 deletions packages/kernel-test/src/vats/peer-rejection-bootstrap.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
import { E } from '@endo/eventual-send';
import { makeDefaultExo } from '@metamask/kernel-utils/exo';

/**
* Bootstrap vat for testing peer-rejection propagation.
* Receives a `peer` root reference that may be a rejected kernel promise
* (e.g. if the peer vat failed to launch), and logs whether calls resolve
* or reject so integration tests can inspect the outcome.
*
* @returns The root object for this vat.
*/
// eslint-disable-next-line @typescript-eslint/explicit-function-return-type
export function buildRootObject() {
return makeDefaultExo('root', {
async bootstrap({ peer }: { peer: unknown }) {
await E(peer as object)
.ping()
// eslint-disable-next-line no-console
.then(() => console.log('peer resolved'))
.catch((error: unknown) => {
const message =
error instanceof Error ? error.message : String(error);
// eslint-disable-next-line no-console
console.log(`peer rejected: ${message}`);
});
},
});
}
2 changes: 2 additions & 0 deletions packages/ocap-kernel/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Added

- Launch all vats in a subcluster concurrently during `launchSubcluster`, reducing startup latency from serial to parallel; failed peer vats receive a rejected kernel promise observable via `E(roots.peer).method()` pipelining ([#983](https://github.com/MetaMask/ocap-kernel/pull/983))

- Add `fetch`, `Request`, `Headers`, and `Response` to available vat endowments ([#942](https://github.com/MetaMask/ocap-kernel/pull/942))
- Add `VatConfig.network: { allowedHosts: string[] }`; requesting `'fetch'` without it rejects `initVat`
- Integrate Snaps attenuated endowment factories into vat globals ([#937](https://github.com/MetaMask/ocap-kernel/pull/937))
Expand Down
3 changes: 3 additions & 0 deletions packages/ocap-kernel/src/Kernel.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,8 @@ const mocks = vi.hoisted(() => {
}

waitForCrank = vi.fn().mockResolvedValue(undefined);

resolvePromises = vi.fn();
}

class RemoteManager {
Expand Down Expand Up @@ -299,6 +301,7 @@ describe('Kernel', () => {
subclusterId: 's1',
bootstrapResult: { body: '{"result":"ok"}', slots: [] },
rootKref: expect.stringMatching(/^ko\d+$/u),
vatRootKrefs: { alice: expect.stringMatching(/^ko\d+$/u) },
});
});
});
Expand Down
9 changes: 6 additions & 3 deletions packages/ocap-kernel/src/store/methods/vat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ export function getVatMethods(ctx: StoreContext) {
getKernelPromise,
addPromiseSubscriber,
} = getPromiseMethods(ctx);
const { initKernelObject } = getObjectMethods(ctx);
const { initKernelObject, getObjectRefCount } = getObjectMethods(ctx);
const { addCListEntry } = getCListMethods(ctx);
const { incrementRefCount, decrementRefCount } = getRefCountMethods(ctx);

Expand Down Expand Up @@ -261,8 +261,11 @@ export function getVatMethods(ctx: StoreContext) {
const { vatSlot } = getReachableAndVatSlot(vatID, kref);
ctx.kv.delete(getSlotKey(vatID, kref));
ctx.kv.delete(getSlotKey(vatID, vatSlot));
// Decrease refcounts that belonged to the terminating vat
decrementRefCount(kref, 'cleanup|export|baseline');
// Skip baseline decrement if GC already zeroed reachable via dropImports.
const { reachable } = getObjectRefCount(kref);
if (reachable > 0) {
decrementRefCount(kref, 'cleanup|export|baseline');
}
ctx.maybeFreeKrefs.add(kref);
work.exports += 1;
}
Expand Down
2 changes: 2 additions & 0 deletions packages/ocap-kernel/src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -746,6 +746,8 @@ export type SubclusterLaunchResult = {
rootKref: KRef;
/** The CapData result of calling bootstrap() on the root object, if any. */
bootstrapResult: CapData<KRef> | undefined;
/** Map from vat name to root kref for all successfully launched vats. */
vatRootKrefs: Record<string, KRef>;
};

const RemoteCommsDisconnectedStruct = object({
Expand Down
60 changes: 59 additions & 1 deletion packages/ocap-kernel/src/vats/SubclusterManager.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ import type { Mocked } from 'vitest';
import { describe, it, expect, vi, beforeEach } from 'vitest';

import type { KernelQueue } from '../KernelQueue.ts';
import { kser } from '../liveslots/kernel-marshal.ts';
import { kser, makeKernelError } from '../liveslots/kernel-marshal.ts';
import type { KernelStore } from '../store/index.ts';
import type {
VatId,
Expand Down Expand Up @@ -65,10 +65,15 @@ describe('SubclusterManager', () => {
getRootObject: vi.fn(),
deleteVatConfig: vi.fn(),
markVatAsTerminated: vi.fn(),
initKernelPromise: vi
.fn()
.mockReturnValue(['kp1', { state: 'unresolved', subscribers: [] }]),
setPromiseDecider: vi.fn(),
} as unknown as Mocked<KernelStore>;

mockKernelQueue = {
waitForCrank: vi.fn().mockResolvedValue(undefined),
resolvePromises: vi.fn(),
} as unknown as Mocked<KernelQueue>;

mockVatManager = {
Expand Down Expand Up @@ -116,14 +121,17 @@ describe('SubclusterManager', () => {
'testVat',
's1',
);
// queueMessage targets the real root KRef once all vats have launched
expect(mockQueueMessage).toHaveBeenCalledWith('ko1', 'bootstrap', [
{ testVat: expect.anything() },
{},
]);
expect(mockKernelQueue.resolvePromises).not.toHaveBeenCalled();
expect(result).toStrictEqual({
subclusterId: 's1',
rootKref: 'ko1',
bootstrapResult: { body: '{"result":"ok"}', slots: [] },
vatRootKrefs: { testVat: 'ko1' },
});
});

Expand Down Expand Up @@ -152,6 +160,43 @@ describe('SubclusterManager', () => {
'bob',
's1',
);
// bootstrap receives real ko<N> refs for both vats; no kernel promises allocated
expect(mockQueueMessage).toHaveBeenCalledWith('ko1', 'bootstrap', [
{ alice: expect.anything(), bob: expect.anything() },
{},
]);
expect(mockKernelQueue.resolvePromises).not.toHaveBeenCalled();
});

it('rejects peer kernel promise when a non-bootstrap vat fails to launch', async () => {
const config: ClusterConfig = {
bootstrap: 'alice',
vats: {
alice: { sourceSpec: 'alice.js' },
bob: { sourceSpec: 'bob.js' },
},
};
const bobError = new Error('bob exploded');
mockVatManager.launchVat
.mockResolvedValueOnce('ko1' as KRef)
.mockRejectedValueOnce(bobError);

await expect(subclusterManager.launchSubcluster(config)).rejects.toThrow(
'bob exploded',
);

// bootstrap receives alice's real ko1 root ref as queueMessage target
expect(mockQueueMessage).toHaveBeenCalledWith('ko1', 'bootstrap', [
{ alice: expect.anything(), bob: expect.anything() },
{},
]);
// initKernelPromise called once — only for bob's rejected promise
expect(mockKernelStore.initKernelPromise).toHaveBeenCalledTimes(1);
// bob's promise rejected with a VAT_TERMINATED kernel error
expect(mockKernelQueue.resolvePromises).toHaveBeenCalledTimes(1);
expect(mockKernelQueue.resolvePromises).toHaveBeenCalledWith('kernel', [
['kp1', true, makeKernelError('VAT_TERMINATED', 'bob exploded')],
]);
});

it('includes unrestricted kernel services when specified', async () => {
Expand Down Expand Up @@ -303,6 +348,12 @@ describe('SubclusterManager', () => {
});

mockVatManager.launchVat.mockRejectedValue(new Error('vat boom'));
// Service lookup now happens before vat launch, so the IO channel
// service must be registered for the test to reach the vat launch step.
(mockGetKernelService as ReturnType<typeof vi.fn>).mockReturnValue({
kref: 'ko99',
systemOnly: false,
});

const config: ClusterConfig = {
bootstrap: 'testVat',
Expand Down Expand Up @@ -365,6 +416,12 @@ describe('SubclusterManager', () => {
});

mockVatManager.launchVat.mockRejectedValue(new Error('launch boom'));
// Service lookup now happens before vat launch, so the IO channel
// service must be registered for the test to reach the vat launch step.
(mockGetKernelService as ReturnType<typeof vi.fn>).mockReturnValue({
kref: 'ko99',
systemOnly: false,
});

const config: ClusterConfig = {
bootstrap: 'testVat',
Expand Down Expand Up @@ -430,6 +487,7 @@ describe('SubclusterManager', () => {
subclusterId: 's1',
rootKref: 'ko1',
bootstrapResult,
vatRootKrefs: { testVat: 'ko1' },
});
});

Expand Down
Loading
Loading