fix(plugins): stop delivering telemetry events twice, and stop losing parked ones (#7324)

This commit is contained in:
Saket Aryan
2026-09-18 13:43:33 +05:30
committed by GitHub
parent 012cd32c3a
commit 3362999095
9 changed files with 2243 additions and 139 deletions
@@ -100,6 +100,16 @@ BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
# A batch is only discarded once it has genuinely been retried this many times.
MAX_CLAIM_ATTEMPTS = 3
# Parked claims drained per run, after the live spool. Bounded so a long backlog
# cannot turn one flush into an unbounded send loop.
MAX_PARKED_PER_RUN = 3
# Added to the wait before a released claim becomes reclaimable, per attempt
# already spent. Releasing straight to "reclaimable now" let two senders burn the
# whole budget within seconds of one another on a single momentary failure, and
# discard a batch a retry a minute later would have delivered.
RETRY_COOLDOWN_SECONDS = 60
def is_enabled() -> bool:
@@ -399,38 +409,201 @@ def spawn_flush() -> bool:
return False
def _claim_name(attempt: int = 0) -> str:
"""Claim filename. The attempt count rides in the name so the 7-day expiry
only ever discards a batch that was actually retried and failed."""
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
def _claim_attempt(claim: Path) -> int:
"""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
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
def _touch(path: Path) -> None:
"""Refresh mtime so a claim's age measures time since it was claimed.
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
after a quiet minute inherited the spool's last-write time and looked
abandoned the instant it was made. A second sender would then take it over
while the first was still posting, and both would deliver the batch.
"""
try:
os.utime(path, None)
except OSError:
pass
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
claim = directory / _claim_name()
spool = _spool_path()
try:
spool.replace(claim)
_touch(claim)
return claim
except OSError:
pass
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove files nothing else will ever pick up again.
*.partial is a temp file orphaned by a crash between write and rename.
*.corrupt is a batch quarantined for undecodable content. No glob in this
module matches either, so without this they accumulate on disk for the life
of the install.
Quarantined batches are kept far longer than debris: they are the only
evidence left of events that could not be delivered, and someone diagnosing
a report of missing telemetry has to be able to find one.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending")):
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
for quarantined in directory.glob("telemetry-*.corrupt"):
try:
if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS:
quarantined.unlink()
except OSError:
continue
# The same reasoning covers *.tmp. _write_identity and _install_salt both
# create one and unlink it in a finally, which a SIGKILL skips, and no glob
# in this module matches the leftovers either.
for temporary in directory.glob("telemetry-*.tmp"):
try:
if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS:
temporary.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
Kept separate from the live spool so flush() can drain both in one run.
Previously parked batches were only reachable when no spool existed at all,
and because sessions keep recording there usually was one — so a batch
parked by a failed send waited until the 7-day expiry deleted it unsent,
even though its own presence is what started the sender.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS:
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# 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:
pass
continue
if age < CLAIM_STALE_SECONDS:
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
try:
orphan.replace(claim)
_touch(claim)
return claim
except OSError:
continue
return None
def _safe_mtime(path: Path) -> float:
try:
return path.stat().st_mtime
except OSError:
return 0.0
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
"""Persist the unsent remainder, atomically, and refresh the lease.
Called after every successful batch. Two jobs: a retry resumes where the
send stopped instead of re-posting from the top, and the rewrite doubles as
the lease heartbeat, so a slow sender does not have its claim stolen
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
"""
if not remaining:
try:
claim.unlink()
except OSError:
pass
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
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
except OSError:
try:
temporary.unlink()
except OSError:
pass
return False
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
"""Persist the remainder and drop the lease, because this sender has given up.
Distinct from the per-batch heartbeat: heartbeating on the way out would
make an abandoned batch look actively owned for a further
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
threshold lets the next flush pick it up immediately, while the attempt
count in the filename still bounds how many times that can happen.
"""
if not _rewrite_claim(claim, remaining):
return
try:
# Backdate past the stale threshold so the next flush can pick it up,
# minus a cooldown that grows with the attempts already spent. Clamped so
# the mtime never lands in the future, which would read as a live lease.
cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS)
released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown
os.utime(claim, (released, released))
except OSError:
pass
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
@@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]:
def flush() -> int:
"""Drain claimed spools to PostHog and return the number of events sent."""
"""Drain the live spool, then any parked claims, and return events sent."""
if not is_enabled():
return 0
claim = _claim_spool()
sent, delivered = _drain(_claim_spool())
if not delivered:
# The network is failing. Retrying other batches now would only burn
# their attempt budget against the same broken connection.
return sent
# 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:
break
count, delivered = _drain(parked)
sent += count
if not delivered:
break
return sent
def _drain(claim: Path | None) -> tuple[int, bool]:
"""Post one claimed batch file, recording progress after every batch.
Returns (events sent, whether everything was delivered).
"""
if claim is None:
return 0
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except ValueError:
# UnicodeDecodeError from a torn write: the content is unrecoverable, so
# 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.
# Reported as delivered because there is nothing left to deliver and the
# rest of the run should continue.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
except OSError:
return 0
# Could not read it, which is not the same as having nothing to send.
# The file is left exactly where it is: a vanished or briefly unreadable
# claim is retryable, and quarantining it here would discard events over
# a transient filesystem error. Reported as undelivered so the run stops
# instead of counting a batch nothing was posted from as delivered.
return 0, False
events = []
for line in lines:
try:
@@ -497,11 +715,18 @@ def flush() -> int:
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
return 0, True
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
@@ -520,10 +745,13 @@ def flush() -> int:
sent = 0
for start in range(0, len(events), BATCH_SIZE):
chunk = events[start : start + BATCH_SIZE]
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
# Carried through from record() so a resend can be collapsed.
"uuid": event.get("uuid"),
"timestamp": event.get("timestamp"),
"properties": {
# Fallback only: events recorded by a build before source
@@ -535,16 +763,24 @@ def flush() -> int:
**(event.get("properties") or {}),
},
}
for event in events[start : start + BATCH_SIZE]
for event in chunk
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
# Keep only what has not been delivered, and release the lease.
# Previously the whole file was kept and the retry re-posted every
# batch, including the ones that had already arrived.
_release_claim(claim, events[start:])
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. 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
def main() -> int:
@@ -0,0 +1,446 @@
"""Delivery semantics of the telemetry spool: no duplicates, no starvation.
These run against a built host's core in-process (not a subprocess) because they
need to inject failures into ``_post``. The identity tests next door cover the
uninitialised-process case that needs a real interpreter.
"""
from __future__ import annotations
import importlib
import json
import os
import sys
import time
from pathlib import Path
import pytest
CORE_ROOT = Path(__file__).resolve().parents[1]
REPOSITORY_ROOT = CORE_ROOT.parents[1]
HOST_CORE = REPOSITORY_ROOT / "integrations" / "claude-code-plugin" / "core"
pytestmark = pytest.mark.skipif(not HOST_CORE.exists(), reason="claude-code-plugin is not built")
@pytest.fixture()
def telemetry(tmp_path, monkeypatch):
# CI runs this directory and claude-code-plugin/tests in ONE pytest process,
# and that suite's conftest sets MEM0_TELEMETRY=false at import, process-wide.
# Without this the whole file silently no-ops: record() returns early and
# every assertion sees an empty spool. Do not rely on ambient env.
monkeypatch.setenv("MEM0_TELEMETRY", "true")
monkeypatch.setenv("MEM0_CODE_DATA_DIR", str(tmp_path / "data"))
monkeypatch.syspath_prepend(str(HOST_CORE))
# Save and RESTORE rather than delete. claude-code-plugin/tests/conftest.py
# imports memory_core once at collection and calls configure_harness() on it;
# dropping the module left a later re-import with default harness config, so
# tests in that suite failed depending on collection order.
names = ("telemetry", "memory_core", "_harness_id")
saved = {name: sys.modules.get(name) for name in names}
for name in names:
sys.modules.pop(name, None)
module = importlib.import_module("telemetry")
monkeypatch.setattr(module, "resolve_distinct_id", lambda: ("tester@example.com", ""))
try:
yield module
finally:
for name in names:
sys.modules.pop(name, None)
if saved[name] is not None:
sys.modules[name] = saved[name]
def _delivered(payloads):
return [event for payload in payloads if "batch" in payload for event in payload["batch"]]
def test_a_partial_failure_does_not_redeliver_what_already_arrived(telemetry):
"""Defect 2a: flush kept the whole claim on failure and retried from the top.
150 events across two batches, the second failing, previously delivered 250.
"""
for index in range(150):
telemetry.record("search", index=index)
sent: list[dict] = []
calls = {"n": 0}
def flaky(payload, url):
calls["n"] += 1
if calls["n"] == 2: # second batch fails
return False
sent.append(payload)
return True
telemetry._post = flaky
telemetry.flush()
telemetry._post = lambda payload, url: sent.append(payload) or True
telemetry.flush()
events = _delivered(sent)
assert len(events) == 150
assert len({event["uuid"] for event in events}) == 150
def test_a_fresh_claim_is_not_immediately_stealable(telemetry):
"""Defect 2b: rename preserves mtime, so a claim inherited the spool's age.
With the last write older than the stale threshold, a claim made now looked
abandoned the instant it existed and a second sender took it over.
"""
telemetry.record("search")
spool = telemetry._spool_path()
old = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(spool, (old, old))
first = telemetry._claim_spool()
assert first is not None
# A second sender starting right now must find nothing to take.
assert telemetry._claim_parked(first.parent) is None
def test_a_live_final_attempt_is_not_deleted_by_another_sender(telemetry):
"""Review finding: exhaustion was judged before liveness, so owners lost batches.
Claiming a parked file bumps its attempt count and refreshes its mtime. Once
the count reaches the budget, the owner draining it looked exhausted to every
other sender, which unlinked the file out from under it. Everything in that
batch was gone, which is precisely the loss this PR exists to stop.
"""
telemetry.record("search", reason="owned-by-the-first-sender")
spool = telemetry._spool_path()
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(spool, (stale, stale))
claim = telemetry._claim_spool()
assert claim is not None
# Walk it to the final attempt, ageing it each round so it can be re-claimed.
# _claim_spool hands back a0 and _release_claim keeps the name, so it takes
# one full round per attempt to reach the budget.
for _ in range(telemetry.MAX_CLAIM_ATTEMPTS):
# Carry the marker through each rewrite so the final assertion proves the
# events survived, not merely that some file with the right name did.
telemetry._release_claim(claim, [{"event": "code.search", "uuid": "owned-by-the-first-sender"}])
parked = sorted(claim.parent.glob("telemetry-*.sending"))
assert parked, "the batch was dropped while still inside its budget"
os.utime(parked[0], (stale, stale))
claim = telemetry._claim_parked(claim.parent)
assert claim is not None
assert telemetry._claim_attempt(claim) >= telemetry.MAX_CLAIM_ATTEMPTS
assert claim.exists()
# The owner is draining it right now: fresh mtime, live lease.
second_sender = telemetry._claim_parked(claim.parent)
assert second_sender is None, "a second sender took a batch under a live lease"
assert claim.exists(), "a second sender deleted a batch its owner was draining"
assert "owned-by-the-first-sender" in claim.read_text(encoding="utf-8")
def test_an_exhausted_batch_is_still_discarded_once_its_lease_lapses(telemetry):
"""The liveness check must defer the cleanup, not cancel it.
Guards the obvious over-correction: skipping live claims is only safe if an
abandoned one at the same attempt count is still reaped on a later run.
"""
telemetry.record("search")
spool = telemetry._spool_path()
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(spool, (stale, stale))
claim = telemetry._claim_spool()
assert claim is not None
exhausted = claim.parent / telemetry._claim_name(telemetry.MAX_CLAIM_ATTEMPTS)
claim.replace(exhausted)
os.utime(exhausted, (stale, stale))
assert telemetry._claim_parked(exhausted.parent) is None
assert not exhausted.exists(), "an abandoned exhausted batch was left behind forever"
def test_a_parked_batch_is_drained_behind_the_live_spool(telemetry):
"""Defect 6: parked claims were only reachable when no spool existed.
Because sessions keep recording there usually was one, so a batch parked by
a failed send waited until the 7-day expiry deleted it unsent — even though
its own presence is what starts the sender.
"""
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
old = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(parked[0], (old, old))
telemetry.record("fresh")
sent: list[dict] = []
telemetry._post = lambda payload, url: sent.append(payload) or True
telemetry.flush()
names = {event["event"] for event in _delivered(sent)}
assert names == {"code.parked", "code.fresh"}
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
assert telemetry._claim_attempt(parked[0]) < telemetry.MAX_CLAIM_ATTEMPTS
sent: list[dict] = []
telemetry._post = lambda payload, url: sent.append(payload) or True
telemetry.flush()
assert [event["event"] for event in _delivered(sent)] == ["code.parked"]
def test_progress_is_recorded_after_every_batch(telemetry):
"""A crash repeats at most one batch, not the whole file."""
for index in range(250):
telemetry.record("search", index=index)
calls = {"n": 0}
def die_after_two(payload, url):
calls["n"] += 1
if calls["n"] > 2:
return False
return True
telemetry._post = die_after_two
telemetry.flush()
parked = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending"))
assert len(parked) == 1
remaining = parked[0].read_text(encoding="utf-8").strip().splitlines()
# Two batches of 100 landed; only the last 50 should still be pending.
assert len(remaining) == 50
assert json.loads(remaining[0])["properties"]["index"] == 200
def test_the_heartbeat_actually_refreshes_the_lease(telemetry):
"""The claim rewrite doubles as the lease heartbeat.
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.
"""
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
directory = telemetry.memory_core.data_dir()
for _ in range(telemetry.MAX_CLAIM_ATTEMPTS + 3):
telemetry.flush()
# Attempts now carry a cooldown, so a released claim is not instantly
# reclaimable. Age it to stand in for the wall time a real retry waits;
# without this the loop spins inside one cooldown and proves nothing.
for parked in directory.glob("telemetry-*.sending"):
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(parked, (stale, stale))
leftover = list(directory.glob("telemetry-*.sending"))
assert leftover == [], f"batch never given up on: {[p.name for p in leftover]}"
def test_a_batch_that_cannot_be_read_is_not_counted_as_delivered(telemetry):
"""Review finding: a read failure reported 'everything delivered'.
Nothing was posted, so calling it delivered lets flush() carry on to other
claims as though this batch had arrived, and hides the failure from the one
signal that says the run went badly. It also must not quarantine: a briefly
unreadable file is retryable, and moving it to .corrupt discards the events
over a transient filesystem error, because nothing ever re-globs .corrupt.
"""
telemetry.record("search")
spool = telemetry._spool_path()
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(spool, (stale, stale))
claim = telemetry._claim_spool()
assert claim is not None
original = Path.read_text
def unreadable(self, *args, **kwargs):
if self == claim:
raise OSError(5, "I/O error")
return original(self, *args, **kwargs)
Path.read_text = unreadable
try:
sent, delivered = telemetry._drain(claim)
finally:
Path.read_text = original
assert sent == 0
assert delivered is False, "an unread batch was reported as delivered"
assert claim.exists(), "a transient read error discarded the batch"
assert not list(claim.parent.glob("*.corrupt")), "quarantined over a transient error"
def test_undecodable_content_is_still_quarantined_and_the_run_continues(telemetry):
"""The other half: genuinely unrecoverable content must not block the run.
Guards the over-correction. If every read problem returned undelivered, one
torn file would stop every later claim on every flush, forever.
"""
telemetry.record("search")
spool = telemetry._spool_path()
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(spool, (stale, stale))
claim = telemetry._claim_spool()
assert claim is not None
claim.write_bytes(b"\xff\xfe torn \x00 write")
sent, delivered = telemetry._drain(claim)
assert (sent, delivered) == (0, True)
assert not claim.exists()
assert list(claim.parent.glob("*.corrupt")), "unrecoverable content was not quarantined"
def test_retries_are_spread_over_real_time_not_burned_at_once(telemetry):
"""Review finding: releasing straight to reclaimable spent the budget instantly.
Two senders hitting one momentary failure could walk a batch from attempt 0
to the limit within seconds and discard it, when a retry a minute later would
have delivered. Each release now has to age past a cooldown that grows with
the attempts already spent.
"""
telemetry.record("doomed")
telemetry._post = lambda payload, url: False
directory = telemetry.memory_core.data_dir()
telemetry.flush()
parked = list(directory.glob("telemetry-*.sending"))
assert parked, "the batch was discarded on its first failure"
assert telemetry._claim_attempt(parked[0]) == 0
# Second sender, immediately: the cooldown has not elapsed, so it must not
# be able to spend another attempt.
telemetry.flush()
still = list(directory.glob("telemetry-*.sending"))
assert len(still) == 1
assert telemetry._claim_attempt(still[0]) <= 1, "burned attempts without waiting"
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()
def test_quarantined_batches_are_eventually_collected(telemetry):
"""Nothing re-globs .corrupt, so without a sweep they live on disk forever.
Kept much longer than .partial debris on purpose: a quarantined batch is the
only remaining evidence of events that could not be delivered.
"""
directory = telemetry.memory_core.data_dir()
directory.mkdir(parents=True, exist_ok=True)
fresh = directory / "telemetry-1-aaaaaaaa-a0.corrupt"
old = directory / "telemetry-2-bbbbbbbb-a0.corrupt"
for path in (fresh, old):
path.write_text("torn", encoding="utf-8")
expired = time.time() - (telemetry.CLAIM_EXPIRY_SECONDS + 60)
os.utime(old, (expired, expired))
telemetry._sweep_debris(directory)
assert fresh.exists(), "a recent quarantine was discarded before anyone could look at it"
assert not old.exists(), "an expired quarantine was left on disk forever"
def test_temp_files_orphaned_by_a_kill_are_collected(telemetry):
"""_write_identity and _install_salt unlink in a finally, which SIGKILL skips."""
directory = telemetry.memory_core.data_dir()
directory.mkdir(parents=True, exist_ok=True)
orphan = directory / "telemetry-salt.999.tmp"
orphan.write_text("abandoned", encoding="utf-8")
stale = time.time() - (telemetry.CLAIM_STALE_SECONDS + 60)
os.utime(orphan, (stale, stale))
telemetry._sweep_debris(directory)
assert not orphan.exists(), "a killed process left a temp file on disk forever"
+255 -19
View File
@@ -100,6 +100,16 @@ BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
# A batch is only discarded once it has genuinely been retried this many times.
MAX_CLAIM_ATTEMPTS = 3
# Parked claims drained per run, after the live spool. Bounded so a long backlog
# cannot turn one flush into an unbounded send loop.
MAX_PARKED_PER_RUN = 3
# Added to the wait before a released claim becomes reclaimable, per attempt
# already spent. Releasing straight to "reclaimable now" let two senders burn the
# whole budget within seconds of one another on a single momentary failure, and
# discard a batch a retry a minute later would have delivered.
RETRY_COOLDOWN_SECONDS = 60
def is_enabled() -> bool:
@@ -399,38 +409,201 @@ def spawn_flush() -> bool:
return False
def _claim_name(attempt: int = 0) -> str:
"""Claim filename. The attempt count rides in the name so the 7-day expiry
only ever discards a batch that was actually retried and failed."""
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
def _claim_attempt(claim: Path) -> int:
"""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
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
def _touch(path: Path) -> None:
"""Refresh mtime so a claim's age measures time since it was claimed.
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
after a quiet minute inherited the spool's last-write time and looked
abandoned the instant it was made. A second sender would then take it over
while the first was still posting, and both would deliver the batch.
"""
try:
os.utime(path, None)
except OSError:
pass
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
claim = directory / _claim_name()
spool = _spool_path()
try:
spool.replace(claim)
_touch(claim)
return claim
except OSError:
pass
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove files nothing else will ever pick up again.
*.partial is a temp file orphaned by a crash between write and rename.
*.corrupt is a batch quarantined for undecodable content. No glob in this
module matches either, so without this they accumulate on disk for the life
of the install.
Quarantined batches are kept far longer than debris: they are the only
evidence left of events that could not be delivered, and someone diagnosing
a report of missing telemetry has to be able to find one.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending")):
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
for quarantined in directory.glob("telemetry-*.corrupt"):
try:
if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS:
quarantined.unlink()
except OSError:
continue
# The same reasoning covers *.tmp. _write_identity and _install_salt both
# create one and unlink it in a finally, which a SIGKILL skips, and no glob
# in this module matches the leftovers either.
for temporary in directory.glob("telemetry-*.tmp"):
try:
if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS:
temporary.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
Kept separate from the live spool so flush() can drain both in one run.
Previously parked batches were only reachable when no spool existed at all,
and because sessions keep recording there usually was one — so a batch
parked by a failed send waited until the 7-day expiry deleted it unsent,
even though its own presence is what started the sender.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS:
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# 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:
pass
continue
if age < CLAIM_STALE_SECONDS:
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
try:
orphan.replace(claim)
_touch(claim)
return claim
except OSError:
continue
return None
def _safe_mtime(path: Path) -> float:
try:
return path.stat().st_mtime
except OSError:
return 0.0
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
"""Persist the unsent remainder, atomically, and refresh the lease.
Called after every successful batch. Two jobs: a retry resumes where the
send stopped instead of re-posting from the top, and the rewrite doubles as
the lease heartbeat, so a slow sender does not have its claim stolen
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
"""
if not remaining:
try:
claim.unlink()
except OSError:
pass
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
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
except OSError:
try:
temporary.unlink()
except OSError:
pass
return False
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
"""Persist the remainder and drop the lease, because this sender has given up.
Distinct from the per-batch heartbeat: heartbeating on the way out would
make an abandoned batch look actively owned for a further
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
threshold lets the next flush pick it up immediately, while the attempt
count in the filename still bounds how many times that can happen.
"""
if not _rewrite_claim(claim, remaining):
return
try:
# Backdate past the stale threshold so the next flush can pick it up,
# minus a cooldown that grows with the attempts already spent. Clamped so
# the mtime never lands in the future, which would read as a live lease.
cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS)
released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown
os.utime(claim, (released, released))
except OSError:
pass
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
@@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]:
def flush() -> int:
"""Drain claimed spools to PostHog and return the number of events sent."""
"""Drain the live spool, then any parked claims, and return events sent."""
if not is_enabled():
return 0
claim = _claim_spool()
sent, delivered = _drain(_claim_spool())
if not delivered:
# The network is failing. Retrying other batches now would only burn
# their attempt budget against the same broken connection.
return sent
# 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:
break
count, delivered = _drain(parked)
sent += count
if not delivered:
break
return sent
def _drain(claim: Path | None) -> tuple[int, bool]:
"""Post one claimed batch file, recording progress after every batch.
Returns (events sent, whether everything was delivered).
"""
if claim is None:
return 0
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except ValueError:
# UnicodeDecodeError from a torn write: the content is unrecoverable, so
# 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.
# Reported as delivered because there is nothing left to deliver and the
# rest of the run should continue.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
except OSError:
return 0
# Could not read it, which is not the same as having nothing to send.
# The file is left exactly where it is: a vanished or briefly unreadable
# claim is retryable, and quarantining it here would discard events over
# a transient filesystem error. Reported as undelivered so the run stops
# instead of counting a batch nothing was posted from as delivered.
return 0, False
events = []
for line in lines:
try:
@@ -497,11 +715,18 @@ def flush() -> int:
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
return 0, True
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
@@ -520,10 +745,13 @@ def flush() -> int:
sent = 0
for start in range(0, len(events), BATCH_SIZE):
chunk = events[start : start + BATCH_SIZE]
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
# Carried through from record() so a resend can be collapsed.
"uuid": event.get("uuid"),
"timestamp": event.get("timestamp"),
"properties": {
# Fallback only: events recorded by a build before source
@@ -535,16 +763,24 @@ def flush() -> int:
**(event.get("properties") or {}),
},
}
for event in events[start : start + BATCH_SIZE]
for event in chunk
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
# Keep only what has not been delivered, and release the lease.
# Previously the whole file was kept and the retry re-posted every
# batch, including the ones that had already arrived.
_release_claim(claim, events[start:])
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. 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
def main() -> int:
+255 -19
View File
@@ -100,6 +100,16 @@ BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
# A batch is only discarded once it has genuinely been retried this many times.
MAX_CLAIM_ATTEMPTS = 3
# Parked claims drained per run, after the live spool. Bounded so a long backlog
# cannot turn one flush into an unbounded send loop.
MAX_PARKED_PER_RUN = 3
# Added to the wait before a released claim becomes reclaimable, per attempt
# already spent. Releasing straight to "reclaimable now" let two senders burn the
# whole budget within seconds of one another on a single momentary failure, and
# discard a batch a retry a minute later would have delivered.
RETRY_COOLDOWN_SECONDS = 60
def is_enabled() -> bool:
@@ -399,38 +409,201 @@ def spawn_flush() -> bool:
return False
def _claim_name(attempt: int = 0) -> str:
"""Claim filename. The attempt count rides in the name so the 7-day expiry
only ever discards a batch that was actually retried and failed."""
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
def _claim_attempt(claim: Path) -> int:
"""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
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
def _touch(path: Path) -> None:
"""Refresh mtime so a claim's age measures time since it was claimed.
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
after a quiet minute inherited the spool's last-write time and looked
abandoned the instant it was made. A second sender would then take it over
while the first was still posting, and both would deliver the batch.
"""
try:
os.utime(path, None)
except OSError:
pass
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
claim = directory / _claim_name()
spool = _spool_path()
try:
spool.replace(claim)
_touch(claim)
return claim
except OSError:
pass
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove files nothing else will ever pick up again.
*.partial is a temp file orphaned by a crash between write and rename.
*.corrupt is a batch quarantined for undecodable content. No glob in this
module matches either, so without this they accumulate on disk for the life
of the install.
Quarantined batches are kept far longer than debris: they are the only
evidence left of events that could not be delivered, and someone diagnosing
a report of missing telemetry has to be able to find one.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending")):
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
for quarantined in directory.glob("telemetry-*.corrupt"):
try:
if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS:
quarantined.unlink()
except OSError:
continue
# The same reasoning covers *.tmp. _write_identity and _install_salt both
# create one and unlink it in a finally, which a SIGKILL skips, and no glob
# in this module matches the leftovers either.
for temporary in directory.glob("telemetry-*.tmp"):
try:
if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS:
temporary.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
Kept separate from the live spool so flush() can drain both in one run.
Previously parked batches were only reachable when no spool existed at all,
and because sessions keep recording there usually was one — so a batch
parked by a failed send waited until the 7-day expiry deleted it unsent,
even though its own presence is what started the sender.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS:
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# 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:
pass
continue
if age < CLAIM_STALE_SECONDS:
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
try:
orphan.replace(claim)
_touch(claim)
return claim
except OSError:
continue
return None
def _safe_mtime(path: Path) -> float:
try:
return path.stat().st_mtime
except OSError:
return 0.0
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
"""Persist the unsent remainder, atomically, and refresh the lease.
Called after every successful batch. Two jobs: a retry resumes where the
send stopped instead of re-posting from the top, and the rewrite doubles as
the lease heartbeat, so a slow sender does not have its claim stolen
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
"""
if not remaining:
try:
claim.unlink()
except OSError:
pass
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
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
except OSError:
try:
temporary.unlink()
except OSError:
pass
return False
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
"""Persist the remainder and drop the lease, because this sender has given up.
Distinct from the per-batch heartbeat: heartbeating on the way out would
make an abandoned batch look actively owned for a further
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
threshold lets the next flush pick it up immediately, while the attempt
count in the filename still bounds how many times that can happen.
"""
if not _rewrite_claim(claim, remaining):
return
try:
# Backdate past the stale threshold so the next flush can pick it up,
# minus a cooldown that grows with the attempts already spent. Clamped so
# the mtime never lands in the future, which would read as a live lease.
cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS)
released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown
os.utime(claim, (released, released))
except OSError:
pass
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
@@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]:
def flush() -> int:
"""Drain claimed spools to PostHog and return the number of events sent."""
"""Drain the live spool, then any parked claims, and return events sent."""
if not is_enabled():
return 0
claim = _claim_spool()
sent, delivered = _drain(_claim_spool())
if not delivered:
# The network is failing. Retrying other batches now would only burn
# their attempt budget against the same broken connection.
return sent
# 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:
break
count, delivered = _drain(parked)
sent += count
if not delivered:
break
return sent
def _drain(claim: Path | None) -> tuple[int, bool]:
"""Post one claimed batch file, recording progress after every batch.
Returns (events sent, whether everything was delivered).
"""
if claim is None:
return 0
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except ValueError:
# UnicodeDecodeError from a torn write: the content is unrecoverable, so
# 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.
# Reported as delivered because there is nothing left to deliver and the
# rest of the run should continue.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
except OSError:
return 0
# Could not read it, which is not the same as having nothing to send.
# The file is left exactly where it is: a vanished or briefly unreadable
# claim is retryable, and quarantining it here would discard events over
# a transient filesystem error. Reported as undelivered so the run stops
# instead of counting a batch nothing was posted from as delivered.
return 0, False
events = []
for line in lines:
try:
@@ -497,11 +715,18 @@ def flush() -> int:
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
return 0, True
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
@@ -520,10 +745,13 @@ def flush() -> int:
sent = 0
for start in range(0, len(events), BATCH_SIZE):
chunk = events[start : start + BATCH_SIZE]
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
# Carried through from record() so a resend can be collapsed.
"uuid": event.get("uuid"),
"timestamp": event.get("timestamp"),
"properties": {
# Fallback only: events recorded by a build before source
@@ -535,16 +763,24 @@ def flush() -> int:
**(event.get("properties") or {}),
},
}
for event in events[start : start + BATCH_SIZE]
for event in chunk
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
# Keep only what has not been delivered, and release the lease.
# Previously the whole file was kept and the retry re-posted every
# batch, including the ones that had already arrived.
_release_claim(claim, events[start:])
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. 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
def main() -> int:
@@ -199,9 +199,11 @@ def test_a_stale_claim_is_reclaimed(isolated_env, monkeypatch):
telemetry.record("search")
orphan = telemetry._claim_spool()
assert orphan is not None
monkeypatch.setattr(
telemetry.time, "time", lambda: orphan.stat().st_mtime + telemetry.CLAIM_STALE_SECONDS + 1
)
# Frozen rather than re-stat'd per call: flush() drains the live spool and
# then looks for parked claims in the same run, so by the second look this
# file no longer exists.
stale_now = orphan.stat().st_mtime + telemetry.CLAIM_STALE_SECONDS + 1
monkeypatch.setattr(telemetry.time, "time", lambda: stale_now)
with patch.object(telemetry, "_post", lambda payload, url: True):
assert telemetry.flush() == 1
@@ -211,9 +213,13 @@ def test_an_expired_claim_is_dropped(isolated_env, monkeypatch):
telemetry.record("search")
orphan = telemetry._claim_spool()
assert orphan is not None
monkeypatch.setattr(
telemetry.time, "time", lambda: orphan.stat().st_mtime + telemetry.CLAIM_EXPIRY_SECONDS + 1
)
expired_now = orphan.stat().st_mtime + telemetry.CLAIM_EXPIRY_SECONDS + 1
monkeypatch.setattr(telemetry.time, "time", lambda: expired_now)
# Expiry now only discards a batch that was genuinely retried and failed,
# so age alone is not enough — age it past the attempt budget too.
retried = orphan.parent / orphan.name.replace("-a0.", f"-a{telemetry.MAX_CLAIM_ATTEMPTS}.")
orphan.replace(retried)
os.utime(retried, (expired_now, expired_now - telemetry.CLAIM_EXPIRY_SECONDS - 1))
assert telemetry._claim_spool() is None
assert not list(memory_core.data_dir().glob("telemetry-*.sending"))
+255 -19
View File
@@ -100,6 +100,16 @@ BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
# A batch is only discarded once it has genuinely been retried this many times.
MAX_CLAIM_ATTEMPTS = 3
# Parked claims drained per run, after the live spool. Bounded so a long backlog
# cannot turn one flush into an unbounded send loop.
MAX_PARKED_PER_RUN = 3
# Added to the wait before a released claim becomes reclaimable, per attempt
# already spent. Releasing straight to "reclaimable now" let two senders burn the
# whole budget within seconds of one another on a single momentary failure, and
# discard a batch a retry a minute later would have delivered.
RETRY_COOLDOWN_SECONDS = 60
def is_enabled() -> bool:
@@ -399,38 +409,201 @@ def spawn_flush() -> bool:
return False
def _claim_name(attempt: int = 0) -> str:
"""Claim filename. The attempt count rides in the name so the 7-day expiry
only ever discards a batch that was actually retried and failed."""
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
def _claim_attempt(claim: Path) -> int:
"""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
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
def _touch(path: Path) -> None:
"""Refresh mtime so a claim's age measures time since it was claimed.
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
after a quiet minute inherited the spool's last-write time and looked
abandoned the instant it was made. A second sender would then take it over
while the first was still posting, and both would deliver the batch.
"""
try:
os.utime(path, None)
except OSError:
pass
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
claim = directory / _claim_name()
spool = _spool_path()
try:
spool.replace(claim)
_touch(claim)
return claim
except OSError:
pass
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove files nothing else will ever pick up again.
*.partial is a temp file orphaned by a crash between write and rename.
*.corrupt is a batch quarantined for undecodable content. No glob in this
module matches either, so without this they accumulate on disk for the life
of the install.
Quarantined batches are kept far longer than debris: they are the only
evidence left of events that could not be delivered, and someone diagnosing
a report of missing telemetry has to be able to find one.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending")):
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
for quarantined in directory.glob("telemetry-*.corrupt"):
try:
if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS:
quarantined.unlink()
except OSError:
continue
# The same reasoning covers *.tmp. _write_identity and _install_salt both
# create one and unlink it in a finally, which a SIGKILL skips, and no glob
# in this module matches the leftovers either.
for temporary in directory.glob("telemetry-*.tmp"):
try:
if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS:
temporary.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
Kept separate from the live spool so flush() can drain both in one run.
Previously parked batches were only reachable when no spool existed at all,
and because sessions keep recording there usually was one — so a batch
parked by a failed send waited until the 7-day expiry deleted it unsent,
even though its own presence is what started the sender.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS:
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# 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:
pass
continue
if age < CLAIM_STALE_SECONDS:
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
try:
orphan.replace(claim)
_touch(claim)
return claim
except OSError:
continue
return None
def _safe_mtime(path: Path) -> float:
try:
return path.stat().st_mtime
except OSError:
return 0.0
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
"""Persist the unsent remainder, atomically, and refresh the lease.
Called after every successful batch. Two jobs: a retry resumes where the
send stopped instead of re-posting from the top, and the rewrite doubles as
the lease heartbeat, so a slow sender does not have its claim stolen
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
"""
if not remaining:
try:
claim.unlink()
except OSError:
pass
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
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
except OSError:
try:
temporary.unlink()
except OSError:
pass
return False
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
"""Persist the remainder and drop the lease, because this sender has given up.
Distinct from the per-batch heartbeat: heartbeating on the way out would
make an abandoned batch look actively owned for a further
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
threshold lets the next flush pick it up immediately, while the attempt
count in the filename still bounds how many times that can happen.
"""
if not _rewrite_claim(claim, remaining):
return
try:
# Backdate past the stale threshold so the next flush can pick it up,
# minus a cooldown that grows with the attempts already spent. Clamped so
# the mtime never lands in the future, which would read as a live lease.
cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS)
released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown
os.utime(claim, (released, released))
except OSError:
pass
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
@@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]:
def flush() -> int:
"""Drain claimed spools to PostHog and return the number of events sent."""
"""Drain the live spool, then any parked claims, and return events sent."""
if not is_enabled():
return 0
claim = _claim_spool()
sent, delivered = _drain(_claim_spool())
if not delivered:
# The network is failing. Retrying other batches now would only burn
# their attempt budget against the same broken connection.
return sent
# 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:
break
count, delivered = _drain(parked)
sent += count
if not delivered:
break
return sent
def _drain(claim: Path | None) -> tuple[int, bool]:
"""Post one claimed batch file, recording progress after every batch.
Returns (events sent, whether everything was delivered).
"""
if claim is None:
return 0
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except ValueError:
# UnicodeDecodeError from a torn write: the content is unrecoverable, so
# 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.
# Reported as delivered because there is nothing left to deliver and the
# rest of the run should continue.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
except OSError:
return 0
# Could not read it, which is not the same as having nothing to send.
# The file is left exactly where it is: a vanished or briefly unreadable
# claim is retryable, and quarantining it here would discard events over
# a transient filesystem error. Reported as undelivered so the run stops
# instead of counting a batch nothing was posted from as delivered.
return 0, False
events = []
for line in lines:
try:
@@ -497,11 +715,18 @@ def flush() -> int:
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
return 0, True
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
@@ -520,10 +745,13 @@ def flush() -> int:
sent = 0
for start in range(0, len(events), BATCH_SIZE):
chunk = events[start : start + BATCH_SIZE]
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
# Carried through from record() so a resend can be collapsed.
"uuid": event.get("uuid"),
"timestamp": event.get("timestamp"),
"properties": {
# Fallback only: events recorded by a build before source
@@ -535,16 +763,24 @@ def flush() -> int:
**(event.get("properties") or {}),
},
}
for event in events[start : start + BATCH_SIZE]
for event in chunk
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
# Keep only what has not been delivered, and release the lease.
# Previously the whole file was kept and the retry re-posted every
# batch, including the ones that had already arrived.
_release_claim(claim, events[start:])
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. 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
def main() -> int:
+255 -19
View File
@@ -100,6 +100,16 @@ BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
# A batch is only discarded once it has genuinely been retried this many times.
MAX_CLAIM_ATTEMPTS = 3
# Parked claims drained per run, after the live spool. Bounded so a long backlog
# cannot turn one flush into an unbounded send loop.
MAX_PARKED_PER_RUN = 3
# Added to the wait before a released claim becomes reclaimable, per attempt
# already spent. Releasing straight to "reclaimable now" let two senders burn the
# whole budget within seconds of one another on a single momentary failure, and
# discard a batch a retry a minute later would have delivered.
RETRY_COOLDOWN_SECONDS = 60
def is_enabled() -> bool:
@@ -399,38 +409,201 @@ def spawn_flush() -> bool:
return False
def _claim_name(attempt: int = 0) -> str:
"""Claim filename. The attempt count rides in the name so the 7-day expiry
only ever discards a batch that was actually retried and failed."""
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
def _claim_attempt(claim: Path) -> int:
"""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
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
def _touch(path: Path) -> None:
"""Refresh mtime so a claim's age measures time since it was claimed.
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
after a quiet minute inherited the spool's last-write time and looked
abandoned the instant it was made. A second sender would then take it over
while the first was still posting, and both would deliver the batch.
"""
try:
os.utime(path, None)
except OSError:
pass
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
claim = directory / _claim_name()
spool = _spool_path()
try:
spool.replace(claim)
_touch(claim)
return claim
except OSError:
pass
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove files nothing else will ever pick up again.
*.partial is a temp file orphaned by a crash between write and rename.
*.corrupt is a batch quarantined for undecodable content. No glob in this
module matches either, so without this they accumulate on disk for the life
of the install.
Quarantined batches are kept far longer than debris: they are the only
evidence left of events that could not be delivered, and someone diagnosing
a report of missing telemetry has to be able to find one.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending")):
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
for quarantined in directory.glob("telemetry-*.corrupt"):
try:
if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS:
quarantined.unlink()
except OSError:
continue
# The same reasoning covers *.tmp. _write_identity and _install_salt both
# create one and unlink it in a finally, which a SIGKILL skips, and no glob
# in this module matches the leftovers either.
for temporary in directory.glob("telemetry-*.tmp"):
try:
if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS:
temporary.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
Kept separate from the live spool so flush() can drain both in one run.
Previously parked batches were only reachable when no spool existed at all,
and because sessions keep recording there usually was one — so a batch
parked by a failed send waited until the 7-day expiry deleted it unsent,
even though its own presence is what started the sender.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS:
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# 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:
pass
continue
if age < CLAIM_STALE_SECONDS:
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
try:
orphan.replace(claim)
_touch(claim)
return claim
except OSError:
continue
return None
def _safe_mtime(path: Path) -> float:
try:
return path.stat().st_mtime
except OSError:
return 0.0
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
"""Persist the unsent remainder, atomically, and refresh the lease.
Called after every successful batch. Two jobs: a retry resumes where the
send stopped instead of re-posting from the top, and the rewrite doubles as
the lease heartbeat, so a slow sender does not have its claim stolen
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
"""
if not remaining:
try:
claim.unlink()
except OSError:
pass
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
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
except OSError:
try:
temporary.unlink()
except OSError:
pass
return False
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
"""Persist the remainder and drop the lease, because this sender has given up.
Distinct from the per-batch heartbeat: heartbeating on the way out would
make an abandoned batch look actively owned for a further
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
threshold lets the next flush pick it up immediately, while the attempt
count in the filename still bounds how many times that can happen.
"""
if not _rewrite_claim(claim, remaining):
return
try:
# Backdate past the stale threshold so the next flush can pick it up,
# minus a cooldown that grows with the attempts already spent. Clamped so
# the mtime never lands in the future, which would read as a live lease.
cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS)
released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown
os.utime(claim, (released, released))
except OSError:
pass
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
@@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]:
def flush() -> int:
"""Drain claimed spools to PostHog and return the number of events sent."""
"""Drain the live spool, then any parked claims, and return events sent."""
if not is_enabled():
return 0
claim = _claim_spool()
sent, delivered = _drain(_claim_spool())
if not delivered:
# The network is failing. Retrying other batches now would only burn
# their attempt budget against the same broken connection.
return sent
# 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:
break
count, delivered = _drain(parked)
sent += count
if not delivered:
break
return sent
def _drain(claim: Path | None) -> tuple[int, bool]:
"""Post one claimed batch file, recording progress after every batch.
Returns (events sent, whether everything was delivered).
"""
if claim is None:
return 0
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except ValueError:
# UnicodeDecodeError from a torn write: the content is unrecoverable, so
# 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.
# Reported as delivered because there is nothing left to deliver and the
# rest of the run should continue.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
except OSError:
return 0
# Could not read it, which is not the same as having nothing to send.
# The file is left exactly where it is: a vanished or briefly unreadable
# claim is retryable, and quarantining it here would discard events over
# a transient filesystem error. Reported as undelivered so the run stops
# instead of counting a batch nothing was posted from as delivered.
return 0, False
events = []
for line in lines:
try:
@@ -497,11 +715,18 @@ def flush() -> int:
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
return 0, True
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
@@ -520,10 +745,13 @@ def flush() -> int:
sent = 0
for start in range(0, len(events), BATCH_SIZE):
chunk = events[start : start + BATCH_SIZE]
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
# Carried through from record() so a resend can be collapsed.
"uuid": event.get("uuid"),
"timestamp": event.get("timestamp"),
"properties": {
# Fallback only: events recorded by a build before source
@@ -535,16 +763,24 @@ def flush() -> int:
**(event.get("properties") or {}),
},
}
for event in events[start : start + BATCH_SIZE]
for event in chunk
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
# Keep only what has not been delivered, and release the lease.
# Previously the whole file was kept and the retry re-posted every
# batch, including the ones that had already arrived.
_release_claim(claim, events[start:])
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. 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
def main() -> int:
+255 -19
View File
@@ -100,6 +100,16 @@ BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
# A batch is only discarded once it has genuinely been retried this many times.
MAX_CLAIM_ATTEMPTS = 3
# Parked claims drained per run, after the live spool. Bounded so a long backlog
# cannot turn one flush into an unbounded send loop.
MAX_PARKED_PER_RUN = 3
# Added to the wait before a released claim becomes reclaimable, per attempt
# already spent. Releasing straight to "reclaimable now" let two senders burn the
# whole budget within seconds of one another on a single momentary failure, and
# discard a batch a retry a minute later would have delivered.
RETRY_COOLDOWN_SECONDS = 60
def is_enabled() -> bool:
@@ -399,38 +409,201 @@ def spawn_flush() -> bool:
return False
def _claim_name(attempt: int = 0) -> str:
"""Claim filename. The attempt count rides in the name so the 7-day expiry
only ever discards a batch that was actually retried and failed."""
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
def _claim_attempt(claim: Path) -> int:
"""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
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
def _touch(path: Path) -> None:
"""Refresh mtime so a claim's age measures time since it was claimed.
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
after a quiet minute inherited the spool's last-write time and looked
abandoned the instant it was made. A second sender would then take it over
while the first was still posting, and both would deliver the batch.
"""
try:
os.utime(path, None)
except OSError:
pass
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
claim = directory / _claim_name()
spool = _spool_path()
try:
spool.replace(claim)
_touch(claim)
return claim
except OSError:
pass
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove files nothing else will ever pick up again.
*.partial is a temp file orphaned by a crash between write and rename.
*.corrupt is a batch quarantined for undecodable content. No glob in this
module matches either, so without this they accumulate on disk for the life
of the install.
Quarantined batches are kept far longer than debris: they are the only
evidence left of events that could not be delivered, and someone diagnosing
a report of missing telemetry has to be able to find one.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending")):
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
for quarantined in directory.glob("telemetry-*.corrupt"):
try:
if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS:
quarantined.unlink()
except OSError:
continue
# The same reasoning covers *.tmp. _write_identity and _install_salt both
# create one and unlink it in a finally, which a SIGKILL skips, and no glob
# in this module matches the leftovers either.
for temporary in directory.glob("telemetry-*.tmp"):
try:
if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS:
temporary.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
Kept separate from the live spool so flush() can drain both in one run.
Previously parked batches were only reachable when no spool existed at all,
and because sessions keep recording there usually was one — so a batch
parked by a failed send waited until the 7-day expiry deleted it unsent,
even though its own presence is what started the sender.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS:
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# 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:
pass
continue
if age < CLAIM_STALE_SECONDS:
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
try:
orphan.replace(claim)
_touch(claim)
return claim
except OSError:
continue
return None
def _safe_mtime(path: Path) -> float:
try:
return path.stat().st_mtime
except OSError:
return 0.0
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
"""Persist the unsent remainder, atomically, and refresh the lease.
Called after every successful batch. Two jobs: a retry resumes where the
send stopped instead of re-posting from the top, and the rewrite doubles as
the lease heartbeat, so a slow sender does not have its claim stolen
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
"""
if not remaining:
try:
claim.unlink()
except OSError:
pass
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
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
except OSError:
try:
temporary.unlink()
except OSError:
pass
return False
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
"""Persist the remainder and drop the lease, because this sender has given up.
Distinct from the per-batch heartbeat: heartbeating on the way out would
make an abandoned batch look actively owned for a further
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
threshold lets the next flush pick it up immediately, while the attempt
count in the filename still bounds how many times that can happen.
"""
if not _rewrite_claim(claim, remaining):
return
try:
# Backdate past the stale threshold so the next flush can pick it up,
# minus a cooldown that grows with the attempts already spent. Clamped so
# the mtime never lands in the future, which would read as a live lease.
cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS)
released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown
os.utime(claim, (released, released))
except OSError:
pass
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
@@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]:
def flush() -> int:
"""Drain claimed spools to PostHog and return the number of events sent."""
"""Drain the live spool, then any parked claims, and return events sent."""
if not is_enabled():
return 0
claim = _claim_spool()
sent, delivered = _drain(_claim_spool())
if not delivered:
# The network is failing. Retrying other batches now would only burn
# their attempt budget against the same broken connection.
return sent
# 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:
break
count, delivered = _drain(parked)
sent += count
if not delivered:
break
return sent
def _drain(claim: Path | None) -> tuple[int, bool]:
"""Post one claimed batch file, recording progress after every batch.
Returns (events sent, whether everything was delivered).
"""
if claim is None:
return 0
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except ValueError:
# UnicodeDecodeError from a torn write: the content is unrecoverable, so
# 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.
# Reported as delivered because there is nothing left to deliver and the
# rest of the run should continue.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
except OSError:
return 0
# Could not read it, which is not the same as having nothing to send.
# The file is left exactly where it is: a vanished or briefly unreadable
# claim is retryable, and quarantining it here would discard events over
# a transient filesystem error. Reported as undelivered so the run stops
# instead of counting a batch nothing was posted from as delivered.
return 0, False
events = []
for line in lines:
try:
@@ -497,11 +715,18 @@ def flush() -> int:
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
return 0, True
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
@@ -520,10 +745,13 @@ def flush() -> int:
sent = 0
for start in range(0, len(events), BATCH_SIZE):
chunk = events[start : start + BATCH_SIZE]
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
# Carried through from record() so a resend can be collapsed.
"uuid": event.get("uuid"),
"timestamp": event.get("timestamp"),
"properties": {
# Fallback only: events recorded by a build before source
@@ -535,16 +763,24 @@ def flush() -> int:
**(event.get("properties") or {}),
},
}
for event in events[start : start + BATCH_SIZE]
for event in chunk
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
# Keep only what has not been delivered, and release the lease.
# Previously the whole file was kept and the retry re-posted every
# batch, including the ones that had already arrived.
_release_claim(claim, events[start:])
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. 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
def main() -> int:
+255 -19
View File
@@ -100,6 +100,16 @@ BATCH_SIZE = 100
SEND_TIMEOUT = 5
CLAIM_STALE_SECONDS = 120
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
# A batch is only discarded once it has genuinely been retried this many times.
MAX_CLAIM_ATTEMPTS = 3
# Parked claims drained per run, after the live spool. Bounded so a long backlog
# cannot turn one flush into an unbounded send loop.
MAX_PARKED_PER_RUN = 3
# Added to the wait before a released claim becomes reclaimable, per attempt
# already spent. Releasing straight to "reclaimable now" let two senders burn the
# whole budget within seconds of one another on a single momentary failure, and
# discard a batch a retry a minute later would have delivered.
RETRY_COOLDOWN_SECONDS = 60
def is_enabled() -> bool:
@@ -399,38 +409,201 @@ def spawn_flush() -> bool:
return False
def _claim_name(attempt: int = 0) -> str:
"""Claim filename. The attempt count rides in the name so the 7-day expiry
only ever discards a batch that was actually retried and failed."""
return f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}-a{attempt}.sending"
def _claim_attempt(claim: Path) -> int:
"""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
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
def _touch(path: Path) -> None:
"""Refresh mtime so a claim's age measures time since it was claimed.
``Path.replace`` is ``os.rename``, which preserves mtime — so a claim created
after a quiet minute inherited the spool's last-write time and looked
abandoned the instant it was made. A second sender would then take it over
while the first was still posting, and both would deliver the batch.
"""
try:
os.utime(path, None)
except OSError:
pass
def _claim_spool() -> Path | None:
"""Rename the spool aside so exactly one sender owns each batch."""
directory = memory_core.data_dir()
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
claim = directory / _claim_name()
spool = _spool_path()
try:
spool.replace(claim)
_touch(claim)
return claim
except OSError:
pass
return _claim_parked(directory)
def _sweep_debris(directory: Path) -> None:
"""Remove files nothing else will ever pick up again.
*.partial is a temp file orphaned by a crash between write and rename.
*.corrupt is a batch quarantined for undecodable content. No glob in this
module matches either, so without this they accumulate on disk for the life
of the install.
Quarantined batches are kept far longer than debris: they are the only
evidence left of events that could not be delivered, and someone diagnosing
a report of missing telemetry has to be able to find one.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending")):
for debris in directory.glob("telemetry-*.partial"):
try:
if now - debris.stat().st_mtime > CLAIM_STALE_SECONDS:
debris.unlink()
except OSError:
continue
for quarantined in directory.glob("telemetry-*.corrupt"):
try:
if now - quarantined.stat().st_mtime > CLAIM_EXPIRY_SECONDS:
quarantined.unlink()
except OSError:
continue
# The same reasoning covers *.tmp. _write_identity and _install_salt both
# create one and unlink it in a finally, which a SIGKILL skips, and no glob
# in this module matches the leftovers either.
for temporary in directory.glob("telemetry-*.tmp"):
try:
if now - temporary.stat().st_mtime > CLAIM_STALE_SECONDS:
temporary.unlink()
except OSError:
continue
def _claim_parked(directory: Path) -> Path | None:
"""Take the oldest abandoned claim, if any lease has actually expired.
Kept separate from the live spool so flush() can drain both in one run.
Previously parked batches were only reachable when no spool existed at all,
and because sessions keep recording there usually was one — so a batch
parked by a failed send waited until the 7-day expiry deleted it unsent,
even though its own presence is what started the sender.
"""
now = time.time()
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
try:
age = now - orphan.stat().st_mtime
except OSError:
continue
if age > CLAIM_EXPIRY_SECONDS:
if age < CLAIM_STALE_SECONDS:
# Someone else holds a live lease on it. This check has to come
# first. Claiming a file bumps its attempt count and refreshes its
# mtime, so a sender that has just taken the final attempt looks
# exhausted to everyone else while it is actively draining. Judging
# exhaustion before liveness let a second sender unlink a batch out
# from under its owner, losing every event in it.
continue
# 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:
pass
continue
if age < CLAIM_STALE_SECONDS:
continue
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
try:
orphan.replace(claim)
_touch(claim)
return claim
except OSError:
continue
return None
def _safe_mtime(path: Path) -> float:
try:
return path.stat().st_mtime
except OSError:
return 0.0
def _rewrite_claim(claim: Path, remaining: list[dict[str, Any]]) -> bool:
"""Persist the unsent remainder, atomically, and refresh the lease.
Called after every successful batch. Two jobs: a retry resumes where the
send stopped instead of re-posting from the top, and the rewrite doubles as
the lease heartbeat, so a slow sender does not have its claim stolen
mid-flight. Interval is one batch, well inside CLAIM_STALE_SECONDS.
"""
if not remaining:
try:
claim.unlink()
except OSError:
pass
return True
temporary = claim.with_suffix(f".{os.getpid()}.partial")
try:
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
except OSError:
try:
temporary.unlink()
except OSError:
pass
return False
def _release_claim(claim: Path, remaining: list[dict[str, Any]]) -> None:
"""Persist the remainder and drop the lease, because this sender has given up.
Distinct from the per-batch heartbeat: heartbeating on the way out would
make an abandoned batch look actively owned for a further
CLAIM_STALE_SECONDS, delaying the retry for no reason. Ageing it past the
threshold lets the next flush pick it up immediately, while the attempt
count in the filename still bounds how many times that can happen.
"""
if not _rewrite_claim(claim, remaining):
return
try:
# Backdate past the stale threshold so the next flush can pick it up,
# minus a cooldown that grows with the attempts already spent. Clamped so
# the mtime never lands in the future, which would read as a live lease.
cooldown = min(_claim_attempt(claim) * RETRY_COOLDOWN_SECONDS, CLAIM_STALE_SECONDS)
released = time.time() - CLAIM_STALE_SECONDS - 1 + cooldown
os.utime(claim, (released, released))
except OSError:
pass
def _resolve_email(key: str) -> str:
"""Trade the API key for the account email so events join other Mem0 surfaces."""
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
@@ -478,16 +651,61 @@ def resolve_distinct_id() -> tuple[str, str]:
def flush() -> int:
"""Drain claimed spools to PostHog and return the number of events sent."""
"""Drain the live spool, then any parked claims, and return events sent."""
if not is_enabled():
return 0
claim = _claim_spool()
sent, delivered = _drain(_claim_spool())
if not delivered:
# The network is failing. Retrying other batches now would only burn
# their attempt budget against the same broken connection.
return sent
# 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:
break
count, delivered = _drain(parked)
sent += count
if not delivered:
break
return sent
def _drain(claim: Path | None) -> tuple[int, bool]:
"""Post one claimed batch file, recording progress after every batch.
Returns (events sent, whether everything was delivered).
"""
if claim is None:
return 0
return 0, True
try:
lines = claim.read_text(encoding="utf-8").splitlines()
except ValueError:
# UnicodeDecodeError from a torn write: the content is unrecoverable, so
# 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.
# Reported as delivered because there is nothing left to deliver and the
# rest of the run should continue.
try:
claim.replace(claim.with_suffix(".corrupt"))
except OSError:
try:
claim.unlink()
except OSError:
pass
return 0, True
except OSError:
return 0
# Could not read it, which is not the same as having nothing to send.
# The file is left exactly where it is: a vanished or briefly unreadable
# claim is retryable, and quarantining it here would discard events over
# a transient filesystem error. Reported as undelivered so the run stops
# instead of counting a batch nothing was posted from as delivered.
return 0, False
events = []
for line in lines:
try:
@@ -497,11 +715,18 @@ def flush() -> int:
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
return 0, True
distinct_id, aliased_anonymous_id = resolve_distinct_id()
if aliased_anonymous_id:
@@ -520,10 +745,13 @@ def flush() -> int:
sent = 0
for start in range(0, len(events), BATCH_SIZE):
chunk = events[start : start + BATCH_SIZE]
batch = [
{
"event": event["event"],
"distinct_id": distinct_id,
# Carried through from record() so a resend can be collapsed.
"uuid": event.get("uuid"),
"timestamp": event.get("timestamp"),
"properties": {
# Fallback only: events recorded by a build before source
@@ -535,16 +763,24 @@ def flush() -> int:
**(event.get("properties") or {}),
},
}
for event in events[start : start + BATCH_SIZE]
for event in chunk
]
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
return sent
sent += len(batch)
try:
claim.unlink()
except OSError:
pass
return sent
# Keep only what has not been delivered, and release the lease.
# Previously the whole file was kept and the retry re-posted every
# batch, including the ones that had already arrived.
_release_claim(claim, events[start:])
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. 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
def main() -> int: