From 00c75cea335b2f90ceed648ab02542d6d2043b9e Mon Sep 17 00:00:00 2001 From: Ben Barclay Date: Wed, 26 Aug 2026 15:45:02 +1000 Subject: [PATCH] feat(telemetry): send exported packages to the ingest service MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Steps 4, 5 and 7 of the shared-metrics exporter: the send logic, the consent gate, and backoff plus multi-process claiming. These arrive together because the sender is not correct without all three. Contract handling: 202 marks sent; 400 is permanent and never retried; 429 honours Retry-After (clamped to a day so a bogus value cannot park a package); 5xx, timeouts and transport errors retry three times in-process with 1s/5s/25s full-jitter backoff, then defer to a later pass. Consent is gated on the package's PERIOD, not its creation time. A period is split across packages created on different days, so a created-at gate would send a period's tail while dropping its head and silently undercount the opt-in day — data that looks complete and is wrong. The opt-in day is recorded once and never moves, so toggling sending off and on does not re-open the pre-consent backlog. Rows are claimed in a write transaction, which is what stops two Hermes processes sharing one database from sending the same package twice. next_attempt_at persists backoff across restarts, so a hard-down service is not retried on every task completion. The body is recomputed from payload_json rather than stored a second time: json.dumps is deterministic here (verified against the real outbox — 11 of 11 files reproduce byte-for-byte), and the only mutable input, the derived identity, is frozen on the row at first attempt. That keeps retries byte-identical across a salt rotation for ~36 bytes instead of a duplicate ~11 KB payload. The outbox directory is never written to or deleted from. A 202 updates SQLite only, because those files are the user's 30-day local history and retention already owns their lifecycle. Tests: 33. Two of them caught real defects in this commit — an unreadable row aborted the claim transaction and blocked every package behind it, and the compression assertions were passing through an injected fake that bypassed the code under test. --- .../observability/shared_metrics_sender.py | 385 +++++++++++++++ .../hermes_cli/test_shared_metrics_sender.py | 449 ++++++++++++++++++ 2 files changed, 834 insertions(+) create mode 100644 hermes_cli/observability/shared_metrics_sender.py create mode 100644 tests/hermes_cli/test_shared_metrics_sender.py diff --git a/hermes_cli/observability/shared_metrics_sender.py b/hermes_cli/observability/shared_metrics_sender.py new file mode 100644 index 0000000000..9086c3b359 --- /dev/null +++ b/hermes_cli/observability/shared_metrics_sender.py @@ -0,0 +1,385 @@ +"""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. + +Two properties are load-bearing and easy to get wrong: + +**The outbox directory is the user's local history, not a queue.** Packages +are pruned by age; a ``202`` marks send state in SQLite and never deletes a +file. See Appendix A.7 of ``docs/observability/relay-shared-metrics.md``. + +**Consent is gated on the package's PERIOD, not its creation time.** One +period is split across packages created on different days, so a created-at +gate would send a period's tail while dropping its head and silently +undercount the opt-in day. +""" + +from __future__ import annotations + +import gzip +import json +import logging +import random +import sqlite3 +import time +import urllib.error +import urllib.request +from dataclasses import dataclass +from datetime import datetime, timezone + +from hermes_cli.sqlite_util import write_txn + +from .shared_metrics_identity import ( + current_salt, + derive_install_id, + substitute_install_id, +) + +logger = logging.getLogger(__name__) + +#: Contract recommends timing out at 30s and treating a timeout as retryable. +REQUEST_TIMEOUT_SECONDS = 30 + +#: In-process attempts per package per pass, then the package waits for a +#: later pass. Backoff is 1s/5s/25s with full jitter. +MAX_ATTEMPTS = 3 +_BACKOFF_BASE_SECONDS = 1 +_BACKOFF_FACTOR = 5 + +#: Contract recommends gzip above roughly this size. +GZIP_THRESHOLD_BYTES = 4096 + +#: Packages per pass. Bounds work on an interactive hook even after an outage. +MAX_PACKAGES_PER_PASS = 20 + +#: 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 + +OPT_IN_PERIOD_KEY = "send_opt_in_period" + + +def _utc_now() -> datetime: + return datetime.now(timezone.utc) + + +def _isoformat(value: datetime) -> str: + return value.astimezone(timezone.utc).isoformat().replace("+00:00", "Z") + + +@dataclass +class SendOutcome: + """What one pass did. Returned for tests and diagnostics.""" + + sent: int = 0 + rejected: int = 0 + deferred: int = 0 + skipped_not_due: int = 0 + + +class _Response: + __slots__ = ("status", "retry_after", "body") + + def __init__(self, status: int, retry_after: str | None, body: str) -> None: + self.status = status + self.retry_after = retry_after + self.body = body + + +def _post(endpoint: str, payload: bytes, *, timeout: int) -> _Response: + """POST one package. Raises on transport failure; never on HTTP status.""" + headers = { + "Content-Type": "application/json", + "User-Agent": "hermes-agent-shared-metrics/1", + } + body = payload + if len(payload) > GZIP_THRESHOLD_BYTES: + body = gzip.compress(payload) + headers["Content-Encoding"] = "gzip" + + request = urllib.request.Request( + endpoint, data=body, headers=headers, method="POST" + ) + try: + with urllib.request.urlopen(request, timeout=timeout) as response: + return _Response( + response.status, + response.headers.get("Retry-After"), + response.read(2048).decode("utf-8", "replace"), + ) + except urllib.error.HTTPError as exc: + # An HTTP error status is a normal contract outcome, not a failure. + return _Response( + exc.code, + exc.headers.get("Retry-After") if exc.headers else None, + exc.read(2048).decode("utf-8", "replace") if exc.fp else "", + ) + + +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. + return max(1, min(int(float(value)), 86_400)) + except (TypeError, ValueError): + return default + + +def opt_in_period(connection: sqlite3.Connection, *, now: datetime | None = None) -> str: + """Return the opt-in day (UTC date), recording it on first use. + + Must run inside a write transaction. The value is written once and then + never moves, so turning sending off and on again does not re-open the + pre-consent backlog. + """ + row = connection.execute( + "SELECT value FROM telemetry_state WHERE key = ?", (OPT_IN_PERIOD_KEY,) + ).fetchone() + if row is not None: + return str(row[0]) + today = (now or _utc_now()).date().isoformat() + connection.execute( + "INSERT OR IGNORE INTO telemetry_state(key, value) VALUES (?, ?)", + (OPT_IN_PERIOD_KEY, today), + ) + return today + + +class SharedMetricsSender: + """Sends exported packages, one bounded pass at a time.""" + + def __init__( + self, + store, + endpoint: str, + *, + post=_post, + sleep=time.sleep, + now=_utc_now, + max_attempts: int = MAX_ATTEMPTS, + ) -> None: + self._store = store + self._endpoint = endpoint + self._post = post + self._sleep = sleep + self._now = now + self._max_attempts = max_attempts + + # -- selection --------------------------------------------------------- + + def _claim(self, connection: sqlite3.Connection, now: datetime) -> list[dict]: + """Atomically take ownership of the packages this pass will try. + + Claiming inside the write transaction is what stops two Hermes + processes sharing one database from sending the same package twice. + Duplicates would be harmless (the service dedupes by package_id and + the bytes are identical) but they waste the user's bandwidth. + """ + period = opt_in_period(connection, now=now) + stamp = _isoformat(now) + rows = connection.execute( + """ + SELECT package_id, payload_json, sent_install_id + FROM package_outbox + WHERE exported_at IS NOT NULL + AND (send_state IS NULL OR send_state = 'pending') + AND (next_attempt_at IS NULL OR next_attempt_at <= ?) + AND substr(period_start, 1, 10) >= ? + ORDER BY created_at, package_id + LIMIT ? + """, + (stamp, period, MAX_PACKAGES_PER_PASS), + ).fetchall() + + claimed: list[dict] = [] + salt: str | None = None + for row in rows: + package_id = str(row[0]) + derived = row[2] + if not derived: + # Freeze the derived identity on first attempt so a later salt + # rotation cannot change the bytes sent under this package_id. + if salt is None: + salt = current_salt(connection, now=now) + try: + payload = json.loads(row[1]) + install_id = str(payload.get("install_id", "")) + except (TypeError, ValueError): + # A row we cannot parse can never be sent. Mark it and move + # on: one unreadable package must not block every other + # package behind it, and aborting here would roll back the + # whole claim transaction. + logger.warning( + "Shared-metrics package %s is unreadable; not sending", + package_id, + ) + connection.execute( + """ + UPDATE package_outbox + SET send_state = 'rejected', last_error = 'unreadable payload' + WHERE package_id = ? + """, + (package_id,), + ) + continue + derived = derive_install_id(install_id, salt) + connection.execute( + "UPDATE package_outbox SET sent_install_id = ? WHERE package_id = ?", + (derived, package_id), + ) + connection.execute( + """ + UPDATE package_outbox + SET send_state = 'pending', + send_attempts = send_attempts + 1, + next_attempt_at = ? + WHERE package_id = ? + """, + # Hold the row for the duration of this pass; success or a + # real backoff overwrite this immediately below. + (_isoformat(now), package_id), + ) + claimed.append( + { + "package_id": package_id, + "payload_json": str(row[1]), + "derived": str(derived), + } + ) + return claimed + + # -- transmission ------------------------------------------------------ + + def _body(self, payload_json: str, derived: 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, and the + only mutable input (the derived id) is frozen in the row. + """ + 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, **columns) -> None: + assignments = ", ".join(f"{name} = ?" for name in columns) + with self._store._connection() as connection: + with write_txn(connection): + connection.execute( + f"UPDATE package_outbox SET {assignments} WHERE package_id = ?", + (*columns.values(), package_id), + ) + + def _defer(self, package_id: str, delay_seconds: int, reason: str) -> None: + retry_at = self._now().timestamp() + delay_seconds + self._mark( + package_id, + send_state="pending", + 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'.""" + package_id = package["package_id"] + body = self._body(package["payload_json"], package["derived"]) + + for attempt in range(1, self._max_attempts + 1): + try: + 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) + return "deferred" + self._sleep(self._backoff(attempt)) + continue + + if response.status == 202: + self._mark( + package_id, + send_state="sent", + sent_at=_isoformat(self._now()), + last_error=None, + ) + return "sent" + + if response.status == 400: + # Permanent per the contract. Keep the file (it is the user's + # history) but never try again. + logger.warning( + "Telemetry package %s rejected as malformed; not retrying", + package_id, + ) + self._mark( + package_id, + send_state="rejected", + last_error=response.body[:500], + ) + return "rejected" + + if response.status == 429: + self._defer( + package_id, + _retry_after_seconds(response.retry_after, _FAILURE_BACKOFF_SECONDS), + "rate limited", + ) + 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) + return "deferred" + self._sleep(self._backoff(attempt)) + + self._defer(package_id, _FAILURE_BACKOFF_SECONDS, "attempts exhausted") + return "deferred" + + @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) + + # -- entry point ------------------------------------------------------- + + def send_pending(self) -> SendOutcome: + """Run one bounded pass. Never raises.""" + outcome = SendOutcome() + try: + now = self._now() + with self._store._connection() as connection: + with write_txn(connection): + claimed = self._claim(connection, now) + except Exception: + logger.warning("Unable to select shared-metrics packages", exc_info=True) + return outcome + + for package in claimed: + 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 + else: + outcome.deferred += 1 + return outcome diff --git a/tests/hermes_cli/test_shared_metrics_sender.py b/tests/hermes_cli/test_shared_metrics_sender.py new file mode 100644 index 0000000000..d0fec15bfa --- /dev/null +++ b/tests/hermes_cli/test_shared_metrics_sender.py @@ -0,0 +1,449 @@ +"""Tests for the shared-metrics sender. + +Covers the four contract responses, the period-based consent gate, frozen +identity across rotation, transactional claiming, and the invariant that +matters most: a package file is never deleted, because the outbox is the +user's local history rather than a send queue. +""" + +from __future__ import annotations + +import json +import sqlite3 +from datetime import datetime, timedelta, timezone + +import pytest + +from hermes_cli.observability.shared_metrics import SharedMetricsStore +from hermes_cli.observability.shared_metrics_sender import ( + MAX_PACKAGES_PER_PASS, + OPT_IN_PERIOD_KEY, + SharedMetricsSender, + opt_in_period, +) + +INSTALL_ID = "12a73e97-4de9-4766-830d-9ca1192c0420" +NOW = datetime(2026, 8, 26, 12, 0, tzinfo=timezone.utc) +ENDPOINT = "https://telemetry.test/v1/telemetry" + + +class FakeResponse: + def __init__(self, status, retry_after=None, body=""): + self.status = status + self.retry_after = retry_after + self.body = body + + +class FakeTransport: + """Records every POST and replays a scripted sequence of responses.""" + + def __init__(self, *responses): + self._responses = list(responses) + self.calls = [] + + def __call__(self, endpoint, payload, *, timeout): + self.calls.append({"endpoint": endpoint, "payload": payload, "timeout": timeout}) + if not self._responses: + return FakeResponse(202) + item = self._responses.pop(0) + if isinstance(item, Exception): + raise item + return item + + @property + def bodies(self): + return [json.loads(c["payload"].decode("utf-8")) for c in self.calls] + + +@pytest.fixture +def store(tmp_path): + return SharedMetricsStore( + database_path=tmp_path / "metrics.sqlite3", + outbox_directory=tmp_path / "outbox", + ) + + +def _add_package(store, package_id, period_day, *, exported=True, install_id=INSTALL_ID): + payload = { + "schema_version": "hermes.shared_metrics.v2", + "package_id": package_id, + "install_id": install_id, + "period_start": f"{period_day}T00:00:00Z", + "period_end": f"{period_day}T23:59:59Z", + "metrics": [{"name": "hermes.client.active", "type": "counter", "value": 1}], + } + with store._connection() as connection: + connection.execute( + """ + INSERT INTO package_outbox( + package_id, period_start, period_end, payload_json, + created_at, exported_at + ) VALUES (?, ?, ?, ?, ?, ?) + """, + ( + package_id, + f"{period_day}T00:00:00Z", + f"{period_day}T23:59:59Z", + json.dumps(payload), + f"{period_day}T01:00:00Z", + f"{period_day}T01:00:01Z" if exported else None, + ), + ) + path = store.outbox_directory / f"{package_id}.json" + path.write_text(json.dumps(payload, indent=2, sort_keys=True)) + return path + + +def _row(store, package_id): + with store._connection() as connection: + row = connection.execute( + """ + SELECT send_state, sent_at, send_attempts, next_attempt_at, + last_error, sent_install_id + FROM package_outbox WHERE package_id = ? + """, + (package_id,), + ).fetchone() + return dict( + send_state=row[0], + sent_at=row[1], + send_attempts=row[2], + next_attempt_at=row[3], + last_error=row[4], + sent_install_id=row[5], + ) + + +def _sender(store, transport, **kwargs): + return SharedMetricsSender( + store, + ENDPOINT, + post=transport, + sleep=lambda _s: None, + now=lambda: NOW, + **kwargs, + ) + + +class TestContractResponses: + def test_202_marks_sent(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(202)) + outcome = _sender(store, transport).send_pending() + assert outcome.sent == 1 + row = _row(store, "pkg-1") + assert row["send_state"] == "sent" + assert row["sent_at"] is not None + + def test_400_is_permanent_and_never_retried(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(400, body='{"error":"invalid_envelope"}')) + outcome = _sender(store, transport).send_pending() + assert outcome.rejected == 1 + assert len(transport.calls) == 1, "a 400 must not be retried" + assert _row(store, "pkg-1")["send_state"] == "rejected" + + # A later pass must not pick it up again. + transport2 = FakeTransport(FakeResponse(202)) + _sender(store, transport2).send_pending() + assert transport2.calls == [] + + def test_429_defers_using_retry_after(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(429, retry_after="120")) + outcome = _sender(store, transport).send_pending() + assert outcome.deferred == 1 + assert len(transport.calls) == 1, "429 waits rather than burning attempts" + row = _row(store, "pkg-1") + assert row["send_state"] == "pending" + assert row["next_attempt_at"] == "2026-08-26T12:02:00Z" + + def test_429_without_retry_after_still_defers(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(429)) + _sender(store, transport).send_pending() + assert _row(store, "pkg-1")["next_attempt_at"] > "2026-08-26T12:00:00Z" + + def test_absurd_retry_after_is_clamped(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(429, retry_after="99999999")) + _sender(store, transport).send_pending() + # clamped to 24h, not years + assert _row(store, "pkg-1")["next_attempt_at"] <= "2026-08-27T12:00:00Z" + + def test_5xx_retries_then_defers(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport( + FakeResponse(503), FakeResponse(503), FakeResponse(503) + ) + outcome = _sender(store, transport).send_pending() + assert outcome.deferred == 1 + assert len(transport.calls) == 3, "three in-process attempts" + assert _row(store, "pkg-1")["send_state"] == "pending" + + def test_5xx_then_success_within_the_same_pass(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(503), FakeResponse(202)) + outcome = _sender(store, transport).send_pending() + assert outcome.sent == 1 + assert len(transport.calls) == 2 + + def test_transport_failure_is_retryable(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport( + OSError("offline"), OSError("offline"), FakeResponse(202) + ) + outcome = _sender(store, transport).send_pending() + assert outcome.sent == 1 + + def test_persistent_offline_defers_without_raising(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(*[OSError("offline")] * 3) + outcome = _sender(store, transport).send_pending() + assert outcome.deferred == 1 + assert "OSError" in _row(store, "pkg-1")["last_error"] + + +class TestConsentGate: + def test_packages_from_before_opt_in_are_never_sent(self, store): + _add_package(store, "old", "2026-08-20") + _add_package(store, "new", "2026-08-26") + transport = FakeTransport(FakeResponse(202)) + _sender(store, transport).send_pending() + assert [b["package_id"] for b in transport.bodies] == ["new"] + + def test_a_period_straddling_opt_in_day_is_sent_whole(self, store): + """The head/tail bug: both packages for the opt-in period must go.""" + _add_package(store, "head", "2026-08-26") + _add_package(store, "tail", "2026-08-26") # created later, same period + transport = FakeTransport(FakeResponse(202), FakeResponse(202)) + _sender(store, transport).send_pending() + assert sorted(b["package_id"] for b in transport.bodies) == ["head", "tail"] + + def test_opt_in_day_is_recorded_once_and_does_not_move(self, store): + with store._connection() as connection: + first = opt_in_period(connection, now=NOW) + later = opt_in_period(connection, now=NOW + timedelta(days=10)) + assert first == later == "2026-08-26" + + def test_opt_in_day_is_persisted(self, store): + with store._connection() as connection: + opt_in_period(connection, now=NOW) + value = connection.execute( + "SELECT value FROM telemetry_state WHERE key = ?", (OPT_IN_PERIOD_KEY,) + ).fetchone()[0] + assert value == "2026-08-26" + + def test_unexported_packages_are_skipped(self, store): + _add_package(store, "pending-export", "2026-08-26", exported=False) + transport = FakeTransport(FakeResponse(202)) + _sender(store, transport).send_pending() + assert transport.calls == [] + + +class TestIdentity: + def test_install_id_is_never_transmitted(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(202)) + _sender(store, transport).send_pending() + raw = transport.calls[0]["payload"].decode("utf-8") + assert INSTALL_ID not in raw + assert transport.bodies[0]["install_id"] != INSTALL_ID + + def test_derived_id_is_frozen_on_the_row(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(503), FakeResponse(202)) + _sender(store, transport).send_pending() + assert _row(store, "pkg-1")["sent_install_id"] == transport.bodies[0]["install_id"] + + def test_retries_send_identical_bytes(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(503), FakeResponse(503), FakeResponse(202)) + _sender(store, transport).send_pending() + payloads = {c["payload"] for c in transport.calls} + assert len(payloads) == 1, "a resend must be byte-identical per the contract" + + def test_only_install_id_differs_from_the_stored_package(self, store): + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(202)) + _sender(store, transport).send_pending() + sent = transport.bodies[0] + with store._connection() as connection: + stored = json.loads( + connection.execute( + "SELECT payload_json FROM package_outbox WHERE package_id = 'pkg-1'" + ).fetchone()[0] + ) + assert set(sent) == set(stored) + for key in stored: + if key != "install_id": + assert sent[key] == stored[key] + + +class TestOutboxIsNotAQueue: + def test_a_sent_package_file_is_not_deleted(self, store): + path = _add_package(store, "pkg-1", "2026-08-26") + _sender(store, FakeTransport(FakeResponse(202))).send_pending() + assert path.exists(), "the outbox is the user's history, not a send queue" + + def test_a_rejected_package_file_is_not_deleted(self, store): + path = _add_package(store, "pkg-1", "2026-08-26") + _sender(store, FakeTransport(FakeResponse(400))).send_pending() + assert path.exists() + + def test_the_package_row_survives_sending(self, store): + _add_package(store, "pkg-1", "2026-08-26") + _sender(store, FakeTransport(FakeResponse(202))).send_pending() + with store._connection() as connection: + assert connection.execute( + "SELECT COUNT(*) FROM package_outbox WHERE package_id = 'pkg-1'" + ).fetchone()[0] == 1 + + +class TestClaimingAndBounds: + def test_a_sent_package_is_not_resent(self, store): + _add_package(store, "pkg-1", "2026-08-26") + _sender(store, FakeTransport(FakeResponse(202))).send_pending() + second = FakeTransport(FakeResponse(202)) + _sender(store, second).send_pending() + assert second.calls == [] + + def test_a_deferred_package_is_skipped_until_due(self, store): + _add_package(store, "pkg-1", "2026-08-26") + _sender(store, FakeTransport(FakeResponse(429, retry_after="600"))).send_pending() + second = FakeTransport(FakeResponse(202)) + _sender(store, second).send_pending() + assert second.calls == [], "backoff must survive within the same process" + + def test_a_deferred_package_is_retried_once_due(self, store): + _add_package(store, "pkg-1", "2026-08-26") + _sender(store, FakeTransport(FakeResponse(429, retry_after="60"))).send_pending() + + later = SharedMetricsSender( + store, + ENDPOINT, + post=(transport := FakeTransport(FakeResponse(202))), + sleep=lambda _s: None, + now=lambda: NOW + timedelta(minutes=5), + ) + later.send_pending() + assert len(transport.calls) == 1 + + def test_attempts_are_counted(self, store): + _add_package(store, "pkg-1", "2026-08-26") + _sender(store, FakeTransport(FakeResponse(429))).send_pending() + assert _row(store, "pkg-1")["send_attempts"] == 1 + + def test_a_pass_is_bounded(self, store): + for i in range(MAX_PACKAGES_PER_PASS + 5): + _add_package(store, f"pkg-{i:02d}", "2026-08-26") + transport = FakeTransport(*[FakeResponse(202)] * 40) + outcome = _sender(store, transport).send_pending() + assert outcome.sent == MAX_PACKAGES_PER_PASS + + def test_two_concurrent_passes_do_not_double_send(self, store): + """Claiming is what stops two Hermes processes duplicating work.""" + _add_package(store, "pkg-1", "2026-08-26") + + seen = [] + + def transport(endpoint, payload, *, timeout): + seen.append(payload) + # A second sender runs while the first is mid-flight. + SharedMetricsSender( + store, + ENDPOINT, + post=lambda *a, **k: (_ for _ in ()).throw( + AssertionError("second pass must not claim a held package") + ), + sleep=lambda _s: None, + now=lambda: NOW, + ).send_pending() + return FakeResponse(202) + + _sender(store, transport).send_pending() + assert len(seen) == 1 + + +class TestResilience: + def test_a_corrupt_row_does_not_stop_the_pass(self, store): + _add_package(store, "good", "2026-08-26") + with store._connection() as connection: + connection.execute( + """ + INSERT INTO package_outbox( + package_id, period_start, period_end, payload_json, + created_at, exported_at + ) VALUES ('bad', '2026-08-26T00:00:00Z', '2026-08-26T23:59:59Z', + 'not json', '2026-08-26T00:00:00Z', '2026-08-26T01:00:00Z') + """ + ) + transport = FakeTransport(*[FakeResponse(202)] * 5) + outcome = _sender(store, transport).send_pending() + assert outcome.sent >= 1 + + def test_send_pending_never_raises_on_a_broken_database(self, store, tmp_path): + store.database_path.write_text("this is not a database") + outcome = _sender(store, FakeTransport(FakeResponse(202))).send_pending() + assert outcome.sent == 0 + + +class TestCompression: + """Compression lives in the real transport, so exercise _post directly.""" + + def _captured_request(self, payload: bytes): + import urllib.request + + from hermes_cli.observability import shared_metrics_sender as mod + + captured = {} + + class FakeConn: + status = 202 + headers = {} + + def read(self, _n=None): + return b"{}" + + def __enter__(self): + return self + + def __exit__(self, *a): + return False + + def fake_urlopen(request, timeout=None): + captured["data"] = request.data + captured["headers"] = {k.lower(): v for k, v in request.headers.items()} + return FakeConn() + + original = urllib.request.urlopen + urllib.request.urlopen = fake_urlopen + try: + mod._post(ENDPOINT, payload, timeout=5) + finally: + urllib.request.urlopen = original + return captured + + def test_large_payloads_are_gzipped(self): + payload = json.dumps({"filler": "x" * 20000}).encode("utf-8") + captured = self._captured_request(payload) + assert captured["data"][:2] == b"\x1f\x8b", "gzip magic bytes" + assert captured["headers"].get("Content-encoding".lower()) == "gzip" + + def test_gzip_actually_shrinks_the_body(self): + payload = json.dumps({"filler": "x" * 20000}).encode("utf-8") + captured = self._captured_request(payload) + assert len(captured["data"]) < len(payload) + + def test_gzip_round_trips_to_the_original_bytes(self): + import gzip as gziplib + + payload = json.dumps({"filler": "x" * 20000}).encode("utf-8") + captured = self._captured_request(payload) + assert gziplib.decompress(captured["data"]) == payload + + def test_small_payloads_are_sent_plain(self): + payload = b'{"small": true}' + captured = self._captured_request(payload) + assert captured["data"] == payload + assert "content-encoding" not in captured["headers"]