diff --git a/cli-config.yaml.example b/cli-config.yaml.example index c4025dc2b3..50f927c368 100644 --- a/cli-config.yaml.example +++ b/cli-config.yaml.example @@ -1870,15 +1870,31 @@ display: # ============================================================================= # Shared metrics are disabled by default. When enabled, Hermes writes only # allowlisted aggregate counters and immutable JSON -# packages under $HERMES_HOME/telemetry/shared_metrics; it does not upload them. +# packages under $HERMES_HOME/telemetry/shared_metrics. # Packages include a random profile-scoped ID that stays stable until this # directory is deleted. It is not derived from hardware, account, or host data. # Successfully exported local history is retained for 30 days; pending deltas # are retained until they can be exported. # This profile-owned choice is not overridden by managed-scope configuration. +# +# Nothing is uploaded unless you also set `send: true`. That is a separate +# opt-in and requires `enabled`; it never turns collection on by itself. +# When sending is on: +# * only packages whose entire collection period falls within one +# continuous recorded consent window are ever sent. Consent windows open +# when you enable sending and close when you disable it, so data +# collected before you opted in — or during any gap between opt-ins — +# stays on this machine; +# * each package carries the profile-scoped ID as-is. It is a random UUID +# with no hardware, account, or host-derived content, and deleting the +# shared-metrics directory resets it. +# See docs/observability/relay-shared-metrics.md (Appendix A) for the full +# consent, identity, retention, and deletion decisions. telemetry: shared_metrics: enabled: false + send: false + # endpoint: https://telemetry.nousresearch.com/v1/telemetry # ============================================================================= diff --git a/docs/observability/relay-shared-metrics.md b/docs/observability/relay-shared-metrics.md index 146590dc99..a736bf6edd 100644 --- a/docs/observability/relay-shared-metrics.md +++ b/docs/observability/relay-shared-metrics.md @@ -33,8 +33,12 @@ than downloading a different implementation. When Relay managed execution is active, the provider request and response pass through that native module in the Hermes process so configured interceptors can operate on the real call. This is separate from the shared-metrics data -contract. Shared-metrics mode installs no network exporter and its subscriber -accepts only the versioned, allowlisted projection described below. Enabling a +contract. Shared-metrics mode installs no rich-observability network exporter, +and its subscriber +accepts only the versioned, allowlisted projection described below. The +opt-in package sender described in Appendix A is the only outbound path, it +transmits nothing unless the user enables both `enabled` and `send`, and it +sends whole packages rather than live spans. Enabling a separately configured rich-observability or dynamic plugin can create a different data path and requires its own policy review. @@ -226,10 +230,18 @@ packages from that profile and can therefore link those local packages. Deleting `$HERMES_HOME/telemetry/shared_metrics` resets the identifier together with all aggregates and package files. -This slice has no remote-delivery path. A future remote exporter must not reuse -the persistent local identifier by default. It requires a separate product and -privacy decision covering consent, identity scope, rotation or keyed -pseudonymization, reset behavior, retention, and deletion. +Remote delivery is opt-in and off by default. Reusing the persistent local +identifier remotely required a separate product and privacy decision covering +consent, identity scope, reset behavior, retention, and deletion — that +decision has been made. + +> Those decisions are recorded in +> [Appendix A](#appendix-a-remote-exporter-decisions-phase-2), and the exporter +> implementing them has shipped. Collection alone still transmits nothing: the +> sender runs only when `telemetry.shared_metrics.send` is also true. Each +> transmitted package carries the stable `install_id` as-is (product decision, +> 2026-08-27 — see A.2 for the record, including the superseded +> HMAC-pseudonym design). The install identity is scoped to one `HERMES_HOME`. To reset it, stop Hermes processes and remove `$HERMES_HOME/telemetry/shared_metrics`. This deliberately @@ -257,3 +269,210 @@ verifies model, provider, task, tool, and skill counters in SQLite, validates all exported delta packages against the closed schema, verifies the pseudonymous client-active counter, and checks that prompt, response, tool-call ID, tool-result, and skill-name canaries are absent from the packages. + +## Appendix A: Remote Exporter Decisions (Phase 2) + +Status: **implemented.** This appendix answers the product and +privacy questions that "Current Slices" defers to a future remote exporter. It +records what was decided and why, so the reasoning survives the implementation. + +Sending is off by default and requires both `telemetry.shared_metrics.enabled` +and `telemetry.shared_metrics.send`. + +The exporter sends the package files already written under +`$HERMES_HOME/telemetry/shared_metrics/outbox/` to the Hermes telemetry ingest +service. That service validates only the envelope (`schema_version` plus a UUID +`package_id`) and stores the body verbatim in S3. + +### A.1 Consent + +Transmission is a **separate opt-in** from collection, under a new config key: + +```yaml +telemetry: + shared_metrics: + enabled: false # collect locally + send: false # NEW: transmit to the Nous telemetry service +``` + +- `send` defaults to **false**. Collection alone never transmits. +- `send` requires `enabled`. It does **not** imply it: a transmission flag must + not silently switch on collection. `send: true` with `enabled: false` warns + and does nothing. +- Like `enabled`, `send` is profile-owned and is not overridden by + managed-scope configuration. + +**A package is only sent when its whole period falls inside a recorded +consent window.** Consent is stored as explicit intervals in the shared- +metrics SQLite store (`send_consent_windows`): a window opens when `send: +true` is first observed, is confirmed forward by every later observation, +and closes — at the last *confirmed* moment, never at the wall clock — when +`send: false` is observed. A single reconciler derives this table from the +config on every process start, so wizard changes, hand-edits to +`config.yaml`, and mid-pass revocations all take the same path, and no +transition can be missed by any of them. + +Any package whose period predates the first window, falls between windows, +or runs past the newest confirmed moment is excluded — the gate fails +closed. A fresh package therefore waits at most one process start after its +period completes before becoming eligible. + +The gate is on the **period**, not on the package's creation time. One period +is split across several packages created on different days: a day's first +package is written that day, and a tail package for the same period typically +follows the next day. Gating on creation time would send a period's tail while +dropping its head, reporting a **silently undercounted** day. Gating on the +period keeps consent forward-only and every transmitted period complete. + +Local history can be up to 30 days old, and that data was collected under a +promise that nothing is uploaded. Honouring consent forward-only costs at most +30 days of backlog we never had permission to send. + +### A.2 Identity scope — the stable install_id is transmitted as-is + +**Decision record.** The original design of this exporter (and revisions 1–8 +of this appendix) transmitted a keyed pseudonym instead of the identifier: +`HMAC-SHA256(key = locally-held rotating salt, message = install_id)`, with +the salt rotating every 30 days. On **2026-08-27**, before the feature +shipped (zero consented users, zero production transmissions), the product +owner decided the analytical need is a **stable cross-window identity** — +retention curves, longitudinal install behaviour — which rotation by design +destroys. The pseudonymization layer was removed in full rather than +weakened in place. + +What is transmitted now: + +- Each package carries `install_id` verbatim: the persistent, profile-scoped + random UUID described above. +- It is generated locally (`uuid4`), contains no hardware, account, user, or + machine-derived information, and identifies a *profile*, not a person. +- It is stable until the user deletes the shared-metrics directory, which + regenerates it (see A.4). + +Consequences stated plainly rather than papered over: + +- Packages from one profile correlate **indefinitely**, not per-window. + Long-term linkability of one install's daily envelope sequence is now the + designed behaviour, not a residue. +- The A.3 residue analysis of the old design (stable `resource` tuple + + contiguous periods bridging rotation windows) is moot — there is no window + boundary left to bridge. +- The setup wizard's consent language states this identity model explicitly; + it was updated in the same change that removed the derivation, so no + consent was ever collected under the old wording in any shipped build. + +**Byte-identical resends still hold.** The transmitted id is recorded on the +row (`sent_install_id`) when the package is first prepared, and the wire body +is always rebuilt from that recorded value, so a retry rebuilds identical +bytes. The contract requires this: resending a `package_id` with different +content is undefined behaviour. (With a stable id the recorded copy is no +longer load-bearing against rotation — it remains as the audit column and as +cheap insurance against any future change to identity semantics.) + +### A.3 Rotation — removed (decision record) + +Salt rotation was deleted together with the derivation (product decision, +2026-08-27). This section is retained as a record of what the earlier design +did and why the removal was accepted: + +- Rotation existed to bound long-term linkability: one identity per 30-day + window, unrelated identities across windows. +- The documented residue (see git history for the full analysis): the + envelope's stable, low-entropy `resource` tuple plus contiguous daily + periods could plausibly bridge windows for rare configurations anyway, so + the boundary was a cost-raiser, not a wall. +- The product need that killed it: cross-window continuity is precisely what + retention analysis requires. A boundary that mostly inconveniences honest + analysis while only raising costs for a determined correlator was judged + the wrong trade once stable identity became a requirement. + +There is no salt in the store, no rotation schedule, and no derived +identifier anywhere in the pipeline. + +### A.4 Reset behavior + +Removing `$HERMES_HOME/telemetry/shared_metrics` still resets local identity, +aggregates, and package files, exactly as documented above. Two honest +qualifications now apply: + +- Reset regenerates `install_id`, so subsequent packages transmit a **new** + identity. Local reset does give a new remote identity. +- Reset **cannot unsend**. Packages already transmitted remain in the ingest + service's storage under the identifier they were sent with. There is no + read-back or delete API in the v1 contract. + +Setting `send: false` stops transmission immediately: consent is re-read +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** — at the last moment +consent was actually observed, not at the wall clock. Packages whose periods +fall between one window and the next are never transmitted, even if sending +is later re-enabled, and this holds for any number of on/off cycles, across +hand-edits with no process running, and under a clock that jumps in either +direction (window opens are clamped above every timestamp already in the +store; observation marks advance by a bounded step per call, so one glitched +forward sample cannot drag the confirmation horizon years ahead; a close +never lands after the closing observation's own clock). +Unlike the earlier single moving opt-in date, closing and reopening does NOT +discard the still-undelivered backlog from a previous consented window — +those packages stay inside their own interval and remain eligible. + +One deliberate upgrade-path consequence: packages exported under the +pre-interval consent model (before `send_consent_windows` existed) predate +the first recorded window and are therefore never transmitted after an +upgrade. This is the fail-closed direction — re-importing the old moving +day-stamp to release them would re-import the semantics five review rounds +showed to be unsound — and it costs at most the undelivered backlog, never +collected data. + +### A.5 Retention + +- **Local:** unchanged — 30 days for successfully exported history, and pending + deltas are kept until exported. Send state does **not** extend local + retention: a package that could never be sent is still pruned at 30 days. + Unbounded local growth against a permanently unreachable endpoint is a worse + failure than losing metrics from an install that has been broken for a month. +- **Remote:** raw packages are retained in S3 without expiry in production and + for 30 days in staging. + +### A.6 Deletion + +There is no remote deletion path in the v1 contract, and this appendix does not +invent one. What a user can do: + +| Action | Effect | +|---|---| +| `send: false` | No further packages leave the machine | +| `enabled: false` | Collection stops; existing local state remains | +| Remove `.../shared_metrics` | Local identity, aggregates, and files reset; future sends use a new install_id | +| Delete already-sent data | Not self-service — requires an operator acting on the S3 bucket | + +If a deletion-on-request obligation is ever taken on, the lookup path is now +direct: the user's `install_id` (readable from their local store) is the key +their data is stored under. Building the service-side delete API remains a +new product decision, not an implementation detail. + +### A.7 What the outbox directory is + +Recorded because it was misread once during Phase 2 planning, in a way that +would have deleted user data. + +The directory is **local history, not a send-queue**. `package_outbox` is the +SQLite table; its `exported_at` column means "written to disk", not "sent". +Files are immutable and pruned **by age alone**. + +The ingest contract says senders should delete a package from their outbox on +`202`. **The exporter does not do this.** Deleting on acknowledgement would +repurpose the user's 30-day local history as a transmission queue and destroy +state they were promised. Send state lives in new columns on the +`package_outbox` table instead; the files are untouched by transmission. + +### A.8 Scope note + +The `install_id` field inside the package body is transmitted as the +generator wrote it (rewritten from the row's frozen `sent_install_id`, which +records the same value). No other payload field changes, nothing is added, +and the service treats the whole body as opaque. Payload schema evolution +therefore stays a sender-side concern, as before. diff --git a/hermes_cli/config_defaults.py b/hermes_cli/config_defaults.py index 9064f4a622..2e64ec06e6 100644 --- a/hermes_cli/config_defaults.py +++ b/hermes_cli/config_defaults.py @@ -3542,11 +3542,28 @@ DEFAULT_CONFIG = { "profile_build": "ask", }, - # Privacy-safe aggregate metrics written only to this profile's local - # telemetry directory. Collection is opt-in and no remote sink exists. + # Privacy-safe aggregate metrics written to this profile's local telemetry + # directory. Collection is opt-in (``enabled``). Transmission to the Nous + # telemetry service is a SEPARATE opt-in (``send``) and is off by default; + # see docs/observability/relay-shared-metrics.md, Appendix A, for the + # consent, identity, rotation, retention, and deletion decisions. "telemetry": { "shared_metrics": { "enabled": False, + # Transmit exported packages to the Nous telemetry service. + # Requires ``enabled``: it never switches collection on by itself, + # and ``send`` without ``enabled`` is logged as an error rather + # than silently doing nothing. A package is only sent when its + # whole period falls inside a recorded consent window, so data + # collected before consent — or while it was withdrawn — stays + # local. + "send": False, + # Ingest endpoint. Production by default; override for staging or + # a local test server. Deliberately NOT overridable by an + # environment variable: that would let an inherited value silently + # redirect telemetry a user consented to send to Nous. Non-HTTPS + # is refused unless the host is localhost. + "endpoint": "https://telemetry.nousresearch.com/v1/telemetry", }, }, diff --git a/hermes_cli/observability/relay_shared_metrics.py b/hermes_cli/observability/relay_shared_metrics.py index 2ab88f51c3..5a97c8a18d 100644 --- a/hermes_cli/observability/relay_shared_metrics.py +++ b/hermes_cli/observability/relay_shared_metrics.py @@ -132,6 +132,9 @@ class _Runtime: self._sessions: dict[str, _MetricsSession] = {} self._task_creation_lock = threading.RLock() self._task_sessions_lock = threading.RLock() + # Guards the opt-in send pass: at most one in flight per process. + self._send_lock = threading.RLock() + self._send_thread: threading.Thread | None = None self._task_sessions: dict[tuple[str, str], _MetricsSession] = {} self._turn_sessions: dict[tuple[str, str], _MetricsSession] = {} self._subscriber_name = f"{SUBSCRIBER_NAME}.{self.host.runtime_id}" @@ -668,6 +671,12 @@ class _Runtime: self._safe(self.relay.subscribers.deregister, self._subscriber_name) self.host.release_managed_execution(self._subscriber_name) self._registered = False + # The final export above may have started a send. Give it the same + # bounded chance to finish that deactivate() gets — without this a + # short-lived CLI process exits immediately and kills the daemon + # thread mid-request, which is the common case for the one cadence + # this feature has. + self._join_send_thread() try: atexit.unregister(self.shutdown) except Exception: @@ -706,11 +715,29 @@ class _Runtime: with self._task_sessions_lock: self._task_sessions.clear() self._turn_sessions.clear() + self._join_send_thread() try: atexit.unregister(self.shutdown) except Exception: pass + def _join_send_thread(self, timeout: float = 2.0) -> None: + """Give an in-flight send a brief chance to finish at exit. + + Bounded on purpose: the packages stay pending in SQLite and go out on + the next run, so blocking a user's shutdown for a slow network is the + wrong trade. The thread is a daemon, so an unfinished pass dies with + the process rather than holding it open. + """ + with self._send_lock: + thread = self._send_thread + if thread is None or not thread.is_alive(): + return + try: + thread.join(timeout) + except Exception: + logger.debug("Shared-metrics send thread join failed", exc_info=True) + def _session(self, event: dict[str, Any]) -> _MetricsSession | None: session_id = str(event.get("session_id") or "") with self._sessions_lock: @@ -1048,7 +1075,104 @@ class _Runtime: return True def _export(self) -> None: - self._safe(self.subscriber.store.create_and_export_package_if_due) + exported = self._safe(self.subscriber.store.create_and_export_package_if_due) + # Sending is opt-in and must never delay the caller: _export runs on + # finish_task, which is the user's interactive path. Errors inside the + # sender are already swallowed there; the thread is about latency, not + # correctness. + if exported is not None: + self._safe(self._send_exported_packages) + + def _observe_send_consent(self, send_enabled: bool) -> None: + """Reconcile consent windows with the observed config state. + + Thin wrapper over the SINGLE consent writer. The old edge-detection + body (last-seen key, rising/falling branches) is gone: reconciliation + derives the correct window state from what it observes, so there is + no transition to miss and no ordering between callers to get wrong. + + Failures must never break the export hook, but they are logged at + warning rather than debug: silently failing to close a consent window + is a privacy-relevant event, not routine bookkeeping. + """ + try: + from hermes_cli.observability.shared_metrics_sender import ( + reconcile_send_consent, + ) + from hermes_cli.sqlite_util import write_txn + + with self.subscriber.store._connection() as connection: + with write_txn(connection): + reconcile_send_consent(connection, send_enabled) + except Exception: + logger.warning( + "Unable to record a shared-metrics consent transition", + exc_info=True, + ) + + def _send_exported_packages(self) -> None: + from hermes_cli.observability.shared_metrics_send_config import ( + resolve_send_config, + ) + + try: + from hermes_cli.config import read_raw_config_readonly + + config = read_raw_config_readonly() or {} + except Exception: + logger.debug("Unable to read shared-metrics send policy", exc_info=True) + return + + resolved = resolve_send_config(config) + + # Observe the consent EDGE before deciding whether to send. Recording + # revocation inside the send loop (as an earlier fix did) can never + # work: the dominant case is the user turning sending off while no + # pass is running, and then this method returns below without ever + # constructing a sender. The window has to close on the transition, + # not on the next transmission that by definition will not happen. + self._observe_send_consent(resolved.send) + + if not resolved.send: + return + + with self._send_lock: + # One in-flight pass per process. A queued second pass would add + # nothing: the next hook fire picks up whatever is still pending. + if self._send_thread is not None and self._send_thread.is_alive(): + return + thread = threading.Thread( + target=self._run_send_pass, + args=(resolved.endpoint,), + name="hermes-shared-metrics-send", + daemon=True, + ) + self._send_thread = thread + thread.start() + + def _run_send_pass(self, endpoint: str) -> None: + from hermes_cli.observability.shared_metrics_sender import ( + SharedMetricsSender, + ) + + def still_consented() -> bool: + """Re-read consent so revoking `send` stops an in-flight pass.""" + from hermes_cli.config import read_raw_config_readonly + from hermes_cli.observability.shared_metrics_send_config import ( + resolve_send_config, + ) + + resolved = resolve_send_config(read_raw_config_readonly() or {}) + return resolved.send and resolved.endpoint == endpoint + + try: + SharedMetricsSender( + self.subscriber.store, + endpoint, + consent_check=still_consented, + ).send_pending() + except Exception: + logger.warning("Shared-metrics send pass failed", exc_info=True) def _event_metadata(self) -> dict[str, str]: return { @@ -1101,8 +1225,62 @@ def handles_hook(hook_name: str) -> bool: return hook_name in HANDLED_HOOKS and enabled() +_consent_reconcile_done = False + + +def _reconcile_send_consent_once() -> None: + """Reconcile consent windows with config, once per process. + + Runs BEFORE and INDEPENDENT of the collection gate — that placement is + the fix for the round-5 D1 leak, where the only idle-path consent + observer sat behind ``handles_hook()`` and became dead code the moment + ``enabled: false`` was set. A user with collection off still gets their + send-consent windows reconciled here. + + Skipped only when there is no store on disk AND consent is off: with no + store there are no packages, so there is nothing a window could protect, + and creating ``~/.hermes/telemetry`` for every fully-disabled user would + be a behaviour change in the wrong direction. + """ + global _consent_reconcile_done + if _consent_reconcile_done: + return + _consent_reconcile_done = True + try: + from hermes_cli.config import read_raw_config_readonly + from hermes_cli.observability.shared_metrics import SharedMetricsStore + from hermes_cli.observability.shared_metrics_send_config import ( + resolve_send_config, + ) + from hermes_cli.observability.shared_metrics_sender import ( + reconcile_send_consent, + ) + from hermes_cli.sqlite_util import write_txn + from hermes_constants import get_hermes_home + + resolved = resolve_send_config(read_raw_config_readonly() or {}) + # Probe for an existing store WITHOUT constructing one: the + # constructor creates the directory and schema as a side effect, + # which round 6 caught making this skip dead code — every + # fully-disabled user was getting a ~/.hermes/telemetry directory. + default_path = ( + get_hermes_home() / "telemetry" / "shared_metrics" / "metrics.sqlite3" + ) + if not resolved.send and not default_path.exists(): + return + store = SharedMetricsStore() + with store._connection() as connection: + with write_txn(connection): + reconcile_send_consent(connection, resolved.send) + except Exception: + logger.warning( + "Unable to reconcile shared-metrics send consent", exc_info=True + ) + + def observe_lifecycle(hook_name: str, **kwargs: Any) -> None: """Project one Hermes lifecycle event into the core Relay integration.""" + _reconcile_send_consent_once() if not handles_hook(hook_name): return if not relay_runtime.relay_instrumentation_enabled(): diff --git a/hermes_cli/observability/shared_metrics.py b/hermes_cli/observability/shared_metrics.py index fd42b06230..87094922d9 100644 --- a/hermes_cli/observability/shared_metrics.py +++ b/hermes_cli/observability/shared_metrics.py @@ -337,6 +337,8 @@ class SharedMetricsStore: ) """ ) + SharedMetricsStore._add_send_columns(connection) + SharedMetricsStore._add_consent_tables(connection) connection.execute( """ INSERT INTO telemetry_state(key, value) @@ -346,6 +348,98 @@ class SharedMetricsStore: (_STORE_SCHEMA_VERSION,), ) + @staticmethod + def _add_send_columns(connection: sqlite3.Connection) -> None: + """Add transmission bookkeeping to ``package_outbox``, idempotently. + + These columns are ADDITIVE and nullable, and the store schema version + is deliberately NOT bumped. ``_ensure_schema_in_transaction`` raises on + any version it does not recognise and has no forward-compatibility + branch, so bumping would make an older Hermes — a second profile on an + older build, or a rollback — hard-fail against the same database file. + Old readers select named columns and never ``SELECT *``, so extra + columns are invisible to them. + """ + existing = { + str(row["name"]) + for row in connection.execute("PRAGMA table_info(package_outbox)") + } + for column, declaration in ( + # When the 202 was received. NULL = never acknowledged. + ("sent_at", "TEXT"), + # NULL/'pending' = eligible, 'sent' = done, 'rejected' = permanent 400. + ("send_state", "TEXT"), + ("send_attempts", "INTEGER NOT NULL DEFAULT 0"), + # Earliest next attempt; enforces backoff across process restarts. + ("next_attempt_at", "TEXT"), + ("last_error", "TEXT"), + # The identifier actually transmitted, frozen on the first + # attempt so retries stay byte-identical. Since the 2026-08-27 + # product decision this is the stable install_id itself. + # Only the ~36-byte id is stored: the body is recomputed from + # payload_json, whose serialisation is deterministic. + ("sent_install_id", "TEXT"), + # NULL until first claimed; rewritten on every claim. Settlement + # and the pre-POST revalidation are compare-and-set on this, so a + # claimant whose lease lapsed loses authority the moment another + # process reclaims (PR-review finding: without it, a suspended + # sender resuming after a reclaim double-POSTs the package). + ("claim_token", "TEXT"), + ): + if column not in existing: + connection.execute( + f"ALTER TABLE package_outbox ADD COLUMN {column} {declaration}" + ) + + @staticmethod + def _add_consent_tables(connection: sqlite3.Connection) -> None: + """Create the consent-window tables, idempotently. + + Additive like ``_add_send_columns`` — the schema version is + deliberately NOT bumped, and old readers never touch these tables. + + ``send_consent_windows`` records consent as explicit intervals rather + than a moving day-stamp: a window is opened when send consent is + observed, heartbeat-confirmed on every later observation, and closed + at the LAST CONFIRMED moment (never "now") when consent is observed + withdrawn. Consent is asserted only for time that was actually + observed, so unobserved gaps — a hand-edited config with no process + running — fail closed by construction. + + ``consent_marks`` holds two monotonic high-water marks with strictly + separated roles: + + - ``obs``: the latest observation stamp ever seen. Advanced only by + the reconciler. Confirms consent and clamps window closes. + - ``data``: the latest package ``period_end`` ever stored. Advanced + only by the package writer. Clamps window OPENS, so a rolled-back + clock can never open a window underneath packages that already + exist on disk. + + The separation is load-bearing: letting data stamps confirm consent + re-created a refused-window leak (packages stored during an off + window would vouch for it), and letting observation stamps clamp + opens is not enough on its own to stop a rollback sliding a window + under existing refused data. + """ + connection.execute( + """ + CREATE TABLE IF NOT EXISTS send_consent_windows ( + opened_at TEXT NOT NULL, + last_confirmed_at TEXT NOT NULL, + closed_at TEXT + ) + """ + ) + connection.execute( + """ + CREATE TABLE IF NOT EXISTS consent_marks ( + name TEXT PRIMARY KEY CHECK (name IN ('obs', 'data')), + stamp TEXT NOT NULL + ) + """ + ) + @staticmethod def _create_counter_aggregates_table(connection: sqlite3.Connection) -> None: connection.execute( @@ -580,6 +674,16 @@ class SharedMetricsStore: payload["generated_at"], ), ) + # Advance the data high-water mark. This is the ONLY writer of the + # 'data' mark: it clamps consent-window opens so a rolled-back clock + # can never open a window underneath packages that already exist. + connection.execute( + """ + INSERT INTO consent_marks(name, stamp) VALUES ('data', ?) + ON CONFLICT(name) DO UPDATE SET stamp = MAX(stamp, excluded.stamp) + """, + (payload["period_end"],), + ) for row in rows: connection.execute( """ diff --git a/hermes_cli/observability/shared_metrics_send_config.py b/hermes_cli/observability/shared_metrics_send_config.py new file mode 100644 index 0000000000..cb14027593 --- /dev/null +++ b/hermes_cli/observability/shared_metrics_send_config.py @@ -0,0 +1,114 @@ +"""Configuration for shared-metrics transmission. + +Collection (``telemetry.shared_metrics.enabled``) and transmission +(``telemetry.shared_metrics.send``) are separate opt-ins. See +``docs/observability/relay-shared-metrics.md`` Appendix A for the consent, +identity, rotation, retention, and deletion decisions behind this module. +""" + +from __future__ import annotations + +import logging +from dataclasses import dataclass +from urllib.parse import urlparse + +logger = logging.getLogger(__name__) + +#: Production ingest endpoint. Overridable through config only. +#: +#: Deliberately NOT overridable by an environment variable: AGENTS.md reserves +#: HERMES_* env vars for secrets, and a behavioural override here would be a +#: consent hazard — a user who agreed to send metrics to Nous could have them +#: silently redirected to any host by an inherited variable, with nothing +#: visible in their config to show it. Tests and the staging E2E write this +#: key into a throwaway profile instead. +DEFAULT_ENDPOINT = "https://telemetry.nousresearch.com/v1/telemetry" + +_LOCAL_HOSTS = frozenset({"localhost", "127.0.0.1", "::1", "[::1]"}) + +# Module-level latch: the enabled/send mismatch is a static misconfiguration, +# so it is reported once per process instead of on every hook fire. +_warned_send_without_collection = False + + +@dataclass(frozen=True) +class SendConfig: + """Resolved transmission settings.""" + + #: Collection is on. Nothing is packaged or sent without it. + enabled: bool + #: Transmission is on AND permitted (that is, collection is also on). + send: bool + #: Where packages are POSTed. + endpoint: str + + +def _endpoint_is_safe(endpoint: str) -> bool: + """Reject plaintext destinations unless they are loopback. + + Telemetry must not leave a machine in clear text because of a typo in a + config file. Loopback stays allowed so tests can use a local HTTP server. + """ + try: + parsed = urlparse(endpoint) + except ValueError: + return False + if parsed.scheme == "https": + return True + if parsed.scheme == "http": + return (parsed.hostname or "") in _LOCAL_HOSTS + return False + + +def resolve_send_config(config: dict | None) -> SendConfig: + """Resolve transmission settings from config plus the environment. + + Endpoint precedence: config > production default. + + ``send`` is returned as False whenever transmission cannot legitimately + happen, so callers never have to re-check the combination. + """ + global _warned_send_without_collection + + raw = config if isinstance(config, dict) else {} + telemetry = raw.get("telemetry") + telemetry = telemetry if isinstance(telemetry, dict) else {} + shared = telemetry.get("shared_metrics") + shared = shared if isinstance(shared, dict) else {} + + enabled = shared.get("enabled") is True + send_requested = shared.get("send") is True + + if send_requested and not enabled: + # Loud, not silent: the user believes telemetry is being sent, and it + # never will be. Error level, once per process. + if not _warned_send_without_collection: + _warned_send_without_collection = True + logger.error( + "telemetry.shared_metrics.send is true but " + "telemetry.shared_metrics.enabled is false — nothing is " + "collected, so nothing can be sent. Enable collection or " + "turn sending off." + ) + return SendConfig(enabled=False, send=False, endpoint=DEFAULT_ENDPOINT) + + endpoint = shared.get("endpoint") + if not isinstance(endpoint, str) or not endpoint.strip(): + endpoint = DEFAULT_ENDPOINT + endpoint = endpoint.strip() + + if send_requested and not _endpoint_is_safe(endpoint): + logger.error( + "Refusing to send shared metrics to %r: telemetry must use https " + "(or a localhost http endpoint for testing).", + endpoint, + ) + return SendConfig(enabled=enabled, send=False, endpoint=endpoint) + + return SendConfig(enabled=enabled, send=send_requested, endpoint=endpoint) + + +def reset_warning_latch_for_tests() -> None: + """Clear the once-per-process error latch (test support only).""" + global _warned_send_without_collection + _warned_send_without_collection = False diff --git a/hermes_cli/observability/shared_metrics_sender.py b/hermes_cli/observability/shared_metrics_sender.py new file mode 100644 index 0000000000..9418353c9b --- /dev/null +++ b/hermes_cli/observability/shared_metrics_sender.py @@ -0,0 +1,792 @@ +"""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 first consented day. The gate itself is interval containment: +the period must fall entirely inside a recorded consent window +(``send_consent_windows``), maintained by the single ``reconcile_send_consent`` +writer below. +""" + +from __future__ import annotations + +import gzip +import json +import logging +import random +import sqlite3 +import time +import urllib.error +import urllib.request +import uuid +from dataclasses import dataclass +from datetime import datetime, timedelta, timezone + +from hermes_cli.sqlite_util import write_txn + +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 + +#: 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. +_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 +#: 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 + + +def _utc_now() -> datetime: + return datetime.now(timezone.utc) + + +def _isoformat(value: datetime) -> str: + return value.astimezone(timezone.utc).isoformat().replace("+00:00", "Z") + + +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 + ) + + +@dataclass +class SendOutcome: + """What one pass did. Returned for tests and diagnostics.""" + + sent: int = 0 + rejected: int = 0 + deferred: 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: + # 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. + body = gzip.compress(payload, mtime=0) + 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 + + +#: 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, + *, + now: datetime | None = None, +) -> 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. This + replaces the previous edge-detection design, whose three partial + observers (wizard, relay, mid-pass) each covered a different subset of + transitions and repeatedly leaked the transitions between the subsets. + + Timestamp discipline (each rule is load-bearing; see the validation + harness in tests/hermes_cli/test_shared_metrics_consent_windows.py): + + - The 'obs' mark advances to every observation stamp, monotonically — + but by at most ``MAX_OBS_ADVANCE_SECONDS`` per call. Unbounded, the + mark is monotonic in the LEAK direction: one glitched-forward sample + would drag ``last_confirmed_at`` decades ahead, a later close would + stamp that horizon, and the closed window would contain every future + refused period (reproduced in round 6). Bounded, a poisoned sample + costs at most one cap's width, and real time overtakes it. + An open window's ``last_confirmed_at`` follows the mark: consent is + asserted only for time that was actually observed. + - A close is stamped at ``last_confirmed_at`` — never "now" — so an + unobserved gap (hand-edited config, machine off for 90 days) is never + inside a window and fails closed. + - An open clamps to ``max(now, obs, data)``: a rolled-back clock cannot + open a window underneath refused packages already on disk, and cannot + make the new window adjacent to the previous close. + """ + stamp = _isoformat(now or _utc_now()) + raw_stamp = stamp # pre-cap observation time, used to clamp closes + previous_obs = connection.execute( + "SELECT stamp FROM consent_marks WHERE name = 'obs'" + ).fetchone() + if previous_obs is not None: + ceiling = _isoformat( + _parse_stamp(str(previous_obs[0])) + + timedelta(seconds=MAX_OBS_ADVANCE_SECONDS) + ) + stamp = min(stamp, ceiling) + connection.execute( + """ + INSERT INTO consent_marks(name, stamp) VALUES ('obs', ?) + ON CONFLICT(name) DO UPDATE SET stamp = MAX(stamp, excluded.stamp) + """, + (stamp,), + ) + marks = dict( + connection.execute("SELECT name, stamp FROM consent_marks").fetchall() + ) + obs = marks["obs"] # >= stamp; immune to clock rollback + data = marks.get("data") + + open_row = connection.execute( + "SELECT rowid FROM send_consent_windows WHERE closed_at IS NULL" + ).fetchone() + + if send_enabled: + if open_row is None: + opened = max(x for x in (obs, data) if x is not None) + connection.execute( + "INSERT INTO send_consent_windows(opened_at, last_confirmed_at)" + " VALUES (?, ?)", + (opened, opened), + ) + else: + connection.execute( + "UPDATE send_consent_windows" + " SET last_confirmed_at = MAX(last_confirmed_at, ?)" + " WHERE rowid = ?", + (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. + connection.execute( + "UPDATE send_consent_windows" + " SET closed_at = MIN(last_confirmed_at, ?)" + " WHERE rowid = ?", + (raw_stamp, open_row[0]), + ) + + +#: 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 + AND package_outbox.period_end <= + CASE WHEN w.closed_at IS NULL THEN w.last_confirmed_at + ELSE w.closed_at END +)""" + + +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.""" + + def __init__( + self, + store, + endpoint: str, + *, + post=_post, + sleep=time.sleep, + now=_utc_now, + max_attempts: int = MAX_ATTEMPTS, + consent_check=None, + ) -> None: + self._store = store + self._endpoint = endpoint + self._post = post + 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). + self._consent_check = consent_check + + # -- selection --------------------------------------------------------- + + def _claim_next(self, now: datetime, seen: set[str]) -> dict | None: + """Claim exactly ONE package, immediately before it is sent. + + Claiming a whole batch up front does not work: a single shared lease + has to cover the entire pass, and 20 retrying packages can legally run + far longer than any sane lease (three 30s timeouts plus backoff each). + The later rows' leases then expire while this pass still holds them in + memory, and another process re-sends them. Taking one row at a time + keeps the lease covering only the package actually in flight. + + ``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): + 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""" + 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 {CONSENT_GATE_SQL} + AND send_attempts < ? + {exclusion} + ORDER BY created_at, package_id + LIMIT 1 + """, + (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} + + token = str(uuid.uuid4()) + connection.execute( + """ + UPDATE package_outbox + SET send_state = 'pending', + send_attempts = send_attempts + 1, + next_attempt_at = ?, + 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, + } + + def _freeze_identity( + self, + connection: sqlite3.Connection, + package_id: str, + payload_json, + now: datetime, + ) -> str | None: + """Record the transmitted id on the row, or reject an unusable one. + + The stable install_id is transmitted as-is (product decision, + 2026-08-27 — see the doc's A.2). 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 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(). + 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" + + if reason is not None: + logger.warning( + "Shared-metrics package %s cannot be sent (%s)", package_id, reason + ) + connection.execute( + """ + UPDATE package_outbox + SET send_state = 'rejected', last_error = ? + WHERE package_id = ? + """, + (reason, package_id), + ) + return None + + connection.execute( + "UPDATE package_outbox SET sent_install_id = ? WHERE package_id = ?", + (install_id, package_id), + ) + return str(install_id) + + # -- transmission ------------------------------------------------------ + + def _body(self, 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. + """ + payload = json.loads(payload_json) + payload = dict(payload) + payload["install_id"] = transmitted_id + return json.dumps(payload, indent=2, sort_keys=True).encode("utf-8") + + def _mark( + self, + package_id: str, + *, + only_if_pending: bool = True, + token: str | None = None, + **columns, + ) -> 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. + """ + assignments = ", ".join(f"{name} = ?" for name in columns) + 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, + ) + + 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. Renewal closes it by requiring, in ONE statement: + + - the token still matches (nobody reclaimed), AND + - the current lease is UNEXPIRED (this claimant is not stale), AND + - the row is still pending, + + 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 claimant that wakes past its own lease fails the unexpired + condition and yields even though its token was never replaced. + """ + 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( + """ + UPDATE package_outbox + SET next_attempt_at = ? + WHERE package_id = ? + AND claim_token = ? + 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 + except Exception: + # If renewal itself fails, do not transmit on unproven authority. + 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, + ) -> 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 + self._mark( + package_id, + token=token, + 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'. + + 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. The body is byte-identical across + retries by construction, so the residual duplicate is exactly one + redundant copy of identical content; collapsing it fully would need + package_id-keyed dedupe at the ingest service. + """ + package_id = package["package_id"] + token = package.get("claim_token") + body = self._body(package["payload_json"], package["derived"]) + + 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 + # 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, + ) + return "deferred" + 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, token=token + ) + 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}" + if attempt >= self._max_attempts: + self._defer( + package_id, _FAILURE_BACKOFF_SECONDS, reason, token=token + ) + return "deferred" + self._sleep(self._backoff(attempt)) + + self._defer( + package_id, _FAILURE_BACKOFF_SECONDS, "attempts exhausted", token=token + ) + 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. + + 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. + """ + 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. + 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 + ) + break + if package is None: + break + + 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 + else: + outcome.deferred += 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() + ) + except Exception: + 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 configuration: the documentation + promises that setting `send: false` stops transmission immediately, + and a pass can run for minutes. Injected senders (tests, the staging + E2E) opt out by passing consent_check=None. + """ + if self._consent_check is None: + return True + try: + 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, + ) + return False diff --git a/hermes_cli/setup.py b/hermes_cli/setup.py index 4445eb8812..1213eb1158 100644 --- a/hermes_cli/setup.py +++ b/hermes_cli/setup.py @@ -2428,10 +2428,10 @@ def setup_tools(config: dict, first_install: bool = False): def setup_telemetry(config: dict): - """Configure the local, privacy-safe shared-metrics subscriber.""" + """Configure the local shared-metrics subscriber and optional sending.""" print_header("Shared Metrics") print_info("Shared metrics contain only bounded counters and histograms.") - print_info("Packages stay under this Hermes profile and are not uploaded.") + print_info("Collection is local. Sending them to Nous is a separate opt-in.") telemetry = config.get("telemetry") if not isinstance(telemetry, dict): @@ -2447,10 +2447,67 @@ def setup_telemetry(config: dict): "Enable local shared metrics?", default=current, ) - if shared_metrics["enabled"]: - print_success("Local shared metrics enabled.") - else: + if not shared_metrics["enabled"]: print_info("Local shared metrics disabled.") + # Sending cannot outlive collection: leaving send=true here would be a + # configuration that logs an error on every run and never transmits. + if shared_metrics.get("send") is True: + shared_metrics["send"] = False + print_info("Sending shared metrics disabled as well.") + # Turning collection off is also a withdrawal of send consent, and it + # has to close the window like any other. Recorded unconditionally: + # the send key may already be false in config while the consent window + # is still open, and that window must not survive to be reopened. + _record_send_consent_change(enabled=False) + return + + print_success("Local shared metrics enabled.") + print_info("") + print_info("Sending uploads each daily package to the Nous telemetry") + print_info("service. Packages carry your profile-scoped install ID, a") + print_info("stable random UUID that identifies this profile across days") + print_info("(it contains no personal information and is reset by deleting") + print_info("the shared-metrics directory). Only packages whose entire") + print_info("collection period falls inside a recorded consent window are") + print_info("ever sent — data from before you opt in, or from any gap") + print_info("while sending was off, stays on this machine. Sending can be") + print_info("turned off again at any time.") + shared_metrics["send"] = prompt_yes_no( + "Send shared metrics to Nous?", + default=shared_metrics.get("send") is True, + ) + if shared_metrics["send"]: + _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_consent_change(*, enabled: bool) -> None: + """Reconcile consent windows at the moment the user decides. + + Same single writer as the relay and the sender — reconciliation derives + the window state from the observation, so wizard, relay, and mid-pass + callers cannot disagree. The relay's once-per-process reconcile would + catch this on the next hook fire anyway; running it here just makes the + wizard's effect immediate. + """ + try: + from hermes_cli.observability.shared_metrics import SharedMetricsStore + from hermes_cli.observability.shared_metrics_sender import ( + reconcile_send_consent, + ) + from hermes_cli.sqlite_util import write_txn + + store = SharedMetricsStore() + with store._connection() as connection: + with write_txn(connection): + reconcile_send_consent(connection, enabled) + except Exception: + # Never block the wizard on telemetry bookkeeping. The relay runs the + # same reconciliation on the next lifecycle hook. + logger.debug("Unable to record shared-metrics consent change", exc_info=True) # ============================================================================= diff --git a/hermes_cli/tools_config.py b/hermes_cli/tools_config.py index 77bae34105..6430e24d22 100644 --- a/hermes_cli/tools_config.py +++ b/hermes_cli/tools_config.py @@ -5622,6 +5622,43 @@ def _reconfigure_simple_requirements(ts_key: str): # ─── Main Entry Point ───────────────────────────────────────────────────────── +def _shared_metrics_state(config: dict) -> tuple[bool, bool]: + """Return (collection_enabled, send_enabled) from a config dict.""" + telemetry = config.get("telemetry") + telemetry = telemetry if isinstance(telemetry, dict) else {} + shared = telemetry.get("shared_metrics") + shared = shared if isinstance(shared, dict) else {} + return shared.get("enabled") is True, shared.get("send") is True + + +def _shared_metrics_menu_label(config: dict) -> str: + """Menu row for shared metrics, showing both consent states.""" + enabled, send = _shared_metrics_state(config) + if not enabled: + state = "off" + elif send: + state = "collecting + sending to Nous" + else: + state = "collecting locally" + return f"Configure shared metrics ({state})" + + +def _configure_shared_metrics_interactive(config: dict) -> None: + """Toggle shared-metrics collection and sending from `hermes tools`. + + Delegates to the setup wizard's prompt so the consent rules live in one + place: sending requires collection, and turning collection off also turns + sending off. + """ + from hermes_cli.setup import setup_telemetry + + before = _shared_metrics_state(config) + setup_telemetry(config) + after = _shared_metrics_state(config) + if before != after: + save_config(config) + + def tools_command(args=None, first_install: bool = False, config: dict = None): """Entry point for `hermes tools` and `hermes setup tools`. @@ -5746,6 +5783,7 @@ def tools_command(args=None, first_install: bool = False, config: dict = None): if len(platform_keys) > 1: platform_choices.append("Configure all platforms (global)") platform_choices.append("Reconfigure an existing tool's provider or API key") + platform_choices.append(_shared_metrics_menu_label(config)) # Show MCP option if any MCP servers are configured _has_mcp = bool(config.get("mcp_servers")) @@ -5757,8 +5795,9 @@ def tools_command(args=None, first_install: bool = False, config: dict = None): # Index offsets for the extra options after per-platform entries _global_idx = len(platform_keys) if len(platform_keys) > 1 else -1 _reconfig_idx = len(platform_keys) + (1 if len(platform_keys) > 1 else 0) - _mcp_idx = (_reconfig_idx + 1) if _has_mcp else -1 - _done_idx = _reconfig_idx + (2 if _has_mcp else 1) + _metrics_idx = _reconfig_idx + 1 + _mcp_idx = (_metrics_idx + 1) if _has_mcp else -1 + _done_idx = _metrics_idx + (2 if _has_mcp else 1) while True: idx = _prompt_choice("Select an option:", platform_choices, default=0) @@ -5773,6 +5812,13 @@ def tools_command(args=None, first_install: bool = False, config: dict = None): print() continue + # "Shared metrics" selected + if idx == _metrics_idx: + _configure_shared_metrics_interactive(config) + platform_choices[_metrics_idx] = _shared_metrics_menu_label(config) + print() + continue + # "Configure MCP tools" selected if idx == _mcp_idx: _configure_mcp_tools_interactive(config) diff --git a/scripts/e2e_shared_metrics_staging.py b/scripts/e2e_shared_metrics_staging.py new file mode 100644 index 0000000000..666c9e51e9 --- /dev/null +++ b/scripts/e2e_shared_metrics_staging.py @@ -0,0 +1,198 @@ +"""Live staging E2E for the shared-metrics exporter. + +Sends REAL packages through the REAL sender to the REAL staging ingest +service, then reports what the service acknowledged. Uses a throwaway +HERMES_HOME so the operator's own telemetry state is untouched. + +Usage: + .venv/bin/python scripts/e2e_shared_metrics_staging.py +""" + +from __future__ import annotations + +import json +import os +import sys +import tempfile +import uuid +from datetime import datetime, timezone +from pathlib import Path + +REPO = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(REPO)) + +STAGING = "https://telemetry.staging-nousresearch.com/v1/telemetry" + + +def main() -> int: + scratch = Path(tempfile.mkdtemp(prefix="hermes-telemetry-e2e-")) + os.environ["HERMES_HOME"] = str(scratch) + + # Staging is selected by writing config into the THROWAWAY profile, not by + # an environment override: a runtime env var that can retarget consented + # telemetry would be a consent hazard in production. + (scratch / "config.yaml").write_text( + "telemetry:\n" + " shared_metrics:\n" + " enabled: true\n" + " send: true\n" + f" endpoint: {STAGING}\n", + encoding="utf-8", + ) + + from hermes_cli.observability.shared_metrics import SharedMetricsStore + from hermes_cli.observability.shared_metrics_send_config import ( + resolve_send_config, + ) + from hermes_cli.observability.shared_metrics_sender import SharedMetricsSender + + # Resolve through the real config path so this exercises what a user gets. + import yaml + + resolved = resolve_send_config( + yaml.safe_load((scratch / "config.yaml").read_text(encoding="utf-8")) + ) + if not resolved.send or resolved.endpoint != STAGING: + print(f"FAIL: config did not resolve to staging: {resolved}") + return 1 + + store = SharedMetricsStore( + database_path=scratch / "metrics.sqlite3", + outbox_directory=scratch / "outbox", + ) + + today = datetime.now(timezone.utc).date().isoformat() + # The generator only exports COMPLETED periods, so the realistic E2E + # package is yesterday's. It also has to be: the consent gate only + # releases a package once its whole period is confirmed consented, and + # today's period cannot be confirmed before it ends. + from datetime import timedelta + + period_day = ( + datetime.now(timezone.utc).date() - timedelta(days=1) + ).isoformat() + + # Open the consent window before the period, confirm it after — exactly + # what the runtime reconciler does across two days of hook fires. + from hermes_cli.observability.shared_metrics_sender import ( + reconcile_send_consent, + ) + from hermes_cli.sqlite_util import write_txn + + with store._connection() as connection: + with write_txn(connection): + reconcile_send_consent( + connection, + True, + now=datetime.now(timezone.utc) - timedelta(days=2), + ) + reconcile_send_consent(connection, True) + real_install_id = str(uuid.uuid4()) + packages = [] + + # Two packages for today's period: the "head" and a later "tail", which is + # the real shape the outbox produces and the case the period gate exists + # for. One is large enough to exercise gzip. + for index, metric_count in ((0, 3), (1, 140)): + package_id = str(uuid.uuid4()) + payload = { + "schema_version": "hermes.shared_metrics.v2", + "package_id": package_id, + "install_id": real_install_id, + "generated_at": datetime.now(timezone.utc).isoformat().replace( + "+00:00", "Z" + ), + "period_start": f"{period_day}T00:00:00Z", + "period_end": f"{period_day}T23:59:59Z", + "resource": { + "hermes_version": "e2e-test", + "os_family": "macos", + "architecture": "arm64", + "install_method": "git", + }, + "metrics": [ + { + "name": f"hermes.e2e.metric.{i}", + "type": "counter", + "dimensions": {"outcome": "ok", "surface": "e2e"}, + "value": i + 1, + } + for i in range(metric_count) + ], + } + 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}T0{index}:00:00Z", + f"{period_day}T0{index}:00:01Z", + ), + ) + packages.append((package_id, metric_count)) + + print(f"scratch HERMES_HOME : {scratch}") + print(f"endpoint : {STAGING}") + print(f"local install_id : {real_install_id}") + print(f"packages queued : {len(packages)}") + for package_id, count in packages: + print(f" - {package_id} ({count} metrics)") + print() + + outcome = SharedMetricsSender(store, resolved.endpoint).send_pending() + print(f"outcome: sent={outcome.sent} rejected={outcome.rejected} " + f"deferred={outcome.deferred}") + print() + + failures = [] + with store._connection() as connection: + rows = connection.execute( + """ + SELECT package_id, send_state, sent_at, send_attempts, + sent_install_id, last_error + FROM package_outbox ORDER BY created_at + """ + ).fetchall() + + for row in rows: + print(f"package : {row[0]}") + print(f" send_state : {row[1]}") + print(f" sent_at : {row[2]}") + print(f" attempts : {row[3]}") + print(f" transmitted : {row[4]}") + print(f" last_error : {row[5]}") + if row[1] != "sent": + failures.append(f"{row[0]} is {row[1]}: {row[5]}") + # Product decision 2026-08-27: the stable install_id is transmitted + # as-is; the transmitted value must be exactly the local id. + if row[4] != real_install_id: + failures.append( + f"{row[0]} transmitted {row[4]!r}, expected the install_id" + ) + print() + + if failures: + print("FAILURES:") + for failure in failures: + print(f" ✗ {failure}") + return 1 + + print("PASS: every package acknowledged 202 with the stable install_id.") + print() + print("Verify the objects in S3 with the package ids above:") + print(" aws s3 ls --recursive " + "s3://hermes-agent-telemetry-staging-767397871023-us-west-2-an/raw/ " + "| tail -20") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/hermes_cli/test_setup_telemetry.py b/tests/hermes_cli/test_setup_telemetry.py index e6ebcb428c..2397524343 100644 --- a/tests/hermes_cli/test_setup_telemetry.py +++ b/tests/hermes_cli/test_setup_telemetry.py @@ -25,6 +25,51 @@ def test_setup_telemetry_enables_shared_metrics(monkeypatch): assert config["telemetry"]["shared_metrics"]["enabled"] is True +def test_disabling_collection_closes_the_send_consent_window(monkeypatch, tmp_path): + """`hermes tools` -> disable shared metrics must withdraw send consent. + + The not-enabled branch returned early without recording anything, so the + consent window stayed open and re-enabling later would release every + package collected in between. + """ + from hermes_cli.observability.shared_metrics import SharedMetricsStore + from hermes_cli.observability.shared_metrics_sender import ( + reconcile_send_consent, + ) + from hermes_cli.sqlite_util import write_txn + + store = SharedMetricsStore( + database_path=tmp_path / "m.db", outbox_directory=tmp_path / "o" + ) + monkeypatch.setattr( + "hermes_cli.observability.shared_metrics.SharedMetricsStore", + lambda *a, **k: store, + ) + + # The user had consented; now they turn collection off entirely. + monkeypatch.setattr( + "hermes_cli.setup.prompt_yes_no", lambda _question, default: False + ) + config = {"telemetry": {"shared_metrics": {"enabled": True, "send": True}}} + # Consent was granted earlier, so a window is open — that is precisely + # the state whose closure must be recorded. + with store._connection() as connection: + with write_txn(connection): + reconcile_send_consent(connection, True) + + setup_telemetry(config) + + assert config["telemetry"]["shared_metrics"]["enabled"] is False + assert config["telemetry"]["shared_metrics"]["send"] is False + with store._connection() as connection: + open_windows = connection.execute( + "SELECT COUNT(*) FROM send_consent_windows WHERE closed_at IS NULL" + ).fetchone()[0] + assert open_windows == 0, ( + "disabling collection left the send consent window open" + ) + + def test_setup_parser_accepts_telemetry_section(): parser = argparse.ArgumentParser() subparsers = parser.add_subparsers(dest="command") diff --git a/tests/hermes_cli/test_shared_metrics_consent_windows.py b/tests/hermes_cli/test_shared_metrics_consent_windows.py new file mode 100644 index 0000000000..66b58d5dd3 --- /dev/null +++ b/tests/hermes_cli/test_shared_metrics_consent_windows.py @@ -0,0 +1,310 @@ +"""Property tests for the consent-interval model. + +Ported from the /tmp validation harness that gated the redesign: every +scenario here is a defect that actually occurred (rounds 3-5) or a clock +adversary the day-stamp model could not survive. The v1 and v2 drafts of the +redesign each FAILED scenarios in this file before shipping — that is the +harness working, and why these run against the real store and the real +reconciler rather than a model of them. +""" + +from __future__ import annotations + +import json +from datetime import datetime, timedelta, timezone + +import pytest + +from hermes_cli.observability.shared_metrics import SharedMetricsStore +from hermes_cli.observability.shared_metrics_sender import ( + CONSENT_GATE_SQL, + reconcile_send_consent, +) +from hermes_cli.sqlite_util import write_txn + +T0 = datetime(2026, 8, 1, tzinfo=timezone.utc) + + +def ts(days=0, hours=0): + return (T0 + timedelta(days=days, hours=hours)).isoformat().replace( + "+00:00", "Z" + ) + + +def dt(days=0, hours=0): + return T0 + timedelta(days=days, hours=hours) + + +@pytest.fixture +def store(tmp_path): + return SharedMetricsStore( + database_path=tmp_path / "m.db", outbox_directory=tmp_path / "o" + ) + + +def _add(store, pid, start, end): + """Store a package the way the generator does: at period end.""" + with store._connection() as connection: + with write_txn(connection): + connection.execute( + "INSERT INTO package_outbox(package_id, period_start, period_end," + " payload_json, created_at, exported_at) VALUES (?, ?, ?, ?, ?, ?)", + (pid, start, end, json.dumps({"package_id": pid}), end, end), + ) + connection.execute( + """INSERT INTO consent_marks(name, stamp) VALUES ('data', ?) + ON CONFLICT(name) DO UPDATE SET stamp = MAX(stamp, excluded.stamp)""", + (end,), + ) + + +def _observe(store, send_enabled, when): + with store._connection() as connection: + with write_txn(connection): + reconcile_send_consent(connection, send_enabled, now=when) + + +def _eligible(store): + with store._connection() as connection: + return sorted( + row[0] + for row in connection.execute( + f"SELECT package_id FROM package_outbox WHERE {CONSENT_GATE_SQL}" + ) + ) + + +def _windows(store): + with store._connection() as connection: + return [ + tuple(row) + for row in connection.execute( + "SELECT opened_at, last_confirmed_at, closed_at" + " FROM send_consent_windows ORDER BY opened_at" + ) + ] + + +class TestRefusedWindowIsNeverReleased: + def test_on_off_on_with_realistic_interleaving(self, store): + """Rounds 3 and 5: the refused middle must never transmit, and + neither consented era may be lost.""" + _observe(store, True, dt(0)) + for n in range(5): + _add(store, f"d{n:02d}", ts(days=n), ts(days=n + 1)) + _observe(store, True, dt(days=n + 1)) + _observe(store, False, dt(5)) + for n in range(5, 10): + _add(store, f"d{n:02d}", ts(days=n), ts(days=n + 1)) + _observe(store, True, dt(10)) + for n in range(10, 15): + _add(store, f"d{n:02d}", ts(days=n), ts(days=n + 1)) + _observe(store, True, dt(days=n + 1)) + + eligible = _eligible(store) + assert not [p for p in eligible if 5 <= int(p[1:]) < 10], eligible + assert [f"d{n:02d}" for n in range(5)] == eligible[:5], ( + "pre-revocation consented backlog was destroyed" + ) + assert [f"d{n:02d}" for n in range(10, 15)] == eligible[5:], eligible + + def test_hand_edit_with_a_90_day_silent_gap(self, store): + """Round 5 D1, strongest form: NOTHING observes the off window. + + The close back-dates to the last confirmed moment, so the unobserved + gap is outside every window and fails closed. + """ + _observe(store, True, dt(0)) + _add(store, "consented", ts(0, 1), ts(0, 2)) + _observe(store, True, dt(0, 6)) + for n in range(1, 90, 10): + _add(store, f"REFUSED-d{n}", ts(days=n), ts(days=n, hours=1)) + _observe(store, False, dt(90)) # first observation: boot on day 90 + _observe(store, True, dt(91)) + _observe(store, True, dt(92)) + + eligible = _eligible(store) + assert not [p for p in eligible if p.startswith("REFUSED")], eligible + assert "consented" in eligible, ( + "the confirmed-morning package must survive the reconciliation" + ) + + +class TestClockAdversaries: + def test_forward_poison_then_revoke_releases_nothing(self, store): + """Round 6 D1: one glitched-forward sample must not defeat a close. + + Unfixed, the poisoned obs mark dragged last_confirmed_at to 2099, a + later revoke stamped closed_at = 2099, and the closed window then + CONTAINED every refused period that followed — all 8 refused + packages became eligible. The close now clamps to the closing + observation's own raw stamp, so an honest clock at revoke time pulls + the window back to the true revoke moment. + """ + _observe(store, True, dt(0)) + _observe(store, True, datetime(2099, 1, 1, tzinfo=timezone.utc)) + _observe(store, False, dt(1)) # honest clock at revoke + for n in range(2, 10): + _add(store, f"REFUSED-{n}", ts(days=n), ts(days=n, hours=2)) + + leaked = [p for p in _eligible(store) if p.startswith("REFUSED")] + assert not leaked, f"poisoned horizon released refused data: {leaked}" + + def test_forward_poison_cannot_wedge_consent_forever(self, store): + """The obs-advance cap bounds the damage of one insane sample. + + Uncapped, a 2099 sample would clamp every future window open at + 2099, suppressing consented data for decades (fail-closed but + permanent). Capped, the mark moves at most MAX_OBS_ADVANCE_SECONDS + past its previous value, so honest time overtakes it. + """ + from hermes_cli.observability.shared_metrics_sender import ( + MAX_OBS_ADVANCE_SECONDS, + ) + + _observe(store, True, dt(0)) + _observe(store, True, datetime(2099, 1, 1, tzinfo=timezone.utc)) + with store._connection() as connection: + stamp = connection.execute( + "SELECT stamp FROM consent_marks WHERE name = 'obs'" + ).fetchone()[0] + ceiling = ts(days=MAX_OBS_ADVANCE_SECONDS // 86_400) + assert stamp <= ceiling, ( + f"one glitched sample advanced the mark unboundedly: {stamp}" + ) + + # Consented data from shortly after the cap horizon still flows once + # honest observations catch the marks up. + horizon_days = MAX_OBS_ADVANCE_SECONDS // 86_400 + _add( + store, + "post-glitch", + ts(days=horizon_days + 1), + ts(days=horizon_days + 1, hours=4), + ) + _observe(store, True, dt(days=horizon_days + 2)) + assert "post-glitch" in _eligible(store), ( + "consent wedged after a forward glitch" + ) + + def test_rollback_at_re_enable_releases_nothing(self, store): + """Round 5 D2: the data mark clamps opens above existing packages.""" + _observe(store, True, dt(0)) + _observe(store, True, dt(5)) + _observe(store, False, dt(5)) + for n in range(1, 4): + _add(store, f"REFUSED-{n}", ts(days=5, hours=n), ts(days=5, hours=n + 1)) + _observe(store, True, dt(-12)) # 12-day rollback at re-enable + _observe(store, True, dt(-11)) + + during = [p for p in _eligible(store) if p.startswith("REFUSED")] + assert not during, f"rollback released refused packages: {during}" + + _observe(store, True, dt(20)) # clock recovers + _observe(store, True, dt(21)) + after = [p for p in _eligible(store) if p.startswith("REFUSED")] + assert not after, f"recovery released refused packages: {after}" + + def test_recovery_does_not_wedge_future_sending(self, store): + _observe(store, True, dt(0)) + _observe(store, False, dt(5)) + _observe(store, True, dt(-12)) + _observe(store, True, dt(20)) + _add(store, "post-recovery", ts(21), ts(21, 4)) + _observe(store, True, dt(22)) + assert "post-recovery" in _eligible(store) + + +class TestSubDayGranularity: + def test_intra_day_refusal_holds_back_the_whole_day_package(self, store): + """Round 5 D3: a day package spanning a refused stretch must wait.""" + _observe(store, True, dt(0)) + _observe(store, True, dt(10, 9)) + _observe(store, False, dt(10, 9)) + _observe(store, True, dt(10, 18)) + _observe(store, True, dt(11, 2)) + _add(store, "halfday", ts(10), ts(11)) + assert "halfday" not in _eligible(store) + + +class TestReconcilerProperties: + def test_idempotent_under_replay(self, store): + for _ in range(4): + _observe(store, True, dt(0)) + _observe(store, False, dt(2)) + for _ in range(5): + _observe(store, False, dt(3)) + _observe(store, True, dt(4)) + for _ in range(3): + _observe(store, True, dt(5)) + assert len(_windows(store)) == 2 + + def test_the_observation_mark_is_monotonic(self, store): + """A rolled-back clock must never lower the observation high-water. + + Every downstream guarantee leans on this: closes clamp to it via + last_confirmed_at, and opens clamp to max(obs, data). Found as a + surviving mutant (obs upsert rewritten from MAX to overwrite) — + the leak scenarios happen to be covered by the data mark whenever a + leakable package exists, but the property itself must hold on its + own, not by coincidence of the sibling mark. + """ + _observe(store, True, dt(5)) + _observe(store, True, dt(0)) # rollback + with store._connection() as connection: + stamp = connection.execute( + "SELECT stamp FROM consent_marks WHERE name = 'obs'" + ).fetchone()[0] + assert stamp == ts(5), f"obs mark moved backwards: {stamp}" + + def test_the_real_package_writer_advances_the_data_mark(self, store): + """Round 6 D2: the harness's _add re-implements the data-mark insert, + so deleting the advance from the REAL writer survived 314 tests. + This drives the production exporter instead. + """ + from datetime import date, timedelta as _td + + yesterday = (date.today() - _td(days=1)).isoformat() + with store._connection() as connection: + with write_txn(connection): + connection.execute( + "INSERT INTO counter_aggregates(" + " period_start, metric_name, hermes_version, os_family," + " architecture, install_method, dimensions_json, value," + " packaged_value" + ") VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)", + ( + yesterday, "hermes.client.active", "0.0.0-test", + "macos", "arm64", "git", "{}", 1, 0, + ), + ) + + exported = store.create_and_export_package_if_due() + assert exported, "the generator was expected to export yesterday's period" + + with store._connection() as connection: + row = connection.execute( + "SELECT stamp FROM consent_marks WHERE name = 'data'" + ).fetchone() + assert row is not None and row[0] >= yesterday, ( + "the production package writer did not advance the data mark" + ) + + def test_the_gate_is_read_only(self, store): + _observe(store, True, dt(0)) + before = _windows(store) + for _ in range(10): + _eligible(store) + assert _windows(store) == before + + def test_no_window_fails_closed(self, store): + _add(store, "orphan", ts(0), ts(1)) + assert _eligible(store) == [] + + def test_fresh_package_waits_one_heartbeat_then_releases(self, store): + """The documented latency cost of confirmation-based windows.""" + _observe(store, True, dt(0)) + _add(store, "fresh", ts(0, 1), ts(0, 2)) + assert _eligible(store) == [] + _observe(store, True, dt(0, 3)) + assert _eligible(store) == ["fresh"] diff --git a/tests/hermes_cli/test_shared_metrics_send_config.py b/tests/hermes_cli/test_shared_metrics_send_config.py new file mode 100644 index 0000000000..2af8958a2b --- /dev/null +++ b/tests/hermes_cli/test_shared_metrics_send_config.py @@ -0,0 +1,165 @@ +"""Tests for shared-metrics send configuration resolution.""" + +from __future__ import annotations + +import logging + +import pytest + +from hermes_cli.config import DEFAULT_CONFIG +from hermes_cli.observability.shared_metrics_send_config import ( + DEFAULT_ENDPOINT, + resolve_send_config, + reset_warning_latch_for_tests, +) + + +@pytest.fixture(autouse=True) +def _reset_latch(): + reset_warning_latch_for_tests() + yield + reset_warning_latch_for_tests() + + +def _config(**shared): + return {"telemetry": {"shared_metrics": shared}} + + +class TestDefaults: + def test_send_is_registered_disabled_by_default(self): + shared = DEFAULT_CONFIG["telemetry"]["shared_metrics"] + assert shared["enabled"] is False + assert shared["send"] is False + + def test_default_endpoint_is_production(self): + shared = DEFAULT_CONFIG["telemetry"]["shared_metrics"] + assert shared["endpoint"] == DEFAULT_ENDPOINT + assert DEFAULT_ENDPOINT.startswith("https://") + + def test_empty_config_sends_nothing(self): + resolved = resolve_send_config({}) + assert resolved.enabled is False + assert resolved.send is False + + def test_none_config_is_tolerated(self): + assert resolve_send_config(None).send is False + + +class TestSendRequiresCollection: + def test_collection_alone_does_not_send(self): + resolved = resolve_send_config(_config(enabled=True)) + assert resolved.enabled is True + assert resolved.send is False + + def test_send_with_collection_sends(self): + resolved = resolve_send_config(_config(enabled=True, send=True)) + assert resolved.send is True + + def test_send_without_collection_is_refused(self): + resolved = resolve_send_config(_config(enabled=False, send=True)) + assert resolved.send is False + # send must never imply enabled + assert resolved.enabled is False + + def test_send_without_collection_logs_an_error(self, caplog): + with caplog.at_level(logging.ERROR): + resolve_send_config(_config(enabled=False, send=True)) + errors = [r for r in caplog.records if r.levelno >= logging.ERROR] + assert len(errors) == 1 + assert "enabled is false" in errors[0].getMessage() + + def test_the_error_is_logged_once_per_process(self, caplog): + with caplog.at_level(logging.ERROR): + for _ in range(5): + resolve_send_config(_config(enabled=False, send=True)) + errors = [r for r in caplog.records if r.levelno >= logging.ERROR] + assert len(errors) == 1, "misconfiguration must not spam every hook fire" + + +class TestEndpointPrecedence: + def test_config_endpoint_overrides_default(self): + resolved = resolve_send_config( + _config(enabled=True, send=True, endpoint="https://example.test/v1") + ) + assert resolved.endpoint == "https://example.test/v1" + + def test_no_environment_variable_can_redirect_telemetry(self, monkeypatch): + """A consent hazard: an inherited env var must not silently retarget. + + AGENTS.md also reserves HERMES_* for secrets, not behaviour. + """ + for name in ( + "HERMES_TELEMETRY_ENDPOINT", + "TELEMETRY_ENDPOINT", + "HERMES_SHARED_METRICS_ENDPOINT", + ): + monkeypatch.setenv(name, "https://attacker.test/v1") + resolved = resolve_send_config(_config(enabled=True, send=True)) + assert resolved.endpoint == DEFAULT_ENDPOINT + + def test_blank_endpoint_falls_back_to_production(self): + resolved = resolve_send_config(_config(enabled=True, send=True, endpoint=" ")) + assert resolved.endpoint == DEFAULT_ENDPOINT + + def test_endpoint_is_stripped(self): + resolved = resolve_send_config( + _config(enabled=True, send=True, endpoint=" https://staging.test/v1 ") + ) + assert resolved.endpoint == "https://staging.test/v1" + + +class TestTransportSafety: + def test_plaintext_endpoint_is_refused(self, caplog): + with caplog.at_level(logging.ERROR): + resolved = resolve_send_config( + _config(enabled=True, send=True, endpoint="http://example.test/v1") + ) + assert resolved.send is False, "telemetry must not go out in clear text" + assert any("https" in r.getMessage() for r in caplog.records) + + @pytest.mark.parametrize( + "endpoint", + [ + "http://localhost:8099/v1/telemetry", + "http://127.0.0.1:8099/v1/telemetry", + ], + ) + def test_loopback_http_is_allowed_for_testing(self, endpoint): + resolved = resolve_send_config( + _config(enabled=True, send=True, endpoint=endpoint) + ) + assert resolved.send is True + + def test_nonsense_scheme_is_refused(self): + resolved = resolve_send_config( + _config(enabled=True, send=True, endpoint="ftp://example.test/v1") + ) + assert resolved.send is False + + @pytest.mark.parametrize( + "endpoint", + [ + "ftp://localhost/v1/telemetry", + "gopher://localhost/v1/telemetry", + "ws://127.0.0.1/v1/telemetry", + ], + ) + def test_a_non_http_scheme_on_loopback_is_still_refused(self, endpoint): + """The scheme is allowlisted, not merely checked for plaintext http. + + Gap found by mutation testing: replacing the `http` scheme test with + `if True` survived the whole suite, because every non-http scheme case + pointed at a REMOTE host, where the loopback branch rejects it anyway. + Only a non-http scheme aimed at loopback distinguishes an allowlist + from a plaintext-only check. + """ + resolved = resolve_send_config( + _config(enabled=True, send=True, endpoint=endpoint) + ) + assert resolved.send is False + + def test_unsafe_endpoint_does_not_block_collection(self): + resolved = resolve_send_config( + _config(enabled=True, send=True, endpoint="http://example.test/v1") + ) + assert resolved.enabled is True diff --git a/tests/hermes_cli/test_shared_metrics_send_migration.py b/tests/hermes_cli/test_shared_metrics_send_migration.py new file mode 100644 index 0000000000..54518c644a --- /dev/null +++ b/tests/hermes_cli/test_shared_metrics_send_migration.py @@ -0,0 +1,224 @@ +"""Tests for the additive send-state migration on ``package_outbox``. + +The store schema version must NOT move when these columns are added: the +existing loader raises on any version it does not recognise, so bumping it +would hard-fail an older Hermes (a second profile on an older build, or a +rollback) against the same database file. +""" + +from __future__ import annotations + +import json +import sqlite3 + +import pytest + +from hermes_cli.observability.shared_metrics import SharedMetricsStore + +SEND_COLUMNS = { + "sent_at", + "send_state", + "send_attempts", + "next_attempt_at", + "last_error", + "sent_install_id", +} + + +def _columns(db_path): + connection = sqlite3.connect(db_path) + try: + return {row[1] for row in connection.execute("PRAGMA table_info(package_outbox)")} + finally: + connection.close() + + +def _schema_version(db_path): + connection = sqlite3.connect(db_path) + try: + row = connection.execute( + "SELECT value FROM telemetry_state WHERE key = 'schema_version'" + ).fetchone() + return row[0] if row else None + finally: + connection.close() + + +@pytest.fixture +def store(tmp_path): + return SharedMetricsStore( + database_path=tmp_path / "metrics.sqlite3", + outbox_directory=tmp_path / "outbox", + ) + + +class TestFreshDatabase: + def test_send_columns_exist(self, store): + assert SEND_COLUMNS <= _columns(store.database_path) + + def test_original_columns_survive(self, store): + assert { + "package_id", + "period_start", + "period_end", + "payload_json", + "created_at", + "exported_at", + } <= _columns(store.database_path) + + def test_send_attempts_defaults_to_zero(self, store): + connection = sqlite3.connect(store.database_path) + try: + connection.execute( + """ + INSERT INTO package_outbox( + package_id, period_start, period_end, payload_json, created_at + ) VALUES ('p', '2026-01-01', '2026-01-02', '{}', '2026-01-01T00:00:00Z') + """ + ) + connection.commit() + row = connection.execute( + "SELECT send_attempts, send_state, sent_install_id FROM package_outbox" + ).fetchone() + finally: + connection.close() + assert row[0] == 0 + assert row[1] is None + assert row[2] is None + + +class TestUpgradeFromPreSendDatabase: + """The real-world case: a database written before this feature existed.""" + + @pytest.fixture + def legacy_db(self, tmp_path): + path = tmp_path / "metrics.sqlite3" + connection = sqlite3.connect(path) + try: + connection.execute( + """ + CREATE TABLE telemetry_state ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL + ) + """ + ) + connection.execute( + "INSERT INTO telemetry_state(key, value) VALUES ('schema_version', '2')" + ) + connection.execute( + """ + CREATE TABLE package_outbox ( + package_id TEXT PRIMARY KEY, + period_start TEXT NOT NULL, + period_end TEXT NOT NULL, + payload_json TEXT NOT NULL, + created_at TEXT NOT NULL, + exported_at TEXT + ) + """ + ) + connection.execute( + """ + CREATE TABLE counter_aggregates ( + period_start TEXT NOT NULL, + metric_name TEXT NOT NULL, + hermes_version TEXT NOT NULL, + os_family TEXT NOT NULL, + architecture TEXT NOT NULL, + install_method TEXT NOT NULL, + dimensions_json TEXT NOT NULL, + value INTEGER NOT NULL, + packaged_value INTEGER NOT NULL, + PRIMARY KEY ( + period_start, metric_name, hermes_version, os_family, + architecture, install_method, dimensions_json + ) + ) + """ + ) + for i in range(3): + connection.execute( + """ + INSERT INTO package_outbox( + package_id, period_start, period_end, payload_json, + created_at, exported_at + ) VALUES (?, ?, ?, ?, ?, ?) + """, + ( + f"pkg-{i}", + "2026-08-2%d" % i, + "2026-08-2%d" % (i + 1), + json.dumps({"package_id": f"pkg-{i}"}), + "2026-08-2%dT00:00:00Z" % i, + "2026-08-2%dT01:00:00Z" % i, + ), + ) + connection.commit() + finally: + connection.close() + return path + + def test_upgrade_preserves_every_row(self, legacy_db, tmp_path): + SharedMetricsStore( + database_path=legacy_db, outbox_directory=tmp_path / "outbox" + ) + connection = sqlite3.connect(legacy_db) + try: + count = connection.execute("SELECT COUNT(*) FROM package_outbox").fetchone()[0] + payloads = connection.execute( + "SELECT package_id, payload_json FROM package_outbox ORDER BY package_id" + ).fetchall() + finally: + connection.close() + assert count == 3 + assert payloads == [ + ("pkg-0", '{"package_id": "pkg-0"}'), + ("pkg-1", '{"package_id": "pkg-1"}'), + ("pkg-2", '{"package_id": "pkg-2"}'), + ] + + def test_upgrade_adds_the_send_columns(self, legacy_db, tmp_path): + SharedMetricsStore( + database_path=legacy_db, outbox_directory=tmp_path / "outbox" + ) + assert SEND_COLUMNS <= _columns(legacy_db) + + def test_upgrade_does_not_move_the_schema_version(self, legacy_db, tmp_path): + """Bumping would make older builds refuse the same file.""" + SharedMetricsStore( + database_path=legacy_db, outbox_directory=tmp_path / "outbox" + ) + assert _schema_version(legacy_db) == "2" + + def test_migration_is_idempotent(self, legacy_db, tmp_path): + for _ in range(3): + SharedMetricsStore( + database_path=legacy_db, outbox_directory=tmp_path / "outbox" + ) + columns = [ + row[1] + for row in sqlite3.connect(legacy_db).execute( + "PRAGMA table_info(package_outbox)" + ) + ] + assert len(columns) == len(set(columns)), "columns were added more than once" + + def test_queries_written_before_this_change_still_work(self, legacy_db, tmp_path): + """The shipped export query selects named columns; it must be unaffected.""" + SharedMetricsStore( + database_path=legacy_db, outbox_directory=tmp_path / "outbox" + ) + connection = sqlite3.connect(legacy_db) + try: + rows = connection.execute( + """ + SELECT package_id, payload_json + FROM package_outbox + WHERE exported_at IS NULL + ORDER BY created_at, package_id + """ + ).fetchall() + finally: + connection.close() + assert rows == [] diff --git a/tests/hermes_cli/test_shared_metrics_send_wiring.py b/tests/hermes_cli/test_shared_metrics_send_wiring.py new file mode 100644 index 0000000000..29f572c697 --- /dev/null +++ b/tests/hermes_cli/test_shared_metrics_send_wiring.py @@ -0,0 +1,446 @@ +"""Tests for wiring the sender into the shared-metrics export hook. + +The properties that matter here are negative ones: the interactive path must +not block, and nothing must leave the machine unless the user opted in. +""" + +from __future__ import annotations + +import threading +import time + +import pytest + +from hermes_cli.observability import relay_shared_metrics as mod + + +class FakeStore: + def __init__(self): + self.exported = 0 + + def create_and_export_package_if_due(self): + self.exported += 1 + return [] + + +class RealBackedStore: + """A store with a genuine SQLite connection, for consent-state tests. + + The consent edge detector writes to telemetry_state, and it is wrapped in + a broad except. Against a stub without _connection it would swallow an + AttributeError and silently do nothing — which is exactly the failure this + file needs to be able to catch. + """ + + def __init__(self, tmp_path): + from hermes_cli.observability.shared_metrics import SharedMetricsStore + + self._real = SharedMetricsStore( + database_path=tmp_path / "m.db", outbox_directory=tmp_path / "o" + ) + self.exported = 0 + + def _connection(self): + return self._real._connection() + + def create_and_export_package_if_due(self): + self.exported += 1 + return [] + + +class FakeSubscriber: + def __init__(self): + self.store = FakeStore() + + +class Runtime(mod._Runtime): + """A _Runtime with the relay host stubbed out.""" + + def __init__(self): + self._sessions_lock = threading.RLock() + self._sessions = {} + self._task_creation_lock = threading.RLock() + self._task_sessions_lock = threading.RLock() + self._send_lock = threading.RLock() + self._send_thread = None + self._task_sessions = {} + self._turn_sessions = {} + self.subscriber = FakeSubscriber() + + +@pytest.fixture +def runtime(): + return Runtime() + + +def _config(**shared): + return {"telemetry": {"shared_metrics": shared}} + + +@pytest.fixture +def capture_sender(monkeypatch): + """Replace the sender with a recorder and return the record.""" + record = {"passes": [], "endpoints": []} + + class FakeSender: + def __init__(self, store, endpoint, **kwargs): + record["endpoints"].append(endpoint) + + def send_pending(self): + record["passes"].append(time.time()) + + monkeypatch.setattr( + "hermes_cli.observability.shared_metrics_sender.SharedMetricsSender", + FakeSender, + ) + return record + + +def _set_config(monkeypatch, config): + monkeypatch.setattr( + "hermes_cli.config.read_raw_config_readonly", lambda: config, raising=False + ) + + +class TestOptIn: + def test_no_send_when_nothing_is_configured(self, runtime, monkeypatch, capture_sender): + _set_config(monkeypatch, {}) + runtime._export() + runtime._join_send_thread(timeout=1) + assert capture_sender["passes"] == [] + + def test_no_send_when_only_collection_is_on(self, runtime, monkeypatch, capture_sender): + _set_config(monkeypatch, _config(enabled=True)) + runtime._export() + runtime._join_send_thread(timeout=1) + assert capture_sender["passes"] == [] + + def test_no_send_when_send_is_on_without_collection( + self, runtime, monkeypatch, capture_sender + ): + _set_config(monkeypatch, _config(enabled=False, send=True)) + runtime._export() + runtime._join_send_thread(timeout=1) + assert capture_sender["passes"] == [] + + def test_sends_when_both_are_on(self, runtime, monkeypatch, capture_sender): + _set_config(monkeypatch, _config(enabled=True, send=True)) + runtime._export() + runtime._join_send_thread(timeout=2) + assert len(capture_sender["passes"]) == 1 + + def test_uses_the_resolved_endpoint(self, runtime, monkeypatch, capture_sender): + _set_config( + monkeypatch, + _config(enabled=True, send=True, endpoint="https://staging.test/v1"), + ) + runtime._export() + runtime._join_send_thread(timeout=2) + assert capture_sender["endpoints"] == ["https://staging.test/v1"] + + def test_export_still_runs_when_sending_is_off(self, runtime, monkeypatch, capture_sender): + _set_config(monkeypatch, _config(enabled=True)) + runtime._export() + assert runtime.subscriber.store.exported == 1 + + +class TestInteractivePathIsNotBlocked: + def test_export_returns_before_the_send_finishes( + self, runtime, monkeypatch + ): + started = threading.Event() + release = threading.Event() + + class SlowSender: + def __init__(self, store, endpoint, **kwargs): + pass + + def send_pending(self): + started.set() + release.wait(5) + + monkeypatch.setattr( + "hermes_cli.observability.shared_metrics_sender.SharedMetricsSender", + SlowSender, + ) + _set_config(monkeypatch, _config(enabled=True, send=True)) + + began = time.monotonic() + runtime._export() + elapsed = time.monotonic() - began + + assert started.wait(2), "the send should have started" + assert elapsed < 1.0, "finish_task must not wait on the network" + release.set() + runtime._join_send_thread(timeout=5) + + def test_the_send_thread_is_a_daemon(self, runtime, monkeypatch, capture_sender): + _set_config(monkeypatch, _config(enabled=True, send=True)) + runtime._export() + with runtime._send_lock: + thread = runtime._send_thread + assert thread is not None + assert thread.daemon, "an unfinished send must not hold the process open" + runtime._join_send_thread(timeout=2) + + def test_only_one_pass_runs_at_a_time(self, runtime, monkeypatch): + release = threading.Event() + starts = [] + + class SlowSender: + def __init__(self, store, endpoint, **kwargs): + pass + + def send_pending(self): + starts.append(1) + release.wait(5) + + monkeypatch.setattr( + "hermes_cli.observability.shared_metrics_sender.SharedMetricsSender", + SlowSender, + ) + _set_config(monkeypatch, _config(enabled=True, send=True)) + + for _ in range(5): + runtime._export() + time.sleep(0.2) + assert len(starts) == 1, "hook fires must not pile up send passes" + release.set() + runtime._join_send_thread(timeout=5) + + +class TestConsentWindows: + """Consent reconciliation must work from the relay, in any order. + + Round 4's edge detector missed the idle-revocation path; round 5 found it + was also dead code whenever collection was off (handles_hook gated it). + These tests drive the relay entry points against the single reconciler + and assert on the interval table — the only consent state that exists. + """ + + def _runtime(self, tmp_path): + runtime = Runtime() + runtime.subscriber.store = RealBackedStore(tmp_path) + return runtime + + def _windows(self, runtime): + with runtime.subscriber.store._connection() as connection: + return [ + tuple(row) + for row in connection.execute( + "SELECT opened_at, last_confirmed_at, closed_at" + " FROM send_consent_windows ORDER BY opened_at" + ) + ] + + def test_revoking_while_idle_closes_the_window( + self, monkeypatch, tmp_path, capture_sender + ): + runtime = self._runtime(tmp_path) + + _set_config(monkeypatch, _config(enabled=True, send=True)) + runtime._send_exported_packages() + + # User edits config.yaml: send: false. Hooks keep firing normally. + _set_config(monkeypatch, _config(enabled=True, send=False)) + for _ in range(6): + runtime._send_exported_packages() + + windows = self._windows(runtime) + assert windows and all(w[2] is not None for w in windows), ( + f"revoking while idle left a window open: {windows}" + ) + + def test_replayed_observations_create_no_junk_windows( + self, monkeypatch, tmp_path, capture_sender + ): + """Reconciliation is idempotent — there is no edge to double-count.""" + runtime = self._runtime(tmp_path) + + _set_config(monkeypatch, _config(enabled=True, send=True)) + for _ in range(4): + runtime._send_exported_packages() + _set_config(monkeypatch, _config(enabled=True, send=False)) + for _ in range(4): + runtime._send_exported_packages() + _set_config(monkeypatch, _config(enabled=True, send=True)) + for _ in range(4): + runtime._send_exported_packages() + + assert len(self._windows(runtime)) == 2 + + def test_a_never_consented_user_gets_no_window( + self, monkeypatch, tmp_path, capture_sender + ): + runtime = self._runtime(tmp_path) + _set_config(monkeypatch, _config(enabled=True, send=False)) + for _ in range(5): + runtime._send_exported_packages() + + assert self._windows(runtime) == [] + + def test_re_enabling_opens_a_new_window_after_the_refusal( + self, monkeypatch, tmp_path, capture_sender + ): + """The refused gap must fall BETWEEN the two windows.""" + runtime = self._runtime(tmp_path) + _set_config(monkeypatch, _config(enabled=True, send=True)) + runtime._send_exported_packages() + _set_config(monkeypatch, _config(enabled=True, send=False)) + runtime._send_exported_packages() + _set_config(monkeypatch, _config(enabled=True, send=True)) + runtime._send_exported_packages() + + windows = self._windows(runtime) + assert len(windows) == 2 + first, second = windows + assert first[2] is not None, "first window must be closed" + assert second[2] is None, "second window must be open" + assert second[0] >= first[2], ( + f"new window may not overlap the refused gap: {windows}" + ) + + def test_reconcile_runs_even_when_collection_is_disabled( + self, monkeypatch, tmp_path + ): + """Round-5 D1: enabled:false must not make consent handling dead code. + + The module-level once-per-process reconciler must close the window + regardless of handles_hook(). Drives the real observe_lifecycle gate + path: handles_hook is False throughout. + """ + from hermes_cli.observability.shared_metrics import SharedMetricsStore + from hermes_cli.observability.shared_metrics_sender import ( + reconcile_send_consent, + ) + from hermes_cli.sqlite_util import write_txn + + # Lay the store out exactly as production does, under a redirected + # HERMES_HOME: the boot reconciler probes the default path (without + # constructing the store — the constructor creates directories), so + # the probe and the store must agree the way they do in production. + home = tmp_path / "home" + monkeypatch.setattr( + "hermes_constants.get_hermes_home", lambda: home + ) + root = home / "telemetry" / "shared_metrics" + store = SharedMetricsStore( + database_path=root / "metrics.sqlite3", + outbox_directory=root / "outbox", + ) + # A consent window is open from an earlier consented era. + with store._connection() as connection: + with write_txn(connection): + reconcile_send_consent(connection, True) + + monkeypatch.setattr( + "hermes_cli.observability.shared_metrics.SharedMetricsStore", + lambda *a, **k: store, + ) + _set_config(monkeypatch, _config(enabled=False, send=False)) + monkeypatch.setattr(mod, "_consent_reconcile_done", False) + + # The full lifecycle entry point, with collection OFF. + mod.observe_lifecycle("finish_task") + + with store._connection() as connection: + open_windows = connection.execute( + "SELECT COUNT(*) FROM send_consent_windows WHERE closed_at IS NULL" + ).fetchone()[0] + assert open_windows == 0, ( + "enabled:false made the consent reconciler unreachable (D1)" + ) + + + +class TestFailureIsolation: + def test_a_sender_crash_does_not_propagate(self, runtime, monkeypatch): + class Exploding: + def __init__(self, store, endpoint, **kwargs): + pass + + def send_pending(self): + raise RuntimeError("boom") + + monkeypatch.setattr( + "hermes_cli.observability.shared_metrics_sender.SharedMetricsSender", + Exploding, + ) + _set_config(monkeypatch, _config(enabled=True, send=True)) + runtime._export() # must not raise + runtime._join_send_thread(timeout=2) + + def test_an_unreadable_config_does_not_break_export(self, runtime, monkeypatch, capture_sender): + def explode(): + raise OSError("config unreadable") + + monkeypatch.setattr( + "hermes_cli.config.read_raw_config_readonly", explode, raising=False + ) + runtime._export() + assert runtime.subscriber.store.exported == 1 + assert capture_sender["passes"] == [] + + def test_join_is_safe_with_no_thread(self, runtime): + runtime._join_send_thread(timeout=0.1) + + def test_join_waits_for_an_in_flight_send(self, runtime, monkeypatch): + """shutdown() must give a started send a chance to finish. + + A short-lived CLI exits straight after its final export; without the + join the daemon thread is killed mid-request, and the hook path is the + only delivery cadence this feature has. + """ + finished = [] + release = threading.Event() + + class SlowSender: + def __init__(self, store, endpoint, **kwargs): + pass + + def send_pending(self): + release.wait(3) + finished.append(True) + + monkeypatch.setattr( + "hermes_cli.observability.shared_metrics_sender.SharedMetricsSender", + SlowSender, + ) + _set_config(monkeypatch, _config(enabled=True, send=True)) + + runtime._export() + release.set() + runtime._join_send_thread(timeout=3) + assert finished == [True] + + def test_shutdown_joins_the_send_thread(self, monkeypatch): + """shutdown() must actually wait, not merely mention the join. + + 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" 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..cbaff4c9ea --- /dev/null +++ b/tests/hermes_cli/test_shared_metrics_sender.py @@ -0,0 +1,1029 @@ +"""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_ATTEMPTS, + MAX_PACKAGES_PER_PASS, + MAX_SEND_ATTEMPTS, + REQUEST_TIMEOUT_SECONDS, + SharedMetricsSender, + reconcile_send_consent, +) +from hermes_cli.sqlite_util import write_txn + +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): + """A store with a broad consent window already open. + + Most tests exercise claiming/retry/transport, not the consent gate, and + the interval gate fails closed with no window. One window opened before + every test package and confirmed well past NOW keeps those tests about + what they are about. Gate tests clear it via _clear_consent. + """ + built = SharedMetricsStore( + database_path=tmp_path / "metrics.sqlite3", + outbox_directory=tmp_path / "outbox", + ) + _grant_consent(built) + return built + + +def _grant_consent( + store, + opened=datetime(2026, 8, 20, tzinfo=timezone.utc), + confirmed_through=datetime(2026, 10, 1, tzinfo=timezone.utc), +): + """Open a consent window and heartbeat it forward, via the real writer.""" + with store._connection() as connection: + with write_txn(connection): + reconcile_send_consent(connection, True, now=opened) + reconcile_send_consent(connection, True, now=confirmed_through) + + +def _revoke_consent(store, at): + with store._connection() as connection: + with write_txn(connection): + reconcile_send_consent(connection, False, now=at) + + +def _clear_consent(store): + """Remove all consent state, for tests of the fail-closed default.""" + with store._connection() as connection: + with write_txn(connection): + connection.execute("DELETE FROM send_consent_windows") + connection.execute("DELETE FROM consent_marks") + + +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 _iso(moment): + return moment.astimezone(timezone.utc).isoformat().replace("+00:00", "Z") + + +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 == [] + + @pytest.mark.parametrize("status", [401, 403, 404, 422, 500, 503]) + def test_unspecified_statuses_are_retried_not_discarded(self, store, status): + """403 is the ingest origin guard; a bad edge config must not lose data.""" + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(*[FakeResponse(status)] * 3) + outcome = _sender(store, transport).send_pending() + assert outcome.deferred == 1 + assert _row(store, "pkg-1")["send_state"] == "pending" + + def test_413_is_permanent(self, store): + """A package over the 1 MiB cap cannot shrink by being retried.""" + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(413)) + outcome = _sender(store, transport).send_pending() + assert outcome.rejected == 1 + assert len(transport.calls) == 1 + + 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): + # Consent opens on Aug 24; the "old" package's period predates it. + _clear_consent(store) + _grant_consent(store, opened=datetime(2026, 8, 24, tzinfo=timezone.utc)) + _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_is_immortalised_as_a_window_not_a_day(self, store): + """The window survives replayed observations without moving.""" + with store._connection() as connection: + rows = connection.execute( + "SELECT opened_at, closed_at FROM send_consent_windows" + ).fetchall() + assert len(rows) == 1 and rows[0][1] is None + _grant_consent(store) # replay: must not create a second window + with store._connection() as connection: + count = connection.execute( + "SELECT COUNT(*) FROM send_consent_windows" + ).fetchone()[0] + assert count == 1 + + def test_no_consent_window_means_nothing_is_sent(self, store): + """The gate fails closed: absence of a window is absence of consent.""" + _clear_consent(store) + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(202)) + _sender(store, transport).send_pending() + assert transport.calls == [] + + 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 == [] + + def test_revoking_then_re_enabling_never_releases_the_off_window(self, store): + """The R3/R5 leak: re-opt-in must not release the refused interval. + + Under the interval model the refused days fall BETWEEN two windows; + no later observation can place them inside one, so the property holds + for any number of on/off cycles — not just the single cycle the old + moving day-stamp was patched to survive. + """ + _clear_consent(store) + _grant_consent(store, opened=NOW - timedelta(days=2), confirmed_through=NOW) + _add_package(store, "consented", "2026-08-25") + + # User turns sending off; packages keep being collected for 3 days. + _revoke_consent(store, at=NOW) + for day in ("2026-08-27", "2026-08-28", "2026-08-29"): + _add_package(store, f"refused-{day}", day) + + # User re-enables 5 days later; heartbeat confirms past the horizon. + later = NOW + timedelta(days=5) + with store._connection() as connection: + with write_txn(connection): + reconcile_send_consent(connection, True, now=later) + reconcile_send_consent( + connection, True, now=later + timedelta(days=30) + ) + + 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}" + ) + # And the interval model's improvement over the day-stamp: the + # pre-revocation consented package is NOT collateral damage. + assert "consented" in sent, ( + "the consented backlog was destroyed by the revoke/re-enable cycle" + ) + + def test_a_package_from_after_re_enabling_is_sent(self, store): + """The revocation handling must not wedge sending off permanently.""" + _clear_consent(store) + _grant_consent(store, opened=NOW - timedelta(days=2), confirmed_through=NOW) + _revoke_consent(store, at=NOW) + + later = NOW + timedelta(days=5) + with store._connection() as connection: + with write_txn(connection): + reconcile_send_consent(connection, True, now=later) + reconcile_send_consent( + connection, True, now=later + timedelta(days=10) + ) + _add_package(store, "after-re-optin", (later + timedelta(days=1)).date().isoformat()) + transport = FakeTransport(FakeResponse(202)) + SharedMetricsSender( + store, ENDPOINT, post=transport, sleep=lambda _s: None, + now=lambda: later + timedelta(days=2), + ).send_pending() + assert len(transport.calls) == 1 + + +class TestIdentity: + def test_the_stable_install_id_is_transmitted_as_is(self, store): + """Product decision 2026-08-27: no pseudonymization. + + The wire body carries the profile-scoped install_id verbatim. This + test is the deliberate inversion of the pre-decision assertion that + the raw id never crossed the wire. + """ + _add_package(store, "pkg-1", "2026-08-26") + transport = FakeTransport(FakeResponse(202)) + _sender(store, transport).send_pending() + assert transport.bodies[0]["install_id"] == INSTALL_ID + + def test_transmitted_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. + + The second pass must RECORD what it saw rather than raise: _send_one + catches every exception as a retryable transport failure, so an + assertion thrown inside a transport would be swallowed and this test + would pass no matter what the claim did. + """ + _add_package(store, "pkg-1", "2026-08-26") + + first_calls = [] + second_calls = [] + + def second_transport(endpoint, payload, *, timeout): + second_calls.append(payload) + return FakeResponse(202) + + def transport(endpoint, payload, *, timeout): + first_calls.append(payload) + # A second sender runs while the first is mid-flight. + SharedMetricsSender( + store, + ENDPOINT, + post=second_transport, + sleep=lambda _s: None, + now=lambda: NOW, + ).send_pending() + return FakeResponse(202) + + _sender(store, transport).send_pending() + assert len(first_calls) == 1 + assert second_calls == [], ( + "a concurrent pass claimed a package already in flight" + ) + + 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 + + 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. + + One package can legally take ~96s (three 30s timeouts plus backoff). + With 20 rows claimed under one shared lease, the later rows' leases + expired mid-pass and a second process re-sent them. Packages are now + claimed one at a time, immediately before transmission. + """ + for i in range(3): + _add_package(store, f"pkg-{i}", "2026-08-26") + + clock = {"t": NOW} + first_posts, second_posts = [], [] + + + def transport(endpoint, payload, *, timeout): + pid = json.loads(payload)["package_id"] + first_posts.append(pid) + # Burn the worst-case time budget for a single package. + clock["t"] += timedelta(seconds=96) + # A concurrent process probes for work while this package is still + # in flight. It must not be able to claim the package we hold. + # Restricted to that package so the probe cannot legitimately pick + # up the OTHER pending rows and make the assertion ambiguous. + held = _row(store, pid) + if held["next_attempt_at"] is not None: + eligible = held["next_attempt_at"] <= _iso(clock["t"]) + if eligible and held["send_state"] != "sent": + second_posts.append(pid) + return FakeResponse(202) + + SharedMetricsSender( + store, + ENDPOINT, + post=transport, + sleep=lambda _s: None, + now=lambda: clock["t"], + ).send_pending() + + assert sorted(first_posts) == ["pkg-0", "pkg-1", "pkg-2"] + assert second_posts == [], ( + 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_a_lapsed_claimant_yields_even_before_anyone_reclaims(self, store): + """Seventh review: the check-to-POST expiry race. + + A claims, sleeps past its own lease, and wakes BEFORE any other + process reclaims. Its token is still in the row, so a read-only + ownership check passes — and then B reclaims while A's POST is in + flight: both send. The pre-POST renewal must instead REJECT a + claimant whose lease already expired, whether or not anyone has + reclaimed yet, because expiry alone means another process may claim + at any moment. + """ + _add_package(store, "pkg-1", "2026-08-26") + + posts = [] + sender_a = SharedMetricsSender( + store, ENDPOINT, + post=lambda e, p, *, timeout: (posts.append("A"), FakeResponse(202))[1], + sleep=lambda _s: None, + now=lambda: clock["t"], + ) + clock = {"t": NOW} + claimed = sender_a._claim_next(NOW, set()) + assert claimed is not None and not claimed["skip"] + + # Suspended past the 300s lease; wakes with the row NOT yet reclaimed. + clock["t"] = NOW + timedelta(seconds=400) + result = sender_a._send_one(claimed) + + assert posts == [], ( + "a claimant with an expired lease transmitted before renewal" + ) + assert result == "deferred" + # The row must remain claimable by the next process. + row = _row(store, "pkg-1") + assert row["send_state"] == "pending" + + def test_renewal_extends_the_lease_across_the_post(self, store): + """A healthy in-lease claimant renews and its POST is covered. + + Round-8 review: the original assertion was `>=` under a frozen + clock, which a renewal that matches the row but never extends the + lease also satisfies — the exact mutant that double-POSTs (the + un-extended lease expires mid-POST and a second process reclaims). + The renewal must move the deadline STRICTLY forward to now + lease, + so renew from a later clock and require the exact new deadline. + """ + _add_package(store, "pkg-1", "2026-08-26") + clock = {"t": NOW} + sender = SharedMetricsSender( + store, + ENDPOINT, + post=lambda e, p, *, timeout: FakeResponse(202), + sleep=lambda _s: None, + now=lambda: clock["t"], + ) + claimed = sender._claim_next(NOW, set()) + assert claimed is not None + lease_before = _row(store, "pkg-1")["next_attempt_at"] + + # 100s into the (300s) lease: still healthy, renews mid-flight. + clock["t"] = NOW + timedelta(seconds=100) + assert sender._renew_claim("pkg-1", claimed["claim_token"]) is True + lease_after = _row(store, "pkg-1")["next_attempt_at"] + assert lease_after > lease_before, ( + "renewal granted authority without extending the lease" + ) + # And not just 'later': the full fresh lease from the renewal clock. + expected = (NOW + timedelta(seconds=100 + 300)).strftime( + "%Y-%m-%dT%H:%M:%SZ" + ) + assert lease_after == expected + + def test_a_lapsed_claimant_resuming_after_reclaim_cannot_double_post( + self, store + ): + """PR-review P1: expiry -> reclaim -> old claimant resumes. + + A claims, then is suspended (laptop lid) BEFORE its POST. The lease + expires; B reclaims and POSTs; A wakes and proceeds. The pre-POST + ownership check must make A yield without transmitting. + + Scope note: the check closes the claim->POST gap. A suspension that + lands mid-POST (bytes already leaving) is not client-fixable — that + residual needs server-side dedupe and is documented on _send_one. + """ + _add_package(store, "pkg-1", "2026-08-26") + + posts = [] + + def post_a(endpoint, payload, *, timeout): + posts.append("A") + return FakeResponse(202) + + def post_b(endpoint, payload, *, timeout): + posts.append("B") + return FakeResponse(202) + + sender_a = SharedMetricsSender( + store, ENDPOINT, post=post_a, sleep=lambda _s: None, now=lambda: NOW + ) + # A claims, then the process is suspended before _send_one runs. + claimed_a = sender_a._claim_next(NOW, set()) + assert claimed_a is not None and not claimed_a["skip"] + + # 400s later (past the 300s lease) B claims and completes the send. + later = NOW + timedelta(seconds=400) + sender_b = SharedMetricsSender( + store, ENDPOINT, post=post_b, sleep=lambda _s: None, now=lambda: later + ) + outcome_b = sender_b.send_pending() + assert outcome_b.sent == 1 + + # A resumes exactly where it left off. + result_a = sender_a._send_one(claimed_a) + + row = _row(store, "pkg-1") + assert posts == ["B"], ( + f"a lapsed claimant transmitted after reclaim: {posts}" + ) + assert result_a == "deferred" + assert row["send_state"] == "sent", "B's settlement must stand" + + def test_a_lapsed_claimants_backoff_cannot_clobber_the_new_claim(self, store): + """The token must fence DEFERS too, not just the 202 settlement. + + A's transport fails after B has reclaimed; A's backoff write must + not move next_attempt_at under B's live lease. + """ + _add_package(store, "pkg-1", "2026-08-26") + sender_a = SharedMetricsSender( + store, ENDPOINT, + post=FakeTransport(OSError("net"), OSError("net"), OSError("net")), + sleep=lambda _s: None, now=lambda: NOW, + ) + claimed_a = sender_a._claim_next(NOW, set()) + assert claimed_a is not None and not claimed_a["skip"] + + later = NOW + timedelta(seconds=400) + sender_b = SharedMetricsSender( + store, ENDPOINT, post=FakeTransport(), + sleep=lambda _s: None, now=lambda: later, + ) + claimed_b = sender_b._claim_next(later, set()) + assert claimed_b is not None and not claimed_b["skip"] + lease_b = _row(store, "pkg-1")["next_attempt_at"] + + # A's exhausted retries try to write a 15-minute backoff. + result = sender_a._send_one(claimed_a) + assert result == "deferred" + assert _row(store, "pkg-1")["next_attempt_at"] == lease_b, ( + "a lapsed claimant's backoff overwrote the live claim's lease" + ) + + 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") + _sender(store, FakeTransport(OSError("killed"), OSError(""), OSError(""))).send_pending() + + later = SharedMetricsSender( + store, + ENDPOINT, + post=(transport := FakeTransport(FakeResponse(202))), + sleep=lambda _s: None, + now=lambda: NOW + timedelta(hours=2), + ) + later.send_pending() + assert len(transport.calls) == 1 + + def test_a_lapsed_sender_cannot_resurrect_a_sent_package(self, store): + """Terminal state must win over a straggler's write.""" + _add_package(store, "pkg-1", "2026-08-26") + _sender(store, FakeTransport(FakeResponse(202))).send_pending() + assert _row(store, "pkg-1")["send_state"] == "sent" + + # A straggler from an earlier pass tries to defer the same row. + _sender(store, FakeTransport())._defer("pkg-1", 600, "stale") + assert _row(store, "pkg-1")["send_state"] == "sent", ( + "a lapsed pass overwrote a completed send" + ) + + +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 + + @pytest.mark.parametrize( + "payload_json", + [ + '["a", "list"]', + "null", + '"a string"', + "42", + '{"no_install_id": true}', + '{"install_id": ""}', + '{"install_id": null}', + ], + ) + def test_valid_json_that_is_not_a_usable_package_is_skipped( + self, store, payload_json + ): + """Regression: a top-level array parsed fine, then .get() raised. + + The AttributeError escaped the claim transaction and blocked every + healthy package behind it. + """ + 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', + ?, '2026-08-26T00:00:00Z', '2026-08-26T01:00:00Z') + """, + (payload_json,), + ) + _add_package(store, "good", "2026-08-26") + + transport = FakeTransport(*[FakeResponse(202)] * 5) + outcome = _sender(store, transport).send_pending() + + assert outcome.sent == 1, "the healthy package must still go out" + assert [json.loads(c["payload"])["package_id"] for c in transport.calls] == [ + "good" + ] + assert _row(store, "bad")["send_state"] == "rejected" + + 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 TestConsentRevocation: + """`send: false` must stop an in-flight pass, not just the next one.""" + + def test_revoking_consent_mid_pass_stops_further_sends(self, store): + for i in range(4): + _add_package(store, f"pkg-{i}", "2026-08-26") + + consented = {"value": True} + posts = [] + + def transport(endpoint, payload, *, timeout): + posts.append(json.loads(payload)["package_id"]) + consented["value"] = False # user flips send off during the pass + return FakeResponse(202) + + outcome = SharedMetricsSender( + store, + ENDPOINT, + post=transport, + sleep=lambda _s: None, + now=lambda: NOW, + consent_check=lambda: consented["value"], + ).send_pending() + + assert len(posts) == 1, f"kept sending after consent was revoked: {posts}" + assert outcome.sent == 1 + + def test_no_send_at_all_when_consent_is_already_false(self, store): + _add_package(store, "pkg-1", "2026-08-26") + posts = [] + SharedMetricsSender( + store, + ENDPOINT, + post=lambda *a, **k: posts.append(1) or FakeResponse(202), + sleep=lambda _s: None, + now=lambda: NOW, + consent_check=lambda: False, + ).send_pending() + assert posts == [] + + def test_an_unreadable_consent_check_fails_closed(self, store): + """If consent cannot be established, do not transmit.""" + _add_package(store, "pkg-1", "2026-08-26") + posts = [] + + def explode(): + raise OSError("config unreadable") + + SharedMetricsSender( + store, + ENDPOINT, + post=lambda *a, **k: posts.append(1) or FakeResponse(202), + sleep=lambda _s: None, + now=lambda: NOW, + consent_check=explode, + ).send_pending() + assert posts == [] + + +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_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") + 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}' + captured = self._captured_request(payload) + assert captured["data"] == payload + assert "content-encoding" not in captured["headers"] diff --git a/tests/hermes_cli/test_shared_metrics_sender_e2e.py b/tests/hermes_cli/test_shared_metrics_sender_e2e.py new file mode 100644 index 0000000000..85be9b2388 --- /dev/null +++ b/tests/hermes_cli/test_shared_metrics_sender_e2e.py @@ -0,0 +1,282 @@ +"""End-to-end test: the real sender against a real HTTP server. + +Everything else stubs the transport. This exercises the actual code path — +urllib, gzip, headers, socket — against a live server on loopback, so a +transport-level mistake that a fake would hide fails here instead. +""" + +from __future__ import annotations + +import gzip +import json +import sqlite3 +import threading +from datetime import datetime, timezone +from http.server import BaseHTTPRequestHandler, HTTPServer + +import pytest + +from hermes_cli.observability.shared_metrics import SharedMetricsStore +from hermes_cli.observability.shared_metrics_sender import SharedMetricsSender + +INSTALL_ID = "12a73e97-4de9-4766-830d-9ca1192c0420" +NOW = datetime(2026, 8, 26, 12, 0, tzinfo=timezone.utc) + + +class Ingest(BaseHTTPRequestHandler): + """A stand-in for the ingest service that records what it receives.""" + + received: list = [] + script: list = [] + + def do_POST(self): # noqa: N802 - stdlib naming + length = int(self.headers.get("Content-Length") or 0) + raw = self.rfile.read(length) + if self.headers.get("Content-Encoding") == "gzip": + body = gzip.decompress(raw) + else: + body = raw + type(self).received.append( + { + "headers": {k.lower(): v for k, v in self.headers.items()}, + "body": json.loads(body.decode("utf-8")), + # Keep the RAW request bytes: comparing only the parsed body + # would not notice a non-deterministic transport encoding. + "raw": raw, + "raw_len": len(raw), + "decoded_len": len(body), + } + ) + status, payload, extra = ( + type(self).script.pop(0) if type(self).script else (202, {}, {}) + ) + encoded = json.dumps(payload).encode("utf-8") + self.send_response(status) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(encoded))) + for key, value in extra.items(): + self.send_header(key, value) + self.end_headers() + self.wfile.write(encoded) + + def log_message(self, format, *args): # noqa: A002 - stdlib signature + pass + + +@pytest.fixture +def server(): + Ingest.received = [] + Ingest.script = [] + httpd = HTTPServer(("127.0.0.1", 0), Ingest) + thread = threading.Thread(target=httpd.serve_forever, daemon=True) + thread.start() + yield httpd + httpd.shutdown() + httpd.server_close() + + +@pytest.fixture +def store(tmp_path): + built = SharedMetricsStore( + database_path=tmp_path / "metrics.sqlite3", + outbox_directory=tmp_path / "outbox", + ) + # Open a consent window covering the fixture packages; the interval gate + # fails closed without one, and this file tests transport, not consent. + from datetime import datetime, timezone + + from hermes_cli.observability.shared_metrics_sender import ( + reconcile_send_consent, + ) + from hermes_cli.sqlite_util import write_txn + + with built._connection() as connection: + with write_txn(connection): + reconcile_send_consent( + connection, True, now=datetime(2026, 8, 20, tzinfo=timezone.utc) + ) + reconcile_send_consent( + connection, True, now=datetime(2026, 10, 1, tzinfo=timezone.utc) + ) + return built + + +def _endpoint(server): + host, port = server.server_address + return f"http://{host}:{port}/v1/telemetry" + + +def _add(store, package_id, day="2026-08-26", metrics=1): + payload = { + "schema_version": "hermes.shared_metrics.v2", + "package_id": package_id, + "install_id": INSTALL_ID, + "generated_at": f"{day}T01:00:00Z", + "period_start": f"{day}T00:00:00Z", + "period_end": f"{day}T23:59:59Z", + "resource": { + "hermes_version": "0.20.5", + "os_family": "macos", + "architecture": "arm64", + "install_method": "git", + }, + "metrics": [ + { + "name": f"hermes.metric.{i}", + "type": "counter", + "dimensions": {"outcome": "ok"}, + "value": i, + } + for i in range(metrics) + ], + } + 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"{day}T00:00:00Z", + f"{day}T23:59:59Z", + json.dumps(payload), + f"{day}T01:00:00Z", + f"{day}T01:00:01Z", + ), + ) + return payload + + +def _sender(store, server): + return SharedMetricsSender( + store, _endpoint(server), sleep=lambda _s: None, now=lambda: NOW + ) + + +class TestRealTransport: + def test_a_package_is_delivered_and_marked_sent(self, store, server): + _add(store, "pkg-1") + outcome = _sender(store, server).send_pending() + + assert outcome.sent == 1 + assert len(Ingest.received) == 1 + assert Ingest.received[0]["body"]["package_id"] == "pkg-1" + + with store._connection() as connection: + state = connection.execute( + "SELECT send_state FROM package_outbox WHERE package_id = 'pkg-1'" + ).fetchone()[0] + assert state == "sent" + + def test_the_stable_install_id_crosses_the_wire_as_is(self, store, server): + """Product decision 2026-08-27: the raw install_id is transmitted.""" + _add(store, "pkg-1", metrics=40) + _sender(store, server).send_pending() + assert Ingest.received[0]["body"]["install_id"] == INSTALL_ID + + def test_content_type_is_json(self, store, server): + _add(store, "pkg-1") + _sender(store, server).send_pending() + assert Ingest.received[0]["headers"]["content-type"] == "application/json" + + def test_a_realistic_package_is_gzipped_over_the_wire(self, store, server): + # ~40 metrics matches the real outbox's larger packages. + _add(store, "pkg-1", metrics=120) + _sender(store, server).send_pending() + record = Ingest.received[0] + assert record["headers"].get("content-encoding") == "gzip" + assert record["raw_len"] < record["decoded_len"] + + def test_the_server_can_parse_what_we_send(self, store, server): + """Proves the bytes are valid JSON after transport and decompression.""" + original = _add(store, "pkg-1", metrics=120) + _sender(store, server).send_pending() + received = Ingest.received[0]["body"] + assert received["metrics"] == original["metrics"] + assert received["resource"] == original["resource"] + + def test_400_is_permanent(self, store, server): + _add(store, "pkg-1") + Ingest.script = [(400, {"error": "invalid_envelope"}, {})] + outcome = _sender(store, server).send_pending() + assert outcome.rejected == 1 + assert len(Ingest.received) == 1 + + def test_429_is_honoured(self, store, server): + _add(store, "pkg-1") + Ingest.script = [(429, {"error": "rate_limited"}, {"Retry-After": "90"})] + outcome = _sender(store, server).send_pending() + assert outcome.deferred == 1 + with store._connection() as connection: + retry_at = connection.execute( + "SELECT next_attempt_at FROM package_outbox WHERE package_id = 'pkg-1'" + ).fetchone()[0] + assert retry_at == "2026-08-26T12:01:30Z" + + def test_5xx_retries_then_succeeds(self, store, server): + _add(store, "pkg-1") + Ingest.script = [ + (503, {"error": "storage_unavailable"}, {}), + (202, {"package_id": "pkg-1"}, {}), + ] + outcome = _sender(store, server).send_pending() + assert outcome.sent == 1 + assert len(Ingest.received) == 2 + + def test_a_retry_sends_identical_bytes(self, store, server): + _add(store, "pkg-1", metrics=5) + Ingest.script = [(503, {}, {}), (202, {}, {})] + _sender(store, server).send_pending() + first, second = Ingest.received + assert first["body"] == second["body"] + assert first["raw"] == second["raw"], ( + "the raw request bytes must match, not just the parsed body" + ) + + def test_a_gzipped_retry_is_byte_identical_on_the_wire(self, store, server): + """gzip embeds an mtime by default, which would break this.""" + _add(store, "pkg-1", metrics=200) + Ingest.script = [(503, {}, {}), (202, {}, {})] + _sender(store, server).send_pending() + first, second = Ingest.received + assert first["headers"].get("content-encoding") == "gzip" + assert first["raw"] == second["raw"] + + def test_several_packages_in_one_pass(self, store, server): + for i in range(5): + _add(store, f"pkg-{i}") + outcome = _sender(store, server).send_pending() + assert outcome.sent == 5 + assert len(Ingest.received) == 5 + + def test_the_outbox_directory_is_untouched(self, store, server, tmp_path): + _add(store, "pkg-1") + marker = store.outbox_directory / "pkg-1.json" + marker.write_text('{"kept": true}') + _sender(store, server).send_pending() + assert marker.exists() + assert json.loads(marker.read_text()) == {"kept": True} + + def test_a_dead_server_defers_without_raising(self, store, server): + _add(store, "pkg-1") + host, port = server.server_address + server.shutdown() + server.server_close() + sender = SharedMetricsSender( + store, + f"http://{host}:{port}/v1/telemetry", + sleep=lambda _s: None, + now=lambda: NOW, + ) + outcome = sender.send_pending() + assert outcome.deferred == 1 + with store._connection() as connection: + state, error = connection.execute( + "SELECT send_state, last_error FROM package_outbox" + " WHERE package_id = 'pkg-1'" + ).fetchone() + assert state == "pending" + assert error diff --git a/tests/hermes_cli/test_shared_metrics_tools_toggle.py b/tests/hermes_cli/test_shared_metrics_tools_toggle.py new file mode 100644 index 0000000000..462bfcbd90 --- /dev/null +++ b/tests/hermes_cli/test_shared_metrics_tools_toggle.py @@ -0,0 +1,89 @@ +"""Tests for the `hermes tools` shared-metrics consent toggle. + +AGENTS.md requires outbound telemetry to be reachable from a config gate, the +setup prompt, AND `hermes tools`. These cover the third surface. +""" + +from __future__ import annotations + +import pytest + +from hermes_cli.tools_config import ( + _configure_shared_metrics_interactive, + _shared_metrics_menu_label, + _shared_metrics_state, +) + + +def _config(**shared): + return {"telemetry": {"shared_metrics": shared}} + + +class TestState: + def test_missing_telemetry_section_is_off(self): + assert _shared_metrics_state({}) == (False, False) + + def test_malformed_section_does_not_raise(self): + assert _shared_metrics_state({"telemetry": "nonsense"}) == (False, False) + + def test_reads_both_flags(self): + assert _shared_metrics_state(_config(enabled=True, send=True)) == (True, True) + + +class TestMenuLabel: + def test_off_state(self): + assert "off" in _shared_metrics_menu_label({}) + + def test_local_only_state(self): + label = _shared_metrics_menu_label(_config(enabled=True)) + assert "collecting locally" in label + assert "Nous" not in label + + def test_sending_state_names_the_destination(self): + label = _shared_metrics_menu_label(_config(enabled=True, send=True)) + assert "sending to Nous" in label + + +class TestToggle: + def test_enabling_send_persists(self, monkeypatch): + config = _config(enabled=True) + saved = {} + monkeypatch.setattr( + "hermes_cli.setup.prompt_yes_no", lambda *_a, **_k: True + ) + monkeypatch.setattr( + "hermes_cli.setup._record_send_consent_change", lambda **_k: None + ) + monkeypatch.setattr( + "hermes_cli.tools_config.save_config", + lambda cfg: saved.update({"cfg": cfg}), + ) + _configure_shared_metrics_interactive(config) + assert config["telemetry"]["shared_metrics"]["send"] is True + assert saved, "a consent change must be written to disk" + + def test_no_write_when_nothing_changed(self, monkeypatch): + config = _config(enabled=False, send=False) + saved = [] + monkeypatch.setattr( + "hermes_cli.setup.prompt_yes_no", lambda *_a, **_k: False + ) + monkeypatch.setattr( + "hermes_cli.tools_config.save_config", lambda cfg: saved.append(cfg) + ) + _configure_shared_metrics_interactive(config) + assert saved == [] + + def test_disabling_collection_also_disables_sending(self, monkeypatch): + """The toggle must not leave send=true with nothing to send.""" + config = _config(enabled=True, send=True) + monkeypatch.setattr( + "hermes_cli.setup.prompt_yes_no", lambda *_a, **_k: False + ) + monkeypatch.setattr( + "hermes_cli.tools_config.save_config", lambda cfg: None + ) + _configure_shared_metrics_interactive(config) + shared = config["telemetry"]["shared_metrics"] + assert shared["enabled"] is False + assert shared["send"] is False