mirror of
https://github.com/MoonshotAI/kimi-code.git
synced 2026-08-14 11:16:19 +00:00
Merge ee132d60cd into f79fde2b90
This commit is contained in:
commit
4f842b8ac6
19 changed files with 1441 additions and 280 deletions
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<typeof vi.fn>;
|
||||
};
|
||||
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();
|
||||
|
|
|
|||
|
|
@ -2,11 +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.
|
||||
* 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';
|
||||
|
|
@ -14,10 +16,16 @@ import { release } from 'node:os';
|
|||
|
||||
import type { ServicesAccessor } from '#/_base/di/instantiation';
|
||||
import { onUnexpectedError } from '#/_base/errors/unexpectedError';
|
||||
import { abortError, createDeadlineAbortSignal, isAbortError } 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,
|
||||
|
|
@ -28,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;
|
||||
|
|
@ -48,7 +57,8 @@ export interface CloudAppenderOptions {
|
|||
readonly retryBackoffsMs?: readonly number[];
|
||||
readonly requestTimeoutMs?: number;
|
||||
readonly sleep?: (ms: number, signal?: AbortSignal) => Promise<void>;
|
||||
readonly now?: () => number;
|
||||
readonly replayMaxFiles?: number;
|
||||
readonly replayTimeoutMs?: number;
|
||||
}
|
||||
|
||||
export interface CloudAppenderHostOptions {
|
||||
|
|
@ -66,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,
|
||||
});
|
||||
|
|
@ -74,25 +86,39 @@ 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[] = [];
|
||||
private flushTimer: ReturnType<typeof setInterval> | null = null;
|
||||
private readonly lifecycleController = new AbortController();
|
||||
private acceptingEvents = true;
|
||||
private started = false;
|
||||
private replayPending = false;
|
||||
private replayPromise: Promise<boolean> | null = null;
|
||||
private flushPromise: Promise<void> | null = null;
|
||||
private shutdownPromise: Promise<void> | null = null;
|
||||
|
||||
constructor(options: CloudAppenderOptions) {
|
||||
this.deviceId = options.deviceId;
|
||||
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,
|
||||
|
|
@ -100,11 +126,11 @@ export class CloudAppender implements ITelemetryAppender {
|
|||
retryBackoffsMs: options.retryBackoffsMs,
|
||||
requestTimeoutMs: options.requestTimeoutMs,
|
||||
sleep: options.sleep,
|
||||
now: options.now,
|
||||
});
|
||||
}
|
||||
|
||||
track(event: string, properties?: TelemetryProperties): void {
|
||||
if (!this.acceptingEvents) return;
|
||||
const eventSessionId = properties?.['sessionId'];
|
||||
const enriched: EnrichedCloudEvent = {
|
||||
event_id: randomUUID().replaceAll('-', ''),
|
||||
|
|
@ -136,20 +162,46 @@ export class CloudAppender implements ITelemetryAppender {
|
|||
}
|
||||
}
|
||||
|
||||
async flush(): Promise<void> {
|
||||
if (this.buffer.length === 0) return;
|
||||
const events = this.buffer;
|
||||
this.buffer = [];
|
||||
await this.transport.send(events);
|
||||
flush(): Promise<void> {
|
||||
if (this.flushPromise !== null) return this.flushPromise;
|
||||
if (this.buffer.length === 0 && !this.replayPending && this.replayPromise === null) {
|
||||
return Promise.resolve();
|
||||
}
|
||||
const flush = this.drainBuffer();
|
||||
this.flushPromise = flush;
|
||||
void flush.then(
|
||||
() => {
|
||||
this.clearFlush(flush);
|
||||
},
|
||||
() => {
|
||||
this.clearFlush(flush);
|
||||
},
|
||||
);
|
||||
return flush;
|
||||
}
|
||||
|
||||
async shutdown(): Promise<void> {
|
||||
this.stopPeriodicFlush();
|
||||
await this.flush();
|
||||
shutdown(options: TelemetryShutdownOptions = {}): Promise<void> {
|
||||
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.replayPending = true;
|
||||
void this.ensureReplay();
|
||||
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 +215,112 @@ export class CloudAppender implements ITelemetryAppender {
|
|||
}
|
||||
|
||||
async retryDiskEvents(): Promise<void> {
|
||||
await this.transport.retryDiskEvents();
|
||||
this.replayPending = true;
|
||||
await this.ensureReplay();
|
||||
}
|
||||
|
||||
private async drainBuffer(): Promise<void> {
|
||||
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 {
|
||||
await this.handoffEvents(events);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private clearFlush(flush: Promise<void>): void {
|
||||
if (this.flushPromise === flush) this.flushPromise = null;
|
||||
}
|
||||
|
||||
private async shutdownOwnedWork(): Promise<void> {
|
||||
await this.flush();
|
||||
}
|
||||
|
||||
private ensureReplay(): Promise<boolean> {
|
||||
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<void> {
|
||||
while (this.buffer.length > 0) {
|
||||
const events = this.buffer;
|
||||
this.buffer = [];
|
||||
await this.handoffEvents(events);
|
||||
}
|
||||
}
|
||||
|
||||
private async handoffEvents(events: readonly EnrichedCloudEvent[]): Promise<void> {
|
||||
try {
|
||||
await this.spool.put(events);
|
||||
} catch (storageError) {
|
||||
this.buffer = [...events, ...this.buffer];
|
||||
throw storageError;
|
||||
}
|
||||
}
|
||||
|
||||
private async replayDiskEvents(signal: AbortSignal): Promise<boolean> {
|
||||
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) {
|
||||
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<typeof setTimeout> | 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);
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1,15 +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 the `telemetry` 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 { 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;
|
||||
|
||||
|
|
@ -36,7 +38,6 @@ export interface CloudPayload {
|
|||
}
|
||||
|
||||
export interface CloudTransportOptions {
|
||||
readonly storage: IFileSystemStorageService;
|
||||
readonly deviceId: string;
|
||||
readonly endpoint?: string;
|
||||
readonly getAccessToken?: () => string | null | Promise<string | null>;
|
||||
|
|
@ -44,25 +45,16 @@ export interface CloudTransportOptions {
|
|||
readonly retryBackoffsMs?: readonly number[];
|
||||
readonly requestTimeoutMs?: number;
|
||||
readonly sleep?: (ms: number, signal?: AbortSignal) => Promise<void>;
|
||||
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';
|
||||
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<string | null>) | null;
|
||||
|
|
@ -70,10 +62,8 @@ export class CloudTransport {
|
|||
private readonly retryBackoffsMs: readonly number[];
|
||||
private readonly requestTimeoutMs: number;
|
||||
private readonly sleepImpl: (ms: number, signal?: AbortSignal) => Promise<void>;
|
||||
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;
|
||||
|
|
@ -81,21 +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<void> {
|
||||
async send(
|
||||
events: readonly EnrichedCloudEvent[],
|
||||
signal?: AbortSignal,
|
||||
retryBackoffsMs: readonly number[] = this.retryBackoffsMs,
|
||||
): Promise<void> {
|
||||
if (events.length === 0) return;
|
||||
let savedToDisk = false;
|
||||
const saveEventsToDisk = async (): Promise<void> => {
|
||||
if (savedToDisk) return;
|
||||
await this.saveToDisk(events);
|
||||
savedToDisk = true;
|
||||
};
|
||||
if (signal?.aborted === true) {
|
||||
await saveEventsToDisk();
|
||||
throw abortError();
|
||||
}
|
||||
if (signal?.aborted === true) throw abortError();
|
||||
|
||||
let payload: CloudPayload;
|
||||
try {
|
||||
|
|
@ -104,88 +88,53 @@ export class CloudTransport {
|
|||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
for (let attempt = 0; attempt <= this.retryBackoffsMs.length; attempt++) {
|
||||
let lastError: unknown;
|
||||
for (let attempt = 0; attempt <= retryBackoffsMs.length; attempt++) {
|
||||
try {
|
||||
await this.sendAttempt(payload, signal);
|
||||
return;
|
||||
} catch (error) {
|
||||
if (isSignalAborted(signal) || (isAbortError(error) && !(error instanceof TransientCloudError))) {
|
||||
throw error;
|
||||
}
|
||||
lastError = error;
|
||||
if (!(error instanceof TransientCloudError)) break;
|
||||
const backoff = retryBackoffsMs[attempt];
|
||||
if (backoff === undefined) break;
|
||||
const sleep = Promise.resolve().then(() => this.sleepImpl(backoff, signal));
|
||||
try {
|
||||
await this.sendHttp(payload, signal);
|
||||
return;
|
||||
} catch (error) {
|
||||
if (isSignalAborted(signal) || isAbortError(error)) {
|
||||
await saveEventsToDisk();
|
||||
throw error;
|
||||
}
|
||||
if (!(error instanceof TransientCloudError)) {
|
||||
break;
|
||||
}
|
||||
const backoff = this.retryBackoffsMs[attempt];
|
||||
if (backoff === undefined) break;
|
||||
await this.sleepImpl(backoff, signal);
|
||||
await (signal === undefined ? sleep : abortable(sleep, signal));
|
||||
} catch (sleepError) {
|
||||
if (isSignalAborted(signal) || isAbortError(sleepError)) throw sleepError;
|
||||
lastError = sleepError;
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
throw lastError instanceof Error
|
||||
? lastError
|
||||
: new TransientCloudError('telemetry delivery failed');
|
||||
}
|
||||
|
||||
private async sendAttempt(payload: CloudPayload, signal?: AbortSignal): Promise<void> {
|
||||
const source = signal ?? new AbortController().signal;
|
||||
const deadline = createDeadlineAbortSignal(source, this.requestTimeoutMs);
|
||||
try {
|
||||
await this.sendHttp(payload, deadline.signal);
|
||||
} catch (error) {
|
||||
if (isSignalAborted(signal) || isAbortError(error)) {
|
||||
await saveEventsToDisk();
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
await saveEventsToDisk();
|
||||
}
|
||||
|
||||
async saveToDisk(events: readonly EnrichedCloudEvent[]): Promise<void> {
|
||||
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(): Promise<void> {
|
||||
const keys = await this.storage.list(TELEMETRY_SCOPE, FAILED_PREFIX);
|
||||
const now = this.now();
|
||||
for (const key of keys) {
|
||||
if (!key.startsWith(FAILED_PREFIX) || !key.endsWith(JSONL_SUFFIX)) continue;
|
||||
const createdAt = parseFailedTimestamp(key);
|
||||
if (createdAt === undefined || 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 {
|
||||
await this.sendHttp(payload);
|
||||
await this.storage.delete(TELEMETRY_SCOPE, key);
|
||||
} catch (error) {
|
||||
if (error instanceof TransientCloudError) continue;
|
||||
}
|
||||
if (deadline.timedOut()) throw new TransientCloudError('telemetry request timed out');
|
||||
throw error;
|
||||
} finally {
|
||||
deadline.clear();
|
||||
}
|
||||
}
|
||||
|
||||
private async readJsonl(key: string): Promise<EnrichedCloudEvent[]> {
|
||||
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);
|
||||
}
|
||||
return events;
|
||||
}
|
||||
|
||||
private async sendHttp(payload: CloudPayload, signal?: AbortSignal): Promise<void> {
|
||||
const token = this.getAccessToken === null ? null : await this.getAccessToken();
|
||||
private async sendHttp(payload: CloudPayload, signal: AbortSignal): Promise<void> {
|
||||
const tokenRequest =
|
||||
this.getAccessToken === null
|
||||
? Promise.resolve(null)
|
||||
: Promise.resolve().then(() => this.getAccessToken?.() ?? null);
|
||||
const token = await abortable(tokenRequest, signal);
|
||||
const headers: Record<string, string> = {
|
||||
'Content-Type': 'application/json',
|
||||
};
|
||||
|
|
@ -206,36 +155,25 @@ export class CloudTransport {
|
|||
private async post(
|
||||
payload: CloudPayload,
|
||||
headers: Record<string, string>,
|
||||
signal?: AbortSignal,
|
||||
signal: AbortSignal,
|
||||
): Promise<Response> {
|
||||
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,48 +244,22 @@ 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<Response> {
|
||||
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<void> {
|
||||
if (signal?.aborted === true) return Promise.reject(abortError());
|
||||
return new Promise((resolve, reject) => {
|
||||
const timer = setTimeout(resolve, ms);
|
||||
timer.unref?.();
|
||||
let timer: ReturnType<typeof setTimeout>;
|
||||
const onAbort = (): void => {
|
||||
clearTimeout(timer);
|
||||
reject(abortError());
|
||||
};
|
||||
timer = setTimeout(() => {
|
||||
signal?.removeEventListener('abort', onAbort);
|
||||
resolve();
|
||||
}, ms);
|
||||
timer.unref?.();
|
||||
signal?.addEventListener('abort', onAbort, { once: true });
|
||||
});
|
||||
}
|
||||
|
|
@ -355,7 +267,3 @@ function abortableSleep(ms: number, signal?: AbortSignal): Promise<void> {
|
|||
function isSignalAborted(signal?: AbortSignal): boolean {
|
||||
return signal?.aborted === true;
|
||||
}
|
||||
|
||||
function abortError(): DOMException {
|
||||
return new DOMException('The operation was aborted.', 'AbortError');
|
||||
}
|
||||
|
|
|
|||
|
|
@ -24,12 +24,22 @@ export type TelemetryProperties = Readonly<Record<string, TelemetryPrimitive>>;
|
|||
|
||||
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> | void;
|
||||
shutdown?(): Promise<void> | void;
|
||||
shutdown?(options?: TelemetryShutdownOptions): Promise<void> | void;
|
||||
}
|
||||
|
||||
export interface ITelemetryAppenderRegistration extends IDisposable {
|
||||
shutdown(options?: TelemetryShutdownOptions): Promise<void>;
|
||||
}
|
||||
|
||||
export interface TelemetryServiceOptions {
|
||||
|
|
@ -51,12 +61,15 @@ 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<void>;
|
||||
setAppender(appender: ITelemetryAppender, options?: TelemetryShutdownOptions): Promise<void>;
|
||||
setEnabled(enabled: boolean): void;
|
||||
flush(): Promise<void>;
|
||||
shutdown(): Promise<void>;
|
||||
shutdown(options?: TelemetryShutdownOptions): Promise<void>;
|
||||
}
|
||||
|
||||
export const nullTelemetryAppender: ITelemetryAppender = {
|
||||
|
|
@ -73,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 () => {},
|
||||
|
|
|
|||
|
|
@ -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,15 +19,23 @@ import type {
|
|||
import {
|
||||
ITelemetryService,
|
||||
type ITelemetryAppender,
|
||||
type ITelemetryAppenderRegistration,
|
||||
nullTelemetryAppender,
|
||||
type TelemetryContextPatch,
|
||||
type TelemetryProperties,
|
||||
type TelemetryShutdownOptions,
|
||||
} from './telemetry';
|
||||
|
||||
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<ITelemetryAppender, Promise<void>>();
|
||||
private shutdownPromise: Promise<void> | null = null;
|
||||
private context: TelemetryProperties = {};
|
||||
private enabled = true;
|
||||
|
||||
|
|
@ -39,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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -63,17 +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<void> {
|
||||
const retirement = this.retirements.get(appender);
|
||||
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.retireAppender(appender, options);
|
||||
}
|
||||
|
||||
setAppender(appender: ITelemetryAppender): void {
|
||||
async setAppender(
|
||||
appender: ITelemetryAppender,
|
||||
options?: TelemetryShutdownOptions,
|
||||
): Promise<void> {
|
||||
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.retireAppender(candidate, options)));
|
||||
}
|
||||
|
||||
setEnabled(enabled: boolean): void {
|
||||
|
|
@ -82,18 +109,69 @@ export class TelemetryService implements ITelemetryService {
|
|||
|
||||
async flush(): Promise<void> {
|
||||
await Promise.all(
|
||||
this.appenders.map((appender) =>
|
||||
Promise.resolve(appender.flush?.()).catch(onUnexpectedError),
|
||||
),
|
||||
this.appenders.map((appender) => this.invokeAppender(() => appender.flush?.())),
|
||||
);
|
||||
}
|
||||
|
||||
async shutdown(): Promise<void> {
|
||||
await Promise.all(
|
||||
this.appenders.map((appender) =>
|
||||
Promise.resolve(appender.shutdown?.()).catch(onUnexpectedError),
|
||||
),
|
||||
);
|
||||
shutdown(options?: TelemetryShutdownOptions): Promise<void> {
|
||||
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);
|
||||
} else if (options !== undefined) {
|
||||
for (const appender of this.retirements.keys()) {
|
||||
void this.retireAppender(appender, options);
|
||||
}
|
||||
}
|
||||
return this.shutdownPromise;
|
||||
}
|
||||
|
||||
private startAppender(appender: ITelemetryAppender): void {
|
||||
try {
|
||||
appender.start?.();
|
||||
} catch (error) {
|
||||
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 retireAppender(
|
||||
appender: ITelemetryAppender,
|
||||
options?: TelemetryShutdownOptions,
|
||||
): Promise<void> {
|
||||
const retirement = this.retirements.get(appender);
|
||||
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;
|
||||
}
|
||||
|
||||
private invokeAppender(operation: () => Promise<void> | void | undefined): Promise<void> {
|
||||
return Promise.resolve().then(operation).catch(onUnexpectedError);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -127,16 +205,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<void> {
|
||||
return this.root.removeAppender(appender, options);
|
||||
}
|
||||
|
||||
setAppender(appender: ITelemetryAppender): void {
|
||||
this.root.setAppender(appender);
|
||||
setAppender(
|
||||
appender: ITelemetryAppender,
|
||||
options?: TelemetryShutdownOptions,
|
||||
): Promise<void> {
|
||||
return this.root.setAppender(appender, options);
|
||||
}
|
||||
|
||||
setEnabled(enabled: boolean): void {
|
||||
|
|
@ -147,8 +231,8 @@ class TelemetryContextView implements ITelemetryService {
|
|||
return this.root.flush();
|
||||
}
|
||||
|
||||
shutdown(): Promise<void> {
|
||||
return this.root.shutdown();
|
||||
shutdown(options?: TelemetryShutdownOptions): Promise<void> {
|
||||
return this.root.shutdown(options);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
191
packages/agent-core-v2/src/app/telemetry/telemetrySpoolStore.ts
Normal file
191
packages/agent-core-v2/src/app/telemetry/telemetrySpoolStore.ts
Normal file
|
|
@ -0,0 +1,191 @@
|
|||
/**
|
||||
* `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<void>;
|
||||
recoverable(limit: number): Promise<readonly TelemetrySpoolEntry[]>;
|
||||
acknowledge(key: string): Promise<void>;
|
||||
}
|
||||
|
||||
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_';
|
||||
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 maxFileBytes: number;
|
||||
private readonly now: () => number;
|
||||
private writes: Promise<void> = 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;
|
||||
}
|
||||
|
||||
put(events: readonly EnrichedCloudEvent[]): Promise<void> {
|
||||
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<readonly TelemetrySpoolEntry[]> {
|
||||
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<void> {
|
||||
return this.storage.delete(TELEMETRY_SCOPE, key);
|
||||
}
|
||||
|
||||
private async enforcePolicy(): Promise<readonly SpoolFile[]> {
|
||||
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));
|
||||
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<void> {
|
||||
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)),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
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 createFile(): SpoolFile {
|
||||
const createdAt = this.now();
|
||||
return {
|
||||
key: `${FAILED_PREFIX}${createdAt}_${randomBytes(6).toString('hex')}${JSONL_SUFFIX}`,
|
||||
createdAt,
|
||||
};
|
||||
}
|
||||
|
||||
private async readJsonl(key: string): Promise<readonly EnrichedCloudEvent[]> {
|
||||
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;
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
|
|
@ -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 () => {},
|
||||
|
|
|
|||
|
|
@ -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 () => {},
|
||||
|
|
|
|||
|
|
@ -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(),
|
||||
|
|
|
|||
|
|
@ -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(),
|
||||
|
|
|
|||
|
|
@ -1,8 +1,24 @@
|
|||
import { mkdtempSync, readdirSync, rmSync } from 'node:fs';
|
||||
/**
|
||||
* `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.
|
||||
*/
|
||||
|
||||
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';
|
||||
|
||||
import { afterEach, beforeEach, describe, expect, it } from 'vitest';
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
|
||||
|
||||
import {
|
||||
resetUnexpectedErrorHandler,
|
||||
|
|
@ -10,12 +26,15 @@ import {
|
|||
} from '#/_base/errors/unexpectedError';
|
||||
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 { TelemetrySpoolStore } from '#/app/telemetry/telemetrySpoolStore';
|
||||
|
||||
import { stubBootstrap } from '../bootstrap/stubs';
|
||||
|
||||
interface CapturedRequest {
|
||||
readonly url: string;
|
||||
readonly headers: Record<string, string>;
|
||||
readonly signal?: AbortSignal;
|
||||
readonly body: {
|
||||
readonly user_id: string;
|
||||
readonly events: readonly Record<string, unknown>[];
|
||||
|
|
@ -26,10 +45,15 @@ type Responder = (req: CapturedRequest) => Response | Promise<Response>;
|
|||
|
||||
function makeFetch(responder: Responder): typeof fetch {
|
||||
return (async (input: unknown, init: unknown) => {
|
||||
const requestInit = init as { headers: Record<string, string>; body: string };
|
||||
const requestInit = init as {
|
||||
headers: Record<string, string>;
|
||||
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,12 +68,33 @@ function statusResponse(status: number): Response {
|
|||
return new Response(null, { status });
|
||||
}
|
||||
|
||||
function deferred<T>(): {
|
||||
readonly promise: Promise<T>;
|
||||
readonly resolve: (value: T) => void;
|
||||
} {
|
||||
let resolve!: (value: T) => void;
|
||||
const promise = new Promise<T>((done) => {
|
||||
resolve = done;
|
||||
});
|
||||
return { promise, resolve };
|
||||
}
|
||||
|
||||
function baseOptions(
|
||||
overrides: Partial<CloudAppenderOptions> & { homeDir?: string } = {},
|
||||
overrides: Partial<CloudAppenderOptions> & {
|
||||
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',
|
||||
|
|
@ -58,6 +103,104 @@ function baseOptions(
|
|||
};
|
||||
}
|
||||
|
||||
function listFailedSpoolFiles(homeDir: string): string[] {
|
||||
return readdirSync(join(homeDir, 'telemetry-v2')).filter((file) =>
|
||||
file.startsWith('failed_'),
|
||||
);
|
||||
}
|
||||
|
||||
function readFirstFailedEvent(homeDir: string): Record<string, unknown> {
|
||||
const file = listFailedSpoolFiles(homeDir)[0] as string;
|
||||
const persisted = readFileSync(join(homeDir, 'telemetry-v2', file), 'utf8');
|
||||
return JSON.parse(persisted.trim()) as Record<string, unknown>;
|
||||
}
|
||||
|
||||
function readAllFailedEvents(homeDir: string): Record<string, unknown>[] {
|
||||
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<string, unknown>),
|
||||
);
|
||||
}
|
||||
|
||||
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;
|
||||
|
||||
|
|
@ -202,6 +345,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<void>();
|
||||
const appender = new CloudAppender(
|
||||
baseOptions({
|
||||
homeDir,
|
||||
fetchImpl: makeFetch(() => {
|
||||
requestStarted.resolve();
|
||||
return new Promise<Response>(() => {});
|
||||
}),
|
||||
}),
|
||||
);
|
||||
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<Response>();
|
||||
const firstRequestStarted = deferred<void>();
|
||||
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<Response>();
|
||||
const requestStarted = deferred<void>();
|
||||
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<void>();
|
||||
let requestSignal: AbortSignal | undefined;
|
||||
const appender = new CloudAppender(
|
||||
baseOptions({
|
||||
homeDir,
|
||||
flushThreshold: 1,
|
||||
fetchImpl: makeFetch((request) => {
|
||||
requestSignal = request.signal;
|
||||
requestStarted.resolve();
|
||||
return new Promise<Response>(() => {});
|
||||
}),
|
||||
}),
|
||||
);
|
||||
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 +526,121 @@ 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({
|
||||
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('shutdown cancellation interrupts token lookup and durably hands off the batch', async () => {
|
||||
const token = deferred<string | null>();
|
||||
const tokenStarted = deferred<void>();
|
||||
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<void>();
|
||||
const sleep = deferred<void>();
|
||||
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(
|
||||
|
|
@ -255,15 +674,162 @@ 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(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('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<void>();
|
||||
const replayResponse = deferred<Response>();
|
||||
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('keeps the newest durable batches when the spool reaches its file cap', 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(
|
||||
readdirSync(join(homeDir, 'telemetry')).filter((f) => f.startsWith('failed_')),
|
||||
).toHaveLength(0);
|
||||
readAllFailedEvents(homeDir)
|
||||
.map((event) => event['event'])
|
||||
.toSorted(),
|
||||
).toEqual(['second', 'third']);
|
||||
});
|
||||
|
||||
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 () => {
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
},
|
||||
|
|
|
|||
|
|
@ -1,3 +1,10 @@
|
|||
/**
|
||||
* `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';
|
||||
|
||||
import { LifecycleScope, ScopeActivation, _clearScopedRegistryForTests, registerScopedService } from '#/_base/di/scope';
|
||||
|
|
@ -6,21 +13,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;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -28,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);
|
||||
|
|
@ -36,6 +54,14 @@ function telemetryWithAppenders(...appenders: ITelemetryAppender[]): TelemetrySe
|
|||
return svc;
|
||||
}
|
||||
|
||||
function deferred(): { readonly promise: Promise<void>; readonly resolve: () => void } {
|
||||
let resolve!: () => void;
|
||||
const promise = new Promise<void>((done) => {
|
||||
resolve = done;
|
||||
});
|
||||
return { promise, resolve };
|
||||
}
|
||||
|
||||
describe('TelemetryService (unit)', () => {
|
||||
it('noop by default — does not throw', () => {
|
||||
const svc = new TelemetryService();
|
||||
|
|
@ -45,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({
|
||||
|
|
@ -57,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' });
|
||||
|
|
@ -72,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');
|
||||
|
|
@ -101,16 +127,81 @@ 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('registration shutdown tightens retirement started by disposal', async () => {
|
||||
const retired = deferred();
|
||||
const shutdownOptions: Array<TelemetryShutdownOptions | undefined> = [];
|
||||
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();
|
||||
|
||||
svc.addAppender(appender);
|
||||
|
||||
expect(appender.startCalls).toBe(1);
|
||||
});
|
||||
|
||||
it('removeAppender stops delivery to that appender', () => {
|
||||
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);
|
||||
});
|
||||
|
||||
it('setAppender starts the replacement appender', () => {
|
||||
const appender = new CapturingAppender();
|
||||
const svc = new TelemetryService();
|
||||
|
||||
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);
|
||||
|
|
@ -141,7 +232,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' } }]);
|
||||
|
|
@ -165,6 +256,86 @@ 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('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('service shutdown forwards a lifecycle budget to retirement started by disposal', async () => {
|
||||
const retired = deferred();
|
||||
const shutdownOptions: Array<TelemetryShutdownOptions | undefined> = [];
|
||||
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();
|
||||
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);
|
||||
|
|
@ -189,6 +360,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() {},
|
||||
|
|
@ -202,6 +390,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() {},
|
||||
|
|
@ -214,6 +417,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)', () => {
|
||||
|
|
|
|||
|
|
@ -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 () => {},
|
||||
|
|
|
|||
|
|
@ -287,9 +287,9 @@ function telemetryStub(events: Array<{ event: string; properties: Record<string,
|
|||
},
|
||||
withContext: () => telemetryStub(events),
|
||||
setContext: () => {},
|
||||
addAppender: () => ({ dispose: () => {} }),
|
||||
removeAppender: () => {},
|
||||
setAppender: () => {},
|
||||
addAppender: () => ({ dispose: () => {}, shutdown: async () => {} }),
|
||||
removeAppender: async () => {},
|
||||
setAppender: async () => {},
|
||||
setEnabled: () => {},
|
||||
flush: async () => {},
|
||||
shutdown: async () => {},
|
||||
|
|
|
|||
|
|
@ -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,28 +61,23 @@ 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 };
|
||||
return { registration };
|
||||
}
|
||||
|
||||
export async function shutdownServerTelemetry(
|
||||
telemetry: ServerTelemetry,
|
||||
deadlineMs = Date.now() + TELEMETRY_SHUTDOWN_TIMEOUT_MS,
|
||||
): Promise<void> {
|
||||
telemetry.registration?.dispose();
|
||||
if (telemetry.appender === undefined) return;
|
||||
if (telemetry.registration === undefined) return;
|
||||
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||
try {
|
||||
await Promise.race([
|
||||
telemetry.appender.shutdown(),
|
||||
telemetry.registration.shutdown({ deadlineMs }),
|
||||
new Promise<void>((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 {
|
||||
|
|
|
|||
|
|
@ -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 = {
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
});
|
||||
|
|
@ -113,13 +114,19 @@ 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();
|
||||
const spool = await app.accessor.get(IFileSystemStorageService).list('telemetry-v2');
|
||||
expect(spool).toHaveLength(1);
|
||||
} finally {
|
||||
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);
|
||||
});
|
||||
|
||||
|
|
@ -131,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);
|
||||
},
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue