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
98 changes: 98 additions & 0 deletions packages/acp-bridge/src/bridge.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ import {
extractErrorMessage,
extractErrorCode,
} from './bridge.js';
import { NdJsonQueueLimitError } from './ndJsonStream.js';
import {
BridgeChannelClosedError,
BridgeTimeoutError,
Expand Down Expand Up @@ -12748,6 +12749,10 @@ describe('createAcpSessionBridge', () => {
failures[1]!.resolve(
Object.assign(new Error('bounded transport queue failed'), {
code: 'ndjson_queue_limit_exceeded',
budget: 'INVALID budget <script>',
maxQueuedBytes: 1.9,
requiredBytes: -5,
availableBytes: Number.NaN,
}),
);
await expect(restore).rejects.toBeInstanceOf(BridgeChannelClosedError);
Expand All @@ -12764,10 +12769,103 @@ describe('createAcpSessionBridge', () => {
'ndjson_frame_too_large',
}),
);
expect(event).toHaveBeenCalledWith(
'channel.exited',
expect.objectContaining({
'qwen-code.daemon.channel.transport_error_code':
'ndjson_queue_limit_exceeded',
'qwen-code.daemon.channel.transport_error_detail':
'unknown:required=0:available=?:cap=1',
}),
);
});
const frameExitAttributes = event.mock.calls.find(
([name, attributes]) =>
name === 'channel.exited' &&
(attributes as Record<string, unknown>)[
'qwen-code.daemon.channel.transport_error_code'
] === 'ndjson_frame_too_large',
)?.[1] as Record<string, unknown> | undefined;
expect(frameExitAttributes).not.toHaveProperty(
'qwen-code.daemon.channel.transport_error_detail',
);
bridge.killAllSync();
});

it('records which budget fired on a transport-guard channel exit', async () => {
const neverPrompt = deferred<never>();
const handle = makeChannel({ promptImpl: () => neverPrompt.promise });
const failure = deferred<unknown>();
handle.channel = {
...handle.channel,
transportFailed: failure.promise,
kill: () => deferred<never>().promise,
};
const event = vi.fn();
const channelLifecycle = vi.fn();
const telemetry: BridgeTelemetry = {
captureContext: () => undefined,
runWithContext: async (_captured, fn) => await fn(),
withSpan: async (_operation, _attributes, fn) => await fn(),
event,
injectPromptContext: (request) => request,
metrics: {
sessionLifecycle: vi.fn(),
channelLifecycle,
promptQueueWait: vi.fn(),
promptDuration: vi.fn(),
cancelled: vi.fn(),
},
};
const bridge = makeBridge({
channelFactory: async () => handle.channel,
telemetry,
});
const stderr = vi
.spyOn(process.stderr, 'write')
.mockImplementation(() => true);
try {
const first = await bridge.spawnOrAttach({ workspaceCwd: WS_A });
const prompt = bridge.sendPrompt(first.sessionId, {
sessionId: first.sessionId,
prompt: [{ type: 'text', text: 'stay pending' }],
});
await vi.waitFor(() => expect(handle.agent.promptCalls).toHaveLength(1));
failure.resolve(
new NdJsonQueueLimitError(
'prepared_response',
256,
67108864,
262144,
0,
),
);
await expect(prompt).rejects.toBeInstanceOf(BridgeChannelClosedError);
handle.crash({ exitCode: null, signalCode: 'SIGTERM' });
await vi.waitFor(() =>
expect(event).toHaveBeenCalledWith(
'channel.exited',
expect.objectContaining({
'qwen-code.daemon.channel.transport_failed': true,
'qwen-code.daemon.channel.transport_failure_initiated_teardown': true,
'qwen-code.daemon.channel.transport_error_code':
'ndjson_queue_limit_exceeded',
'qwen-code.daemon.channel.transport_error_detail':
'prepared_response:required=262144:available=0:cap=67108864',
}),
),
);
expect(stderr).toHaveBeenCalledWith(
expect.stringContaining(
'transport_detail=prepared_response:required=262144:available=0:cap=67108864',
),
);
} finally {
stderr.mockRestore();
await bridge.shutdown();
}
});

