diff --git a/integrations/agent-plugin-core/typescript/src/telemetry.ts b/integrations/agent-plugin-core/typescript/src/telemetry.ts index 51c546403..2b6c5286e 100644 --- a/integrations/agent-plugin-core/typescript/src/telemetry.ts +++ b/integrations/agent-plugin-core/typescript/src/telemetry.ts @@ -101,6 +101,7 @@ export function createTelemetry(config: TelemetryConfig) { let consecutiveFailures = 0; let retryNotBefore = 0; let exitFlushAttempted = false; + let flushing = false; const flushThreshold = config.flushThreshold ?? 10; const maxQueueSize = config.maxQueueSize ?? 100; @@ -121,6 +122,12 @@ export function createTelemetry(config: TelemetryConfig) { }); 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 @@ -129,6 +136,7 @@ export function createTelemetry(config: TelemetryConfig) { if (!force && Date.now() < retryNotBefore) return; const batch = queue; queue = []; + flushing = true; try { await deliver(batch); consecutiveFailures = 0; @@ -155,6 +163,8 @@ export function createTelemetry(config: TelemetryConfig) { // 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; } } diff --git a/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts b/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts index 576163f1e..cf9671682 100644 --- a/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts +++ b/integrations/agent-plugin-core/typescript/tests/telemetry.test.ts @@ -344,3 +344,28 @@ test("the exit flush is attempted once, not until the budget is spent", async () 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 ae4d9e471..a34d6c08b 100644 --- a/integrations/openclaw/cli/config-file.ts +++ b/integrations/openclaw/cli/config-file.ts @@ -137,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, }; } 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/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(); + }); +});