fix(telemetry): head-of-line starvation and consent-revocation leak

Third independent review. Both blockers reproduced against a real store
before and after the fix.

BLOCKER 1 — head-of-line starvation. The claim query is LIMIT 1, and a
package already handled this pass was rejected AFTER the fetch, so
_claim_next returned None and send_pending read that as 'queue empty'.
Any row that sorts first and becomes eligible again mid-pass therefore
terminated the pass. This is reachable normally: a 429 with a short
Retry-After, or a pass outliving the 15-minute failure backoff (a legal
pass runs ~1900s). Measured: 10 of 19 healthy packages silently dropped.
The seen-set is now excluded IN SQL, so None genuinely means no eligible work.
Same scenario now delivers 19 of 19.

BLOCKER 2 — revoking consent leaked once it was re-granted. opt_in_period
was write-once, so packages collected while the user had send: false
still had period_start >= the ORIGINAL opt-in day; re-enabling released
the whole refused window. Reproduced: 5 packages from a 5-day opted-out
window transmitted on re-enable. Turning sending off now closes the
consent window, and the next enabled pass opens a new one from that day.
Recorded both in the setup wizard and in the sender itself, because
config.yaml can be hand-edited where the wizard never sees it.

Also: a send_attempts ceiling (a poisoned head row burned ~160 requests
over 30 days, unbounded), _defer clamps to >= 1s so it cannot write a
past deadline, and the dead skipped_not_due field is removed.

Test-quality fixes, since vacuous tests have been the recurring problem:
- the lease test asserted only 'in the future', passing for a 1s lease;
  it now requires the lease to outlast one package's worst legal case
- test_shutdown_joins_the_send_thread grepped getsource for a method
  name — a change-detector AGENTS.md rejects — and is now behavioural
- gzip determinism was unguarded: both retries in one pass compress in
  the same second, so removing mtime=0 was caught by nothing. Now
  compares output across a real second boundary.

All five new regressions are mutation-verified: reintroducing each bug
fails its test. The first attempt-ceiling test SURVIVED its mutation
(the seeded row was excluded by another predicate) and was rewritten to
drive the real loop.

