Direct resync stale SSE clients
This commit is contained in:
@@ -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
|
||||
|
||||
+39
-10
@@ -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",
|
||||
|
||||
Reference in New Issue
Block a user