Merge branch 'pr3/spool-delivery' into pr4/install-marker-and-identity

This commit is contained in:
Saket Aryan
2026-09-15 00:33:09 +05:30
9 changed files with 982 additions and 142 deletions
@@ -49,6 +49,7 @@ except ImportError:
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
_salt_cache: str = ""
_harness: str = _DEFAULT_HARNESS
_source_tag: str = _DEFAULT_SOURCE_TAG
_PRIVATE_KEYS = {
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _salt_path() -> Path:
return memory_core.data_dir() / "telemetry-salt"
def _install_salt() -> str:
"""Random per-install salt, created on first use and kept in the identity file."""
identity = _read_identity()
salt = identity.get("salt")
if not salt:
salt = uuid.uuid4().hex
identity["salt"] = salt
_write_identity(identity)
return salt
"""Random per-install salt, created once and memoized for the process.
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
the identity file. Three reasons, all of which produced wrong data when this
lived in the identity dict:
- Hooks are short-lived separate processes firing on every tool call, and
people run more than one agent window. A read-modify-write would let each
process mint its own salt, so one repository would hash several ways in the
window before a writer won.
- resolve_distinct_id holds a copy of the identity dict across a network call
to /v1/ping/, so whichever write landed second erased the other's key —
losing either the salt (repo_hash changes mid-stream) or the email (a
second $identify, splitting the person).
- Touching the identity file from record() would create it, and is_first_run
keys off that file, so recording an event would silently suppress the
install event.
On a read-only or full data directory the fallback is derived from the data
directory path: stable for the machine rather than random per call, so the
failure mode is a weaker salt and not unbounded cardinality in PostHog.
"""
global _salt_cache
if _salt_cache:
return _salt_cache
path = _salt_path()
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
try:
with os.fdopen(handle, "w", encoding="utf-8") as stream:
stream.write(uuid.uuid4().hex)
except OSError:
pass
except FileExistsError:
pass
except OSError:
# Cannot persist. Stable-per-machine beats random-per-call.
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
except OSError:
_salt_cache = ""
if not _salt_cache:
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
def _scoped_digest(value: str, length: int = 16) -> str:
@@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
def _claim_attempt(claim: Path) -> int:
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
Anchored on field position, not on a leading "a": the legacy shape is
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
otherwise parse as attempt 1234567 and be discarded unsent on the first
flush after an upgrade.
"""
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
tail = stem.rsplit("-", 1)[-1]
parts = stem.split("-")
if len(parts) != 4:
return 0
tail = parts[3]
if tail.startswith("a") and tail[1:].isdigit():
return int(tail[1:])
return 0
@@ -481,6 +536,21 @@ def _claim_spool() -> Path | None:
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove temp files orphaned by a crash between write and rename.
Neither glob in this module matches *.partial, so nothing else would ever
clean them up.
"""
now = time.time()
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
@@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
# for files that never carried an attempt marker.
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
@@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
temporary.write_text(
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
encoding="utf-8",
)
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
# fsync before the rename: without it the rename can land while the
# bytes have not, and the claim comes back empty or truncated after a
# crash. _drain then reads zero events and unlinks it.
with open(temporary, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
temporary.replace(claim)
_touch(claim)
return True
@@ -664,6 +742,7 @@ def flush() -> int:
# Parked batches used to starve behind the live spool indefinitely. Bounded
# per run so a long backlog cannot turn one flush into an unbounded loop.
directory = memory_core.data_dir()
_sweep_debris(directory)
for _ in range(MAX_PARKED_PER_RUN):
parked = _claim_parked(directory)
if parked is None:
@@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except OSError:
except (OSError, ValueError):
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
# rather than retry: flush() runs from a bare `finally:` in
# flush_worker, so raising here also skips the handoff cleanup, and an
# undecodable file would otherwise be re-read on every flush forever.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
events = []
for line in lines:
@@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
# Only delete when the file really is empty. A non-empty file that
# parses to nothing is a torn write, and its contents are the unsent
# remainder — deleting it is the data loss this PR exists to prevent.
try:
claim.unlink()
empty = claim.stat().st_size == 0
except OSError:
empty = True
try:
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
except OSError:
pass
return 0, True
@@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return sent, False
sent += len(chunk)
# Record progress and refresh the lease after each successful batch, so
# a crash repeats at most one batch instead of the entire file.
_rewrite_claim(claim, events[start + len(chunk) :])
# a crash repeats at most one batch instead of the entire file. If the
# rewrite fails the claim still holds delivered events, so stop rather
# than carry on as though progress were recorded — continuing is how the
# duplicate delivery this PR fixes would come back.
if not _rewrite_claim(claim, events[start + len(chunk) :]):
_release_claim(claim, events[start + len(chunk) :])
return sent, False
return sent, True
@@ -112,16 +112,19 @@ def test_a_parked_batch_is_drained_behind_the_live_spool(telemetry):
assert names == {"code.parked", "code.fresh"}
def test_an_untried_batch_is_not_expired_by_age_alone(telemetry):
"""Expiry should discard what failed, not what never got a turn."""
def test_a_batch_is_retried_until_the_budget_is_spent_not_discarded(telemetry):
"""Expiry discards what failed repeatedly, not what merely sat for a while.
The budget is the attempt count, because age cannot be one: every re-claim
touches the mtime and every release backdates it, so age never accumulates.
"""
telemetry.record("parked")
telemetry._post = lambda payload, url: False
telemetry.flush()
parked = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending"))
assert len(parked) == 1
ancient = time.time() - (telemetry.CLAIM_EXPIRY_SECONDS + 60)
os.utime(parked[0], (ancient, ancient))
assert telemetry._claim_attempt(parked[0]) < telemetry.MAX_CLAIM_ATTEMPTS
sent: list[dict] = []
telemetry._post = lambda payload, url: sent.append(payload) or True
@@ -154,11 +157,88 @@ def test_progress_is_recorded_after_every_batch(telemetry):
assert json.loads(remaining[0])["properties"]["index"] == 200
def test_the_heartbeat_stays_well_inside_the_lease(telemetry):
def test_the_heartbeat_actually_refreshes_the_lease(telemetry):
"""The claim rewrite doubles as the lease heartbeat.
_post makes a single attempt with SEND_TIMEOUT and no retry, so a heartbeat
lands at least that often. If a retry loop is ever added to _post, this is
the assertion that catches a sender losing its claim mid-flight.
Previously asserted `SEND_TIMEOUT * 4 < CLAIM_STALE_SECONDS`, which compares
two constants and executes none of the code under test. Drive the real
rewrite and watch the mtime move instead.
"""
assert telemetry.SEND_TIMEOUT * 4 < telemetry.CLAIM_STALE_SECONDS
for index in range(150):
telemetry.record("search", index=index)
claim = telemetry._claim_spool()
assert claim is not None
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(claim, (stale, stale))
assert time.time() - claim.stat().st_mtime > telemetry.CLAIM_STALE_SECONDS
telemetry._rewrite_claim(claim, [{"event": "code.x", "properties": {}}])
assert time.time() - claim.stat().st_mtime < telemetry.CLAIM_STALE_SECONDS
def test_an_undeliverable_batch_is_eventually_given_up_on(telemetry):
"""Expiry has to be reachable from a state the state machine can produce.
It was not: every re-claim touched the mtime and every release backdated it
by a fixed amount, so age hovered near the stale threshold and the 7-day
expiry never fired. An undeliverable batch lived on disk forever, and
spawn_flush saw it and started a sender on every hook.
"""
telemetry.record("doomed")
telemetry._post = lambda payload, url: False
for _ in range(telemetry.MAX_CLAIM_ATTEMPTS + 3):
telemetry.flush()
leftover = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending"))
assert leftover == [], f"batch never given up on: {[p.name for p in leftover]}"
def test_a_legacy_claim_filename_is_not_mistaken_for_a_huge_attempt_count(telemetry):
"""The old shape is telemetry-<pid>-<hex>.sending, and hex can start with 'a'."""
assert telemetry._claim_attempt(Path("telemetry-999-deadbeef.sending")) == 0
assert telemetry._claim_attempt(Path("telemetry-999-a1234567.sending")) == 0
assert telemetry._claim_attempt(Path("telemetry-999-deadbeef-a2.sending")) == 2
def test_a_torn_claim_is_quarantined_not_deleted(telemetry):
"""A non-empty file that parses to nothing is the remainder, not garbage."""
telemetry.record("search")
claim = telemetry._claim_spool()
claim.write_bytes(b"\xff\xfe not utf-8 at all")
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(claim, (stale, stale))
sent = telemetry.flush()
assert sent == 0
assert not claim.exists()
quarantined = list(telemetry.memory_core.data_dir().glob("*.corrupt"))
assert len(quarantined) == 1, "torn claim was destroyed instead of kept"
def test_a_failed_rewrite_stops_instead_of_redelivering(telemetry):
"""Ignoring the rewrite result reintroduced the duplicates this PR fixes."""
for index in range(250):
telemetry.record("search", index=index)
telemetry._rewrite_claim = lambda claim, remaining: False
delivered = []
telemetry._post = lambda payload, url: delivered.extend(payload.get("batch", [])) or True
telemetry.flush()
assert len(delivered) == 100, f"kept going after a failed rewrite: {len(delivered)}"
def test_partial_files_are_swept(telemetry):
"""Nothing else globs *.partial, so a crash mid-rename orphans one forever."""
data_dir = telemetry.memory_core.data_dir()
data_dir.mkdir(parents=True, exist_ok=True)
debris = data_dir / "telemetry-1-abc-a0.1.partial"
debris.write_text("x", encoding="utf-8")
old = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(debris, (old, old))
telemetry.flush()
assert not debris.exists()
+121 -19
View File
@@ -49,6 +49,7 @@ except ImportError:
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
_salt_cache: str = ""
_harness: str = _DEFAULT_HARNESS
_source_tag: str = _DEFAULT_SOURCE_TAG
_PRIVATE_KEYS = {
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _salt_path() -> Path:
return memory_core.data_dir() / "telemetry-salt"
def _install_salt() -> str:
"""Random per-install salt, created on first use and kept in the identity file."""
identity = _read_identity()
salt = identity.get("salt")
if not salt:
salt = uuid.uuid4().hex
identity["salt"] = salt
_write_identity(identity)
return salt
"""Random per-install salt, created once and memoized for the process.
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
the identity file. Three reasons, all of which produced wrong data when this
lived in the identity dict:
- Hooks are short-lived separate processes firing on every tool call, and
people run more than one agent window. A read-modify-write would let each
process mint its own salt, so one repository would hash several ways in the
window before a writer won.
- resolve_distinct_id holds a copy of the identity dict across a network call
to /v1/ping/, so whichever write landed second erased the other's key —
losing either the salt (repo_hash changes mid-stream) or the email (a
second $identify, splitting the person).
- Touching the identity file from record() would create it, and is_first_run
keys off that file, so recording an event would silently suppress the
install event.
On a read-only or full data directory the fallback is derived from the data
directory path: stable for the machine rather than random per call, so the
failure mode is a weaker salt and not unbounded cardinality in PostHog.
"""
global _salt_cache
if _salt_cache:
return _salt_cache
path = _salt_path()
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
try:
with os.fdopen(handle, "w", encoding="utf-8") as stream:
stream.write(uuid.uuid4().hex)
except OSError:
pass
except FileExistsError:
pass
except OSError:
# Cannot persist. Stable-per-machine beats random-per-call.
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
except OSError:
_salt_cache = ""
if not _salt_cache:
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
def _scoped_digest(value: str, length: int = 16) -> str:
@@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
def _claim_attempt(claim: Path) -> int:
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
Anchored on field position, not on a leading "a": the legacy shape is
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
otherwise parse as attempt 1234567 and be discarded unsent on the first
flush after an upgrade.
"""
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
tail = stem.rsplit("-", 1)[-1]
parts = stem.split("-")
if len(parts) != 4:
return 0
tail = parts[3]
if tail.startswith("a") and tail[1:].isdigit():
return int(tail[1:])
return 0
@@ -481,6 +536,21 @@ def _claim_spool() -> Path | None:
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove temp files orphaned by a crash between write and rename.
Neither glob in this module matches *.partial, so nothing else would ever
clean them up.
"""
now = time.time()
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
@@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
# for files that never carried an attempt marker.
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
@@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
temporary.write_text(
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
encoding="utf-8",
)
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
# fsync before the rename: without it the rename can land while the
# bytes have not, and the claim comes back empty or truncated after a
# crash. _drain then reads zero events and unlinks it.
with open(temporary, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
temporary.replace(claim)
_touch(claim)
return True
@@ -664,6 +742,7 @@ def flush() -> int:
# Parked batches used to starve behind the live spool indefinitely. Bounded
# per run so a long backlog cannot turn one flush into an unbounded loop.
directory = memory_core.data_dir()
_sweep_debris(directory)
for _ in range(MAX_PARKED_PER_RUN):
parked = _claim_parked(directory)
if parked is None:
@@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except OSError:
except (OSError, ValueError):
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
# rather than retry: flush() runs from a bare `finally:` in
# flush_worker, so raising here also skips the handoff cleanup, and an
# undecodable file would otherwise be re-read on every flush forever.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
events = []
for line in lines:
@@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
# Only delete when the file really is empty. A non-empty file that
# parses to nothing is a torn write, and its contents are the unsent
# remainder — deleting it is the data loss this PR exists to prevent.
try:
claim.unlink()
empty = claim.stat().st_size == 0
except OSError:
empty = True
try:
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
except OSError:
pass
return 0, True
@@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return sent, False
sent += len(chunk)
# Record progress and refresh the lease after each successful batch, so
# a crash repeats at most one batch instead of the entire file.
_rewrite_claim(claim, events[start + len(chunk) :])
# a crash repeats at most one batch instead of the entire file. If the
# rewrite fails the claim still holds delivered events, so stop rather
# than carry on as though progress were recorded — continuing is how the
# duplicate delivery this PR fixes would come back.
if not _rewrite_claim(claim, events[start + len(chunk) :]):
_release_claim(claim, events[start + len(chunk) :])
return sent, False
return sent, True
+121 -19
View File
@@ -49,6 +49,7 @@ except ImportError:
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
_salt_cache: str = ""
_harness: str = _DEFAULT_HARNESS
_source_tag: str = _DEFAULT_SOURCE_TAG
_PRIVATE_KEYS = {
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _salt_path() -> Path:
return memory_core.data_dir() / "telemetry-salt"
def _install_salt() -> str:
"""Random per-install salt, created on first use and kept in the identity file."""
identity = _read_identity()
salt = identity.get("salt")
if not salt:
salt = uuid.uuid4().hex
identity["salt"] = salt
_write_identity(identity)
return salt
"""Random per-install salt, created once and memoized for the process.
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
the identity file. Three reasons, all of which produced wrong data when this
lived in the identity dict:
- Hooks are short-lived separate processes firing on every tool call, and
people run more than one agent window. A read-modify-write would let each
process mint its own salt, so one repository would hash several ways in the
window before a writer won.
- resolve_distinct_id holds a copy of the identity dict across a network call
to /v1/ping/, so whichever write landed second erased the other's key —
losing either the salt (repo_hash changes mid-stream) or the email (a
second $identify, splitting the person).
- Touching the identity file from record() would create it, and is_first_run
keys off that file, so recording an event would silently suppress the
install event.
On a read-only or full data directory the fallback is derived from the data
directory path: stable for the machine rather than random per call, so the
failure mode is a weaker salt and not unbounded cardinality in PostHog.
"""
global _salt_cache
if _salt_cache:
return _salt_cache
path = _salt_path()
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
try:
with os.fdopen(handle, "w", encoding="utf-8") as stream:
stream.write(uuid.uuid4().hex)
except OSError:
pass
except FileExistsError:
pass
except OSError:
# Cannot persist. Stable-per-machine beats random-per-call.
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
except OSError:
_salt_cache = ""
if not _salt_cache:
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
def _scoped_digest(value: str, length: int = 16) -> str:
@@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
def _claim_attempt(claim: Path) -> int:
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
Anchored on field position, not on a leading "a": the legacy shape is
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
otherwise parse as attempt 1234567 and be discarded unsent on the first
flush after an upgrade.
"""
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
tail = stem.rsplit("-", 1)[-1]
parts = stem.split("-")
if len(parts) != 4:
return 0
tail = parts[3]
if tail.startswith("a") and tail[1:].isdigit():
return int(tail[1:])
return 0
@@ -481,6 +536,21 @@ def _claim_spool() -> Path | None:
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove temp files orphaned by a crash between write and rename.
Neither glob in this module matches *.partial, so nothing else would ever
clean them up.
"""
now = time.time()
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
@@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
# for files that never carried an attempt marker.
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
@@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
temporary.write_text(
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
encoding="utf-8",
)
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
# fsync before the rename: without it the rename can land while the
# bytes have not, and the claim comes back empty or truncated after a
# crash. _drain then reads zero events and unlinks it.
with open(temporary, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
temporary.replace(claim)
_touch(claim)
return True
@@ -664,6 +742,7 @@ def flush() -> int:
# Parked batches used to starve behind the live spool indefinitely. Bounded
# per run so a long backlog cannot turn one flush into an unbounded loop.
directory = memory_core.data_dir()
_sweep_debris(directory)
for _ in range(MAX_PARKED_PER_RUN):
parked = _claim_parked(directory)
if parked is None:
@@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except OSError:
except (OSError, ValueError):
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
# rather than retry: flush() runs from a bare `finally:` in
# flush_worker, so raising here also skips the handoff cleanup, and an
# undecodable file would otherwise be re-read on every flush forever.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
events = []
for line in lines:
@@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
# Only delete when the file really is empty. A non-empty file that
# parses to nothing is a torn write, and its contents are the unsent
# remainder — deleting it is the data loss this PR exists to prevent.
try:
claim.unlink()
empty = claim.stat().st_size == 0
except OSError:
empty = True
try:
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
except OSError:
pass
return 0, True
@@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return sent, False
sent += len(chunk)
# Record progress and refresh the lease after each successful batch, so
# a crash repeats at most one batch instead of the entire file.
_rewrite_claim(claim, events[start + len(chunk) :])
# a crash repeats at most one batch instead of the entire file. If the
# rewrite fails the claim still holds delivered events, so stop rather
# than carry on as though progress were recorded — continuing is how the
# duplicate delivery this PR fixes would come back.
if not _rewrite_claim(claim, events[start + len(chunk) :]):
_release_claim(claim, events[start + len(chunk) :])
return sent, False
return sent, True
@@ -315,3 +315,49 @@ def test_spawn_flush_does_nothing_without_a_spool(isolated_env):
with patch.object(telemetry.subprocess, "Popen") as popen:
assert telemetry.spawn_flush() is True
popen.assert_called_once()
def test_salt_is_stable_across_processes(isolated_env):
"""Hooks are separate short-lived processes; one repo must hash one way.
An unlocked read-modify-write let each process mint its own salt, so a
repository hashed several ways in the window before one writer won.
"""
import subprocess as sp
core = str(Path(__file__).resolve().parents[1] / "core")
script = (
f"import sys; sys.path.insert(0, {core!r})\n"
"import telemetry\n"
"print(telemetry._install_salt())"
)
env = {**os.environ, "MEM0_CODE_DATA_DIR": str(memory_core.data_dir())}
salts = {
sp.run([sys.executable, "-c", script], capture_output=True, text=True, env=env).stdout.strip()
for _ in range(4)
}
assert len(salts) == 1, f"one repo hashed {len(salts)} ways: {salts}"
def test_salt_does_not_touch_the_identity_file(isolated_env):
"""The identity file is is_first_run's marker and the sender's email store.
Writing the salt into it would create it from record(), suppressing the
install event, and would race resolve_distinct_id, which holds a stale copy
of that dict across a network call.
"""
telemetry._install_salt()
assert not telemetry._identity_path().exists()
def test_salt_is_stable_when_it_cannot_be_persisted(isolated_env, monkeypatch):
"""A read-only data dir must degrade to a weaker salt, not to random-per-call.
Random per call is unbounded cardinality in PostHog, which is worse than no
salt at all.
"""
telemetry._salt_cache = ""
monkeypatch.setattr(telemetry.os, "open", lambda *a, **k: (_ for _ in ()).throw(OSError("read-only")))
first = telemetry._install_salt()
telemetry._salt_cache = ""
assert telemetry._install_salt() == first
+121 -19
View File
@@ -49,6 +49,7 @@ except ImportError:
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
_salt_cache: str = ""
_harness: str = _DEFAULT_HARNESS
_source_tag: str = _DEFAULT_SOURCE_TAG
_PRIVATE_KEYS = {
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _salt_path() -> Path:
return memory_core.data_dir() / "telemetry-salt"
def _install_salt() -> str:
"""Random per-install salt, created on first use and kept in the identity file."""
identity = _read_identity()
salt = identity.get("salt")
if not salt:
salt = uuid.uuid4().hex
identity["salt"] = salt
_write_identity(identity)
return salt
"""Random per-install salt, created once and memoized for the process.
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
the identity file. Three reasons, all of which produced wrong data when this
lived in the identity dict:
- Hooks are short-lived separate processes firing on every tool call, and
people run more than one agent window. A read-modify-write would let each
process mint its own salt, so one repository would hash several ways in the
window before a writer won.
- resolve_distinct_id holds a copy of the identity dict across a network call
to /v1/ping/, so whichever write landed second erased the other's key —
losing either the salt (repo_hash changes mid-stream) or the email (a
second $identify, splitting the person).
- Touching the identity file from record() would create it, and is_first_run
keys off that file, so recording an event would silently suppress the
install event.
On a read-only or full data directory the fallback is derived from the data
directory path: stable for the machine rather than random per call, so the
failure mode is a weaker salt and not unbounded cardinality in PostHog.
"""
global _salt_cache
if _salt_cache:
return _salt_cache
path = _salt_path()
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
try:
with os.fdopen(handle, "w", encoding="utf-8") as stream:
stream.write(uuid.uuid4().hex)
except OSError:
pass
except FileExistsError:
pass
except OSError:
# Cannot persist. Stable-per-machine beats random-per-call.
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
except OSError:
_salt_cache = ""
if not _salt_cache:
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
def _scoped_digest(value: str, length: int = 16) -> str:
@@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
def _claim_attempt(claim: Path) -> int:
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
Anchored on field position, not on a leading "a": the legacy shape is
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
otherwise parse as attempt 1234567 and be discarded unsent on the first
flush after an upgrade.
"""
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
tail = stem.rsplit("-", 1)[-1]
parts = stem.split("-")
if len(parts) != 4:
return 0
tail = parts[3]
if tail.startswith("a") and tail[1:].isdigit():
return int(tail[1:])
return 0
@@ -481,6 +536,21 @@ def _claim_spool() -> Path | None:
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove temp files orphaned by a crash between write and rename.
Neither glob in this module matches *.partial, so nothing else would ever
clean them up.
"""
now = time.time()
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
@@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
# for files that never carried an attempt marker.
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
@@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
temporary.write_text(
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
encoding="utf-8",
)
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
# fsync before the rename: without it the rename can land while the
# bytes have not, and the claim comes back empty or truncated after a
# crash. _drain then reads zero events and unlinks it.
with open(temporary, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
temporary.replace(claim)
_touch(claim)
return True
@@ -664,6 +742,7 @@ def flush() -> int:
# Parked batches used to starve behind the live spool indefinitely. Bounded
# per run so a long backlog cannot turn one flush into an unbounded loop.
directory = memory_core.data_dir()
_sweep_debris(directory)
for _ in range(MAX_PARKED_PER_RUN):
parked = _claim_parked(directory)
if parked is None:
@@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except OSError:
except (OSError, ValueError):
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
# rather than retry: flush() runs from a bare `finally:` in
# flush_worker, so raising here also skips the handoff cleanup, and an
# undecodable file would otherwise be re-read on every flush forever.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
events = []
for line in lines:
@@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
# Only delete when the file really is empty. A non-empty file that
# parses to nothing is a torn write, and its contents are the unsent
# remainder — deleting it is the data loss this PR exists to prevent.
try:
claim.unlink()
empty = claim.stat().st_size == 0
except OSError:
empty = True
try:
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
except OSError:
pass
return 0, True
@@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return sent, False
sent += len(chunk)
# Record progress and refresh the lease after each successful batch, so
# a crash repeats at most one batch instead of the entire file.
_rewrite_claim(claim, events[start + len(chunk) :])
# a crash repeats at most one batch instead of the entire file. If the
# rewrite fails the claim still holds delivered events, so stop rather
# than carry on as though progress were recorded — continuing is how the
# duplicate delivery this PR fixes would come back.
if not _rewrite_claim(claim, events[start + len(chunk) :]):
_release_claim(claim, events[start + len(chunk) :])
return sent, False
return sent, True
+121 -19
View File
@@ -49,6 +49,7 @@ except ImportError:
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
_salt_cache: str = ""
_harness: str = _DEFAULT_HARNESS
_source_tag: str = _DEFAULT_SOURCE_TAG
_PRIVATE_KEYS = {
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _salt_path() -> Path:
return memory_core.data_dir() / "telemetry-salt"
def _install_salt() -> str:
"""Random per-install salt, created on first use and kept in the identity file."""
identity = _read_identity()
salt = identity.get("salt")
if not salt:
salt = uuid.uuid4().hex
identity["salt"] = salt
_write_identity(identity)
return salt
"""Random per-install salt, created once and memoized for the process.
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
the identity file. Three reasons, all of which produced wrong data when this
lived in the identity dict:
- Hooks are short-lived separate processes firing on every tool call, and
people run more than one agent window. A read-modify-write would let each
process mint its own salt, so one repository would hash several ways in the
window before a writer won.
- resolve_distinct_id holds a copy of the identity dict across a network call
to /v1/ping/, so whichever write landed second erased the other's key —
losing either the salt (repo_hash changes mid-stream) or the email (a
second $identify, splitting the person).
- Touching the identity file from record() would create it, and is_first_run
keys off that file, so recording an event would silently suppress the
install event.
On a read-only or full data directory the fallback is derived from the data
directory path: stable for the machine rather than random per call, so the
failure mode is a weaker salt and not unbounded cardinality in PostHog.
"""
global _salt_cache
if _salt_cache:
return _salt_cache
path = _salt_path()
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
try:
with os.fdopen(handle, "w", encoding="utf-8") as stream:
stream.write(uuid.uuid4().hex)
except OSError:
pass
except FileExistsError:
pass
except OSError:
# Cannot persist. Stable-per-machine beats random-per-call.
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
except OSError:
_salt_cache = ""
if not _salt_cache:
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
def _scoped_digest(value: str, length: int = 16) -> str:
@@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
def _claim_attempt(claim: Path) -> int:
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
Anchored on field position, not on a leading "a": the legacy shape is
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
otherwise parse as attempt 1234567 and be discarded unsent on the first
flush after an upgrade.
"""
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
tail = stem.rsplit("-", 1)[-1]
parts = stem.split("-")
if len(parts) != 4:
return 0
tail = parts[3]
if tail.startswith("a") and tail[1:].isdigit():
return int(tail[1:])
return 0
@@ -481,6 +536,21 @@ def _claim_spool() -> Path | None:
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove temp files orphaned by a crash between write and rename.
Neither glob in this module matches *.partial, so nothing else would ever
clean them up.
"""
now = time.time()
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
@@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
# for files that never carried an attempt marker.
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
@@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
temporary.write_text(
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
encoding="utf-8",
)
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
# fsync before the rename: without it the rename can land while the
# bytes have not, and the claim comes back empty or truncated after a
# crash. _drain then reads zero events and unlinks it.
with open(temporary, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
temporary.replace(claim)
_touch(claim)
return True
@@ -664,6 +742,7 @@ def flush() -> int:
# Parked batches used to starve behind the live spool indefinitely. Bounded
# per run so a long backlog cannot turn one flush into an unbounded loop.
directory = memory_core.data_dir()
_sweep_debris(directory)
for _ in range(MAX_PARKED_PER_RUN):
parked = _claim_parked(directory)
if parked is None:
@@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except OSError:
except (OSError, ValueError):
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
# rather than retry: flush() runs from a bare `finally:` in
# flush_worker, so raising here also skips the handoff cleanup, and an
# undecodable file would otherwise be re-read on every flush forever.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
events = []
for line in lines:
@@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
# Only delete when the file really is empty. A non-empty file that
# parses to nothing is a torn write, and its contents are the unsent
# remainder — deleting it is the data loss this PR exists to prevent.
try:
claim.unlink()
empty = claim.stat().st_size == 0
except OSError:
empty = True
try:
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
except OSError:
pass
return 0, True
@@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return sent, False
sent += len(chunk)
# Record progress and refresh the lease after each successful batch, so
# a crash repeats at most one batch instead of the entire file.
_rewrite_claim(claim, events[start + len(chunk) :])
# a crash repeats at most one batch instead of the entire file. If the
# rewrite fails the claim still holds delivered events, so stop rather
# than carry on as though progress were recorded — continuing is how the
# duplicate delivery this PR fixes would come back.
if not _rewrite_claim(claim, events[start + len(chunk) :]):
_release_claim(claim, events[start + len(chunk) :])
return sent, False
return sent, True
+121 -19
View File
@@ -49,6 +49,7 @@ except ImportError:
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
_salt_cache: str = ""
_harness: str = _DEFAULT_HARNESS
_source_tag: str = _DEFAULT_SOURCE_TAG
_PRIVATE_KEYS = {
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _salt_path() -> Path:
return memory_core.data_dir() / "telemetry-salt"
def _install_salt() -> str:
"""Random per-install salt, created on first use and kept in the identity file."""
identity = _read_identity()
salt = identity.get("salt")
if not salt:
salt = uuid.uuid4().hex
identity["salt"] = salt
_write_identity(identity)
return salt
"""Random per-install salt, created once and memoized for the process.
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
the identity file. Three reasons, all of which produced wrong data when this
lived in the identity dict:
- Hooks are short-lived separate processes firing on every tool call, and
people run more than one agent window. A read-modify-write would let each
process mint its own salt, so one repository would hash several ways in the
window before a writer won.
- resolve_distinct_id holds a copy of the identity dict across a network call
to /v1/ping/, so whichever write landed second erased the other's key —
losing either the salt (repo_hash changes mid-stream) or the email (a
second $identify, splitting the person).
- Touching the identity file from record() would create it, and is_first_run
keys off that file, so recording an event would silently suppress the
install event.
On a read-only or full data directory the fallback is derived from the data
directory path: stable for the machine rather than random per call, so the
failure mode is a weaker salt and not unbounded cardinality in PostHog.
"""
global _salt_cache
if _salt_cache:
return _salt_cache
path = _salt_path()
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
try:
with os.fdopen(handle, "w", encoding="utf-8") as stream:
stream.write(uuid.uuid4().hex)
except OSError:
pass
except FileExistsError:
pass
except OSError:
# Cannot persist. Stable-per-machine beats random-per-call.
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
except OSError:
_salt_cache = ""
if not _salt_cache:
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
def _scoped_digest(value: str, length: int = 16) -> str:
@@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
def _claim_attempt(claim: Path) -> int:
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
Anchored on field position, not on a leading "a": the legacy shape is
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
otherwise parse as attempt 1234567 and be discarded unsent on the first
flush after an upgrade.
"""
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
tail = stem.rsplit("-", 1)[-1]
parts = stem.split("-")
if len(parts) != 4:
return 0
tail = parts[3]
if tail.startswith("a") and tail[1:].isdigit():
return int(tail[1:])
return 0
@@ -481,6 +536,21 @@ def _claim_spool() -> Path | None:
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove temp files orphaned by a crash between write and rename.
Neither glob in this module matches *.partial, so nothing else would ever
clean them up.
"""
now = time.time()
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
@@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
# for files that never carried an attempt marker.
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
@@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
temporary.write_text(
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
encoding="utf-8",
)
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
# fsync before the rename: without it the rename can land while the
# bytes have not, and the claim comes back empty or truncated after a
# crash. _drain then reads zero events and unlinks it.
with open(temporary, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
temporary.replace(claim)
_touch(claim)
return True
@@ -664,6 +742,7 @@ def flush() -> int:
# Parked batches used to starve behind the live spool indefinitely. Bounded
# per run so a long backlog cannot turn one flush into an unbounded loop.
directory = memory_core.data_dir()
_sweep_debris(directory)
for _ in range(MAX_PARKED_PER_RUN):
parked = _claim_parked(directory)
if parked is None:
@@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except OSError:
except (OSError, ValueError):
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
# rather than retry: flush() runs from a bare `finally:` in
# flush_worker, so raising here also skips the handoff cleanup, and an
# undecodable file would otherwise be re-read on every flush forever.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
events = []
for line in lines:
@@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
# Only delete when the file really is empty. A non-empty file that
# parses to nothing is a torn write, and its contents are the unsent
# remainder — deleting it is the data loss this PR exists to prevent.
try:
claim.unlink()
empty = claim.stat().st_size == 0
except OSError:
empty = True
try:
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
except OSError:
pass
return 0, True
@@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return sent, False
sent += len(chunk)
# Record progress and refresh the lease after each successful batch, so
# a crash repeats at most one batch instead of the entire file.
_rewrite_claim(claim, events[start + len(chunk) :])
# a crash repeats at most one batch instead of the entire file. If the
# rewrite fails the claim still holds delivered events, so stop rather
# than carry on as though progress were recorded — continuing is how the
# duplicate delivery this PR fixes would come back.
if not _rewrite_claim(claim, events[start + len(chunk) :]):
_release_claim(claim, events[start + len(chunk) :])
return sent, False
return sent, True
+121 -19
View File
@@ -49,6 +49,7 @@ except ImportError:
_PLATFORM_SOURCE = "MEM0_PLUGIN"
_PLATFORM_APPLICATION = ""
_salt_cache: str = ""
_harness: str = _DEFAULT_HARNESS
_source_tag: str = _DEFAULT_SOURCE_TAG
_PRIVATE_KEYS = {
@@ -121,15 +122,60 @@ def _digest(value: str, length: int = 16) -> str:
return hashlib.sha256(value.encode("utf-8")).hexdigest()[:length]
def _salt_path() -> Path:
return memory_core.data_dir() / "telemetry-salt"
def _install_salt() -> str:
"""Random per-install salt, created on first use and kept in the identity file."""
identity = _read_identity()
salt = identity.get("salt")
if not salt:
salt = uuid.uuid4().hex
identity["salt"] = salt
_write_identity(identity)
return salt
"""Random per-install salt, created once and memoized for the process.
Deliberately its own file, claimed with O_CREAT|O_EXCL, rather than a key in
the identity file. Three reasons, all of which produced wrong data when this
lived in the identity dict:
- Hooks are short-lived separate processes firing on every tool call, and
people run more than one agent window. A read-modify-write would let each
process mint its own salt, so one repository would hash several ways in the
window before a writer won.
- resolve_distinct_id holds a copy of the identity dict across a network call
to /v1/ping/, so whichever write landed second erased the other's key —
losing either the salt (repo_hash changes mid-stream) or the email (a
second $identify, splitting the person).
- Touching the identity file from record() would create it, and is_first_run
keys off that file, so recording an event would silently suppress the
install event.
On a read-only or full data directory the fallback is derived from the data
directory path: stable for the machine rather than random per call, so the
failure mode is a weaker salt and not unbounded cardinality in PostHog.
"""
global _salt_cache
if _salt_cache:
return _salt_cache
path = _salt_path()
try:
path.parent.mkdir(parents=True, exist_ok=True)
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
try:
with os.fdopen(handle, "w", encoding="utf-8") as stream:
stream.write(uuid.uuid4().hex)
except OSError:
pass
except FileExistsError:
pass
except OSError:
# Cannot persist. Stable-per-machine beats random-per-call.
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
try:
_salt_cache = path.read_text(encoding="utf-8").strip()
except OSError:
_salt_cache = ""
if not _salt_cache:
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
return _salt_cache
def _scoped_digest(value: str, length: int = 16) -> str:
@@ -445,9 +491,18 @@ def _claim_name(attempt: int = 0) -> str:
def _claim_attempt(claim: Path) -> int:
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape."""
"""Attempts recorded in a claim filename; 0 for the pre-attempt-count shape.
Anchored on field position, not on a leading "a": the legacy shape is
``telemetry-<pid>-<hex>.sending`` and a hex id such as ``a1234567`` would
otherwise parse as attempt 1234567 and be discarded unsent on the first
flush after an upgrade.
"""
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
tail = stem.rsplit("-", 1)[-1]
parts = stem.split("-")
if len(parts) != 4:
return 0
tail = parts[3]
if tail.startswith("a") and tail[1:].isdigit():
return int(tail[1:])
return 0
@@ -481,6 +536,21 @@ def _claim_spool() -> Path | None:
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove temp files orphaned by a crash between write and rename.
Neither glob in this module matches *.partial, so nothing else would ever
clean them up.
"""
now = time.time()
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
@@ -496,7 +566,11 @@ def _claim_parked(directory: Path) -> Path | None:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
# Attempts, not age. Every re-claim touches the mtime and every release
# backdates it by a fixed amount, so age is pinned near the stale
# threshold and never reaches the expiry. Age stays only as a backstop
# for files that never carried an attempt marker.
if _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS or age > CLAIM_EXPIRY_SECONDS:
try:
orphan.unlink()
except OSError:
@@ -538,10 +612,14 @@ def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
temporary.write_text(
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
encoding="utf-8",
)
payload = "".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining)
# fsync before the rename: without it the rename can land while the
# bytes have not, and the claim comes back empty or truncated after a
# crash. _drain then reads zero events and unlinks it.
with open(temporary, "w", encoding="utf-8") as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
temporary.replace(claim)
_touch(claim)
return True
@@ -664,6 +742,7 @@ def flush() -> int:
# Parked batches used to starve behind the live spool indefinitely. Bounded
# per run so a long backlog cannot turn one flush into an unbounded loop.
directory = memory_core.data_dir()
_sweep_debris(directory)
for _ in range(MAX_PARKED_PER_RUN):
parked = _claim_parked(directory)
if parked is None:
@@ -684,7 +763,18 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except OSError:
except (OSError, ValueError):
# ValueError covers UnicodeDecodeError from a torn write. Quarantine
# rather than retry: flush() runs from a bare `finally:` in
# flush_worker, so raising here also skips the handoff cleanup, and an
# undecodable file would otherwise be re-read on every flush forever.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
events = []
for line in lines:
@@ -695,8 +785,15 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
if isinstance(value, dict) and value.get("event"):
events.append(value)
if not events:
# Only delete when the file really is empty. A non-empty file that
# parses to nothing is a torn write, and its contents are the unsent
# remainder — deleting it is the data loss this PR exists to prevent.
try:
claim.unlink()
empty = claim.stat().st_size == 0
except OSError:
empty = True
try:
claim.replace(claim.with_suffix(".corrupt")) if not empty else claim.unlink()
except OSError:
pass
return 0, True
@@ -746,8 +843,13 @@ def _drain(claim: Path | None) -> tuple[int, bool]:
return sent, False
sent += len(chunk)
# Record progress and refresh the lease after each successful batch, so
# a crash repeats at most one batch instead of the entire file.
_rewrite_claim(claim, events[start + len(chunk) :])
# a crash repeats at most one batch instead of the entire file. If the
# rewrite fails the claim still holds delivered events, so stop rather
# than carry on as though progress were recorded — continuing is how the
# duplicate delivery this PR fixes would come back.
if not _rewrite_claim(claim, events[start + len(chunk) :]):
_release_claim(claim, events[start + len(chunk) :])
return sent, False
return sent, True