feat(telemetry): send exported packages to the ingest service
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.
This commit is contained in:
@@ -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
|
||||
@@ -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"]
|
||||
Reference in New Issue
Block a user