diff --git a/.changeset/calm-ravens-share.md b/.changeset/calm-ravens-share.md new file mode 100644 index 0000000000..a07ebd72c6 --- /dev/null +++ b/.changeset/calm-ravens-share.md @@ -0,0 +1,5 @@ +--- +"@moonshot-ai/kimi-code": patch +--- + +Share one telemetry shutdown deadline across CLI-owned pipelines. diff --git a/apps/kimi-code/src/cli/run-prompt.ts b/apps/kimi-code/src/cli/run-prompt.ts index abee299622..43842ae7c3 100644 --- a/apps/kimi-code/src/cli/run-prompt.ts +++ b/apps/kimi-code/src/cli/run-prompt.ts @@ -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'; /** @@ -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(); } })()); diff --git a/apps/kimi-code/src/cli/sub/web/run.ts b/apps/kimi-code/src/cli/sub/web/run.ts index 9318b85395..2328b50cd7 100644 --- a/apps/kimi-code/src/cli/sub/web/run.ts +++ b/apps/kimi-code/src/cli/sub/web/run.ts @@ -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'; @@ -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, @@ -58,7 +63,7 @@ const WEB_ASSETS_DIR = 'dist-web'; interface RoutedServer { readonly address: string; readonly logger: ServerLogger; - close(): Promise; + close(options?: { readonly telemetryDeadlineMs?: number }): Promise; } export interface WebCliOptions extends ServerCliOptions { @@ -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); } @@ -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 }); diff --git a/apps/kimi-code/src/cli/telemetry-shutdown.ts b/apps/kimi-code/src/cli/telemetry-shutdown.ts new file mode 100644 index 0000000000..9c5083aaae --- /dev/null +++ b/apps/kimi-code/src/cli/telemetry-shutdown.ts @@ -0,0 +1,41 @@ +export interface TelemetryShutdownDeadline { + readonly expiresAtMs: number; + remainingMs(): number; + run(pipeline: (remainingMs: number) => Promise): Promise; +} + +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. + } + } + }, + }; +} diff --git a/apps/kimi-code/src/cli/v2/run-v2-print.ts b/apps/kimi-code/src/cli/v2/run-v2-print.ts index ff0c96c590..2959ebda44 100644 --- a/apps/kimi-code/src/cli/v2/run-v2-print.ts +++ b/apps/kimi-code/src/cli/v2/run-v2-print.ts @@ -80,6 +80,10 @@ import { raceWithTimeout, requireConfiguredModel, } from '../run-prompt'; +import { + createTelemetryShutdownDeadline, + formatTelemetryShutdownWarning, +} from '../telemetry-shutdown'; import { createKimiCodeHostIdentity } from '../version'; import { resolveOutputFormat } from '../options'; @@ -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(); } diff --git a/apps/kimi-code/test/cli/run-prompt.test.ts b/apps/kimi-code/test/cli/run-prompt.test.ts index 4bad127d89..d9679735a7 100644 --- a/apps/kimi-code/test/cli/run-prompt.test.ts +++ b/apps/kimi-code/test/cli/run-prompt.test.ts @@ -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 { diff --git a/apps/kimi-code/test/cli/telemetry.test.ts b/apps/kimi-code/test/cli/telemetry.test.ts index 490aa83ec1..b6654a92f3 100644 --- a/apps/kimi-code/test/cli/telemetry.test.ts +++ b/apps/kimi-code/test/cli/telemetry.test.ts @@ -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(); + }); +}); diff --git a/apps/kimi-code/test/cli/v2-run-print.test.ts b/apps/kimi-code/test/cli/v2-run-print.test.ts index e6a5e5ce3e..6e1883b60d 100644 --- a/apps/kimi-code/test/cli/v2-run-print.test.ts +++ b/apps/kimi-code/test/cli/v2-run-print.test.ts @@ -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; + }; + 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(); diff --git a/apps/kimi-code/test/cli/web/web.test.ts b/apps/kimi-code/test/cli/web/web.test.ts index 1b51bfc535..ce06c45b1f 100644 --- a/apps/kimi-code/test/cli/web/web.test.ts +++ b/apps/kimi-code/test/cli/web/web.test.ts @@ -515,6 +515,93 @@ describe('`kimi web` option threading', () => { }); }); +describe('`kimi web` foreground shutdown', () => { + it('passes no fresh budget to host telemetry after engine shutdown expires', async () => { + vi.resetModules(); + const startServer = vi.fn(); + const shutdownTelemetry = vi.fn(async () => {}); + const logger = { info: vi.fn(), error: vi.fn() }; + vi.doMock('@moonshot-ai/kap-server', async () => { + const actual = await vi.importActual( + '@moonshot-ai/kap-server', + ); + return { + ...actual, + createServerLogger: vi.fn(() => logger), + startServer, + }; + }); + vi.doMock('@moonshot-ai/kimi-telemetry', async () => { + const actual = await vi.importActual( + '@moonshot-ai/kimi-telemetry', + ); + return { ...actual, shutdownTelemetry, track: vi.fn() }; + }); + vi.doMock('#/cli/telemetry', () => ({ initializeServerTelemetry: vi.fn() })); + + let nowMs = 10_000; + const serverClose = vi.fn(async () => { + nowMs = 13_001; + throw new Error('engine telemetry unavailable'); + }); + startServer.mockResolvedValue({ + host: '127.0.0.1', + port: 58_627, + close: serverClose, + }); + + const previousListeners = new Set(process.listeners('SIGINT')); + const exit = vi.spyOn(process, 'exit').mockImplementation((() => undefined) as never); + const now = vi.spyOn(Date, 'now').mockImplementation(() => nowMs); + let signalHandler: (() => void) | undefined; + try { + const { startServerForeground } = await import('#/cli/sub/web/run'); + let markReady: (() => void) | undefined; + const ready = new Promise((resolve) => { + markReady = resolve; + }); + void startServerForeground( + { + host: '127.0.0.1', + port: 0, + logLevel: 'silent', + debugEndpoints: false, + insecureNoTls: true, + allowRemoteShutdown: false, + allowRemoteTerminals: false, + dangerousBypassAuth: false, + allowedHosts: [], + }, + { onReady: () => markReady?.() }, + ); + await ready; + + signalHandler = process + .listeners('SIGINT') + .find((listener) => !previousListeners.has(listener)) as (() => void) | undefined; + signalHandler?.(); + + await vi.waitFor(() => expect(exit).toHaveBeenCalledWith(0)); + expect(serverClose).toHaveBeenCalledWith({ telemetryDeadlineMs: 13_000 }); + expect(shutdownTelemetry).toHaveBeenCalledWith({ timeoutMs: 0 }); + expect(logger.error).toHaveBeenCalledWith( + expect.objectContaining({ + err: expect.objectContaining({ message: 'engine telemetry unavailable' }), + }), + 'telemetry shutdown error', + ); + } finally { + if (signalHandler !== undefined) process.removeListener('SIGINT', signalHandler); + now.mockRestore(); + exit.mockRestore(); + vi.doUnmock('@moonshot-ai/kap-server'); + vi.doUnmock('@moonshot-ai/kimi-telemetry'); + vi.doUnmock('#/cli/telemetry'); + vi.resetModules(); + } + }); +}); + describe('shared parsers stay strict', () => { it('rejects out-of-range --port', async () => { const { parsePort } = await import('#/cli/sub/web/shared'); diff --git a/packages/kap-server/src/start.ts b/packages/kap-server/src/start.ts index 79610a3d95..1ee0c374b8 100644 --- a/packages/kap-server/src/start.ts +++ b/packages/kap-server/src/start.ts @@ -162,7 +162,12 @@ export interface RunningServer { readonly authTokenService: IAuthTokenService; readonly host: string; readonly port: number; - close(): Promise; + close(options?: ServerCloseOptions): Promise; +} + +export interface ServerCloseOptions { + /** Absolute Unix timestamp shared with other host-owned telemetry pipelines. */ + readonly telemetryDeadlineMs?: number; } const DEFAULT_HOST = '127.0.0.1'; @@ -335,13 +340,13 @@ export async function startServer(opts: ServerStartOptions = {}): Promise => { + const close = async (options: ServerCloseOptions = {}): Promise => { await app.close(); authFailureLimiter?.dispose(); modelCatalogRefreshScheduler.dispose(); // Telemetry is best-effort and must never prevent core or instance cleanup. try { - await shutdownServerTelemetry(telemetry); + await shutdownServerTelemetry(telemetry, options.telemetryDeadlineMs); } catch (error) { logger.warn( { err: error instanceof Error ? error.message : String(error) }, diff --git a/packages/kap-server/test/boot.test.ts b/packages/kap-server/test/boot.test.ts index e68800fe12..34ed22183c 100644 --- a/packages/kap-server/test/boot.test.ts +++ b/packages/kap-server/test/boot.test.ts @@ -29,6 +29,8 @@ import { listenWithPortRetry, type RunningServer, startServer } from '../src/sta import { getServerVersion } from '../src/version'; import { authedFetch } from './helpers/auth'; +const EXPIRED_TELEMETRY_DEADLINE_MAX_CLOSE_MS = 1_500; + describe('server-v2 boot', () => { let server: RunningServer | undefined; let home: string | undefined; @@ -230,6 +232,33 @@ describe('server-v2 boot', () => { expect(() => core.accessor.get(IBootstrapService)).toThrow(); expect(await listLiveServerInstances(home)).toEqual([]); }); + + it('passes an expired host telemetry deadline to owned engine telemetry', async () => { + home = await mkdtemp(join(tmpdir(), 'kimi-server-v2-telemetry-deadline-')); + const auth = { + _serviceBrand: undefined, + getCachedAccessToken: () => new Promise(() => {}), + } as unknown as IOAuthToolkit; + + server = await startServer({ + host: '127.0.0.1', + port: 0, + homeDir: home, + logLevel: 'silent', + telemetry: true, + seeds: [[IOAuthToolkit, auth]], + }); + server.core.accessor.get(ITelemetryService).track('server_probe'); + + const closeStartedAt = Date.now(); + await server.close({ telemetryDeadlineMs: Date.now() }); + server = undefined; + + expect(Date.now() - closeStartedAt).toBeLessThan( + EXPIRED_TELEMETRY_DEADLINE_MAX_CLOSE_MS, + ); + expect(await listLiveServerInstances(home)).toEqual([]); + }); }); function silentLogger() {