fix(plugins): flush on exit regardless of backoff, and prefer the backlog over new events
Findings from an independent review of this branch. The exit-time flush was gated by its own cooldown. beforeExit called the same flush() that opens with a retryNotBefore check, so after any failed delivery a process exiting inside the 2s to 60s window sent nothing and the queue died with it. That is precisely the loss this branch exists to stop, and timer.unref makes beforeExit often the only remaining chance. flush(force) now skips the cooldown and beforeExit passes it. Overflow kept the newest and evicted the batch being retried, which threw away exactly the events the retry exists to save. Both the failure path and capture() now keep the backlog and drop the new event instead, matching Python's record(), which refuses new events once the spool is full. That change needs a bound, so delivery now gives up after five attempts, as the Python core does. Without one a payload the server will never accept would be retried for the whole session and, with the backlog now preferred, would hold the queue against everything behind it. Events carry a capture-time timestamp. They now sit through backoff and across an entire outage, so without one PostHog records them at whatever moment delivery happened to succeed. It also matters for the uuid dedupe, whose key includes the event date. openclaw wiped a resolved account on any capture without an apiKey: the fingerprint comparison was `undefined === ""`, so a keyless call looked like a key change. Guarded on a real key being present. Masked today because every call site supplies one, which is why only a test found it. Two smaller ones. opencode's PostHog source was "plugin", which named no particular plugin and matched no vocabulary; it is now OPENCODE_PLUGIN like every other surface, and a saved insight filtering source = "plugin" needs repointing. And an em dash in opencode's published description had been rewritten to a — escape by my own json.dumps when adding the test script; restored. One openclaw test asserted only not.toThrow() under a name claiming it checked the identity, and under the fingerprint gate the path it exercised no longer uses the email at all. It now pins the real condition. core 35, openclaw 432, opencode 35, pi-agent 89, deepseek 47. The wire e2e still passes 9 of 9 against a real local server, and the identity e2e 7 of 7. Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb
This commit is contained in:
@@ -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<string, unknown>[] = [];
|
||||
@@ -109,9 +114,13 @@ export function createTelemetry(config: TelemetryConfig) {
|
||||
if (!response.ok) throw new Error(`posthog responded ${response.status}`);
|
||||
});
|
||||
|
||||
async function flush(): Promise<void> {
|
||||
async function flush(force = false): Promise<void> {
|
||||
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<string, unknown> = {}): Record<string, unknown> | 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?.();
|
||||
|
||||
@@ -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<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();
|
||||
});
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -55,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", () => {
|
||||
|
||||
@@ -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": {
|
||||
|
||||
@@ -15,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);
|
||||
@@ -32,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)", () => {
|
||||
|
||||
@@ -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}`,
|
||||
|
||||
Reference in New Issue
Block a user