refactor(telemetry): tighten alias helpers and gate on MEM0_TELEMETRY
Review pass on the PostHog identity-stitching fix:
- Gate _maybe_alias_anon_to_email and _maybeAliasAnonToEmail on telemetry
being enabled. Previous version still read ~/.mem0/config.json and
persisted telemetry.aliased_to even when MEM0_TELEMETRY=false, which
meant re-enabling telemetry later would skip identity merging forever.
Python checks client_telemetry.posthog; TS exports a new
isTelemetryEnabled() getter from telemetry.ts.
- Trim capture_identify properties to {$anon_distinct_id, client_source}.
The full platform metadata duplicated capture_event's payload and is
not needed for $identify; PostHog enriches the person profile from
subsequent events anyway.
- Flatten config.ts: drop the four-helper structure (loadNodeModules /
configPath / loadRawConfig / read+write) into one getNodeFs() + a
shared loadConfig(). Same behaviour, ~30 fewer lines, easier to read.
- Drop the tautological Mem0AnonIds shape test and the
void createMockFetch hack at the end of the TS test file.
- Add tests for the new telemetry-disabled gate (Python + TS).
Logic traced end-to-end: 114 Python tests + 590 TS tests pass. Ruff and
Prettier clean. Zero new typecheck errors.
This commit is contained in:
@@ -2,12 +2,10 @@
|
||||
* Best-effort read/write of ~/.mem0/config.json from the TS SDK.
|
||||
*
|
||||
* Used to stitch PostHog identities: the OSS Python SDK and the Python CLI
|
||||
* each persist an anonymous distinct_id here, and the TS MemoryClient needs
|
||||
* to read those on init so it can fire $identify and merge them into the
|
||||
* email-based identity.
|
||||
* each persist an anonymous distinct_id here, and the TS MemoryClient reads
|
||||
* those on init to fire $identify and merge them into the email identity.
|
||||
*
|
||||
* Node-only. In browsers (or any environment without `process.versions.node`)
|
||||
* every function returns null/no-ops without attempting `fs` access.
|
||||
* Node-only. Browsers (no `process.versions.node`) no-op.
|
||||
*/
|
||||
|
||||
export interface Mem0AnonIds {
|
||||
@@ -16,62 +14,38 @@ export interface Mem0AnonIds {
|
||||
aliasedTo?: string;
|
||||
}
|
||||
|
||||
function isNode(): boolean {
|
||||
try {
|
||||
return (
|
||||
typeof process !== "undefined" &&
|
||||
!!process.versions &&
|
||||
!!process.versions.node
|
||||
);
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
async function loadNodeModules(): Promise<{
|
||||
interface NodeFs {
|
||||
fs: typeof import("fs");
|
||||
path: typeof import("path");
|
||||
os: typeof import("os");
|
||||
} | null> {
|
||||
if (!isNode()) return null;
|
||||
configPath: string;
|
||||
}
|
||||
|
||||
async function getNodeFs(): Promise<NodeFs | null> {
|
||||
if (typeof process === "undefined" || !process.versions?.node) return null;
|
||||
try {
|
||||
const [fs, path, os] = await Promise.all([
|
||||
import("fs"),
|
||||
import("path"),
|
||||
import("os"),
|
||||
]);
|
||||
const fsMod = (fs as any).default ?? fs;
|
||||
const pathMod = (path as any).default ?? path;
|
||||
const osMod = (os as any).default ?? os;
|
||||
const dir = process.env.MEM0_DIR || pathMod.join(osMod.homedir(), ".mem0");
|
||||
return {
|
||||
fs: fs.default ?? fs,
|
||||
path: path.default ?? path,
|
||||
os: os.default ?? os,
|
||||
fs: fsMod,
|
||||
path: pathMod,
|
||||
configPath: pathMod.join(dir, "config.json"),
|
||||
};
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
async function configPath(): Promise<{
|
||||
path: string;
|
||||
modules: NonNullable<Awaited<ReturnType<typeof loadNodeModules>>>;
|
||||
} | null> {
|
||||
const modules = await loadNodeModules();
|
||||
if (!modules) return null;
|
||||
function loadConfig(node: NodeFs): Record<string, any> | null {
|
||||
try {
|
||||
const dir =
|
||||
process.env.MEM0_DIR || modules.path.join(modules.os.homedir(), ".mem0");
|
||||
return { path: modules.path.join(dir, "config.json"), modules };
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
async function loadRawConfig(): Promise<Record<string, unknown> | null> {
|
||||
const resolved = await configPath();
|
||||
if (!resolved) return null;
|
||||
try {
|
||||
if (!resolved.modules.fs.existsSync(resolved.path)) return null;
|
||||
const raw = resolved.modules.fs.readFileSync(resolved.path, "utf8");
|
||||
const parsed = JSON.parse(raw);
|
||||
if (!node.fs.existsSync(node.configPath)) return null;
|
||||
const parsed = JSON.parse(node.fs.readFileSync(node.configPath, "utf8"));
|
||||
return parsed && typeof parsed === "object" ? parsed : null;
|
||||
} catch {
|
||||
return null;
|
||||
@@ -79,40 +53,41 @@ async function loadRawConfig(): Promise<Record<string, unknown> | null> {
|
||||
}
|
||||
|
||||
export async function readMem0AnonIds(): Promise<Mem0AnonIds | null> {
|
||||
const config = await loadRawConfig();
|
||||
const node = await getNodeFs();
|
||||
if (!node) return null;
|
||||
const config = loadConfig(node);
|
||||
if (!config) return null;
|
||||
const telemetry =
|
||||
config.telemetry && typeof config.telemetry === "object"
|
||||
? (config.telemetry as Record<string, unknown>)
|
||||
? config.telemetry
|
||||
: {};
|
||||
const oss = typeof config.user_id === "string" ? config.user_id : undefined;
|
||||
const cli =
|
||||
typeof telemetry.anonymous_id === "string"
|
||||
? telemetry.anonymous_id
|
||||
: undefined;
|
||||
const aliasedTo =
|
||||
typeof telemetry.aliased_to === "string" ? telemetry.aliased_to : undefined;
|
||||
return { oss, cli, aliasedTo };
|
||||
return {
|
||||
oss: typeof config.user_id === "string" ? config.user_id : undefined,
|
||||
cli:
|
||||
typeof telemetry.anonymous_id === "string"
|
||||
? telemetry.anonymous_id
|
||||
: undefined,
|
||||
aliasedTo:
|
||||
typeof telemetry.aliased_to === "string"
|
||||
? telemetry.aliased_to
|
||||
: undefined,
|
||||
};
|
||||
}
|
||||
|
||||
export async function markMem0Aliased(email: string): Promise<void> {
|
||||
const resolved = await configPath();
|
||||
if (!resolved) return;
|
||||
const node = await getNodeFs();
|
||||
if (!node) return;
|
||||
try {
|
||||
const dirname = resolved.modules.path.dirname(resolved.path);
|
||||
resolved.modules.fs.mkdirSync(dirname, { recursive: true });
|
||||
const config = (await loadRawConfig()) ?? {};
|
||||
node.fs.mkdirSync(node.path.dirname(node.configPath), { recursive: true });
|
||||
const config = loadConfig(node) ?? {};
|
||||
const telemetry =
|
||||
config.telemetry && typeof config.telemetry === "object"
|
||||
? (config.telemetry as Record<string, unknown>)
|
||||
? config.telemetry
|
||||
: {};
|
||||
telemetry.aliased_to = email;
|
||||
config.telemetry = telemetry;
|
||||
resolved.modules.fs.writeFileSync(
|
||||
resolved.path,
|
||||
JSON.stringify(config, null, 4),
|
||||
);
|
||||
node.fs.writeFileSync(node.configPath, JSON.stringify(config, null, 4));
|
||||
} catch {
|
||||
// Read-only filesystem (Lambda, container) — alias is best-effort.
|
||||
// Best-effort: read-only filesystems and unwritable paths just skip.
|
||||
}
|
||||
}
|
||||
|
||||
@@ -20,7 +20,12 @@ import {
|
||||
CreateMemoryExportPayload,
|
||||
GetMemoryExportPayload,
|
||||
} from "./mem0.types";
|
||||
import { captureClientEvent, generateHash, telemetry } from "./telemetry";
|
||||
import {
|
||||
captureClientEvent,
|
||||
generateHash,
|
||||
isTelemetryEnabled,
|
||||
telemetry,
|
||||
} from "./telemetry";
|
||||
import { readMem0AnonIds, markMem0Aliased } from "./config";
|
||||
import { camelToSnake, camelToSnakeKeys, snakeToCamelKeys } from "./utils";
|
||||
import { createExceptionFromResponse, MemoryError } from "../common/exceptions";
|
||||
@@ -136,6 +141,7 @@ export default class MemoryClient {
|
||||
}
|
||||
|
||||
private async _maybeAliasAnonToEmail(): Promise<void> {
|
||||
if (!isTelemetryEnabled()) return;
|
||||
try {
|
||||
const email = this.telemetryId;
|
||||
if (!email || !email.includes("@")) return;
|
||||
|
||||
@@ -81,6 +81,10 @@ class UnifiedTelemetry implements TelemetryClient {
|
||||
}
|
||||
}
|
||||
|
||||
function isTelemetryEnabled(): boolean {
|
||||
return MEM0_TELEMETRY;
|
||||
}
|
||||
|
||||
const telemetry = new UnifiedTelemetry(POSTHOG_API_KEY, POSTHOG_HOST);
|
||||
|
||||
async function captureClientEvent(
|
||||
@@ -110,4 +114,4 @@ async function captureClientEvent(
|
||||
);
|
||||
}
|
||||
|
||||
export { telemetry, captureClientEvent, generateHash };
|
||||
export { telemetry, captureClientEvent, generateHash, isTelemetryEnabled };
|
||||
|
||||
@@ -9,8 +9,8 @@ import * as os from "os";
|
||||
import * as path from "path";
|
||||
import { MemoryClient } from "../mem0";
|
||||
import { telemetry } from "../telemetry";
|
||||
import { readMem0AnonIds, markMem0Aliased, Mem0AnonIds } from "../config";
|
||||
import { createMockFetch, TEST_API_KEY } from "./helpers";
|
||||
import { readMem0AnonIds, markMem0Aliased } from "../config";
|
||||
import { TEST_API_KEY } from "./helpers";
|
||||
import { setupMockFetch, installConsoleSuppression } from "./setup";
|
||||
|
||||
installConsoleSuppression();
|
||||
@@ -285,16 +285,42 @@ describe("MemoryClient — _maybeAliasAnonToEmail", () => {
|
||||
(client as any)._maybeAliasAnonToEmail(),
|
||||
).resolves.toBeUndefined();
|
||||
});
|
||||
});
|
||||
|
||||
// ─── Type smoke test ──────────────────────────────────────────
|
||||
test("noop when telemetry disabled — no fs read, no fs write, no events", async () => {
|
||||
fs.writeFileSync(
|
||||
path.join(tmpHome, "config.json"),
|
||||
JSON.stringify({ user_id: "oss-uuid" }),
|
||||
);
|
||||
const fetchMock = setupMockFetch();
|
||||
|
||||
describe("Mem0AnonIds shape", () => {
|
||||
test("interface accepts partial fields", () => {
|
||||
const a: Mem0AnonIds = { oss: "x" };
|
||||
const b: Mem0AnonIds = { cli: "y", aliasedTo: "u@x.com" };
|
||||
expect(a.oss).toBe("x");
|
||||
expect(b.aliasedTo).toBe("u@x.com");
|
||||
jest.resetModules();
|
||||
const original = process.env.MEM0_TELEMETRY;
|
||||
process.env.MEM0_TELEMETRY = "false";
|
||||
try {
|
||||
const { MemoryClient: ColdClient } = await import("../mem0");
|
||||
const client = Object.create(ColdClient.prototype);
|
||||
client.apiKey = TEST_API_KEY;
|
||||
client.host = "https://api.mem0.ai";
|
||||
client.telemetryId = "test@example.com";
|
||||
await client._maybeAliasAnonToEmail();
|
||||
} finally {
|
||||
if (original === undefined) delete process.env.MEM0_TELEMETRY;
|
||||
else process.env.MEM0_TELEMETRY = original;
|
||||
jest.resetModules();
|
||||
}
|
||||
|
||||
const identifyCalls = (fetchMock.mock.calls as any[]).filter(
|
||||
([, init]: [string, RequestInit]) => {
|
||||
if (!init?.body) return false;
|
||||
return JSON.parse(init.body as string).event === "$identify";
|
||||
},
|
||||
);
|
||||
expect(identifyCalls.length).toBe(0);
|
||||
|
||||
const written = JSON.parse(
|
||||
fs.readFileSync(path.join(tmpHome, "config.json"), "utf8"),
|
||||
);
|
||||
expect(written.telemetry?.aliased_to).toBeUndefined();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -303,11 +329,9 @@ describe("Mem0AnonIds shape", () => {
|
||||
describe("config.ts in browser-like environment", () => {
|
||||
test("readMem0AnonIds returns null when not Node", async () => {
|
||||
const originalProcess = global.process;
|
||||
// Simulate a browser: no `process` global.
|
||||
// @ts-expect-error force-undefining global
|
||||
// @ts-expect-error force-undefining global to simulate a browser
|
||||
delete global.process;
|
||||
try {
|
||||
// Re-import to pick up the missing global.
|
||||
jest.resetModules();
|
||||
const { readMem0AnonIds: browserRead } = await import("../config");
|
||||
expect(await browserRead()).toBeNull();
|
||||
@@ -317,7 +341,3 @@ describe("config.ts in browser-like environment", () => {
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
// ─── Suppress unused-import warnings ──────────────────────────
|
||||
// Some helpers above are intentionally re-exported for clarity in tests.
|
||||
void createMockFetch;
|
||||
|
||||
+6
-5
@@ -34,13 +34,14 @@ ENTITY_PARAMS = frozenset({"user_id", "agent_id", "app_id", "run_id"})
|
||||
|
||||
|
||||
def _maybe_alias_anon_to_email(user_email):
|
||||
"""Stitch prior anonymous PostHog identities to the resolved email.
|
||||
"""Fire $identify per prior anon ID so PostHog merges them into email.
|
||||
|
||||
Reads ~/.mem0/config.json for both anon IDs (OSS user_id, CLI
|
||||
telemetry.anonymous_id), fires one $identify per anon that hasn't been
|
||||
aliased yet, then sets telemetry.aliased_to so this only happens once
|
||||
per (anon_id, email) pair. Best-effort — never raises.
|
||||
Idempotent via telemetry.aliased_to — only writes the flag when telemetry
|
||||
is actually enabled, so disabling/re-enabling MEM0_TELEMETRY still works.
|
||||
Best-effort — never raises.
|
||||
"""
|
||||
if client_telemetry.posthog is None:
|
||||
return
|
||||
if not user_email or "@" not in user_email:
|
||||
return
|
||||
try:
|
||||
|
||||
@@ -114,29 +114,17 @@ class AnonymousTelemetry:
|
||||
_logger.debug("Failed to capture telemetry event %r: %s", event_name, e)
|
||||
|
||||
def capture_identify(self, anon_id, email):
|
||||
"""Send a PostHog $identify event so the anon person merges into the
|
||||
identified person. Distinct_id is the new identity (email); the prior
|
||||
anonymous distinct_id is passed as $anon_distinct_id per PostHog spec.
|
||||
|
||||
Never raises.
|
||||
"""
|
||||
"""Fire $identify with $anon_distinct_id so PostHog merges anon_id into email."""
|
||||
if self.posthog is None:
|
||||
return
|
||||
if not anon_id or not email or anon_id == email:
|
||||
return
|
||||
properties = {
|
||||
"$anon_distinct_id": anon_id,
|
||||
"client_source": "python",
|
||||
"client_version": mem0.__version__,
|
||||
"python_version": sys.version,
|
||||
"os": sys.platform,
|
||||
"os_version": platform.version(),
|
||||
"os_release": platform.release(),
|
||||
"processor": platform.processor(),
|
||||
"machine": platform.machine(),
|
||||
}
|
||||
try:
|
||||
self.posthog.capture(distinct_id=email, event="$identify", properties=properties)
|
||||
self.posthog.capture(
|
||||
distinct_id=email,
|
||||
event="$identify",
|
||||
properties={"$anon_distinct_id": anon_id, "client_source": "python"},
|
||||
)
|
||||
except Exception as e:
|
||||
_logger.debug("Failed to capture $identify for %r: %s", email, e)
|
||||
|
||||
|
||||
@@ -297,6 +297,23 @@ class TestMaybeAliasAnonToEmail:
|
||||
client_main._maybe_alias_anon_to_email("not-an-email")
|
||||
telemetry.capture_identify.assert_not_called()
|
||||
|
||||
def test_skips_when_telemetry_disabled(self):
|
||||
"""When client_telemetry.posthog is None (MEM0_TELEMETRY=false), do nothing —
|
||||
no fs read, no fs write, no event. Re-enabling telemetry later must still alias."""
|
||||
from mem0.client import main as client_main
|
||||
|
||||
disabled = MagicMock()
|
||||
disabled.posthog = None
|
||||
with (
|
||||
patch.object(client_main, "client_telemetry", disabled),
|
||||
patch.object(client_main, "read_anon_ids") as read,
|
||||
patch.object(client_main, "mark_aliased") as mark,
|
||||
):
|
||||
client_main._maybe_alias_anon_to_email("user@example.com")
|
||||
read.assert_not_called()
|
||||
mark.assert_not_called()
|
||||
disabled.capture_identify.assert_not_called()
|
||||
|
||||
def test_does_not_raise_on_telemetry_failure(self):
|
||||
from mem0.client import main as client_main
|
||||
|
||||
|
||||
Reference in New Issue
Block a user