Skip to content
Merged
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
137 changes: 111 additions & 26 deletions apps/desktop/src/main/__tests__/runtime-host-observation-preload.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,30 +18,127 @@
*/

import assert from 'node:assert/strict';
import { EventEmitter } from 'node:events';
import { createRequire } from 'node:module';
import { fileURLToPath } from 'node:url';
import { runInNewContext } from 'node:vm';
import test from 'node:test';
import { deferred } from '@maka/core/test-only/async-primitives';
import { build } from 'esbuild';
import type { MakaBridge } from '../../preload/bridge-contract.js';

const owner = {
hostId: 'owner-host', targetEpoch: 'owner-epoch', profileId: 'local',
profileName: 'Local', profileKind: 'local', profileAccess: 'owner', readiness: 'ready',
};

test('observation readiness includes its active seed even when the invoke reply overtakes event IPC', async () => {
const owner = {
hostId: 'owner-host', targetEpoch: 'owner-epoch', profileId: 'local',
profileName: 'Local', profileKind: 'local', profileAccess: 'owner', readiness: 'ready',
};
// Deliberately never deliver event IPC: only the invoke response arrives.
const { bridge } = await preloadHarness(async (channel) => {
if (channel === 'sessions:observe') return {
kind: 'ready',
value: [{
type: 'text_delta', id: 'seed-1', turnId: 'turn-1', messageId: 'message-1',
ts: 1, startOffset: 0, text: 'All output accumulated while away',
}],
};
throw new Error('Unexpected channel: ' + channel);
});
const order: string[] = [];
let unsubscribe = () => {};
await new Promise<void>((resolve, reject) => {
unsubscribe = bridge.sessions.subscribeEvents(
JSON.stringify([owner.hostId, 'session-1']),
(event) => { if (event.type === 'text_delta') order.push(event.text); },
() => { order.push('ready'); resolve(); },
undefined,
reject,
);
});
assert.deepEqual(order, ['All output accumulated while away', 'ready']);
unsubscribe();
});

test('cancelled Session observation removes preload listeners without publishing readiness or errors', async () => {
const started = deferred<void>();
const observation = deferred<{ kind: 'cancelled' }>();
const { bridge, events } = await preloadHarness(async (channel) => {
if (channel === 'sessions:observe') {
started.resolve();
return observation.promise;
}
throw new Error('Unexpected channel: ' + channel);
});
const callbacks: string[] = [];
const unsubscribe = bridge.sessions.subscribeEvents(
JSON.stringify([owner.hostId, 'session-1']),
() => callbacks.push('event'),
() => callbacks.push('ready'),
() => callbacks.push('seed'),
() => callbacks.push('error'),
);
try {
await started.promise;
assert.equal(events.listenerCount('sessions:event:session-1'), 1);
assert.equal(events.listenerCount('sessions:observation-seed'), 1);
observation.resolve({ kind: 'cancelled' });
await new Promise<void>((resolve) => setImmediate(resolve));

assert.equal(events.listenerCount('sessions:event:session-1'), 0);
assert.equal(events.listenerCount('sessions:observation-seed'), 0);
events.emit('sessions:event:session-1', {}, owner, {
type: 'text_delta', id: 'late-1', turnId: 'turn-1', messageId: 'message-1',
ts: 1, startOffset: 0, text: 'Late output',
});
events.emit('sessions:observation-seed', {}, owner, { sessionId: 'session-1', phase: 'ready' });
assert.deepEqual(callbacks, []);
} finally {
observation.resolve({ kind: 'cancelled' });
unsubscribe();
}
});

test('cancelled transcript open rejects and removes its preload listener', async () => {
const started = deferred<string>();
const transcript = deferred<{ kind: 'cancelled' }>();
const { bridge, events } = await preloadHarness(async (channel, ...args) => {
if (channel === 'session-local:transcript') return null;
if (channel === 'sessions:transcript:open') {
started.resolve(`sessions:transcript:${args[2]}`);
return transcript.promise;
}
if (channel === 'sessions:transcript:close') return;
throw new Error('Unexpected channel: ' + channel);
});
let cancel = () => {};
const opening = bridge.transcripts.open(
JSON.stringify([owner.hostId, 'session-1']),
() => assert.fail('A cancelled transcript must not deliver a batch'),
(requestCancellation) => { cancel = requestCancellation; },
);
try {
const channel = await started.promise;
assert.equal(events.listenerCount(channel), 1);
const rejection = assert.rejects(opening, /Desktop transcript open was cancelled/);
transcript.resolve({ kind: 'cancelled' });
await rejection;
assert.equal(events.listenerCount(channel), 0);
} finally {
transcript.resolve({ kind: 'cancelled' });
cancel();
await opening.catch(() => undefined);
}
});