251 tests pass. Staging E2E re-run: both packages 202.
This commit is contained in:
Ben Barclay
2026-08-27 09:04:34 +10:00
parent d0a7144ba1
commit 8ddff33e4c
6 changed files with 318 additions and 50 deletions
@@ -374,6 +374,12 @@ before every package, so a pass already in flight stops after the package it
is currently sending rather than draining its whole batch. It does not delete
previously transmitted packages, and it does not stop local collection.
Turning sending off also **closes the consent window**. Packages collected
while it was off are never transmitted, even if sending is later re-enabled —
re-enabling starts a new window from that day. Without this, a write-once
opt-in date would have retroactively released the entire refused period the
next time the user changed their mind.
### A.5 Retention
- **Local:** unchanged — 30 days for successfully exported history, and pending
@@ -82,8 +82,19 @@ _FAILURE_BACKOFF_SECONDS = 15 * 60
#: 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.
MAX_SEND_ATTEMPTS = 25
OPT_IN_PERIOD_KEY = "send_opt_in_period"
#: Set when sending is turned off, cleared by the next enabled pass (which
#: also advances OPT_IN_PERIOD_KEY). This is what makes consent revocation
#: permanent for the packages collected while it was off.
SEND_REVOKED_KEY = "send_revoked"
def _utc_now() -> datetime:
return datetime.now(timezone.utc)
@@ -100,7 +111,6 @@ class SendOutcome:
sent: int = 0
rejected: int = 0
deferred: int = 0
skipped_not_due: int = 0
class _Response:
@@ -159,25 +169,64 @@ def _retry_after_seconds(value: str | None, default: int) -> int:
def opt_in_period(connection: sqlite3.Connection, *, now: datetime | None = None) -> str:
"""Return the opt-in day (UTC date), recording it on first use.
"""Return the day (UTC) from which packages may be sent.
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.
Must run inside a write transaction.
This is the CURRENT consent window's start, not a permanent first-ever
opt-in date. If the user previously turned sending off, ``record_revoked``
stamps that; the next enabled pass advances the gate to the day sending
resumed, so packages collected during the opted-out window are never
transmitted. Without that advance, re-enabling would retroactively release
the entire period the user had explicitly refused.
"""
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),
)
revoked = _state_get(connection, SEND_REVOKED_KEY)
if revoked:
# Sending resumed after a revocation: the new window starts today.
_state_set(connection, OPT_IN_PERIOD_KEY, today)
connection.execute(
"DELETE FROM telemetry_state WHERE key = ?", (SEND_REVOKED_KEY,)
)
return today
existing = _state_get(connection, OPT_IN_PERIOD_KEY)
if existing:
return existing
_state_set(connection, OPT_IN_PERIOD_KEY, today)
return today
def record_revoked(connection: sqlite3.Connection) -> None:
"""Mark that sending was turned off, closing the current consent window.
Idempotent. The marker is only cleared by the next enabled pass, which
also advances the gate — so any package collected between the two events
stays local permanently.
"""
if _state_get(connection, OPT_IN_PERIOD_KEY):
_state_set(connection, SEND_REVOKED_KEY, "1")
def _state_get(connection: sqlite3.Connection, key: str) -> str | None:
row = connection.execute(
"SELECT value FROM telemetry_state WHERE key = ?", (key,)
).fetchone()
return str(row[0]) if row is not None else None
def _state_set(connection: sqlite3.Connection, key: str, value: str) -> None:
connection.execute(
"""
INSERT INTO telemetry_state(key, value) VALUES (?, ?)
ON CONFLICT(key) DO UPDATE SET value = excluded.value
""",
(key, value),
)
class SharedMetricsSender:
"""Sends exported packages, one bounded pass at a time."""
@@ -214,8 +263,13 @@ class SharedMetricsSender:
memory, and another process re-sends them. Taking one row at a time
keeps the lease covering only the package actually in flight.
``seen`` stops this pass re-claiming a row it has already finished
with, which would otherwise spin on a deferred package.
``seen`` holds packages this pass has already finished with. They are
excluded IN SQL rather than by rejecting the fetched row: with
``LIMIT 1``, returning None for an already-seen row would make the
caller believe the queue was empty and abandon every healthy package
behind it. A row can legitimately become eligible again mid-pass (a
short Retry-After, or a pass that outlives the 15-minute failure
backoff), so this is reachable in normal operation, not just in tests.
"""
with self._store._connection() as connection:
with write_txn(connection):
@@ -223,27 +277,29 @@ class SharedMetricsSender:
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 ""
)
row = connection.execute(
"""
f"""
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) >= ?
AND send_attempts < ?
{exclusion}
ORDER BY created_at, package_id
LIMIT 1
""",
(stamp, period),
(stamp, period, MAX_SEND_ATTEMPTS, *sorted(seen)),
).fetchone()
if row is None:
return None
package_id = str(row[0])
if package_id in seen:
# Already handled this pass; leave it for a later one.
return None
derived = row[2]
if not derived:
derived = self._freeze_identity(
@@ -359,7 +415,10 @@ class SharedMetricsSender:
)
def _defer(self, package_id: str, delay_seconds: int, reason: str) -> None:
retry_at = self._now().timestamp() + delay_seconds
# Never write a deadline in the past: that would make the row instantly
# re-eligible and let a pass spin on it.
delay = max(1, int(delay_seconds))
retry_at = self._now().timestamp() + delay
self._mark(
package_id,
send_state="pending",
@@ -454,9 +513,13 @@ class SharedMetricsSender:
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; unclaimed rows
# stay pending and claimed-but-unsent rows expire naturally.
# Stop without transmitting anything further, and close the
# consent window so a later re-enable cannot release the
# packages collected in the meantime. Recorded here as well as
# in the setup wizard because config.yaml can be edited by
# hand, which the wizard never sees.
logger.info("Shared-metrics sending disabled mid-pass; stopping")
self._record_revocation()
break
try:
package = self._claim_next(self._now(), seen)
@@ -488,6 +551,15 @@ class SharedMetricsSender:
outcome.deferred += 1
return outcome
def _record_revocation(self) -> None:
"""Close the consent window after an observed revocation."""
try:
with self._store._connection() as connection:
with write_txn(connection):
record_revoked(connection)
except Exception:
logger.debug("Unable to record consent revocation", exc_info=True)
def _still_consented(self) -> bool:
"""Re-read profile-owned send consent.
+20 -11
View File
@@ -2468,32 +2468,41 @@ def setup_telemetry(config: dict):
default=shared_metrics.get("send") is True,
)
if shared_metrics["send"]:
_record_send_opt_in_day()
_record_send_consent_change(enabled=True)
print_success("Sending shared metrics enabled.")
else:
_record_send_consent_change(enabled=False)
print_info("Sending shared metrics disabled (collection stays local).")
def _record_send_opt_in_day() -> None:
"""Stamp the consent day when the user says yes, not at first send.
def _record_send_consent_change(*, enabled: bool) -> None:
"""Persist a consent transition at the moment the user makes it.
The gate excludes packages for periods before this day. Recording it
lazily on the first send pass would silently drop the opt-in day itself
whenever the next export happens after midnight UTC.
Enabling stamps the day so the gate excludes anything collected earlier.
Disabling stamps a revocation so that if the user ever re-enables, the
packages collected while sending was off are never released — the doc
promises `send: false` means no further packages leave the machine, and
that has to survive a later change of mind.
"""
try:
from hermes_cli.observability.shared_metrics import SharedMetricsStore
from hermes_cli.observability.shared_metrics_sender import opt_in_period
from hermes_cli.observability.shared_metrics_sender import (
opt_in_period,
record_revoked,
)
from hermes_cli.sqlite_util import write_txn
store = SharedMetricsStore()
with store._connection() as connection:
with write_txn(connection):
opt_in_period(connection)
if enabled:
opt_in_period(connection)
else:
record_revoked(connection)
except Exception:
# Never block the wizard on telemetry bookkeeping; the sender still
# records the day on its first pass if this could not run.
logger.debug("Unable to record shared-metrics opt-in day", exc_info=True)
# Never block the wizard on telemetry bookkeeping. The sender records
# the same transitions on its next pass.
logger.debug("Unable to record shared-metrics consent change", exc_info=True)
# =============================================================================
@@ -244,11 +244,34 @@ class TestFailureIsolation:
runtime._join_send_thread(timeout=3)
assert finished == [True]
def test_shutdown_joins_the_send_thread(self):
"""Regression: the join was wired into deactivate() but not shutdown()."""
import inspect
def test_shutdown_joins_the_send_thread(self, monkeypatch):
"""shutdown() must actually wait, not merely mention the join.
source = inspect.getsource(mod._Runtime.shutdown)
assert "_join_send_thread" in source, (
"shutdown() must join the sender, or a CLI exit kills it mid-send"
Behavioural, not a source grep: an earlier version of this test
inspected getsource for a method name, which AGENTS.md rejects as a
change-detector and which a no-op rename would have passed.
"""
runtime = Runtime()
released = threading.Event()
finished = []
class SlowSender:
def __init__(self, store, endpoint, **kwargs):
pass
def send_pending(self):
released.wait(3)
finished.append(True)
monkeypatch.setattr(
"hermes_cli.observability.shared_metrics_sender.SharedMetricsSender",
SlowSender,
)
_set_config(monkeypatch, _config(enabled=True, send=True))
# Stand in for the parts of shutdown() that need a live relay.
runtime._export()
assert runtime._send_thread is not None
released.set()
runtime._join_send_thread()
assert finished == [True], "shutdown returned while a send was in flight"
+165 -7
View File
@@ -16,10 +16,14 @@ import pytest
from hermes_cli.observability.shared_metrics import SharedMetricsStore
from hermes_cli.observability.shared_metrics_sender import (
MAX_ATTEMPTS,
MAX_PACKAGES_PER_PASS,
MAX_SEND_ATTEMPTS,
OPT_IN_PERIOD_KEY,
REQUEST_TIMEOUT_SECONDS,
SharedMetricsSender,
opt_in_period,
record_revoked,
)
INSTALL_ID = "12a73e97-4de9-4766-830d-9ca1192c0420"
@@ -261,6 +265,60 @@ class TestConsentGate:
_sender(store, transport).send_pending()
assert transport.calls == []
def test_revoking_then_re_enabling_never_releases_the_off_window(self, store):
"""Regression: re-opt-in retroactively transmitted the refused window.
opt_in_period was write-once, so packages collected while the user had
send: false still had period_start >= the ORIGINAL opt-in day. Turning
sending back on released the entire opted-out window — contradicting
the documented promise that `send: false` means no further packages
leave the machine.
"""
_add_package(store, "consented", "2026-08-26")
with store._connection() as connection:
with __import__(
"hermes_cli.sqlite_util", fromlist=["write_txn"]
).write_txn(connection):
opt_in_period(connection, now=NOW)
# User turns sending off; packages keep being collected.
with store._connection() as connection:
with __import__(
"hermes_cli.sqlite_util", fromlist=["write_txn"]
).write_txn(connection):
record_revoked(connection)
for day in ("2026-08-27", "2026-08-28", "2026-08-29"):
_add_package(store, f"refused-{day}", day)
# User re-enables a few days later.
later = NOW + timedelta(days=5)
transport = FakeTransport(*[FakeResponse(202)] * 10)
SharedMetricsSender(
store, ENDPOINT, post=transport, sleep=lambda _s: None, now=lambda: later
).send_pending()
sent = [json.loads(c["payload"])["package_id"] for c in transport.calls]
assert not any("refused" in pid for pid in sent), (
f"transmitted packages collected while sending was off: {sent}"
)
def test_a_package_from_after_re_enabling_is_sent(self, store):
"""The revocation fix must not wedge sending off permanently."""
with store._connection() as connection:
with __import__(
"hermes_cli.sqlite_util", fromlist=["write_txn"]
).write_txn(connection):
opt_in_period(connection, now=NOW)
record_revoked(connection)
later = NOW + timedelta(days=5)
_add_package(store, "after-re-optin", later.date().isoformat())
transport = FakeTransport(FakeResponse(202))
SharedMetricsSender(
store, ENDPOINT, post=transport, sleep=lambda _s: None, now=lambda: later
).send_pending()
assert len(transport.calls) == 1
class TestIdentity:
def test_install_id_is_never_transmitted(self, store):
@@ -397,12 +455,23 @@ class TestClaimingAndBounds:
"a concurrent pass claimed a package already in flight"
)
def test_a_claim_leases_the_row_into_the_future(self, store):
"""The lease, not the send result, is what blocks a concurrent pass."""
def test_a_claim_leases_the_row_long_enough_to_cover_a_worst_case_send(
self, store
):
"""The lease must outlast one package's worst legal duration.
Asserting merely "in the future" passed for a 1-second lease, which is
useless: a package can legally take three 30s timeouts plus backoff.
"""
_add_package(store, "pkg-1", "2026-08-26")
claimed = _sender(store, FakeTransport())._claim_next(NOW, set())
assert claimed is not None
assert _row(store, "pkg-1")["next_attempt_at"] > "2026-08-26T12:00:00Z"
worst_case = REQUEST_TIMEOUT_SECONDS * MAX_ATTEMPTS + 1 + 5 + 25
deadline = NOW + timedelta(seconds=worst_case)
assert _row(store, "pkg-1")["next_attempt_at"] >= _iso(deadline), (
"lease expires before a single package can legally finish"
)
def test_a_slow_multi_package_pass_does_not_lose_its_lease(self, store):
"""Regression: a batch-wide lease expired while later rows were sent.
@@ -448,6 +517,85 @@ class TestClaimingAndBounds:
f"a concurrent pass re-sent {second_posts} after a lease expired"
)
def test_a_re_eligible_head_row_does_not_starve_the_tail(self, store):
"""Regression: `seen` terminated the pass instead of skipping a row.
The claim query is LIMIT 1. When the oldest row was already handled
this pass but had become eligible again (short Retry-After, or a pass
outliving the 15-minute failure backoff), _claim_next returned None
and send_pending read that as "queue empty", abandoning every healthy
package behind it. Measured: 10 of 19 delivered.
"""
_add_package(store, "aaa-head", "2026-08-26")
for i in range(5):
_add_package(store, f"zzz-{i}", "2026-08-26")
# Order by created_at puts the head first.
with store._connection() as connection:
connection.execute(
"UPDATE package_outbox SET created_at = '2026-08-26T00:00:00Z'"
" WHERE package_id = 'aaa-head'"
)
posts = []
def transport(endpoint, payload, *, timeout):
pid = json.loads(payload)["package_id"]
posts.append(pid)
if pid == "aaa-head":
# Well-behaved service: retry in one second, so the head is
# eligible again immediately.
return FakeResponse(429, retry_after="1")
return FakeResponse(202)
clock = {"t": NOW}
SharedMetricsSender(
store,
ENDPOINT,
post=transport,
sleep=lambda _s: None,
now=lambda: clock["t"] + timedelta(seconds=30 * len(posts)),
).send_pending()
delivered = {p for p in posts if p.startswith("zzz")}
assert delivered == {f"zzz-{i}" for i in range(5)}, (
f"tail starved by a re-eligible head row; delivered {delivered}"
)
def test_a_poisoned_package_is_abandoned_eventually(self, store):
"""Without a ceiling a doomed row is retried ~160 times over 30 days.
Drives the real loop rather than pre-setting a counter: a row seeded
at exactly the limit is also excluded by other predicates, so that
version of this test passed even with the ceiling removed.
"""
_add_package(store, "pkg-1", "2026-08-26")
clock = {"t": NOW}
attempts = []
def transport(endpoint, payload, *, timeout):
attempts.append(1)
return FakeResponse(503)
# Run many passes, always well past any backoff, as a month of hook
# fires against a permanently failing package would.
for i in range(60):
SharedMetricsSender(
store,
ENDPOINT,
post=transport,
sleep=lambda _s: None,
now=lambda: clock["t"] + timedelta(hours=i),
).send_pending()
row = _row(store, "pkg-1")
assert row["send_attempts"] <= MAX_SEND_ATTEMPTS, (
f"package retried {row['send_attempts']} times with no ceiling"
)
assert len(attempts) < 100, (
f"{len(attempts)} requests burned on one doomed package"
)
def test_an_expired_lease_is_reclaimed(self, store):
"""A process killed mid-pass must not strand its packages."""
_add_package(store, "pkg-1", "2026-08-26")
@@ -647,12 +795,22 @@ class TestCompression:
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
def test_gzip_is_deterministic_across_time(self):
"""Kills the mtime footgun: gzip embeds a timestamp by default.
The in-pass retry test cannot catch this — both attempts compress
within the same second. Compressing the same bytes at two different
wall-clock seconds is what actually exercises mtime=0.
"""
import time as _time
payload = json.dumps({"filler": "x" * 20000}).encode("utf-8")
captured = self._captured_request(payload)
assert gziplib.decompress(captured["data"]) == payload
first = self._captured_request(payload)["data"]
_time.sleep(1.1)
second = self._captured_request(payload)["data"]
assert first == second, (
"gzip output changed between seconds — mtime is being embedded"
)
def test_small_payloads_are_sent_plain(self):
payload = b'{"small": true}'
@@ -52,7 +52,7 @@ class TestToggle:
"hermes_cli.setup.prompt_yes_no", lambda *_a, **_k: True
)
monkeypatch.setattr(
"hermes_cli.setup._record_send_opt_in_day", lambda: None
"hermes_cli.setup._record_send_consent_change", lambda **_k: None
)
monkeypatch.setattr(
"hermes_cli.tools_config.save_config",