diff --git a/integrations/agent-plugin-core/typescript/src/telemetry.ts b/integrations/agent-plugin-core/typescript/src/telemetry.ts index 0beffe6b9..cfaa73034 100644 --- a/integrations/agent-plugin-core/typescript/src/telemetry.ts +++ b/integrations/agent-plugin-core/typescript/src/telemetry.ts @@ -84,6 +84,11 @@ export function errorKind(error: unknown): string { // 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; +// Attempts before a batch is given up on, mirroring the Python core's budget. +// Without one, a payload the server will never accept is retried for the whole +// session and, now that the backlog is preferred over new events, would block +// everything behind it. +const MAX_DELIVERY_ATTEMPTS = 5; export function createTelemetry(config: TelemetryConfig) { let queue: Record[] = []; @@ -109,9 +114,13 @@ export function createTelemetry(config: TelemetryConfig) { if (!response.ok) throw new Error(`posthog responded ${response.status}`); }); - async function flush(): Promise { + async function flush(force = false): Promise { if (!queue.length) return; - if (Date.now() < retryNotBefore) 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 = []; try { @@ -124,16 +133,25 @@ export function createTelemetry(config: TelemetryConfig) { // 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; + if (consecutiveFailures >= MAX_DELIVERY_ATTEMPTS) { + // Give up on this batch so it cannot hold the queue for the session. + 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); } } function beforeExit(): void { - void flush(); + void flush(true); } function build(event: string, properties: Record = {}): Record | null { @@ -148,6 +166,11 @@ export function createTelemetry(config: TelemetryConfig) { // 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 ?? {}), @@ -170,8 +193,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?.(); diff --git a/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts b/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts index 5a8ca0439..4ca80b6fb 100644 --- a/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts +++ b/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts @@ -197,7 +197,10 @@ test("repeated failures back off instead of retrying every flush", async () => { telemetry.resetForTesting(); }); -test("a long outage costs the oldest events, not unbounded memory", async () => { +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, @@ -211,7 +214,66 @@ test("a long outage costs the oldest events, not unbounded memory", async () => telemetry.capture("d"); telemetry.capture("e"); - assert.ok(telemetry.queueForTesting().length <= 3, "queue grew past maxQueueSize"); + const events = telemetry.queueForTesting().map((e) => (e as any).event); + assert.ok(events.length <= 3, "queue grew past maxQueueSize"); + assert.deepEqual(events.slice(0, 2), ["a", "b"], "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(); }); diff --git a/integrations/openclaw/telemetry.ts b/integrations/openclaw/telemetry.ts index f0f3c3fc7..7b4c4879f 100644 --- a/integrations/openclaw/telemetry.ts +++ b/integrations/openclaw/telemetry.ts @@ -54,7 +54,11 @@ function distinctId(apiKey?: string): string { if (auth.keyFingerprint === keyFingerprint(apiKey)) { return createHash("sha256").update(auth.userEmail).digest("hex"); } - clearResolvedAccount(); + // 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. + if (apiKey) clearResolvedAccount(); } } catch { // Fall through to the API key or anonymous identity. diff --git a/integrations/openclaw/tests/telemetry.test.ts b/integrations/openclaw/tests/telemetry.test.ts index 8c6f6f33c..ad7a9bbe0 100644 --- a/integrations/openclaw/tests/telemetry.test.ts +++ b/integrations/openclaw/tests/telemetry.test.ts @@ -55,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", () => { diff --git a/integrations/opencode-plugin/package.json b/integrations/opencode-plugin/package.json index 6722b4c0e..acdbaf469 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 \u2014 add, search, and manage memories across sessions", + "description": "Mem0 persistent memory plugin for OpenCode — add, search, and manage memories across sessions", "main": "dist/index.js", "types": "index.d.ts", "exports": { diff --git a/integrations/opencode-plugin/telemetry.test.ts b/integrations/opencode-plugin/telemetry.test.ts index e32629784..79dc3ffdb 100644 --- a/integrations/opencode-plugin/telemetry.test.ts +++ b/integrations/opencode-plugin/telemetry.test.ts @@ -15,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); @@ -32,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)", () => { diff --git a/integrations/opencode-plugin/telemetry.ts b/integrations/opencode-plugin/telemetry.ts index 98a3bcb84..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}`,