diff --git a/apps/web/src/lib/code-reviews/terminal-reason-from-failure.ts b/apps/web/src/lib/code-reviews/terminal-reason-from-failure.ts index 53f79ca41a..ee47654ee7 100644 --- a/apps/web/src/lib/code-reviews/terminal-reason-from-failure.ts +++ b/apps/web/src/lib/code-reviews/terminal-reason-from-failure.ts @@ -66,6 +66,7 @@ const FAILURE_CODE_REASONS = { assistant_error: 'assistant_failed', missing_assistant_reply: 'assistant_no_reply', payment_required: 'billing', + kilo_output_limit: 'assistant_failed', user_interrupt: 'user_cancelled', container_shutdown: 'container_shutdown', system_interrupt: 'interrupted', diff --git a/packages/db/src/schema.ts b/packages/db/src/schema.ts index 2835eb744e..0075bfcc67 100644 --- a/packages/db/src/schema.ts +++ b/packages/db/src/schema.ts @@ -5428,6 +5428,7 @@ export type CloudAgentSessionRunFailureCode = | 'wrapper_error_after_activity' | 'missing_assistant_reply' | 'payment_required' + | 'kilo_output_limit' | 'user_interrupt' | 'container_shutdown' | 'system_interrupt' diff --git a/packages/worker-utils/src/cloud-agent-failure.test.ts b/packages/worker-utils/src/cloud-agent-failure.test.ts index 1ce98df5e4..1d88e0cb7c 100644 --- a/packages/worker-utils/src/cloud-agent-failure.test.ts +++ b/packages/worker-utils/src/cloud-agent-failure.test.ts @@ -44,6 +44,18 @@ describe('classifyCloudAgentFailure', () => { }); }); + it.each([ + 'kilo_output_limit', + 'wrapper_no_output', + 'wrapper_ping_timeout', + 'wrapper_disconnected', + ] as const)('classifies %s as platform wrapper_liveness', code => { + expect(classifyCloudAgentFailure({ source: 'run', stage: 'agent_activity', code })).toEqual({ + responsibility: 'platform', + reason: 'wrapper_liveness', + }); + }); + it('preserves ambiguous assistant and source-control failures as unknown', () => { expect( classifyCloudAgentFailure({ diff --git a/packages/worker-utils/src/cloud-agent-failure.ts b/packages/worker-utils/src/cloud-agent-failure.ts index e19030edd7..c73b5b2ecc 100644 --- a/packages/worker-utils/src/cloud-agent-failure.ts +++ b/packages/worker-utils/src/cloud-agent-failure.ts @@ -28,6 +28,7 @@ export const CLOUD_AGENT_FAILURE_CODES = [ 'wrapper_error_after_activity', 'missing_assistant_reply', 'payment_required', + 'kilo_output_limit', 'user_interrupt', 'container_shutdown', 'system_interrupt', @@ -234,6 +235,7 @@ export function classifyCloudAgentFailure( case 'wrapper_error_before_activity': case 'wrapper_error_after_activity': case 'missing_assistant_reply': + case 'kilo_output_limit': return classified('platform', 'wrapper_liveness'); case 'assistant_error': case 'payment_required': diff --git a/packages/worker-utils/src/cloud-agent-queue-report.ts b/packages/worker-utils/src/cloud-agent-queue-report.ts index 81595ca45a..bb7f6eb73b 100644 --- a/packages/worker-utils/src/cloud-agent-queue-report.ts +++ b/packages/worker-utils/src/cloud-agent-queue-report.ts @@ -31,10 +31,15 @@ export const CloudAgentRunFailureClassifications = [ { failureStage: 'post_dispatch_no_activity', failureCode: 'missing_assistant_reply' }, { failureStage: 'post_dispatch_no_activity', failureCode: 'payment_required' }, { failureStage: 'post_dispatch_no_activity', failureCode: 'model_missing' }, + { failureStage: 'post_dispatch_no_activity', failureCode: 'kilo_output_limit' }, { failureStage: 'agent_activity', failureCode: 'assistant_error' }, { failureStage: 'agent_activity', failureCode: 'payment_required' }, { failureStage: 'agent_activity', failureCode: 'model_missing' }, { failureStage: 'agent_activity', failureCode: 'wrapper_error_after_activity' }, + { failureStage: 'agent_activity', failureCode: 'wrapper_no_output' }, + { failureStage: 'agent_activity', failureCode: 'wrapper_ping_timeout' }, + { failureStage: 'agent_activity', failureCode: 'wrapper_disconnected' }, + { failureStage: 'agent_activity', failureCode: 'kilo_output_limit' }, { failureStage: 'interruption', failureCode: 'user_interrupt' }, { failureStage: 'interruption', failureCode: 'container_shutdown' }, { failureStage: 'interruption', failureCode: 'system_interrupt' }, diff --git a/services/cloud-agent-next/src/session/safe-failure-projection.ts b/services/cloud-agent-next/src/session/safe-failure-projection.ts index 1a20dc284d..943c92b2eb 100644 --- a/services/cloud-agent-next/src/session/safe-failure-projection.ts +++ b/services/cloud-agent-next/src/session/safe-failure-projection.ts @@ -50,6 +50,7 @@ const GENERIC_FAILURE_MESSAGES = { wrapper_error_after_activity: 'Agent wrapper failed while processing the message', missing_assistant_reply: 'No assistant reply was produced', payment_required: 'Assistant request failed: insufficient credits', + kilo_output_limit: 'Assistant response hit the output length limit', user_interrupt: 'The message was interrupted by the user', container_shutdown: 'The agent container shut down', system_interrupt: 'The message was interrupted', diff --git a/services/cloud-agent-next/src/session/wrapper-supervisor.test.ts b/services/cloud-agent-next/src/session/wrapper-supervisor.test.ts index 6baac57da2..7701ef838d 100644 --- a/services/cloud-agent-next/src/session/wrapper-supervisor.test.ts +++ b/services/cloud-agent-next/src/session/wrapper-supervisor.test.ts @@ -971,6 +971,27 @@ describe('WrapperSupervisor', () => { expect(harness.events.map(event => event.streamEventType)).toEqual(['cloud.message.failed']); }); + it('preserves wrapper_no_output after agent activity', async () => { + const acceptedAt = 2_000; + const noOutputDeadlineAt = acceptedAt + WRAPPER_NO_OUTPUT_TIMEOUT_MS; + const harness = createHarness([ + liveRuntimeState({ noOutputDeadlineAt, nextPingAt: noOutputDeadlineAt + 1 }), + OWNED_WRAPPER_LEASE, + ]); + await putSessionMessageState(harness.storage, { + ...acceptedMessage(), + agentActivityObservedAt: acceptedAt + 1_000, + }); + + await harness.supervisor.runMaintenance(noOutputDeadlineAt); + + await expect(getSessionMessageState(harness.storage, MESSAGE_ID)).resolves.toMatchObject({ + status: 'failed', + failureStage: 'agent_activity', + failureCode: 'wrapper_no_output', + }); + }); + it('terminates an unresponsive wrapper on ping timeout before no-output expires', async () => { const pingDeadlineAt = 92_000; const noOutputDeadlineAt = 332_000; @@ -989,6 +1010,49 @@ describe('WrapperSupervisor', () => { }); }); + it('preserves wrapper_ping_timeout after agent activity', async () => { + const pingDeadlineAt = 92_000; + const noOutputDeadlineAt = 332_000; + const harness = createHarness([ + liveRuntimeState({ pingDeadlineAt, noOutputDeadlineAt }), + OWNED_WRAPPER_LEASE, + ]); + await putSessionMessageState(harness.storage, { + ...acceptedMessage(), + agentActivityObservedAt: 9_000, + }); + + await harness.supervisor.runMaintenance(pingDeadlineAt); + + await expect(getSessionMessageState(harness.storage, MESSAGE_ID)).resolves.toMatchObject({ + status: 'failed', + failureStage: 'agent_activity', + failureCode: 'wrapper_ping_timeout', + }); + }); + + it('preserves kilo_output_limit after agent activity', async () => { + const harness = createHarness([liveRuntimeState(), OWNED_WRAPPER_LEASE]); + await putSessionMessageState(harness.storage, { + ...acceptedMessage(), + agentActivityObservedAt: 9_000, + }); + + await harness.supervisor.onTerminalEvent({ + wrapperRunId: WRAPPER_RUN_ID, + status: 'failed', + errorSource: 'assistant', + failureCode: 'kilo_output_limit', + error: 'Assistant response hit the output length limit', + }); + + await expect(getSessionMessageState(harness.storage, MESSAGE_ID)).resolves.toMatchObject({ + status: 'failed', + failureStage: 'agent_activity', + failureCode: 'kilo_output_limit', + }); + }); + it('defers liveness failure while disconnect grace is active for the current connection', async () => { const pingDeadlineAt = 92_000; const noOutputDeadlineAt = 332_000; diff --git a/services/cloud-agent-next/src/session/wrapper-supervisor.ts b/services/cloud-agent-next/src/session/wrapper-supervisor.ts index 8dee60044b..d9682770aa 100644 --- a/services/cloud-agent-next/src/session/wrapper-supervisor.ts +++ b/services/cloud-agent-next/src/session/wrapper-supervisor.ts @@ -12,6 +12,7 @@ import type { MessageSettlementOutbox } from './message-settlement-outbox.js'; import { classifyAssistantFailure, classifyAssistantFailureMessage, + genericFailureMessage, } from './safe-failure-projection.js'; import { countPendingSessionMessages, type SessionQueueStorage } from './pending-messages.js'; import type { SessionMessageQueue } from './session-message-queue.js'; @@ -649,7 +650,7 @@ export function createWrapperSupervisor( error, completionSource: 'wrapper_failure', failureStage: activityObserved ? 'agent_activity' : 'post_dispatch_no_activity', - failureCode: activityObserved ? 'wrapper_error_after_activity' : failureCode, + failureCode, }); } await messageSettlementOutbox.releaseWrapperTerminalWaitForIdleBatch(); @@ -728,7 +729,7 @@ export function createWrapperSupervisor( error: 'Wrapper disconnected', completionSource: 'wrapper_failure', failureStage: activityObserved ? 'agent_activity' : 'post_dispatch_no_activity', - failureCode: activityObserved ? 'wrapper_error_after_activity' : 'wrapper_disconnected', + failureCode: 'wrapper_disconnected', }); } await clearWrapperRuntimeIdentity( @@ -1431,17 +1432,23 @@ export function createWrapperSupervisor( if (status === 'failed') { if (errorSource === 'assistant') { const assistantFailure = classifyAssistantFailure(error); + const failureCode = + terminalFailureCode ?? assistantFailure.terminalCode ?? 'assistant_error'; + const usesExplicitTerminalCode = terminalFailureCode === 'kilo_output_limit'; await messageSettlementOutbox.terminalizeSessionMessageOnce(message.messageId, { kind: 'failed', reason: 'assistant_error', error: error ?? 'Assistant request failed', completionSource: 'wrapper_failure', failureStage: 'agent_activity', - failureCode: - terminalFailureCode ?? assistantFailure.terminalCode ?? 'assistant_error', - assistantFailureReason: assistantFailure.reason, - providerOwnership: assistantFailure.providerOwnership, - safeFailureMessage: assistantFailure.safeMessage, + failureCode, + ...(usesExplicitTerminalCode + ? { safeFailureMessage: genericFailureMessage(failureCode) } + : { + assistantFailureReason: assistantFailure.reason, + providerOwnership: assistantFailure.providerOwnership, + safeFailureMessage: assistantFailure.safeMessage, + }), ...(persistedModelNotFoundDiagnostics ? { modelNotFoundRuntimeDiagnostics: persistedModelNotFoundDiagnostics } : {}), diff --git a/services/cloud-agent-next/src/shared/protocol.ts b/services/cloud-agent-next/src/shared/protocol.ts index d049acb544..d873ed0784 100644 --- a/services/cloud-agent-next/src/shared/protocol.ts +++ b/services/cloud-agent-next/src/shared/protocol.ts @@ -53,7 +53,11 @@ export type IngestEvent = { data: unknown; }; -export const WrapperTerminalFailureCodes = ['payment_required', 'model_missing'] as const; +export const WrapperTerminalFailureCodes = [ + 'payment_required', + 'model_missing', + 'kilo_output_limit', +] as const; export type WrapperTerminalFailureCode = (typeof WrapperTerminalFailureCodes)[number]; /** diff --git a/services/cloud-agent-next/src/telemetry/queue-reports.ts b/services/cloud-agent-next/src/telemetry/queue-reports.ts index 76cfb5ff69..fd70cb4c52 100644 --- a/services/cloud-agent-next/src/telemetry/queue-reports.ts +++ b/services/cloud-agent-next/src/telemetry/queue-reports.ts @@ -47,6 +47,7 @@ const FAILED_RUN_DIAGNOSTIC_MESSAGES: Partial< wrapper_error_after_activity: 'Wrapper failed after agent activity', missing_assistant_reply: 'No assistant reply was produced', payment_required: 'Model request failed: insufficient credits', + kilo_output_limit: 'Assistant response hit the output length limit', unclassified: 'Run failed without a classified cause', }; diff --git a/services/cloud-agent-next/src/websocket/ingest.test.ts b/services/cloud-agent-next/src/websocket/ingest.test.ts index 310a17029b..7fc931ebc2 100644 --- a/services/cloud-agent-next/src/websocket/ingest.test.ts +++ b/services/cloud-agent-next/src/websocket/ingest.test.ts @@ -1233,30 +1233,65 @@ describe('createIngestHandler', () => { }); it.each([ - { failureCode: 'payment_required' as const, error: 'Insufficient credits' }, - { failureCode: 'model_missing' as const, error: 'Model not found' }, - ])('forwards $failureCode wrapper failures to the session coordinator', async failure => { - const doContext = createNewPathDOContext(); - const handler = createIngestHandler( - createFakeState(), - createFakeEventQueries(), - SESSION_ID, - vi.fn(), - doContext - ); - const ws = createFakeWebSocket(makeNewPathAttachment()); + { + failureCode: 'payment_required' as const, + error: 'Insufficient credits', + safeMessage: 'Assistant request failed: insufficient credits', + }, + { + failureCode: 'model_missing' as const, + error: 'Model not found', + safeMessage: 'Assistant request failed: model not found', + }, + { + failureCode: 'kilo_output_limit' as const, + error: 'Assistant response hit the output length limit', + safeMessage: 'Assistant response hit the output length limit', + }, + ])( + 'forwards $failureCode wrapper failures to the session coordinator', + async ({ safeMessage, ...failure }) => { + const doContext = createNewPathDOContext(); + const eventQueries = createFakeEventQueries(); + const broadcast = vi.fn(); + const handler = createIngestHandler( + createFakeState(), + eventQueries, + SESSION_ID, + broadcast, + doContext + ); + const ws = createFakeWebSocket(makeNewPathAttachment()); - await handler.handleIngestMessage( - ws, - makeStreamMessage('error', { fatal: true, ...failure }) - ); + await handler.handleIngestMessage( + ws, + makeStreamMessage('error', { fatal: true, errorSource: 'assistant', ...failure }) + ); - expect(doContext.handleWrapperTerminalEvent).toHaveBeenCalledWith({ - wrapperRunId: WRAPPER_RUN_ID, - status: 'failed', - ...failure, - }); - }); + expect(doContext.handleWrapperTerminalEvent).toHaveBeenCalledWith({ + wrapperRunId: WRAPPER_RUN_ID, + status: 'failed', + errorSource: 'assistant', + ...failure, + }); + expect(broadcast).toHaveBeenCalledWith( + expect.objectContaining({ + stream_event_type: 'cloud.status', + payload: JSON.stringify({ cloudStatus: { type: 'error', message: safeMessage } }), + }) + ); + expect(eventQueries.insert).toHaveBeenCalledWith( + expect.objectContaining({ + payload: JSON.stringify({ + fatal: true, + errorSource: 'assistant', + error: safeMessage, + message: safeMessage, + }), + }) + ); + } + ); it('does NOT terminalize on wrapper complete event (new path)', async () => { const state = createFakeState(); diff --git a/services/cloud-agent-next/src/websocket/ingest.ts b/services/cloud-agent-next/src/websocket/ingest.ts index 1b0b90695a..578abe106d 100644 --- a/services/cloud-agent-next/src/websocket/ingest.ts +++ b/services/cloud-agent-next/src/websocket/ingest.ts @@ -26,6 +26,7 @@ import { } from '../session/ingest-handlers/index.js'; import { WrapperTerminalFailureCodes, + type WrapperTerminalFailureCode, type CompleteEventData, type KilocodeEventData, type CloudStatusData, @@ -37,6 +38,7 @@ import type { TerminalizeParams } from '../session/session-message-state.js'; import { classifyAssistantFailure, classifyAssistantFailureMessage, + genericFailureMessage, } from '../session/safe-failure-projection.js'; import { parseModelNotFoundRuntimeDiagnostics } from '../shared/runtime-model-diagnostics.js'; import { @@ -153,6 +155,22 @@ function sanitizeKilocodeEventData(data: unknown): unknown { return data; } +/** + * The wrapper sends fixed text for kilo_output_limit that text classification + * cannot recognize, so resolve it from the structured code. + * payment_required/model_missing keep text classification: the model-not-found + * diagnostics gate keys off the classified safe message. + */ +function safeAssistantFailureMessage( + failureCode: WrapperTerminalFailureCode | undefined, + rawError: unknown +): string { + if (failureCode === 'kilo_output_limit') { + return genericFailureMessage(failureCode); + } + return classifyAssistantFailureMessage(rawError); +} + function sanitizePublicEventData(eventType: string, data: unknown): unknown { if (eventType === 'kilocode') return sanitizeKilocodeEventData(data); @@ -162,7 +180,7 @@ function sanitizePublicEventData(eventType: string, data: unknown): unknown { const rawError = parsed.data.error ?? parsed.data.message; const safeMessage = parsed.data.errorSource === 'assistant' - ? classifyAssistantFailureMessage(rawError) + ? safeAssistantFailureMessage(parsed.data.failureCode, rawError) : 'Agent wrapper failed'; return { fatal: parsed.data.fatal, @@ -923,7 +941,7 @@ export function createIngestHandler( const fatalMessage = errorData.error ?? errorData.message ?? 'Fatal error'; const safeFatalMessage = errorData.errorSource === 'assistant' - ? classifyAssistantFailureMessage(fatalMessage) + ? safeAssistantFailureMessage(errorData.failureCode, fatalMessage) : 'Agent wrapper failed'; const shouldForwardModelNotFoundDiagnostics = errorData.errorSource === 'assistant' && diff --git a/services/cloud-agent-next/test/unit/wrapper/terminal-failure-attribution.test.ts b/services/cloud-agent-next/test/unit/wrapper/terminal-failure-attribution.test.ts new file mode 100644 index 0000000000..cec9d35d3b --- /dev/null +++ b/services/cloud-agent-next/test/unit/wrapper/terminal-failure-attribution.test.ts @@ -0,0 +1,294 @@ +/** + * Unit tests for Tier-1 failure attribution in createConnectionManager: + * output-limit termination (finish === "length" / MessageOutputLengthError). + */ + +import { describe, expect, it, vi, beforeEach, afterEach } from 'vitest'; +import { + createConnectionManager, + type ConnectionCallbacks, +} from '../../../wrapper/src/connection.js'; +import { WrapperState, type SessionContext } from '../../../wrapper/src/state.js'; +import type { KiloEvent, WrapperKiloClient } from '../../../wrapper/src/kilo-api.js'; + +if (typeof CloseEvent === 'undefined') { + const g = globalThis as Record; + g.CloseEvent = class extends Event { + code: number; + reason: string; + wasClean: boolean; + constructor(type: string, init?: { code?: number; reason?: string; wasClean?: boolean }) { + super(type); + this.code = init?.code ?? 0; + this.reason = init?.reason ?? ''; + this.wasClean = init?.wasClean ?? false; + } + }; +} + +if (typeof MessageEvent === 'undefined') { + const g = globalThis as Record; + g.MessageEvent = class extends Event { + data: unknown; + constructor(type: string, init?: { data?: unknown }) { + super(type); + this.data = init?.data; + } + }; +} + +class MockWebSocket { + static CONNECTING = 0; + static OPEN = 1; + static CLOSING = 2; + static CLOSED = 3; + static instances: MockWebSocket[] = []; + + readyState = MockWebSocket.CONNECTING; + onopen: ((event: Event) => void) | null = null; + onclose: ((event: CloseEvent) => void) | null = null; + onerror: ((event: Event) => void) | null = null; + onmessage: ((event: MessageEvent) => void) | null = null; + + sent: string[] = []; + url: string; + + constructor(url: string, _options?: unknown) { + this.url = url; + MockWebSocket.instances.push(this); + } + + send(data: string): void { + this.sent.push(data); + } + + close(_code?: number, _reason?: string): void { + this.readyState = MockWebSocket.CLOSED; + } + + simulateOpen(): void { + this.readyState = MockWebSocket.OPEN; + this.onopen?.(new Event('open')); + } + + static reset(): void { + MockWebSocket.instances = []; + } + + static get latest(): MockWebSocket | undefined { + return MockWebSocket.instances[MockWebSocket.instances.length - 1]; + } +} + +const ROOT_SESSION_ID = 'kilo_sess_456'; +const CHILD_SESSION_ID = 'kilo_sess_child'; +const ASSISTANT_MESSAGE_ID = 'assistant_msg_root_123'; + +const createSessionContext = (overrides: Partial = {}): SessionContext => ({ + kiloSessionId: ROOT_SESSION_ID, + ingestUrl: 'wss://ingest.example.com/ingest', + ingestToken: 'token_secret', + workerAuthToken: 'kilo_token_789', + wrapperRunId: 'run_test', + wrapperGeneration: 7, + wrapperConnectionId: 'conn_test', + ...overrides, +}); + +const createCallbacks = (): ConnectionCallbacks & { + onTerminalError: ReturnType; + onCompletionSignal: ReturnType; + onSessionIdle: ReturnType; + onDisconnect: ReturnType; +} => ({ + onTerminalError: vi.fn(), + onCommand: vi.fn(), + onDisconnect: vi.fn(), + onCompletionSignal: vi.fn(), + onSessionIdle: vi.fn(), + onSseEvent: vi.fn(), +}); + +const createMockKiloClient = (overrides: Partial = {}): WrapperKiloClient => ({ + createSession: vi.fn().mockResolvedValue({ id: 'kilo_sess' }), + getSession: vi.fn().mockResolvedValue({ id: 'kilo_sess' }), + sendPromptAsync: vi.fn().mockResolvedValue(undefined), + abortSession: vi.fn().mockResolvedValue(true), + summarizeSession: vi.fn().mockResolvedValue(true), + sendCommand: vi.fn().mockResolvedValue(undefined), + answerPermission: vi.fn().mockResolvedValue(true), + answerQuestion: vi.fn().mockResolvedValue(true), + rejectQuestion: vi.fn().mockResolvedValue(true), + generateCommitMessage: vi.fn().mockResolvedValue({ message: 'test commit' }), + getSessionStatuses: vi.fn().mockResolvedValue({}), + getQuestions: vi.fn().mockResolvedValue([]), + getPermissions: vi.fn().mockResolvedValue([]), + getNetworkWaits: vi.fn().mockResolvedValue([]), + resumeNetworkWait: vi.fn().mockResolvedValue(true), + listEffectiveModels: vi.fn().mockResolvedValue([]), + subscribeEvents: vi.fn().mockResolvedValue({ + stream: (async function* () { + await new Promise(() => {}); + })(), + }), + serverUrl: 'http://127.0.0.1:0', + ...overrides, +}); + +function createEventStream(events: KiloEvent[]): AsyncIterable { + return (async function* () { + for (const event of events) { + yield event; + } + await new Promise(() => {}); + })(); +} + +function stubFetch(): void { + vi.stubGlobal( + 'fetch', + vi.fn().mockImplementation(() => { + const stream = new ReadableStream({ + start() {}, + }); + return Promise.resolve( + new Response(stream, { + status: 200, + headers: { 'Content-Type': 'text/event-stream' }, + }) + ); + }) + ); +} + +async function openConnection( + manager: ReturnType +): Promise { + const openPromise = manager.open(); + const ws = MockWebSocket.latest!; + ws.simulateOpen(); + await openPromise; + return ws; +} + +function rootAssistantUpdated(overrides: { + finish?: string; + sessionID?: string; + error?: { name: string }; +}): KiloEvent { + return { + type: 'message.updated', + properties: { + info: { + id: ASSISTANT_MESSAGE_ID, + role: 'assistant', + sessionID: overrides.sessionID ?? ROOT_SESSION_ID, + finish: overrides.finish, + time: { completed: 1_716_200_000_000 }, + ...(overrides.error ? { error: overrides.error } : {}), + }, + }, + }; +} + +function rootIdle(): KiloEvent { + return { type: 'session.idle', properties: { sessionID: ROOT_SESSION_ID } }; +} + +describe('terminal failure attribution', () => { + let state: WrapperState; + let callbacks: ReturnType; + + beforeEach(() => { + vi.useFakeTimers(); + MockWebSocket.reset(); + vi.stubGlobal('WebSocket', MockWebSocket); + stubFetch(); + state = new WrapperState(); + state.bindSession(createSessionContext()); + callbacks = createCallbacks(); + }); + + afterEach(() => { + vi.useRealTimers(); + vi.unstubAllGlobals(); + }); + + async function runEvents(events: KiloEvent[], session: SessionContext = createSessionContext()) { + state = new WrapperState(); + state.bindSession(session); + const kiloClient = createMockKiloClient({ + subscribeEvents: vi.fn().mockResolvedValue({ + stream: createEventStream(events), + }), + }); + const manager = createConnectionManager(state, { kiloClient }, callbacks); + await openConnection(manager); + await vi.advanceTimersByTimeAsync(0); + await vi.advanceTimersByTimeAsync(0); + } + + it('raises kilo_output_limit for finish=length on the root session', async () => { + await runEvents([rootAssistantUpdated({ finish: 'length' })]); + + expect(callbacks.onTerminalError).toHaveBeenCalledWith({ + code: 'kilo_output_limit', + message: 'Assistant response hit the output length limit', + errorSource: 'assistant', + }); + expect(callbacks.onCompletionSignal).not.toHaveBeenCalled(); + }); + + it('raises kilo_output_limit for MessageOutputLengthError on the root session', async () => { + await runEvents([rootAssistantUpdated({ error: { name: 'MessageOutputLengthError' } })]); + + expect(callbacks.onTerminalError).toHaveBeenCalledWith({ + code: 'kilo_output_limit', + message: 'Assistant response hit the output length limit', + errorSource: 'assistant', + }); + expect(callbacks.onCompletionSignal).not.toHaveBeenCalled(); + }); + + it('does not raise kilo_output_limit for finish=length on a child session', async () => { + await runEvents([ + rootAssistantUpdated({ finish: 'length', sessionID: CHILD_SESSION_ID }), + rootIdle(), + ]); + + expect(callbacks.onTerminalError).not.toHaveBeenCalled(); + expect(callbacks.onCompletionSignal).toHaveBeenCalled(); + }); + + it('completes finish=stop with no visible text as success', async () => { + await runEvents([rootAssistantUpdated({ finish: 'stop' }), rootIdle()]); + + expect(callbacks.onTerminalError).not.toHaveBeenCalled(); + expect(callbacks.onCompletionSignal).toHaveBeenCalled(); + expect(callbacks.onSessionIdle).toHaveBeenCalled(); + }); + + it('skips output-limit detection while finalizing', async () => { + state = new WrapperState(); + state.bindSession(createSessionContext()); + state.acceptMessage('msg_1', { + autoCommit: false, + condenseOnComplete: true, + }); + expect(state.beginFinalizing()).toBe(true); + + const kiloClient = createMockKiloClient({ + subscribeEvents: vi.fn().mockResolvedValue({ + stream: createEventStream([rootAssistantUpdated({ finish: 'length' }), rootIdle()]), + }), + }); + const manager = createConnectionManager(state, { kiloClient }, callbacks); + await openConnection(manager); + await vi.advanceTimersByTimeAsync(0); + await vi.advanceTimersByTimeAsync(0); + + expect(callbacks.onTerminalError).not.toHaveBeenCalled(); + expect(callbacks.onCompletionSignal).toHaveBeenCalled(); + expect(callbacks.onSessionIdle).toHaveBeenCalled(); + }); +}); diff --git a/services/cloud-agent-next/wrapper/src/connection.ts b/services/cloud-agent-next/wrapper/src/connection.ts index d906b3e3d0..c92fd1f09a 100644 --- a/services/cloud-agent-next/wrapper/src/connection.ts +++ b/services/cloud-agent-next/wrapper/src/connection.ts @@ -230,10 +230,19 @@ export function isSessionIdleEvent( function isAssistantCompletionSignal(info: unknown): boolean { if (!isRecord(info) || info.role !== 'assistant') return false; + if (isOutputLimitTermination(info)) return false; const time = isRecord(info.time) ? info.time : undefined; return typeof time?.completed === 'number' || (info.error !== undefined && info.error !== null); } +function isOutputLimitTermination(info: Record): boolean { + if (info.finish === 'length') return true; + const error = info.error; + return isRecord(error) && error.name === 'MessageOutputLengthError'; +} + +const KILO_OUTPUT_LIMIT_MESSAGE = 'Assistant response hit the output length limit'; + // --------------------------------------------------------------------------- // Types // --------------------------------------------------------------------------- @@ -1220,6 +1229,18 @@ export function createConnectionManager( const currentSessionId = state.currentSession?.kiloSessionId; if (!currentSessionId || msgSessionId === currentSessionId) { state.setLastAssistantMessageId(messageInfo.id); + // finish === "length" must win over completion-signal arming so a + // terminalized run does not also arm post-completion waiters. + // Skip while finalizing: post-completion turns (e.g. condense) must + // not rewrite a successful run into kilo_output_limit. + if (isOutputLimitTermination(messageInfo) && !state.isFinalizing) { + callbacks.onTerminalError({ + code: 'kilo_output_limit', + message: KILO_OUTPUT_LIMIT_MESSAGE, + errorSource: 'assistant', + }); + return; + } } } if (isAssistantCompletionSignal(messageInfo)) {