it('does not publish a channel whose transport fails during initialize', async () => {
const failure = deferred<unknown>();
const transportError = Object.assign(new Error('frame failed'), {
Expand Down
41 changes: 40 additions & 1 deletion packages/acp-bridge/src/bridge.ts
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,7 @@ import {
StandaloneSessionSpawnError,
} from './bridgeErrors.js';
import type { BridgeChannelUnavailableReason } from './bridgeErrors.js';
import type { NdJsonQueueLimitError } from './ndJsonStream.js';
import {
resolveSessionRestoreTimeoutMs,
restoreRetryAfterSeconds,
Expand Down Expand Up @@ -319,6 +320,31 @@ function safeTransportFailureCode(error: unknown): string | undefined {
: undefined;
}

function safeTransportFailureDetail(error: unknown): string | undefined {
if (!isRecord(error) || error['code'] !== 'ndjson_queue_limit_exceeded') {
return undefined;
}
Comment thread
yiliang114 marked this conversation as resolved.
const queueError = error as Partial<NdJsonQueueLimitError>;
const budget =
typeof queueError.budget === 'string' &&
/^[a-z0-9_.-]{1,64}$/iu.test(queueError.budget)
? queueError.budget
: 'unknown';
const numbers: string[] = [];
for (const value of [
queueError.requiredBytes,
queueError.availableBytes,
queueError.maxQueuedBytes,
]) {
Comment thread
yiliang114 marked this conversation as resolved.
numbers.push(
typeof value === 'number' && Number.isFinite(value)
? String(Math.max(0, Math.floor(value)))
: '?',
);
}
return `${budget}:required=${numbers[0]}:available=${numbers[1]}:cap=${numbers[2]}`;
Comment thread
yiliang114 marked this conversation as resolved.
}

function sessionSourceRequestMeta(
sourceType: string | undefined,
sourceId: string | undefined,
Expand Down Expand Up @@ -962,6 +988,12 @@ interface ChannelInfo {
transportFailureInitiatedTeardown: boolean;
/** Safe bounded code retained for telemetry; never the raw error message. */
transportFailureCode?: string;
/**
* Bounded queue-budget detail for `ndjson_queue_limit_exceeded` transport
* failures: which budget fired plus required/available/cap bytes. Derived
* from typed error fields only; never the raw error message.
*/
transportFailureDetail?: string;
/**
* Cached channel-close race for workspace-scoped status requests. Workspace
* status can be polled frequently by dashboards, so keep one promise per
Expand Down Expand Up @@ -4601,6 +4633,7 @@ export function createAcpSessionBridge(opts: BridgeOptions): AcpSessionBridge {
}
info.transportFailed = true;
info.transportFailureCode = safeTransportFailureCode(error);
info.transportFailureDetail = safeTransportFailureDetail(error);
info.isDying = true;
info.channelLiveness?.stop();
clearInFlightExtensionRefreshes(info.connection);
Expand Down Expand Up @@ -4707,12 +4740,18 @@ export function createAcpSessionBridge(opts: BridgeOptions): AcpSessionBridge {
info.transportFailureCode,
}
: {}),
...(info.transportFailureDetail
? {
'qwen-code.daemon.channel.transport_error_detail':
info.transportFailureDetail,
}
: {}),
...(exitInfo?.signalCode
? { 'qwen-code.daemon.channel.signal': exitInfo.signalCode }
: {}),
});
writeStderrLine(
`qwen serve: channel exited (code=${exitInfo?.exitCode ?? 'none'}, signal=${exitInfo?.signalCode ?? 'none'}, transport=${info.transportFailed ? (info.transportFailureCode ?? 'failed') : 'ok'}, ${sessions.length} session(s) torn down)`,
`qwen serve: channel exited (code=${exitInfo?.exitCode ?? 'none'}, signal=${exitInfo?.signalCode ?? 'none'}, transport=${info.transportFailed ? (info.transportFailureCode ?? 'failed') : 'ok'}${info.transportFailureDetail ? `, transport_detail=${info.transportFailureDetail}` : ''}, ${sessions.length} session(s) torn down)`,
Comment thread
yiliang114 marked this conversation as resolved.
);
}
for (const sid of sessions) {
Expand Down
15 changes: 12 additions & 3 deletions packages/acp-bridge/src/ndJsonStream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -406,7 +406,10 @@ describe('ndJsonStream', () => {
);
await vi.waitFor(() =>
expect(onTransportError).toHaveBeenCalledWith(
expect.any(NdJsonQueueLimitError),
expect.objectContaining({
code: 'ndjson_queue_limit_exceeded',
budget: 'decoded',
}),
),
);
expect(onTransportError).toHaveBeenCalledOnce();
Expand Down Expand Up @@ -767,7 +770,10 @@ describe('ndJsonStream', () => {

await vi.waitFor(() =>
expect(onTransportError).toHaveBeenCalledWith(
expect.any(NdJsonQueueLimitError),
expect.objectContaining({
code: 'ndjson_queue_limit_exceeded',
budget: 'inbound_request',
}),
),
);
clearInterval(timer);
Expand Down Expand Up @@ -1036,7 +1042,10 @@ describe('ndJsonStream', () => {
method: 'agent/request',
params: {},
}),
).rejects.toBeInstanceOf(NdJsonQueueLimitError);
).rejects.toMatchObject({
code: 'ndjson_queue_limit_exceeded',
budget: 'outbound_request',
});
expect(onTransportError).toHaveBeenCalledOnce();
writer.releaseLock();
await stream.readable.cancel();
Expand Down
13 changes: 12 additions & 1 deletion packages/acp-bridge/src/ndJsonStream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -80,17 +80,25 @@ export class NdJsonFrameTooLargeError extends Error {
}
}

export type NdJsonQueueBudget =
| 'decoded'
| 'inbound_request'
| 'outbound_request'
| 'prepared_response'
| 'outbound_operation';

export class NdJsonQueueLimitError extends Error {
readonly code = 'ndjson_queue_limit_exceeded';

constructor(
readonly budget: NdJsonQueueBudget,
readonly maxQueuedMessages: number,
readonly maxQueuedBytes: number,
readonly requiredBytes: number,
readonly availableBytes: number,
) {
super(
`NDJSON decoded queue is full ` +
`NDJSON ${budget} queue limit exceeded ` +
`(required ${requiredBytes} bytes, available ${availableBytes} bytes)`,
);
this.name = 'NdJsonQueueLimitError';
Expand Down Expand Up @@ -290,6 +298,7 @@ function createBoundedReadable(
): Promise<void> => {
const queueLimitError = (available: number) =>
new NdJsonQueueLimitError(
'decoded',
limits.maxQueuedMessages,
limits.maxQueuedBytes,
queueCharge,
Expand Down Expand Up @@ -733,6 +742,7 @@ class BoundedInboundRequestLedger {
frameBytes > availableBytes
) {
throw new NdJsonQueueLimitError(
'inbound_request',
this.limits.maxQueuedMessages,
this.limits.maxQueuedBytes,
frameBytes,
Expand Down Expand Up @@ -774,6 +784,7 @@ class BoundedOutstandingRequestLedger {
frameBytes > availableBytes
) {
throw new NdJsonQueueLimitError(
'outbound_request',
this.limits.maxQueuedMessages,
this.limits.maxQueuedBytes,
frameBytes,
Expand Down
33 changes: 30 additions & 3 deletions packages/acp-bridge/src/spawnChannel.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -463,17 +463,44 @@ describe('createSpawnChannelFactory env policy', () => {
).not.toThrow();
expect(() =>
channel.transportGuard?.reservePreparedResponse(third),
).toThrow('NDJSON decoded queue is full');
).toThrow('NDJSON prepared_response queue limit exceeded');

await expect(channel.transportFailed).resolves.toMatchObject({
code: 'ndjson_queue_limit_exceeded',
budget: 'prepared_response',
});
await vi.waitFor(() =>
expect(child.kill).toHaveBeenCalledWith(expectedTreeFallbackSignal()),
);
writer.releaseLock();
});

it('attributes outbound operation budget failures', async () => {
const child = createFakeChildProcess();
mockSpawn.mockReturnValue(child);
const channel = await createSpawnChannelFactory({
pipeLimits: {
maxFrameBytes: 4096,
maxQueuedMessages: 1,
maxQueuedBytes: 4096,
},
})('/tmp/project');

const release = channel.transportGuard?.reserveOutboundOperation([
'first',
{},
]);
expect(() =>
channel.transportGuard?.reserveOutboundOperation(['second', {}]),
).toThrow('NDJSON outbound_operation queue limit exceeded');
await expect(channel.transportFailed).resolves.toMatchObject({
code: 'ndjson_queue_limit_exceeded',
budget: 'outbound_operation',
maxQueuedMessages: 1,
});
release?.();
});

it('stops estimating a large response once its byte budget is exceeded', async () => {
const child = createFakeChildProcess();
mockSpawn.mockReturnValue(child);
Expand Down Expand Up @@ -505,7 +532,7 @@ describe('createSpawnChannelFactory env policy', () => {

expect(() =>
channel.transportGuard?.reservePreparedResponse(response),
).toThrow('NDJSON decoded queue is full');
).toThrow('NDJSON prepared_response queue limit exceeded');
expect(elementReads).toBeLessThan(1_000);
await expect(channel.transportFailed).resolves.toMatchObject({
code: 'ndjson_queue_limit_exceeded',
Expand All @@ -528,7 +555,7 @@ describe('createSpawnChannelFactory env policy', () => {

expect(() =>
channel.transportGuard?.reservePreparedResponse(response),
).toThrow('NDJSON decoded queue is full');
).toThrow('NDJSON prepared_response queue limit exceeded');
await expect(channel.transportFailed).resolves.toMatchObject({
code: 'ndjson_queue_limit_exceeded',
});
Expand Down
2 changes: 2 additions & 0 deletions packages/acp-bridge/src/spawnChannel.ts
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,7 @@ class PreparedResponseBudget {
charge > availableBytes
) {
throw new NdJsonQueueLimitError(
'prepared_response',
this.limits.maxQueuedMessages,
this.limits.maxQueuedBytes,
charge,
Expand Down Expand Up @@ -180,6 +181,7 @@ class OutboundOperationBudget {
charge > availableBytes
) {
throw new NdJsonQueueLimitError(
'outbound_operation',
this.limits.maxQueuedMessages,
this.limits.maxQueuedBytes,
charge,
Expand Down
Loading