From c05556f3cfd4fc409c94603f6330356bb6e87c70 Mon Sep 17 00:00:00 2001 From: Saket Aryan Date: Thu, 17 Sep 2026 18:43:27 +0530 Subject: [PATCH] fix(plugins): stop the TypeScript telemetry losing events and misattributing accounts The Python plugin telemetry was hardened across #7322 to #7326. The TypeScript side has the same defect classes and was not touched, because the two share no code: agent-plugin-core/python generates into six bundles, agent-plugin-core/ typescript is a separate core each plugin wraps. Fixing one surfaced nothing about the other, which is how this survived. Three fixes, all confirmed by running the code rather than reading it. A failed delivery deleted the batch. The core detached the queue before the await and swallowed the error, so one blip destroyed the events with nothing recording that it happened. Probed: two events in, delivery throws, queue goes to zero, no retry ever. The batch is now put back, bounded by maxQueueSize and biased to the newest so a long outage costs the oldest events rather than unbounded memory, and repeated failures back off to a ceiling instead of retrying every flush against a host that is blocking us. Every event now carries a uuid stamped at capture, which is what makes the retry safe: PostHog collapses anything it already accepted. Deliberately no disk spool, and that is written into the code so it reads as a decision. Python spools because its hooks are per-tool-call processes that exit immediately. These plugins live inside a host for a whole session, so re-queueing covers the same transient failures without the claim and lease machinery that took three review rounds to get right on the Python side. What it leaves uncovered is narrow: a session that both starts and ends offline. openclaw used a cached email forever. The refresh was guarded by !hasEmail, so after an API key change every event kept reporting under the previous account. The email is now bound to a fingerprint of the key it was resolved for and only used while those agree; a mismatch forgets the account and re-resolves. The resolution latch is per key rather than once per process, so a key changed mid-session is actually looked up. A row with an email and no fingerprint, which is what an upgrade from the current version looks like, is verified rather than adopted, matching the decision reached on #7325. opencode hashed the project id unsalted, which is reversible for anyone who can enumerate project ids. Salted with the API key rather than a stored per-install value: it is already in play, it is high entropy, and it needs no new file and so no write race to get wrong. Per account rather than per machine, which also keeps joins working across machines, and it resets on key rotation consistently with distinctId, which already did. Also wired opencode's tests into CI. They existed and nothing ran them, so the regression test asked for on #7322 would not have gated anything. Verified: agent-plugin-core/ts 30, openclaw 431, opencode 35, pi-agent 89, deepseek 47. The new tests were each checked against the unfixed code first; the core ones fail 3 of 3 and the openclaw one fails without the fingerprint gate. Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb --- .github/workflows/opencode-plugin-checks.yml | 3 + .../typescript/src/telemetry.ts | 33 ++++++- .../typescript/tests/telemetry.test.ts | 92 +++++++++++++++++++ integrations/openclaw/cli/config-file.ts | 17 ++++ integrations/openclaw/telemetry.ts | 44 +++++++-- integrations/openclaw/tests/telemetry.test.ts | 54 ++++++++++- integrations/opencode-plugin/package.json | 3 +- .../opencode-plugin/telemetry.test.ts | 44 ++++++++- integrations/opencode-plugin/telemetry.ts | 23 ++++- 9 files changed, 294 insertions(+), 19 deletions(-) diff --git a/.github/workflows/opencode-plugin-checks.yml b/.github/workflows/opencode-plugin-checks.yml index 4fb95fc6c..a93d741eb 100644 --- a/.github/workflows/opencode-plugin-checks.yml +++ b/.github/workflows/opencode-plugin-checks.yml @@ -32,6 +32,9 @@ jobs: - name: Type check run: bun run type-check + - name: Test + run: bun test + - name: Build run: bun run build diff --git a/integrations/agent-plugin-core/typescript/src/telemetry.ts b/integrations/agent-plugin-core/typescript/src/telemetry.ts index c6a7544ce..656ed9858 100644 --- a/integrations/agent-plugin-core/typescript/src/telemetry.ts +++ b/integrations/agent-plugin-core/typescript/src/telemetry.ts @@ -1,3 +1,5 @@ +import { randomUUID } from "node:crypto"; + import { redactSecrets } from "./lifecycle.ts"; const POSTHOG_API_KEY = "phc_hgJkUVJFYtmaJqrvf6CYN67TIQ8yhXAkWzUn9AMU4yX"; @@ -74,9 +76,20 @@ export function errorKind(error: unknown): string { return error instanceof Error ? error.constructor.name : "other"; } +// Delivery is retried in memory, not spooled to disk, and that is a decision +// rather than an omission. The Python core spools because its hooks are separate +// processes that fire per tool call and exit immediately, so nothing survives +// without a file. These plugins are loaded into a host that lives for a whole +// session, so re-queueing covers the same transient failures without the claim +// and lease machinery a correct cross-process spool needs. What that leaves +// uncovered is narrow: a session that both starts and ends with no connectivity. +const RETRY_BACKOFF_CEILING_MS = 60_000; + export function createTelemetry(config: TelemetryConfig) { let queue: Record[] = []; let timer: ReturnType | undefined; + let consecutiveFailures = 0; + let retryNotBefore = 0; const flushThreshold = config.flushThreshold ?? 10; const maxQueueSize = config.maxQueueSize ?? 100; @@ -91,12 +104,24 @@ export function createTelemetry(config: TelemetryConfig) { async function flush(): Promise { if (!queue.length) return; + if (Date.now() < retryNotBefore) return; const batch = queue; queue = []; try { await deliver(batch); + consecutiveFailures = 0; + retryNotBefore = 0; } catch { - // Telemetry must never affect plugin behavior. + // Put it back. Detaching the batch and swallowing the error deleted the + // events outright, so any blip silently dropped telemetry with nothing + // recording that it had happened. Every event carries a uuid, so a retry + // that duplicates one PostHog already accepted is collapsed there. + // + // Bounded by maxQueueSize and biased to the newest, matching capture(): + // a long outage costs the oldest events rather than unbounded memory. + queue = [...batch, ...queue].slice(-maxQueueSize); + consecutiveFailures += 1; + retryNotBefore = Date.now() + Math.min(2 ** consecutiveFailures * 1_000, RETRY_BACKOFF_CEILING_MS); } } @@ -112,6 +137,10 @@ export function createTelemetry(config: TelemetryConfig) { return { event: config.eventName?.(event) ?? event, distinct_id: distinctId, + // Stamped once, at capture. This is what makes retrying safe: a batch + // re-sent after a failure carries the same ids, so PostHog collapses + // anything it already accepted instead of counting it twice. + uuid: randomUUID(), properties: { ...safeProperties(properties), ...safeProperties(config.commonProperties ?? {}), @@ -149,6 +178,8 @@ export function createTelemetry(config: TelemetryConfig) { function resetForTesting(): void { queue = []; + consecutiveFailures = 0; + retryNotBefore = 0; if (timer) clearInterval(timer); timer = undefined; process.off("beforeExit", beforeExit); diff --git a/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts b/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts index e04794614..3e5ab9fba 100644 --- a/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts +++ b/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts @@ -122,3 +122,95 @@ test("error classification does not expose messages", () => { assert.equal(errorKind(new Error("request timeout")), "timeout"); assert.equal(errorKind(new Error("fetch failed")), "network"); }); + +test("a failed delivery keeps the batch instead of deleting it", async () => { + // The defect: the queue was detached before the await and the error swallowed, + // so one blip destroyed the events with nothing recording that it happened. + const attempts: Record[][] = []; + let failNext = true; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", + flushThreshold: 1000, + delivery: async (batch) => { + attempts.push(batch); + if (failNext) throw new Error("network down"); + }, + }); + + telemetry.capture("one"); + telemetry.capture("two"); + await telemetry.flush(); + + assert.equal(attempts.length, 1); + assert.equal(telemetry.queueForTesting().length, 2, "events were dropped on failure"); + + failNext = false; + // Backoff is in force, so wait it out the way wall time would. + await new Promise((resolve) => setTimeout(resolve, 2_100)); + await telemetry.flush(); + + assert.equal(attempts.length, 2, "never retried"); + assert.equal(telemetry.queueForTesting().length, 0); + telemetry.resetForTesting(); +}); + +test("a retried event carries the same uuid so PostHog can collapse it", async () => { + const attempts: Record[][] = []; + let failNext = true; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", + flushThreshold: 1000, + delivery: async (batch) => { + attempts.push(batch); + if (failNext) throw new Error("network down"); + }, + }); + + telemetry.capture("once"); + await telemetry.flush(); + failNext = false; + await new Promise((resolve) => setTimeout(resolve, 2_100)); + await telemetry.flush(); + + assert.equal(attempts.length, 2); + const first = attempts[0][0].uuid; + assert.ok(first, "events carry no uuid, so a retry would double count"); + assert.equal(attempts[1][0].uuid, first, "retry minted a new uuid"); + telemetry.resetForTesting(); +}); + +test("repeated failures back off instead of retrying every flush", async () => { + let calls = 0; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", + flushThreshold: 1000, + delivery: async () => { calls += 1; throw new Error("blocked"); }, + }); + + telemetry.capture("one"); + await telemetry.flush(); + await telemetry.flush(); + await telemetry.flush(); + + assert.equal(calls, 1, "a blocked host was hammered on every flush"); + assert.equal(telemetry.queueForTesting().length, 1, "the event was lost while backing off"); + telemetry.resetForTesting(); +}); + +test("a long outage costs the oldest events, not unbounded memory", async () => { + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", + flushThreshold: 1000, maxQueueSize: 3, + delivery: async () => { throw new Error("down"); }, + }); + + telemetry.capture("a"); + telemetry.capture("b"); + await telemetry.flush(); + telemetry.capture("c"); + telemetry.capture("d"); + telemetry.capture("e"); + + assert.ok(telemetry.queueForTesting().length <= 3, "queue grew past maxQueueSize"); + telemetry.resetForTesting(); +}); diff --git a/integrations/openclaw/cli/config-file.ts b/integrations/openclaw/cli/config-file.ts index 99453cb1b..ae4d9e471 100644 --- a/integrations/openclaw/cli/config-file.ts +++ b/integrations/openclaw/cli/config-file.ts @@ -35,6 +35,8 @@ export interface PluginAuthConfig { autoCapture?: boolean; topK?: number; anonymousTelemetryId?: string; + /** SHA-256 prefix of the API key userEmail was resolved for. */ + keyFingerprint?: string; } // ============================================================================ @@ -236,6 +238,21 @@ export function getBaseUrl(): string { return auth.baseUrl || DEFAULT_BASE_URL; } +/** Forget the resolved account, so the next capture re-resolves for the current key. */ +export function clearResolvedAccount(): void { + const full = readFullConfig() as any; + const cfg = full?.plugins?.entries?.[PLUGIN_ID]?.config; + if (!cfg) return; + let changed = false; + for (const key of ["userEmail", "keyFingerprint"]) { + if (key in cfg) { + delete cfg[key]; + changed = true; + } + } + if (changed) writeFullConfig(full); +} + /** Remove anonymousTelemetryId from config (after PostHog aliasing) */ export function clearAnonymousTelemetryId(): void { const full = readFullConfig() as any; diff --git a/integrations/openclaw/telemetry.ts b/integrations/openclaw/telemetry.ts index fc1c6b99d..f0f3c3fc7 100644 --- a/integrations/openclaw/telemetry.ts +++ b/integrations/openclaw/telemetry.ts @@ -1,14 +1,20 @@ import { createHash, randomUUID } from "node:crypto"; import { createTelemetry } from "../agent-plugin-core/typescript/src/telemetry.ts"; -import { clearAnonymousTelemetryId, getBaseUrl, readPluginAuth, writePluginAuth } from "./cli/config-file.ts"; +import { + clearAnonymousTelemetryId, + clearResolvedAccount, + getBaseUrl, + readPluginAuth, + writePluginAuth, +} from "./cli/config-file.ts"; declare const __OPENCLAW_PLUGIN_VERSION__: string; export const PLUGIN_VERSION: string = __OPENCLAW_PLUGIN_VERSION__; let cachedAnonymousId: string | undefined; let aliasCheckDone = false; -let emailResolutionAttempted = false; +let resolutionAttemptedFor = ""; let currentDistinctId = ""; function enabled(): boolean { @@ -33,10 +39,23 @@ function anonymousId(): string { return (cachedAnonymousId = created); } +/** SHA-256 prefix of the key an account was resolved for. */ +function keyFingerprint(apiKey?: string): string { + return apiKey ? createHash("sha256").update(apiKey).digest("hex").slice(0, 16) : ""; +} + function distinctId(apiKey?: string): string { try { - const email = readPluginAuth().userEmail; - if (email) return createHash("sha256").update(email).digest("hex"); + const auth = readPluginAuth(); + if (auth.userEmail) { + // Only when it belongs to the key in hand. Without this check a cached + // email was used forever: switch to a different account and every event + // kept reporting under the previous one, with nothing to notice it by. + if (auth.keyFingerprint === keyFingerprint(apiKey)) { + return createHash("sha256").update(auth.userEmail).digest("hex"); + } + clearResolvedAccount(); + } } catch { // Fall through to the API key or anonymous identity. } @@ -66,8 +85,11 @@ function identifyAnonymous(id: string): void { } function resolveEmail(apiKey: string): void { - if (emailResolutionAttempted) return; - emailResolutionAttempted = true; + // Latched per key, not once per process. A single boolean meant a key changed + // mid-session was never looked up, so the fallback identity stuck until restart. + const fingerprint = keyFingerprint(apiKey); + if (resolutionAttemptedFor === fingerprint) return; + resolutionAttemptedFor = fingerprint; fetch(`${getBaseUrl().replace(/\/+$/, "")}/v1/ping/`, { method: "GET", headers: { Authorization: `Token ${apiKey}`, "Content-Type": "application/json" }, @@ -76,7 +98,7 @@ function resolveEmail(apiKey: string): void { .then((response) => response.json()) .then((data: any) => { if (!data?.user_email) return; - writePluginAuth({ userEmail: data.user_email }); + writePluginAuth({ userEmail: data.user_email, keyFingerprint: fingerprint }); const oldId = createHash("sha256").update(apiKey).digest("hex"); const newId = createHash("sha256").update(data.user_email).digest("hex"); for (const event of telemetry.queueForTesting()) { @@ -96,13 +118,15 @@ export function captureEvent( if (!enabled()) return; try { currentDistinctId = distinctId(context?.apiKey); - let hasEmail = false; + let resolvedForThisKey = false; try { - hasEmail = Boolean(readPluginAuth().userEmail); + const auth = readPluginAuth(); + resolvedForThisKey = + Boolean(auth.userEmail) && auth.keyFingerprint === keyFingerprint(context?.apiKey); } catch { // Resolve it below when possible. } - if (context?.apiKey && !hasEmail) resolveEmail(context.apiKey); + if (context?.apiKey && !resolvedForThisKey) resolveEmail(context.apiKey); identifyAnonymous(currentDistinctId); telemetry.capture(eventName, { mode: context?.mode, diff --git a/integrations/openclaw/tests/telemetry.test.ts b/integrations/openclaw/tests/telemetry.test.ts index 9a7bee716..8c6f6f33c 100644 --- a/integrations/openclaw/tests/telemetry.test.ts +++ b/integrations/openclaw/tests/telemetry.test.ts @@ -3,15 +3,30 @@ import { describe, it, expect, vi, beforeEach, afterEach } from "vitest"; // Mock config-file before importing telemetry vi.mock("../cli/config-file.ts", () => ({ readPluginAuth: vi.fn().mockReturnValue({}), + writePluginAuth: vi.fn(), + clearAnonymousTelemetryId: vi.fn(), + clearResolvedAccount: vi.fn(), + getBaseUrl: vi.fn().mockReturnValue("https://api.mem0.ai"), })); import { captureEvent } from "../telemetry.ts"; -import { readPluginAuth } from "../cli/config-file.ts"; +import { clearResolvedAccount, readPluginAuth } from "../cli/config-file.ts"; + +/** sha256(key).slice(0, 16), the shape telemetry.ts stores. */ +async function fingerprintOf(apiKey: string): Promise { + const { createHash } = await import("node:crypto"); + return createHash("sha256").update(apiKey).digest("hex").slice(0, 16); +} describe("telemetry", () => { let fetchSpy: ReturnType; beforeEach(() => { + // Call history has to be cleared per test, not just restored: the mocks are + // module-level vi.fn()s, so without this one test's calls are visible to the + // next and assertions on "was not called" pass or fail by ordering. + vi.clearAllMocks(); + (readPluginAuth as ReturnType).mockReturnValue({}); // Reset telemetry enabled state (globalThis as any).__mem0_telemetry_override = undefined; fetchSpy = vi.fn().mockResolvedValue({ ok: true }); @@ -52,6 +67,43 @@ describe("telemetry", () => { expect(() => captureEvent("test_event", {}, {})).not.toThrow(); }); + it("keeps using a cached email only while it belongs to the current key", async () => { + (readPluginAuth as ReturnType).mockReturnValue({ + userEmail: "person@example.com", + keyFingerprint: await fingerprintOf("key-a"), + }); + + captureEvent("test_event", {}, { apiKey: "key-a" }); + + expect(clearResolvedAccount).not.toHaveBeenCalled(); + }); + + it("forgets the account when the API key changes", async () => { + // The defect: the cached email was used forever, so events after an account + // switch kept reporting under the previous account. + (readPluginAuth as ReturnType).mockReturnValue({ + userEmail: "person@example.com", + keyFingerprint: await fingerprintOf("key-a"), + }); + + captureEvent("test_event", {}, { apiKey: "key-b" }); + + expect(clearResolvedAccount).toHaveBeenCalled(); + }); + + it("re-resolves for a key it has not looked up before", async () => { + (readPluginAuth as ReturnType).mockReturnValue({ + userEmail: "person@example.com", + keyFingerprint: await fingerprintOf("key-a"), + }); + + captureEvent("test_event", {}, { apiKey: "key-c" }); + + // The resolution latch is per key, not once per process, so a key changed + // mid-session is actually looked up instead of sticking to the fallback. + expect(fetchSpy).toHaveBeenCalled(); + }); + it("handles readPluginAuth errors gracefully", () => { (readPluginAuth as ReturnType).mockImplementationOnce(() => { throw new Error("config read failed"); diff --git a/integrations/opencode-plugin/package.json b/integrations/opencode-plugin/package.json index 79f68c4a8..6722b4c0e 100644 --- a/integrations/opencode-plugin/package.json +++ b/integrations/opencode-plugin/package.json @@ -2,7 +2,7 @@ "name": "@mem0/opencode-plugin", "version": "0.3.0", "type": "module", - "description": "Mem0 persistent memory plugin for OpenCode — add, search, and manage memories across sessions", + "description": "Mem0 persistent memory plugin for OpenCode \u2014 add, search, and manage memories across sessions", "main": "dist/index.js", "types": "index.d.ts", "exports": { @@ -38,6 +38,7 @@ "build": "bun build opencode-mem0.ts --outdir dist --target bun --format esm --entry-naming index.[ext]", "dev": "bun build opencode-mem0.ts --outdir dist --target bun --format esm --entry-naming index.[ext] --watch", "type-check": "tsc --noEmit", + "test": "bun test", "prepack": "bun run build", "postpack": "" }, diff --git a/integrations/opencode-plugin/telemetry.test.ts b/integrations/opencode-plugin/telemetry.test.ts index 393150977..e32629784 100644 --- a/integrations/opencode-plugin/telemetry.test.ts +++ b/integrations/opencode-plugin/telemetry.test.ts @@ -1,3 +1,5 @@ +import { createHash } from "node:crypto"; + import { afterEach, describe, expect, test } from "bun:test"; import { buildEvent, captureEvent, isTelemetryEnabled } from "./telemetry"; @@ -56,9 +58,12 @@ describe("opencode telemetry", () => { expect(typeof props.os_version).toBe("string"); }); - test("project_hash is sha256(projectId) when a project id is supplied", async () => { + test("project_hash is a salted digest of the project id", async () => { + // Previously asserted the bare sha256(projectId), which is the defect: that + // digest is reversible by anyone who can guess a project id. Salted with the + // API key, which is already in play here and is high entropy. const { createHash } = await import("node:crypto"); - const expected = createHash("sha256").update("acme-repo").digest("hex"); + const expected = createHash("sha256").update(`${KEY}:acme-repo`).digest("hex"); const props = buildEvent("session_start", {}, KEY, "acme-repo")! .properties as Record; expect(props.project_hash).toBe(expected); @@ -76,3 +81,38 @@ describe("opencode telemetry", () => { } }); }); + +describe("project_hash salting", () => { + const PROJECT = "my-project"; + + test("is not a bare digest of the project id", () => { + // The defect: an unsalted SHA-256 over a guessable identifier is reversible + // by anyone who can enumerate project ids. + const unsalted = createHash("sha256").update(PROJECT).digest("hex"); + const payload = buildEvent("session_start", {}, KEY, PROJECT) as Record; + + expect(payload.properties.project_hash).toBeDefined(); + expect(payload.properties.project_hash).not.toBe(unsalted); + }); + + test("differs per account for the same project", () => { + const a = buildEvent("session_start", {}, "m0-account-a", PROJECT) as Record; + const b = buildEvent("session_start", {}, "m0-account-b", PROJECT) as Record; + + expect(a.properties.project_hash).not.toBe(b.properties.project_hash); + }); + + test("is stable for one account, so joins still work", () => { + const first = buildEvent("session_start", {}, KEY, PROJECT) as Record; + const second = buildEvent("session_end", {}, KEY, PROJECT) as Record; + + expect(first.properties.project_hash).toBe(second.properties.project_hash); + }); + + test("is omitted rather than unsalted when there is no key", () => { + const payload = buildEvent("session_start", {}, undefined, PROJECT); + + // No key means no event at all, so there is no unsalted hash to leak. + expect(payload).toBeNull(); + }); +}); diff --git a/integrations/opencode-plugin/telemetry.ts b/integrations/opencode-plugin/telemetry.ts index 8d0861d97..98a3bcb84 100644 --- a/integrations/opencode-plugin/telemetry.ts +++ b/integrations/opencode-plugin/telemetry.ts @@ -40,8 +40,23 @@ function distinctId(apiKey: string): string { return createHash("sha256").update(apiKey).digest("hex").slice(0, 32); } -function projectHash(projectId?: string): Record { - return projectId ? { project_hash: createHash("sha256").update(projectId).digest("hex") } : {}; +/** + * Salted so the hash is not enumerable. + * + * An unsalted SHA-256 of a project id is reversible by anyone who can guess the + * id, which for a project identifier is a small space. The API key is the salt: + * it is already in play here (distinctId is a digest of it), it is high entropy, + * and using it needs no per-install file and so no write race to get wrong. The + * hash is therefore per account rather than per machine, which also keeps joins + * working for one user across machines. It resets when the key rotates, which is + * consistent, because distinctId resets with it. + * + * Both are omitted without a key. An event cannot be built without a distinctId + * anyway, so this costs nothing. + */ +function projectHash(projectId?: string, apiKey?: string): Record { + if (!projectId || !apiKey) return {}; + return { project_hash: createHash("sha256").update(`${apiKey}:${projectId}`).digest("hex") }; } export function buildEvent( @@ -51,7 +66,7 @@ export function buildEvent( projectId?: string, ): Record | null { currentDistinctId = apiKey ? distinctId(apiKey) : ""; - const event = telemetry.build(eventType, { ...properties, ...projectHash(projectId) }); + const event = telemetry.build(eventType, { ...properties, ...projectHash(projectId, apiKey) }); return event ? { api_key: POSTHOG_API_KEY, ...event } : null; } @@ -62,5 +77,5 @@ export function captureEvent( projectId?: string, ): void { currentDistinctId = apiKey ? distinctId(apiKey) : ""; - telemetry.capture(eventType, { ...properties, ...projectHash(projectId) }); + telemetry.capture(eventType, { ...properties, ...projectHash(projectId, apiKey) }); }