Skip to content
Closed
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
5 changes: 5 additions & 0 deletions .changeset/flush-stream-json-output.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@moonshot-ai/kimi-code": patch
---

Respect stream JSON backpressure without allowing a stalled consumer to block cleanup or exit indefinitely.
41 changes: 38 additions & 3 deletions apps/kimi-code/src/cli/prompt-render.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,41 @@ interface RetryingEventLike {
export interface PromptOutput {
readonly columns?: number | undefined;
write(chunk: string): boolean;
once?(event: 'drain', listener: () => void): unknown;
}

interface PromptOutputDrainState {
pending: Promise<void> | undefined;
}

const outputDrainStates = new WeakMap<PromptOutput, PromptOutputDrainState>();

export function writePromptOutput(output: PromptOutput, chunk: string): void {
if (output.write(chunk) || output.once === undefined) return;

let state = outputDrainStates.get(output);
if (state === undefined) {
state = { pending: undefined };
outputDrainStates.set(output, state);
}
if (state.pending !== undefined) return;

let resolveDrain!: () => void;
state.pending = new Promise<void>((resolve) => {
resolveDrain = resolve;
});
output.once('drain', () => {
state.pending = undefined;
resolveDrain();
});
}

export async function drainPromptOutput(output: PromptOutput): Promise<void> {
while (true) {
const pending = outputDrainStates.get(output)?.pending;
if (pending === undefined) return;
await pending;
}
}

const PROMPT_BLOCK_BULLET = '• ';
Expand Down Expand Up @@ -262,7 +297,7 @@ export class PromptJsonWriter implements PromptTurnWriter {
private writeJsonLine(
message: PromptJsonAssistantMessage | PromptJsonToolMessage | PromptJsonRetryMetaMessage,
): void {
this.stdout.write(`${JSON.stringify(message)}\n`);
writePromptOutput(this.stdout, `${JSON.stringify(message)}\n`);
}
}

Expand Down Expand Up @@ -381,7 +416,7 @@ export function writeExperimentalVersion(
type: 'system.version',
version,
};
stdout.write(`${JSON.stringify(message)}\n`);
writePromptOutput(stdout, `${JSON.stringify(message)}\n`);
return;
}
stderr.write(`kimi version ${version}\n`);
Expand All @@ -403,7 +438,7 @@ export function writeResumeHint(
command,
content,
};
stdout.write(`${JSON.stringify(message)}\n`);
writePromptOutput(stdout, `${JSON.stringify(message)}\n`);
return;
}
stderr.write(`${content}\n`);
Expand Down
27 changes: 22 additions & 5 deletions apps/kimi-code/src/cli/run-prompt.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,11 @@ import {
} from '@moonshot-ai/kimi-code-sdk';
import { resolve } from 'pathe';

import { CLI_SHUTDOWN_TIMEOUT_MS, PROMPT_CLEANUP_TIMEOUT_MS } from '#/constant/app';
import {
CLI_SHUTDOWN_TIMEOUT_MS,
HEADLESS_STDIO_DRAIN_TIMEOUT_MS,
PROMPT_CLEANUP_TIMEOUT_MS,
} from '#/constant/app';

import { isKimiV2Enabled } from './experimental-v2';
import { resolveOutputFormat } from './options';
Expand All @@ -29,7 +33,13 @@ import {
type HeadlessGoalCreate,
} from './goal-prompt';
import type { PromptHarness, PromptSession } from './prompt-session';
import { PromptJsonWriter, PromptTranscriptWriter, writeResumeHint } from './prompt-render';
import {
drainPromptOutput,
PromptJsonWriter,
PromptTranscriptWriter,
writePromptOutput,
writeResumeHint,
} from './prompt-render';
import { createCliTelemetryBootstrap, initializeCliTelemetry } from './telemetry';
import { createKimiCodeHostIdentity } from './version';

