Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions apps/pythinker-code/src/utils/region.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,12 +9,15 @@ export interface PythinkerRegionProfile {
readonly telemetryEndpoint: string;
}

// Pythinker runs a single telemetry host, so every region reports to it.
const TELEMETRY_ENDPOINT = 'https://telemetry-logs.pythinker.com/v1/event';

const PROFILES: Record<PythinkerRegion, PythinkerRegionProfile> = {
'mainland-cn': {
telemetryEndpoint: 'https://telemetry-logs.pythinker.com/v1/event',
telemetryEndpoint: TELEMETRY_ENDPOINT,
},
global: {
telemetryEndpoint: 'https://telemetry-logs.pythinker.ai/v1/event',
telemetryEndpoint: TELEMETRY_ENDPOINT,
},
};

Expand Down
4 changes: 0 additions & 4 deletions packages/agent-core-v2/src/app/telemetry/cloudAppender.ts
Original file line number Diff line number Diff line change
Expand Up @@ -84,10 +84,6 @@ export class CloudAppender implements ITelemetryAppender {
storage: options.storage,
deviceId: options.deviceId,
endpoint: options.endpoint,
homeDir: options.bootstrap.homeDir,
readMarker:
(options.bootstrap.getEnv('PYTHINKER_CODE_REGION_MARKER') ??
process.env['PYTHINKER_CODE_REGION_MARKER']) !== 'off',
getAccessToken: options.getAccessToken,
fetchImpl: options.fetchImpl,
retryBackoffsMs: options.retryBackoffsMs,
Expand Down
18 changes: 1 addition & 17 deletions packages/agent-core-v2/src/app/telemetry/cloudTransport.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,4 @@
import { randomBytes } from 'node:crypto';
import { readFileSync } from 'node:fs';
import { join } from 'node:path';

import { isAbortError } from '#/_base/utils/abort';
import type { IFileSystemStorageService } from '#/persistence/interface/storage';
Expand Down Expand Up @@ -33,8 +31,6 @@ export interface CloudTransportOptions {
readonly storage: IFileSystemStorageService;
readonly deviceId: string;
readonly endpoint?: string;
readonly homeDir?: string;
readonly readMarker?: boolean;
readonly getAccessToken?: () => string | null | Promise<string | null>;
readonly fetchImpl?: typeof fetch;
readonly retryBackoffsMs?: readonly number[];
Expand All @@ -50,25 +46,13 @@ 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 GLOBAL_TELEMETRY_ENDPOINT = 'https://telemetry-logs.pythinker.ai/v1/event';
const TELEMETRY_SCOPE = 'telemetry';
const FAILED_PREFIX = 'failed_';
const JSONL_SUFFIX = '.jsonl';

const textEncoder = new TextEncoder();
const textDecoder = new TextDecoder();

function defaultTelemetryEndpoint(homeDir?: string, readMarker = true): string {
if (!readMarker || homeDir === undefined) return TELEMETRY_ENDPOINT;
try {
return readFileSync(join(homeDir, 'region'), 'utf8').trim() === 'global'
? GLOBAL_TELEMETRY_ENDPOINT
: TELEMETRY_ENDPOINT;
} catch {
return TELEMETRY_ENDPOINT;
}
}

export class CloudTransport {
private readonly storage: IFileSystemStorageService;
private readonly deviceId: string;
Expand All @@ -83,7 +67,7 @@ export class CloudTransport {
constructor(options: CloudTransportOptions) {
this.storage = options.storage;
this.deviceId = options.deviceId;
this.endpoint = options.endpoint ?? defaultTelemetryEndpoint(options.homeDir, options.readMarker);
this.endpoint = options.endpoint ?? TELEMETRY_ENDPOINT;
this.getAccessToken = options.getAccessToken ?? null;
this.fetchImpl = options.fetchImpl ?? globalThis.fetch.bind(globalThis);
this.retryBackoffsMs = options.retryBackoffsMs ?? RETRY_BACKOFFS_MS;
Expand Down
50 changes: 1 addition & 49 deletions packages/agent-core-v2/test/app/telemetry/cloudAppender.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -111,7 +111,7 @@ describe('CloudAppender', () => {
expect(typeof event?.['timestamp']).toBe('number');
});

it('reads the install marker from the bootstrapped home for the default endpoint', async () => {
it('reports to the single telemetry host whatever the install marker says', async () => {
writeFileSync(join(homeDir, 'region'), 'global\n');
const requests: CapturedRequest[] = [];
const appender = new CloudAppender(
Expand All @@ -127,58 +127,10 @@ describe('CloudAppender', () => {
appender.track('tool.call', { name: 'bash' });
await appender.flush();

expect(requests).toHaveLength(1);
expect(requests[0]?.url).toBe('https://telemetry-logs.pythinker.ai/v1/event');
});

it('honors the marker opt-out from the bootstrap env bag (no process.env needed)', async () => {
writeFileSync(join(homeDir, 'region'), 'global\n');
const requests: CapturedRequest[] = [];
const appender = new CloudAppender(
baseOptions({
homeDir,
bootstrapEnv: { PYTHINKER_CODE_REGION_MARKER: 'off' },
fetchImpl: makeFetch((req) => {
requests.push(req);
return okResponse();
}),
}),
);

appender.track('tool.call', { name: 'bash' });
await appender.flush();

expect(requests).toHaveLength(1);
expect(requests[0]?.url).toBe('https://telemetry-logs.pythinker.com/v1/event');
});

it('honors PYTHINKER_CODE_REGION_MARKER=off so embedded servers ignore the install marker', async () => {
writeFileSync(join(homeDir, 'region'), 'global\n');
const savedMarkerFlag = process.env['PYTHINKER_CODE_REGION_MARKER'];
process.env['PYTHINKER_CODE_REGION_MARKER'] = 'off';
try {
const requests: CapturedRequest[] = [];
const appender = new CloudAppender(
baseOptions({
homeDir,
fetchImpl: makeFetch((req) => {
requests.push(req);
return okResponse();
}),
}),
);

appender.track('tool.call', { name: 'bash' });
await appender.flush();

expect(requests).toHaveLength(1);
expect(requests[0]?.url).toBe('https://telemetry-logs.pythinker.com/v1/event');
} finally {
if (savedMarkerFlag === undefined) delete process.env['PYTHINKER_CODE_REGION_MARKER'];
else process.env['PYTHINKER_CODE_REGION_MARKER'] = savedMarkerFlag;
}
});

it('applies setContext sessionId and model updates to subsequent events', async () => {
const requests: CapturedRequest[] = [];
const appender = new CloudAppender(
Expand Down
2 changes: 1 addition & 1 deletion packages/agent-gateway/test/setup.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ process.env['PYTHINKER_CODE_EXPERIMENTAL_SEARCH_WORKER'] = 'false';
process.env['PYTHINKER_CODE_EXPERIMENTAL_PERSISTENCE_MINIDB_READMODEL'] = 'false';

const realFetch = globalThis.fetch.bind(globalThis);
const TELEMETRY_HOSTS = new Set(['telemetry-logs.pythinker.com', 'telemetry-logs.pythinker.ai']);
const TELEMETRY_HOSTS = new Set(['telemetry-logs.pythinker.com']);
globalThis.fetch = ((input: string | URL | Request, init?: RequestInit) => {
const url = typeof input === 'string' ? input : input instanceof URL ? input.href : input.url;
if (TELEMETRY_HOSTS.has(new URL(url).hostname)) {
Expand Down
Loading