refactor(hermes_cli): shared_metrics_sender — reuse store time helpers, unify defer paths, txn ctx manager

This commit is contained in:
Teknium
2026-09-02 20:48:41 -07:00
parent 5bd425061c
commit 6837d2e41b

View File

@@ -1,10 +1,8 @@
"""Transmit exported shared-metrics packages to the Nous telemetry service.
Implements the sender side of the ingest contract (see the telemetry repo's ``CONTRACT.md``):
* ``202`` — durably stored. Mark sent. * ``400`` — permanently malformed. Never retry. * ``429`` —
keep, retry after ``Retry-After``. * ``5xx`` / timeout / connection error — keep, retry with
backoff.
Sender side of the ingest contract (telemetry repo ``CONTRACT.md``): ``202`` durably stored,
mark sent; ``400`` permanently malformed, never retry; ``429`` keep, retry after
``Retry-After``; ``5xx`` / timeout / connection error keep, retry with backoff.
"""
from __future__ import annotations
@@ -18,11 +16,14 @@ import time
import urllib.error
import urllib.request
import uuid
from contextlib import contextmanager
from dataclasses import dataclass
from datetime import datetime, timedelta, timezone
from hermes_cli.sqlite_util import write_txn
from .shared_metrics import _isoformat, _utc_now
logger = logging.getLogger(__name__)
#: Contract recommends timing out at 30s and treating a timeout as retryable.
@@ -41,49 +42,40 @@ GZIP_THRESHOLD_BYTES = 4096
MAX_PACKAGES_PER_PASS = 20
#: How long a claimed row is held. The claim writes a LEASE INTO THE FUTURE:
#: selection requires `next_attempt_at <= now`, so for the length of the lease
#: no other process can take the package.
#:
#: This must exceed the worst case for ONE package — three 30s request
#: timeouts plus 1s+5s of backoff, about 96s — which is why packages are
#: claimed one at a time, immediately before being sent. An earlier revision
#: claimed up to 20 rows under a single shared lease; a full batch can legally
#: run ~1900s, so the later rows' leases expired while the pass still held
#: them in memory and another process re-sent them.
#: selection requires `next_attempt_at <= now`, so no other process can take the
#: package for the length of the lease. Must exceed the worst case for ONE package
#: (three 30s timeouts plus 1s+5s backoff, ~96s) — which is why packages are claimed
#: one at a time right before sending: a batch of 20 under one lease can legally run
#: ~1900s, so later rows' leases expired while still held and got re-sent.
_CLAIM_LEASE_SECONDS = 300
#: Floor applied after a pass fails to deliver, so a hard-down service is not
#: retried on every task completion.
_FAILURE_BACKOFF_SECONDS = 15 * 60
#: Statuses that are permanent per the ingest contract. Deliberately narrow:
#: 400 means the envelope is malformed and will never validate. 413 is added
#: because a package over the service's 1 MiB cap cannot shrink on retry.
#: Everything else — including 403 from the origin guard and 404 from a bad
#: path — is retried, because those are usually deployment or edge
#: Permanent statuses per the ingest contract. Deliberately narrow: 400 means the
#: envelope will never validate; 413 because a package over the 1 MiB cap cannot
#: shrink on retry. Everything else — including 403 from the origin guard and 404
#: from a bad path — is retried, since those are usually deployment/edge
#: misconfiguration that resolves without the package changing.
_PERMANENT_STATUSES = frozenset({400, 413})
#: Attempts after which a package is abandoned. Without a ceiling a
#: permanently-poisoned row is retried until 30-day retention deletes it —
#: measured at ~160 requests — which wastes the user's bandwidth and keeps a
#: doomed package at the head of the queue.
#: Attempts after which a package is abandoned. Without a ceiling a poisoned row
#: is retried until 30-day retention deletes it (~160 requests), wasting the user's
#: bandwidth and pinning a doomed package at the head of the queue.
MAX_SEND_ATTEMPTS = 25
def _utc_now() -> datetime:
return datetime.now(timezone.utc)
def _isoformat(value: datetime) -> str:
return value.astimezone(timezone.utc).isoformat().replace("+00:00", "Z")
#: Maximum distance one reconcile call can advance the 'obs' mark. Honest heartbeats
#: arrive hours apart so the cap never binds normally; a machine off for months catches
#: up in a few hook fires (fail-closed latency only). It bounds FORWARD clock poison:
#: without it one glitched sample (NTP flap reading 2099) drags the mark — and every
#: window open and confirmation horizon — decades ahead, a refused-data leak.
MAX_OBS_ADVANCE_SECONDS = 30 * 24 * 3600
def _parse_stamp(value: str) -> datetime:
"""Parse a stamp this module itself wrote (Z-suffixed ISO-8601, UTC)."""
return datetime.fromisoformat(value.replace("Z", "+00:00")).astimezone(
timezone.utc
)
return datetime.fromisoformat(value.replace("Z", "+00:00")).astimezone(timezone.utc)
@dataclass
@@ -112,17 +104,11 @@ def _post(endpoint: str, payload: bytes, *, timeout: int) -> _Response:
}
body = payload
if len(payload) > GZIP_THRESHOLD_BYTES:
# mtime=0: gzip embeds a timestamp by default, which would make two
# sends of one package differ on the wire. The service decompresses
# before storing so it would not change what lands in S3, but a
# deterministic body keeps "a resend is byte-identical" true at the
# transport layer too, and makes the property testable.
# mtime=0 keeps two sends of one package byte-identical on the wire too.
body = gzip.compress(payload, mtime=0)
headers["Content-Encoding"] = "gzip"
request = urllib.request.Request(
endpoint, data=body, headers=headers, method="POST"
)
request = urllib.request.Request(endpoint, data=body, headers=headers, method="POST")
try:
with urllib.request.urlopen(request, timeout=timeout) as response:
return _Response(
@@ -143,25 +129,13 @@ def _retry_after_seconds(value: str | None, default: int) -> int:
if not value:
return default
try:
# Contract sends seconds. Clamp so a hostile or bogus value cannot
# park a package for years, and never go below one second.
# Contract sends seconds. Clamp so a bogus value cannot park a package for
# years, and never go below one second.
return max(1, min(int(float(value)), 86_400))
except (TypeError, ValueError):
return default
#: Maximum distance one reconcile call can advance the 'obs' mark. Honest
#: heartbeats arrive hours apart at most, so the cap never binds in normal
#: operation; a machine legitimately off for months catches up in a few
#: hook fires (fail-closed latency only). What it bounds is FORWARD clock
#: poison: without it, a single glitched sample (NTP flap reading 2099)
#: permanently drags the mark — and with it every window open and every
#: confirmation horizon — decades ahead, which round 6 reproduced as a
#: refused-data leak. Capped, one insane sample moves the mark at most
#: this far, and real time overtakes it again.
MAX_OBS_ADVANCE_SECONDS = 30 * 24 * 3600
def reconcile_send_consent(
connection: sqlite3.Connection,
send_enabled: bool,
@@ -170,9 +144,8 @@ def reconcile_send_consent(
) -> None:
"""Reconcile the consent-window table with the observed config state.
THE ONLY writer of consent state. Must run inside a write transaction. A pure function of
(config, now, store): call it from anywhere, any number of times, in any order — the resulting
windows are the same.
THE ONLY writer of consent state. Must run inside a write transaction. A pure function
of (config, now, store): call it from anywhere, any number of times, in any order.
"""
stamp = _isoformat(now or _utc_now())
raw_stamp = stamp # pre-cap observation time, used to clamp closes
@@ -181,8 +154,7 @@ def reconcile_send_consent(
).fetchone()
if previous_obs is not None:
ceiling = _isoformat(
_parse_stamp(str(previous_obs[0]))
+ timedelta(seconds=MAX_OBS_ADVANCE_SECONDS)
_parse_stamp(str(previous_obs[0])) + timedelta(seconds=MAX_OBS_ADVANCE_SECONDS)
)
stamp = min(stamp, ceiling)
connection.execute(
@@ -192,9 +164,7 @@ def reconcile_send_consent(
""",
(stamp,),
)
marks = dict(
connection.execute("SELECT name, stamp FROM consent_marks").fetchall()
)
marks = dict(connection.execute("SELECT name, stamp FROM consent_marks").fetchall())
obs = marks["obs"] # >= stamp; immune to clock rollback
data = marks.get("data")
@@ -218,17 +188,13 @@ def reconcile_send_consent(
(obs, open_row[0]),
)
elif open_row is not None:
# Close at the last CONFIRMED moment, but never after the closing
# observation's own raw stamp. The two clamps serve different
# adversaries and both are load-bearing:
# - min with last_confirmed_at: an unobserved gap (machine off,
# hand-edited config) is never asserted as consented (v1's leak).
# - min with the RAW stamp (pre-cap, pre-MAX): if last_confirmed_at
# was poisoned by a glitched-forward sample, an honest clock at
# revoke time pulls the close back to the true revoke moment, so
# the refused era that follows falls OUTSIDE the closed window
# (round 6's D1 leak). A rolled-back clock at close time only
# closes EARLIER — fail-closed.
# Close at the last CONFIRMED moment, but never after the closing observation's
# own raw stamp. Both clamps are load-bearing: min with last_confirmed_at means an
# unobserved gap (machine off, hand-edited config) is never asserted as consented;
# min with the RAW (pre-cap, pre-MAX) stamp means a last_confirmed_at poisoned by a
# glitched-forward sample is pulled back to the true revoke moment by an honest
# clock, so the refused era falls OUTSIDE the window. A rolled-back clock only
# closes EARLIER — fail-closed.
connection.execute(
"UPDATE send_consent_windows"
" SET closed_at = MIN(last_confirmed_at, ?)"
@@ -237,10 +203,10 @@ def reconcile_send_consent(
)
#: Claim-time consent predicate: the package's period must fall entirely
#: inside SOME recorded consent window. An open window vouches only up to its
#: last confirmed moment, so a package whose period runs past it waits for
#: the next reconcile heartbeat (fail-closed; released within one hook fire).
#: Claim-time consent predicate: the package's period must fall entirely inside SOME
#: recorded consent window. An open window vouches only up to its last confirmed moment,
#: so a package whose period runs past it waits for the next reconcile heartbeat
#: (fail-closed; released within one hook fire).
CONSENT_GATE_SQL = """EXISTS (
SELECT 1 FROM send_consent_windows w
WHERE package_outbox.period_start >= w.opened_at
@@ -270,36 +236,38 @@ class SharedMetricsSender:
self._sleep = sleep
self._now = now
self._max_attempts = max_attempts
# Called before every package. None disables the check for callers
# that have already established consent out of band (tests, E2E).
# Called before every package. None disables the check for callers that have
# already established consent out of band (tests, E2E).
self._consent_check = consent_check
@contextmanager
def _write(self):
"""One store connection inside a write transaction."""
with self._store._connection() as connection:
with write_txn(connection):
yield connection
# -- selection ---------------------------------------------------------
def _claim_next(self, now: datetime, seen: set[str]) -> dict | None:
"""Claim exactly ONE package, immediately before it is sent.
Claiming a batch up front fails: one shared lease would have to cover the whole pass, and 20
retrying packages can outlive any sane lease, so later rows expire and another process
re-sends them. ``seen`` (packages this pass finished with) is excluded IN SQL: with LIMIT 1,
returning None for a seen row would look like an empty queue and abandon everything behind
it, and rows can legitimately become eligible again mid-pass.
Claiming a batch up front fails: one shared lease would have to cover the whole
pass, and 20 retrying packages can outlive any sane lease. ``seen`` (packages this
pass finished with) is excluded IN SQL: with LIMIT 1, returning None for a seen row
would look like an empty queue and abandon everything behind it, and rows can
legitimately become eligible again mid-pass.
"""
with self._store._connection() as connection:
with write_txn(connection):
stamp = _isoformat(now)
lease_until = now + timedelta(seconds=_CLAIM_LEASE_SECONDS)
with self._write() as connection:
stamp = _isoformat(now)
lease_until = now + timedelta(seconds=_CLAIM_LEASE_SECONDS)
placeholders = ",".join("?" for _ in seen)
exclusion = (
f" AND package_id NOT IN ({placeholders})" if seen else ""
)
# Consent is a READ here — the claim must never mutate the
# window table. The old design's opt_in_period() call at this
# exact spot meant selecting a row could rewrite what was
# permitted to be sent (and did, under a rolled-back clock).
row = connection.execute(
f"""
placeholders = ",".join("?" for _ in seen)
exclusion = f" AND package_id NOT IN ({placeholders})" if seen else ""
# Consent is a READ here — the claim must never mutate the window table, or
# selecting a row could rewrite what is permitted to be sent.
row = connection.execute(
f"""
SELECT package_id, payload_json, sent_install_id
FROM package_outbox
WHERE exported_at IS NOT NULL
@@ -311,25 +279,20 @@ class SharedMetricsSender:
ORDER BY created_at, package_id
LIMIT 1
""",
(stamp, MAX_SEND_ATTEMPTS, *sorted(seen)),
).fetchone()
if row is None:
return None
(stamp, MAX_SEND_ATTEMPTS, *sorted(seen)),
).fetchone()
if row is None:
return None
package_id = str(row[0])
derived = row[2]
if not derived:
derived = self._freeze_identity(
connection, package_id, row[1], now
)
if derived is None:
# Unusable row, already marked rejected. Signal the
# caller to continue rather than stop.
return {"package_id": package_id, "skip": True}
package_id = str(row[0])
derived = row[2] or self._freeze_identity(connection, package_id, row[1])
if derived is None:
# Unusable row, already marked rejected. Tell the caller to continue.
return {"package_id": package_id, "skip": True}
token = str(uuid.uuid4())
connection.execute(
"""
token = str(uuid.uuid4())
connection.execute(
"""
UPDATE package_outbox
SET send_state = 'pending',
send_attempts = send_attempts + 1,
@@ -337,58 +300,49 @@ class SharedMetricsSender:
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.
# 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,
}
# Lease INTO THE FUTURE (see _CLAIM_LEASE_SECONDS). 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,
}
@staticmethod
def _freeze_identity(
self,
connection: sqlite3.Connection,
package_id: str,
payload_json,
now: datetime,
connection: sqlite3.Connection, package_id: str, payload_json
) -> str | None:
"""Record the transmitted id on the row, or reject an unusable one.
"""Record the transmitted install_id on the row, or reject an unusable one.
What remains of "freezing" is the validation and the audit column: ``sent_install_id``
records exactly what the wire will carry, and rejecting unusable rows here rather than
raising matters because an exception rolls back the claim transaction and blocks every
``sent_install_id`` records exactly what the wire will carry. Rejecting here rather
than raising matters: an exception rolls back the claim transaction and blocks every
healthy package behind this one.
"""
reason = None
install_id = None
try:
payload = json.loads(payload_json)
except (TypeError, ValueError):
reason = "unreadable payload"
else:
# Valid JSON is not enough: a top-level array, string, number or
# null parses cleanly and then has no .get().
# Valid JSON is not enough: a top-level array/string/number parses cleanly.
if not isinstance(payload, dict):
reason = f"payload is {type(payload).__name__}, expected object"
else:
install_id = payload.get("install_id")
if not isinstance(install_id, str) or not install_id.strip():
reason = "payload has no usable install_id"
reason = (
None
if isinstance(install_id, str) and install_id.strip()
else "payload has no usable install_id"
)
if reason is not None:
logger.warning(
"Shared-metrics package %s cannot be sent (%s)", package_id, reason
)
logger.warning("Shared-metrics package %s cannot be sent (%s)", package_id, reason)
connection.execute(
"""
UPDATE package_outbox
@@ -407,16 +361,15 @@ class SharedMetricsSender:
# -- transmission ------------------------------------------------------
def _body(self, payload_json: str, transmitted_id: str) -> bytes:
@staticmethod
def _body(payload_json: str, transmitted_id: str) -> bytes:
"""Rebuild the exact bytes to send.
The payload is recomputed from the stored package rather than kept as a second copy:
json.dumps with these options is deterministic. The install_id is written from the frozen
``sent_install_id`` column rather than trusted implicitly, keeping "a resend is byte-
identical" anchored to one recorded value.
Recomputed from the stored package (json.dumps with these options is deterministic)
with install_id taken from the frozen ``sent_install_id`` column, keeping "a resend
is byte-identical" anchored to one recorded value.
"""
payload = json.loads(payload_json)
payload = dict(payload)
payload = dict(json.loads(payload_json))
payload["install_id"] = transmitted_id
return json.dumps(payload, indent=2, sort_keys=True).encode("utf-8")
@@ -430,54 +383,40 @@ class SharedMetricsSender:
) -> 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.
Guarded on send_state so a pass whose lease lapsed cannot resurrect a row another
process already finished (overwriting 'sent' back to 'pending' would re-send). With
``token`` the write is also compare-and-set on claim_token: a superseded claimant
writes zero rows and its settlement/backoff/error silently lose to the newer claim.
"""
assignments = ", ".join(f"{name} = ?" for name in columns)
predicate = (
" AND (send_state IS NULL OR send_state = 'pending')"
if only_if_pending
else ""
)
predicate = " AND (send_state IS NULL OR send_state = 'pending')" 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}",
params,
)
with self._write() as connection:
connection.execute(
f"UPDATE package_outbox SET {assignments} WHERE package_id = ?{predicate}",
params,
)
def _renew_claim(self, package_id: str, token: str | None) -> bool:
"""Atomically re-assert ownership and extend the lease. CAS, one row.
A read-only ownership check is not enough: a claimant whose lease expired while suspended
can pass the check (its token is still in the row if no one reclaimed yet) and then POST
while another process legitimately reclaims — the check-to-POST expiry race a seventh review
reproduced.
and only then pushing next_attempt_at a fresh lease into the future, so the upcoming POST
(30s timeout, well under the 300s lease) runs entirely inside renewed authority. rowcount ==
1 is the only grant.
A read-only ownership check is not enough: a claimant whose lease expired while
suspended still has its token in the row if nobody reclaimed yet, and would POST
while another process legitimately reclaims (check-to-POST expiry race). Pushing a
fresh lease in the same CAS keeps the upcoming POST (30s, well under the 300s lease)
entirely inside renewed authority. rowcount == 1 is the only grant.
"""
if token is None:
return False
try:
now = self._now()
lease_until = now + timedelta(seconds=_CLAIM_LEASE_SECONDS)
with self._store._connection() as connection:
with write_txn(connection):
cursor = connection.execute(
"""
with self._write() as connection:
cursor = connection.execute(
"""
UPDATE package_outbox
SET next_attempt_at = ?
WHERE package_id = ?
@@ -485,176 +424,120 @@ class SharedMetricsSender:
AND (send_state IS NULL OR send_state = 'pending')
AND next_attempt_at > ?
""",
(
_isoformat(lease_until),
package_id,
token,
_isoformat(now),
),
)
return cursor.rowcount == 1
(_isoformat(lease_until), package_id, token, _isoformat(now)),
)
return cursor.rowcount == 1
except Exception:
# If renewal itself fails, do not transmit on unproven authority.
logger.warning(
"Unable to renew shared-metrics claim", exc_info=True
)
logger.warning("Unable to renew shared-metrics claim", exc_info=True)
return False
def _defer(
self,
package_id: str,
delay_seconds: int,
reason: str,
*,
token: str | None = None,
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
# deliberately unreachable today and no test can distinguish it. It
# stays because a past deadline would make the row instantly
# re-eligible and let a pass spin on it — a cheap guard against a
# future caller that forgets.
delay = max(1, int(delay_seconds))
retry_at = self._now().timestamp() + delay
# Clamp to >= 1s: a past deadline would make the row instantly re-eligible and let
# a pass spin on it. Unreachable today (Retry-After is pre-clamped, other callers
# pass positive constants) — kept as a guard against a future caller.
retry_at = self._now().timestamp() + max(1, int(delay_seconds))
self._mark(
package_id,
token=token,
send_state="pending",
next_attempt_at=_isoformat(
datetime.fromtimestamp(retry_at, tz=timezone.utc)
),
next_attempt_at=_isoformat(datetime.fromtimestamp(retry_at, tz=timezone.utc)),
last_error=reason[:500],
)
def _send_one(self, package: dict) -> str:
"""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.
Delivery is at-least-once. The pre-POST renewal plus token-fenced writes close the
claim->POST and settle-after-reclaim gaps, but a suspension landing MID-POST can
still duplicate: no client-side check can revoke a request in flight.
"""
package_id = package["package_id"]
token = package.get("claim_token")
body = self._body(package["payload_json"], package["derived"])
def defer(delay: int, reason: str) -> str:
self._defer(package_id, delay, reason, token=token)
return "deferred"
for attempt in range(1, self._max_attempts + 1):
# Atomically renew the claim before EVERY external POST. The
# renewal is compare-and-set on (token, pending, lease unexpired)
# and extends the lease past the request, so a suspended-then-
# resumed claimant whose lease lapsed yields here even if nobody
# has reclaimed yet — a read-only ownership check passed in that
# state and still double-sent (check-to-POST expiry race). The
# ingest key is minute-prefixed, so duplicates become distinct
# Renew before EVERY external POST (CAS on token, pending, lease unexpired), so a
# suspended-then-resumed claimant whose lease lapsed yields even before anyone
# reclaims. The ingest key is minute-prefixed, so duplicates become distinct
# stored objects, not overwrites.
if not self._renew_claim(package_id, token):
logger.info(
"Shared-metrics claim on %s superseded or expired; yielding",
package_id,
"Shared-metrics claim on %s superseded or expired; yielding", package_id
)
return "deferred"
try:
response = self._post(
self._endpoint, body, timeout=REQUEST_TIMEOUT_SECONDS
)
response = self._post(self._endpoint, body, timeout=REQUEST_TIMEOUT_SECONDS)
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, token=token
else:
if response.status == 202:
self._mark(
package_id, token=token, send_state="sent",
sent_at=_isoformat(self._now()), last_error=None,
)
return "deferred"
self._sleep(self._backoff(attempt))
continue
if response.status == 202:
self._mark(
package_id,
token=token,
send_state="sent",
sent_at=_isoformat(self._now()),
last_error=None,
)
return "sent"
if response.status in _PERMANENT_STATUSES:
# Only statuses the contract (or the envelope schema) makes
# terminal. Everything else retries: 403 in particular is the
# ingest service's origin guard, which returns 403 during an
# edge/Transform-Rule misconfiguration — treating that as
# permanent would discard every package sent during the
# incident instead of retrying after recovery.
logger.warning(
"Telemetry package %s rejected with HTTP %s; not retrying",
package_id,
response.status,
)
self._mark(
package_id,
token=token,
send_state="rejected",
last_error=f"HTTP {response.status}: {response.body[:400]}",
)
return "rejected"
if response.status == 429:
self._defer(
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}"
return "sent"
if response.status in _PERMANENT_STATUSES:
# Only contract-terminal statuses. 403 in particular is the origin guard
# during an edge misconfiguration — treating it as permanent would
# discard every package sent during the incident.
logger.warning(
"Telemetry package %s rejected with HTTP %s; not retrying",
package_id,
response.status,
)
self._mark(
package_id, token=token, send_state="rejected",
last_error=f"HTTP {response.status}: {response.body[:400]}",
)
return "rejected"
if response.status == 429:
return defer(
_retry_after_seconds(response.retry_after, _FAILURE_BACKOFF_SECONDS),
"rate limited",
)
# 5xx and anything unexpected: retryable.
reason = f"HTTP {response.status}"
if attempt >= self._max_attempts:
self._defer(
package_id, _FAILURE_BACKOFF_SECONDS, reason, token=token
)
return "deferred"
return defer(_FAILURE_BACKOFF_SECONDS, reason)
self._sleep(self._backoff(attempt))
self._defer(
package_id, _FAILURE_BACKOFF_SECONDS, "attempts exhausted", token=token
)
return "deferred"
return defer(_FAILURE_BACKOFF_SECONDS, "attempts exhausted")
@staticmethod
def _backoff(attempt: int) -> float:
"""1s, 5s, 25s with full jitter."""
ceiling = _BACKOFF_BASE_SECONDS * (_BACKOFF_FACTOR ** (attempt - 1))
return random.uniform(0, ceiling)
return random.uniform(0, _BACKOFF_BASE_SECONDS * (_BACKOFF_FACTOR ** (attempt - 1)))
# -- entry point -------------------------------------------------------
def send_pending(self) -> SendOutcome:
"""Run one bounded pass. Never raises.
Claims and sends ONE package at a time so each row's lease only has to cover its own
transmission, and re-checks consent before every send so revoking `send` mid-pass stops the
remaining packages.
Claims and sends ONE package at a time so each row's lease only covers its own
transmission, and re-checks consent before every send so revoking `send` mid-pass
stops the remaining packages.
"""
outcome = SendOutcome()
seen: set[str] = set()
for _ in range(MAX_PACKAGES_PER_PASS):
if not self._still_consented():
# The user turned sending off while this pass was running.
# Stop without transmitting anything further, and reconcile
# so the window closes at its last confirmed moment. This is
# the same single writer every other observation point uses —
# not a separate recording mechanism.
# Stop transmitting and reconcile through the single consent writer so
# the window closes at its last confirmed moment.
logger.info("Shared-metrics sending disabled mid-pass; stopping")
self._reconcile(send_enabled=False)
break
try:
package = self._claim_next(self._now(), seen)
except Exception:
logger.warning(
"Unable to select shared-metrics packages", exc_info=True
)
logger.warning("Unable to select shared-metrics packages", exc_info=True)
break
if package is None:
break
@@ -662,42 +545,29 @@ class SharedMetricsSender:
seen.add(package["package_id"])
if package.get("skip"):
# Unusable row already marked rejected during the claim.
outcome.rejected += 1
continue
try:
result = self._send_one(package)
except Exception:
logger.warning("Unable to send shared-metrics package", exc_info=True)
outcome.deferred += 1
continue
if result == "sent":
outcome.sent += 1
elif result == "rejected":
outcome.rejected += 1
result = "rejected"
else:
outcome.deferred += 1
try:
result = self._send_one(package)
except Exception:
logger.warning("Unable to send shared-metrics package", exc_info=True)
result = "deferred"
setattr(outcome, result, getattr(outcome, result) + 1)
return outcome
def _reconcile(self, *, send_enabled: bool) -> None:
"""Run the single consent writer from within a pass."""
try:
with self._store._connection() as connection:
with write_txn(connection):
reconcile_send_consent(
connection, send_enabled, now=self._now()
)
with self._write() as connection:
reconcile_send_consent(connection, send_enabled, now=self._now())
except Exception:
logger.warning(
"Unable to reconcile shared-metrics consent", exc_info=True
)
logger.warning("Unable to reconcile shared-metrics consent", exc_info=True)
def _still_consented(self) -> bool:
"""Re-read profile-owned send consent.
Consent is a boundary, not cached config: docs promise ``send: false`` stops transmission
immediately, and a pass can run for minutes. Injected senders opt out via
consent_check=None.
Consent is a boundary, not cached config: docs promise ``send: false`` stops
transmission immediately, and a pass can run for minutes.
"""
if self._consent_check is None:
return True
@@ -705,8 +575,5 @@ class SharedMetricsSender:
return bool(self._consent_check())
except Exception:
# Fail CLOSED: if consent cannot be established, do not transmit.
logger.warning(
"Unable to confirm shared-metrics send consent; stopping",
exc_info=True,
)
logger.warning("Unable to confirm shared-metrics send consent; stopping", exc_info=True)
return False