Skip to content
Open
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/calm-ravens-share.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@moonshot-ai/kimi-code": patch
---

Share one telemetry shutdown deadline across CLI-owned pipelines.
14 changes: 13 additions & 1 deletion apps/kimi-code/src/cli/run-prompt.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,10 @@ import {
import type { PromptHarness, PromptSession } from './prompt-session';
import { PromptJsonWriter, PromptTranscriptWriter, writeResumeHint } from './prompt-render';
import { createCliTelemetryBootstrap, initializeCliTelemetry } from './telemetry';
import {
createTelemetryShutdownDeadline,
formatTelemetryShutdownWarning,
} from './telemetry-shutdown';
import { createKimiCodeHostIdentity } from './version';

/**
Expand Down Expand Up @@ -153,7 +157,15 @@ export async function runPrompt(
try {
await restorePromptSessionPermission();
} finally {
await shutdownTelemetry({ timeoutMs: CLI_SHUTDOWN_TIMEOUT_MS });
const telemetryDeadline = createTelemetryShutdownDeadline(
CLI_SHUTDOWN_TIMEOUT_MS,
(error) => {
stderr.write(formatTelemetryShutdownWarning(error));
},
);
await telemetryDeadline.run((remainingMs) =>
shutdownTelemetry({ timeoutMs: remainingMs }),
);
await harness.close();
}
})());
Expand Down
37 changes: 25 additions & 12 deletions apps/kimi-code/src/cli/sub/web/run.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,11 @@ import { existsSync } from 'node:fs';
import { join } from 'node:path';

import { hostRequestHeadersSeed } from '@moonshot-ai/agent-core-v2';
import { createServerLogger, startServer, type ServerLogger } from '@moonshot-ai/kap-server';
import {
createServerLogger,
startServer,
type ServerLogger,
} from '@moonshot-ai/kap-server';
import { shutdownTelemetry, track } from '@moonshot-ai/kimi-telemetry';
import chalk from 'chalk';
import { type Command } from 'commander';
Expand All @@ -24,6 +28,7 @@ import { openUrl as defaultOpenUrl } from '#/utils/open-url';
import { getDataDir } from '#/utils/paths';

import { initializeServerTelemetry } from '../../telemetry';
import { createTelemetryShutdownDeadline } from '../../telemetry-shutdown';
import {
buildKimiDefaultHeaders,
getHostPackageRoot,
Expand Down Expand Up @@ -58,7 +63,7 @@ const WEB_ASSETS_DIR = 'dist-web';
interface RoutedServer {
readonly address: string;
readonly logger: ServerLogger;
close(): Promise<void>;
close(options?: { readonly telemetryDeadlineMs?: number }): Promise<void>;
}

export interface WebCliOptions extends ServerCliOptions {
Expand Down Expand Up @@ -250,15 +255,23 @@ async function runServerInProcess(
if (stopping) return;
stopping = true;
running?.logger.info({ reason }, 'server shutting down');
try {
await running?.close();
await shutdownTelemetry({ timeoutMs: CLI_SHUTDOWN_TIMEOUT_MS });
} catch (error) {
running?.logger.error(
{ err: error instanceof Error ? error : new Error(String(error)) },
'server shutdown error',
);
}
const telemetryDeadline = createTelemetryShutdownDeadline(
CLI_SHUTDOWN_TIMEOUT_MS,
(error) => {
running?.logger.error(
{ err: error instanceof Error ? error : new Error(String(error)) },
'telemetry shutdown error',
);
},
);
await telemetryDeadline.run(async () =>
running?.close({
telemetryDeadlineMs: telemetryDeadline.expiresAtMs,
}),
);
await telemetryDeadline.run((remainingMs) =>
shutdownTelemetry({ timeoutMs: remainingMs }),
);
process.exit(0);
}

Expand Down Expand Up @@ -301,7 +314,7 @@ async function runServerInProcess(
running = {
address: `http://${v2.host}:${v2.port}`,
logger,
close: () => v2.close(),
close: (closeOptions) => v2.close(closeOptions),
};

track('server_started', { daemon: false });
Expand Down
41 changes: 41 additions & 0 deletions apps/kimi-code/src/cli/telemetry-shutdown.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
export interface TelemetryShutdownDeadline {
readonly expiresAtMs: number;
remainingMs(): number;
run(pipeline: (remainingMs: number) => Promise<void>): Promise<void>;
}

const TELEMETRY_SHUTDOWN_WARNING_PREFIX = 'Warning: telemetry shutdown failed:';

export function formatTelemetryShutdownWarning(error: unknown): string {
const message = error instanceof Error ? error.message : String(error);
return `${TELEMETRY_SHUTDOWN_WARNING_PREFIX} ${message}\n`;
}

/**
* One process-exit budget shared by every telemetry pipeline owned by a CLI
* entrypoint. Pipelines remain independent: each receives the current budget,
* and a failure never prevents the next owner from attempting its cleanup.
*/
export function createTelemetryShutdownDeadline(
timeoutMs: number,
onError: (error: unknown) => void,
now: () => number = Date.now,
): TelemetryShutdownDeadline {
const expiresAtMs = now() + timeoutMs;
const remainingMs = (): number => Math.max(0, expiresAtMs - now());
return {
expiresAtMs,
remainingMs,
run: async (pipeline) => {
try {
await pipeline(remainingMs());
} catch (error) {
try {
onError(error);
} catch {
// A broken stderr/logger boundary must not block later cleanup.
}
}
},
};
}
15 changes: 14 additions & 1 deletion apps/kimi-code/src/cli/v2/run-v2-print.ts
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,10 @@ import {
raceWithTimeout,
requireConfiguredModel,
} from '../run-prompt';
import {
createTelemetryShutdownDeadline,
formatTelemetryShutdownWarning,
} from '../telemetry-shutdown';
import { createKimiCodeHostIdentity } from '../version';

import { resolveOutputFormat } from '../options';
Expand Down Expand Up @@ -171,8 +175,17 @@ export async function runV2Print(
try {
await restorePermission();
} finally {
const telemetryDeadline = createTelemetryShutdownDeadline(
CLI_SHUTDOWN_TIMEOUT_MS,
(error) => {
stderr.write(formatTelemetryShutdownWarning(error));
},
);
if (telemetryService !== undefined) {
await raceWithTimeout(telemetryService.shutdown(), CLI_SHUTDOWN_TIMEOUT_MS);
const service = telemetryService;
await telemetryDeadline.run((remainingMs) =>
raceWithTimeout(service.shutdown(), remainingMs),
);
}
app.dispose();
}
Expand Down
17 changes: 17 additions & 0 deletions apps/kimi-code/test/cli/run-prompt.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -325,6 +325,23 @@ describe('runPrompt', () => {
).rejects.toThrow('close failed');
});

it('reports telemetry shutdown failure without changing a successful result', async () => {
const stdout = writer();
const stderr = writer();
mocks.shutdownTelemetry.mockRejectedValueOnce(
new Error('telemetry unavailable'),
);

await expect(
runPrompt(opts(), '1.2.3-test', { stdout, stderr, process: fakeProcess() }),
).resolves.toBeUndefined();

expect(stderr.text()).toContain(
'Warning: telemetry shutdown failed: telemetry unavailable',
);
expect(mocks.harnessClose).toHaveBeenCalledOnce();
});

it('ignores a cleanup rejection that lands after the timeout', async () => {
vi.useFakeTimers();
try {
Expand Down
87 changes: 87 additions & 0 deletions apps/kimi-code/test/cli/telemetry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -117,3 +117,90 @@ describe('initializeServerTelemetry', () => {
);
});
});

describe('telemetry shutdown deadline', () => {
it('gives each pipeline only the budget remaining from one deadline', async () => {
const { createTelemetryShutdownDeadline } = await import(
'#/cli/telemetry-shutdown'
);
let nowMs = 1_000;
const observedBudgets: number[] = [];
const deadline = createTelemetryShutdownDeadline(
3_000,
vi.fn(),
() => nowMs,
);

await deadline.run(async (remainingMs) => {
observedBudgets.push(remainingMs);
nowMs = 2_250;
});
await deadline.run(async (remainingMs) => {
observedBudgets.push(remainingMs);
});

expect(observedBudgets).toEqual([3_000, 1_750]);
});

it('still invokes a later pipeline with no budget after the deadline expires', async () => {
const { createTelemetryShutdownDeadline } = await import(
'#/cli/telemetry-shutdown'
);
let nowMs = 1_000;
const pipeline = vi.fn(async (_remainingMs: number) => {});
const deadline = createTelemetryShutdownDeadline(
3_000,
vi.fn(),
() => nowMs,
);
nowMs = 4_001;

await deadline.run(pipeline);

expect(pipeline).toHaveBeenCalledWith(0);
});

it('continues to a later pipeline after an earlier one rejects', async () => {
const { createTelemetryShutdownDeadline } = await import(
'#/cli/telemetry-shutdown'
);
const laterPipeline = vi.fn(async (_remainingMs: number) => {});
const reportError = vi.fn();
const deadline = createTelemetryShutdownDeadline(
3_000,
reportError,
() => 1_000,
);

await expect(
deadline.run(async () => {
throw new Error('telemetry unavailable');
}),
).resolves.toBeUndefined();
await deadline.run(laterPipeline);

expect(laterPipeline).toHaveBeenCalledOnce();
expect(reportError).toHaveBeenCalledWith(
expect.objectContaining({ message: 'telemetry unavailable' }),
);
});

it('continues to a later pipeline when the error reporter also throws', async () => {
const { createTelemetryShutdownDeadline } = await import(
'#/cli/telemetry-shutdown'
);
const laterPipeline = vi.fn(async (_remainingMs: number) => {});
const deadline = createTelemetryShutdownDeadline(3_000, () => {
throw new Error('stderr unavailable');
});

await expect(
deadline.run(async () => {
throw new Error('telemetry unavailable');
}),
).resolves.toBeUndefined();
await deadline.run(laterPipeline);

expect(laterPipeline).toHaveBeenCalledOnce();
});
});
22 changes: 22 additions & 0 deletions apps/kimi-code/test/cli/v2-run-print.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -267,6 +267,28 @@ describe('runV2Print', () => {
expect(app.dispose).toHaveBeenCalled();
});

it('preserves a successful print result when telemetry shutdown rejects', async () => {
const stdout = writer();
const stderr = writer();
const { app, agent, appServices } = makeFakeHarness();
const telemetry = appServices.get(ITelemetryService) as {
shutdown: ReturnType<typeof vi.fn>;
};
telemetry.shutdown.mockRejectedValue(new Error('telemetry unavailable'));
mocks.bootstrap.mockReturnValue({ app });
mocks.ensureMainAgent.mockResolvedValue(agent);

await expect(
runV2Print(opts() as never, '1.2.3-test', { stdout, stderr }),
).resolves.toBeUndefined();

expect(telemetry.shutdown).toHaveBeenCalledOnce();
expect(app.dispose).toHaveBeenCalledOnce();
expect(stderr.text()).toContain(
'Warning: telemetry shutdown failed: telemetry unavailable',
);
});

it('seeds explicit skill dirs from --skillsDir into bootstrap', async () => {
const stdout = writer();
const stderr = writer();
Expand Down
Loading