fix(telemetry): fence send authority on a per-claim token
Responds to the independent PR review (andrexibiza). Both P1s were
checked against current HEAD rather than taken on authority - the
review was written against 613849c190, before the interval-model
consent replacement landed.
P1-1 (same-UTC-day revoke/re-enable releases refused data): already
fixed by the interval model. The reviewer's exact reproduction - opt in
06:00, revoke 12:00, package collected 18:00, re-enable 20:00 same day
- was re-run at HEAD: the off-window package stays local, and a full-day
aggregate straddling the revocation boundary also stays local (period
containment, timestamp precision). The consent-windows harness already
pins both. The reviewer's related ask that consent-ledger persistence
failures fail closed also holds structurally now: reconciliation derives
state rather than recording transitions, so a lost write means a shorter
confirmed horizon - less is released, never more.
P1-2 (lease has no owner) was VALID at head. Reproduced exactly as
described: A claims, is suspended past the 300s lease, B reclaims and
POSTs, A resumes and POSTs again - and the ingest key is minute-
prefixed, so the duplicate lands as a DISTINCT stored object, making
this worse than a benign idempotent overwrite.
Fix: every claim now mints a claim_token (additive nullable column,
schema version unchanged). Ownership is revalidated immediately before
every external POST, and every settlement, rejection, and backoff write
is compare-and-set on (package_id, claim_token, pending). A lapsed
claimant that resumes yields without transmitting, and its stale
backoff cannot move next_attempt_at under the live claim's lease.
Two deterministic regressions ship with it: expiry -> reclaim -> resume
(the reviewer's schedule), and the subtler stale-backoff-clobber case.
Honest scope, documented on _send_one: delivery remains at-least-once.
The token closes the claim->POST gap; a suspension landing mid-POST
(bytes already on the wire) is not client-revocable. The residual
duplicate is byte-identical content; collapsing it fully needs
package_id-keyed dedupe at the ingest service.
275 tests pass; ruff + footguns clean; staging E2E 202.
This commit is contained in:
@@ -378,6 +378,12 @@ class SharedMetricsStore:
|
||||
# Only the ~36-byte id is stored: the body is recomputed from
|
||||
# payload_json, whose serialisation is deterministic.
|
||||
("sent_install_id", "TEXT"),
|
||||
# NULL until first claimed; rewritten on every claim. Settlement
|
||||
# and the pre-POST revalidation are compare-and-set on this, so a
|
||||
# claimant whose lease lapsed loses authority the moment another
|
||||
# process reclaims (PR-review finding: without it, a suspended
|
||||
# sender resuming after a reclaim double-POSTs the package).
|
||||
("claim_token", "TEXT"),
|
||||
):
|
||||
if column not in existing:
|
||||
connection.execute(
|
||||
|
||||
@@ -33,6 +33,7 @@ import sqlite3
|
||||
import time
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
import uuid
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
@@ -396,24 +397,31 @@ class SharedMetricsSender:
|
||||
# caller to continue rather than stop.
|
||||
return {"package_id": package_id, "skip": True}
|
||||
|
||||
token = str(uuid.uuid4())
|
||||
connection.execute(
|
||||
"""
|
||||
UPDATE package_outbox
|
||||
SET send_state = 'pending',
|
||||
send_attempts = send_attempts + 1,
|
||||
next_attempt_at = ?
|
||||
next_attempt_at = ?,
|
||||
claim_token = ?
|
||||
WHERE package_id = ?
|
||||
""",
|
||||
# Lease INTO THE FUTURE: selection requires
|
||||
# next_attempt_at <= now, so no other process can take
|
||||
# this row while it is in flight. Success or a real
|
||||
# backoff overwrites it; if this process dies, it expires.
|
||||
(_isoformat(lease_until), package_id),
|
||||
# The token is this claim's identity: a reclaim after
|
||||
# expiry mints a new one, and every later write by THIS
|
||||
# claimant is compare-and-set against it, so a lapsed
|
||||
# claimant that resumes cannot settle or transmit.
|
||||
(_isoformat(lease_until), token, package_id),
|
||||
)
|
||||
return {
|
||||
"package_id": package_id,
|
||||
"payload_json": str(row[1]),
|
||||
"derived": str(derived),
|
||||
"claim_token": token,
|
||||
"skip": False,
|
||||
}
|
||||
|
||||
@@ -479,12 +487,25 @@ class SharedMetricsSender:
|
||||
payload = substitute_install_id(json.loads(payload_json), derived)
|
||||
return json.dumps(payload, indent=2, sort_keys=True).encode("utf-8")
|
||||
|
||||
def _mark(self, package_id: str, *, only_if_pending: bool = True, **columns) -> None:
|
||||
def _mark(
|
||||
self,
|
||||
package_id: str,
|
||||
*,
|
||||
only_if_pending: bool = True,
|
||||
token: str | None = None,
|
||||
**columns,
|
||||
) -> None:
|
||||
"""Write send state for one package.
|
||||
|
||||
Guarded on send_state so a pass whose lease lapsed cannot resurrect a
|
||||
row another process has already finished: without this, a slow sender
|
||||
could overwrite 'sent' back to 'pending' and cause a re-send.
|
||||
|
||||
When ``token`` is given, the write is additionally compare-and-set on
|
||||
claim_token: it lands only if THIS claim is still the current one. A
|
||||
claimant that lapsed and was superseded writes zero rows — its
|
||||
settlement, backoff, and error strings all silently lose to the
|
||||
newer claim's, which is the correct outcome.
|
||||
"""
|
||||
assignments = ", ".join(f"{name} = ?" for name in columns)
|
||||
predicate = (
|
||||
@@ -492,15 +513,48 @@ class SharedMetricsSender:
|
||||
if only_if_pending
|
||||
else ""
|
||||
)
|
||||
params: list = [*columns.values(), package_id]
|
||||
if token is not None:
|
||||
predicate += " AND claim_token = ?"
|
||||
params.append(token)
|
||||
with self._store._connection() as connection:
|
||||
with write_txn(connection):
|
||||
connection.execute(
|
||||
f"UPDATE package_outbox SET {assignments} "
|
||||
f"WHERE package_id = ?{predicate}",
|
||||
(*columns.values(), package_id),
|
||||
params,
|
||||
)
|
||||
|
||||
def _defer(self, package_id: str, delay_seconds: int, reason: str) -> None:
|
||||
def _still_owns(self, package_id: str, token: str | None) -> bool:
|
||||
"""Return whether this pass's claim on the row is still current."""
|
||||
if token is None:
|
||||
# Defensive: a package dict without a token (not produced by
|
||||
# _claim_next today) gets no authority rather than unlimited.
|
||||
return False
|
||||
try:
|
||||
with self._store._connection() as connection:
|
||||
row = connection.execute(
|
||||
"SELECT 1 FROM package_outbox"
|
||||
" WHERE package_id = ? AND claim_token = ?"
|
||||
" AND (send_state IS NULL OR send_state = 'pending')",
|
||||
(package_id, token),
|
||||
).fetchone()
|
||||
return row is not None
|
||||
except Exception:
|
||||
# If the check itself fails, do not transmit on stale authority.
|
||||
logger.warning(
|
||||
"Unable to verify shared-metrics claim ownership", exc_info=True
|
||||
)
|
||||
return False
|
||||
|
||||
def _defer(
|
||||
self,
|
||||
package_id: str,
|
||||
delay_seconds: int,
|
||||
reason: str,
|
||||
*,
|
||||
token: str | None = None,
|
||||
) -> None:
|
||||
# Defence in depth: no current caller can pass a non-positive delay
|
||||
# (Retry-After is already clamped to [1, 86400] when parsed, and every
|
||||
# other call site passes a positive constant), so this clamp is
|
||||
@@ -512,6 +566,7 @@ class SharedMetricsSender:
|
||||
retry_at = self._now().timestamp() + delay
|
||||
self._mark(
|
||||
package_id,
|
||||
token=token,
|
||||
send_state="pending",
|
||||
next_attempt_at=_isoformat(
|
||||
datetime.fromtimestamp(retry_at, tz=timezone.utc)
|
||||
@@ -520,11 +575,33 @@ class SharedMetricsSender:
|
||||
)
|
||||
|
||||
def _send_one(self, package: dict) -> str:
|
||||
"""Try one package. Returns 'sent', 'rejected', or 'deferred'."""
|
||||
"""Try one package. Returns 'sent', 'rejected', or 'deferred'.
|
||||
|
||||
Delivery is at-least-once. The pre-POST ownership check plus the
|
||||
token-fenced writes close the claim->POST and settle-after-reclaim
|
||||
gaps, but a suspension landing MID-POST (bytes already on the wire
|
||||
when the machine sleeps) can still duplicate: no client-side check
|
||||
can revoke a request in flight. The body is byte-identical across
|
||||
retries by construction, so the residual duplicate is exactly one
|
||||
redundant copy of identical content; collapsing it fully would need
|
||||
package_id-keyed dedupe at the ingest service.
|
||||
"""
|
||||
package_id = package["package_id"]
|
||||
token = package.get("claim_token")
|
||||
body = self._body(package["payload_json"], package["derived"])
|
||||
|
||||
for attempt in range(1, self._max_attempts + 1):
|
||||
# Revalidate ownership immediately before the external POST. The
|
||||
# claim can lapse between claiming and here — a suspended laptop,
|
||||
# a GC pause, a long gzip — and another process may have
|
||||
# reclaimed and transmitted. Without this check the resumed
|
||||
# claimant POSTs a duplicate; the ingest key is minute-prefixed,
|
||||
# so duplicates become distinct stored objects, not overwrites.
|
||||
if not self._still_owns(package_id, token):
|
||||
logger.info(
|
||||
"Shared-metrics claim on %s superseded; yielding", package_id
|
||||
)
|
||||
return "deferred"
|
||||
try:
|
||||
response = self._post(
|
||||
self._endpoint, body, timeout=REQUEST_TIMEOUT_SECONDS
|
||||
@@ -532,7 +609,9 @@ class SharedMetricsSender:
|
||||
except Exception as exc: # transport failure: offline, DNS, TLS
|
||||
reason = f"{type(exc).__name__}: {exc}"
|
||||
if attempt >= self._max_attempts:
|
||||
self._defer(package_id, _FAILURE_BACKOFF_SECONDS, reason)
|
||||
self._defer(
|
||||
package_id, _FAILURE_BACKOFF_SECONDS, reason, token=token
|
||||
)
|
||||
return "deferred"
|
||||
self._sleep(self._backoff(attempt))
|
||||
continue
|
||||
@@ -540,6 +619,7 @@ class SharedMetricsSender:
|
||||
if response.status == 202:
|
||||
self._mark(
|
||||
package_id,
|
||||
token=token,
|
||||
send_state="sent",
|
||||
sent_at=_isoformat(self._now()),
|
||||
last_error=None,
|
||||
@@ -560,6 +640,7 @@ class SharedMetricsSender:
|
||||
)
|
||||
self._mark(
|
||||
package_id,
|
||||
token=token,
|
||||
send_state="rejected",
|
||||
last_error=f"HTTP {response.status}: {response.body[:400]}",
|
||||
)
|
||||
@@ -570,17 +651,22 @@ class SharedMetricsSender:
|
||||
package_id,
|
||||
_retry_after_seconds(response.retry_after, _FAILURE_BACKOFF_SECONDS),
|
||||
"rate limited",
|
||||
token=token,
|
||||
)
|
||||
return "deferred"
|
||||
|
||||
# 5xx and anything unexpected: retryable.
|
||||
reason = f"HTTP {response.status}"
|
||||
if attempt >= self._max_attempts:
|
||||
self._defer(package_id, _FAILURE_BACKOFF_SECONDS, reason)
|
||||
self._defer(
|
||||
package_id, _FAILURE_BACKOFF_SECONDS, reason, token=token
|
||||
)
|
||||
return "deferred"
|
||||
self._sleep(self._backoff(attempt))
|
||||
|
||||
self._defer(package_id, _FAILURE_BACKOFF_SECONDS, "attempts exhausted")
|
||||
self._defer(
|
||||
package_id, _FAILURE_BACKOFF_SECONDS, "attempts exhausted", token=token
|
||||
)
|
||||
return "deferred"
|
||||
|
||||
@staticmethod
|
||||
|
||||
@@ -649,6 +649,87 @@ class TestClaimingAndBounds:
|
||||
f"{len(attempts)} requests burned on one doomed package"
|
||||
)
|
||||
|
||||
def test_a_lapsed_claimant_resuming_after_reclaim_cannot_double_post(
|
||||
self, store
|
||||
):
|
||||
"""PR-review P1: expiry -> reclaim -> old claimant resumes.
|
||||
|
||||
A claims, then is suspended (laptop lid) BEFORE its POST. The lease
|
||||
expires; B reclaims and POSTs; A wakes and proceeds. The pre-POST
|
||||
ownership check must make A yield without transmitting.
|
||||
|
||||
Scope note: the check closes the claim->POST gap. A suspension that
|
||||
lands mid-POST (bytes already leaving) is not client-fixable — that
|
||||
residual needs server-side dedupe and is documented on _send_one.
|
||||
"""
|
||||
_add_package(store, "pkg-1", "2026-08-26")
|
||||
|
||||
posts = []
|
||||
|
||||
def post_a(endpoint, payload, *, timeout):
|
||||
posts.append("A")
|
||||
return FakeResponse(202)
|
||||
|
||||
def post_b(endpoint, payload, *, timeout):
|
||||
posts.append("B")
|
||||
return FakeResponse(202)
|
||||
|
||||
sender_a = SharedMetricsSender(
|
||||
store, ENDPOINT, post=post_a, sleep=lambda _s: None, now=lambda: NOW
|
||||
)
|
||||
# A claims, then the process is suspended before _send_one runs.
|
||||
claimed_a = sender_a._claim_next(NOW, set())
|
||||
assert claimed_a is not None and not claimed_a["skip"]
|
||||
|
||||
# 400s later (past the 300s lease) B claims and completes the send.
|
||||
later = NOW + timedelta(seconds=400)
|
||||
sender_b = SharedMetricsSender(
|
||||
store, ENDPOINT, post=post_b, sleep=lambda _s: None, now=lambda: later
|
||||
)
|
||||
outcome_b = sender_b.send_pending()
|
||||
assert outcome_b.sent == 1
|
||||
|
||||
# A resumes exactly where it left off.
|
||||
result_a = sender_a._send_one(claimed_a)
|
||||
|
||||
row = _row(store, "pkg-1")
|
||||
assert posts == ["B"], (
|
||||
f"a lapsed claimant transmitted after reclaim: {posts}"
|
||||
)
|
||||
assert result_a == "deferred"
|
||||
assert row["send_state"] == "sent", "B's settlement must stand"
|
||||
|
||||
def test_a_lapsed_claimants_backoff_cannot_clobber_the_new_claim(self, store):
|
||||
"""The token must fence DEFERS too, not just the 202 settlement.
|
||||
|
||||
A's transport fails after B has reclaimed; A's backoff write must
|
||||
not move next_attempt_at under B's live lease.
|
||||
"""
|
||||
_add_package(store, "pkg-1", "2026-08-26")
|
||||
sender_a = SharedMetricsSender(
|
||||
store, ENDPOINT,
|
||||
post=FakeTransport(OSError("net"), OSError("net"), OSError("net")),
|
||||
sleep=lambda _s: None, now=lambda: NOW,
|
||||
)
|
||||
claimed_a = sender_a._claim_next(NOW, set())
|
||||
assert claimed_a is not None and not claimed_a["skip"]
|
||||
|
||||
later = NOW + timedelta(seconds=400)
|
||||
sender_b = SharedMetricsSender(
|
||||
store, ENDPOINT, post=FakeTransport(),
|
||||
sleep=lambda _s: None, now=lambda: later,
|
||||
)
|
||||
claimed_b = sender_b._claim_next(later, set())
|
||||
assert claimed_b is not None and not claimed_b["skip"]
|
||||
lease_b = _row(store, "pkg-1")["next_attempt_at"]
|
||||
|
||||
# A's exhausted retries try to write a 15-minute backoff.
|
||||
result = sender_a._send_one(claimed_a)
|
||||
assert result == "deferred"
|
||||
assert _row(store, "pkg-1")["next_attempt_at"] == lease_b, (
|
||||
"a lapsed claimant's backoff overwrote the live claim's lease"
|
||||
)
|
||||
|
||||
def test_an_expired_lease_is_reclaimed(self, store):
|
||||
"""A process killed mid-pass must not strand its packages."""
|
||||
_add_package(store, "pkg-1", "2026-08-26")
|
||||
|
||||
Reference in New Issue
Block a user