kimi-code / packages /agent-core-v2 /test /agent /llmRequester /llmRequesterService.test.ts
SaylorTwift's picture
SaylorTwift HF Staff
Add files using upload-large-folder tool
f0634fb verified
Raw History Blame Contribute Delete
49.2 kB
import { createControlledPromise } from '@antfu/utils';
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import { SyncDescriptor } from '#/_base/di/descriptors';
import { DisposableStore, toDisposable } from '#/_base/di/lifecycle';
import { TestInstantiationService } from '#/_base/di/test';
import type { AgentContext } from '#/agent/agentContext/agentContext';
import { IAgentContextMemoryService } from '#/agent/contextMemory/contextMemory';
import type { ContextMessage } from '#/agent/contextMemory/types';
import {
IAgentContextProjectorService,
type MediaStripSnapshot,
type ProjectionPolicy,
} from '#/agent/contextProjector/contextProjector';
import { AgentContextProjectorService } from '#/agent/contextProjector/contextProjectorService';
import { AgentLLMRequesterService, KIMI_CODE_INFINITE_RETRY_ENV } from '#/agent/llmRequester/llmRequesterService';
import { IAgentLLMRequesterService } from '#/agent/llmRequester/llmRequester';
import { createMachineRequester } from '#/agent/loop/machine/requester';
import { IBootstrapService } from '#/app/bootstrap/bootstrap';
import {
createTurnMachine,
type AssistantEntry,
type TurnEvent,
type TurnInput,
type TurnLlmEvent,
} from '#human/agent/turn';
import { UNKNOWN_CAPABILITY } from '#human/llm/capability';
import type { LlmModel } from '#human/llm/model';
import type { LlmRequester } from '#human/llm/requester/requester';
import { createActor, emit, setup } from '#human/xstate2';
import { ISessionTokenCountingService } from '#/session/tokenCounting/sessionTokenCounting';
import { IAgentProfileService } from '#/agent/profile/profile';
import { IAgentStateService } from '#/agent/state/agentState';
import { AgentStateService } from '#/agent/state/agentStateService';
import { IAgentToolRegistryService } from '#/agent/toolRegistry/toolRegistry';
import { IAgentToolSelectService } from '#/agent/toolSelect/toolSelect';
import { IAgentMediaResolverService } from '#/agent/media/mediaResolver';
import { ISessionUsageService } from '#/session/usage/sessionUsage';
import { IConfigService } from '#/app/config/config';
import type { Event2 } from '#/app/event/event2';
import { IEventBus } from '#/app/event/eventBus';
import {
APIConnectionError,
APIContextOverflowError,
APIEmptyResponseError,
APIProviderQuotaExhaustedError,
APIProviderRateLimitError,
APIRequestTooLargeError,
APIStatusError,
} from '#/llm-adapter/contract/errors';
import { emptyUsage, type TokenUsage } from '#human/llm/usage';
import { type Message } from '#/llm-adapter/contract/message';
import { isToolCall, type StreamedMessagePart, type ToolCall } from '#human/llm/message';
import type { ThinkingEffort } from '#human/llm/thinking';
import type { ModelCapability } from '#/llm-adapter/contract/capability';
import { IModelCatalog, type Model } from '#/llm-adapter/model/catalog';
import { IModelService } from '#/llm-adapter/model/model';
import {
type ModelRequestEvent,
type ModelRequestInput,
type ModelRequester,
} from '#/llm-adapter/model/model-requester';
import { ITelemetryService } from '#/app/telemetry/telemetry';
import { ILogService } from '#/_base/log/log';
import { Error2, ErrorCodes } from '#/errors';
import { IEventDispatcher } from '#/state/eventDispatcher';
import type { WireRecord } from '#/wire/record';
import { recordingTelemetry, type TelemetryRecord } from '../../app/telemetry/stubs';
import { stubBootstrap } from '../../app/bootstrap/stubs';
import {
recordingWireLog,
registerTestAgentWire,
registerTestEventDispatcher,
} from '../../wire/stubs';
const turnHarnessModel: LlmModel = {
provider: 'test',
model: 'test-model',
capability: UNKNOWN_CAPABILITY,
};
function createTurnHarness(requester: LlmRequester) {
return setup({
types: {
input: {} as TurnInput,
context: {} as { turnInput: TurnInput },
events: {} as TurnEvent,
emitted: {} as TurnLlmEvent,
},
actors: { turn: createTurnMachine(requester) },
}).createMachine({
id: 'turn-harness',
initial: 'running',
context: ({ input }) => ({ turnInput: input }),
states: {
running: {
invoke: {
src: 'turn',
input: ({ context }) => context.turnInput,
onDone: { target: 'completed' },
},
on: {
'*': {
actions: emit(({ event }) => event as TurnLlmEvent),
},
},
},
completed: { type: 'final' },
},
});
}
const capabilities: ModelCapability = {
image_in: false,
video_in: false,
audio_in: false,
thinking: false,
tool_use: false,
max_context_tokens: 1000,
};
const history: Message[] = [
{ role: 'user', content: [{ type: 'text', text: 'hello' }], toolCalls: [] },
];
type ProjectionKind = 'normal' | 'strict' | 'degraded' | 'stripped';
function classifyProjectionPolicy(policy: ProjectionPolicy | undefined): ProjectionKind {
if (typeof policy?.media === 'object') return 'stripped';
if (policy?.media === 'degraded') return 'degraded';
if (policy?.structure === 'strict') return 'strict';
return 'normal';
}
function recordProjectionCalls(
project: (
messages: readonly ContextMessage[],
policy: ProjectionPolicy | undefined,
) => readonly Message[] = (messages) => messages,
): {
projector: Pick<IAgentContextProjectorService, 'project'>;
calls: ProjectionKind[];
} {
const calls: ProjectionKind[] = [];
return {
projector: {
project: (messages: readonly ContextMessage[], policy) => {
calls.push(classifyProjectionPolicy(policy));
return project(messages, policy);
},
},
calls,
};
}
function createRequester(
calls: { value: number },
firstCallError?: Error | null,
subsequentCallErrors: readonly Error[] = [],
capturedInputs?: ModelRequestInput[],
): ModelRequester {
const model: Model = {
id: 'm',
name: 'wire-model',
aliases: [],
protocol: 'anthropic',
baseUrl: 'https://example.test',
headers: {},
capabilities,
maxContextSize: 1000,
alwaysThinking: false,
providerName: 'p',
};
return {
model,
request: async function* (input) {
calls.value += 1;
capturedInputs?.push(input);
const error =
calls.value === 1
? firstCallError === null
? undefined
: (firstCallError ??
new APIStatusError(400, 'messages: `tool_use` ids must be unique'))
: subsequentCallErrors[calls.value - 2];
if (error !== undefined) throw error;
yield {
type: 'finish',
message: { role: 'assistant', content: [{ type: 'text', text: 'ok' }], toolCalls: [] },
providerFinishReason: 'completed',
rawFinishReason: 'stop',
id: 'resp-1',
};
},
};
}
let disposables: DisposableStore;
beforeEach(() => {
disposables = new DisposableStore();
});
afterEach(() => disposables.dispose());
function createService(
requester: ModelRequester,
projector:
| (Pick<IAgentContextProjectorService, 'project'> &
Partial<Pick<IAgentContextProjectorService, 'captureMediaStripSnapshot'>>)
| undefined,
options: {
readonly thinkingLevel?: ThinkingEffort;
readonly mediaResolver?: Partial<IAgentMediaResolverService>;
readonly contextMessages?: Message[];
readonly env?: Record<string, string>;
} = {},
) {
const ix = disposables.add(new TestInstantiationService());
ix.stub(IBootstrapService, stubBootstrap('/tmp/kimi-code-llm-requester-test', options.env ?? {}));
const thinkingLevel = options.thinkingLevel ?? 'off';
const profile: Partial<IAgentProfileService> = {
hasProvider: () => true,
resolveModelContext: () => ({
modelAlias: 'm',
modelCapabilities: capabilities,
maxOutputSize: undefined,
alwaysThinking: undefined,
thinkingLevel,
reservedContextSize: undefined,
compactionTriggerRatio: undefined,
compactionMaxAttempts: undefined,
}),
resolveRequestParams: () => ({}),
getSystemPrompt: () => 'system',
data: () => ({
cwd: '',
modelAlias: 'm',
modelCapabilities: capabilities,
thinkingLevel,
systemPrompt: 'system',
}),
};
const measuredCalls: { readonly messages: number; readonly usage: TokenUsage }[] = [];
const tokenCounting = {
get: () => ({ size: 0, measured: 0, estimated: 0 }),
measured: (
_agent: AgentContext,
input: readonly Message[],
_output: readonly Message[],
usage: TokenUsage,
) => {
measuredCalls.push({ messages: input.length, usage });
},
};
const usage = { record: () => Promise.resolve(), status: () => ({}) };
const context = {
get: () => options.contextMessages ?? history,
};
const tools = { list: () => [] };
const config: Partial<IConfigService> = {
get: (() => undefined) as IConfigService['get'],
};
const log = { info: () => undefined, warn: () => undefined };
const telemetryRecords: TelemetryRecord[] = [];
const telemetry = recordingTelemetry(telemetryRecords);
const toolSelect: Partial<IAgentToolSelectService> = {
enabled: () => false,
shapeTools: (entries) => entries,
shapeHistory: (messages) => messages,
};
const testSnapshot = Object.freeze({}) as MediaStripSnapshot;
const events: Event2[] = [];
const eventBus: IEventBus = {
_serviceBrand: undefined,
publish: (event) => events.push(event),
subscribe: () => toDisposable(() => {}),
};
ix.stub(IAgentContextMemoryService, context);
ix.stub(IAgentToolSelectService, toolSelect);
ix.stub(IAgentMediaResolverService, options.mediaResolver ?? {
resolve: async (messages) => messages,
displayPaths: async () => new Map(),
});
if (projector === undefined) {
ix.set(
IAgentContextProjectorService,
new SyncDescriptor(AgentContextProjectorService),
);
} else {
ix.stub(IAgentContextProjectorService, {
captureMediaStripSnapshot: () => testSnapshot,
...projector,
});
}
ix.stub(ISessionTokenCountingService, tokenCounting);
ix.stub(IAgentToolRegistryService, tools);
ix.stub(IAgentProfileService, profile);
ix.stub(ISessionUsageService, usage);
ix.stub(IConfigService, config);
ix.stub(ILogService, log);
ix.stub(ITelemetryService, telemetry);
ix.stub(IModelCatalog, {
_serviceBrand: undefined,
get: () => requester.model,
getRequester: () => requester,
findByName: () => [],
});
ix.stub(IModelService, {
get: () => undefined,
});
const records: WireRecord[] = [];
registerTestAgentWire(ix, 'wire/llm-requester', {
log: recordingWireLog(records),
eventBus,
});
registerTestEventDispatcher(ix);
ix.set(IAgentStateService, new AgentStateService());
ix.set(IAgentLLMRequesterService, new SyncDescriptor(AgentLLMRequesterService));
return {
service: ix.get(IAgentLLMRequesterService),
dispatcher: ix.get(IEventDispatcher),
records,
events,
telemetry,
telemetryRecords,
measuredCalls,
};
}
describe('AgentLLMRequesterService measured anchors', () => {
it('skips the measured anchor when the stream reports no usage', async () => {
const { service, measuredCalls } = createService(createRequester({ value: 0 }), undefined);
await service.request();
expect(measuredCalls).toHaveLength(0);
});
it('writes the measured anchor from the reported usage', async () => {
const requester = createRequester({ value: 0 });
const base = requester.request.bind(requester);
requester.request = async function* (input, signal, options) {
yield {
type: 'usage',
usage: { inputOther: 40, output: 2, inputCacheRead: 0, inputCacheCreation: 0 },
model: 'wire-model',
};
yield* base(input, signal, options);
};
const { service, measuredCalls } = createService(requester, undefined);
await service.request();
expect(measuredCalls).toHaveLength(1);
expect(measuredCalls[0]?.usage.inputOther).toBe(40);
});
});
describe('AgentLLMRequesterService Anthropic effort diagnostics', () => {
it('warns and sends when the effort is not listed by the model', async () => {
const calls = { value: 0 };
const requester = createRequester(calls, null);
Object.defineProperty(requester.model, 'supportEfforts', { value: ['max'] });
const { service, events } = createService(requester, undefined, { thinkingLevel: 'high' });
const result = await service.request();
expect(result.message.content).toEqual([{ type: 'text', text: 'ok' }]);
expect(calls.value).toBe(1);
expect(events.filter((event) => event.type === 'warning')).toEqual([
expect.objectContaining({
type: 'warning',
code: 'anthropic-thinking-effort-not-listed',
message:
'Thinking effort "high" is not listed for model "wire-model" (known: max). The configured value will be sent unchanged to the Anthropic-compatible backend.',
}),
]);
});
});
describe('AgentLLMRequesterService strict resend', () => {
it('resends once with strict projection after a recoverable structural 400', async () => {
const calls = { value: 0 };
const projection = recordProjectionCalls();
const { service } = createService(createRequester(calls), projection.projector);
const result = await service.request();
expect(result.message.content).toEqual([{ type: 'text', text: 'ok' }]);
expect(result.usage).toEqual(emptyUsage());
expect(calls.value).toBe(2);
expect(projection.calls).toEqual(['normal', 'strict']);
});
it('does not resend for non-recoverable errors', async () => {
const requester = createRequester({ value: 0 });
Object.defineProperty(requester, 'request', {
value: async function* () {
const events: ModelRequestEvent[] = [];
for (const event of events) yield event;
throw new APIStatusError(401, 'unauthorized');
},
});
const projection = recordProjectionCalls();
const { service } = createService(requester, projection.projector);
await expect(service.request()).rejects.toMatchObject({
statusCode: 401,
});
expect(projection.calls).toEqual(['normal']);
});
});
describe('AgentLLMRequesterService infinite retry', () => {
afterEach(() => {
vi.useRealTimers();
});
it('retries every request error while KIMI_CODE_INFINITE_RETRY is set', async () => {
vi.useFakeTimers();
const calls = { value: 0 };
const requester = createRequester(calls, new APIStatusError(400, 'endpoint broken'), [
new APIStatusError(404, 'model not found'),
new APIConnectionError('socket hang up'),
new APIProviderQuotaExhaustedError('quota exhausted'),
]);
const { service } = createService(requester, undefined, {
env: { [KIMI_CODE_INFINITE_RETRY_ENV]: '1' },
});
const promise = service.request();
await vi.runAllTimersAsync();
const finish = await promise;
expect(calls.value).toBe(5);
expect(finish.message.content).toEqual([{ type: 'text', text: 'ok' }]);
});
it('honors the provider retry-after delay while retrying indefinitely', async () => {
const calls = { value: 0 };
const requester = createRequester(calls, new APIProviderRateLimitError('slow down', null, 1));
const { service } = createService(requester, undefined, {
env: { [KIMI_CODE_INFINITE_RETRY_ENV]: '1' },
});
const startedAt = Date.now();
await service.request();
expect(calls.value).toBe(2);
expect(Date.now() - startedAt).toBeLessThan(500);
});
it('stops retrying when the caller aborts during the backoff wait', async () => {
vi.useFakeTimers();
const calls = { value: 0 };
const requester = createRequester(calls, new APIStatusError(400, 'endpoint broken'));
const { service } = createService(requester, undefined, {
env: { [KIMI_CODE_INFINITE_RETRY_ENV]: '1' },
});
const controller = new AbortController();
setTimeout(() => controller.abort(new Error('stop')), 100);
const promise = service.request({}, undefined, controller.signal);
const assertion = expect(promise).rejects.toThrow('stop');
await vi.runAllTimersAsync();
await assertion;
expect(calls.value).toBe(1);
});
it('keeps deterministic projection recovery ahead of infinite retry', async () => {
vi.useFakeTimers();
const calls = { value: 0 };
const requester = createRequester(calls, new APIRequestTooLargeError(413, 'Request Entity Too Large'));
const { service } = createService(requester, undefined, {
env: { [KIMI_CODE_INFINITE_RETRY_ENV]: '1' },
});
await service.request();
expect(calls.value).toBe(2);
});
it('lets context overflow reach deterministic recovery instead of retrying', async () => {
vi.useFakeTimers();
const calls = { value: 0 };
const requester = createRequester(
calls,
new APIContextOverflowError(400, 'context length exceeded'),
);
const { service } = createService(requester, undefined, {
env: { [KIMI_CODE_INFINITE_RETRY_ENV]: '1' },
});
await expect(service.request()).rejects.toBeInstanceOf(APIContextOverflowError);
expect(calls.value).toBe(1);
});
it('retries operation requests indefinitely', async () => {
vi.useFakeTimers();
const calls = { value: 0 };
const requester = createRequester(calls, new APIStatusError(400, 'endpoint broken'), [
new APIStatusError(404, 'model not found'),
]);
const { service } = createService(requester, undefined, {
env: { [KIMI_CODE_INFINITE_RETRY_ENV]: '1' },
});
const promise = service.request({
source: { type: 'operation', requestKind: 'full_compaction' },
});
await vi.runAllTimersAsync();
await promise;
expect(calls.value).toBe(3);
});
it('does not retry when the switch is unset', async () => {
vi.useFakeTimers();
const calls = { value: 0 };
const requester = createRequester(calls, new APIStatusError(400, 'endpoint broken'));
const { service } = createService(requester, undefined);
await expect(service.request()).rejects.toMatchObject({ statusCode: 400 });
expect(calls.value).toBe(1);
});
});
describe('AgentLLMRequesterService media-stripped resend', () => {
const IMAGE_FORMAT_400 = new APIStatusError(
400,
'unsupported image format: image/avif is not supported',
);
it('resends once with the media-stripped projection after an image-format 400', async () => {
const calls = { value: 0 };
const projection = recordProjectionCalls();
const { service } = createService(createRequester(calls, IMAGE_FORMAT_400), projection.projector);
const result = await service.request();
expect(result.message.content).toEqual([{ type: 'text', text: 'ok' }]);
expect(calls.value).toBe(2);
expect(projection.calls).toEqual(['normal', 'stripped']);
});
it('keeps later steps of the same turn on the stripped projection', async () => {
const calls = { value: 0 };
const projection = recordProjectionCalls();
const { service } = createService(createRequester(calls, IMAGE_FORMAT_400), projection.projector);
await service.request({ source: { type: 'turn', turnId: 1, step: 1 } });
expect(calls.value).toBe(2);
expect(projection.calls).toEqual(['normal', 'stripped']);
await service.request({ source: { type: 'turn', turnId: 1, step: 2 } });
expect(calls.value).toBe(3);
expect(projection.calls).toEqual(['normal', 'stripped', 'stripped']);
});
it('does not resend for an unrelated 400', async () => {
const calls = { value: 0 };
const projection = recordProjectionCalls();
const { service } = createService(
createRequester(calls, new APIStatusError(400, 'some other validation problem')),
projection.projector,
);
await expect(service.request()).rejects.toMatchObject({ statusCode: 400 });
expect(calls.value).toBe(1);
expect(projection.calls).toEqual(['normal']);
});
it('warns the user when media are stripped from the retried request', async () => {
const calls = { value: 0 };
const projection = recordProjectionCalls((messages, policy) =>
typeof policy?.media === 'object' ? history : messages,
);
const contextMessages: Message[] = [
{
role: 'user',
content: [{ type: 'image_url', imageUrl: { url: 'data:image/png;base64,IMAGE' } }],
toolCalls: [],
},
];
const { service, dispatcher, events } = createService(
createRequester(calls, IMAGE_FORMAT_400),
projection.projector,
{ contextMessages },
);
await service.request();
await dispatcher.flush();
expect(events.filter((event) => event.type === 'warning')).toEqual([
expect.objectContaining({
type: 'warning',
code: 'media-stripped',
message:
'Provider rejected the media in the request; all media were omitted and the request was retried.',
}),
]);
});
});
describe('AgentLLMRequesterService media-degraded resend', () => {
const BODY_TOO_LARGE_413 = new APIRequestTooLargeError(413, 'Request Entity Too Large');
it('resends once with the media-degraded projection after an HTTP 413', async () => {
const calls = { value: 0 };
const projection = recordProjectionCalls();
const { service } = createService(
createRequester(
calls,
new Error2(ErrorCodes.PROVIDER_API_ERROR, 'Provider request failed', {
cause: BODY_TOO_LARGE_413,
}),
),
projection.projector,
);
const result = await service.request();
expect(result.message.content).toEqual([{ type: 'text', text: 'ok' }]);
expect(calls.value).toBe(2);
expect(projection.calls).toEqual(['normal', 'degraded']);
});
it('attaches display paths to degraded older media', async () => {
const calls = { value: 0 };
const capturedInputs: ModelRequestInput[] = [];
const imageMessage = (url: string): Message => ({
role: 'user',
content: [{ type: 'image_url', imageUrl: { url } }],
toolCalls: [],
});
const { service } = createService(
createRequester(calls, BODY_TOO_LARGE_413, [], capturedInputs),
undefined,
{
mediaResolver: {
resolve: async (messages) => messages,
displayPaths: async () => new Map([['kimi-file://f_old', '/session/media/f_old.png']]),
},
},
);
await service.request({
messages: [
imageMessage('kimi-file://f_old'),
imageMessage('kimi-file://f_keep1'),
imageMessage('kimi-file://f_keep2'),
],
source: { type: 'turn', turnId: 1, step: 1 },
});
expect(calls.value).toBe(2);
const parts = capturedInputs[1]!.messages.flatMap((message) => message.content);
const urls = parts
.filter((part) => part.type === 'image_url')
.map((part) => part.imageUrl.url);
expect(urls).toEqual(['kimi-file://f_keep1', 'kimi-file://f_keep2']);
const texts = parts.filter((part) => part.type === 'text').map((part) => part.text);
expect(texts).toContain('<image path="/session/media/f_old.png"></image>');
});
it('falls back to media-stripped when the media-degraded request still receives 413', async () => { const calls = { value: 0 };
const projection = recordProjectionCalls();
const { service } = createService(
createRequester(calls, BODY_TOO_LARGE_413, [BODY_TOO_LARGE_413]),
projection.projector,
);
const result = await service.request({ source: { type: 'turn', turnId: 1, step: 1 } });
expect(result.message.content).toEqual([{ type: 'text', text: 'ok' }]);
expect(calls.value).toBe(3);
expect(projection.calls).toEqual(['normal', 'degraded', 'stripped']);
});
it('records repeated-413 recovery projections on the sticky later request', async () => {
const calls = { value: 0 };
const { service, dispatcher, records } = createService(
createRequester(calls, BODY_TOO_LARGE_413, [BODY_TOO_LARGE_413]),
{
project: (messages: readonly ContextMessage[]) => messages,
},
);
await service.request({ source: { type: 'turn', turnId: 1, step: 1 } });
await service.request({ source: { type: 'turn', turnId: 1, step: 2 } });
await dispatcher.flush();
expect(
records
.filter((record) => record.type === 'llm.request')
.map((record) => record['projection']),
).toEqual([undefined, 'media-degraded', 'media-stripped', 'media-stripped']);
});
it('keeps new recovery media visible on later snapshot-stripped steps', async () => {
const calls = { value: 0 };
const capturedInputs: ModelRequestInput[] = [];
const oldUrl = 'data:image/png;base64,REJECTED';
const newUrl = 'data:image/png;base64,SMALL';
const imageMessage = (url: string, id: string): Message => ({
role: 'user',
content: [{ type: 'image_url', imageUrl: { url, id } }],
toolCalls: [],
});
const { service } = createService(
createRequester(
calls,
BODY_TOO_LARGE_413,
[BODY_TOO_LARGE_413],
capturedInputs,
),
undefined,
);
await service.request({
messages: [imageMessage(oldUrl, 'rejected-id')],
source: { type: 'turn', turnId: 1, step: 1 },
});
await service.request({
messages: [
imageMessage(oldUrl, 'rejected-id'),
imageMessage(newUrl, 'recovery-id'),
],
source: { type: 'turn', turnId: 1, step: 2 },
});
const visibleUrls = capturedInputs
.at(-1)
?.messages.flatMap((message) => message.content)
.filter((part) => part.type === 'image_url')
.map((part) => part.imageUrl.url);
expect(visibleUrls).toEqual([newUrl]);
});
it('stops after the media-stripped request also receives 413', async () => {
const calls = { value: 0 };
const projection = recordProjectionCalls();
const { service } = createService(
createRequester(calls, BODY_TOO_LARGE_413, [BODY_TOO_LARGE_413, BODY_TOO_LARGE_413]),
projection.projector,
);
await expect(
service.request({ source: { type: 'turn', turnId: 1, step: 1 } }),
).rejects.toBe(BODY_TOO_LARGE_413);
expect(calls.value).toBe(3);
expect(projection.calls).toEqual(['normal', 'degraded', 'stripped']);
});
it('keeps later steps of the same turn on the degraded projection', async () => {
const calls = { value: 0 };
const projection = recordProjectionCalls();
const { service } = createService(createRequester(calls, BODY_TOO_LARGE_413), projection.projector);
await service.request({ source: { type: 'turn', turnId: 1, step: 1 } });
expect(calls.value).toBe(2);
expect(projection.calls).toEqual(['normal', 'degraded']);
await service.request({ source: { type: 'turn', turnId: 1, step: 2 } });
expect(calls.value).toBe(3);
expect(projection.calls).toEqual(['normal', 'degraded', 'degraded']);
});
it('does not resend for a plain 400 or a non-413 status', async () => {
for (const error of [
new APIStatusError(400, 'max_tokens must be positive'),
new APIStatusError(422, 'unprocessable'),
]) {
const calls = { value: 0 };
const projection = recordProjectionCalls();
const { service } = createService(createRequester(calls, error), projection.projector);
await expect(service.request()).rejects.toBe(error);
expect(calls.value).toBe(1);
expect(projection.calls).toEqual(['normal']);
}
});
it('does not warn when the degraded projection leaves the request unchanged', async () => {
const calls = { value: 0 };
const projection = recordProjectionCalls();
const { service, dispatcher, events } = createService(
createRequester(calls, BODY_TOO_LARGE_413),
projection.projector,
);
await service.request();
await dispatcher.flush();
expect(events.filter((event) => event.type === 'warning')).toEqual([]);
});
it('warns for each escalation when the degraded resend is also rejected as too large', async () => {
const calls = { value: 0 };
const projection = recordProjectionCalls((messages, policy) => {
if (policy?.media === 'degraded') {
const message = messages[0]!;
return [{ ...message, content: message.content.slice(-2) }];
}
return typeof policy?.media === 'object' ? history : messages;
});
const contextMessages: Message[] = [
{
role: 'user',
content: ['ONE', 'TWO', 'THREE'].map((data) => ({
type: 'image_url',
imageUrl: { url: `data:image/png;base64,${data}` },
})),
toolCalls: [],
},
];
const { service, dispatcher, events } = createService(
createRequester(calls, BODY_TOO_LARGE_413, [BODY_TOO_LARGE_413]),
projection.projector,
{ contextMessages },
);
await service.request({ source: { type: 'turn', turnId: 1, step: 1 } });
await dispatcher.flush();
expect(events.filter((event) => event.type === 'warning')).toEqual([
expect.objectContaining({ type: 'warning', code: 'media-degraded' }),
expect.objectContaining({ type: 'warning', code: 'media-stripped' }),
]);
});
});
describe('AgentLLMRequesterService combined recovery projections', () => {
const BODY_TOO_LARGE_413 = new APIRequestTooLargeError(413, 'Request Entity Too Large');
const IMAGE_FORMAT_400 = new APIStatusError(
400,
'unsupported image format: image/avif is not supported',
);
const STRUCTURAL_400 = new APIStatusError(400, 'messages: `tool_use` ids must be unique');
function createPolicyRecordingProjector(policies: {
policies: (ProjectionPolicy | undefined)[];
}): Pick<IAgentContextProjectorService, 'project'> {
return {
project: (messages: readonly ContextMessage[], policy) => {
policies.policies.push(policy);
return messages;
},
};
}
it('accumulates media repairs on top of strict across repeated rejections', async () => {
const calls = { value: 0 };
const policies: (ProjectionPolicy | undefined)[] = [];
const { service, dispatcher, records } = createService(
createRequester(calls, STRUCTURAL_400, [BODY_TOO_LARGE_413, BODY_TOO_LARGE_413]),
createPolicyRecordingProjector({ policies }),
);
await service.request({ source: { type: 'turn', turnId: 1, step: 1 } });
expect(calls.value).toBe(4);
expect(policies).toEqual([
undefined,
{ structure: 'strict' },
{ structure: 'strict', media: 'degraded' },
{ structure: 'strict', media: { strip: expect.anything() } },
]);
await dispatcher.flush();
expect(
records.filter((record) => record.type === 'llm.request').map((record) => record['projection']),
).toEqual([undefined, 'strict', 'strict-media-degraded', 'strict-media-stripped']);
});
it('strips rejected images on top of strict after an image-format rejection on the strict resend', async () => {
const calls = { value: 0 };
const policies: (ProjectionPolicy | undefined)[] = [];
const { service } = createService(
createRequester(calls, STRUCTURAL_400, [IMAGE_FORMAT_400]),
createPolicyRecordingProjector({ policies }),
);
await service.request();
expect(calls.value).toBe(3);
expect(policies.map((policy) => policy?.structure)).toEqual([undefined, 'strict', 'strict']);
expect(typeof policies[2]?.media).toBe('object');
});
it('applies the strict repair on top of degraded media without repeating the media warning', async () => {
const calls = { value: 0 };
const policies: (ProjectionPolicy | undefined)[] = [];
const contextMessages: Message[] = [
{
role: 'user',
content: ['ONE', 'TWO', 'THREE'].map((data) => ({
type: 'image_url',
imageUrl: { url: `data:image/png;base64,${data}` },
})),
toolCalls: [],
},
];
const projector = {
project: (messages: readonly ContextMessage[], policy: ProjectionPolicy | undefined) => {
policies.push(policy);
if (policy?.media !== 'degraded') return messages;
const message = messages[0]!;
return [{ ...message, content: message.content.slice(policy.structure === 'strict' ? -1 : -2) }];
},
};
const { service, dispatcher, events } = createService(
createRequester(calls, BODY_TOO_LARGE_413, [STRUCTURAL_400]),
projector,
{ contextMessages },
);
await service.request();
await dispatcher.flush();
expect(calls.value).toBe(3);
expect(policies).toEqual([
undefined,
{ media: 'degraded' },
{ structure: 'strict', media: 'degraded' },
]);
expect(events.filter((event) => event.type === 'warning')).toEqual([
expect.objectContaining({ type: 'warning', code: 'media-degraded' }),
]);
});
});
describe('AgentLLMRequesterService trace id', () => {
const passthroughProjector = {
project: (messages: readonly ContextMessage[]) => messages,
};
function createTracedRequester(traceId: string | null): ModelRequester {
const model: Model = {
id: 'm',
name: 'wire-model',
aliases: [],
protocol: 'openai',
baseUrl: 'https://example.test',
headers: {},
capabilities,
maxContextSize: 1000,
alwaysThinking: false,
providerName: 'p',
};
return {
model,
request: async function* (_input, _signal, requestOptions) {
requestOptions?.onTraceId?.(traceId);
yield {
type: 'finish',
message: { role: 'assistant', content: [{ type: 'text', text: 'ok' }], toolCalls: [] },
providerFinishReason: 'completed',
rawFinishReason: 'stop',
id: 'resp-1',
traceId: traceId ?? undefined,
};
},
};
}
it('exposes the request trace and returns it on finish', async () => {
const requester = createTracedRequester('trace-req-1');
const headersArrived = createControlledPromise<void>();
const releaseStream = createControlledPromise<void>();
Object.defineProperty(requester, 'request', {
value: async function* (_input: unknown, _signal: unknown, requestOptions: {
onTraceId?: (traceId: string | null) => void;
}) {
requestOptions.onTraceId?.('trace-req-1');
headersArrived.resolve();
await releaseStream;
yield {
type: 'finish',
message: { role: 'assistant', content: [{ type: 'text', text: 'ok' }], toolCalls: [] },
providerFinishReason: 'completed',
rawFinishReason: 'stop',
id: 'resp-1',
traceId: 'trace-req-1',
} satisfies ModelRequestEvent;
},
});
const { service } = createService(requester, passthroughProjector);
const request = service.start({ source: { type: 'turn', turnId: 1, step: 1 } });
await headersArrived;
expect(request.trace.traceId).toBe('trace-req-1');
releaseStream.resolve();
const finish = await request.result;
expect(finish.traceId).toBe('trace-req-1');
expect(request.trace.traceId).toBe('trace-req-1');
});
it('reports an absent trace before a request that returns none', async () => {
const { service } = createService(createTracedRequester(null), passthroughProjector);
const request = service.start();
const finish = await request.result;
expect(finish.traceId).toBeUndefined();
expect(request.trace.traceId).toBeUndefined();
});
it('attaches trace_id, turn_id and step_no to api_error from the failed request', async () => {
const requester = createTracedRequester(null);
Object.defineProperty(requester, 'request', {
value: async function* () {
const events: ModelRequestEvent[] = [];
for (const event of events) yield event;
throw new APIStatusError(500, 'boom', 'req-1', null, 'trace-fail-1');
},
});
const { service, telemetryRecords } = createService(requester, passthroughProjector);
const request = service.start({ source: { type: 'turn', turnId: 3, step: 2 } });
await expect(request.result).rejects.toMatchObject({ statusCode: 500 });
expect(telemetryRecords).toContainEqual({
event: 'api_error',
properties: expect.objectContaining({
error_type: '5xx_server',
trace_id: 'trace-fail-1',
turn_id: 3,
step_no: 2,
}),
});
expect(request.trace.traceId).toBe('trace-fail-1');
});
it('keeps the header-captured trace when the request fails after headers arrived', async () => {
const requester = createTracedRequester(null);
Object.defineProperty(requester, 'request', {
value: async function* (...args: unknown[]) {
const requestOptions = args[2] as
| { onTraceId?: (traceId: string | null) => void }
| undefined;
requestOptions?.onTraceId?.('trace-mid-stream');
const events: ModelRequestEvent[] = [];
for (const event of events) yield event;
throw new APIEmptyResponseError('no content, no tool calls');
},
});
const { service, telemetryRecords } = createService(requester, passthroughProjector);
const request = service.start({ source: { type: 'turn', turnId: 4, step: 1 } });
await expect(request.result).rejects.toThrow();
const apiError = telemetryRecords.find((record) => record.event === 'api_error');
expect(apiError?.properties?.['trace_id']).toBe('trace-mid-stream');
expect(request.trace.traceId).toBe('trace-mid-stream');
});
it('clears the previous physical request trace before a projection retry', async () => {
const requester = createTracedRequester(null);
let attempts = 0;
Object.defineProperty(requester, 'request', {
value: async function* (...args: unknown[]) {
const events: ModelRequestEvent[] = [];
for (const event of events) yield event;
attempts += 1;
const requestOptions = args[2] as
| { onTraceId?: (traceId: string | null) => void }
| undefined;
if (attempts === 1) {
requestOptions?.onTraceId?.('trace-first-projection');
throw new APIRequestTooLargeError(413, 'retry with degraded media');
}
throw new APIConnectionError('socket hang up');
},
});
const { service, telemetryRecords } = createService(requester, passthroughProjector);
const request = service.start();
await expect(request.result).rejects.toThrow('socket hang up');
expect(attempts).toBe(2);
expect(request.trace.traceId).toBeUndefined();
expect(
telemetryRecords.find((record) => record.event === 'api_error')?.properties?.['trace_id'],
).toBeUndefined();
});
it('mirrors the request trace into the ambient telemetry context', async () => {
const { service, telemetry } = createService(
createTracedRequester('trace-ambient-1'),
passthroughProjector,
);
await service.request({ source: { type: 'turn', turnId: 1, step: 1 } });
expect(telemetry.getContext()['trace_id']).toBe('trace-ambient-1');
});
it('clears the ambient trace when the next turn request starts without one', async () => {
let nextTrace: string | null = 'trace-ambient-2';
const requester = createTracedRequester(null);
Object.defineProperty(requester, 'request', {
value: async function* (_input: unknown, _signal: unknown, requestOptions: {
onTraceId?: (traceId: string | null) => void;
}) {
requestOptions?.onTraceId?.(nextTrace);
yield {
type: 'finish',
message: { role: 'assistant', content: [{ type: 'text', text: 'ok' }], toolCalls: [] },
providerFinishReason: 'completed',
rawFinishReason: 'stop',
id: 'resp-1',
traceId: nextTrace ?? undefined,
} satisfies ModelRequestEvent;
},
});
const { service, telemetry } = createService(requester, passthroughProjector);
await service.request({ source: { type: 'turn', turnId: 1, step: 1 } });
expect(telemetry.getContext()['trace_id']).toBe('trace-ambient-2');
nextTrace = null;
await service.request({ source: { type: 'turn', turnId: 1, step: 2 } });
expect(telemetry.getContext()['trace_id']).toBeUndefined();
});
it('mirrors the failing request trace into the ambient telemetry context', async () => {
const requester = createTracedRequester(null);
Object.defineProperty(requester, 'request', {
value: async function* () {
const events: ModelRequestEvent[] = [];
for (const event of events) yield event;
throw new APIStatusError(500, 'boom', 'req-1', null, 'trace-fail-ambient');
},
});
const { service, telemetry } = createService(requester, passthroughProjector);
await expect(
service.request({ source: { type: 'turn', turnId: 1, step: 1 } }),
).rejects.toMatchObject({ statusCode: 500 });
expect(telemetry.getContext()['trace_id']).toBe('trace-fail-ambient');
});
it('keeps the ambient trace untouched for operation requests', async () => {
const { service, telemetry } = createService(
createTracedRequester('trace-operation-1'),
passthroughProjector,
);
telemetry.setContext({ trace_id: 'trace-turn-1' });
await service.request({ source: { type: 'operation', requestKind: 'full_compaction' } });
expect(telemetry.getContext()['trace_id']).toBe('trace-turn-1');
});
});
describe('AgentLLMRequesterService media resolver wiring', () => {
it('resolves the projected messages through the DI-injected media resolver', async () => {
const requester = createRequester({ value: 0 }, null);
const resolve = vi.fn(async (messages: readonly Message[], _requester: ModelRequester) => messages);
const { service } = createService(requester, undefined, {
mediaResolver: { resolve },
});
await service.request();
expect(resolve).toHaveBeenCalledTimes(1);
expect(resolve.mock.calls[0]?.[1]).toBe(requester);
});
});
function createScriptedRequester(
script: { ids: string[]; error?: Error }[],
): ModelRequester {
const base = createRequester({ value: 0 });
let callIndex = 0;
return {
model: base.model,
request: async function* () {
const step = script[Math.min(callIndex++, script.length - 1)]!;
if (step.error !== undefined) {
if (step.ids.length > 0) {
yield {
type: 'part',
part: {
type: 'function',
id: step.ids[0]!,
name: 'Bash',
arguments: null,
_streamIndex: 0,
},
} satisfies ModelRequestEvent;
}
throw step.error;
}
const toolCalls: ToolCall[] = [];
for (const [index, id] of step.ids.entries()) {
yield {
type: 'part',
part: { type: 'function', id, name: 'Bash', arguments: null, _streamIndex: index },
} satisfies ModelRequestEvent;
yield {
type: 'part',
part: { type: 'tool_call_part', argumentsPart: '{"command":"ls"}', index },
} satisfies ModelRequestEvent;
toolCalls.push({ type: 'function', id, name: 'Bash', arguments: '{"command":"ls"}' });
}
yield {
type: 'finish',
message: { role: 'assistant', content: [], toolCalls },
providerFinishReason: 'completed',
rawFinishReason: 'stop',
id: 'resp-1',
} satisfies ModelRequestEvent;
},
};
}
describe('AgentLLMRequesterService tool call id normalization', () => {
it('passes provider-unique ids through unchanged', async () => {
const parts: StreamedMessagePart[] = [];
const { service } = createService(
createScriptedRequester([{ ids: ['call_1', 'call_2'] }]),
undefined,
);
const result = await service.request({}, (part) => {
parts.push(part);
});
expect(result.message.toolCalls.map((c) => c.id)).toEqual(['call_1', 'call_2']);
expect(parts.filter(isToolCall).map((p) => p.id)).toEqual(['call_1', 'call_2']);
});
it('rewrites an id repeated across responses and keeps streamed parts consistent', async () => {
const parts: StreamedMessagePart[] = [];
const { service } = createService(
createScriptedRequester([{ ids: ['Bash_0'] }, { ids: ['Bash_0'] }]),
undefined,
);
const first = await service.request({}, (part) => {
parts.push(part);
});
const second = await service.request({}, (part) => {
parts.push(part);
});
expect(first.message.toolCalls[0]!.id).toBe('Bash_0');
expect(second.message.toolCalls[0]).toMatchObject({ id: 'Bash_0__2', rawId: 'Bash_0' });
expect(parts.filter(isToolCall).map((p) => [p.id, p.rawId])).toEqual([
['Bash_0', undefined],
['Bash_0__2', 'Bash_0'],
]);
});
it('rewrites duplicates within a single response', async () => {
const { service } = createService(
createScriptedRequester([{ ids: ['Bash_0', 'Bash_0'] }]),
undefined,
);
const result = await service.request();
expect(result.message.toolCalls.map((c) => [c.id, c.rawId])).toEqual([
['Bash_0', undefined],
['Bash_0__2', 'Bash_0'],
]);
});
it('rolls claims back when the attempt fails mid-stream', async () => {
const { service } = createService(
createScriptedRequester([
{ ids: ['Bash_9'], error: new Error('stream boom') },
{ ids: ['Bash_9'] },
]),
undefined,
);
await expect(service.request()).rejects.toThrow('stream boom');
const retry = await service.request();
expect(retry.message.toolCalls[0]!.id).toBe('Bash_9');
});
it('rewrites an id that already exists in the restored context', async () => {
const { service } = createService(
createScriptedRequester([{ ids: ['Bash_0'] }]),
undefined,
{
contextMessages: [
{
role: 'assistant',
content: [],
toolCalls: [{ type: 'function', id: 'Bash_0', name: 'Bash', arguments: '{}' }],
},
],
},
);
const result = await service.request();
expect(result.message.toolCalls[0]!.id).toBe('Bash_0__2');
});
});
describe('AgentLLMRequesterService attempt retry notification', () => {
afterEach(() => {
vi.useRealTimers();
});
it('notifies before resending with a repaired projection', async () => {
const calls = { value: 0 };
const { service } = createService(createRequester(calls), undefined);
const onAttemptRetry = vi.fn();
const result = await service.request({ onAttemptRetry });
expect(result.message.content).toEqual([{ type: 'text', text: 'ok' }]);
expect(calls.value).toBe(2);
expect(onAttemptRetry).toHaveBeenCalledTimes(1);
});
it('notifies before each indefinite-retry backoff', async () => {
vi.useFakeTimers();
const calls = { value: 0 };
const requester = createRequester(calls, new APIConnectionError('socket hang up'), [
new APIConnectionError('socket hang up again'),
]);
const { service } = createService(requester, undefined, {
env: { [KIMI_CODE_INFINITE_RETRY_ENV]: '1' },
});
const onAttemptRetry = vi.fn();
const promise = service.request({ onAttemptRetry });
await vi.runAllTimersAsync();
await promise;
expect(calls.value).toBe(3);
expect(onAttemptRetry).toHaveBeenCalledTimes(2);
});
it('does not notify when the error is final', async () => {
const calls = { value: 0 };
const { service } = createService(
createRequester(calls, new APIStatusError(400, 'max_tokens must be positive')),
undefined,
);
const onAttemptRetry = vi.fn();
await expect(service.request({ onAttemptRetry })).rejects.toMatchObject({ statusCode: 400 });
expect(onAttemptRetry).not.toHaveBeenCalled();
});
});
describe('turn machine stream state across service-internal retries', () => {
afterEach(() => {
vi.useRealTimers();
});
it('discards the interrupted attempt stream when the service retries below the turn', async () => {
vi.useFakeTimers();
const { service } = createService(
createScriptedRequester([
{ ids: ['call_a'], error: new APIConnectionError('terminated') },
{ ids: ['call_b'] },
]),
undefined,
{ env: { [KIMI_CODE_INFINITE_RETRY_ENV]: '1' } },
);
const machineRequester = createMachineRequester(service);
const doneEntries: AssistantEntry[] = [];
const actor = createActor(createTurnHarness(machineRequester.requester), {
input: { request: { model: turnHarnessModel }, history: [] },
});
actor.on('llm.done', (event) => doneEntries.push(event.entry));
actor.start();
await vi.runAllTimersAsync();
for (let index = 0; index < 10; index += 1) {
await vi.advanceTimersByTimeAsync(0);
}
expect(doneEntries).toHaveLength(1);
expect(doneEntries[0]?.message.toolCalls.map((toolCall) => toolCall.id)).toEqual(['call_b']);
});
});