diff --git a/routstr/nostr/analytics_runtime.py b/routstr/nostr/analytics_runtime.py index 79f7e69f..154f7a76 100644 --- a/routstr/nostr/analytics_runtime.py +++ b/routstr/nostr/analytics_runtime.py @@ -15,6 +15,7 @@ from ..core.terminal_outcomes import ( from ..core.vault import decrypt from .analytics_v2_delivery import ( AnalyticsV2Delivery, + AnalyticsV2DeliveryError, AnalyticsV2Producer, DeliveryStateSnapshot, SharingDisabledError, @@ -57,6 +58,7 @@ class AnalyticsCoordinator: self._relays: tuple[str, ...] = () self._writer_started = False self._retry_at = 0.0 + self._paused_reason: str | None = None self._closed = False async def prepare_startup(self) -> None: @@ -120,10 +122,8 @@ class AnalyticsCoordinator: else: try: provider_d = await resolve_provider_id_strict(pubkey, list(relays)) - except Exception: - await self._stop_public(disable=state.sharing_enabled) - self._retry_at = time.monotonic() + 60 - logger.exception("Stats need a stable provider identity before sharing") + except ValueError as error: + await self._pause(str(error), disable=state.sharing_enabled) return identity = (private_key, pubkey, provider_d) if ( @@ -134,6 +134,13 @@ class AnalyticsCoordinator: and not self._task.done() ): return + try: + delivery = AnalyticsV2Delivery(create_session, operator_relays=list(relays)) + except AnalyticsV2DeliveryError as error: + # Checked before activation, so a node that cannot deliver never + # opens and closes sharing on every retry. + await self._pause(str(error), disable=state.sharing_enabled) + return identity_changed = state.identity_pubkey is not None and ( state.identity_pubkey != pubkey or state.provider_d != provider_d @@ -165,15 +172,14 @@ class AnalyticsCoordinator: public_key_hex=pubkey, provider_d=provider_d, ) - self._delivery = AnalyticsV2Delivery( - create_session, operator_relays=list(relays) - ) + self._delivery = delivery self._identity = identity self._relays = relays self._task = asyncio.create_task( - run_analytics_v2_publisher(producer, self._delivery), + run_analytics_v2_publisher(producer, delivery), name="analytics-v2-publisher", ) + self._paused_reason = None except SharingDisabledError: return except Exception: @@ -181,6 +187,14 @@ class AnalyticsCoordinator: self._retry_at = time.monotonic() + 60 raise + async def _pause(self, reason: str, *, disable: bool) -> None: + await self._stop_public(disable=disable) + self._retry_at = time.monotonic() + 60 + # Retried every minute; warn again only when the cause changes. + if reason != self._paused_reason: + self._paused_reason = reason + logger.warning("Stats sharing paused", extra={"reason": reason}) + async def _stop_task(self) -> None: task, self._task = self._task, None if task is not None: diff --git a/tests/unit/test_analytics_runtime.py b/tests/unit/test_analytics_runtime.py index 2ca9282e..7b8b1163 100644 --- a/tests/unit/test_analytics_runtime.py +++ b/tests/unit/test_analytics_runtime.py @@ -3,6 +3,7 @@ from __future__ import annotations import asyncio import importlib import json +import logging import time from collections.abc import AsyncIterator from contextlib import asynccontextmanager @@ -386,3 +387,48 @@ async def test_upgrade_preserves_existing_sharing_choice_across_restart( assert node.events == [] finally: await restarted.close() + + +@pytest.mark.asyncio +async def test_relays_without_a_public_target_never_activate_sharing( + node: Any, monkeypatch: Any +) -> None: + monkeypatch.setattr(runtime.settings, "enable_analytics_sharing", True) + monkeypatch.setattr(runtime.settings, "relays", ["ws://umbrel.local:4848"]) + await node.coordinator.prepare_startup() + for _ in range(2): + node.coordinator._retry_at = 0 + await node.coordinator.sync_once() + + state = await runtime.get_analytics_v2_delivery_state(node.sessions) + assert (state.sharing_enabled, state.generation) == (False, 0) + assert node.coordinator._task is None + assert node.writer.running + + +@pytest.mark.asyncio +async def test_missing_provider_identity_warns_once( + node: Any, monkeypatch: Any, caplog: pytest.LogCaptureFixture +) -> None: + async def unlisted(*args: object) -> str: + raise ValueError( + "Configure PROVIDER_ID to select one provider for public stats" + ) + + monkeypatch.setattr(runtime.settings, "enable_analytics_sharing", True) + monkeypatch.setattr(runtime.settings, "provider_id", "") + monkeypatch.setattr(runtime, "resolve_provider_id_strict", unlisted) + runtime_logger = logging.getLogger(runtime.__name__) + runtime_logger.addHandler(caplog.handler) + try: + await node.coordinator.prepare_startup() + for _ in range(3): + node.coordinator._retry_at = 0 + await node.coordinator.sync_once() + finally: + runtime_logger.removeHandler(caplog.handler) + + records = [r for r in caplog.records if r.name == runtime.__name__] + assert [r.levelno for r in records] == [logging.WARNING] + assert "PROVIDER_ID" in records[0].__dict__["reason"] + assert node.coordinator._task is None