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) }); }