Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 8fabc87f6e | |||
| 7900a5d8d1 | |||
| cbf97c6077 | |||
| 44c6a07ddf | |||
| 6f207d2e28 | |||
| c05556f3cf |
@@ -32,6 +32,9 @@ jobs:
|
||||
- name: Type check
|
||||
run: bun run type-check
|
||||
|
||||
- name: Test
|
||||
run: bun test
|
||||
|
||||
- name: Build
|
||||
run: bun run build
|
||||
|
||||
|
||||
@@ -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<string, unknown>[] = [];
|
||||
let timer: ReturnType<typeof setInterval> | 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<string, unknown>[]) => {
|
||||
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<void> {
|
||||
async function flush(force = false): Promise<void> {
|
||||
// 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<string, unknown> = {}): Record<string, unknown> | 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);
|
||||
|
||||
@@ -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<string, unknown>[][] = [];
|
||||
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<string, unknown>[][] = [];
|
||||
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<string, unknown>[][] = [];
|
||||
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();
|
||||
});
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -148,6 +148,7 @@ const ALLOWED_KEYS = [
|
||||
"mode",
|
||||
"apiKey",
|
||||
"anonymousTelemetryId",
|
||||
"keyFingerprint",
|
||||
"baseUrl",
|
||||
"userId",
|
||||
"userEmail",
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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");
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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<string> {
|
||||
const { createHash } = await import("node:crypto");
|
||||
return createHash("sha256").update(apiKey).digest("hex").slice(0, 16);
|
||||
}
|
||||
|
||||
describe("telemetry", () => {
|
||||
let fetchSpy: ReturnType<typeof vi.fn>;
|
||||
|
||||
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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).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<typeof vi.fn>).mockImplementationOnce(() => {
|
||||
throw new Error("config read failed");
|
||||
|
||||
@@ -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": ""
|
||||
},
|
||||
|
||||
@@ -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<string, unknown>;
|
||||
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<string, unknown>;
|
||||
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<string, unknown>;
|
||||
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<string, any>;
|
||||
|
||||
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<string, any>;
|
||||
const b = buildEvent("session_start", {}, "m0-account-b", PROJECT) as Record<string, any>;
|
||||
|
||||
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<string, any>;
|
||||
const second = buildEvent("session_end", {}, KEY, PROJECT) as Record<string, any>;
|
||||
|
||||
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();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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<string, string> {
|
||||
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<string, string> {
|
||||
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<string, unknown> | 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) });
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user