feat(telemetry): opt-in shared-metrics exporter (#95278)
feat(telemetry): opt-in shared-metrics exporter
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
# =============================================================================
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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",
|
||||
},
|
||||
},
|
||||
|
||||
|
||||
@@ -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():
|
||||
|
||||
@@ -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(
|
||||
"""
|
||||
|
||||
114
hermes_cli/observability/shared_metrics_send_config.py
Normal file
114
hermes_cli/observability/shared_metrics_send_config.py
Normal file
@@ -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
|
||||
792
hermes_cli/observability/shared_metrics_sender.py
Normal file
792
hermes_cli/observability/shared_metrics_sender.py
Normal file
@@ -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
|
||||
@@ -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)
|
||||
|
||||
|
||||
# =============================================================================
|
||||
|
||||
@@ -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)
|
||||
|
||||
198
scripts/e2e_shared_metrics_staging.py
Normal file
198
scripts/e2e_shared_metrics_staging.py
Normal file
@@ -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())
|
||||
@@ -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")
|
||||
|
||||
310
tests/hermes_cli/test_shared_metrics_consent_windows.py
Normal file
310
tests/hermes_cli/test_shared_metrics_consent_windows.py
Normal file
@@ -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"]
|
||||
165
tests/hermes_cli/test_shared_metrics_send_config.py
Normal file
165
tests/hermes_cli/test_shared_metrics_send_config.py
Normal file
@@ -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
|
||||
224
tests/hermes_cli/test_shared_metrics_send_migration.py
Normal file
224
tests/hermes_cli/test_shared_metrics_send_migration.py
Normal file
@@ -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 == []
|
||||
446
tests/hermes_cli/test_shared_metrics_send_wiring.py
Normal file
446
tests/hermes_cli/test_shared_metrics_send_wiring.py
Normal file
@@ -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"
|
||||
1029
tests/hermes_cli/test_shared_metrics_sender.py
Normal file
1029
tests/hermes_cli/test_shared_metrics_sender.py
Normal file
File diff suppressed because it is too large
Load Diff
282
tests/hermes_cli/test_shared_metrics_sender_e2e.py
Normal file
282
tests/hermes_cli/test_shared_metrics_sender_e2e.py
Normal file
@@ -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
|
||||
89
tests/hermes_cli/test_shared_metrics_tools_toggle.py
Normal file
89
tests/hermes_cli/test_shared_metrics_tools_toggle.py
Normal file
@@ -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
|
||||
Reference in New Issue
Block a user