refactor(hermes_cli): shared_metrics store — _write() txn ctx, resource-column table, private-path helper

This commit is contained in:
Teknium
2026-09-02 20:33:40 -07:00
parent 4b121f6124
commit da5cc8d1fe

View File

@@ -32,6 +32,8 @@ _SCHEMA_BUSY_TIMEOUT_MS = 5_000
_LOCAL_HISTORY_RETENTION_DAYS = 30
_ACTIVE_INSTALL_STATE_KEY = "client_active_recorded_at"
_ACTIVE_INSTALL_INTERVAL = timedelta(hours=24)
# Column order of the client resource in every counter_aggregates statement.
_RESOURCE_COLUMNS = ("hermes_version", "os_family", "architecture", "install_method")
logger = logging.getLogger(__name__)
@@ -44,6 +46,25 @@ def _isoformat(value: datetime) -> str:
return value.astimezone(timezone.utc).isoformat().replace("+00:00", "Z")
def _compact_json(value: Any) -> str:
return json.dumps(value, sort_keys=True, separators=(",", ":"))
def _row_resource(row: sqlite3.Row) -> dict[str, str]:
return {column: row[column] for column in _RESOURCE_COLUMNS}
def _ensure_private(path: Path, mode: int) -> None:
if mode == 0o700:
path.mkdir(parents=True, exist_ok=True, mode=mode)
else:
path.touch(mode=mode, exist_ok=True)
try:
path.chmod(mode)
except OSError:
pass
class SharedMetricsStore:
"""Persist allowlisted counters and export immutable delta packages."""
@@ -55,16 +76,12 @@ class SharedMetricsStore:
root = get_hermes_home() / "telemetry" / "shared_metrics"
self.database_path = database_path or root / "metrics.sqlite3"
self.outbox_directory = outbox_directory or root / "outbox"
self._ensure_private_directory(self.database_path.parent)
self._ensure_private_directory(self.outbox_directory)
self._ensure_private_file(self.database_path)
_ensure_private(self.database_path.parent, 0o700)
_ensure_private(self.outbox_directory, 0o700)
_ensure_private(self.database_path, 0o600)
self._ensure_schema()
def record_model_call(
self,
dimensions: dict[str, str],
resource: dict[str, str],
) -> None:
def record_model_call(self, dimensions: dict[str, str], resource: dict[str, str]) -> None:
"""Increment the terminal model-call counter for the current UTC day."""
self.record_counter(MODEL_ROUTE_METRIC, dimensions, resource)
@@ -73,72 +90,57 @@ class SharedMetricsStore:
dimensions: dict[str, str] = {}
self._validate_counter(CLIENT_ACTIVE_METRIC, dimensions, resource)
now = _utc_now()
with self._connection() as connection:
with write_txn(connection):
row = connection.execute(
"SELECT value FROM telemetry_state WHERE key = ?",
(_ACTIVE_INSTALL_STATE_KEY,),
).fetchone()
if row is not None:
last_recorded = self._parse_state_timestamp(row["value"])
if last_recorded is not None and last_recorded > now:
# A wall-clock correction must not suppress activity until
# the stale future timestamp plus another full interval.
connection.execute(
"""
with self._write() as connection:
row = connection.execute(
"SELECT value FROM telemetry_state WHERE key = ?",
(_ACTIVE_INSTALL_STATE_KEY,),
).fetchone()
last_recorded = self._parse_state_timestamp(row["value"]) if row is not None else None
if last_recorded is not None:
if last_recorded > now:
# A wall-clock correction must not suppress activity until the
# stale future timestamp plus another full interval.
connection.execute(
"""
UPDATE telemetry_state
SET value = ?
WHERE key = ?
""",
(_isoformat(now), _ACTIVE_INSTALL_STATE_KEY),
)
return False
if (
last_recorded is not None
and now < last_recorded + _ACTIVE_INSTALL_INTERVAL
):
return False
(_isoformat(now), _ACTIVE_INSTALL_STATE_KEY),
)
return False
if now < last_recorded + _ACTIVE_INSTALL_INTERVAL:
return False
self._install_id(connection)
self._record_counter_in_transaction(
connection,
CLIENT_ACTIVE_METRIC,
dimensions,
resource,
period_start=now.date().isoformat(),
)
connection.execute(
"""
self._install_id(connection)
self._record_counter_in_transaction(
connection, CLIENT_ACTIVE_METRIC, dimensions, resource,
period_start=now.date().isoformat(),
)
connection.execute(
"""
INSERT INTO telemetry_state(key, value)
VALUES (?, ?)
ON CONFLICT(key) DO UPDATE SET value = excluded.value
""",
(_ACTIVE_INSTALL_STATE_KEY, _isoformat(now)),
)
(_ACTIVE_INSTALL_STATE_KEY, _isoformat(now)),
)
return True
def record_counter(
self,
metric_name: str,
dimensions: dict[str, str],
resource: dict[str, str],
self, metric_name: str, dimensions: dict[str, str], resource: dict[str, str]
) -> None:
"""Increment one allowlisted counter for the current UTC day."""
self._validate_counter(metric_name, dimensions, resource)
with self._connection() as connection:
self._record_counter_in_transaction(
connection,
metric_name,
dimensions,
resource,
connection, metric_name, dimensions, resource,
period_start=_utc_now().date().isoformat(),
)
@staticmethod
def _validate_counter(
metric_name: str,
dimensions: dict[str, str],
resource: dict[str, str],
metric_name: str, dimensions: dict[str, str], resource: dict[str, str]
) -> None:
if metric_name not in COUNTER_METRICS:
raise ValueError(f"Unsupported shared metric: {metric_name}")
@@ -156,11 +158,6 @@ class SharedMetricsStore:
*,
period_start: str,
) -> None:
dimensions_json = json.dumps(
dimensions,
sort_keys=True,
separators=(",", ":"),
)
connection.execute(
"""
INSERT INTO counter_aggregates(
@@ -188,18 +185,14 @@ class SharedMetricsStore:
(
period_start,
metric_name,
resource["hermes_version"],
resource["os_family"],
resource["architecture"],
resource["install_method"],
dimensions_json,
*(resource[column] for column in _RESOURCE_COLUMNS),
_compact_json(dimensions),
),
)
def create_and_export_package(self) -> list[Path]:
"""Commit one pending delta package, then atomically export the outbox."""
pending_periods = self._pending_period_count()
for _ in range(pending_periods):
for _ in range(self._pending_period_count()):
if self._create_package() is None:
break
return self._export_and_prune()
@@ -214,10 +207,7 @@ class SharedMetricsStore:
try:
self._prune_expired_history()
except Exception:
logger.warning(
"Unable to prune expired shared-metrics history",
exc_info=True,
)
logger.warning("Unable to prune expired shared-metrics history", exc_info=True)
return exported
def counter_snapshot(self) -> list[dict[str, Any]]:
@@ -250,12 +240,7 @@ class SharedMetricsStore:
{
"period_start": row["period_start"],
"metric_name": row["metric_name"],
"resource": {
"hermes_version": row["hermes_version"],
"os_family": row["os_family"],
"architecture": row["architecture"],
"install_method": row["install_method"],
},
"resource": _row_resource(row),
"dimensions": json.loads(row["dimensions_json"]),
"value": row["value"],
"packaged_value": row["packaged_value"],
@@ -265,14 +250,9 @@ class SharedMetricsStore:
@contextmanager
def _connection(
self,
*,
busy_timeout_ms: int = _BUSY_TIMEOUT_MS,
self, *, busy_timeout_ms: int = _BUSY_TIMEOUT_MS
) -> Iterator[sqlite3.Connection]:
connection = sqlite3.connect(
self.database_path,
timeout=busy_timeout_ms / 1000,
)
connection = sqlite3.connect(self.database_path, timeout=busy_timeout_ms / 1000)
try:
connection.row_factory = sqlite3.Row
connection.execute(f"PRAGMA busy_timeout={busy_timeout_ms}")
@@ -281,27 +261,17 @@ class SharedMetricsStore:
finally:
connection.close()
@staticmethod
def _ensure_private_directory(path: Path) -> None:
path.mkdir(parents=True, exist_ok=True, mode=0o700)
try:
path.chmod(0o700)
except OSError:
pass
@staticmethod
def _ensure_private_file(path: Path) -> None:
path.touch(mode=0o600, exist_ok=True)
try:
path.chmod(0o600)
except OSError:
pass
@contextmanager
def _write(self, *, busy_timeout_ms: int = _BUSY_TIMEOUT_MS) -> Iterator[sqlite3.Connection]:
"""A connection inside one write transaction."""
with self._connection(busy_timeout_ms=busy_timeout_ms) as connection:
with write_txn(connection):
yield connection
def _ensure_schema(self) -> None:
with self._connection(busy_timeout_ms=_SCHEMA_BUSY_TIMEOUT_MS) as connection:
# Serialize first-run creation and upgrades across Hermes processes.
with write_txn(connection):
self._ensure_schema_in_transaction(connection)
# Serialize first-run creation and upgrades across Hermes processes.
with self._write(busy_timeout_ms=_SCHEMA_BUSY_TIMEOUT_MS) as connection:
self._ensure_schema_in_transaction(connection)
@staticmethod
def _ensure_schema_in_transaction(connection: sqlite3.Connection) -> None:
@@ -352,11 +322,10 @@ class SharedMetricsStore:
def _add_send_columns(connection: sqlite3.Connection) -> None:
"""Add transmission bookkeeping to ``package_outbox``, idempotently.
The columns are ADDITIVE and nullable, and the schema version is deliberately NOT bumped:
``_ensure_schema_in_transaction`` raises on unknown versions with no forward-compat branch,
so a bump would hard-fail an older Hermes (second profile, rollback) against the same
database. Old readers select named columns, never ``SELECT *``, so extra columns are
invisible to them.
Additive and nullable; the schema version is deliberately NOT bumped because
``_ensure_schema_in_transaction`` raises on unknown versions, so a bump would
hard-fail an older Hermes (second profile, rollback) sharing the database. Old
readers select named columns, never ``SELECT *``.
"""
existing = {
str(row["name"])
@@ -371,17 +340,12 @@ class SharedMetricsStore:
# 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.
# The install_id actually transmitted, frozen on the first attempt so
# retries stay byte-identical; the body is recomputed from payload_json.
("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).
# Rewritten on every claim. Settlement and the pre-POST revalidation are
# compare-and-set on it, so a claimant whose lease lapsed loses authority
# the moment another process reclaims (else a suspended sender double-POSTs).
("claim_token", "TEXT"),
):
if column not in existing:
@@ -391,15 +355,11 @@ class SharedMetricsStore:
@staticmethod
def _add_consent_tables(connection: sqlite3.Connection) -> None:
"""Create the consent-window tables, idempotently.
"""Create the consent-window tables, idempotently (additive; version not bumped).
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.
``send_consent_windows`` records consent as explicit intervals rather than a moving
day-stamp: opened when send consent is observed, heartbeat-confirmed on every later
observation, closed at the LAST CONFIRMED moment (never "now") on withdrawal.
"""
connection.execute(
"""
@@ -481,19 +441,14 @@ class SharedMetricsStore:
connection.execute("DROP TABLE counter_aggregates_v1")
def _install_id(self, connection: sqlite3.Connection) -> str:
row = connection.execute(
"SELECT value FROM telemetry_state WHERE key = 'install_id'"
).fetchone()
if row is not None:
return str(row["value"])
candidate = str(uuid.uuid4())
connection.execute(
"INSERT OR IGNORE INTO telemetry_state(key, value) VALUES ('install_id', ?)",
(candidate,),
)
row = connection.execute(
"SELECT value FROM telemetry_state WHERE key = 'install_id'"
).fetchone()
query = "SELECT value FROM telemetry_state WHERE key = 'install_id'"
row = connection.execute(query).fetchone()
if row is None:
connection.execute(
"INSERT OR IGNORE INTO telemetry_state(key, value) VALUES ('install_id', ?)",
(str(uuid.uuid4()),),
)
row = connection.execute(query).fetchone()
if row is None:
raise RuntimeError("Unable to create the shared-metrics install identity")
return str(row["value"])
@@ -504,9 +459,7 @@ class SharedMetricsStore:
parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
except (TypeError, ValueError):
return None
if parsed.tzinfo is None:
return None
return parsed.astimezone(timezone.utc)
return None if parsed.tzinfo is None else parsed.astimezone(timezone.utc)
def _pending_period_count(self) -> int:
with self._connection() as connection:
@@ -535,34 +488,29 @@ class SharedMetricsStore:
def _create_pending_packages_if_due(self) -> None:
now = _utc_now()
with self._connection() as connection:
with write_txn(connection):
# Gate on the committed package, not its file write, so a failed
# outbox export can be retried without packaging deltas twice.
package_created_today = connection.execute(
"""
with self._write() as connection:
# Gate on the committed package, not its file write, so a failed
# outbox export can be retried without packaging deltas twice.
package_created_today = connection.execute(
"""
SELECT 1
FROM package_outbox
WHERE substr(created_at, 1, 10) >= ?
LIMIT 1
""",
(now.date().isoformat(),),
).fetchone()
if package_created_today is not None:
return
while self._create_package_in_transaction(connection, now) is not None:
pass
(now.date().isoformat(),),
).fetchone()
if package_created_today is not None:
return
while self._create_package_in_transaction(connection, now) is not None:
pass
def _create_package(self) -> dict[str, Any] | None:
now = _utc_now()
with self._connection() as connection:
with write_txn(connection):
return self._create_package_in_transaction(connection, now)
with self._write() as connection:
return self._create_package_in_transaction(connection, _utc_now())
def _create_package_in_transaction(
self,
connection: sqlite3.Connection,
now: datetime,
self, connection: sqlite3.Connection, now: datetime
) -> dict[str, Any] | None:
period_row = connection.execute(
"""
@@ -583,10 +531,12 @@ class SharedMetricsStore:
LIMIT 1
"""
).fetchone()
period_value = period_row["period_start"] if period_row is not None else None
if not period_value:
if period_row is None or not period_row["period_start"]:
return None
period_value = period_row["period_start"]
resource = _row_resource(period_row)
resource_values = tuple(resource[column] for column in _RESOURCE_COLUMNS)
rows = connection.execute(
"""
SELECT metric_name, dimensions_json, value, packaged_value
@@ -599,42 +549,21 @@ class SharedMetricsStore:
AND value > packaged_value
ORDER BY metric_name, dimensions_json
""",
(
period_value,
period_row["hermes_version"],
period_row["os_family"],
period_row["architecture"],
period_row["install_method"],
),
(period_value, *resource_values),
).fetchall()
period_start = datetime.fromisoformat(str(period_value)).replace(
tzinfo=timezone.utc
)
period_end = period_start + timedelta(days=1)
package_id = str(uuid.uuid4())
resource = {
"hermes_version": period_row["hermes_version"],
"os_family": period_row["os_family"],
"architecture": period_row["architecture"],
"install_method": period_row["install_method"],
}
period_start = datetime.fromisoformat(str(period_value)).replace(tzinfo=timezone.utc)
if not client_resource_is_valid(resource):
raise ValueError("Unsupported shared-metrics client resource")
payload = {
"schema_version": _PACKAGE_SCHEMA_VERSION,
"package_id": package_id,
"package_id": str(uuid.uuid4()),
"install_id": self._install_id(connection),
"period_start": _isoformat(period_start),
"period_end": _isoformat(period_end),
"period_end": _isoformat(period_start + timedelta(days=1)),
"generated_at": _isoformat(now),
"resource": resource,
"metrics": [self._package_metric(row) for row in rows],
}
payload_json = json.dumps(
payload,
sort_keys=True,
separators=(",", ":"),
)
connection.execute(
"""
INSERT INTO package_outbox(
@@ -646,16 +575,16 @@ class SharedMetricsStore:
) VALUES (?, ?, ?, ?, ?)
""",
(
package_id,
payload["package_id"],
payload["period_start"],
payload["period_end"],
payload_json,
_compact_json(payload),
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.
# 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', ?)
@@ -676,15 +605,7 @@ class SharedMetricsStore:
AND install_method = ?
AND dimensions_json = ?
""",
(
period_value,
row["metric_name"],
period_row["hermes_version"],
period_row["os_family"],
period_row["architecture"],
period_row["install_method"],
row["dimensions_json"],
),
(period_value, row["metric_name"], *resource_values, row["dimensions_json"]),
)
return payload
@@ -719,11 +640,7 @@ class SharedMetricsStore:
package_id = str(row["package_id"])
path = self.outbox_directory / f"{package_id}.json"
atomic_json_write(
path,
json.loads(row["payload_json"]),
indent=2,
sort_keys=True,
mode=0o600,
path, json.loads(row["payload_json"]), indent=2, sort_keys=True, mode=0o600
)
with self._connection() as connection:
connection.execute(
@@ -739,11 +656,8 @@ class SharedMetricsStore:
def _prune_expired_history(self, *, now: datetime | None = None) -> None:
"""Remove exported local history after the bounded retention window."""
cutoff = (now or _utc_now()) - timedelta(
days=_LOCAL_HISTORY_RETENTION_DAYS
)
cutoff = (now or _utc_now()) - timedelta(days=_LOCAL_HISTORY_RETENTION_DAYS)
cutoff_timestamp = _isoformat(cutoff)
cutoff_period = cutoff.date().isoformat()
with self._connection() as connection:
rows = connection.execute(
"""
@@ -760,32 +674,27 @@ class SharedMetricsStore:
for row in rows:
package_id = str(row["package_id"])
try:
(self.outbox_directory / f"{package_id}.json").unlink(
missing_ok=True
)
(self.outbox_directory / f"{package_id}.json").unlink(missing_ok=True)
except OSError:
logger.warning(
"Unable to prune expired shared-metrics package %s",
package_id,
exc_info=True,
"Unable to prune expired shared-metrics package %s", package_id, exc_info=True
)
continue
removable_package_ids.append(package_id)
with self._connection() as connection:
with write_txn(connection):
for package_id in removable_package_ids:
connection.execute(
"""
with self._write() as connection:
for package_id in removable_package_ids:
connection.execute(
"""
DELETE FROM package_outbox
WHERE package_id = ?
AND exported_at IS NOT NULL
AND exported_at < ?
""",
(package_id, cutoff_timestamp),
)
connection.execute(
"""
(package_id, cutoff_timestamp),
)
connection.execute(
"""
DELETE FROM counter_aggregates
WHERE period_start < ?
AND value = packaged_value
@@ -797,5 +706,5 @@ class SharedMetricsStore:
= counter_aggregates.period_start
)
""",
(cutoff_period,),
)
(cutoff.date().isoformat(),),
)