diff --git a/tests/test_sse_replay.py b/tests/test_sse_replay.py index bca3d1d0..bdbc86f6 100644 --- a/tests/test_sse_replay.py +++ b/tests/test_sse_replay.py @@ -136,6 +136,33 @@ def test_replay_limit_is_bounded(): assert sse_router._bounded_replay_limit(25, city_count=5) == 25 +def test_replay_gap_direct_resync_policy(): + assert ( + sse_router._should_direct_resync( + since_revision=1000, + latest_revision=1200, + limit=60, + ) + is False + ) + assert ( + sse_router._should_direct_resync( + since_revision=1000, + latest_revision=1300, + limit=60, + ) + is True + ) + assert ( + sse_router._should_direct_resync( + since_revision=0, + latest_revision=1300, + limit=60, + ) + is False + ) + + def test_legacy_high_replay_limit_is_clamped_by_city_count(monkeypatch): captured = {} @@ -178,6 +205,51 @@ def test_legacy_high_replay_limit_is_clamped_by_city_count(monkeypatch): assert captured["resync_limit"] == 120 +def test_stale_client_gets_direct_resync_without_replay_scan(monkeypatch): + calls = {"replay": 0, "requires_resync": 0} + + class FakeStore: + def latest_revision(self): + return 2000 + + def replay_events(self, *, cities, since_revision, limit): + calls["replay"] += 1 + return [] + + def replay_requires_resync(self, *, cities, since_revision, replay_count, limit): + calls["requires_resync"] += 1 + return False + + async def finite_stream( + user_id, + *, + cities=None, + replay_events=None, + connected_revision=0, + resync_event=None, + ): + yield sse_router.sse_manager._format_event( + {"type": "connected", "revision": connected_revision} + ) + if resync_event: + yield sse_router.sse_manager._format_event(resync_event) + + monkeypatch.setattr(sse_router, "event_store", FakeStore()) + monkeypatch.setattr(sse_router.sse_manager, "event_stream", finite_stream) + + response = TestClient(app).get( + "/api/events?cities=ankara,buenos%20aires,istanbul,jeddah,seoul" + "&since_revision=100&replay_limit=500" + ) + + assert response.status_code == 200 + assert calls == {"replay": 0, "requires_resync": 0} + events = _decode_sse_events(response.text) + assert events[-1]["type"] == "resync_required" + assert events[-1]["reason"] == "replay_gap_too_large" + assert events[-1]["latest_revision"] == 2000 + + def test_ingest_patch_uses_external_fanout_without_direct_broadcast(monkeypatch): class FakeExternalStore: uses_external_live_fanout = True diff --git a/web/routers/sse_router.py b/web/routers/sse_router.py index 2001dadc..ee100139 100644 --- a/web/routers/sse_router.py +++ b/web/routers/sse_router.py @@ -22,6 +22,7 @@ _live_subscription_started = False SSE_REPLAY_BASE_LIMIT = 60 SSE_REPLAY_EVENTS_PER_CITY = 24 SSE_REPLAY_MAX_LIMIT = 240 +SSE_REPLAY_DIRECT_RESYNC_FACTOR = 3 def _parse_cities_param(cities: str) -> Set[str]: @@ -46,6 +47,21 @@ def _bounded_replay_limit(value: int, *, city_count: int = 0) -> int: return max(1, min(MAX_REPLAY_LIMIT, route_limit, limit)) +def _should_direct_resync( + *, + since_revision: int, + latest_revision: int, + limit: int, +) -> bool: + since = max(0, int(since_revision or 0)) + latest = max(0, int(latest_revision or 0)) + if since <= 0 or latest <= since: + return False + gap = latest - since + threshold = max(SSE_REPLAY_MAX_LIMIT, max(1, int(limit or 1)) * SSE_REPLAY_DIRECT_RESYNC_FACTOR) + return gap > threshold + + def _ensure_live_subscription() -> None: starter = getattr(event_store, "start_live_subscription", None) if not callable(starter): @@ -82,23 +98,36 @@ async def sse_events( if since_revision is not None: try: - replay_events = event_store.replay_events( - cities=city_set, - since_revision=max(0, int(since_revision)), - limit=limit, - ) - if event_store.replay_requires_resync( - cities=city_set, - since_revision=max(0, int(since_revision)), - replay_count=len(replay_events), + since = max(0, int(since_revision)) + if _should_direct_resync( + since_revision=since, + latest_revision=latest_revision, limit=limit, ): resync_event = { "type": "resync_required", - "reason": "replay_window_exceeded", + "reason": "replay_gap_too_large", "latest_revision": latest_revision, "ts": int(time.time() * 1000), } + else: + replay_events = event_store.replay_events( + cities=city_set, + since_revision=since, + limit=limit, + ) + if event_store.replay_requires_resync( + cities=city_set, + since_revision=since, + replay_count=len(replay_events), + limit=limit, + ): + resync_event = { + "type": "resync_required", + "reason": "replay_window_exceeded", + "latest_revision": latest_revision, + "ts": int(time.time() * 1000), + } except Exception: resync_event = { "type": "resync_required",