From a8d36343128c0410a3e71626b8edfd31f875bd17 Mon Sep 17 00:00:00 2001 From: Saket Aryan Date: Fri, 18 Sep 2026 14:13:22 +0530 Subject: [PATCH] fix(plugins): stop the TypeScript telemetry losing events and misattributing accounts (#7358) --- .github/workflows/opencode-plugin-checks.yml | 3 + .../typescript/src/telemetry.ts | 99 ++++++- .../typescript/tests/telemetry.test.ts | 247 ++++++++++++++++++ integrations/openclaw/cli/config-file.ts | 20 ++ integrations/openclaw/config.ts | 1 + integrations/openclaw/openclaw.plugin.json | 4 + integrations/openclaw/telemetry.ts | 64 ++++- .../openclaw/tests/config-file.test.ts | 40 +++ integrations/openclaw/tests/config.test.ts | 15 ++ integrations/openclaw/tests/telemetry.test.ts | 80 +++++- integrations/opencode-plugin/package.json | 1 + .../opencode-plugin/telemetry.test.ts | 51 +++- integrations/opencode-plugin/telemetry.ts | 27 +- 13 files changed, 623 insertions(+), 29 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..2b6c5286e 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,34 +76,108 @@ 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; +// Consecutive failed flushes before the queue is dropped. Deliberately NOT the +// same thing as Python's budget, which rides in the claim filename and so +// follows one batch: this counter lives in the closure and counts the outage, +// not the payload. Events captured between attempts join the same queue and go +// with it. Per-batch accounting would need an attempt count on every event, and +// the queue is already bounded, so the simpler rule is the one in force here. +// Without any bound a payload the server will never accept is retried for the +// whole session and, now that the backlog is preferred over new events, holds +// the queue against everything behind it. +const MAX_DELIVERY_ATTEMPTS = 5; + export function createTelemetry(config: TelemetryConfig) { let queue: Record[] = []; let timer: ReturnType | undefined; + let consecutiveFailures = 0; + let retryNotBefore = 0; + let exitFlushAttempted = false; + let flushing = false; const flushThreshold = config.flushThreshold ?? 10; const maxQueueSize = config.maxQueueSize ?? 100; const deliver = config.delivery ?? (async (batch: Record[]) => { - await fetch(POSTHOG_BATCH_URL, { + const response = await fetch(POSTHOG_BATCH_URL, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ api_key: POSTHOG_API_KEY, batch }), signal: AbortSignal.timeout(3_000), }); + // fetch only rejects on a network-level failure. Without this check a 500, + // a 503 or a 429 resolved normally and the batch was counted as delivered + // and dropped, which is the likelier outage than a refused connection. + // Any non-2xx is retried, matching the Python core: the backoff and the + // queue bound contain a payload that will never be accepted, because the + // re-queued batch sits at the front and is the first thing evicted. + if (!response.ok) throw new Error(`posthog responded ${response.status}`); }); - async function flush(): Promise { + async function flush(force = false): Promise { + // One at a time. Two overlapping flushes each detach the queue and each + // prepend their own batch back on failure, so the later batch lands in front + // of the earlier one and the truncation then drops the OLDER events first, + // inverting the priority the failure path exists to establish. A second + // caller returns immediately; the queue waits for the next flush. + if (flushing) return; if (!queue.length) return; + // `force` skips the cooldown. beforeExit is the last chance this process + // gets, and gating it on the same backoff meant that after any failure the + // exit flush did nothing and the queue died with the process, which is the + // loss this whole mechanism exists to prevent. + if (!force && Date.now() < retryNotBefore) return; const batch = queue; queue = []; + flushing = true; 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. + // + consecutiveFailures += 1; + if (consecutiveFailures >= MAX_DELIVERY_ATTEMPTS) { + // Give up on the queue so a failing outage cannot hold it for the + // session. This drops whatever is queued now, which includes events + // captured during the outage, not only the batch that kept failing. + consecutiveFailures = 0; + retryNotBefore = 0; + return; + } + // Keep the FRONT on overflow, so the batch being retried survives and a + // new event is what gets dropped. Matches the Python core, where record() + // refuses new events once the spool is full rather than evicting the + // backlog. Keeping the newest would throw away exactly the events this + // retry exists to save. + queue = [...batch, ...queue].slice(0, maxQueueSize); + retryNotBefore = Date.now() + Math.min(2 ** consecutiveFailures * 1_000, RETRY_BACKOFF_CEILING_MS); + } finally { + flushing = false; } } function beforeExit(): void { - void flush(); + // Once, and only once. Node re-emits beforeExit whenever the handler + // schedules more async work, so an unconditional forced flush looped until + // the attempt budget was spent: five attempts against a 3s delivery timeout + // is fifteen seconds added to the shutdown of whatever editor or CLI is + // hosting this. The backoff used to end that loop after one attempt, and + // removing it for the forced path removed the only thing bounding it. + if (exitFlushAttempted) return; + exitFlushAttempted = true; + void flush(true); } function build(event: string, properties: Record = {}): Record | null { @@ -112,6 +188,15 @@ 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(), + // Capture time, not ingestion time. Events now sit through backoff and + // across a whole outage, so without this PostHog records them whenever + // delivery happened to succeed. It also matters for the uuid dedupe + // above, whose key includes the event date. + timestamp: new Date().toISOString(), properties: { ...safeProperties(properties), ...safeProperties(config.commonProperties ?? {}), @@ -134,8 +219,10 @@ export function createTelemetry(config: TelemetryConfig) { try { const payload = build(event, properties); if (!payload) return; + // Full means drop this event, not evict the backlog. Same rule as the + // failure path above and as Python's record(). + if (queue.length >= maxQueueSize) return; queue.push(payload); - if (queue.length > maxQueueSize) queue = queue.slice(-maxQueueSize); if (!timer) { timer = setInterval(() => void flush(), config.flushIntervalMs ?? 5_000); timer.unref?.(); @@ -149,6 +236,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..cf9671682 100644 --- a/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts +++ b/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts @@ -122,3 +122,250 @@ 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 full queue drops the new event and keeps the batch being retried", async () => { + // Python's record() refuses new events once the spool is full rather than + // evicting the backlog. Keeping the newest here would throw away exactly the + // events the retry exists to save. + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", + flushThreshold: 1000, maxQueueSize: 3, + delivery: async () => { throw new Error("down"); }, + }); + + // Fill past the cap BEFORE the flush, so the re-queue actually has to truncate. + // Capturing only two left the queue empty at re-queue time and the slice on the + // failure path never ran, which is the half that decides the direction. + telemetry.capture("a"); + telemetry.capture("b"); + telemetry.capture("c"); + await telemetry.flush(); + telemetry.capture("d"); + telemetry.capture("e"); + + const events = telemetry.queueForTesting().map((e) => (e as any).event); + assert.equal(events.length, 3, "queue grew past maxQueueSize"); + assert.deepEqual(events, ["a", "b", "c"], "the retried batch was evicted instead of the new events"); + telemetry.resetForTesting(); +}); + +test("the exit-time flush ignores the backoff", async () => { + // beforeExit is the last chance the process gets. Gating it on the same + // cooldown meant that after any failure it did nothing and the queue died. + let attempts = 0; + let failing = true; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + delivery: async () => { attempts += 1; if (failing) throw new Error("down"); }, + }); + + telemetry.capture("a"); + await telemetry.flush(); + assert.equal(attempts, 1); + + failing = false; + await telemetry.flush(); + assert.equal(attempts, 1, "the backoff should still hold for an ordinary flush"); + + await telemetry.flush(true); + assert.equal(attempts, 2, "the exit flush was suppressed by the backoff"); + assert.equal(telemetry.queueForTesting().length, 0); + telemetry.resetForTesting(); +}); + +test("every event carries a capture-time timestamp", async () => { + const sent: Record[][] = []; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + delivery: async (batch) => { sent.push(batch); }, + }); + + telemetry.capture("a"); + const capturedAt = Date.now(); + await new Promise((resolve) => setTimeout(resolve, 50)); + await telemetry.flush(); + + const stamped = sent[0][0].timestamp as string; + assert.ok(stamped, "no timestamp, so PostHog would record delivery time"); + assert.ok(Math.abs(Date.parse(stamped) - capturedAt) < 1_000, "not capture time"); + telemetry.resetForTesting(); +}); + +test("a batch the server will never accept is eventually given up on", async () => { + let attempts = 0; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + delivery: async () => { attempts += 1; throw new Error("permanently bad"); }, + }); + + telemetry.capture("doomed"); + for (let i = 0; i < 8; i += 1) await telemetry.flush(true); + + assert.ok(attempts <= 6, `retried ${attempts} times with no cap`); + assert.equal(telemetry.queueForTesting().length, 0, "a doomed batch held the queue forever"); + telemetry.resetForTesting(); +}); + +test("an HTTP error response is a failure, not a delivery", async () => { + // fetch only rejects on a network-level failure, so a 500 used to resolve + // normally and the batch was dropped as delivered. Exercises the real default + // delivery path rather than an injected one, which is where this hid. + const realFetch = globalThis.fetch; + let calls = 0; + globalThis.fetch = (async () => { + calls += 1; + return new Response("upstream is unwell", { status: 503 }); + }) as typeof fetch; + + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + }); + try { + telemetry.capture("during.outage"); + await telemetry.flush(); + + assert.equal(calls, 1, "never reached the network"); + assert.equal(telemetry.queueForTesting().length, 1, "a 503 was counted as delivered"); + } finally { + globalThis.fetch = realFetch; + telemetry.resetForTesting(); + } +}); + +test("a 2xx is a delivery", async () => { + const realFetch = globalThis.fetch; + globalThis.fetch = (async () => new Response("ok", { status: 200 })) as typeof fetch; + + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + }); + try { + telemetry.capture("fine"); + await telemetry.flush(); + assert.equal(telemetry.queueForTesting().length, 0, "a good response did not clear the queue"); + } finally { + globalThis.fetch = realFetch; + telemetry.resetForTesting(); + } +}); + +test("the exit flush is attempted once, not until the budget is spent", async () => { + // Node re-emits beforeExit whenever the handler schedules async work, so an + // unconditional forced flush looped until MAX_DELIVERY_ATTEMPTS. Against the + // real 3s delivery timeout that is fifteen seconds added to a host's shutdown. + let attempts = 0; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + delivery: async () => { attempts += 1; throw new Error("down"); }, + }); + + telemetry.capture("a"); + const handlers = process.listeners("beforeExit"); + const ours = handlers[handlers.length - 1] as () => void; + ours(); + ours(); + ours(); + await new Promise((resolve) => setTimeout(resolve, 20)); + + assert.equal(attempts, 1, `exit flush ran ${attempts} times`); + telemetry.resetForTesting(); +}); + +test("overlapping flushes do not reorder the backlog behind newer events", async () => { + // Each flush detaches the queue and prepends its own batch back on failure, so + // two in flight at once put the LATER batch in front of the earlier one. The + // truncation then drops the older events first, inverting the priority the + // failure path exists to establish. + let release: (() => void)[] = []; + const telemetry = createTelemetry({ + host: "h", source: "S", version: "1", distinctId: "d", flushThreshold: 1000, + delivery: () => new Promise((_resolve, reject) => { release.push(() => reject(new Error("down"))); }), + }); + + telemetry.capture("first"); + const a = telemetry.flush(); + telemetry.capture("second"); + const b = telemetry.flush(); + + release.forEach((fn) => fn()); + await Promise.all([a, b]); + + const events = telemetry.queueForTesting().map((e) => (e as any).event); + assert.equal(release.length, 1, "a second delivery started while one was in flight"); + assert.deepEqual(events, ["first", "second"], `backlog reordered: ${events.join(",")}`); + telemetry.resetForTesting(); +}); diff --git a/integrations/openclaw/cli/config-file.ts b/integrations/openclaw/cli/config-file.ts index 99453cb1b..a34d6c08b 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; } // ============================================================================ @@ -135,6 +137,9 @@ export function readPluginAuth(): PluginAuthConfig { autoCapture: cfg.autoCapture as boolean | undefined, topK: cfg.topK as number | undefined, anonymousTelemetryId: cfg.anonymousTelemetryId as string | undefined, + // Without this the reader silently drops it, every fingerprint comparison + // fails against undefined, and the resolved email is never used again. + keyFingerprint: cfg.keyFingerprint as string | undefined, }; } @@ -236,6 +241,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/config.ts b/integrations/openclaw/config.ts index 4e43c8c44..09dcbc837 100644 --- a/integrations/openclaw/config.ts +++ b/integrations/openclaw/config.ts @@ -148,6 +148,7 @@ const ALLOWED_KEYS = [ "mode", "apiKey", "anonymousTelemetryId", + "keyFingerprint", "baseUrl", "userId", "userEmail", diff --git a/integrations/openclaw/openclaw.plugin.json b/integrations/openclaw/openclaw.plugin.json index 2d9f46518..02d05be4b 100644 --- a/integrations/openclaw/openclaw.plugin.json +++ b/integrations/openclaw/openclaw.plugin.json @@ -209,6 +209,10 @@ "type": "string", "description": "Persistent anonymous telemetry identifier" }, + "keyFingerprint": { + "type": "string", + "description": "Digest of the API key userEmail was resolved for. Set automatically; a mismatch re-resolves the account." + }, "oss": { "type": "object", "properties": { diff --git a/integrations/openclaw/telemetry.ts b/integrations/openclaw/telemetry.ts index fc1c6b99d..bd77e1656 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,35 @@ 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"); + } + // Only a REAL key that disagrees means the account changed. Without the + // apiKey guard the comparison is `undefined === ""` for any call that + // simply omits the key, so a capture with no context wiped a perfectly + // good account out of openclaw.json. + // + // A row with an email and NO fingerprint is the legacy shape, from an + // install predating this field. Clearing it here deleted a real account + // before anything had replaced it, and if the re-resolve then failed + // because the user was offline the email was gone from disk for good. The + // Python core refuses the same trade: verify, and keep what you have until + // the verification succeeds. resolveEmail below overwrites both fields + // when it does, so there is nothing to clear first. + if (apiKey && auth.keyFingerprint) clearResolvedAccount(); + } } catch { // Fall through to the API key or anonymous identity. } @@ -66,8 +97,16 @@ 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; + const releaseLatch = () => { + // A failed lookup must not pin the fallback identity for the rest of the + // process. Released so the next capture tries again. + if (resolutionAttemptedFor === fingerprint) resolutionAttemptedFor = ""; + }; fetch(`${getBaseUrl().replace(/\/+$/, "")}/v1/ping/`, { method: "GET", headers: { Authorization: `Token ${apiKey}`, "Content-Type": "application/json" }, @@ -76,7 +115,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()) { @@ -84,7 +123,8 @@ function resolveEmail(apiKey: string): void { } }) .catch(() => { - // The API-key hash remains a stable fallback. + // The API-key hash remains a stable fallback, and the next capture retries. + releaseLatch(); }); } @@ -96,13 +136,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/config-file.test.ts b/integrations/openclaw/tests/config-file.test.ts index e059ddff3..ba63180cb 100644 --- a/integrations/openclaw/tests/config-file.test.ts +++ b/integrations/openclaw/tests/config-file.test.ts @@ -237,3 +237,43 @@ describe("getBaseUrl", () => { expect(getBaseUrl()).toBe(DEFAULT_BASE_URL); }); }); + +// --------------------------------------------------------------------------- +// keyFingerprint round trip +// --------------------------------------------------------------------------- + +describe("keyFingerprint survives a write and read", () => { + it("readPluginAuth returns a persisted keyFingerprint", () => { + // It did not. readPluginAuth builds its result field by field, and this one + // was missing, so every fingerprint comparison ran against undefined, the + // resolved email was never used again, and telemetry silently fell back to + // the API key hash. The telemetry tests could not see it because they mock + // this module and their mock returned the field the real reader dropped. + setConfigFile({ + plugins: { + entries: { + "openclaw-mem0": { + config: { + apiKey: "m0-test", + userEmail: "person@example.com", + keyFingerprint: "0123456789abcdef", + }, + }, + }, + }, + }); + + expect(readPluginAuth().keyFingerprint).toBe("0123456789abcdef"); + }); + + it("writePluginAuth persists it where readPluginAuth looks", () => { + setConfigFile({ plugins: { entries: { "openclaw-mem0": { config: {} } } } }); + + writePluginAuth({ userEmail: "person@example.com", keyFingerprint: "abc123" }); + + const written = JSON.parse(mockWriteText.mock.calls.at(-1)![1] as string); + const cfg = written.plugins.entries["openclaw-mem0"].config; + expect(cfg.keyFingerprint).toBe("abc123"); + expect(cfg.userEmail).toBe("person@example.com"); + }); +}); diff --git a/integrations/openclaw/tests/config.test.ts b/integrations/openclaw/tests/config.test.ts index 2c6dbac2e..334c064fb 100644 --- a/integrations/openclaw/tests/config.test.ts +++ b/integrations/openclaw/tests/config.test.ts @@ -485,3 +485,18 @@ describe("mem0ConfigSchema.parse() — apiKey edge cases", () => { expect(cfg.needsSetup).toBe(false); }); }); + +describe("telemetry fingerprint round trip", () => { + it("a config carrying keyFingerprint is accepted by the real schema", () => { + // What writePluginAuth persists after a successful lookup. It was not in + // ALLOWED_KEYS, and assertAllowedKeys throws, so the first successful + // resolve wrote a config that broke every subsequent load of the plugin. + const persisted = { + apiKey: "m0-test", + userEmail: "person@example.com", + keyFingerprint: "0123456789abcdef", + }; + + expect(() => mem0ConfigSchema.parse(persisted)).not.toThrow(); + }); +}); diff --git a/integrations/openclaw/tests/telemetry.test.ts b/integrations/openclaw/tests/telemetry.test.ts index 9a7bee716..ad7a9bbe0 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 }); @@ -40,11 +55,31 @@ describe("telemetry", () => { expect(() => captureEvent("test_event")).not.toThrow(); }); - it("uses userEmail as distinct ID when available", () => { - (readPluginAuth as ReturnType).mockReturnValueOnce({ + it("uses userEmail as distinct ID when it belongs to the current key", async () => { + // Previously asserted only not.toThrow(), which passed whatever the identity + // turned out to be, and under the fingerprint gate the no-context path does + // not use the email at all. Pin the real condition instead. + (readPluginAuth as ReturnType).mockReturnValue({ userEmail: "test@example.com", + keyFingerprint: await fingerprintOf("key-a"), }); - expect(() => captureEvent("test_event")).not.toThrow(); + + captureEvent("test_event", {}, { apiKey: "key-a" }); + + expect(clearResolvedAccount).not.toHaveBeenCalled(); + }); + + it("a capture with no apiKey leaves a resolved account alone", async () => { + // `undefined === ""` made every keyless capture look like a key change, so + // one context-free call wiped a good account out of openclaw.json. + (readPluginAuth as ReturnType).mockReturnValue({ + userEmail: "test@example.com", + keyFingerprint: await fingerprintOf("key-a"), + }); + + captureEvent("test_event"); + + expect(clearResolvedAccount).not.toHaveBeenCalled(); }); it("falls back to a generated anonymous id when no apiKey", () => { @@ -52,6 +87,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..acdbaf469 100644 --- a/integrations/opencode-plugin/package.json +++ b/integrations/opencode-plugin/package.json @@ -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..79dc3ffdb 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"; @@ -13,7 +15,10 @@ describe("opencode telemetry", () => { expect(payload).not.toBeNull(); const props = payload!.properties as Record; expect(payload!.event).toBe("plugin.session_start"); - expect(props.source).toBe("plugin"); + // Was "plugin", which named no particular plugin and matched no vocabulary. + // Now shaped like every other surface. Any saved PostHog insight filtering + // source = "plugin" needs repointing; historical data is untouched. + expect(props.source).toBe("OPENCODE_PLUGIN"); expect(props.platform).toBe("opencode"); expect(props.memory_count).toBe(5); expect(props.$process_person_profile).toBe(false); @@ -30,7 +35,7 @@ describe("opencode telemetry", () => { const props = buildEvent("x", { platform: "HACK", source: "HACK" }, KEY)! .properties as Record; expect(props.platform).toBe("opencode"); - expect(props.source).toBe("plugin"); + expect(props.source).toBe("OPENCODE_PLUGIN"); }); test("returns null without an API key (no anonymous events)", () => { @@ -56,9 +61,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 +84,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..e2316dd88 100644 --- a/integrations/opencode-plugin/telemetry.ts +++ b/integrations/opencode-plugin/telemetry.ts @@ -25,7 +25,9 @@ const PLUGIN_VERSION = (() => { let currentDistinctId = ""; const telemetry = createTelemetry({ host: "opencode", - source: "plugin", + // Shaped like the platform's EventSource values, as every other surface is. + // "plugin" said nothing about which plugin and matched no vocabulary. + source: "OPENCODE_PLUGIN", version: PLUGIN_VERSION, distinctId: () => currentDistinctId, eventName: (event) => `plugin.${event}`, @@ -40,8 +42,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 +68,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 +79,5 @@ export function captureEvent( projectId?: string, ): void { currentDistinctId = apiKey ? distinctId(apiKey) : ""; - telemetry.capture(eventType, { ...properties, ...projectHash(projectId) }); + telemetry.capture(eventType, { ...properties, ...projectHash(projectId, apiKey) }); }