From 2388e1bd7d39a9278045e18918c671518370875a Mon Sep 17 00:00:00 2001 From: 7Sageer <7sageer@djwcb.cn> Date: Mon, 27 Jul 2026 17:41:13 +0800 Subject: [PATCH 1/7] fix(telemetry): make appender shutdown durable Own queued and in-flight flushes through shutdown cancellation, replay recoverable spool data when appenders start, and forward one lifecycle budget through the telemetry facade. Refs #2246 --- .../src/app/telemetry/cloudAppender.ts | 113 ++++++- .../src/app/telemetry/cloudTransport.ts | 76 +++-- .../src/app/telemetry/telemetry.ts | 10 +- .../src/app/telemetry/telemetryService.ts | 19 +- .../test/app/telemetry/cloudAppender.test.ts | 309 +++++++++++++++++- .../app/telemetry/telemetryService.test.ts | 71 +++- packages/kap-server/src/services/telemetry.ts | 9 +- packages/kap-server/test/telemetry.test.ts | 6 +- 8 files changed, 539 insertions(+), 74 deletions(-) diff --git a/packages/agent-core-v2/src/app/telemetry/cloudAppender.ts b/packages/agent-core-v2/src/app/telemetry/cloudAppender.ts index fd0cb5655f..3beb31ecd9 100644 --- a/packages/agent-core-v2/src/app/telemetry/cloudAppender.ts +++ b/packages/agent-core-v2/src/app/telemetry/cloudAppender.ts @@ -5,7 +5,8 @@ * telemetry endpoint through `CloudTransport`, which persists failed events * through the `storage` byte layer. Reads host facts (`clientVersion`, env, * platform/arch) from `IBootstrapService`; `createCloudAppender` assembles - * one from a `ServicesAccessor` so hosts only supply identity facts. + * one from a `ServicesAccessor` so hosts only supply identity facts. Owns + * periodic flush, startup spool replay, and deadline-aware durable shutdown. * App-scoped; independent of `@moonshot-ai/kimi-telemetry`. */ @@ -14,10 +15,16 @@ import { release } from 'node:os'; import type { ServicesAccessor } from '#/_base/di/instantiation'; import { onUnexpectedError } from '#/_base/errors/unexpectedError'; +import { abortError } from '#/_base/utils/abort'; import { IBootstrapService } from '#/app/bootstrap/bootstrap'; import { IFileSystemStorageService } from '#/persistence/interface/storage'; -import type { ITelemetryAppender, TelemetryContextPatch, TelemetryProperties } from './telemetry'; +import type { + ITelemetryAppender, + TelemetryContextPatch, + TelemetryProperties, + TelemetryShutdownOptions, +} from './telemetry'; import { type CloudContext, type CloudPrimitive, @@ -84,6 +91,12 @@ export class CloudAppender implements ITelemetryAppender { private sessionId: string | null; private buffer: EnrichedCloudEvent[] = []; private flushTimer: ReturnType | null = null; + private readonly lifecycleController = new AbortController(); + private acceptingEvents = true; + private started = false; + private startupReplay: Promise | null = null; + private flushPromise: Promise | null = null; + private shutdownPromise: Promise | null = null; constructor(options: CloudAppenderOptions) { this.deviceId = options.deviceId; @@ -105,6 +118,7 @@ export class CloudAppender implements ITelemetryAppender { } track(event: string, properties?: TelemetryProperties): void { + if (!this.acceptingEvents) return; const eventSessionId = properties?.['sessionId']; const enriched: EnrichedCloudEvent = { event_id: randomUUID().replaceAll('-', ''), @@ -136,20 +150,45 @@ export class CloudAppender implements ITelemetryAppender { } } - async flush(): Promise { - if (this.buffer.length === 0) return; - const events = this.buffer; - this.buffer = []; - await this.transport.send(events); + flush(): Promise { + if (this.flushPromise !== null) return this.flushPromise; + if (this.buffer.length === 0) return Promise.resolve(); + const flush = this.drainBuffer(); + this.flushPromise = flush; + void flush.then( + () => { + this.clearFlush(flush); + }, + () => { + this.clearFlush(flush); + }, + ); + return flush; } - async shutdown(): Promise { - this.stopPeriodicFlush(); - await this.flush(); + shutdown(options: TelemetryShutdownOptions = {}): Promise { + const clearDeadline = this.armShutdownDeadline(options); + if (this.shutdownPromise === null) { + this.acceptingEvents = false; + this.stopPeriodicFlush(); + this.shutdownPromise = this.shutdownOwnedWork(); + } + const shutdown = this.shutdownPromise; + void shutdown.then(clearDeadline, clearDeadline); + return shutdown; + } + + start(): void { + if (this.started || !this.acceptingEvents) return; + this.started = true; + this.startupReplay = this.transport + .retryDiskEvents(this.lifecycleController.signal) + .catch(() => {}); + this.startPeriodicFlush(); } startPeriodicFlush(): void { - if (this.flushTimer !== null) return; + if (!this.acceptingEvents || this.flushTimer !== null) return; this.flushTimer = setInterval(() => { void this.flush().catch(() => {}); }, this.flushIntervalMs); @@ -163,7 +202,57 @@ export class CloudAppender implements ITelemetryAppender { } async retryDiskEvents(): Promise { - await this.transport.retryDiskEvents(); + await this.transport.retryDiskEvents(this.lifecycleController.signal); + } + + private async drainBuffer(): Promise { + await (this.startupReplay ?? Promise.resolve()); + while (this.buffer.length > 0) { + const events = this.buffer; + this.buffer = []; + try { + await this.transport.send(events, this.lifecycleController.signal); + } catch (error) { + this.buffer = [...events, ...this.buffer]; + throw error; + } + } + } + + private clearFlush(flush: Promise): void { + if (this.flushPromise === flush) this.flushPromise = null; + } + + private async shutdownOwnedWork(): Promise { + await (this.startupReplay ?? Promise.resolve()); + await this.flush(); + } + + private armShutdownDeadline(options: TelemetryShutdownOptions): () => void { + const abort = (): void => { + if (!this.lifecycleController.signal.aborted) { + this.lifecycleController.abort(abortError('Telemetry shutdown deadline reached')); + } + }; + const signal = options.signal; + if (signal?.aborted === true) { + abort(); + } else { + signal?.addEventListener('abort', abort, { once: true }); + } + let timeout: ReturnType | undefined; + if (options.deadlineMs !== undefined) { + const remainingMs = options.deadlineMs - Date.now(); + if (remainingMs <= 0) { + abort(); + } else { + timeout = setTimeout(abort, remainingMs); + } + } + return () => { + if (timeout !== undefined) clearTimeout(timeout); + signal?.removeEventListener('abort', abort); + }; } } diff --git a/packages/agent-core-v2/src/app/telemetry/cloudTransport.ts b/packages/agent-core-v2/src/app/telemetry/cloudTransport.ts index 0de6499475..7544f166a7 100644 --- a/packages/agent-core-v2/src/app/telemetry/cloudTransport.ts +++ b/packages/agent-core-v2/src/app/telemetry/cloudTransport.ts @@ -2,13 +2,14 @@ * `telemetry` domain (L1) — `CloudTransport`, the HTTP transport behind * `CloudAppender`. Posts enriched events to the telemetry endpoint with Bearer * auth, retry, and a byte-store fallback for failed events, persisted through - * the `storage` byte layer (`IFileSystemStorageService`) under the `telemetry` scope. + * the `storage` byte layer (`IFileSystemStorageService`) under an isolated + * `telemetry-v2` scope. * App-scoped; independent of `@moonshot-ai/kimi-telemetry`. */ import { randomBytes } from 'node:crypto'; -import { isAbortError } from '#/_base/utils/abort'; +import { abortable, isAbortError } from '#/_base/utils/abort'; import type { IFileSystemStorageService } from '#/persistence/interface/storage'; export type CloudPrimitive = boolean | number | string | undefined | null; @@ -54,7 +55,7 @@ export const DISK_EVENT_MAX_AGE_MS = 7 * 24 * 60 * 60 * 1000; export const RETRY_BACKOFFS_MS = [1_000, 4_000, 16_000] as const; const DEFAULT_REQUEST_TIMEOUT_MS = 10_000; -const TELEMETRY_SCOPE = 'telemetry'; +const TELEMETRY_SCOPE = 'telemetry-v2'; const FAILED_PREFIX = 'failed_'; const JSONL_SUFFIX = '.jsonl'; @@ -86,15 +87,9 @@ export class CloudTransport { async send(events: readonly EnrichedCloudEvent[], signal?: AbortSignal): Promise { if (events.length === 0) return; - let savedToDisk = false; - const saveEventsToDisk = async (): Promise => { - if (savedToDisk) return; - await this.saveToDisk(events); - savedToDisk = true; - }; if (signal?.aborted === true) { - await saveEventsToDisk(); - throw abortError(); + await this.saveToDisk(events); + return; } let payload: CloudPayload; @@ -104,32 +99,33 @@ export class CloudTransport { return; } - try { - for (let attempt = 0; attempt <= this.retryBackoffsMs.length; attempt++) { - try { - await this.sendHttp(payload, signal); + for (let attempt = 0; attempt <= this.retryBackoffsMs.length; attempt++) { + try { + const request = this.sendHttp(payload, signal); + await (signal === undefined ? request : abortable(request, signal)); + return; + } catch (error) { + if (isSignalAborted(signal) || isAbortError(error)) { + await this.saveToDisk(events); return; - } catch (error) { - if (isSignalAborted(signal) || isAbortError(error)) { - await saveEventsToDisk(); - throw error; - } - if (!(error instanceof TransientCloudError)) { - break; + } + if (!(error instanceof TransientCloudError)) break; + const backoff = this.retryBackoffsMs[attempt]; + if (backoff === undefined) break; + try { + const sleep = this.sleepImpl(backoff, signal); + await (signal === undefined ? sleep : abortable(sleep, signal)); + } catch (sleepError) { + if (isSignalAborted(signal) || isAbortError(sleepError)) { + await this.saveToDisk(events); + return; } - const backoff = this.retryBackoffsMs[attempt]; - if (backoff === undefined) break; - await this.sleepImpl(backoff, signal); + break; } } - } catch (error) { - if (isSignalAborted(signal) || isAbortError(error)) { - await saveEventsToDisk(); - throw error; - } } - await saveEventsToDisk(); + await this.saveToDisk(events); } async saveToDisk(events: readonly EnrichedCloudEvent[]): Promise { @@ -139,13 +135,15 @@ export class CloudTransport { await this.storage.write(TELEMETRY_SCOPE, key, textEncoder.encode(text)); } - async retryDiskEvents(): Promise { + async retryDiskEvents(signal?: AbortSignal): Promise { const keys = await this.storage.list(TELEMETRY_SCOPE, FAILED_PREFIX); const now = this.now(); for (const key of keys) { + if (signal?.aborted === true) throw abortError(); if (!key.startsWith(FAILED_PREFIX) || !key.endsWith(JSONL_SUFFIX)) continue; const createdAt = parseFailedTimestamp(key); - if (createdAt === undefined || now - createdAt > DISK_EVENT_MAX_AGE_MS) { + if (createdAt === undefined) continue; + if (now - createdAt > DISK_EVENT_MAX_AGE_MS) { await this.storage.delete(TELEMETRY_SCOPE, key).catch(() => undefined); continue; } @@ -163,9 +161,11 @@ export class CloudTransport { } try { - await this.sendHttp(payload); + const request = this.sendHttp(payload, signal); + await (signal === undefined ? request : abortable(request, signal)); await this.storage.delete(TELEMETRY_SCOPE, key); } catch (error) { + if (isSignalAborted(signal) || isAbortError(error)) throw error; if (error instanceof TransientCloudError) continue; } } @@ -342,12 +342,16 @@ async function fetchWithTimeout( function abortableSleep(ms: number, signal?: AbortSignal): Promise { if (signal?.aborted === true) return Promise.reject(abortError()); return new Promise((resolve, reject) => { - const timer = setTimeout(resolve, ms); - timer.unref?.(); + let timer: ReturnType; const onAbort = (): void => { clearTimeout(timer); reject(abortError()); }; + timer = setTimeout(() => { + signal?.removeEventListener('abort', onAbort); + resolve(); + }, ms); + timer.unref?.(); signal?.addEventListener('abort', onAbort, { once: true }); }); } diff --git a/packages/agent-core-v2/src/app/telemetry/telemetry.ts b/packages/agent-core-v2/src/app/telemetry/telemetry.ts index f9d61c49d8..74bf1b2067 100644 --- a/packages/agent-core-v2/src/app/telemetry/telemetry.ts +++ b/packages/agent-core-v2/src/app/telemetry/telemetry.ts @@ -24,12 +24,18 @@ export type TelemetryProperties = Readonly>; export type TelemetryContextPatch = TelemetryProperties; +export interface TelemetryShutdownOptions { + readonly signal?: AbortSignal; + readonly deadlineMs?: number; +} + export interface ITelemetryAppender { + start?(): void; track(event: string, properties?: TelemetryProperties): void; withContext?(patch: TelemetryContextPatch): ITelemetryAppender; setContext?(patch: TelemetryContextPatch): void; flush?(): Promise | void; - shutdown?(): Promise | void; + shutdown?(options?: TelemetryShutdownOptions): Promise | void; } export interface TelemetryServiceOptions { @@ -56,7 +62,7 @@ export interface ITelemetryService { setAppender(appender: ITelemetryAppender): void; setEnabled(enabled: boolean): void; flush(): Promise; - shutdown(): Promise; + shutdown(options?: TelemetryShutdownOptions): Promise; } export const nullTelemetryAppender: ITelemetryAppender = { diff --git a/packages/agent-core-v2/src/app/telemetry/telemetryService.ts b/packages/agent-core-v2/src/app/telemetry/telemetryService.ts index 6d0126527d..a860ef8a1c 100644 --- a/packages/agent-core-v2/src/app/telemetry/telemetryService.ts +++ b/packages/agent-core-v2/src/app/telemetry/telemetryService.ts @@ -22,6 +22,7 @@ import { nullTelemetryAppender, type TelemetryContextPatch, type TelemetryProperties, + type TelemetryShutdownOptions, } from './telemetry'; export class TelemetryService implements ITelemetryService { @@ -64,6 +65,7 @@ export class TelemetryService implements ITelemetryService { } addAppender(appender: ITelemetryAppender): IDisposable { + this.startAppender(appender); this.appenders.push(appender); return toDisposable(() => this.removeAppender(appender)); } @@ -73,6 +75,7 @@ export class TelemetryService implements ITelemetryService { } setAppender(appender: ITelemetryAppender): void { + this.startAppender(appender); this.appenders = [appender]; } @@ -88,13 +91,21 @@ export class TelemetryService implements ITelemetryService { ); } - async shutdown(): Promise { + async shutdown(options?: TelemetryShutdownOptions): Promise { await Promise.all( this.appenders.map((appender) => - Promise.resolve(appender.shutdown?.()).catch(onUnexpectedError), + Promise.resolve(appender.shutdown?.(options)).catch(onUnexpectedError), ), ); } + + private startAppender(appender: ITelemetryAppender): void { + try { + appender.start?.(); + } catch (error) { + onUnexpectedError(error); + } + } } class TelemetryContextView implements ITelemetryService { @@ -147,8 +158,8 @@ class TelemetryContextView implements ITelemetryService { return this.root.flush(); } - shutdown(): Promise { - return this.root.shutdown(); + shutdown(options?: TelemetryShutdownOptions): Promise { + return this.root.shutdown(options); } } diff --git a/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts b/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts index 3c3abf3b3b..5c543bfd4f 100644 --- a/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts +++ b/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts @@ -1,4 +1,20 @@ -import { mkdtempSync, readdirSync, rmSync } from 'node:fs'; +/** + * Cloud telemetry lifecycle tests — exercise the real appender, transport, + * and file-storage stack while stubbing only the outbound HTTP boundary. + * Covers batching, durable shutdown, startup replay, privacy, and wire shape. + * Run with `pnpm --filter @moonshot-ai/agent-core-v2 exec vitest run + * test/app/telemetry/cloudAppender.test.ts`. + */ + +import { getEventListeners } from 'node:events'; +import { + mkdtempSync, + mkdirSync, + readFileSync, + readdirSync, + rmSync, + writeFileSync, +} from 'node:fs'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; @@ -10,12 +26,14 @@ import { } from '#/_base/errors/unexpectedError'; import { FileStorageService } from '#/persistence/backends/node-fs/fileStorageService'; import { CloudAppender, type CloudAppenderOptions } from '#/app/telemetry/cloudAppender'; +import { CloudTransport } from '#/app/telemetry/cloudTransport'; import { stubBootstrap } from '../bootstrap/stubs'; interface CapturedRequest { readonly url: string; readonly headers: Record; + readonly signal?: AbortSignal; readonly body: { readonly user_id: string; readonly events: readonly Record[]; @@ -26,10 +44,15 @@ type Responder = (req: CapturedRequest) => Response | Promise; function makeFetch(responder: Responder): typeof fetch { return (async (input: unknown, init: unknown) => { - const requestInit = init as { headers: Record; body: string }; + const requestInit = init as { + headers: Record; + body: string; + signal?: AbortSignal; + }; const req: CapturedRequest = { url: String(input), headers: requestInit.headers, + signal: requestInit.signal, body: JSON.parse(requestInit.body) as CapturedRequest['body'], }; return responder(req); @@ -44,6 +67,17 @@ function statusResponse(status: number): Response { return new Response(null, { status }); } +function deferred(): { + readonly promise: Promise; + readonly resolve: (value: T) => void; +} { + let resolve!: (value: T) => void; + const promise = new Promise((done) => { + resolve = done; + }); + return { promise, resolve }; +} + function baseOptions( overrides: Partial & { homeDir?: string } = {}, ): CloudAppenderOptions { @@ -58,6 +92,18 @@ function baseOptions( }; } +function listFailedSpoolFiles(homeDir: string): string[] { + return readdirSync(join(homeDir, 'telemetry-v2')).filter((file) => + file.startsWith('failed_'), + ); +} + +function readFirstFailedEvent(homeDir: string): Record { + const file = listFailedSpoolFiles(homeDir)[0] as string; + const persisted = readFileSync(join(homeDir, 'telemetry-v2', file), 'utf8'); + return JSON.parse(persisted.trim()) as Record; +} + describe('CloudAppender', () => { let homeDir: string; @@ -202,6 +248,171 @@ describe('CloudAppender', () => { expect(sends).toBe(1); }); + it('shutdown returns the original lifecycle promise when called repeatedly', async () => { + const appender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => okResponse()), + }), + ); + appender.track('once'); + + const first = appender.shutdown(); + const second = appender.shutdown(); + + expect(second).toBe(first); + await first; + }); + + it('a later shutdown cancellation tightens an active lifecycle', async () => { + const requestStarted = deferred(); + const appender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => { + requestStarted.resolve(); + return new Promise(() => {}); + }), + }), + ); + const cancellation = new AbortController(); + appender.track('cancelled_by_later_caller'); + + const first = appender.shutdown(); + await requestStarted.promise; + const second = appender.shutdown({ signal: cancellation.signal }); + cancellation.abort(); + await second; + + expect(second).toBe(first); + const files = listFailedSpoolFiles(homeDir); + expect(files).toHaveLength(1); + }); + + it('track drops events after shutdown begins', async () => { + let requests = 0; + const appender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => { + requests += 1; + return okResponse(); + }), + }), + ); + + await appender.shutdown(); + appender.track('too_late'); + await appender.flush(); + + expect(requests).toBe(0); + }); + + it('flush drains events tracked while an earlier batch is in flight', async () => { + const firstResponse = deferred(); + const firstRequestStarted = deferred(); + const batches: string[][] = []; + const appender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch((request) => { + batches.push(request.body.events.map((event) => String(event['event']))); + if (batches.length === 1) { + firstRequestStarted.resolve(); + return firstResponse.promise; + } + return okResponse(); + }), + }), + ); + + appender.track('first'); + const flushing = appender.flush(); + await firstRequestStarted.promise; + appender.track('second'); + firstResponse.resolve(okResponse()); + await flushing; + + expect(batches).toEqual([['kfc_first'], ['kfc_second']]); + }); + + it('concurrent flush calls share ownership of one batch', async () => { + const response = deferred(); + const requestStarted = deferred(); + let requests = 0; + const appender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => { + requests += 1; + requestStarted.resolve(); + return response.promise; + }), + }), + ); + appender.track('once'); + + const first = appender.flush(); + await requestStarted.promise; + const second = appender.flush(); + response.resolve(okResponse()); + await Promise.all([first, second]); + + expect(second).toBe(first); + expect(requests).toBe(1); + }); + + it('shutdown persists a threshold batch already in flight when cancellation fires', async () => { + const requestStarted = deferred(); + let requestSignal: AbortSignal | undefined; + const appender = new CloudAppender( + baseOptions({ + homeDir, + flushThreshold: 1, + fetchImpl: makeFetch((request) => { + requestSignal = request.signal; + requestStarted.resolve(); + return new Promise(() => {}); + }), + }), + ); + const cancellation = new AbortController(); + + appender.track('threshold_in_flight'); + await requestStarted.promise; + const closing = appender.shutdown({ signal: cancellation.signal }); + cancellation.abort(); + await closing; + + expect(requestSignal?.aborted).toBe(true); + const files = listFailedSpoolFiles(homeDir); + expect(files).toHaveLength(1); + expect(readFirstFailedEvent(homeDir)).toMatchObject({ + event: 'threshold_in_flight', + }); + }); + + it('shutdown persists buffered events when its absolute deadline has elapsed', async () => { + let requests = 0; + const appender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => { + requests += 1; + return okResponse(); + }), + }), + ); + + appender.track('deadline_elapsed'); + await appender.shutdown({ deadlineMs: Date.now() - 1 }); + + expect(requests).toBe(0); + const files = listFailedSpoolFiles(homeDir); + expect(files).toHaveLength(1); + expect(readFirstFailedEvent(homeDir)).toMatchObject({ event: 'deadline_elapsed' }); + }); + it('retries on 5xx and saves to disk after exhausting backoffs', async () => { let attempts = 0; const appender = new CloudAppender( @@ -218,10 +429,42 @@ describe('CloudAppender', () => { await appender.flush(); expect(attempts).toBe(4); - const files = readdirSync(join(homeDir, 'telemetry')).filter((f) => f.startsWith('failed_')); + const files = listFailedSpoolFiles(homeDir); expect(files).toHaveLength(1); }); + it('releases lifecycle abort listeners after a retry backoff completes', async () => { + let attempts = 0; + const transport = new CloudTransport({ + storage: new FileStorageService(homeDir), + deviceId: 'dev', + retryBackoffsMs: [0], + fetchImpl: makeFetch(() => { + attempts += 1; + return attempts === 1 ? statusResponse(500) : okResponse(); + }), + }); + const lifecycle = new AbortController(); + + await transport.send( + [ + { + event_id: 'event-1', + device_id: 'dev', + session_id: null, + event: 'retry_listener_cleanup', + timestamp: 1, + properties: {}, + context: {}, + }, + ], + lifecycle.signal, + ); + + expect(attempts).toBe(2); + expect(getEventListeners(lifecycle.signal, 'abort')).toHaveLength(0); + }); + it('retries a 401 once without the Authorization header', async () => { const seenAuths: (string | undefined)[] = []; const appender = new CloudAppender( @@ -255,15 +498,63 @@ describe('CloudAppender', () => { appender.track('evt'); await appender.flush(); - expect( - readdirSync(join(homeDir, 'telemetry')).filter((f) => f.startsWith('failed_')), - ).toHaveLength(1); + expect(listFailedSpoolFiles(homeDir)).toHaveLength(1); shouldFail = false; await appender.retryDiskEvents(); - expect( - readdirSync(join(homeDir, 'telemetry')).filter((f) => f.startsWith('failed_')), - ).toHaveLength(0); + expect(listFailedSpoolFiles(homeDir)).toHaveLength(0); + }); + + it('start replays recoverable events left by an earlier appender', async () => { + const failingAppender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => statusResponse(500)), + }), + ); + failingAppender.track('persisted_before_restart'); + await failingAppender.flush(); + expect(listFailedSpoolFiles(homeDir)).toHaveLength(1); + + let replayed = 0; + const restartedAppender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => { + replayed += 1; + return okResponse(); + }), + }), + ); + + restartedAppender.start(); + await restartedAppender.shutdown(); + + expect(replayed).toBe(1); + expect(listFailedSpoolFiles(homeDir)).toHaveLength(0); + }); + + it('start leaves legacy telemetry spool files for their owning pipeline', async () => { + const telemetryDir = join(homeDir, 'telemetry'); + mkdirSync(telemetryDir, { recursive: true }); + const legacyFile = 'failed_abcdef123456.jsonl'; + writeFileSync(join(telemetryDir, legacyFile), '{"event":"legacy"}\n'); + let requests = 0; + const appender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => { + requests += 1; + return okResponse(); + }), + }), + ); + + appender.start(); + await appender.shutdown(); + + expect(requests).toBe(0); + expect(readdirSync(telemetryDir)).toContain(legacyFile); }); it('drops non-primitive properties and reports the violation', async () => { diff --git a/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts b/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts index 86528094b3..3a153d5557 100644 --- a/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts +++ b/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts @@ -1,3 +1,9 @@ +/** + * Telemetry facade tests — exercise appender fan-out, context views, error + * isolation, lifecycle-option forwarding, and App-scope registration through + * the public `ITelemetryService` surface. + */ + import { afterEach, beforeEach, describe, expect, it } from 'vitest'; import { LifecycleScope, ScopeActivation, _clearScopedRegistryForTests, registerScopedService } from '#/_base/di/scope'; @@ -6,21 +12,32 @@ import { resetUnexpectedErrorHandler, setUnexpectedErrorHandler, } from '#/_base/errors/unexpectedError'; -import { type ITelemetryAppender, type TelemetryProperties, ITelemetryService } from '#/app/telemetry/telemetry'; +import { + type ITelemetryAppender, + type TelemetryProperties, + type TelemetryShutdownOptions, + ITelemetryService, +} from '#/app/telemetry/telemetry'; import { TelemetryService } from '#/app/telemetry/telemetryService'; class CapturingAppender implements ITelemetryAppender { readonly events: { event: string; properties?: TelemetryProperties }[] = []; + startCalls = 0; flushCalls = 0; shutdownCalls = 0; + shutdownOptions: TelemetryShutdownOptions | undefined; + start(): void { + this.startCalls += 1; + } track(event: string, properties?: TelemetryProperties): void { this.events.push({ event, properties }); } flush(): void { this.flushCalls += 1; } - shutdown(): void { + shutdown(options?: TelemetryShutdownOptions): void { this.shutdownCalls += 1; + this.shutdownOptions = options; } } @@ -101,6 +118,15 @@ describe('TelemetryService (unit)', () => { expect(b.events).toHaveLength(1); }); + it('addAppender starts the registered appender', () => { + const appender = new CapturingAppender(); + const svc = new TelemetryService(); + + svc.addAppender(appender); + + expect(appender.startCalls).toBe(1); + }); + it('removeAppender stops delivery to that appender', () => { const a = new CapturingAppender(); const b = new CapturingAppender(); @@ -111,6 +137,15 @@ describe('TelemetryService (unit)', () => { expect(b.events).toHaveLength(1); }); + it('setAppender starts the replacement appender', () => { + const appender = new CapturingAppender(); + const svc = new TelemetryService(); + + svc.setAppender(appender); + + expect(appender.startCalls).toBe(1); + }); + it('setEnabled(false) drops track; setEnabled(true) resumes', () => { const appender = new CapturingAppender(); const svc = telemetryWithAppenders(appender); @@ -165,6 +200,21 @@ describe('TelemetryService (unit)', () => { expect(b.shutdownCalls).toBe(1); }); + it('shutdown forwards one lifecycle budget to every appender', async () => { + const first = new CapturingAppender(); + const second = new CapturingAppender(); + const svc = telemetryWithAppenders(first, second); + const options = { + signal: new AbortController().signal, + deadlineMs: Date.now() + 25, + }; + + await svc.shutdown(options); + + expect(first.shutdownOptions).toBe(options); + expect(second.shutdownOptions).toBe(options); + }); + it('flush is a no-op for appenders without flush', async () => { const minimal: ITelemetryAppender = { track() {} }; const svc = telemetryWithAppenders(minimal); @@ -189,6 +239,23 @@ describe('TelemetryService (error isolation)', () => { expect(good.events).toEqual([{ event: 'evt', properties: {} }]); }); + it('a throwing appender start does not prevent other appenders from registering', () => { + const bad: ITelemetryAppender = { + start() { + throw new Error('boom'); + }, + track() {}, + }; + const good = new CapturingAppender(); + const svc = new TelemetryService(); + + svc.addAppender(bad); + svc.addAppender(good); + svc.track('evt'); + + expect(good.events).toEqual([{ event: 'evt', properties: {} }]); + }); + it('flush tolerates a rejecting appender and still flushes the rest', async () => { const bad: ITelemetryAppender = { track() {}, diff --git a/packages/kap-server/src/services/telemetry.ts b/packages/kap-server/src/services/telemetry.ts index 25efe53505..85efc3899d 100644 --- a/packages/kap-server/src/services/telemetry.ts +++ b/packages/kap-server/src/services/telemetry.ts @@ -63,13 +63,6 @@ export async function initializeServerTelemetry( getAccessToken: async () => (await auth.getCachedAccessToken()) ?? null, }); const registration = service.addAppender(appender); - try { - // The server is long-lived: flush on a timer, not only at the threshold. - appender.startPeriodicFlush(); - } catch (error) { - registration.dispose(); - throw error; - } return { appender, registration }; } @@ -82,7 +75,7 @@ export async function shutdownServerTelemetry( let timer: ReturnType | undefined; try { await Promise.race([ - telemetry.appender.shutdown(), + telemetry.appender.shutdown({ deadlineMs }), new Promise((resolve) => { timer = setTimeout(resolve, Math.max(0, deadlineMs - Date.now())); }), diff --git a/packages/kap-server/test/telemetry.test.ts b/packages/kap-server/test/telemetry.test.ts index 5892691dfd..82cdaec70f 100644 --- a/packages/kap-server/test/telemetry.test.ts +++ b/packages/kap-server/test/telemetry.test.ts @@ -113,7 +113,11 @@ describe('server telemetry', () => { const telemetry = await initializeServerTelemetry(app, home as string); app.accessor.get(ITelemetryService).track('server_probe'); - await expect(shutdownServerTelemetry(telemetry, Date.now())).resolves.toBeUndefined(); + try { + await expect(shutdownServerTelemetry(telemetry, Date.now())).resolves.toBeUndefined(); + } finally { + await telemetry.appender?.shutdown(); + } }); it('keeps the null appender when config sets telemetry = false', async () => { From 3d2ed37122fc6bf40509dd26992657c34600d94d Mon Sep 17 00:00:00 2001 From: 7Sageer <7sageer@djwcb.cn> Date: Tue, 28 Jul 2026 18:43:16 +0800 Subject: [PATCH 2/7] fix(telemetry): close durable shutdown review gaps --- apps/kimi-code/src/cli/v2/run-v2-print.ts | 8 +- apps/kimi-code/test/cli/v2-run-print.test.ts | 20 ++ .../src/app/telemetry/cloudAppender.ts | 114 +++++++-- .../src/app/telemetry/cloudTransport.ts | 204 +++++----------- .../src/app/telemetry/telemetry.ts | 19 +- .../src/app/telemetry/telemetryService.ts | 112 +++++++-- .../src/app/telemetry/telemetrySpoolStore.ts | 164 +++++++++++++ .../test/agent/mcp/output.test.ts | 6 +- .../test/agent/media/tools/read-media.test.ts | 6 +- .../agent/plan/tools/exit-plan-mode.test.ts | 6 +- .../plan/tools/plan-tools-telemetry.test.ts | 6 +- .../test/app/telemetry/cloudAppender.test.ts | 218 +++++++++++++++++- .../agent-core-v2/test/app/telemetry/stubs.ts | 6 +- .../app/telemetry/telemetryService.test.ts | 123 +++++++++- .../os/backends/node-local/tools/glob.test.ts | 6 +- .../test/session/sessionFs/fsService.test.ts | 6 +- packages/kap-server/src/services/telemetry.ts | 20 +- packages/kap-server/test/boot.test.ts | 2 +- packages/kap-server/test/telemetry.test.ts | 11 +- 19 files changed, 798 insertions(+), 259 deletions(-) create mode 100644 packages/agent-core-v2/src/app/telemetry/telemetrySpoolStore.ts 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..2718300b3c 100644 --- a/apps/kimi-code/src/cli/v2/run-v2-print.ts +++ b/apps/kimi-code/src/cli/v2/run-v2-print.ts @@ -172,7 +172,11 @@ export async function runV2Print( await restorePermission(); } finally { if (telemetryService !== undefined) { - await raceWithTimeout(telemetryService.shutdown(), CLI_SHUTDOWN_TIMEOUT_MS); + const deadlineMs = Date.now() + CLI_SHUTDOWN_TIMEOUT_MS; + await raceWithTimeout( + telemetryService.shutdown({ deadlineMs }), + PROMPT_CLEANUP_TIMEOUT_MS, + ); } app.dispose(); } @@ -189,7 +193,7 @@ export async function runV2Print( // model is reconciled via setContext once resolved. telemetryService = app.accessor.get(ITelemetryService); if (telemetryEnabled) { - telemetryService.setAppender( + await telemetryService.setAppender( createCloudAppender(app.accessor, { deviceId, appName: CLI_USER_AGENT_PRODUCT, 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..36a918acb0 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,26 @@ describe('runV2Print', () => { expect(app.dispose).toHaveBeenCalled(); }); + it('forwards an absolute deadline into telemetry shutdown', async () => { + const stdout = writer(); + const stderr = writer(); + const { app, agent, appServices } = makeFakeHarness(); + const telemetry = appServices.get(ITelemetryService) as { + shutdown: ReturnType; + }; + const startedAt = Date.now(); + + mocks.bootstrap.mockReturnValue({ app }); + mocks.ensureMainAgent.mockResolvedValue(agent); + + await runV2Print(opts() as never, '1.2.3-test', { stdout, stderr }); + + expect(telemetry.shutdown).toHaveBeenCalledOnce(); + const options = telemetry.shutdown.mock.calls[0]?.[0] as { deadlineMs?: number } | undefined; + expect(options?.deadlineMs).toBeGreaterThanOrEqual(startedAt); + expect(options?.deadlineMs).toBeLessThanOrEqual(Date.now() + 3_000); + }); + it('seeds explicit skill dirs from --skillsDir into bootstrap', async () => { const stdout = writer(); const stderr = writer(); diff --git a/packages/agent-core-v2/src/app/telemetry/cloudAppender.ts b/packages/agent-core-v2/src/app/telemetry/cloudAppender.ts index 3beb31ecd9..9ea1e3f6a8 100644 --- a/packages/agent-core-v2/src/app/telemetry/cloudAppender.ts +++ b/packages/agent-core-v2/src/app/telemetry/cloudAppender.ts @@ -2,12 +2,13 @@ * `telemetry` domain (L1) — `CloudAppender`, an `ITelemetryAppender` that * batches events, drops non-primitive properties, redacts PII from string * values, enriches events with common context, and posts them to the - * telemetry endpoint through `CloudTransport`, which persists failed events - * through the `storage` byte layer. Reads host facts (`clientVersion`, env, - * platform/arch) from `IBootstrapService`; `createCloudAppender` assembles - * one from a `ServicesAccessor` so hosts only supply identity facts. Owns - * periodic flush, startup spool replay, and deadline-aware durable shutdown. - * App-scoped; independent of `@moonshot-ai/kimi-telemetry`. + * telemetry endpoint through `CloudTransport`, with durable handoff owned by + * the telemetry spool store, assembled over the `storage` byte layer. Reads + * host facts (`clientVersion`, env, platform/arch) from `IBootstrapService`; + * `createCloudAppender` assembles one from a `ServicesAccessor` so hosts only + * supply identity facts. Owns periodic flush, startup spool replay, and + * deadline-aware durable shutdown. App-scoped; independent of + * `@moonshot-ai/kimi-telemetry`. */ import { randomUUID } from 'node:crypto'; @@ -15,7 +16,7 @@ import { release } from 'node:os'; import type { ServicesAccessor } from '#/_base/di/instantiation'; import { onUnexpectedError } from '#/_base/errors/unexpectedError'; -import { abortError } from '#/_base/utils/abort'; +import { abortError, createDeadlineAbortSignal, isAbortError } from '#/_base/utils/abort'; import { IBootstrapService } from '#/app/bootstrap/bootstrap'; import { IFileSystemStorageService } from '#/persistence/interface/storage'; @@ -35,9 +36,10 @@ import { } from './cloudTransport'; import { resolveCoreVersion } from './coreVersion'; import { cleanTelemetryProperties } from './privacy'; +import { type ITelemetrySpoolStore, TelemetrySpoolStore } from './telemetrySpoolStore'; export interface CloudAppenderOptions { - readonly storage: IFileSystemStorageService; + readonly spool: ITelemetrySpoolStore; readonly bootstrap: IBootstrapService; readonly deviceId: string; readonly sessionId?: string; @@ -55,7 +57,8 @@ export interface CloudAppenderOptions { readonly retryBackoffsMs?: readonly number[]; readonly requestTimeoutMs?: number; readonly sleep?: (ms: number, signal?: AbortSignal) => Promise; - readonly now?: () => number; + readonly replayMaxFiles?: number; + readonly replayTimeoutMs?: number; } export interface CloudAppenderHostOptions { @@ -73,7 +76,9 @@ export function createCloudAppender( host: CloudAppenderHostOptions, ): CloudAppender { return new CloudAppender({ - storage: accessor.get(IFileSystemStorageService), + spool: new TelemetrySpoolStore({ + storage: accessor.get(IFileSystemStorageService), + }), bootstrap: accessor.get(IBootstrapService), ...host, }); @@ -81,12 +86,17 @@ export function createCloudAppender( const DEFAULT_FLUSH_THRESHOLD = 50; const DEFAULT_FLUSH_INTERVAL_MS = 30_000; +const DEFAULT_REPLAY_MAX_FILES = 20; +const DEFAULT_REPLAY_TIMEOUT_MS = 5_000; export class CloudAppender implements ITelemetryAppender { private readonly transport: CloudTransport; + private readonly spool: ITelemetrySpoolStore; private readonly context: CloudContext; private readonly flushThreshold: number; private readonly flushIntervalMs: number; + private readonly replayMaxFiles: number; + private readonly replayTimeoutMs: number; private deviceId: string; private sessionId: string | null; private buffer: EnrichedCloudEvent[] = []; @@ -94,7 +104,8 @@ export class CloudAppender implements ITelemetryAppender { private readonly lifecycleController = new AbortController(); private acceptingEvents = true; private started = false; - private startupReplay: Promise | null = null; + private replayPending = false; + private replayPromise: Promise | null = null; private flushPromise: Promise | null = null; private shutdownPromise: Promise | null = null; @@ -103,9 +114,11 @@ export class CloudAppender implements ITelemetryAppender { this.sessionId = options.sessionId ?? null; this.flushThreshold = options.flushThreshold ?? DEFAULT_FLUSH_THRESHOLD; this.flushIntervalMs = options.flushIntervalMs ?? DEFAULT_FLUSH_INTERVAL_MS; + this.replayMaxFiles = Math.max(0, Math.floor(options.replayMaxFiles ?? DEFAULT_REPLAY_MAX_FILES)); + this.replayTimeoutMs = Math.max(0, options.replayTimeoutMs ?? DEFAULT_REPLAY_TIMEOUT_MS); + this.spool = options.spool; this.context = buildContext(options); this.transport = new CloudTransport({ - storage: options.storage, deviceId: options.deviceId, endpoint: options.endpoint, getAccessToken: options.getAccessToken, @@ -113,7 +126,6 @@ export class CloudAppender implements ITelemetryAppender { retryBackoffsMs: options.retryBackoffsMs, requestTimeoutMs: options.requestTimeoutMs, sleep: options.sleep, - now: options.now, }); } @@ -152,7 +164,9 @@ export class CloudAppender implements ITelemetryAppender { flush(): Promise { if (this.flushPromise !== null) return this.flushPromise; - if (this.buffer.length === 0) return Promise.resolve(); + if (this.buffer.length === 0 && !this.replayPending && this.replayPromise === null) { + return Promise.resolve(); + } const flush = this.drainBuffer(); this.flushPromise = flush; void flush.then( @@ -181,9 +195,8 @@ export class CloudAppender implements ITelemetryAppender { start(): void { if (this.started || !this.acceptingEvents) return; this.started = true; - this.startupReplay = this.transport - .retryDiskEvents(this.lifecycleController.signal) - .catch(() => {}); + this.replayPending = true; + void this.ensureReplay(); this.startPeriodicFlush(); } @@ -202,19 +215,22 @@ export class CloudAppender implements ITelemetryAppender { } async retryDiskEvents(): Promise { - await this.transport.retryDiskEvents(this.lifecycleController.signal); + this.replayPending = true; + await this.ensureReplay(); } private async drainBuffer(): Promise { - await (this.startupReplay ?? Promise.resolve()); + if (!(await this.ensureReplay())) { + await this.handoffBufferedEvents(); + return; + } while (this.buffer.length > 0) { const events = this.buffer; this.buffer = []; try { await this.transport.send(events, this.lifecycleController.signal); - } catch (error) { - this.buffer = [...events, ...this.buffer]; - throw error; + } catch { + await this.handoffEvents(events); } } } @@ -224,10 +240,62 @@ export class CloudAppender implements ITelemetryAppender { } private async shutdownOwnedWork(): Promise { - await (this.startupReplay ?? Promise.resolve()); await this.flush(); } + private ensureReplay(): Promise { + if (!this.replayPending) return Promise.resolve(true); + if (this.replayPromise !== null) return this.replayPromise; + const replay = this.replayDiskEvents(this.lifecycleController.signal).catch(() => false); + this.replayPromise = replay; + void replay.then((complete) => { + this.replayPending = !complete; + if (this.replayPromise === replay) this.replayPromise = null; + }); + return replay; + } + + private async handoffBufferedEvents(): Promise { + while (this.buffer.length > 0) { + const events = this.buffer; + this.buffer = []; + await this.handoffEvents(events); + } + } + + private async handoffEvents(events: readonly EnrichedCloudEvent[]): Promise { + try { + await this.spool.put(events); + } catch (storageError) { + this.buffer = [...events, ...this.buffer]; + throw storageError; + } + } + + private async replayDiskEvents(signal: AbortSignal): Promise { + const deadline = createDeadlineAbortSignal(signal, this.replayTimeoutMs); + try { + const entries = await this.spool.recoverable(this.replayMaxFiles + 1); + let complete = entries.length <= this.replayMaxFiles; + for (const entry of entries.slice(0, this.replayMaxFiles)) { + if (deadline.signal.aborted) { + complete = false; + break; + } + try { + await this.transport.send(entry.events, deadline.signal, []); + await this.spool.acknowledge(entry.key); + } catch (error) { + complete = false; + if (deadline.signal.aborted || isAbortError(error)) break; + } + } + return complete && !deadline.signal.aborted; + } finally { + deadline.clear(); + } + } + private armShutdownDeadline(options: TelemetryShutdownOptions): () => void { const abort = (): void => { if (!this.lifecycleController.signal.aborted) { diff --git a/packages/agent-core-v2/src/app/telemetry/cloudTransport.ts b/packages/agent-core-v2/src/app/telemetry/cloudTransport.ts index 7544f166a7..eccc8df3ce 100644 --- a/packages/agent-core-v2/src/app/telemetry/cloudTransport.ts +++ b/packages/agent-core-v2/src/app/telemetry/cloudTransport.ts @@ -1,16 +1,17 @@ /** - * `telemetry` domain (L1) — `CloudTransport`, the HTTP transport behind - * `CloudAppender`. Posts enriched events to the telemetry endpoint with Bearer - * auth, retry, and a byte-store fallback for failed events, persisted through - * the `storage` byte layer (`IFileSystemStorageService`) under an isolated - * `telemetry-v2` scope. - * App-scoped; independent of `@moonshot-ai/kimi-telemetry`. + * `telemetry` domain (L1) — HTTP delivery transport behind `CloudAppender`. + * + * Owns authentication, request deadlines, retry, and telemetry wire shaping. + * Durable handoff and recovery are coordinated by `CloudAppender` through its + * spool store. App-scoped; independent of `@moonshot-ai/kimi-telemetry`. */ -import { randomBytes } from 'node:crypto'; - -import { abortable, isAbortError } from '#/_base/utils/abort'; -import type { IFileSystemStorageService } from '#/persistence/interface/storage'; +import { + abortable, + abortError, + createDeadlineAbortSignal, + isAbortError, +} from '#/_base/utils/abort'; export type CloudPrimitive = boolean | number | string | undefined | null; @@ -37,7 +38,6 @@ export interface CloudPayload { } export interface CloudTransportOptions { - readonly storage: IFileSystemStorageService; readonly deviceId: string; readonly endpoint?: string; readonly getAccessToken?: () => string | null | Promise; @@ -45,25 +45,16 @@ export interface CloudTransportOptions { readonly retryBackoffsMs?: readonly number[]; readonly requestTimeoutMs?: number; readonly sleep?: (ms: number, signal?: AbortSignal) => Promise; - readonly now?: () => number; } export const TELEMETRY_ENDPOINT = 'https://telemetry-logs.kimi.com/v1/event'; export const SERVER_EVENT_PREFIX = 'kfc_'; export const USER_ID_PREFIX = 'kfc_device_id_'; -export const DISK_EVENT_MAX_AGE_MS = 7 * 24 * 60 * 60 * 1000; export const RETRY_BACKOFFS_MS = [1_000, 4_000, 16_000] as const; const DEFAULT_REQUEST_TIMEOUT_MS = 10_000; -const TELEMETRY_SCOPE = 'telemetry-v2'; -const FAILED_PREFIX = 'failed_'; -const JSONL_SUFFIX = '.jsonl'; - -const textEncoder = new TextEncoder(); -const textDecoder = new TextDecoder(); export class CloudTransport { - private readonly storage: IFileSystemStorageService; private readonly deviceId: string; private readonly endpoint: string; private readonly getAccessToken: (() => string | null | Promise) | null; @@ -71,10 +62,8 @@ export class CloudTransport { private readonly retryBackoffsMs: readonly number[]; private readonly requestTimeoutMs: number; private readonly sleepImpl: (ms: number, signal?: AbortSignal) => Promise; - private readonly now: () => number; constructor(options: CloudTransportOptions) { - this.storage = options.storage; this.deviceId = options.deviceId; this.endpoint = options.endpoint ?? TELEMETRY_ENDPOINT; this.getAccessToken = options.getAccessToken ?? null; @@ -82,15 +71,15 @@ export class CloudTransport { this.retryBackoffsMs = options.retryBackoffsMs ?? RETRY_BACKOFFS_MS; this.requestTimeoutMs = options.requestTimeoutMs ?? DEFAULT_REQUEST_TIMEOUT_MS; this.sleepImpl = options.sleep ?? abortableSleep; - this.now = options.now ?? Date.now; } - async send(events: readonly EnrichedCloudEvent[], signal?: AbortSignal): Promise { + async send( + events: readonly EnrichedCloudEvent[], + signal?: AbortSignal, + retryBackoffsMs: readonly number[] = this.retryBackoffsMs, + ): Promise { if (events.length === 0) return; - if (signal?.aborted === true) { - await this.saveToDisk(events); - return; - } + if (signal?.aborted === true) throw abortError(); let payload: CloudPayload; try { @@ -99,93 +88,53 @@ export class CloudTransport { return; } - for (let attempt = 0; attempt <= this.retryBackoffsMs.length; attempt++) { + let lastError: unknown; + for (let attempt = 0; attempt <= retryBackoffsMs.length; attempt++) { try { - const request = this.sendHttp(payload, signal); - await (signal === undefined ? request : abortable(request, signal)); + await this.sendAttempt(payload, signal); return; } catch (error) { - if (isSignalAborted(signal) || isAbortError(error)) { - await this.saveToDisk(events); - return; + if (isSignalAborted(signal) || (isAbortError(error) && !(error instanceof TransientCloudError))) { + throw error; } + lastError = error; if (!(error instanceof TransientCloudError)) break; - const backoff = this.retryBackoffsMs[attempt]; + const backoff = retryBackoffsMs[attempt]; if (backoff === undefined) break; + const sleep = Promise.resolve().then(() => this.sleepImpl(backoff, signal)); try { - const sleep = this.sleepImpl(backoff, signal); await (signal === undefined ? sleep : abortable(sleep, signal)); } catch (sleepError) { - if (isSignalAborted(signal) || isAbortError(sleepError)) { - await this.saveToDisk(events); - return; - } + if (isSignalAborted(signal) || isAbortError(sleepError)) throw sleepError; + lastError = sleepError; break; } } } - - await this.saveToDisk(events); - } - - async saveToDisk(events: readonly EnrichedCloudEvent[]): Promise { - if (events.length === 0) return; - const key = `${FAILED_PREFIX}${this.now()}_${randomBytes(6).toString('hex')}${JSONL_SUFFIX}`; - const text = events.map((event) => JSON.stringify(event)).join('\n') + '\n'; - await this.storage.write(TELEMETRY_SCOPE, key, textEncoder.encode(text)); - } - - async retryDiskEvents(signal?: AbortSignal): Promise { - const keys = await this.storage.list(TELEMETRY_SCOPE, FAILED_PREFIX); - const now = this.now(); - for (const key of keys) { - if (signal?.aborted === true) throw abortError(); - if (!key.startsWith(FAILED_PREFIX) || !key.endsWith(JSONL_SUFFIX)) continue; - const createdAt = parseFailedTimestamp(key); - if (createdAt === undefined) continue; - if (now - createdAt > DISK_EVENT_MAX_AGE_MS) { - await this.storage.delete(TELEMETRY_SCOPE, key).catch(() => undefined); - continue; - } - - let events: EnrichedCloudEvent[]; - let payload: CloudPayload; - try { - events = await this.readJsonl(key); - payload = buildPayload(events, this.deviceId); - } catch (error) { - if (error instanceof SyntaxError || error instanceof TypeError) { - await this.storage.delete(TELEMETRY_SCOPE, key).catch(() => undefined); - } - continue; - } - - try { - const request = this.sendHttp(payload, signal); - await (signal === undefined ? request : abortable(request, signal)); - await this.storage.delete(TELEMETRY_SCOPE, key); - } catch (error) { - if (isSignalAborted(signal) || isAbortError(error)) throw error; - if (error instanceof TransientCloudError) continue; - } - } + throw lastError instanceof Error + ? lastError + : new TransientCloudError('telemetry delivery failed'); } - private async readJsonl(key: string): Promise { - const bytes = await this.storage.read(TELEMETRY_SCOPE, key); - if (bytes === undefined) return []; - const text = textDecoder.decode(bytes); - const events: EnrichedCloudEvent[] = []; - for (const line of text.split('\n')) { - const trimmed = line.trim(); - if (trimmed.length === 0) continue; - events.push(JSON.parse(trimmed) as EnrichedCloudEvent); + private async sendAttempt(payload: CloudPayload, signal?: AbortSignal): Promise { + const source = signal ?? new AbortController().signal; + const deadline = createDeadlineAbortSignal(source, this.requestTimeoutMs); + try { + await this.sendHttp(payload, deadline.signal); + } catch (error) { + if (deadline.timedOut()) throw new TransientCloudError('telemetry request timed out'); + throw error; + } finally { + deadline.clear(); } - return events; } - private async sendHttp(payload: CloudPayload, signal?: AbortSignal): Promise { - const token = this.getAccessToken === null ? null : await this.getAccessToken(); + private async sendHttp(payload: CloudPayload, signal: AbortSignal): Promise { + const tokenRequest = + this.getAccessToken === null + ? Promise.resolve(null) + : Promise.resolve().then(() => this.getAccessToken?.() ?? null); + const token = await abortable(tokenRequest, signal); const headers: Record = { 'Content-Type': 'application/json', }; @@ -206,36 +155,25 @@ export class CloudTransport { private async post( payload: CloudPayload, headers: Record, - signal?: AbortSignal, + signal: AbortSignal, ): Promise { try { - return await fetchWithTimeout( - this.fetchImpl, - this.endpoint, - { + const request = Promise.resolve().then(() => + this.fetchImpl(this.endpoint, { method: 'POST', headers: { ...headers }, body: JSON.stringify(payload), - }, - this.requestTimeoutMs, - signal, + signal, + }), ); + return await abortable(request, signal); } catch (error) { - if (signal?.aborted === true || isAbortError(error)) throw error; + if (signal.aborted || isAbortError(error)) throw error; throw new TransientCloudError(String(error)); } } } -function parseFailedTimestamp(key: string): number | undefined { - const rest = key.slice(FAILED_PREFIX.length); - const underscore = rest.indexOf('_'); - if (underscore === -1) return undefined; - const raw = rest.slice(0, underscore); - const ts = Number(raw); - return Number.isFinite(ts) ? ts : undefined; -} - export class TransientCloudError extends Error { override readonly name = 'TransientCloudError'; } @@ -306,37 +244,7 @@ function handleStatus(status: number): void { if (status >= 500 || status === 429) { throw new TransientCloudError(`HTTP ${String(status)}`); } - if (status >= 400) { - return; - } -} - -async function fetchWithTimeout( - fetchImpl: typeof fetch, - url: string, - init: RequestInit, - timeoutMs: number, - externalSignal?: AbortSignal, -): Promise { - const controller = new AbortController(); - const abortFromExternal = (): void => { - controller.abort(externalSignal?.reason); - }; - const timeout = setTimeout(() => { - controller.abort(new Error('telemetry request timed out')); - }, timeoutMs); - timeout.unref?.(); - if (externalSignal?.aborted === true) abortFromExternal(); - externalSignal?.addEventListener('abort', abortFromExternal, { once: true }); - try { - return await fetchImpl(url, { - ...init, - signal: controller.signal, - }); - } finally { - clearTimeout(timeout); - externalSignal?.removeEventListener('abort', abortFromExternal); - } + if (status >= 400) return; } function abortableSleep(ms: number, signal?: AbortSignal): Promise { @@ -359,7 +267,3 @@ function abortableSleep(ms: number, signal?: AbortSignal): Promise { function isSignalAborted(signal?: AbortSignal): boolean { return signal?.aborted === true; } - -function abortError(): DOMException { - return new DOMException('The operation was aborted.', 'AbortError'); -} diff --git a/packages/agent-core-v2/src/app/telemetry/telemetry.ts b/packages/agent-core-v2/src/app/telemetry/telemetry.ts index 74bf1b2067..54a3bf0c5b 100644 --- a/packages/agent-core-v2/src/app/telemetry/telemetry.ts +++ b/packages/agent-core-v2/src/app/telemetry/telemetry.ts @@ -38,6 +38,10 @@ export interface ITelemetryAppender { shutdown?(options?: TelemetryShutdownOptions): Promise | void; } +export interface ITelemetryAppenderRegistration extends IDisposable { + shutdown(options?: TelemetryShutdownOptions): Promise; +} + export interface TelemetryServiceOptions { readonly appender?: ITelemetryAppender; readonly appenders?: readonly ITelemetryAppender[]; @@ -57,9 +61,12 @@ export interface ITelemetryService { ): void; withContext(patch: TelemetryContextPatch): ITelemetryService; setContext(patch: TelemetryContextPatch): void; - addAppender(appender: ITelemetryAppender): IDisposable; - removeAppender(appender: ITelemetryAppender): void; - setAppender(appender: ITelemetryAppender): void; + addAppender(appender: ITelemetryAppender): ITelemetryAppenderRegistration; + removeAppender( + appender: ITelemetryAppender, + options?: TelemetryShutdownOptions, + ): Promise; + setAppender(appender: ITelemetryAppender, options?: TelemetryShutdownOptions): Promise; setEnabled(enabled: boolean): void; flush(): Promise; shutdown(options?: TelemetryShutdownOptions): Promise; @@ -79,9 +86,9 @@ export const noopTelemetryService: ITelemetryService = { track2: () => {}, withContext: () => noopTelemetryService, setContext: () => {}, - addAppender: () => ({ dispose: () => {} }), - removeAppender: () => {}, - setAppender: () => {}, + addAppender: () => ({ dispose: () => {}, shutdown: async () => {} }), + removeAppender: async () => {}, + setAppender: async () => {}, setEnabled: () => {}, flush: async () => {}, shutdown: async () => {}, diff --git a/packages/agent-core-v2/src/app/telemetry/telemetryService.ts b/packages/agent-core-v2/src/app/telemetry/telemetryService.ts index a860ef8a1c..66ef6b556b 100644 --- a/packages/agent-core-v2/src/app/telemetry/telemetryService.ts +++ b/packages/agent-core-v2/src/app/telemetry/telemetryService.ts @@ -7,9 +7,9 @@ * App-scoped root. Has no cross-domain collaborators. */ -import { type IDisposable, toDisposable } from '#/_base/di/lifecycle'; import { LifecycleScope, ScopeActivation, registerScopedService } from '#/_base/di/scope'; import { onUnexpectedError } from '#/_base/errors/unexpectedError'; +import { BugIndicatingError } from '#/errors'; import type { StrictPropertyCheck, @@ -19,6 +19,7 @@ import type { import { ITelemetryService, type ITelemetryAppender, + type ITelemetryAppenderRegistration, nullTelemetryAppender, type TelemetryContextPatch, type TelemetryProperties, @@ -29,6 +30,12 @@ export class TelemetryService implements ITelemetryService { declare readonly _serviceBrand: undefined; private appenders: ITelemetryAppender[] = [nullTelemetryAppender]; + private readonly registrations = new Map< + ITelemetryAppender, + ITelemetryAppenderRegistration + >(); + private readonly retirements = new Map>(); + private shutdownPromise: Promise | null = null; private context: TelemetryProperties = {}; private enabled = true; @@ -40,8 +47,8 @@ export class TelemetryService implements ITelemetryService { for (const appender of this.appenders) { try { appender.track(event, merged); - } catch (err) { - onUnexpectedError(err); + } catch (error) { + onUnexpectedError(error); } } } @@ -64,19 +71,36 @@ export class TelemetryService implements ITelemetryService { } } - addAppender(appender: ITelemetryAppender): IDisposable { + addAppender(appender: ITelemetryAppender): ITelemetryAppenderRegistration { + this.assertOpen(); + const existing = this.registrations.get(appender); + if (existing !== undefined) return existing; this.startAppender(appender); this.appenders.push(appender); - return toDisposable(() => this.removeAppender(appender)); + return this.createRegistration(appender); } - removeAppender(appender: ITelemetryAppender): void { + removeAppender( + appender: ITelemetryAppender, + options?: TelemetryShutdownOptions, + ): Promise { + const retirement = this.retirements.get(appender); + if (retirement !== undefined) return retirement; + if (!this.appenders.includes(appender)) return Promise.resolve(); this.appenders = this.appenders.filter((a) => a !== appender); + return this.removeDetachedAppender(appender, options); } - setAppender(appender: ITelemetryAppender): void { - this.startAppender(appender); + async setAppender( + appender: ITelemetryAppender, + options?: TelemetryShutdownOptions, + ): Promise { + this.assertOpen(); + if (!this.appenders.includes(appender)) this.startAppender(appender); + if (!this.registrations.has(appender)) this.createRegistration(appender); + const previous = this.appenders.filter((candidate) => candidate !== appender); this.appenders = [appender]; + await Promise.all(previous.map((candidate) => this.removeDetachedAppender(candidate, options))); } setEnabled(enabled: boolean): void { @@ -85,18 +109,24 @@ export class TelemetryService implements ITelemetryService { async flush(): Promise { await Promise.all( - this.appenders.map((appender) => - Promise.resolve(appender.flush?.()).catch(onUnexpectedError), - ), + this.appenders.map((appender) => this.invokeAppender(() => appender.flush?.())), ); } - async shutdown(options?: TelemetryShutdownOptions): Promise { - await Promise.all( - this.appenders.map((appender) => - Promise.resolve(appender.shutdown?.(options)).catch(onUnexpectedError), - ), - ); + shutdown(options?: TelemetryShutdownOptions): Promise { + if (this.shutdownPromise === null) { + const appenders = this.appenders; + this.appenders = []; + for (const appender of appenders) { + void this.removeDetachedAppender(appender, options); + } + this.shutdownPromise = Promise.all(this.retirements.values()).then(() => undefined); + } else if (options !== undefined) { + for (const appender of this.registrations.keys()) { + void this.invokeAppender(() => appender.shutdown?.(options)); + } + } + return this.shutdownPromise; } private startAppender(appender: ITelemetryAppender): void { @@ -106,6 +136,38 @@ export class TelemetryService implements ITelemetryService { onUnexpectedError(error); } } + + private createRegistration(appender: ITelemetryAppender): ITelemetryAppenderRegistration { + const registration: ITelemetryAppenderRegistration = { + dispose: () => { + void this.removeAppender(appender); + }, + shutdown: (options) => this.removeAppender(appender, options), + }; + this.registrations.set(appender, registration); + return registration; + } + + private assertOpen(): void { + if (this.shutdownPromise !== null) { + throw new BugIndicatingError('Telemetry service has already shut down'); + } + } + + private removeDetachedAppender( + appender: ITelemetryAppender, + options?: TelemetryShutdownOptions, + ): Promise { + const retirement = this.retirements.get(appender); + if (retirement !== undefined) return retirement; + const pending = this.invokeAppender(() => appender.shutdown?.(options)); + this.retirements.set(appender, pending); + return pending; + } + + private invokeAppender(operation: () => Promise | void | undefined): Promise { + return Promise.resolve().then(operation).catch(onUnexpectedError); + } } class TelemetryContextView implements ITelemetryService { @@ -138,16 +200,22 @@ class TelemetryContextView implements ITelemetryService { this.context = { ...this.context, ...patch }; } - addAppender(appender: ITelemetryAppender): IDisposable { + addAppender(appender: ITelemetryAppender): ITelemetryAppenderRegistration { return this.root.addAppender(appender); } - removeAppender(appender: ITelemetryAppender): void { - this.root.removeAppender(appender); + removeAppender( + appender: ITelemetryAppender, + options?: TelemetryShutdownOptions, + ): Promise { + return this.root.removeAppender(appender, options); } - setAppender(appender: ITelemetryAppender): void { - this.root.setAppender(appender); + setAppender( + appender: ITelemetryAppender, + options?: TelemetryShutdownOptions, + ): Promise { + return this.root.setAppender(appender, options); } setEnabled(enabled: boolean): void { diff --git a/packages/agent-core-v2/src/app/telemetry/telemetrySpoolStore.ts b/packages/agent-core-v2/src/app/telemetry/telemetrySpoolStore.ts new file mode 100644 index 0000000000..c8bfdce611 --- /dev/null +++ b/packages/agent-core-v2/src/app/telemetry/telemetrySpoolStore.ts @@ -0,0 +1,164 @@ +/** + * `telemetry` domain (L1) — durable failed-delivery spool access-pattern store. + * + * Owns the telemetry namespace, JSONL wire format, retention, capacity, and + * acknowledgement protocol over the `storage` byte layer. App-scoped and + * assembled privately by `CloudAppender`. + */ + +import { randomBytes } from 'node:crypto'; + +import type { IFileSystemStorageService } from '#/persistence/interface/storage'; + +import type { EnrichedCloudEvent } from './cloudTransport'; + +export interface TelemetrySpoolEntry { + readonly key: string; + readonly events: readonly EnrichedCloudEvent[]; +} + +export interface ITelemetrySpoolStore { + put(events: readonly EnrichedCloudEvent[]): Promise; + recoverable(limit: number): Promise; + acknowledge(key: string): Promise; +} + +export interface TelemetrySpoolStoreOptions { + readonly storage: IFileSystemStorageService; + readonly maxFiles?: number; + readonly now?: () => number; +} + +export const TELEMETRY_SPOOL_MAX_AGE_MS = 7 * 24 * 60 * 60 * 1000; +export const TELEMETRY_SPOOL_MAX_FILES = 256; + +const TELEMETRY_SCOPE = 'telemetry-v2'; +const FAILED_PREFIX = 'failed_'; +const JSONL_SUFFIX = '.jsonl'; +const textEncoder = new TextEncoder(); +const textDecoder = new TextDecoder(); + +export class TelemetrySpoolStore implements ITelemetrySpoolStore { + private readonly storage: IFileSystemStorageService; + private readonly maxFiles: number; + private readonly now: () => number; + private writes: Promise = Promise.resolve(); + + constructor(options: TelemetrySpoolStoreOptions) { + this.storage = options.storage; + this.maxFiles = Math.max(1, Math.floor(options.maxFiles ?? TELEMETRY_SPOOL_MAX_FILES)); + this.now = options.now ?? Date.now; + } + + put(events: readonly EnrichedCloudEvent[]): Promise { + if (events.length === 0) return Promise.resolve(); + const pending = this.writes.then(() => this.putOwned(events)); + this.writes = pending.catch(() => undefined); + return pending; + } + + async recoverable(limit: number): Promise { + await this.writes; + const files = await this.enforcePolicy(); + const entries: TelemetrySpoolEntry[] = []; + const boundedLimit = Math.max(0, Math.floor(limit)); + for (const file of files) { + if (entries.length >= boundedLimit) break; + try { + entries.push({ key: file.key, events: await this.readJsonl(file.key) }); + } catch (error) { + if (error instanceof SyntaxError || error instanceof TypeError) { + await this.acknowledge(file.key).catch(() => undefined); + continue; + } + throw error; + } + } + return entries; + } + + acknowledge(key: string): Promise { + return this.storage.delete(TELEMETRY_SCOPE, key); + } + + private async enforcePolicy(): Promise { + const keys = await this.storage.list(TELEMETRY_SCOPE, FAILED_PREFIX); + const now = this.now(); + const files: SpoolFile[] = []; + const discarded: string[] = []; + for (const key of keys) { + const createdAt = parseCreatedAt(key); + if (createdAt === undefined || now - createdAt > TELEMETRY_SPOOL_MAX_AGE_MS) { + discarded.push(key); + } else { + files.push({ key, createdAt }); + } + } + files.sort((left, right) => left.createdAt - right.createdAt || left.key.localeCompare(right.key)); + await Promise.all(discarded.map((key) => this.acknowledge(key).catch(() => undefined))); + return files; + } + + private async putOwned(events: readonly EnrichedCloudEvent[]): Promise { + const text = events.map((event) => JSON.stringify(event)).join('\n') + '\n'; + const files = await this.enforcePolicy().catch(() => []); + if (files.length < this.maxFiles) { + await this.storage.write(TELEMETRY_SCOPE, this.createKey(), textEncoder.encode(text), { + atomic: true, + }); + return; + } + + const compactedCount = files.length - this.maxFiles + 1; + const compacted = files.slice(0, compactedCount); + const existing = await Promise.all( + compacted.map(async (file) => { + const bytes = await this.storage.read(TELEMETRY_SCOPE, file.key); + if (bytes === undefined) return ''; + const jsonl = textDecoder.decode(bytes); + return jsonl.length === 0 || jsonl.endsWith('\n') ? jsonl : jsonl + '\n'; + }), + ); + await this.storage.write( + TELEMETRY_SCOPE, + this.createKey(), + textEncoder.encode(existing.join('') + text), + { atomic: true }, + ); + await Promise.all(compacted.map((file) => this.acknowledge(file.key).catch(() => undefined))); + } + + private createKey(): string { + return `${FAILED_PREFIX}${this.now()}_${randomBytes(6).toString('hex')}${JSONL_SUFFIX}`; + } + + private async readJsonl(key: string): Promise { + const bytes = await this.storage.read(TELEMETRY_SCOPE, key); + if (bytes === undefined) return []; + const events: EnrichedCloudEvent[] = []; + for (const line of textDecoder.decode(bytes).split('\n')) { + const trimmed = line.trim(); + if (trimmed.length === 0) continue; + const event: unknown = JSON.parse(trimmed); + if (event === null || typeof event !== 'object' || Array.isArray(event)) { + throw new TypeError('telemetry spool entry must be an object'); + } + events.push(event as EnrichedCloudEvent); + } + return events; + } +} + +interface SpoolFile { + readonly key: string; + readonly createdAt: number; +} + +function parseCreatedAt(key: string): number | undefined { + if (!key.startsWith(FAILED_PREFIX) || !key.endsWith(JSONL_SUFFIX)) return undefined; + const rest = key.slice(FAILED_PREFIX.length); + const underscore = rest.indexOf('_'); + if (underscore === -1) return undefined; + const createdAt = Number(rest.slice(0, underscore)); + return Number.isFinite(createdAt) ? createdAt : undefined; +} diff --git a/packages/agent-core-v2/test/agent/mcp/output.test.ts b/packages/agent-core-v2/test/agent/mcp/output.test.ts index e051fd31d9..194fca6815 100644 --- a/packages/agent-core-v2/test/agent/mcp/output.test.ts +++ b/packages/agent-core-v2/test/agent/mcp/output.test.ts @@ -43,9 +43,9 @@ function recordingTelemetry(records: TelemetryRecord[]): ITelemetryService { track2: (event, properties) => telemetry.track(event, properties as TelemetryProperties), withContext: () => telemetry, setContext: () => {}, - addAppender: () => ({ dispose: () => {} }), - removeAppender: () => {}, - setAppender: () => {}, + addAppender: () => ({ dispose: () => {}, shutdown: async () => {} }), + removeAppender: async () => {}, + setAppender: async () => {}, setEnabled: () => {}, flush: async () => {}, shutdown: async () => {}, diff --git a/packages/agent-core-v2/test/agent/media/tools/read-media.test.ts b/packages/agent-core-v2/test/agent/media/tools/read-media.test.ts index 02c9d0f670..91f6f54d33 100644 --- a/packages/agent-core-v2/test/agent/media/tools/read-media.test.ts +++ b/packages/agent-core-v2/test/agent/media/tools/read-media.test.ts @@ -110,9 +110,9 @@ function recordingTelemetry(records: TelemetryRecord[]): ITelemetryService { track2: (event, properties) => telemetry.track(event, properties as TelemetryProperties), withContext: () => telemetry, setContext: () => {}, - addAppender: () => ({ dispose: () => {} }), - removeAppender: () => {}, - setAppender: () => {}, + addAppender: () => ({ dispose: () => {}, shutdown: async () => {} }), + removeAppender: async () => {}, + setAppender: async () => {}, setEnabled: () => {}, flush: async () => {}, shutdown: async () => {}, diff --git a/packages/agent-core-v2/test/agent/plan/tools/exit-plan-mode.test.ts b/packages/agent-core-v2/test/agent/plan/tools/exit-plan-mode.test.ts index 72c2cacdef..0866b7f0e2 100644 --- a/packages/agent-core-v2/test/agent/plan/tools/exit-plan-mode.test.ts +++ b/packages/agent-core-v2/test/agent/plan/tools/exit-plan-mode.test.ts @@ -44,9 +44,9 @@ function recordingTelemetry(): ITelemetryService { track2: vi.fn(), withContext: () => recordingTelemetry(), setContext: () => {}, - addAppender: () => ({ dispose: () => {} }), - removeAppender: () => {}, - setAppender: () => {}, + addAppender: () => ({ dispose: () => {}, shutdown: async () => {} }), + removeAppender: async () => {}, + setAppender: async () => {}, setEnabled: () => {}, flush: () => Promise.resolve(), shutdown: () => Promise.resolve(), diff --git a/packages/agent-core-v2/test/agent/plan/tools/plan-tools-telemetry.test.ts b/packages/agent-core-v2/test/agent/plan/tools/plan-tools-telemetry.test.ts index 6672a56e4e..69c5160199 100644 --- a/packages/agent-core-v2/test/agent/plan/tools/plan-tools-telemetry.test.ts +++ b/packages/agent-core-v2/test/agent/plan/tools/plan-tools-telemetry.test.ts @@ -48,9 +48,9 @@ function recordingTelemetry(): { track2, withContext: () => recordingTelemetry().telemetry, setContext: () => {}, - addAppender: () => ({ dispose: () => {} }), - removeAppender: () => {}, - setAppender: () => {}, + addAppender: () => ({ dispose: () => {}, shutdown: async () => {} }), + removeAppender: async () => {}, + setAppender: async () => {}, setEnabled: () => {}, flush: () => Promise.resolve(), shutdown: () => Promise.resolve(), diff --git a/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts b/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts index 5c543bfd4f..a26d9ec841 100644 --- a/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts +++ b/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts @@ -1,9 +1,8 @@ /** - * Cloud telemetry lifecycle tests — exercise the real appender, transport, - * and file-storage stack while stubbing only the outbound HTTP boundary. - * Covers batching, durable shutdown, startup replay, privacy, and wire shape. - * Run with `pnpm --filter @moonshot-ai/agent-core-v2 exec vitest run - * test/app/telemetry/cloudAppender.test.ts`. + * `telemetry` domain (L1) — cloud appender lifecycle integration coverage. + * + * Exercises batching, durable shutdown, startup replay, privacy, and wire + * shape through the real appender, transport, spool, and file-storage stack. */ import { getEventListeners } from 'node:events'; @@ -18,7 +17,7 @@ import { import { tmpdir } from 'node:os'; import { join } from 'node:path'; -import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import { resetUnexpectedErrorHandler, @@ -27,6 +26,7 @@ import { import { FileStorageService } from '#/persistence/backends/node-fs/fileStorageService'; import { CloudAppender, type CloudAppenderOptions } from '#/app/telemetry/cloudAppender'; import { CloudTransport } from '#/app/telemetry/cloudTransport'; +import { TelemetrySpoolStore } from '#/app/telemetry/telemetrySpoolStore'; import { stubBootstrap } from '../bootstrap/stubs'; @@ -79,11 +79,21 @@ function deferred(): { } function baseOptions( - overrides: Partial & { homeDir?: string } = {}, + overrides: Partial & { + homeDir?: string; + now?: () => number; + spoolMaxFiles?: number; + } = {}, ): CloudAppenderOptions { - const { homeDir: dir = '', storage, ...rest } = overrides; + const { homeDir: dir = '', now, spool, spoolMaxFiles, ...rest } = overrides; return { - storage: storage ?? new FileStorageService(dir), + spool: + spool ?? + new TelemetrySpoolStore({ + storage: new FileStorageService(dir), + maxFiles: spoolMaxFiles, + now, + }), bootstrap: { ...stubBootstrap(), clientVersion: '1.0.0' }, deviceId: 'dev', appName: 'test-app', @@ -104,6 +114,16 @@ function readFirstFailedEvent(homeDir: string): Record { return JSON.parse(persisted.trim()) as Record; } +function readAllFailedEvents(homeDir: string): Record[] { + return listFailedSpoolFiles(homeDir).flatMap((file) => + readFileSync(join(homeDir, 'telemetry-v2', file), 'utf8') + .trim() + .split('\n') + .filter((line) => line.length > 0) + .map((line) => JSON.parse(line) as Record), + ); +} + describe('CloudAppender', () => { let homeDir: string; @@ -436,7 +456,6 @@ describe('CloudAppender', () => { it('releases lifecycle abort listeners after a retry backoff completes', async () => { let attempts = 0; const transport = new CloudTransport({ - storage: new FileStorageService(homeDir), deviceId: 'dev', retryBackoffsMs: [0], fetchImpl: makeFetch(() => { @@ -465,6 +484,86 @@ describe('CloudAppender', () => { expect(getEventListeners(lifecycle.signal, 'abort')).toHaveLength(0); }); + it('shutdown cancellation interrupts token lookup and durably hands off the batch', async () => { + const token = deferred(); + const tokenStarted = deferred(); + let requests = 0; + const appender = new CloudAppender( + baseOptions({ + homeDir, + getAccessToken: () => { + tokenStarted.resolve(); + return token.promise; + }, + fetchImpl: makeFetch(() => { + requests += 1; + return okResponse(); + }), + }), + ); + const cancellation = new AbortController(); + appender.track('token_lookup_cancelled'); + + const closing = appender.shutdown({ signal: cancellation.signal }); + await tokenStarted.promise; + cancellation.abort(); + await closing; + token.resolve(null); + + expect(requests).toBe(0); + const persisted = readFirstFailedEvent(homeDir); + expect(persisted).toMatchObject({ + event: 'token_lookup_cancelled', + }); + + let replayedEventId: unknown; + const restartedAppender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch((request) => { + replayedEventId = request.body.events[0]?.['event_id']; + return okResponse(); + }), + }), + ); + restartedAppender.start(); + await restartedAppender.shutdown(); + + expect(replayedEventId).toBe(persisted['event_id']); + }); + + it('shutdown cancellation interrupts retry sleep and durably hands off the batch', async () => { + const sleepStarted = deferred(); + const sleep = deferred(); + let attempts = 0; + const appender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => { + attempts += 1; + return statusResponse(500); + }), + sleep: () => { + sleepStarted.resolve(); + return sleep.promise; + }, + }), + ); + const cancellation = new AbortController(); + appender.track('retry_sleep_cancelled'); + + const closing = appender.shutdown({ signal: cancellation.signal }); + await sleepStarted.promise; + cancellation.abort(); + await closing; + sleep.resolve(); + + expect(attempts).toBe(1); + expect(readFirstFailedEvent(homeDir)).toMatchObject({ + event: 'retry_sleep_cancelled', + }); + }); + it('retries a 401 once without the Authorization header', async () => { const seenAuths: (string | undefined)[] = []; const appender = new CloudAppender( @@ -534,6 +633,105 @@ describe('CloudAppender', () => { expect(listFailedSpoolFiles(homeDir)).toHaveLength(0); }); + it('flush joins startup replay even when no live events are buffered', async () => { + const failingAppender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => statusResponse(500)), + }), + ); + failingAppender.track('persisted_before_flush'); + await failingAppender.flush(); + + const replayStarted = deferred(); + const replayResponse = deferred(); + const restartedAppender = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => { + replayStarted.resolve(); + return replayResponse.promise; + }), + }), + ); + restartedAppender.start(); + + const settled = vi.fn(); + const flushing = restartedAppender.flush().then(settled); + await replayStarted.promise; + await Promise.resolve(); + + expect(settled).not.toHaveBeenCalled(); + + replayResponse.resolve(okResponse()); + await flushing; + expect(listFailedSpoolFiles(homeDir)).toHaveLength(0); + }); + + it('startup replay limits the number of recovered files in one lifecycle', async () => { + let now = 1; + const failingAppender = new CloudAppender( + baseOptions({ + homeDir, + now: () => now++, + spoolMaxFiles: 10, + fetchImpl: makeFetch(() => statusResponse(500)), + }), + ); + for (const event of ['first', 'second', 'third']) { + failingAppender.track(event); + await failingAppender.flush(); + } + + const replayed: string[] = []; + const restartedAppender = new CloudAppender( + baseOptions({ + homeDir, + now: () => now, + replayMaxFiles: 1, + spoolMaxFiles: 10, + fetchImpl: makeFetch((request) => { + replayed.push(String(request.body.events[0]?.['event'])); + return okResponse(); + }), + }), + ); + + restartedAppender.start(); + restartedAppender.track('live_after_backlog'); + await restartedAppender.shutdown(); + + expect(replayed).toEqual(['kfc_first']); + expect(listFailedSpoolFiles(homeDir)).toHaveLength(3); + expect(readAllFailedEvents(homeDir).map((event) => event['event'])).toContain( + 'live_after_backlog', + ); + }); + + it('caps the durable spool by file count', async () => { + let now = 1; + const appender = new CloudAppender( + baseOptions({ + homeDir, + now: () => now++, + spoolMaxFiles: 2, + fetchImpl: makeFetch(() => statusResponse(500)), + }), + ); + + for (const event of ['first', 'second', 'third']) { + appender.track(event); + await appender.flush(); + } + + expect(listFailedSpoolFiles(homeDir)).toHaveLength(2); + expect( + readAllFailedEvents(homeDir) + .map((event) => event['event']) + .toSorted(), + ).toEqual(['first', 'second', 'third']); + }); + it('start leaves legacy telemetry spool files for their owning pipeline', async () => { const telemetryDir = join(homeDir, 'telemetry'); mkdirSync(telemetryDir, { recursive: true }); diff --git a/packages/agent-core-v2/test/app/telemetry/stubs.ts b/packages/agent-core-v2/test/app/telemetry/stubs.ts index f2534b7309..f90477d825 100644 --- a/packages/agent-core-v2/test/app/telemetry/stubs.ts +++ b/packages/agent-core-v2/test/app/telemetry/stubs.ts @@ -43,9 +43,9 @@ export function recordingTelemetry( setContext(patch: TelemetryContextPatch) { currentContext = { ...currentContext, ...patch }; }, - addAppender: () => ({ dispose: () => {} }), - removeAppender: () => {}, - setAppender: () => {}, + addAppender: () => ({ dispose: () => {}, shutdown: async () => {} }), + removeAppender: async () => {}, + setAppender: async () => {}, setEnabled(next) { enabled = next; }, diff --git a/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts b/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts index 3a153d5557..bb8403f7c6 100644 --- a/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts +++ b/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts @@ -1,7 +1,8 @@ /** - * Telemetry facade tests — exercise appender fan-out, context views, error - * isolation, lifecycle-option forwarding, and App-scope registration through - * the public `ITelemetryService` surface. + * `telemetry` domain (L1) — telemetry facade lifecycle coverage. + * + * Exercises appender fan-out, context views, error isolation, lifecycle-option + * forwarding, and App-scope registration through `ITelemetryService`. */ import { afterEach, beforeEach, describe, expect, it } from 'vitest'; @@ -45,7 +46,7 @@ function telemetryWithAppenders(...appenders: ITelemetryAppender[]): TelemetrySe const svc = new TelemetryService(); const [first, ...rest] = appenders; if (first !== undefined) { - svc.setAppender(first); + void svc.setAppender(first); } for (const appender of rest) { svc.addAppender(appender); @@ -53,6 +54,14 @@ function telemetryWithAppenders(...appenders: ITelemetryAppender[]): TelemetrySe return svc; } +function deferred(): { readonly promise: Promise; readonly resolve: () => void } { + let resolve!: () => void; + const promise = new Promise((done) => { + resolve = done; + }); + return { promise, resolve }; +} + describe('TelemetryService (unit)', () => { it('noop by default — does not throw', () => { const svc = new TelemetryService(); @@ -62,7 +71,7 @@ describe('TelemetryService (unit)', () => { it('merges bound context into tracked properties', () => { const appender = new CapturingAppender(); const svc = new TelemetryService(); - svc.setAppender(appender); + void svc.setAppender(appender); svc.setContext({ sessionId: 's1' }); svc.track('turn.start', { agentId: 'main' }); expect(appender.events[0]).toEqual({ @@ -74,7 +83,7 @@ describe('TelemetryService (unit)', () => { it('withContext merges context and shares the appender', () => { const appender = new CapturingAppender(); const root = new TelemetryService(); - root.setAppender(appender); + void root.setAppender(appender); root.setContext({ sessionId: 's1' }); const child = root.withContext({ agentId: 'main', turnId: 't1' }); child.track('tool.call', { name: 'bash' }); @@ -89,7 +98,7 @@ describe('TelemetryService (unit)', () => { it('per-call properties override bound context on key collision', () => { const appender = new CapturingAppender(); const svc = new TelemetryService(); - svc.setAppender(appender); + void svc.setAppender(appender); svc.setContext({ sessionId: 's1' }); svc.track('evt', { sessionId: 'override' }); expect(appender.events[0]?.properties?.['sessionId']).toBe('override'); @@ -118,6 +127,18 @@ describe('TelemetryService (unit)', () => { expect(b.events).toHaveLength(1); }); + it('an owning registration durably retires its appender', async () => { + const appender = new CapturingAppender(); + const svc = new TelemetryService(); + const registration = svc.addAppender(appender); + + await registration.shutdown(); + svc.track('after_retirement'); + + expect(appender.shutdownCalls).toBe(1); + expect(appender.events).toHaveLength(0); + }); + it('addAppender starts the registered appender', () => { const appender = new CapturingAppender(); const svc = new TelemetryService(); @@ -131,7 +152,7 @@ describe('TelemetryService (unit)', () => { const a = new CapturingAppender(); const b = new CapturingAppender(); const svc = telemetryWithAppenders(a, b); - svc.removeAppender(a); + void svc.removeAppender(a); svc.track('evt'); expect(a.events).toHaveLength(0); expect(b.events).toHaveLength(1); @@ -141,11 +162,22 @@ describe('TelemetryService (unit)', () => { const appender = new CapturingAppender(); const svc = new TelemetryService(); - svc.setAppender(appender); + void svc.setAppender(appender); expect(appender.startCalls).toBe(1); }); + it('setAppender durably retires the appender it replaces', async () => { + const previous = new CapturingAppender(); + const replacement = new CapturingAppender(); + const svc = telemetryWithAppenders(previous); + + await svc.setAppender(replacement); + + expect(previous.shutdownCalls).toBe(1); + expect(replacement.startCalls).toBe(1); + }); + it('setEnabled(false) drops track; setEnabled(true) resumes', () => { const appender = new CapturingAppender(); const svc = telemetryWithAppenders(appender); @@ -176,7 +208,7 @@ describe('TelemetryService (unit)', () => { const child = root.withContext({ agent_id: 'main' }); const appender = new CapturingAppender(); - root.setAppender(appender); + void root.setAppender(appender); child.track('sent'); expect(appender.events).toEqual([{ event: 'sent', properties: { agent_id: 'main' } }]); @@ -215,6 +247,47 @@ describe('TelemetryService (unit)', () => { expect(second.shutdownOptions).toBe(options); }); + it('repeated shutdown shares and awaits retirement started by disposal', async () => { + const retired = deferred(); + let shutdownCalls = 0; + const appender: ITelemetryAppender = { + track() {}, + shutdown() { + shutdownCalls += 1; + return retired.promise; + }, + }; + const svc = new TelemetryService(); + const registration = svc.addAppender(appender); + registration.dispose(); + + const first = svc.shutdown(); + const second = svc.shutdown(); + + expect(second).toBe(first); + expect(shutdownCalls).toBe(0); + await Promise.resolve(); + expect(shutdownCalls).toBe(1); + + retired.resolve(); + await first; + }); + + it('rejects appenders registered after shutdown begins', async () => { + const svc = new TelemetryService(); + await svc.shutdown(); + const appender = new CapturingAppender(); + + expect(() => { + svc.addAppender(appender); + }).toThrow('Telemetry service has already shut down'); + await expect(svc.setAppender(appender)).rejects.toThrow( + 'Telemetry service has already shut down', + ); + + expect(appender.startCalls).toBe(0); + }); + it('flush is a no-op for appenders without flush', async () => { const minimal: ITelemetryAppender = { track() {} }; const svc = telemetryWithAppenders(minimal); @@ -269,6 +342,21 @@ describe('TelemetryService (error isolation)', () => { expect(good.flushCalls).toBe(1); }); + it('flush tolerates a synchronously throwing appender and still flushes the rest', async () => { + const bad: ITelemetryAppender = { + track() {}, + flush() { + throw new Error('boom'); + }, + }; + const good = new CapturingAppender(); + const svc = telemetryWithAppenders(bad, good); + + await expect(svc.flush()).resolves.toBeUndefined(); + + expect(good.flushCalls).toBe(1); + }); + it('shutdown tolerates a rejecting appender and still shuts down the rest', async () => { const bad: ITelemetryAppender = { track() {}, @@ -281,6 +369,21 @@ describe('TelemetryService (error isolation)', () => { await expect(svc.shutdown()).resolves.toBeUndefined(); expect(good.shutdownCalls).toBe(1); }); + + it('shutdown tolerates a synchronously throwing appender and still shuts down the rest', async () => { + const bad: ITelemetryAppender = { + track() {}, + shutdown() { + throw new Error('boom'); + }, + }; + const good = new CapturingAppender(); + const svc = telemetryWithAppenders(bad, good); + + await expect(svc.shutdown()).resolves.toBeUndefined(); + + expect(good.shutdownCalls).toBe(1); + }); }); describe('ITelemetryService (scoped)', () => { diff --git a/packages/agent-core-v2/test/os/backends/node-local/tools/glob.test.ts b/packages/agent-core-v2/test/os/backends/node-local/tools/glob.test.ts index 2ae88c6f41..c38c261fdd 100644 --- a/packages/agent-core-v2/test/os/backends/node-local/tools/glob.test.ts +++ b/packages/agent-core-v2/test/os/backends/node-local/tools/glob.test.ts @@ -161,9 +161,9 @@ function telemetryStub( }, withContext: () => telemetryStub(events), setContext: () => {}, - addAppender: () => ({ dispose: () => {} }), - removeAppender: () => {}, - setAppender: () => {}, + addAppender: () => ({ dispose: () => {}, shutdown: async () => {} }), + removeAppender: async () => {}, + setAppender: async () => {}, setEnabled: () => {}, flush: async () => {}, shutdown: async () => {}, diff --git a/packages/agent-core-v2/test/session/sessionFs/fsService.test.ts b/packages/agent-core-v2/test/session/sessionFs/fsService.test.ts index d53a3bc47f..c14bbe7b92 100644 --- a/packages/agent-core-v2/test/session/sessionFs/fsService.test.ts +++ b/packages/agent-core-v2/test/session/sessionFs/fsService.test.ts @@ -287,9 +287,9 @@ function telemetryStub(events: Array<{ event: string; properties: Record telemetryStub(events), setContext: () => {}, - addAppender: () => ({ dispose: () => {} }), - removeAppender: () => {}, - setAppender: () => {}, + addAppender: () => ({ dispose: () => {}, shutdown: async () => {} }), + removeAppender: async () => {}, + setAppender: async () => {}, setEnabled: () => {}, flush: async () => {}, shutdown: async () => {}, diff --git a/packages/kap-server/src/services/telemetry.ts b/packages/kap-server/src/services/telemetry.ts index 85efc3899d..8ca8826352 100644 --- a/packages/kap-server/src/services/telemetry.ts +++ b/packages/kap-server/src/services/telemetry.ts @@ -9,11 +9,10 @@ */ import { - type CloudAppender, createCloudAppender, IBootstrapService, IConfigService, - type IDisposable, + type ITelemetryAppenderRegistration, IOAuthToolkit, ITelemetryService, type Scope, @@ -32,11 +31,10 @@ const TELEMETRY_DISABLE_ENV_VALUES = new Set(['1', 'true', 't', 'yes', 'y']); * telemetry endpoint must not hold shutdown hostage. */ const TELEMETRY_SHUTDOWN_TIMEOUT_MS = 3_000; +const TELEMETRY_SHUTDOWN_HARD_CAP_GRACE_MS = 1_000; export interface ServerTelemetry { - /** Present only when telemetry is enabled by both config and environment. */ - readonly appender?: CloudAppender; - readonly registration?: IDisposable; + readonly registration?: ITelemetryAppenderRegistration; } function isTelemetryDisabledByEnv(core: Scope): boolean { @@ -63,21 +61,23 @@ export async function initializeServerTelemetry( getAccessToken: async () => (await auth.getCachedAccessToken()) ?? null, }); const registration = service.addAppender(appender); - return { appender, registration }; + return { registration }; } export async function shutdownServerTelemetry( telemetry: ServerTelemetry, deadlineMs = Date.now() + TELEMETRY_SHUTDOWN_TIMEOUT_MS, ): Promise { - telemetry.registration?.dispose(); - if (telemetry.appender === undefined) return; + if (telemetry.registration === undefined) return; let timer: ReturnType | undefined; try { await Promise.race([ - telemetry.appender.shutdown({ deadlineMs }), + telemetry.registration.shutdown({ deadlineMs }), new Promise((resolve) => { - timer = setTimeout(resolve, Math.max(0, deadlineMs - Date.now())); + timer = setTimeout( + resolve, + Math.max(0, deadlineMs - Date.now()) + TELEMETRY_SHUTDOWN_HARD_CAP_GRACE_MS, + ); }), ]); } finally { diff --git a/packages/kap-server/test/boot.test.ts b/packages/kap-server/test/boot.test.ts index e68800fe12..3c0a450d65 100644 --- a/packages/kap-server/test/boot.test.ts +++ b/packages/kap-server/test/boot.test.ts @@ -200,7 +200,7 @@ describe('server-v2 boot', () => { const storage = new InMemoryStorageService(); const write = storage.write.bind(storage); vi.spyOn(storage, 'write').mockImplementation(async (scope, key, data, options) => { - if (scope === 'telemetry') throw new Error('telemetry storage unavailable'); + if (scope === 'telemetry-v2') throw new Error('telemetry storage unavailable'); await write(scope, key, data, options); }); const auth = { diff --git a/packages/kap-server/test/telemetry.test.ts b/packages/kap-server/test/telemetry.test.ts index 82cdaec70f..86cce4d035 100644 --- a/packages/kap-server/test/telemetry.test.ts +++ b/packages/kap-server/test/telemetry.test.ts @@ -11,6 +11,7 @@ import { join } from 'node:path'; import { bootstrap, + IFileSystemStorageService, type ITelemetryAppender, ITelemetryService, IOAuthToolkit, @@ -74,7 +75,7 @@ describe('server telemetry', () => { it('attaches the cloud appender by default and persists the device id', async () => { const app = await bootCore(); const telemetry = await initializeServerTelemetry(app, home as string); - expect(telemetry.appender).toBeDefined(); + expect(telemetry.registration).toBeDefined(); expect(readKimiDeviceId(home as string)).not.toBeNull(); await shutdownServerTelemetry(telemetry); }); @@ -115,15 +116,17 @@ describe('server telemetry', () => { try { await expect(shutdownServerTelemetry(telemetry, Date.now())).resolves.toBeUndefined(); + const spool = await app.accessor.get(IFileSystemStorageService).list('telemetry-v2'); + expect(spool).toHaveLength(1); } finally { - await telemetry.appender?.shutdown(); + await telemetry.registration?.shutdown(); } }); it('keeps the null appender when config sets telemetry = false', async () => { const app = await bootCore('telemetry = false\n'); const telemetry = await initializeServerTelemetry(app, home as string); - expect(telemetry.appender).toBeUndefined(); + expect(telemetry.registration).toBeUndefined(); await shutdownServerTelemetry(telemetry); }); @@ -135,7 +138,7 @@ describe('server telemetry', () => { KIMI_DISABLE_TELEMETRY: value, }); const telemetry = await initializeServerTelemetry(app, home as string); - expect(telemetry.appender).toBeUndefined(); + expect(telemetry.registration).toBeUndefined(); expect(readKimiDeviceId(home as string)).toBeNull(); await shutdownServerTelemetry(telemetry); }, From ee132d60cd0094122f8d7316d8c40d3e65f6f5b3 Mon Sep 17 00:00:00 2001 From: 7Sageer <7sageer@djwcb.cn> Date: Tue, 28 Jul 2026 20:52:41 +0800 Subject: [PATCH 3/7] fix(telemetry): preserve shutdown budgets and spool bounds --- .../src/app/telemetry/telemetryService.ts | 23 +++-- .../src/app/telemetry/telemetrySpoolStore.ts | 79 ++++++++++------ .../test/app/telemetry/cloudAppender.test.ts | 89 +++++++++++++++++-- .../app/telemetry/telemetryService.test.ts | 48 ++++++++++ 4 files changed, 198 insertions(+), 41 deletions(-) diff --git a/packages/agent-core-v2/src/app/telemetry/telemetryService.ts b/packages/agent-core-v2/src/app/telemetry/telemetryService.ts index 66ef6b556b..faa42bc63c 100644 --- a/packages/agent-core-v2/src/app/telemetry/telemetryService.ts +++ b/packages/agent-core-v2/src/app/telemetry/telemetryService.ts @@ -85,10 +85,10 @@ export class TelemetryService implements ITelemetryService { options?: TelemetryShutdownOptions, ): Promise { const retirement = this.retirements.get(appender); - if (retirement !== undefined) return retirement; + if (retirement !== undefined) return this.retireAppender(appender, options); if (!this.appenders.includes(appender)) return Promise.resolve(); this.appenders = this.appenders.filter((a) => a !== appender); - return this.removeDetachedAppender(appender, options); + return this.retireAppender(appender, options); } async setAppender( @@ -100,7 +100,7 @@ export class TelemetryService implements ITelemetryService { if (!this.registrations.has(appender)) this.createRegistration(appender); const previous = this.appenders.filter((candidate) => candidate !== appender); this.appenders = [appender]; - await Promise.all(previous.map((candidate) => this.removeDetachedAppender(candidate, options))); + await Promise.all(previous.map((candidate) => this.retireAppender(candidate, options))); } setEnabled(enabled: boolean): void { @@ -115,15 +115,15 @@ export class TelemetryService implements ITelemetryService { shutdown(options?: TelemetryShutdownOptions): Promise { if (this.shutdownPromise === null) { - const appenders = this.appenders; + const appenders = new Set([...this.appenders, ...this.retirements.keys()]); this.appenders = []; for (const appender of appenders) { - void this.removeDetachedAppender(appender, options); + void this.retireAppender(appender, options); } this.shutdownPromise = Promise.all(this.retirements.values()).then(() => undefined); } else if (options !== undefined) { - for (const appender of this.registrations.keys()) { - void this.invokeAppender(() => appender.shutdown?.(options)); + for (const appender of this.retirements.keys()) { + void this.retireAppender(appender, options); } } return this.shutdownPromise; @@ -154,12 +154,17 @@ export class TelemetryService implements ITelemetryService { } } - private removeDetachedAppender( + private retireAppender( appender: ITelemetryAppender, options?: TelemetryShutdownOptions, ): Promise { const retirement = this.retirements.get(appender); - if (retirement !== undefined) return retirement; + if (retirement !== undefined) { + if (options !== undefined) { + void this.invokeAppender(() => appender.shutdown?.(options)); + } + return retirement; + } const pending = this.invokeAppender(() => appender.shutdown?.(options)); this.retirements.set(appender, pending); return pending; diff --git a/packages/agent-core-v2/src/app/telemetry/telemetrySpoolStore.ts b/packages/agent-core-v2/src/app/telemetry/telemetrySpoolStore.ts index c8bfdce611..c2fbb64480 100644 --- a/packages/agent-core-v2/src/app/telemetry/telemetrySpoolStore.ts +++ b/packages/agent-core-v2/src/app/telemetry/telemetrySpoolStore.ts @@ -26,11 +26,13 @@ export interface ITelemetrySpoolStore { export interface TelemetrySpoolStoreOptions { readonly storage: IFileSystemStorageService; readonly maxFiles?: number; + readonly maxFileBytes?: number; readonly now?: () => number; } export const TELEMETRY_SPOOL_MAX_AGE_MS = 7 * 24 * 60 * 60 * 1000; export const TELEMETRY_SPOOL_MAX_FILES = 256; +export const TELEMETRY_SPOOL_MAX_FILE_BYTES = 256 * 1024; const TELEMETRY_SCOPE = 'telemetry-v2'; const FAILED_PREFIX = 'failed_'; @@ -41,12 +43,17 @@ const textDecoder = new TextDecoder(); export class TelemetrySpoolStore implements ITelemetrySpoolStore { private readonly storage: IFileSystemStorageService; private readonly maxFiles: number; + private readonly maxFileBytes: number; private readonly now: () => number; private writes: Promise = Promise.resolve(); constructor(options: TelemetrySpoolStoreOptions) { this.storage = options.storage; this.maxFiles = Math.max(1, Math.floor(options.maxFiles ?? TELEMETRY_SPOOL_MAX_FILES)); + this.maxFileBytes = Math.max( + 1, + Math.floor(options.maxFileBytes ?? TELEMETRY_SPOOL_MAX_FILE_BYTES), + ); this.now = options.now ?? Date.now; } @@ -95,41 +102,51 @@ export class TelemetrySpoolStore implements ITelemetrySpoolStore { } } files.sort((left, right) => left.createdAt - right.createdAt || left.key.localeCompare(right.key)); + const overflow = files.splice(0, Math.max(0, files.length - this.maxFiles)); + discarded.push(...overflow.map((file) => file.key)); await Promise.all(discarded.map((key) => this.acknowledge(key).catch(() => undefined))); return files; } private async putOwned(events: readonly EnrichedCloudEvent[]): Promise { - const text = events.map((event) => JSON.stringify(event)).join('\n') + '\n'; - const files = await this.enforcePolicy().catch(() => []); - if (files.length < this.maxFiles) { - await this.storage.write(TELEMETRY_SCOPE, this.createKey(), textEncoder.encode(text), { - atomic: true, - }); - return; + const chunks = this.chunk(events).slice(-this.maxFiles); + const files = [...(await this.enforcePolicy())]; + for (const chunk of chunks) { + const file = this.createFile(); + await this.storage.write(TELEMETRY_SCOPE, file.key, chunk, { atomic: true }); + files.push(file); + const overflow = files.splice(0, Math.max(0, files.length - this.maxFiles)); + await Promise.all( + overflow.map((candidate) => this.acknowledge(candidate.key).catch(() => undefined)), + ); } + } - const compactedCount = files.length - this.maxFiles + 1; - const compacted = files.slice(0, compactedCount); - const existing = await Promise.all( - compacted.map(async (file) => { - const bytes = await this.storage.read(TELEMETRY_SCOPE, file.key); - if (bytes === undefined) return ''; - const jsonl = textDecoder.decode(bytes); - return jsonl.length === 0 || jsonl.endsWith('\n') ? jsonl : jsonl + '\n'; - }), - ); - await this.storage.write( - TELEMETRY_SCOPE, - this.createKey(), - textEncoder.encode(existing.join('') + text), - { atomic: true }, - ); - await Promise.all(compacted.map((file) => this.acknowledge(file.key).catch(() => undefined))); + private chunk(events: readonly EnrichedCloudEvent[]): Uint8Array[] { + const chunks: Uint8Array[] = []; + let lines: Uint8Array[] = []; + let byteLength = 0; + for (const event of events) { + const line = textEncoder.encode(`${JSON.stringify(event)}\n`); + if (line.byteLength > this.maxFileBytes) continue; + if (byteLength + line.byteLength > this.maxFileBytes) { + chunks.push(concatBytes(lines, byteLength)); + lines = []; + byteLength = 0; + } + lines.push(line); + byteLength += line.byteLength; + } + if (lines.length > 0) chunks.push(concatBytes(lines, byteLength)); + return chunks; } - private createKey(): string { - return `${FAILED_PREFIX}${this.now()}_${randomBytes(6).toString('hex')}${JSONL_SUFFIX}`; + private createFile(): SpoolFile { + const createdAt = this.now(); + return { + key: `${FAILED_PREFIX}${createdAt}_${randomBytes(6).toString('hex')}${JSONL_SUFFIX}`, + createdAt, + }; } private async readJsonl(key: string): Promise { @@ -162,3 +179,13 @@ function parseCreatedAt(key: string): number | undefined { const createdAt = Number(rest.slice(0, underscore)); return Number.isFinite(createdAt) ? createdAt : undefined; } + +function concatBytes(parts: readonly Uint8Array[], byteLength: number): Uint8Array { + const bytes = new Uint8Array(byteLength); + let offset = 0; + for (const part of parts) { + bytes.set(part, offset); + offset += part.byteLength; + } + return bytes; +} diff --git a/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts b/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts index a26d9ec841..897f5d764c 100644 --- a/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts +++ b/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts @@ -1,8 +1,9 @@ /** - * `telemetry` domain (L1) — cloud appender lifecycle integration coverage. + * `telemetry` domain (L1) — cloud delivery and durable spool coverage. * - * Exercises batching, durable shutdown, startup replay, privacy, and wire - * shape through the real appender, transport, spool, and file-storage stack. + * Exercises spool retention and capacity plus appender batching, durable + * shutdown, startup replay, privacy, and wire shape through the real transport + * and file-storage stack. */ import { getEventListeners } from 'node:events'; @@ -25,7 +26,7 @@ import { } from '#/_base/errors/unexpectedError'; import { FileStorageService } from '#/persistence/backends/node-fs/fileStorageService'; import { CloudAppender, type CloudAppenderOptions } from '#/app/telemetry/cloudAppender'; -import { CloudTransport } from '#/app/telemetry/cloudTransport'; +import { CloudTransport, type EnrichedCloudEvent } from '#/app/telemetry/cloudTransport'; import { TelemetrySpoolStore } from '#/app/telemetry/telemetrySpoolStore'; import { stubBootstrap } from '../bootstrap/stubs'; @@ -124,6 +125,82 @@ function readAllFailedEvents(homeDir: string): Record[] { ); } +function spoolEvent(event: string, padding = ''): EnrichedCloudEvent { + return { + event_id: `${event}-id`, + device_id: 'dev', + session_id: null, + event, + timestamp: 1, + properties: { padding }, + context: {}, + }; +} + +describe('TelemetrySpoolStore', () => { + let homeDir: string; + + beforeEach(() => { + homeDir = mkdtempSync(join(tmpdir(), 'telemetry-spool-store-')); + }); + + afterEach(() => { + rmSync(homeDir, { recursive: true, force: true }); + }); + + it('recovery expires an old event after a newer write reaches the file cap', async () => { + const dayMs = 24 * 60 * 60 * 1000; + let now = 0; + const store = new TelemetrySpoolStore({ + storage: new FileStorageService(homeDir), + maxFiles: 1, + now: () => now, + }); + + await store.put([spoolEvent('old')]); + now = 6 * dayMs; + await store.put([spoolEvent('new')]); + now = 8 * dayMs; + + const entries = await store.recoverable(10); + + expect(entries.flatMap((entry) => entry.events).map((event) => event.event)).toEqual([ + 'new', + ]); + }); + + it('put keeps only the newest chunk when one batch exceeds byte capacity', async () => { + const store = new TelemetrySpoolStore({ + storage: new FileStorageService(homeDir), + maxFiles: 1, + maxFileBytes: 512, + }); + + await store.put([ + spoolEvent('first', 'x'.repeat(300)), + spoolEvent('second', 'x'.repeat(300)), + ]); + + const entries = await store.recoverable(10); + + expect(entries.flatMap((entry) => entry.events).map((event) => event.event)).toEqual([ + 'second', + ]); + }); + + it('put drops an event when its JSONL record exceeds byte capacity', async () => { + const store = new TelemetrySpoolStore({ + storage: new FileStorageService(homeDir), + maxFiles: 1, + maxFileBytes: 128, + }); + + await store.put([spoolEvent('oversized', 'x'.repeat(300))]); + + await expect(store.recoverable(10)).resolves.toEqual([]); + }); +}); + describe('CloudAppender', () => { let homeDir: string; @@ -708,7 +785,7 @@ describe('CloudAppender', () => { ); }); - it('caps the durable spool by file count', async () => { + it('keeps the newest durable batches when the spool reaches its file cap', async () => { let now = 1; const appender = new CloudAppender( baseOptions({ @@ -729,7 +806,7 @@ describe('CloudAppender', () => { readAllFailedEvents(homeDir) .map((event) => event['event']) .toSorted(), - ).toEqual(['first', 'second', 'third']); + ).toEqual(['second', 'third']); }); it('start leaves legacy telemetry spool files for their owning pipeline', async () => { diff --git a/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts b/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts index bb8403f7c6..33d188480a 100644 --- a/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts +++ b/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts @@ -139,6 +139,30 @@ describe('TelemetryService (unit)', () => { expect(appender.events).toHaveLength(0); }); + it('registration shutdown tightens retirement started by disposal', async () => { + const retired = deferred(); + const shutdownOptions: Array = []; + const appender: ITelemetryAppender = { + track() {}, + shutdown(options) { + shutdownOptions.push(options); + return retired.promise; + }, + }; + const svc = new TelemetryService(); + const registration = svc.addAppender(appender); + const options = { deadlineMs: 42 }; + + registration.dispose(); + await Promise.resolve(); + const closing = registration.shutdown(options); + await Promise.resolve(); + + expect(shutdownOptions).toEqual([undefined, options]); + retired.resolve(); + await closing; + }); + it('addAppender starts the registered appender', () => { const appender = new CapturingAppender(); const svc = new TelemetryService(); @@ -273,6 +297,30 @@ describe('TelemetryService (unit)', () => { await first; }); + it('service shutdown forwards a lifecycle budget to retirement started by disposal', async () => { + const retired = deferred(); + const shutdownOptions: Array = []; + const appender: ITelemetryAppender = { + track() {}, + shutdown(options) { + shutdownOptions.push(options); + return retired.promise; + }, + }; + const svc = new TelemetryService(); + const registration = svc.addAppender(appender); + const options = { deadlineMs: 42 }; + + registration.dispose(); + await Promise.resolve(); + const closing = svc.shutdown(options); + await Promise.resolve(); + + expect(shutdownOptions).toEqual([undefined, options]); + retired.resolve(); + await closing; + }); + it('rejects appenders registered after shutdown begins', async () => { const svc = new TelemetryService(); await svc.shutdown(); From a290335ff282cec9b05c50334790cfb9f4675a12 Mon Sep 17 00:00:00 2001 From: 7Sageer <7sageer@djwcb.cn> Date: Tue, 28 Jul 2026 22:08:59 +0800 Subject: [PATCH 4/7] fix(telemetry): harden appender replacement lifecycle --- .../src/app/telemetry/cloudAppender.ts | 6 + .../src/app/telemetry/telemetry.ts | 1 + .../src/app/telemetry/telemetryService.ts | 32 +++++- .../test/app/telemetry/cloudAppender.test.ts | 85 +++++++++++++- .../app/telemetry/telemetryService.test.ts | 104 +++++++++++++++++- packages/node-sdk/src/sdk-rpc-client-v2.ts | 2 +- 6 files changed, 220 insertions(+), 10 deletions(-) diff --git a/packages/agent-core-v2/src/app/telemetry/cloudAppender.ts b/packages/agent-core-v2/src/app/telemetry/cloudAppender.ts index 9ea1e3f6a8..311999536b 100644 --- a/packages/agent-core-v2/src/app/telemetry/cloudAppender.ts +++ b/packages/agent-core-v2/src/app/telemetry/cloudAppender.ts @@ -214,7 +214,13 @@ export class CloudAppender implements ITelemetryAppender { this.flushTimer = null; } + recover(): Promise { + return this.retryDiskEvents(); + } + async retryDiskEvents(): Promise { + const replay = this.replayPromise; + if (replay !== null) await replay; this.replayPending = true; await this.ensureReplay(); } diff --git a/packages/agent-core-v2/src/app/telemetry/telemetry.ts b/packages/agent-core-v2/src/app/telemetry/telemetry.ts index 54a3bf0c5b..879314da65 100644 --- a/packages/agent-core-v2/src/app/telemetry/telemetry.ts +++ b/packages/agent-core-v2/src/app/telemetry/telemetry.ts @@ -31,6 +31,7 @@ export interface TelemetryShutdownOptions { export interface ITelemetryAppender { start?(): void; + recover?(): Promise | void; track(event: string, properties?: TelemetryProperties): void; withContext?(patch: TelemetryContextPatch): ITelemetryAppender; setContext?(patch: TelemetryContextPatch): void; diff --git a/packages/agent-core-v2/src/app/telemetry/telemetryService.ts b/packages/agent-core-v2/src/app/telemetry/telemetryService.ts index faa42bc63c..089dc8bf4d 100644 --- a/packages/agent-core-v2/src/app/telemetry/telemetryService.ts +++ b/packages/agent-core-v2/src/app/telemetry/telemetryService.ts @@ -35,6 +35,7 @@ export class TelemetryService implements ITelemetryService { ITelemetryAppenderRegistration >(); private readonly retirements = new Map>(); + private readonly retiredAppenders = new WeakSet(); private shutdownPromise: Promise | null = null; private context: TelemetryProperties = {}; private enabled = true; @@ -73,6 +74,7 @@ export class TelemetryService implements ITelemetryService { addAppender(appender: ITelemetryAppender): ITelemetryAppenderRegistration { this.assertOpen(); + this.assertReusable(appender); const existing = this.registrations.get(appender); if (existing !== undefined) return existing; this.startAppender(appender); @@ -96,11 +98,25 @@ export class TelemetryService implements ITelemetryService { options?: TelemetryShutdownOptions, ): Promise { this.assertOpen(); - if (!this.appenders.includes(appender)) this.startAppender(appender); + this.assertReusable(appender); + const active = this.appenders.includes(appender); if (!this.registrations.has(appender)) this.createRegistration(appender); const previous = this.appenders.filter((candidate) => candidate !== appender); this.appenders = [appender]; await Promise.all(previous.map((candidate) => this.retireAppender(candidate, options))); + if ( + this.shutdownPromise === null && + this.appenders.includes(appender) && + !this.retiredAppenders.has(appender) + ) { + if (active) { + if (previous.length > 0) { + await this.invokeAppender(() => appender.recover?.()); + } + } else { + this.startAppender(appender); + } + } } setEnabled(enabled: boolean): void { @@ -154,6 +170,12 @@ export class TelemetryService implements ITelemetryService { } } + private assertReusable(appender: ITelemetryAppender): void { + if (this.retiredAppenders.has(appender)) { + throw new BugIndicatingError('Telemetry appender has already shut down'); + } + } + private retireAppender( appender: ITelemetryAppender, options?: TelemetryShutdownOptions, @@ -165,8 +187,16 @@ export class TelemetryService implements ITelemetryService { } return retirement; } + if (this.retiredAppenders.has(appender)) return Promise.resolve(); + this.retiredAppenders.add(appender); + this.registrations.delete(appender); const pending = this.invokeAppender(() => appender.shutdown?.(options)); this.retirements.set(appender, pending); + void pending.then(() => { + if (this.retirements.get(appender) === pending) { + this.retirements.delete(appender); + } + }); return pending; } diff --git a/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts b/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts index 897f5d764c..3249311b97 100644 --- a/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts +++ b/packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts @@ -2,8 +2,8 @@ * `telemetry` domain (L1) — cloud delivery and durable spool coverage. * * Exercises spool retention and capacity plus appender batching, durable - * shutdown, startup replay, privacy, and wire shape through the real transport - * and file-storage stack. + * shutdown, startup and replacement-handoff replay, privacy, and wire shape + * through the real transport and file-storage stack. */ import { getEventListeners } from 'node:events'; @@ -27,6 +27,7 @@ import { import { FileStorageService } from '#/persistence/backends/node-fs/fileStorageService'; import { CloudAppender, type CloudAppenderOptions } from '#/app/telemetry/cloudAppender'; import { CloudTransport, type EnrichedCloudEvent } from '#/app/telemetry/cloudTransport'; +import { TelemetryService } from '#/app/telemetry/telemetryService'; import { TelemetrySpoolStore } from '#/app/telemetry/telemetrySpoolStore'; import { stubBootstrap } from '../bootstrap/stubs'; @@ -148,12 +149,12 @@ describe('TelemetrySpoolStore', () => { rmSync(homeDir, { recursive: true, force: true }); }); - it('recovery expires an old event after a newer write reaches the file cap', async () => { + it('recovery expires entries older than retention while keeping newer entries', async () => { const dayMs = 24 * 60 * 60 * 1000; let now = 0; const store = new TelemetrySpoolStore({ storage: new FileStorageService(homeDir), - maxFiles: 1, + maxFiles: 2, now: () => now, }); @@ -169,6 +170,25 @@ describe('TelemetrySpoolStore', () => { ]); }); + it('put keeps the newest entries when file capacity is exceeded', async () => { + let now = 0; + const store = new TelemetrySpoolStore({ + storage: new FileStorageService(homeDir), + maxFiles: 1, + now: () => now, + }); + + await store.put([spoolEvent('old')]); + now = 1; + await store.put([spoolEvent('new')]); + + const entries = await store.recoverable(10); + + expect(entries.flatMap((entry) => entry.events).map((event) => event.event)).toEqual([ + 'new', + ]); + }); + it('put keeps only the newest chunk when one batch exceeds byte capacity', async () => { const store = new TelemetrySpoolStore({ storage: new FileStorageService(homeDir), @@ -710,6 +730,63 @@ describe('CloudAppender', () => { expect(listFailedSpoolFiles(homeDir)).toHaveLength(0); }); + it('setAppender replays events persisted while the previous appender retires', async () => { + const previous = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => statusResponse(500)), + }), + ); + const replayed: string[] = []; + const replacement = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch((request) => { + replayed.push(String(request.body.events[0]?.['event'])); + return okResponse(); + }), + }), + ); + const service = new TelemetryService(); + await service.setAppender(previous); + service.track('handoff'); + + await service.setAppender(replacement); + await replacement.flush(); + + expect(replayed).toEqual(['kfc_handoff']); + expect(listFailedSpoolFiles(homeDir)).toHaveLength(0); + }); + + it('setAppender recovers durable events when the replacement is already active', async () => { + const previous = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch(() => statusResponse(500)), + }), + ); + const replayed: string[] = []; + const replacement = new CloudAppender( + baseOptions({ + homeDir, + fetchImpl: makeFetch((request) => { + replayed.push(String(request.body.events[0]?.['event'])); + return okResponse(); + }), + }), + ); + const service = new TelemetryService(); + await service.setAppender(previous); + service.addAppender(replacement); + await replacement.flush(); + previous.track('active_handoff'); + + await service.setAppender(replacement); + + expect(replayed).toEqual(['kfc_active_handoff']); + expect(listFailedSpoolFiles(homeDir)).toHaveLength(0); + }); + it('flush joins startup replay even when no live events are buffered', async () => { const failingAppender = new CloudAppender( baseOptions({ diff --git a/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts b/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts index 33d188480a..f3e4e8ea23 100644 --- a/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts +++ b/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts @@ -1,8 +1,9 @@ /** * `telemetry` domain (L1) — telemetry facade lifecycle coverage. * - * Exercises appender fan-out, context views, error isolation, lifecycle-option - * forwarding, and App-scope registration through `ITelemetryService`. + * Exercises appender fan-out, context views, error isolation, ordered + * replacement, terminal appender rejection, lifecycle-option forwarding, and + * App-scope registration through `ITelemetryService`. */ import { afterEach, beforeEach, describe, expect, it } from 'vitest'; @@ -182,11 +183,11 @@ describe('TelemetryService (unit)', () => { expect(b.events).toHaveLength(1); }); - it('setAppender starts the replacement appender', () => { + it('setAppender starts the replacement appender after retirement completes', async () => { const appender = new CapturingAppender(); const svc = new TelemetryService(); - void svc.setAppender(appender); + await svc.setAppender(appender); expect(appender.startCalls).toBe(1); }); @@ -202,6 +203,101 @@ describe('TelemetryService (unit)', () => { expect(replacement.startCalls).toBe(1); }); + it('setAppender starts the replacement only after the previous appender retires', async () => { + const retired = deferred(); + const lifecycle: string[] = []; + const previous: ITelemetryAppender = { + track() {}, + async shutdown() { + lifecycle.push('shutdown-started'); + await retired.promise; + lifecycle.push('shutdown-finished'); + }, + }; + const replacement: ITelemetryAppender = { + start() { + lifecycle.push('replacement-started'); + }, + track() {}, + }; + const svc = new TelemetryService(); + await svc.setAppender(previous); + + const replacing = svc.setAppender(replacement); + await Promise.resolve(); + + expect(lifecycle).toEqual(['shutdown-started']); + + retired.resolve(); + await replacing; + + expect(lifecycle).toEqual([ + 'shutdown-started', + 'shutdown-finished', + 'replacement-started', + ]); + }); + + it('setAppender recovers an active replacement after the previous appender retires', async () => { + const retired = deferred(); + const lifecycle: string[] = []; + const previous: ITelemetryAppender = { + track() {}, + async shutdown() { + lifecycle.push('shutdown-started'); + await retired.promise; + lifecycle.push('shutdown-finished'); + }, + }; + const replacement: ITelemetryAppender = { + start() {}, + recover() { + lifecycle.push('replacement-recovered'); + }, + track() {}, + }; + const svc = new TelemetryService(); + await svc.setAppender(previous); + svc.addAppender(replacement); + + const replacing = svc.setAppender(replacement); + await Promise.resolve(); + + expect(lifecycle).toEqual(['shutdown-started']); + + retired.resolve(); + await replacing; + + expect(lifecycle).toEqual([ + 'shutdown-started', + 'shutdown-finished', + 'replacement-recovered', + ]); + }); + + it('addAppender rejects an appender after its registration retires it', async () => { + const appender = new CapturingAppender(); + const svc = new TelemetryService(); + const registration = svc.addAppender(appender); + await registration.shutdown(); + + expect(() => svc.addAppender(appender)).toThrow( + 'Telemetry appender has already shut down', + ); + }); + + it('setAppender rejects an appender retired by a later replacement', async () => { + const retired = new CapturingAppender(); + const replacement = new CapturingAppender(); + const svc = new TelemetryService(); + await svc.setAppender(retired); + await svc.setAppender(replacement); + + await expect(svc.setAppender(retired)).rejects.toThrow( + 'Telemetry appender has already shut down', + ); + }); + it('setEnabled(false) drops track; setEnabled(true) resumes', () => { const appender = new CapturingAppender(); const svc = telemetryWithAppenders(appender); diff --git a/packages/node-sdk/src/sdk-rpc-client-v2.ts b/packages/node-sdk/src/sdk-rpc-client-v2.ts index 99b4cda7ed..2ae8f1b3ef 100644 --- a/packages/node-sdk/src/sdk-rpc-client-v2.ts +++ b/packages/node-sdk/src/sdk-rpc-client-v2.ts @@ -470,7 +470,7 @@ export class SDKRpcClientV2 extends SDKRpcClientBase { private installEngineTelemetry(client: TelemetryClient | undefined): void { if (client === undefined) return; const telemetry = this.app.accessor.get(ITelemetryService); - telemetry.setAppender(client); + void telemetry.setAppender(client); void this.configReady.then(() => { telemetry.setEnabled(this.engineAccessor.get(IConfigService).get('telemetry') !== false); }); From 7511438a57423be4822c915de200bdcb5542292c Mon Sep 17 00:00:00 2001 From: 7Sageer <7sageer@djwcb.cn> Date: Wed, 29 Jul 2026 00:45:01 +0800 Subject: [PATCH 5/7] docs(agent-core-v2): document telemetry lifecycle handoff --- .agents/skills/agent-core-dev/telemetry.md | 36 ++++++++++++++-------- 1 file changed, 23 insertions(+), 13 deletions(-) diff --git a/.agents/skills/agent-core-dev/telemetry.md b/.agents/skills/agent-core-dev/telemetry.md index 0d04dc5da4..0d8b8b6408 100644 --- a/.agents/skills/agent-core-dev/telemetry.md +++ b/.agents/skills/agent-core-dev/telemetry.md @@ -11,8 +11,9 @@ Telemetry is a **layer-1 root** domain (alongside `log`): the facade lives at `A - `src/app/telemetry/telemetryService.ts`: `TelemetryService` impl + `registerScopedService(LifecycleScope.App, …)`. - `src/app/telemetry/agentTelemetryContext.ts` + `agentTelemetryContextService.ts`: `IAgentTelemetryContextService` — Agent-scoped mutable request context (`mode` / `provider_type` / `protocol` / `turn_id` / `trace_id`) snapshot into turn telemetry at launch. Agent identity (`agent_id`) is not part of it — identity is bound by the Agent-scoped `ITelemetryService` view. - `src/app/telemetry/consoleAppender.ts`: `ConsoleAppender` — echoes events to a log function (dev / debug). -- `src/app/telemetry/cloudAppender.ts`: `CloudAppender` — sanitizes + PII-cleans properties, batches + enriches + posts to the telemetry endpoint. -- `src/app/telemetry/cloudTransport.ts`: `CloudTransport` — HTTP transport behind `CloudAppender`. +- `src/app/telemetry/cloudAppender.ts`: `CloudAppender` — sanitizes + PII-cleans properties, owns batching / replay / shutdown, and posts to the telemetry endpoint. +- `src/app/telemetry/cloudTransport.ts`: `CloudTransport` — cancellable HTTP transport behind `CloudAppender`. +- `src/app/telemetry/telemetrySpoolStore.ts`: `TelemetrySpoolStore` — bounded durable storage and recovery for unsent v2 batches under the isolated `telemetry-v2` namespace. - `src/app/telemetry/privacy.ts`: outbound PII redaction (`cleanTelemetryProperties`) — URLs, emails, tokens, and absolute file paths become `` labels; `node_modules/` tails are kept. ## Emitting events (business services) @@ -48,18 +49,20 @@ An appender is the destination an event is fanned out to. It is **not a DI Servi ```ts export interface ITelemetryAppender { + start?(): void; + recover?(): Promise | void; track(event: string, properties?: TelemetryProperties): void; withContext?(patch: TelemetryContextPatch): ITelemetryAppender; setContext?(patch: TelemetryContextPatch): void; flush?(): Promise | void; - shutdown?(): Promise | void; + shutdown?(options?: TelemetryShutdownOptions): Promise | void; } ``` Built-in appenders: - `ConsoleAppender` — `[telemetry] ` to a log function (default `console.log`); options `prefix` / `pretty` / `log`. -- `CloudAppender` — batches events, enriches with common context (`app_name` / `version` / `platform` / …), and posts to `https://telemetry-logs.kimi.com/v1/event` through `CloudTransport` (Bearer auth, retry, on-disk fallback). Options: `homeDir` / `deviceId` / `sessionId?` / `appName` / `version` / `uiMode?` / `model?` / `getAccessToken?` / `endpoint?` / `flushThreshold?` / `flushIntervalMs?`. +- `CloudAppender` — batches events, enriches with common context (`app_name` / `version` / `platform` / …), and posts to `https://telemetry-logs.kimi.com/v1/event` through `CloudTransport` (Bearer auth + cancellable retry). `createCloudAppender(accessor, hostOptions)` supplies bootstrap facts and the v2 durable spool; host options include `deviceId` / `sessionId?` / `appName` / `uiMode?` / `model?` / `buildSha?` / `getAccessToken?`. ### Registering appenders (bootstrap) @@ -70,21 +73,28 @@ const app = createAppScope(); const telemetry = app.accessor.get(ITelemetryService); telemetry.addAppender(new ConsoleAppender({ prefix: '[dev]' })); // dev echo -telemetry.addAppender(new CloudAppender({ // production - homeDir, deviceId, sessionId, - appName: 'kimi-code', version, uiMode: 'shell', model, - getAccessToken: () => auth.getCachedAccessToken(KIMI_CODE_PROVIDER_NAME), -})); +const registration = telemetry.addAppender(createCloudAppender( // production + app.accessor, + { + deviceId, sessionId, + appName: 'kimi-code', uiMode: 'shell', model, + getAccessToken: () => auth.getCachedAccessToken(KIMI_CODE_PROVIDER_NAME), + }, +)); ``` -`addAppender` returns an `IDisposable` that removes the appender when disposed. `setAppender(appender)` resets to a single appender (mainly for tests). `removeAppender(appender)` drops one. +`addAppender` starts the appender and returns an `ITelemetryAppenderRegistration`. `registration.dispose()` begins best-effort retirement without waiting; use `await registration.shutdown(options)` when the host must wait until the appender is durably retired. `removeAppender(appender, options)` has the same awaited retirement contract. -> There is no production bootstrap wired yet — `TelemetryService` defaults to `[nullTelemetryAppender]`, so `track(...)` is a no-op until `addAppender` is called at startup. +`await setAppender(appender, options)` switches fan-out to one appender, durably retires every previous appender, then starts the replacement. If the replacement was already active, its optional `recover()` hook runs after retirement so it can pick up batches handed to the shared spool during the transition. Shutdown is terminal: a retired appender object cannot be registered again. + +`TelemetryService` defaults to `[nullTelemetryAppender]`, so each host owns registration. Kap-server and the experimental v2 CLI register a `CloudAppender` when telemetry is enabled. ## Lifecycle - `setEnabled(false)` drops `track` (service-level switch); `setEnabled(true)` resumes. `flush` / `shutdown` are unaffected by the switch. -- `flush()` / `shutdown()` fan out to all appenders concurrently; a single rejecting appender is swallowed. Await `shutdown()` before process exit so buffered events (e.g. in `CloudAppender`) are sent. +- `flush()` / `shutdown(options)` fan out to all appenders concurrently; a single rejecting appender is swallowed. Later shutdown calls may tighten an active lifecycle with an `AbortSignal` or earlier `deadlineMs`. +- `CloudAppender` owns one drain across threshold, periodic, explicit-flush, and shutdown triggers. Shutdown rejects new events, cancels network / retry waits at the lifecycle boundary, and does not complete until each unsent batch is delivered or handed to the v2 spool. Startup replay runs before newly buffered events. +- Await `registration.shutdown(options)` or `telemetry.shutdown(options)` before process exit when a buffering appender is registered. ## Red lines (this topic) @@ -92,6 +102,6 @@ telemetry.addAppender(new CloudAppender({ // production - Telemetry is layer-1 root: do not inject any business-domain service into it, and keep the facade at `App` scope (only the ambient context service binds at `Agent`). - Appenders are plain `ITelemetryAppender` objects, not DI Services — register them with `addAppender`, never via `registerScopedService`. - `track` is fire-and-forget and must not throw; appender `track` must be synchronous — buffer and send asynchronously via `flush` / `shutdown`. -- Await `telemetry.shutdown()` before process exit when a buffering appender is registered. +- Await the registration or service shutdown contract before process exit when a buffering appender is registered; do not rely on `dispose()` when durable handoff must complete. - Keep event names stable; register every business event in `events.ts` and emit via `track2` — properties must be JSON-serializable primitives (non-primitives are dropped with a warning by `CloudAppender`). - Agent identity is ambient: agent-scope events go through `defineAgentTelemetryEvent` and get `agent_id` from the scoped telemetry view — do not pass `agent_id` at business call sites (per-event identities such as `subagent_created` and the cron events are the exception). From 25ac4bd9992e0b9e34a5a6bd36f1b7c38a8ea587 Mon Sep 17 00:00:00 2001 From: 7Sageer <7sageer@djwcb.cn> Date: Wed, 29 Jul 2026 01:22:19 +0800 Subject: [PATCH 6/7] fix(telemetry): stage events during appender replacement --- .agents/skills/agent-core-dev/telemetry.md | 2 +- .../src/app/telemetry/telemetryService.ts | 110 +++++++++++++----- .../app/telemetry/telemetryService.test.ts | 101 +++++++++++++--- 3 files changed, 171 insertions(+), 42 deletions(-) diff --git a/.agents/skills/agent-core-dev/telemetry.md b/.agents/skills/agent-core-dev/telemetry.md index 0d8b8b6408..a4ac648add 100644 --- a/.agents/skills/agent-core-dev/telemetry.md +++ b/.agents/skills/agent-core-dev/telemetry.md @@ -85,7 +85,7 @@ const registration = telemetry.addAppender(createCloudAppender( // production `addAppender` starts the appender and returns an `ITelemetryAppenderRegistration`. `registration.dispose()` begins best-effort retirement without waiting; use `await registration.shutdown(options)` when the host must wait until the appender is durably retired. `removeAppender(appender, options)` has the same awaited retirement contract. -`await setAppender(appender, options)` switches fan-out to one appender, durably retires every previous appender, then starts the replacement. If the replacement was already active, its optional `recover()` hook runs after retirement so it can pick up batches handed to the shared spool during the transition. Shutdown is terminal: a retired appender object cannot be registered again. +`await setAppender(appender, options)` stages new events, durably retires every previous appender, then starts the replacement before releasing those events to it. `flush()` and `shutdown()` join that transition. If the replacement was already active, its optional `recover()` hook runs after retirement so it can pick up batches handed to the shared spool before staged events resume. Shutdown is terminal: a retired appender object cannot be registered again. `TelemetryService` defaults to `[nullTelemetryAppender]`, so each host owns registration. Kap-server and the experimental v2 CLI register a `CloudAppender` when telemetry is enabled. diff --git a/packages/agent-core-v2/src/app/telemetry/telemetryService.ts b/packages/agent-core-v2/src/app/telemetry/telemetryService.ts index 089dc8bf4d..96d761ff5e 100644 --- a/packages/agent-core-v2/src/app/telemetry/telemetryService.ts +++ b/packages/agent-core-v2/src/app/telemetry/telemetryService.ts @@ -26,6 +26,11 @@ import { type TelemetryShutdownOptions, } from './telemetry'; +interface PendingTelemetryEvent { + readonly event: string; + readonly properties: TelemetryProperties; +} + export class TelemetryService implements ITelemetryService { declare readonly _serviceBrand: undefined; @@ -36,21 +41,23 @@ export class TelemetryService implements ITelemetryService { >(); private readonly retirements = new Map>(); private readonly retiredAppenders = new WeakSet(); + private appenderTransition: Promise | null = null; + private pendingEvents: PendingTelemetryEvent[] | null = null; private shutdownPromise: Promise | null = null; private context: TelemetryProperties = {}; private enabled = true; track(event: string, properties?: TelemetryProperties): void { - if (!this.enabled) { + if (!this.enabled || this.shutdownPromise !== null) { return; } const merged = { ...this.context, ...properties }; + if (this.pendingEvents !== null) { + this.pendingEvents.push({ event, properties: merged }); + return; + } for (const appender of this.appenders) { - try { - appender.track(event, merged); - } catch (error) { - onUnexpectedError(error); - } + this.trackAppender(appender, event, merged); } } @@ -98,25 +105,44 @@ export class TelemetryService implements ITelemetryService { options?: TelemetryShutdownOptions, ): Promise { this.assertOpen(); + this.assertReusable(appender); + const previousTransition = this.appenderTransition; + const transition = (async () => { + if (previousTransition !== null) await previousTransition; + await this.replaceAppender(appender, options); + })(); + this.appenderTransition = transition; + try { + await transition; + } finally { + if (this.appenderTransition === transition) this.appenderTransition = null; + } + } + + private async replaceAppender( + appender: ITelemetryAppender, + options?: TelemetryShutdownOptions, + ): Promise { this.assertReusable(appender); const active = this.appenders.includes(appender); - if (!this.registrations.has(appender)) this.createRegistration(appender); const previous = this.appenders.filter((candidate) => candidate !== appender); - this.appenders = [appender]; + const pendingEvents: PendingTelemetryEvent[] = []; + this.pendingEvents = pendingEvents; + this.appenders = []; await Promise.all(previous.map((candidate) => this.retireAppender(candidate, options))); - if ( - this.shutdownPromise === null && - this.appenders.includes(appender) && - !this.retiredAppenders.has(appender) - ) { - if (active) { - if (previous.length > 0) { - await this.invokeAppender(() => appender.recover?.()); - } - } else { - this.startAppender(appender); + if (active) { + if (previous.length > 0) { + await this.invokeAppender(() => appender.recover?.()); } + } else { + this.startAppender(appender); } + if (!this.registrations.has(appender)) this.createRegistration(appender); + this.appenders = [appender]; + for (const pending of pendingEvents) { + this.trackAppender(appender, pending.event, pending.properties); + } + if (this.pendingEvents === pendingEvents) this.pendingEvents = null; } setEnabled(enabled: boolean): void { @@ -124,6 +150,8 @@ export class TelemetryService implements ITelemetryService { } async flush(): Promise { + const transition = this.appenderTransition; + if (transition !== null) await transition; await Promise.all( this.appenders.map((appender) => this.invokeAppender(() => appender.flush?.())), ); @@ -131,20 +159,46 @@ export class TelemetryService implements ITelemetryService { shutdown(options?: TelemetryShutdownOptions): Promise { if (this.shutdownPromise === null) { - const appenders = new Set([...this.appenders, ...this.retirements.keys()]); - this.appenders = []; - for (const appender of appenders) { - void this.retireAppender(appender, options); - } - this.shutdownPromise = Promise.all(this.retirements.values()).then(() => undefined); + this.shutdownPromise = this.shutdownAfterTransition(this.appenderTransition, options); } else if (options !== undefined) { - for (const appender of this.retirements.keys()) { - void this.retireAppender(appender, options); - } + const tighten = (): void => { + const appenders = new Set([...this.appenders, ...this.retirements.keys()]); + for (const appender of appenders) { + void this.retireAppender(appender, options); + } + }; + const transition = this.appenderTransition; + if (transition === null) tighten(); + else void transition.then(tighten); } return this.shutdownPromise; } + private async shutdownAfterTransition( + transition: Promise | null, + options?: TelemetryShutdownOptions, + ): Promise { + if (transition !== null) await transition; + const appenders = new Set([...this.appenders, ...this.retirements.keys()]); + this.appenders = []; + for (const appender of appenders) { + void this.retireAppender(appender, options); + } + await Promise.all(this.retirements.values()); + } + + private trackAppender( + appender: ITelemetryAppender, + event: string, + properties: TelemetryProperties, + ): void { + try { + appender.track(event, properties); + } catch (error) { + onUnexpectedError(error); + } + } + private startAppender(appender: ITelemetryAppender): void { try { appender.start?.(); diff --git a/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts b/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts index f3e4e8ea23..988bff318d 100644 --- a/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts +++ b/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts @@ -45,11 +45,7 @@ class CapturingAppender implements ITelemetryAppender { function telemetryWithAppenders(...appenders: ITelemetryAppender[]): TelemetryService { const svc = new TelemetryService(); - const [first, ...rest] = appenders; - if (first !== undefined) { - void svc.setAppender(first); - } - for (const appender of rest) { + for (const appender of appenders) { svc.addAppender(appender); } return svc; @@ -69,10 +65,10 @@ describe('TelemetryService (unit)', () => { expect(() => svc.track('evt', { a: 1 })).not.toThrow(); }); - it('merges bound context into tracked properties', () => { + it('merges bound context into tracked properties', async () => { const appender = new CapturingAppender(); const svc = new TelemetryService(); - void svc.setAppender(appender); + await svc.setAppender(appender); svc.setContext({ sessionId: 's1' }); svc.track('turn.start', { agentId: 'main' }); expect(appender.events[0]).toEqual({ @@ -81,10 +77,10 @@ describe('TelemetryService (unit)', () => { }); }); - it('withContext merges context and shares the appender', () => { + it('withContext merges context and shares the appender', async () => { const appender = new CapturingAppender(); const root = new TelemetryService(); - void root.setAppender(appender); + await root.setAppender(appender); root.setContext({ sessionId: 's1' }); const child = root.withContext({ agentId: 'main', turnId: 't1' }); child.track('tool.call', { name: 'bash' }); @@ -96,10 +92,10 @@ describe('TelemetryService (unit)', () => { }); }); - it('per-call properties override bound context on key collision', () => { + it('per-call properties override bound context on key collision', async () => { const appender = new CapturingAppender(); const svc = new TelemetryService(); - void svc.setAppender(appender); + await svc.setAppender(appender); svc.setContext({ sessionId: 's1' }); svc.track('evt', { sessionId: 'override' }); expect(appender.events[0]?.properties?.['sessionId']).toBe('override'); @@ -238,6 +234,50 @@ describe('TelemetryService (unit)', () => { ]); }); + it('setAppender holds replacement traffic until previous retirement completes', async () => { + const retired = deferred(); + const lifecycle: string[] = []; + const previous: ITelemetryAppender = { + track() {}, + async shutdown() { + lifecycle.push('shutdown-started'); + await retired.promise; + lifecycle.push('shutdown-finished'); + }, + }; + const replacement: ITelemetryAppender = { + start() { + lifecycle.push('replacement-started'); + }, + track(event) { + lifecycle.push(`replacement-tracked:${event}`); + }, + flush() { + lifecycle.push('replacement-flushed'); + }, + }; + const svc = new TelemetryService(); + await svc.setAppender(previous); + + const replacing = svc.setAppender(replacement); + svc.track('during-replacement'); + const flushing = svc.flush(); + await Promise.resolve(); + + expect(lifecycle).toEqual(['shutdown-started']); + + retired.resolve(); + await Promise.all([replacing, flushing]); + + expect(lifecycle).toEqual([ + 'shutdown-started', + 'shutdown-finished', + 'replacement-started', + 'replacement-tracked:during-replacement', + 'replacement-flushed', + ]); + }); + it('setAppender recovers an active replacement after the previous appender retires', async () => { const retired = deferred(); const lifecycle: string[] = []; @@ -323,12 +363,12 @@ describe('TelemetryService (unit)', () => { expect(appender.events).toEqual([{ event: 'sent', properties: { turnId: 't1' } }]); }); - it('withContext view follows root appender changes', () => { + it('withContext view follows root appender changes', async () => { const root = new TelemetryService(); const child = root.withContext({ agent_id: 'main' }); const appender = new CapturingAppender(); - void root.setAppender(appender); + await root.setAppender(appender); child.track('sent'); expect(appender.events).toEqual([{ event: 'sent', properties: { agent_id: 'main' } }]); @@ -352,6 +392,41 @@ describe('TelemetryService (unit)', () => { expect(b.shutdownCalls).toBe(1); }); + it('shutdown durably retires a replacement accepted during transition', async () => { + const retired = deferred(); + const lifecycle: string[] = []; + const previous: ITelemetryAppender = { + track() {}, + shutdown: () => retired.promise, + }; + const replacement: ITelemetryAppender = { + start() { + lifecycle.push('replacement-started'); + }, + track(event) { + lifecycle.push(`replacement-tracked:${event}`); + }, + shutdown() { + lifecycle.push('replacement-shutdown'); + }, + }; + const svc = new TelemetryService(); + await svc.setAppender(previous); + + const replacing = svc.setAppender(replacement); + svc.track('before-shutdown'); + const closing = svc.shutdown(); + retired.resolve(); + + await Promise.all([replacing, closing]); + + expect(lifecycle).toEqual([ + 'replacement-started', + 'replacement-tracked:before-shutdown', + 'replacement-shutdown', + ]); + }); + it('shutdown forwards one lifecycle budget to every appender', async () => { const first = new CapturingAppender(); const second = new CapturingAppender(); From 52997f5e953a2b0c976edb350b7056f1fa6637f2 Mon Sep 17 00:00:00 2001 From: 7Sageer <7sageer@djwcb.cn> Date: Wed, 29 Jul 2026 01:43:59 +0800 Subject: [PATCH 7/7] fix(telemetry): tighten replacement shutdown budgets --- .../src/app/telemetry/telemetryService.ts | 23 ++++++++------ .../app/telemetry/telemetryService.test.ts | 30 +++++++++++++++++++ 2 files changed, 44 insertions(+), 9 deletions(-) diff --git a/packages/agent-core-v2/src/app/telemetry/telemetryService.ts b/packages/agent-core-v2/src/app/telemetry/telemetryService.ts index 96d761ff5e..c6b55f5f7c 100644 --- a/packages/agent-core-v2/src/app/telemetry/telemetryService.ts +++ b/packages/agent-core-v2/src/app/telemetry/telemetryService.ts @@ -159,27 +159,32 @@ export class TelemetryService implements ITelemetryService { shutdown(options?: TelemetryShutdownOptions): Promise { if (this.shutdownPromise === null) { + if (options !== undefined) this.tightenRetirements(options); this.shutdownPromise = this.shutdownAfterTransition(this.appenderTransition, options); } else if (options !== undefined) { - const tighten = (): void => { - const appenders = new Set([...this.appenders, ...this.retirements.keys()]); - for (const appender of appenders) { - void this.retireAppender(appender, options); - } - }; + this.tightenRetirements(options); const transition = this.appenderTransition; - if (transition === null) tighten(); - else void transition.then(tighten); + if (transition !== null) { + void transition.then(() => { + this.tightenRetirements(options); + }); + } } return this.shutdownPromise; } + private tightenRetirements(options: TelemetryShutdownOptions): void { + for (const appender of this.retirements.keys()) { + void this.retireAppender(appender, options); + } + } + private async shutdownAfterTransition( transition: Promise | null, options?: TelemetryShutdownOptions, ): Promise { if (transition !== null) await transition; - const appenders = new Set([...this.appenders, ...this.retirements.keys()]); + const appenders = new Set(this.appenders); this.appenders = []; for (const appender of appenders) { void this.retireAppender(appender, options); diff --git a/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts b/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts index 988bff318d..912337bc64 100644 --- a/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts +++ b/packages/agent-core-v2/test/app/telemetry/telemetryService.test.ts @@ -427,6 +427,36 @@ describe('TelemetryService (unit)', () => { ]); }); + it('shutdown forwards its budget to retirement active during replacement', async () => { + const retired = deferred(); + const shutdownOptions: Array = []; + const previous: ITelemetryAppender = { + track() {}, + shutdown(options) { + shutdownOptions.push(options); + options?.signal?.addEventListener('abort', retired.resolve, { once: true }); + return retired.promise; + }, + }; + const replacement = new CapturingAppender(); + const controller = new AbortController(); + const options = { signal: controller.signal }; + const svc = new TelemetryService(); + await svc.setAppender(previous); + + const replacing = svc.setAppender(replacement); + await Promise.resolve(); + const closing = svc.shutdown(options); + await Promise.resolve(); + + expect(shutdownOptions).toEqual([undefined, options]); + + controller.abort(); + await Promise.all([replacing, closing]); + + expect(replacement.shutdownOptions).toBe(options); + }); + it('shutdown forwards one lifecycle budget to every appender', async () => { const first = new CapturingAppender(); const second = new CapturingAppender();