Expand Down Expand Up @@ -146,15 +156,20 @@ export async function runPrompt(
let restorePromptSessionPermission = async (): Promise<void> => {};
let removeTerminationCleanup: (() => void) | undefined;
let cleanupPromise: Promise<void> | undefined;
let outputDrainAttempted = false;
const cleanupPromptRun = async (): Promise<void> => {
const pending = (cleanupPromise ??= (async () => {
removeTerminationCleanup?.();
setCrashPhase('shutdown');
try {
await restorePromptSessionPermission();
} finally {
await shutdownTelemetry({ timeoutMs: CLI_SHUTDOWN_TIMEOUT_MS });
await harness.close();
try {
await shutdownTelemetry({ timeoutMs: CLI_SHUTDOWN_TIMEOUT_MS });
await harness.close();
} finally {
if (!outputDrainAttempted) await drainPromptOutput(stdout);
}
}
})());
// Bound cleanup so a wedged shutdown step (e.g. a SessionEnd hook, MCP
Expand Down Expand Up @@ -213,6 +228,8 @@ export async function runPrompt(
);
}
writeResumeHint(session.id, outputFormat, stdout, stderr);
outputDrainAttempted = true;
await raceWithTimeout(drainPromptOutput(stdout), HEADLESS_STDIO_DRAIN_TIMEOUT_MS);

withTelemetryContext({ sessionId: session.id }).track('exit', {
duration_ms: Date.now() - startedAt,
Expand Down Expand Up @@ -268,7 +285,7 @@ async function runHeadlessGoal(
unsubscribeGoalEvents();
const snapshot = completedSnapshot ?? (await session.getGoal()).goal;
if (outputFormat === 'stream-json') {
stdout.write(`${JSON.stringify(goalSummaryJson(snapshot))}\n`);
writePromptOutput(stdout, `${JSON.stringify(goalSummaryJson(snapshot))}\n`);
} else {
stderr.write(`${formatGoalSummaryText(snapshot)}\n`);
}
Expand Down
18 changes: 14 additions & 4 deletions apps/kimi-code/src/cli/v2/run-v2-print.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ import { resolve } from 'pathe';
import {
CLI_SHUTDOWN_TIMEOUT_MS,
CLI_USER_AGENT_PRODUCT,
HEADLESS_STDIO_DRAIN_TIMEOUT_MS,
PROMPT_CLEANUP_TIMEOUT_MS,
} from '#/constant/app';

Expand All @@ -85,11 +86,13 @@ import { createKimiCodeHostIdentity } from '../version';
import { resolveOutputFormat } from '../options';
import type { CLIOptions, PromptOutputFormat } from '../options';
import {
drainPromptOutput,
type PromptOutput,
PromptJsonWriter,
type PromptTurnWriter,
PromptTranscriptWriter,
writeExperimentalVersion,
writePromptOutput,
writeResumeHint,
} from '../prompt-render';

Expand Down Expand Up @@ -165,16 +168,21 @@ export async function runV2Print(
let removeTerminationCleanup: (() => void) | undefined;
let cleanupPromise: Promise<void> | undefined;
let telemetryService: ITelemetryService | undefined;
let outputDrainAttempted = false;
const cleanup = async (): Promise<void> => {
const pending = (cleanupPromise ??= (async () => {
removeTerminationCleanup?.();
try {
await restorePermission();
} finally {
if (telemetryService !== undefined) {
await raceWithTimeout(telemetryService.shutdown(), CLI_SHUTDOWN_TIMEOUT_MS);
try {
if (telemetryService !== undefined) {
await raceWithTimeout(telemetryService.shutdown(), CLI_SHUTDOWN_TIMEOUT_MS);
}
app.dispose();
} finally {
if (!outputDrainAttempted) await drainPromptOutput(stdout);
}
app.dispose();
}
})());
await raceWithTimeout(pending, PROMPT_CLEANUP_TIMEOUT_MS);
Expand Down Expand Up @@ -232,6 +240,8 @@ export async function runV2Print(
);
}
writeResumeHint(resolved.session.id, outputFormat, stdout, stderr);
outputDrainAttempted = true;
await raceWithTimeout(drainPromptOutput(stdout), HEADLESS_STDIO_DRAIN_TIMEOUT_MS);

telemetryService.withContext({ sessionId: resolved.session.id }).track2('exit', {
duration_ms: Date.now() - startedAt,
Expand Down Expand Up @@ -535,7 +545,7 @@ async function runNativeGoal(
subscription.dispose();
const snapshot = completedSnapshot ?? goalService.getGoal().goal;
if (outputFormat === 'stream-json') {
stdout.write(`${JSON.stringify(goalSummaryJson(snapshot))}\n`);
writePromptOutput(stdout, `${JSON.stringify(goalSummaryJson(snapshot))}\n`);
} else {
stderr.write(`${formatGoalSummaryText(snapshot)}\n`);
}
Expand Down
Loading