Merge branch 'pr3/spool-delivery' into pr4/install-marker-and-identity
This commit is contained in:
@@ -145,36 +145,49 @@ def _install_salt() -> str:
|
||||
keys off that file, so recording an event would silently suppress the
|
||||
install event.
|
||||
|
||||
On a read-only or full data directory the fallback is derived from the data
|
||||
directory path: stable for the machine rather than random per call, so the
|
||||
failure mode is a weaker salt and not unbounded cardinality in PostHog.
|
||||
Published atomically, and there is deliberately no derived fallback. Creating
|
||||
the file with O_CREAT|O_EXCL and then writing into it leaves a window where
|
||||
the file exists and is empty, and a concurrent hook that reads it in that
|
||||
window gets nothing. Falling back to a digest of the path would hand that
|
||||
process a salt an attacker can compute, memoized for its whole run, which is
|
||||
the privacy control this function exists to provide silently turning itself
|
||||
off under load. The salt is written to a private temp file first and linked
|
||||
into place, so the name either does not exist or already has the full value.
|
||||
|
||||
Returns "" when it genuinely cannot persist. Callers omit the hash entirely
|
||||
rather than emit an unsalted one.
|
||||
"""
|
||||
global _salt_cache
|
||||
if _salt_cache:
|
||||
return _salt_cache
|
||||
|
||||
path = _salt_path()
|
||||
temporary = path.with_name(f"{path.name}.{os.getpid()}.tmp")
|
||||
try:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
handle = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
|
||||
handle = os.open(temporary, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
|
||||
with os.fdopen(handle, "w", encoding="utf-8") as stream:
|
||||
stream.write(uuid.uuid4().hex)
|
||||
stream.flush()
|
||||
os.fsync(stream.fileno())
|
||||
try:
|
||||
with os.fdopen(handle, "w", encoding="utf-8") as stream:
|
||||
stream.write(uuid.uuid4().hex)
|
||||
# Atomic claim: fails if another process already published one.
|
||||
# os.link rather than replace, which would clobber theirs.
|
||||
os.link(temporary, path)
|
||||
except FileExistsError:
|
||||
pass
|
||||
except OSError:
|
||||
pass
|
||||
finally:
|
||||
try:
|
||||
temporary.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
except FileExistsError:
|
||||
pass
|
||||
except OSError:
|
||||
# Cannot persist. Stable-per-machine beats random-per-call.
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
try:
|
||||
_salt_cache = path.read_text(encoding="utf-8").strip()
|
||||
except OSError:
|
||||
_salt_cache = ""
|
||||
if not _salt_cache:
|
||||
_salt_cache = hashlib.sha256(str(path).encode("utf-8")).hexdigest()
|
||||
return _salt_cache
|
||||
|
||||
|
||||
@@ -187,10 +200,17 @@ def _scoped_digest(value: str, length: int = 16) -> str:
|
||||
privacy control without the salt. Salting per install keeps every
|
||||
within-account join the analytics actually use and gives up only
|
||||
cross-machine joins on the same repository, which nothing computes.
|
||||
|
||||
Returns "" when there is no salt, so record() omits the property. An
|
||||
unsalted digest over this input space is close to plaintext, and emitting one
|
||||
under a name that implies it is hashed is worse than sending nothing.
|
||||
"""
|
||||
if not value:
|
||||
return ""
|
||||
return hashlib.sha256(f"{_install_salt()}:{value}".encode("utf-8")).hexdigest()[:length]
|
||||
salt = _install_salt()
|
||||
if not salt:
|
||||
return ""
|
||||
return hashlib.sha256(f"{salt}:{value}".encode("utf-8")).hexdigest()[:length]
|
||||
|
||||
|
||||
def _safe_value(value: Any) -> Any:
|
||||
@@ -418,10 +438,17 @@ def record(
|
||||
os=sys.platform,
|
||||
python_version=platform.python_version(),
|
||||
)
|
||||
# Assigned only when the digest is real. _scoped_digest returns "" when
|
||||
# the salt could not be persisted, and an empty property is worse than an
|
||||
# absent one: it survives the None filter below and reads as a value.
|
||||
if repo is not None:
|
||||
properties["repo_hash"] = _scoped_digest(getattr(repo, "identity", ""))
|
||||
repo_hash = _scoped_digest(getattr(repo, "identity", ""))
|
||||
if repo_hash:
|
||||
properties["repo_hash"] = repo_hash
|
||||
if session_id:
|
||||
properties["session_hash"] = _scoped_digest(session_id)
|
||||
session_hash = _scoped_digest(session_id)
|
||||
if session_hash:
|
||||
properties["session_hash"] = session_hash
|
||||
line = json.dumps(
|
||||
{
|
||||
"event": f"{EVENT_PREFIX}.{event}",
|
||||
@@ -566,6 +593,14 @@ def _claim_parked(directory: Path) -> Path | None:
|
||||
age = now - orphan.stat().st_mtime
|
||||
except OSError:
|
||||
continue
|
||||
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
|
||||
@@ -576,9 +611,6 @@ def _claim_parked(directory: Path) -> Path | None:
|
||||
except OSError:
|
||||
pass
|
||||
continue
|
||||
if age < CLAIM_STALE_SECONDS:
|
||||
# Someone else holds a live lease on it.
|
||||
continue
|
||||
claim = orphan.parent / _claim_name(_claim_attempt(orphan) + 1)
|
||||
try:
|
||||
orphan.replace(claim)
|
||||
|
||||
@@ -104,6 +104,67 @@ def test_a_fresh_claim_is_not_immediately_stealable(telemetry):
|
||||
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.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user