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
334 changes: 333 additions & 1 deletion packages/acp-bridge/src/bridge.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ import { extractErrorMessage, extractErrorCode } from './bridge.js';
import type { ChannelFactory } from './channel.js';
import type { BridgeTelemetry } from './bridgeOptions.js';
import { createInMemoryChannel } from './inMemoryChannel.js';
import type { BridgeEvent } from './eventBus.js';
import { EventBus, type BridgeEvent } from './eventBus.js';
import { ApprovalMode, ShellExecutionService } from '@qwen-code/qwen-code-core';
import {
FakeAgent,
Expand Down Expand Up @@ -1240,6 +1240,255 @@ describe('createAcpSessionBridge', () => {
await bridge.shutdown();
});

it('loads history replay from response metadata when requested', async () => {
const handles: ChannelHandle[] = [];
const factory: ChannelFactory = async () => {
const h = makeChannel({
loadSessionImpl: (p) => {
expect(p._meta).toMatchObject({
'qwen.session.loadReplayMode': 'bulk',
});
return {
_meta: {
keep: 'state',
'qwen.session.loadReplay': {
v: 1,
partial: true,
replayError: 'replay boom',
updates: [
{
sessionUpdate: 'user_message_chunk',
content: { type: 'text', text: 'loaded prompt' },
_meta: { timestamp: 1_700_000_000_000 },
},
{
sessionUpdate: 'agent_message_chunk',
content: { type: 'text', text: 'loaded answer' },
},
],
},
},
};
},
});
handles.push(h);
return h.channel;
};
const bridge = makeBridge({ channelFactory: factory });

const loaded = await bridge.loadSession({
sessionId: 'persisted-bulk-history',
workspaceCwd: WS_A,
historyReplay: 'response',
});

expect(handles[0]?.agent.loadSessionCalls[0]).toMatchObject({
sessionId: 'persisted-bulk-history',
cwd: WS_A,
mcpServers: [],
_meta: { 'qwen.session.loadReplayMode': 'bulk' },
});
expect(loaded.state).toEqual({ _meta: { keep: 'state' } });
expect(loaded.partial).toBe(true);
expect(loaded.replayError).toBe('replay boom');
expect(loaded.lastEventId).toBe(2);
expect(loaded.compactedReplay).toEqual([]);
expect(loaded.liveJournal).toHaveLength(2);
expect(loaded.liveJournal?.[0]?._meta?.['serverTimestamp']).toBe(
1_700_000_000_000,
);

const iterator = bridge
.subscribeEvents(loaded.sessionId, { lastEventId: 0 })
[Symbol.asyncIterator]();
const first = await iterator.next();
expect(first.value).toMatchObject({
type: 'state_resync_required',
data: { reason: 'seeded_replay_not_in_ring' },
});
await iterator.return?.();
await bridge.shutdown();
});

it('rejects oversized response-mode load replay payloads', async () => {
const factory: ChannelFactory = async () =>
makeChannel({
loadSessionImpl: () => ({
_meta: {
'qwen.session.loadReplay': {
v: 1,
updates: Array.from({ length: 10_001 }, () => ({
sessionUpdate: 'agent_message_chunk',
})),
},
},
}),
}).channel;
const bridge = makeBridge({ channelFactory: factory });

await expect(
bridge.loadSession({
sessionId: 'persisted-huge-bulk-history',
workspaceCwd: WS_A,
historyReplay: 'response',
}),
).rejects.toThrow(
'qwen.session.loadReplay updates exceed limit (10001 > 10000)',
);

await bridge.shutdown();
});

it('removes a response-mode restore entry if replay seeding throws', async () => {
const seedSpy = vi
.spyOn(EventBus.prototype, 'seedReplayEvents')
.mockImplementationOnce(() => {
throw new Error('seed boom');
});
const handles: ChannelHandle[] = [];
const factory: ChannelFactory = async () => {
const h = makeChannel({
loadSessionImpl: () => ({
_meta: {
'qwen.session.loadReplay': {
v: 1,
updates: [{ sessionUpdate: 'agent_message_chunk' }],
},
},
}),
});
handles.push(h);
return h.channel;
};
const bridge = makeBridge({ channelFactory: factory });

try {
await expect(
bridge.loadSession({
sessionId: 'persisted-seed-fails',
workspaceCwd: WS_A,
historyReplay: 'response',
}),
).rejects.toThrow('seed boom');
expect(bridge.sessionCount).toBe(0);
expect(handles[0]?.killed).toBe(true);

await expect(
bridge.loadSession({
sessionId: 'persisted-seed-fails',
workspaceCwd: WS_A,
historyReplay: 'response',
}),
).resolves.toMatchObject({
sessionId: 'persisted-seed-fails',
attached: false,
});
expect(
handles.reduce(
(count, handle) => count + handle.agent.loadSessionCalls.length,
0,
),
).toBe(2);
} finally {
seedSpy.mockRestore();
await bridge.shutdown();
}
});

it('reports the bad index when response-mode load replay is invalid', async () => {
const factory: ChannelFactory = async () =>
makeChannel({
loadSessionImpl: () => ({
_meta: {
'qwen.session.loadReplay': {
v: 1,
updates: [
{ sessionUpdate: 'agent_message_chunk' },
{ sessionUpdate: 'future_update_type' },
],
},
},
}),
}).channel;
const bridge = makeBridge({ channelFactory: factory });

await expect(
bridge.loadSession({
sessionId: 'persisted-bad-bulk-history',
workspaceCwd: WS_A,
historyReplay: 'response',
}),
).rejects.toThrow(
/Invalid qwen\.session\.loadReplay update at index 1 .*count=2.*future_update_type/,
);

await bridge.shutdown();
});

it('orders response-mode replay before restore-time buffered events', async () => {
let capturedConn: AgentSideConnection | undefined;
const factory: ChannelFactory = async () => {
const { clientStream, agentStream } = createInMemoryChannel();
const fakeAgent = new FakeAgent({
loadSessionImpl: async (p) => {
void capturedConn!.extNotification(
'qwen/notify/session/mcp-budget-event',
{
v: 1,
sessionId: p.sessionId,
kind: 'budget_warning',
liveCount: 4,
reservedCount: 4,
budget: 4,
thresholdRatio: 0.75,
mode: 'warn',
},
);
await new Promise((r) => setTimeout(r, 20));
return {
_meta: {
'qwen.session.loadReplay': {
v: 1,
updates: [
{
sessionUpdate: 'agent_message_chunk',
content: { type: 'text', text: 'loaded answer' },
},
],
},
},
};
},
});
capturedConn = new AgentSideConnection(() => fakeAgent, agentStream);
return {
stream: clientStream,
exited: new Promise<
| { exitCode: number | null; signalCode: NodeJS.Signals | null }
| undefined
>(() => {}),
kill: async () => {},
killSync: () => {},
};
};
const bridge = makeBridge({ channelFactory: factory });

const loaded = await bridge.loadSession({
sessionId: 'bulk-replay-plus-early-event',
workspaceCwd: WS_A,
historyReplay: 'response',
});

expect(loaded.liveJournal?.map((event) => event.type)).toEqual([
'session_update',
'mcp_budget_warning',
]);
expect(loaded.liveJournal?.[0]?.id).toBe(1);
expect(loaded.liveJournal?.[1]?.id).toBe(2);

await bridge.shutdown();
});

it('buffers load replay events until the restored session is registered', async () => {
let capturedConn: AgentSideConnection | undefined;
const factory: ChannelFactory = async () => {
Expand Down Expand Up @@ -1301,6 +1550,56 @@ describe('createAcpSessionBridge', () => {
await bridge.shutdown();
});

it('keeps streamed load replay if response mode returns no bulk payload', async () => {
let capturedConn: AgentSideConnection | undefined;
const factory: ChannelFactory = async () => {
const { clientStream, agentStream } = createInMemoryChannel();
const fakeAgent = new FakeAgent({
loadSessionImpl: async (p) => {
expect(p._meta).toMatchObject({
'qwen.session.loadReplayMode': 'bulk',
});
await capturedConn!.sessionUpdate({
sessionId: p.sessionId,
update: {
sessionUpdate: 'agent_message_chunk',
content: { type: 'text', text: 'legacy-streamed' },
},
});
return {};
},
});
capturedConn = new AgentSideConnection(() => fakeAgent, agentStream);
return {
stream: clientStream,
exited: new Promise<
| { exitCode: number | null; signalCode: NodeJS.Signals | null }
| undefined
>(() => {}),
kill: async () => {},
killSync: () => {},
};
};
const bridge = makeBridge({ channelFactory: factory });

const loaded = await bridge.loadSession({
sessionId: 'persisted-history-response-fallback',
workspaceCwd: WS_A,
historyReplay: 'response',
});

expect(loaded.liveJournal).toHaveLength(1);
expect(loaded.liveJournal?.[0]?.data).toMatchObject({
sessionId: 'persisted-history-response-fallback',
update: {
sessionUpdate: 'agent_message_chunk',
content: { text: 'legacy-streamed' },
},
});

await bridge.shutdown();
});

it('resumes an existing ACP session without calling session/load', async () => {
const handles: ChannelHandle[] = [];
const factory: ChannelFactory = async () => {
Expand Down Expand Up @@ -1417,6 +1716,39 @@ describe('createAcpSessionBridge', () => {
await bridge.shutdown();
});

it('rejects coalescing load requests with incompatible replay modes', async () => {
let releaseLoad: ((value: LoadSessionResponse) => void) | undefined;
const factory: ChannelFactory = async () =>
makeChannel({
loadSessionImpl: () =>
new Promise<LoadSessionResponse>((resolve) => {
releaseLoad = resolve;
}),
}).channel;
const bridge = makeBridge({ channelFactory: factory });

const first = bridge.loadSession({
sessionId: 'coalesce-replay-mode',
workspaceCwd: WS_A,
});
for (let i = 0; i < 50 && !releaseLoad; i++) {
await new Promise((r) => setTimeout(r, 10));
}
expect(releaseLoad).toBeDefined();

await expect(
bridge.loadSession({
sessionId: 'coalesce-replay-mode',
workspaceCwd: WS_A,
historyReplay: 'response',
}),
).rejects.toBeInstanceOf(RestoreInProgressError);

releaseLoad!({});
await first;
await bridge.shutdown();
});

it('survives spawn-owner disconnect kill while a coalesced restore is mid-flight', async () => {
let releaseLoad: ((value: LoadSessionResponse) => void) | undefined;
const factory: ChannelFactory = async () =>
Expand Down
Loading
Loading