fix(plugins): stop delivering telemetry events twice, and stop losing parked ones
Two defects, one cause: the spool protocol infers ownership instead of holding it, and never records progress. Duplicate delivery after a partial failure. flush() posts the claim in batches of 100 and returns on the first failure, keeping the whole file. The retry then posts every batch again, including the ones that already arrived — 150 recorded events were delivered 250 times. Progress is now written back to the claim after each successful batch, so a retry resumes where the send stopped and a crash repeats at most one batch. Duplicate delivery when two senders overlap. spool.replace(claim) 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 existed. A second sender starting while the first was still posting took it over and sent it too — most likely at session end, when the MCP server's exit sender and the SessionEnd flush worker both drain. Claims are now touched at claim time, and the per-batch rewrite doubles as a lease heartbeat. _post makes one attempt with SEND_TIMEOUT and no retry, so a heartbeat lands well inside the 120s lease; a test asserts that margin so adding a retry loop to _post cannot silently break it. Parked batches starved. _claim_spool only looked at parked .sending files when no spool existed, 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, despite its own presence being what starts the sender in the first place. flush() now drains the live spool and then parked claims in the same run, oldest first, bounded. Expiry applies only after a genuine retry has failed, with the attempt count carried in the filename. A sender that gives up releases its lease rather than heartbeating on the way out, so the next run picks the batch up promptly instead of waiting a full stale window for a batch nobody is working on. A failing send stops the run, so one broken connection cannot burn every parked batch's attempt budget at once. Two existing tests asserted the old lifecycle and are updated in place, each with a comment saying what changed. Claude-Session: https://claude.ai/code/session_01C7tEmH86HAr7GoAAKCEHZb
This commit is contained in:
@@ -99,6 +99,11 @@ BATCH_SIZE = 100
|
|||||||
SEND_TIMEOUT = 5
|
SEND_TIMEOUT = 5
|
||||||
CLAIM_STALE_SECONDS = 120
|
CLAIM_STALE_SECONDS = 120
|
||||||
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
|
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
|
||||||
|
|
||||||
|
|
||||||
def is_enabled() -> bool:
|
def is_enabled() -> bool:
|
||||||
@@ -301,38 +306,139 @@ def spawn_flush() -> bool:
|
|||||||
return False
|
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."""
|
||||||
|
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||||
|
tail = stem.rsplit("-", 1)[-1]
|
||||||
|
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:
|
def _claim_spool() -> Path | None:
|
||||||
"""Rename the spool aside so exactly one sender owns each batch."""
|
"""Rename the spool aside so exactly one sender owns each batch."""
|
||||||
directory = memory_core.data_dir()
|
directory = memory_core.data_dir()
|
||||||
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
|
claim = directory / _claim_name()
|
||||||
spool = _spool_path()
|
spool = _spool_path()
|
||||||
try:
|
try:
|
||||||
spool.replace(claim)
|
spool.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
|
return _claim_parked(directory)
|
||||||
|
|
||||||
|
|
||||||
|
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()
|
now = time.time()
|
||||||
for orphan in sorted(directory.glob("telemetry-*.sending")):
|
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
|
||||||
try:
|
try:
|
||||||
age = now - orphan.stat().st_mtime
|
age = now - orphan.stat().st_mtime
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
if age > CLAIM_EXPIRY_SECONDS:
|
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||||
try:
|
try:
|
||||||
orphan.unlink()
|
orphan.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
continue
|
continue
|
||||||
if age < CLAIM_STALE_SECONDS:
|
if age < CLAIM_STALE_SECONDS:
|
||||||
|
# Someone else holds a live lease on it.
|
||||||
continue
|
continue
|
||||||
|
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
|
||||||
try:
|
try:
|
||||||
orphan.replace(claim)
|
orphan.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
return None
|
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:
|
||||||
|
temporary.write_text(
|
||||||
|
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||||
|
encoding="utf-8",
|
||||||
|
)
|
||||||
|
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:
|
||||||
|
released = time.time() - CLAIM_STALE_SECONDS - 1
|
||||||
|
os.utime(claim, (released, released))
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
def _resolve_email(key: str) -> str:
|
def _resolve_email(key: str) -> str:
|
||||||
"""Trade the API key for the account email so events join other Mem0 surfaces."""
|
"""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/"
|
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
|
||||||
@@ -380,16 +486,40 @@ def resolve_distinct_id() -> tuple[str, str]:
|
|||||||
|
|
||||||
|
|
||||||
def flush() -> int:
|
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():
|
if not is_enabled():
|
||||||
return 0
|
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()
|
||||||
|
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:
|
if claim is None:
|
||||||
return 0
|
return 0, True
|
||||||
try:
|
try:
|
||||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||||
except OSError:
|
except OSError:
|
||||||
return 0
|
return 0, True
|
||||||
events = []
|
events = []
|
||||||
for line in lines:
|
for line in lines:
|
||||||
try:
|
try:
|
||||||
@@ -403,7 +533,7 @@ def flush() -> int:
|
|||||||
claim.unlink()
|
claim.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
return 0
|
return 0, True
|
||||||
|
|
||||||
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
||||||
if aliased_anonymous_id:
|
if aliased_anonymous_id:
|
||||||
@@ -422,10 +552,13 @@ def flush() -> int:
|
|||||||
|
|
||||||
sent = 0
|
sent = 0
|
||||||
for start in range(0, len(events), BATCH_SIZE):
|
for start in range(0, len(events), BATCH_SIZE):
|
||||||
|
chunk = events[start : start + BATCH_SIZE]
|
||||||
batch = [
|
batch = [
|
||||||
{
|
{
|
||||||
"event": event["event"],
|
"event": event["event"],
|
||||||
"distinct_id": distinct_id,
|
"distinct_id": distinct_id,
|
||||||
|
# Carried through from record() so a resend can be collapsed.
|
||||||
|
"uuid": event.get("uuid"),
|
||||||
"timestamp": event.get("timestamp"),
|
"timestamp": event.get("timestamp"),
|
||||||
"properties": {
|
"properties": {
|
||||||
# Fallback only: events recorded by a build before source
|
# Fallback only: events recorded by a build before source
|
||||||
@@ -437,16 +570,19 @@ def flush() -> int:
|
|||||||
**(event.get("properties") or {}),
|
**(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):
|
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
|
||||||
return sent
|
# Keep only what has not been delivered, and release the lease.
|
||||||
sent += len(batch)
|
# Previously the whole file was kept and the retry re-posted every
|
||||||
try:
|
# batch, including the ones that had already arrived.
|
||||||
claim.unlink()
|
_release_claim(claim, events[start:])
|
||||||
except OSError:
|
return sent, False
|
||||||
pass
|
sent += len(chunk)
|
||||||
return sent
|
# Record progress and refresh the lease after each successful batch, so
|
||||||
|
# a crash repeats at most one batch instead of the entire file.
|
||||||
|
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||||
|
return sent, True
|
||||||
|
|
||||||
|
|
||||||
def main() -> int:
|
def main() -> int:
|
||||||
|
|||||||
@@ -0,0 +1,164 @@
|
|||||||
|
"""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):
|
||||||
|
monkeypatch.setenv("MEM0_CODE_DATA_DIR", str(tmp_path / "data"))
|
||||||
|
monkeypatch.syspath_prepend(str(HOST_CORE))
|
||||||
|
for name in ("telemetry", "memory_core", "_harness_id"):
|
||||||
|
sys.modules.pop(name, None)
|
||||||
|
module = importlib.import_module("telemetry")
|
||||||
|
monkeypatch.setattr(module, "resolve_distinct_id", lambda: ("tester@example.com", ""))
|
||||||
|
yield module
|
||||||
|
for name in ("telemetry", "memory_core", "_harness_id"):
|
||||||
|
sys.modules.pop(name, None)
|
||||||
|
|
||||||
|
|
||||||
|
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_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_an_untried_batch_is_not_expired_by_age_alone(telemetry):
|
||||||
|
"""Expiry should discard what failed, not what never got a turn."""
|
||||||
|
telemetry.record("parked")
|
||||||
|
telemetry._post = lambda payload, url: False
|
||||||
|
telemetry.flush()
|
||||||
|
|
||||||
|
parked = list(telemetry.memory_core.data_dir().glob("telemetry-*.sending"))
|
||||||
|
assert len(parked) == 1
|
||||||
|
ancient = time.time() - (telemetry.CLAIM_EXPIRY_SECONDS + 60)
|
||||||
|
os.utime(parked[0], (ancient, ancient))
|
||||||
|
|
||||||
|
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_stays_well_inside_the_lease(telemetry):
|
||||||
|
"""The claim rewrite doubles as the lease heartbeat.
|
||||||
|
|
||||||
|
_post makes a single attempt with SEND_TIMEOUT and no retry, so a heartbeat
|
||||||
|
lands at least that often. If a retry loop is ever added to _post, this is
|
||||||
|
the assertion that catches a sender losing its claim mid-flight.
|
||||||
|
"""
|
||||||
|
assert telemetry.SEND_TIMEOUT * 4 < telemetry.CLAIM_STALE_SECONDS
|
||||||
@@ -99,6 +99,11 @@ BATCH_SIZE = 100
|
|||||||
SEND_TIMEOUT = 5
|
SEND_TIMEOUT = 5
|
||||||
CLAIM_STALE_SECONDS = 120
|
CLAIM_STALE_SECONDS = 120
|
||||||
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
|
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
|
||||||
|
|
||||||
|
|
||||||
def is_enabled() -> bool:
|
def is_enabled() -> bool:
|
||||||
@@ -301,38 +306,139 @@ def spawn_flush() -> bool:
|
|||||||
return False
|
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."""
|
||||||
|
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||||
|
tail = stem.rsplit("-", 1)[-1]
|
||||||
|
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:
|
def _claim_spool() -> Path | None:
|
||||||
"""Rename the spool aside so exactly one sender owns each batch."""
|
"""Rename the spool aside so exactly one sender owns each batch."""
|
||||||
directory = memory_core.data_dir()
|
directory = memory_core.data_dir()
|
||||||
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
|
claim = directory / _claim_name()
|
||||||
spool = _spool_path()
|
spool = _spool_path()
|
||||||
try:
|
try:
|
||||||
spool.replace(claim)
|
spool.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
|
return _claim_parked(directory)
|
||||||
|
|
||||||
|
|
||||||
|
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()
|
now = time.time()
|
||||||
for orphan in sorted(directory.glob("telemetry-*.sending")):
|
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
|
||||||
try:
|
try:
|
||||||
age = now - orphan.stat().st_mtime
|
age = now - orphan.stat().st_mtime
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
if age > CLAIM_EXPIRY_SECONDS:
|
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||||
try:
|
try:
|
||||||
orphan.unlink()
|
orphan.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
continue
|
continue
|
||||||
if age < CLAIM_STALE_SECONDS:
|
if age < CLAIM_STALE_SECONDS:
|
||||||
|
# Someone else holds a live lease on it.
|
||||||
continue
|
continue
|
||||||
|
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
|
||||||
try:
|
try:
|
||||||
orphan.replace(claim)
|
orphan.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
return None
|
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:
|
||||||
|
temporary.write_text(
|
||||||
|
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||||
|
encoding="utf-8",
|
||||||
|
)
|
||||||
|
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:
|
||||||
|
released = time.time() - CLAIM_STALE_SECONDS - 1
|
||||||
|
os.utime(claim, (released, released))
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
def _resolve_email(key: str) -> str:
|
def _resolve_email(key: str) -> str:
|
||||||
"""Trade the API key for the account email so events join other Mem0 surfaces."""
|
"""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/"
|
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
|
||||||
@@ -380,16 +486,40 @@ def resolve_distinct_id() -> tuple[str, str]:
|
|||||||
|
|
||||||
|
|
||||||
def flush() -> int:
|
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():
|
if not is_enabled():
|
||||||
return 0
|
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()
|
||||||
|
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:
|
if claim is None:
|
||||||
return 0
|
return 0, True
|
||||||
try:
|
try:
|
||||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||||
except OSError:
|
except OSError:
|
||||||
return 0
|
return 0, True
|
||||||
events = []
|
events = []
|
||||||
for line in lines:
|
for line in lines:
|
||||||
try:
|
try:
|
||||||
@@ -403,7 +533,7 @@ def flush() -> int:
|
|||||||
claim.unlink()
|
claim.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
return 0
|
return 0, True
|
||||||
|
|
||||||
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
||||||
if aliased_anonymous_id:
|
if aliased_anonymous_id:
|
||||||
@@ -422,10 +552,13 @@ def flush() -> int:
|
|||||||
|
|
||||||
sent = 0
|
sent = 0
|
||||||
for start in range(0, len(events), BATCH_SIZE):
|
for start in range(0, len(events), BATCH_SIZE):
|
||||||
|
chunk = events[start : start + BATCH_SIZE]
|
||||||
batch = [
|
batch = [
|
||||||
{
|
{
|
||||||
"event": event["event"],
|
"event": event["event"],
|
||||||
"distinct_id": distinct_id,
|
"distinct_id": distinct_id,
|
||||||
|
# Carried through from record() so a resend can be collapsed.
|
||||||
|
"uuid": event.get("uuid"),
|
||||||
"timestamp": event.get("timestamp"),
|
"timestamp": event.get("timestamp"),
|
||||||
"properties": {
|
"properties": {
|
||||||
# Fallback only: events recorded by a build before source
|
# Fallback only: events recorded by a build before source
|
||||||
@@ -437,16 +570,19 @@ def flush() -> int:
|
|||||||
**(event.get("properties") or {}),
|
**(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):
|
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
|
||||||
return sent
|
# Keep only what has not been delivered, and release the lease.
|
||||||
sent += len(batch)
|
# Previously the whole file was kept and the retry re-posted every
|
||||||
try:
|
# batch, including the ones that had already arrived.
|
||||||
claim.unlink()
|
_release_claim(claim, events[start:])
|
||||||
except OSError:
|
return sent, False
|
||||||
pass
|
sent += len(chunk)
|
||||||
return sent
|
# Record progress and refresh the lease after each successful batch, so
|
||||||
|
# a crash repeats at most one batch instead of the entire file.
|
||||||
|
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||||
|
return sent, True
|
||||||
|
|
||||||
|
|
||||||
def main() -> int:
|
def main() -> int:
|
||||||
|
|||||||
@@ -99,6 +99,11 @@ BATCH_SIZE = 100
|
|||||||
SEND_TIMEOUT = 5
|
SEND_TIMEOUT = 5
|
||||||
CLAIM_STALE_SECONDS = 120
|
CLAIM_STALE_SECONDS = 120
|
||||||
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
|
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
|
||||||
|
|
||||||
|
|
||||||
def is_enabled() -> bool:
|
def is_enabled() -> bool:
|
||||||
@@ -301,38 +306,139 @@ def spawn_flush() -> bool:
|
|||||||
return False
|
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."""
|
||||||
|
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||||
|
tail = stem.rsplit("-", 1)[-1]
|
||||||
|
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:
|
def _claim_spool() -> Path | None:
|
||||||
"""Rename the spool aside so exactly one sender owns each batch."""
|
"""Rename the spool aside so exactly one sender owns each batch."""
|
||||||
directory = memory_core.data_dir()
|
directory = memory_core.data_dir()
|
||||||
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
|
claim = directory / _claim_name()
|
||||||
spool = _spool_path()
|
spool = _spool_path()
|
||||||
try:
|
try:
|
||||||
spool.replace(claim)
|
spool.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
|
return _claim_parked(directory)
|
||||||
|
|
||||||
|
|
||||||
|
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()
|
now = time.time()
|
||||||
for orphan in sorted(directory.glob("telemetry-*.sending")):
|
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
|
||||||
try:
|
try:
|
||||||
age = now - orphan.stat().st_mtime
|
age = now - orphan.stat().st_mtime
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
if age > CLAIM_EXPIRY_SECONDS:
|
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||||
try:
|
try:
|
||||||
orphan.unlink()
|
orphan.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
continue
|
continue
|
||||||
if age < CLAIM_STALE_SECONDS:
|
if age < CLAIM_STALE_SECONDS:
|
||||||
|
# Someone else holds a live lease on it.
|
||||||
continue
|
continue
|
||||||
|
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
|
||||||
try:
|
try:
|
||||||
orphan.replace(claim)
|
orphan.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
return None
|
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:
|
||||||
|
temporary.write_text(
|
||||||
|
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||||
|
encoding="utf-8",
|
||||||
|
)
|
||||||
|
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:
|
||||||
|
released = time.time() - CLAIM_STALE_SECONDS - 1
|
||||||
|
os.utime(claim, (released, released))
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
def _resolve_email(key: str) -> str:
|
def _resolve_email(key: str) -> str:
|
||||||
"""Trade the API key for the account email so events join other Mem0 surfaces."""
|
"""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/"
|
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
|
||||||
@@ -380,16 +486,40 @@ def resolve_distinct_id() -> tuple[str, str]:
|
|||||||
|
|
||||||
|
|
||||||
def flush() -> int:
|
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():
|
if not is_enabled():
|
||||||
return 0
|
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()
|
||||||
|
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:
|
if claim is None:
|
||||||
return 0
|
return 0, True
|
||||||
try:
|
try:
|
||||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||||
except OSError:
|
except OSError:
|
||||||
return 0
|
return 0, True
|
||||||
events = []
|
events = []
|
||||||
for line in lines:
|
for line in lines:
|
||||||
try:
|
try:
|
||||||
@@ -403,7 +533,7 @@ def flush() -> int:
|
|||||||
claim.unlink()
|
claim.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
return 0
|
return 0, True
|
||||||
|
|
||||||
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
||||||
if aliased_anonymous_id:
|
if aliased_anonymous_id:
|
||||||
@@ -422,10 +552,13 @@ def flush() -> int:
|
|||||||
|
|
||||||
sent = 0
|
sent = 0
|
||||||
for start in range(0, len(events), BATCH_SIZE):
|
for start in range(0, len(events), BATCH_SIZE):
|
||||||
|
chunk = events[start : start + BATCH_SIZE]
|
||||||
batch = [
|
batch = [
|
||||||
{
|
{
|
||||||
"event": event["event"],
|
"event": event["event"],
|
||||||
"distinct_id": distinct_id,
|
"distinct_id": distinct_id,
|
||||||
|
# Carried through from record() so a resend can be collapsed.
|
||||||
|
"uuid": event.get("uuid"),
|
||||||
"timestamp": event.get("timestamp"),
|
"timestamp": event.get("timestamp"),
|
||||||
"properties": {
|
"properties": {
|
||||||
# Fallback only: events recorded by a build before source
|
# Fallback only: events recorded by a build before source
|
||||||
@@ -437,16 +570,19 @@ def flush() -> int:
|
|||||||
**(event.get("properties") or {}),
|
**(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):
|
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
|
||||||
return sent
|
# Keep only what has not been delivered, and release the lease.
|
||||||
sent += len(batch)
|
# Previously the whole file was kept and the retry re-posted every
|
||||||
try:
|
# batch, including the ones that had already arrived.
|
||||||
claim.unlink()
|
_release_claim(claim, events[start:])
|
||||||
except OSError:
|
return sent, False
|
||||||
pass
|
sent += len(chunk)
|
||||||
return sent
|
# Record progress and refresh the lease after each successful batch, so
|
||||||
|
# a crash repeats at most one batch instead of the entire file.
|
||||||
|
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||||
|
return sent, True
|
||||||
|
|
||||||
|
|
||||||
def main() -> int:
|
def main() -> int:
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import json
|
import json
|
||||||
|
import os
|
||||||
import sys
|
import sys
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from unittest.mock import patch
|
from unittest.mock import patch
|
||||||
@@ -198,9 +199,11 @@ def test_a_stale_claim_is_reclaimed(isolated_env, monkeypatch):
|
|||||||
telemetry.record("search")
|
telemetry.record("search")
|
||||||
orphan = telemetry._claim_spool()
|
orphan = telemetry._claim_spool()
|
||||||
assert orphan is not None
|
assert orphan is not None
|
||||||
monkeypatch.setattr(
|
# Frozen rather than re-stat'd per call: flush() drains the live spool and
|
||||||
telemetry.time, "time", lambda: orphan.stat().st_mtime + telemetry.CLAIM_STALE_SECONDS + 1
|
# 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):
|
with patch.object(telemetry, "_post", lambda payload, url: True):
|
||||||
assert telemetry.flush() == 1
|
assert telemetry.flush() == 1
|
||||||
@@ -210,9 +213,13 @@ def test_an_expired_claim_is_dropped(isolated_env, monkeypatch):
|
|||||||
telemetry.record("search")
|
telemetry.record("search")
|
||||||
orphan = telemetry._claim_spool()
|
orphan = telemetry._claim_spool()
|
||||||
assert orphan is not None
|
assert orphan is not None
|
||||||
monkeypatch.setattr(
|
expired_now = orphan.stat().st_mtime + telemetry.CLAIM_EXPIRY_SECONDS + 1
|
||||||
telemetry.time, "time", lambda: 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 telemetry._claim_spool() is None
|
||||||
assert not list(memory_core.data_dir().glob("telemetry-*.sending"))
|
assert not list(memory_core.data_dir().glob("telemetry-*.sending"))
|
||||||
|
|
||||||
|
|||||||
@@ -99,6 +99,11 @@ BATCH_SIZE = 100
|
|||||||
SEND_TIMEOUT = 5
|
SEND_TIMEOUT = 5
|
||||||
CLAIM_STALE_SECONDS = 120
|
CLAIM_STALE_SECONDS = 120
|
||||||
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
|
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
|
||||||
|
|
||||||
|
|
||||||
def is_enabled() -> bool:
|
def is_enabled() -> bool:
|
||||||
@@ -301,38 +306,139 @@ def spawn_flush() -> bool:
|
|||||||
return False
|
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."""
|
||||||
|
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||||
|
tail = stem.rsplit("-", 1)[-1]
|
||||||
|
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:
|
def _claim_spool() -> Path | None:
|
||||||
"""Rename the spool aside so exactly one sender owns each batch."""
|
"""Rename the spool aside so exactly one sender owns each batch."""
|
||||||
directory = memory_core.data_dir()
|
directory = memory_core.data_dir()
|
||||||
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
|
claim = directory / _claim_name()
|
||||||
spool = _spool_path()
|
spool = _spool_path()
|
||||||
try:
|
try:
|
||||||
spool.replace(claim)
|
spool.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
|
return _claim_parked(directory)
|
||||||
|
|
||||||
|
|
||||||
|
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()
|
now = time.time()
|
||||||
for orphan in sorted(directory.glob("telemetry-*.sending")):
|
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
|
||||||
try:
|
try:
|
||||||
age = now - orphan.stat().st_mtime
|
age = now - orphan.stat().st_mtime
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
if age > CLAIM_EXPIRY_SECONDS:
|
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||||
try:
|
try:
|
||||||
orphan.unlink()
|
orphan.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
continue
|
continue
|
||||||
if age < CLAIM_STALE_SECONDS:
|
if age < CLAIM_STALE_SECONDS:
|
||||||
|
# Someone else holds a live lease on it.
|
||||||
continue
|
continue
|
||||||
|
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
|
||||||
try:
|
try:
|
||||||
orphan.replace(claim)
|
orphan.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
return None
|
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:
|
||||||
|
temporary.write_text(
|
||||||
|
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||||
|
encoding="utf-8",
|
||||||
|
)
|
||||||
|
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:
|
||||||
|
released = time.time() - CLAIM_STALE_SECONDS - 1
|
||||||
|
os.utime(claim, (released, released))
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
def _resolve_email(key: str) -> str:
|
def _resolve_email(key: str) -> str:
|
||||||
"""Trade the API key for the account email so events join other Mem0 surfaces."""
|
"""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/"
|
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
|
||||||
@@ -380,16 +486,40 @@ def resolve_distinct_id() -> tuple[str, str]:
|
|||||||
|
|
||||||
|
|
||||||
def flush() -> int:
|
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():
|
if not is_enabled():
|
||||||
return 0
|
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()
|
||||||
|
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:
|
if claim is None:
|
||||||
return 0
|
return 0, True
|
||||||
try:
|
try:
|
||||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||||
except OSError:
|
except OSError:
|
||||||
return 0
|
return 0, True
|
||||||
events = []
|
events = []
|
||||||
for line in lines:
|
for line in lines:
|
||||||
try:
|
try:
|
||||||
@@ -403,7 +533,7 @@ def flush() -> int:
|
|||||||
claim.unlink()
|
claim.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
return 0
|
return 0, True
|
||||||
|
|
||||||
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
||||||
if aliased_anonymous_id:
|
if aliased_anonymous_id:
|
||||||
@@ -422,10 +552,13 @@ def flush() -> int:
|
|||||||
|
|
||||||
sent = 0
|
sent = 0
|
||||||
for start in range(0, len(events), BATCH_SIZE):
|
for start in range(0, len(events), BATCH_SIZE):
|
||||||
|
chunk = events[start : start + BATCH_SIZE]
|
||||||
batch = [
|
batch = [
|
||||||
{
|
{
|
||||||
"event": event["event"],
|
"event": event["event"],
|
||||||
"distinct_id": distinct_id,
|
"distinct_id": distinct_id,
|
||||||
|
# Carried through from record() so a resend can be collapsed.
|
||||||
|
"uuid": event.get("uuid"),
|
||||||
"timestamp": event.get("timestamp"),
|
"timestamp": event.get("timestamp"),
|
||||||
"properties": {
|
"properties": {
|
||||||
# Fallback only: events recorded by a build before source
|
# Fallback only: events recorded by a build before source
|
||||||
@@ -437,16 +570,19 @@ def flush() -> int:
|
|||||||
**(event.get("properties") or {}),
|
**(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):
|
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
|
||||||
return sent
|
# Keep only what has not been delivered, and release the lease.
|
||||||
sent += len(batch)
|
# Previously the whole file was kept and the retry re-posted every
|
||||||
try:
|
# batch, including the ones that had already arrived.
|
||||||
claim.unlink()
|
_release_claim(claim, events[start:])
|
||||||
except OSError:
|
return sent, False
|
||||||
pass
|
sent += len(chunk)
|
||||||
return sent
|
# Record progress and refresh the lease after each successful batch, so
|
||||||
|
# a crash repeats at most one batch instead of the entire file.
|
||||||
|
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||||
|
return sent, True
|
||||||
|
|
||||||
|
|
||||||
def main() -> int:
|
def main() -> int:
|
||||||
|
|||||||
@@ -99,6 +99,11 @@ BATCH_SIZE = 100
|
|||||||
SEND_TIMEOUT = 5
|
SEND_TIMEOUT = 5
|
||||||
CLAIM_STALE_SECONDS = 120
|
CLAIM_STALE_SECONDS = 120
|
||||||
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
|
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
|
||||||
|
|
||||||
|
|
||||||
def is_enabled() -> bool:
|
def is_enabled() -> bool:
|
||||||
@@ -301,38 +306,139 @@ def spawn_flush() -> bool:
|
|||||||
return False
|
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."""
|
||||||
|
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||||
|
tail = stem.rsplit("-", 1)[-1]
|
||||||
|
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:
|
def _claim_spool() -> Path | None:
|
||||||
"""Rename the spool aside so exactly one sender owns each batch."""
|
"""Rename the spool aside so exactly one sender owns each batch."""
|
||||||
directory = memory_core.data_dir()
|
directory = memory_core.data_dir()
|
||||||
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
|
claim = directory / _claim_name()
|
||||||
spool = _spool_path()
|
spool = _spool_path()
|
||||||
try:
|
try:
|
||||||
spool.replace(claim)
|
spool.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
|
return _claim_parked(directory)
|
||||||
|
|
||||||
|
|
||||||
|
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()
|
now = time.time()
|
||||||
for orphan in sorted(directory.glob("telemetry-*.sending")):
|
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
|
||||||
try:
|
try:
|
||||||
age = now - orphan.stat().st_mtime
|
age = now - orphan.stat().st_mtime
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
if age > CLAIM_EXPIRY_SECONDS:
|
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||||
try:
|
try:
|
||||||
orphan.unlink()
|
orphan.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
continue
|
continue
|
||||||
if age < CLAIM_STALE_SECONDS:
|
if age < CLAIM_STALE_SECONDS:
|
||||||
|
# Someone else holds a live lease on it.
|
||||||
continue
|
continue
|
||||||
|
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
|
||||||
try:
|
try:
|
||||||
orphan.replace(claim)
|
orphan.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
return None
|
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:
|
||||||
|
temporary.write_text(
|
||||||
|
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||||
|
encoding="utf-8",
|
||||||
|
)
|
||||||
|
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:
|
||||||
|
released = time.time() - CLAIM_STALE_SECONDS - 1
|
||||||
|
os.utime(claim, (released, released))
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
def _resolve_email(key: str) -> str:
|
def _resolve_email(key: str) -> str:
|
||||||
"""Trade the API key for the account email so events join other Mem0 surfaces."""
|
"""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/"
|
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
|
||||||
@@ -380,16 +486,40 @@ def resolve_distinct_id() -> tuple[str, str]:
|
|||||||
|
|
||||||
|
|
||||||
def flush() -> int:
|
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():
|
if not is_enabled():
|
||||||
return 0
|
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()
|
||||||
|
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:
|
if claim is None:
|
||||||
return 0
|
return 0, True
|
||||||
try:
|
try:
|
||||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||||
except OSError:
|
except OSError:
|
||||||
return 0
|
return 0, True
|
||||||
events = []
|
events = []
|
||||||
for line in lines:
|
for line in lines:
|
||||||
try:
|
try:
|
||||||
@@ -403,7 +533,7 @@ def flush() -> int:
|
|||||||
claim.unlink()
|
claim.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
return 0
|
return 0, True
|
||||||
|
|
||||||
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
||||||
if aliased_anonymous_id:
|
if aliased_anonymous_id:
|
||||||
@@ -422,10 +552,13 @@ def flush() -> int:
|
|||||||
|
|
||||||
sent = 0
|
sent = 0
|
||||||
for start in range(0, len(events), BATCH_SIZE):
|
for start in range(0, len(events), BATCH_SIZE):
|
||||||
|
chunk = events[start : start + BATCH_SIZE]
|
||||||
batch = [
|
batch = [
|
||||||
{
|
{
|
||||||
"event": event["event"],
|
"event": event["event"],
|
||||||
"distinct_id": distinct_id,
|
"distinct_id": distinct_id,
|
||||||
|
# Carried through from record() so a resend can be collapsed.
|
||||||
|
"uuid": event.get("uuid"),
|
||||||
"timestamp": event.get("timestamp"),
|
"timestamp": event.get("timestamp"),
|
||||||
"properties": {
|
"properties": {
|
||||||
# Fallback only: events recorded by a build before source
|
# Fallback only: events recorded by a build before source
|
||||||
@@ -437,16 +570,19 @@ def flush() -> int:
|
|||||||
**(event.get("properties") or {}),
|
**(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):
|
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
|
||||||
return sent
|
# Keep only what has not been delivered, and release the lease.
|
||||||
sent += len(batch)
|
# Previously the whole file was kept and the retry re-posted every
|
||||||
try:
|
# batch, including the ones that had already arrived.
|
||||||
claim.unlink()
|
_release_claim(claim, events[start:])
|
||||||
except OSError:
|
return sent, False
|
||||||
pass
|
sent += len(chunk)
|
||||||
return sent
|
# Record progress and refresh the lease after each successful batch, so
|
||||||
|
# a crash repeats at most one batch instead of the entire file.
|
||||||
|
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||||
|
return sent, True
|
||||||
|
|
||||||
|
|
||||||
def main() -> int:
|
def main() -> int:
|
||||||
|
|||||||
@@ -99,6 +99,11 @@ BATCH_SIZE = 100
|
|||||||
SEND_TIMEOUT = 5
|
SEND_TIMEOUT = 5
|
||||||
CLAIM_STALE_SECONDS = 120
|
CLAIM_STALE_SECONDS = 120
|
||||||
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
|
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
|
||||||
|
|
||||||
|
|
||||||
def is_enabled() -> bool:
|
def is_enabled() -> bool:
|
||||||
@@ -301,38 +306,139 @@ def spawn_flush() -> bool:
|
|||||||
return False
|
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."""
|
||||||
|
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||||
|
tail = stem.rsplit("-", 1)[-1]
|
||||||
|
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:
|
def _claim_spool() -> Path | None:
|
||||||
"""Rename the spool aside so exactly one sender owns each batch."""
|
"""Rename the spool aside so exactly one sender owns each batch."""
|
||||||
directory = memory_core.data_dir()
|
directory = memory_core.data_dir()
|
||||||
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
|
claim = directory / _claim_name()
|
||||||
spool = _spool_path()
|
spool = _spool_path()
|
||||||
try:
|
try:
|
||||||
spool.replace(claim)
|
spool.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
|
return _claim_parked(directory)
|
||||||
|
|
||||||
|
|
||||||
|
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()
|
now = time.time()
|
||||||
for orphan in sorted(directory.glob("telemetry-*.sending")):
|
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
|
||||||
try:
|
try:
|
||||||
age = now - orphan.stat().st_mtime
|
age = now - orphan.stat().st_mtime
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
if age > CLAIM_EXPIRY_SECONDS:
|
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||||
try:
|
try:
|
||||||
orphan.unlink()
|
orphan.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
continue
|
continue
|
||||||
if age < CLAIM_STALE_SECONDS:
|
if age < CLAIM_STALE_SECONDS:
|
||||||
|
# Someone else holds a live lease on it.
|
||||||
continue
|
continue
|
||||||
|
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
|
||||||
try:
|
try:
|
||||||
orphan.replace(claim)
|
orphan.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
return None
|
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:
|
||||||
|
temporary.write_text(
|
||||||
|
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||||
|
encoding="utf-8",
|
||||||
|
)
|
||||||
|
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:
|
||||||
|
released = time.time() - CLAIM_STALE_SECONDS - 1
|
||||||
|
os.utime(claim, (released, released))
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
def _resolve_email(key: str) -> str:
|
def _resolve_email(key: str) -> str:
|
||||||
"""Trade the API key for the account email so events join other Mem0 surfaces."""
|
"""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/"
|
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
|
||||||
@@ -380,16 +486,40 @@ def resolve_distinct_id() -> tuple[str, str]:
|
|||||||
|
|
||||||
|
|
||||||
def flush() -> int:
|
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():
|
if not is_enabled():
|
||||||
return 0
|
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()
|
||||||
|
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:
|
if claim is None:
|
||||||
return 0
|
return 0, True
|
||||||
try:
|
try:
|
||||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||||
except OSError:
|
except OSError:
|
||||||
return 0
|
return 0, True
|
||||||
events = []
|
events = []
|
||||||
for line in lines:
|
for line in lines:
|
||||||
try:
|
try:
|
||||||
@@ -403,7 +533,7 @@ def flush() -> int:
|
|||||||
claim.unlink()
|
claim.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
return 0
|
return 0, True
|
||||||
|
|
||||||
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
||||||
if aliased_anonymous_id:
|
if aliased_anonymous_id:
|
||||||
@@ -422,10 +552,13 @@ def flush() -> int:
|
|||||||
|
|
||||||
sent = 0
|
sent = 0
|
||||||
for start in range(0, len(events), BATCH_SIZE):
|
for start in range(0, len(events), BATCH_SIZE):
|
||||||
|
chunk = events[start : start + BATCH_SIZE]
|
||||||
batch = [
|
batch = [
|
||||||
{
|
{
|
||||||
"event": event["event"],
|
"event": event["event"],
|
||||||
"distinct_id": distinct_id,
|
"distinct_id": distinct_id,
|
||||||
|
# Carried through from record() so a resend can be collapsed.
|
||||||
|
"uuid": event.get("uuid"),
|
||||||
"timestamp": event.get("timestamp"),
|
"timestamp": event.get("timestamp"),
|
||||||
"properties": {
|
"properties": {
|
||||||
# Fallback only: events recorded by a build before source
|
# Fallback only: events recorded by a build before source
|
||||||
@@ -437,16 +570,19 @@ def flush() -> int:
|
|||||||
**(event.get("properties") or {}),
|
**(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):
|
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
|
||||||
return sent
|
# Keep only what has not been delivered, and release the lease.
|
||||||
sent += len(batch)
|
# Previously the whole file was kept and the retry re-posted every
|
||||||
try:
|
# batch, including the ones that had already arrived.
|
||||||
claim.unlink()
|
_release_claim(claim, events[start:])
|
||||||
except OSError:
|
return sent, False
|
||||||
pass
|
sent += len(chunk)
|
||||||
return sent
|
# Record progress and refresh the lease after each successful batch, so
|
||||||
|
# a crash repeats at most one batch instead of the entire file.
|
||||||
|
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||||
|
return sent, True
|
||||||
|
|
||||||
|
|
||||||
def main() -> int:
|
def main() -> int:
|
||||||
|
|||||||
@@ -99,6 +99,11 @@ BATCH_SIZE = 100
|
|||||||
SEND_TIMEOUT = 5
|
SEND_TIMEOUT = 5
|
||||||
CLAIM_STALE_SECONDS = 120
|
CLAIM_STALE_SECONDS = 120
|
||||||
CLAIM_EXPIRY_SECONDS = 7 * 24 * 60 * 60
|
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
|
||||||
|
|
||||||
|
|
||||||
def is_enabled() -> bool:
|
def is_enabled() -> bool:
|
||||||
@@ -301,38 +306,139 @@ def spawn_flush() -> bool:
|
|||||||
return False
|
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."""
|
||||||
|
stem = claim.name[: -len(".sending")] if claim.name.endswith(".sending") else claim.name
|
||||||
|
tail = stem.rsplit("-", 1)[-1]
|
||||||
|
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:
|
def _claim_spool() -> Path | None:
|
||||||
"""Rename the spool aside so exactly one sender owns each batch."""
|
"""Rename the spool aside so exactly one sender owns each batch."""
|
||||||
directory = memory_core.data_dir()
|
directory = memory_core.data_dir()
|
||||||
claim = directory / f"telemetry-{os.getpid()}-{uuid.uuid4().hex[:8]}.sending"
|
claim = directory / _claim_name()
|
||||||
spool = _spool_path()
|
spool = _spool_path()
|
||||||
try:
|
try:
|
||||||
spool.replace(claim)
|
spool.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
|
return _claim_parked(directory)
|
||||||
|
|
||||||
|
|
||||||
|
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()
|
now = time.time()
|
||||||
for orphan in sorted(directory.glob("telemetry-*.sending")):
|
for orphan in sorted(directory.glob("telemetry-*.sending"), key=_safe_mtime):
|
||||||
try:
|
try:
|
||||||
age = now - orphan.stat().st_mtime
|
age = now - orphan.stat().st_mtime
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
if age > CLAIM_EXPIRY_SECONDS:
|
if age > CLAIM_EXPIRY_SECONDS and _claim_attempt(orphan) >= MAX_CLAIM_ATTEMPTS:
|
||||||
try:
|
try:
|
||||||
orphan.unlink()
|
orphan.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
continue
|
continue
|
||||||
if age < CLAIM_STALE_SECONDS:
|
if age < CLAIM_STALE_SECONDS:
|
||||||
|
# Someone else holds a live lease on it.
|
||||||
continue
|
continue
|
||||||
|
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
|
||||||
try:
|
try:
|
||||||
orphan.replace(claim)
|
orphan.replace(claim)
|
||||||
|
_touch(claim)
|
||||||
return claim
|
return claim
|
||||||
except OSError:
|
except OSError:
|
||||||
continue
|
continue
|
||||||
return None
|
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:
|
||||||
|
temporary.write_text(
|
||||||
|
"".join(json.dumps(event, separators=(",", ":"), default=str) + "\n" for event in remaining),
|
||||||
|
encoding="utf-8",
|
||||||
|
)
|
||||||
|
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:
|
||||||
|
released = time.time() - CLAIM_STALE_SECONDS - 1
|
||||||
|
os.utime(claim, (released, released))
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
def _resolve_email(key: str) -> str:
|
def _resolve_email(key: str) -> str:
|
||||||
"""Trade the API key for the account email so events join other Mem0 surfaces."""
|
"""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/"
|
url = os.environ.get("MEM0_API_URL", memory_core.DEFAULT_API_URL).rstrip("/") + "/v1/ping/"
|
||||||
@@ -380,16 +486,40 @@ def resolve_distinct_id() -> tuple[str, str]:
|
|||||||
|
|
||||||
|
|
||||||
def flush() -> int:
|
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():
|
if not is_enabled():
|
||||||
return 0
|
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()
|
||||||
|
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:
|
if claim is None:
|
||||||
return 0
|
return 0, True
|
||||||
try:
|
try:
|
||||||
lines = claim.read_text(encoding="utf-8").splitlines()
|
lines = claim.read_text(encoding="utf-8").splitlines()
|
||||||
except OSError:
|
except OSError:
|
||||||
return 0
|
return 0, True
|
||||||
events = []
|
events = []
|
||||||
for line in lines:
|
for line in lines:
|
||||||
try:
|
try:
|
||||||
@@ -403,7 +533,7 @@ def flush() -> int:
|
|||||||
claim.unlink()
|
claim.unlink()
|
||||||
except OSError:
|
except OSError:
|
||||||
pass
|
pass
|
||||||
return 0
|
return 0, True
|
||||||
|
|
||||||
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
distinct_id, aliased_anonymous_id = resolve_distinct_id()
|
||||||
if aliased_anonymous_id:
|
if aliased_anonymous_id:
|
||||||
@@ -422,10 +552,13 @@ def flush() -> int:
|
|||||||
|
|
||||||
sent = 0
|
sent = 0
|
||||||
for start in range(0, len(events), BATCH_SIZE):
|
for start in range(0, len(events), BATCH_SIZE):
|
||||||
|
chunk = events[start : start + BATCH_SIZE]
|
||||||
batch = [
|
batch = [
|
||||||
{
|
{
|
||||||
"event": event["event"],
|
"event": event["event"],
|
||||||
"distinct_id": distinct_id,
|
"distinct_id": distinct_id,
|
||||||
|
# Carried through from record() so a resend can be collapsed.
|
||||||
|
"uuid": event.get("uuid"),
|
||||||
"timestamp": event.get("timestamp"),
|
"timestamp": event.get("timestamp"),
|
||||||
"properties": {
|
"properties": {
|
||||||
# Fallback only: events recorded by a build before source
|
# Fallback only: events recorded by a build before source
|
||||||
@@ -437,16 +570,19 @@ def flush() -> int:
|
|||||||
**(event.get("properties") or {}),
|
**(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):
|
if not _post({"api_key": POSTHOG_API_KEY, "batch": batch}, POSTHOG_BATCH_URL):
|
||||||
return sent
|
# Keep only what has not been delivered, and release the lease.
|
||||||
sent += len(batch)
|
# Previously the whole file was kept and the retry re-posted every
|
||||||
try:
|
# batch, including the ones that had already arrived.
|
||||||
claim.unlink()
|
_release_claim(claim, events[start:])
|
||||||
except OSError:
|
return sent, False
|
||||||
pass
|
sent += len(chunk)
|
||||||
return sent
|
# Record progress and refresh the lease after each successful batch, so
|
||||||
|
# a crash repeats at most one batch instead of the entire file.
|
||||||
|
_rewrite_claim(claim, events[start + len(chunk) :])
|
||||||
|
return sent, True
|
||||||
|
|
||||||
|
|
||||||
def main() -> int:
|
def main() -> int:
|
||||||
|
|||||||
Reference in New Issue
Block a user