async function preloadHarness(invoke: (channel: string, ...args: unknown[]) => Promise<unknown>) {
const events = new EventEmitter();
const ipcRenderer = {
// Deliberately never deliver event IPC: only the invoke response arrives.
on() {}, off() {}, send() {},
async invoke(channel: string) {
on: events.on.bind(events), off: events.off.bind(events), send() {},
async invoke(channel: string, ...args: unknown[]) {
if (channel === 'runtime-host:activeIdentity') return owner;
if (channel === 'runtime-host:identities') return [owner];
if (channel === 'sessions:unobserve') return;
if (channel === 'sessions:observe') return [{
type: 'text_delta', id: 'seed-1', turnId: 'turn-1', messageId: 'message-1',
ts: 1, startOffset: 0, text: 'All output accumulated while away',
}];
throw new Error('Unexpected channel: ' + channel);
return invoke(channel, ...args);
},
};
let bridge: MakaBridge | undefined;
Expand All @@ -61,17 +158,5 @@ test('observation readiness includes its active seed even when the invoke reply
crypto: globalThis.crypto,
});
assert.ok(bridge);
const order: string[] = [];
let unsubscribe = () => {};
await new Promise<void>((resolve, reject) => {
unsubscribe = bridge!.sessions.subscribeEvents(
JSON.stringify([owner.hostId, 'session-1']),
(event) => { if (event.type === 'text_delta') order.push(event.text); },
() => { order.push('ready'); resolve(); },
undefined,
reject,
);
});
assert.deepEqual(order, ['All output accumulated while away', 'ready']);
unsubscribe();
});
return { bridge, events };
}
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import { join } from "node:path";
import test from "node:test";
import type { IpcMain } from "electron";
import { WORKHUB_COORDINATION_SESSION_ID } from '@maka/core/session';
import { deferred } from '@maka/core/test-only/async-primitives';
import { SIDE_CONVERSATION_SESSION_LABEL } from '@maka/core/side-conversation';
import { type AttachmentRef } from '@maka/core/events';
import {
Expand Down Expand Up @@ -61,7 +62,7 @@ test('registers Session observation as one reconnectable operation', () => {

for (const phase of ['connecting', 'seeding'] as const) {
for (const cancellation of ['unobserve', 'renderer destruction'] as const) {
test(`Session observation IPC completes normally after ${cancellation} while ${phase}`, async () => {
test(`Session observation IPC returns cancellation after ${cancellation} while ${phase}`, async () => {
const errors: unknown[] = [];
const observations = new RuntimeHostSessionObservationRegistry((error) => errors.push(error));
const ipc = ipcHarness();
Expand All @@ -85,7 +86,7 @@ for (const phase of ['connecting', 'seeding'] as const) {
else ipc.rendererDestroyed();
finishSeed();

assert.deepEqual(await observing, []);
assert.deepEqual(await observing, { kind: 'cancelled' });
assert.deepEqual(observations.observedSessionIds(), []);
assert.deepEqual(await observations.attach(source), []);
assert.equal(seeds, phase === 'seeding' ? 1 : 0);
Expand Down Expand Up @@ -116,6 +117,195 @@ test('forward transcript paging is an observation operation scoped to the render
assert.equal(calls.length, 1);
});

test('treats pending Session observation teardown as IPC cancellation', async () => {
const observations = new RuntimeHostSessionObservationRegistry();
const started = deferred();
const finishInitialization = deferred();
await observations.attach({
async observe() {
started.resolve();
await finishInitialization.promise;
},
async unobserve() {},
});
const ipc = observationIpcHarness(observations);

const observing = ipc.invoke('sessions:observe', 'session-1', 'observer-1');
try {
await started.promise;
assert.deepEqual(observations.trackedSessionIds(), ['session-1']);
await observations.unobserve('observer-1');
assert.deepEqual(observations.trackedSessionIds(), []);
finishInitialization.resolve();
assert.deepEqual(await observing, { kind: 'cancelled' });
} finally {
finishInitialization.resolve();
await observations.close();
}
});

test('treats pending transcript teardown as IPC cancellation', async () => {
const observations = new RuntimeHostSessionObservationRegistry();
const started = deferred();
const finishInitialization = deferred();
await observations.attach({
async observe() {},
async unobserve() {},
async openTranscript() {
started.resolve();
await finishInitialization.promise;
return {
sessionId: 'session-1',
generation: 'generation-1',
hostEpoch: 'host-epoch-1',
readThroughMessageId: null,
};
},
async loadTranscriptBefore() {},
async loadTranscriptAround() {},
async loadTranscriptAfter() {},
async closeTranscript() {},
});
const ipc = observationIpcHarness(observations);

const opening = ipc.invoke('sessions:transcript:open', 'session-1', 'consumer-1');
try {
await started.promise;
assert.deepEqual(observations.trackedSessionIds(), ['session-1']);
await observations.closeTranscript('consumer-1', 9);
assert.deepEqual(observations.trackedSessionIds(), []);
finishInitialization.resolve();
assert.deepEqual(await opening, { kind: 'cancelled' });
} finally {
finishInitialization.resolve();
await observations.close();
}
});

for (const teardown of ['forgetSession', 'close'] as const) {
test(`${teardown} rejects pending observation IPC instead of silently cancelling`, async () => {
const observations = new RuntimeHostSessionObservationRegistry();
const sessionStarted = deferred();
const transcriptStarted = deferred();
const finishInitialization = deferred();
await observations.attach({
async observe() {
sessionStarted.resolve();
await finishInitialization.promise;
},
async unobserve() {},
async openTranscript() {
transcriptStarted.resolve();
await finishInitialization.promise;
return {
sessionId: 'session-1',
generation: 'generation-1',
hostEpoch: 'host-epoch-1',
readThroughMessageId: null,
};
},
async loadTranscriptBefore() {},
async loadTranscriptAround() {},
async loadTranscriptAfter() {},
async closeTranscript() {},
});
const ipc = observationIpcHarness(observations);
const observing = ipc.invoke('sessions:observe', 'session-1', 'observer-1');
const opening = ipc.invoke('sessions:transcript:open', 'session-1', 'consumer-1');
// Attach rejection handlers before teardown. The late source completion
// must not turn the lost observation into readiness or silent cancellation.
const results = Promise.allSettled([observing, opening]);
try {
await Promise.all([sessionStarted.promise, transcriptStarted.promise]);
if (teardown === 'forgetSession') await observations.forgetSession('session-1');
else await observations.close();
assert.deepEqual(observations.trackedSessionIds(), []);
finishInitialization.resolve();
for (const result of await results) {
assert.equal(result.status, 'rejected');
if (result.status === 'rejected') assert.ok(result.reason instanceof Error);
}
} finally {
finishInitialization.resolve();
await observations.close();
}
});
}

test('preserves genuine Session observation initialization failures', async () => {
const observations = new RuntimeHostSessionObservationRegistry();
const sessionFailure = new Error('seed failed');
const transcriptFailure = new Error('transcript open failed');
await observations.attach({
async observe() {
throw sessionFailure;
},
async unobserve() {},
async openTranscript() {
throw transcriptFailure;
},
async loadTranscriptBefore() {},
async loadTranscriptAround() {},
async loadTranscriptAfter() {},
async closeTranscript() {},
});
const ipc = observationIpcHarness(observations);

await assert.rejects(
ipc.invoke('sessions:observe', 'session-1', 'observer-1'),
(error) => error === sessionFailure,
);
await assert.rejects(
ipc.invoke('sessions:transcript:open', 'session-1', 'consumer-1'),
(error) => error === transcriptFailure,
);
await observations.close();
});

test('returns explicit ready results for Session observation IPC', async () => {
const observations = new RuntimeHostSessionObservationRegistry();
const transcript = {
sessionId: 'session-1',
generation: 'generation-1',
hostEpoch: 'host-epoch-1',
readThroughMessageId: null,
};
await observations.attach({
async observe() {},
async unobserve() {},
async openTranscript() {
return transcript;
},
async loadTranscriptBefore() {},
async loadTranscriptAround() {},
async loadTranscriptAfter() {},
async closeTranscript() {},
});
const ipc = observationIpcHarness(observations);

assert.deepEqual(await ipc.invoke('sessions:observe', 'session-1', 'observer-1'), {
kind: 'ready',
value: [],
});
assert.deepEqual(await ipc.invoke('sessions:transcript:open', 'session-1', 'consumer-1'), {
kind: 'ready',
value: transcript,
});
await observations.close();
});

function observationIpcHarness(observations: RuntimeHostSessionObservationRegistry) {
const ipc = ipcHarness();
registerRuntimeHostSessionObservationIpc(
{
observations,
resolveSideConversation: async () => false,
},
ipc,
);
return ipc;
}

test("keeps synthetic E2E interactions visible through Host hydration and retires their answer", async () => {
const observer = observerWithSnapshot();
const ipc = ipcHarness();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -353,7 +353,7 @@ test('WorkHub tail navigation converges through the preload with a fragmented sp
for (const batch of encodeDesktopTranscriptSnapshot({ ...snapshot, navigationVersion: 0, durable: [] })) {
deliver(batch);
}
return { ...snapshot, readThroughMessageId: null };
return { kind: 'ready', value: { ...snapshot, readThroughMessageId: null } };
}
if (channel === 'sessions:transcript:load-around') {
const request = args[1] as DesktopTranscriptRangeRequest;
Expand Down
Loading