From 6837d2e41bcd19e72a854b1df21a679c055b86eb Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 20:48:41 -0700 Subject: [PATCH] =?UTF-8?q?refactor(hermes=5Fcli):=20shared=5Fmetrics=5Fse?= =?UTF-8?q?nder=20=E2=80=94=20reuse=20store=20time=20helpers,=20unify=20de?= =?UTF-8?q?fer=20paths,=20txn=20ctx=20manager?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../observability/shared_metrics_sender.py | 531 +++++++----------- 1 file changed, 199 insertions(+), 332 deletions(-) diff --git a/hermes_cli/observability/shared_metrics_sender.py b/hermes_cli/observability/shared_metrics_sender.py index 5bcedba764..dab55fdbd1 100644 --- a/hermes_cli/observability/shared_metrics_sender.py +++ b/hermes_cli/observability/shared_metrics_sender.py @@ -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