import math from datetime import UTC, datetime from statistics import median from unittest import mock import pytest from fansifter_common.utils import timezone from freezegun import freeze_time from pytest_mock import MockerFixture from resonance_engine.config import settings from resonance_engine.dsp.enums import DSPClientName, DSPClientStatus, DSPId from resonance_engine.dsp.models import DSPClient from resonance_engine.dsp.stats import RequestStats from resonance_engine.fandata.enums import FanConnectionStatus from resonance_engine.fandata.models import FanCollectionState, FanConnection from resonance_engine.tasks.enums import FanoutSource, TaskStatus from resonance_engine.tasks.models import CollectTask, FanoutTask from resonance_engine.tasks.planner import ( REQUESTS_PER_FAN, SAFETY_FACTOR, compute_fanout_plan, drain_budget_s, plan_fanout, ) from tests.unit.helpers import create_model, create_model_batch, override_settings class TestComputeFanoutPlan: def test_cold_start_seeds_from_nominal(self) -> None: plan = compute_fanout_plan(11) assert plan.fanout_count == int( 11 * SAFETY_FACTOR * settings.fan_fanout_window_s / REQUESTS_PER_FAN ) assert plan.fanout_count_source == "seed" assert plan.rps is None def test_zero_nominal_returns_zero_count(self) -> None: plan = compute_fanout_plan(0) assert plan.fanout_count == 0 assert plan.fanout_count_source == "seed" def test_partial_signals_fall_back_to_seed(self) -> None: # Both rps and reqs_per_fan are needed to size warm; missing/zero → seed. assert compute_fanout_plan(11, rps=10).fanout_count_source == "seed" assert compute_fanout_plan(11, reqs_per_fan=5).fanout_count_source == "seed" for bad in (0, None): plan = compute_fanout_plan(11, rps=bad, reqs_per_fan=5) assert plan.fanout_count_source == "seed" def test_warm_sizes_by_fan_drain_rate(self) -> None: # floor(2 * 0.9 * 600) = 1080 (fans_per_s = 2 fans/s) plan = compute_fanout_plan(100, fans_per_s=2, window_s=600) assert plan.fanout_count == 1080 assert plan.fanout_count_source == "collected" def test_higher_drain_rate_dispatches_more(self) -> None: low = compute_fanout_plan(100, fans_per_s=1, window_s=600).fanout_count high = compute_fanout_plan(100, fans_per_s=2, window_s=600).fanout_count assert high > low def test_slower_drain_shrinks_dispatch(self) -> None: # an overrun lowers fans/span → fewer fans next run (self-correction) fast = compute_fanout_plan(100, fans_per_s=2, window_s=600).fanout_count slow = compute_fanout_plan(100, fans_per_s=1, window_s=600).fanout_count assert slow < fast def test_overran_flag_set_when_drain_exceeds_window(self) -> None: plan = compute_fanout_plan( 100, rps=10, reqs_per_fan=5, window_s=600, last_run_duration_s=1080 ) assert plan.overran is True assert plan.last_run_duration_s == 1080 def test_overrun_damper_caps_growth_at_last_fans(self) -> None: # rate would grow dispatch to 54000, but the last run overran → cap at # what it actually worked (no piling onto a backlog). plan = compute_fanout_plan( 1, fans_per_s=100, last_fans=500, window_s=600, last_run_duration_s=1200, ) assert plan.overran is True assert plan.fanout_count == 500 # capped, not 100 * 0.9 * 600 = 54000 def test_no_damper_when_within_window(self) -> None: # same rate but the last run fit the window → grow freely, no cap. plan = compute_fanout_plan( 1, fans_per_s=100, last_fans=500, window_s=600, last_run_duration_s=300, ) assert plan.overran is False assert plan.fanout_count == 54000 def test_headroom_subtracts_pending(self) -> None: # 54000 sized, 4000 still pending → dispatch only the free capacity. plan = compute_fanout_plan(1, fans_per_s=100, window_s=600, pending=4000) assert plan.fanout_count == 50000 assert plan.pending == 4000 def test_headroom_floors_at_zero_when_saturated(self) -> None: # pending already exceeds a window's capacity → dispatch nothing. plan = compute_fanout_plan(1, fans_per_s=100, window_s=600, pending=60000) assert plan.fanout_count == 0 class TestFanoutControllerConvergence: """Forward-sim of the closed loop: a fixed-capacity fleet, fanouts every window. Proves the fans_per_s controller is bounded and converges toward ~window-fit drain instead of the runaway 300-500% overrun it replaced.""" def test_converges_from_overshoot(self) -> None: window = 300 capacity = 50.0 # fans/sec the fleet drains (fixed) drained_per_tick = capacity * window keep = 3 # _SIGNAL_FANOUTS — median window backlog = 0.0 history: list[tuple[int, float]] = [] # (fans, span), newest last fanout_count = 30_000 # massive initial overshoot spans: list[float] = [] for _ in range(40): # this fanout's fans finish after the backlog ahead of them + itself span = (backlog + fanout_count) / capacity spans.append(span) history.append((fanout_count, span)) history = history[-keep:] # a window elapses before the next tick; the fleet drains what it can backlog = max(0.0, backlog + fanout_count - drained_per_tick) fans_per_s = median([f / s for f, s in history]) plan = compute_fanout_plan( 1, fans_per_s=fans_per_s, last_fans=history[-1][0], window_s=window, last_run_duration_s=history[-1][1], ) fanout_count = plan.fanout_count # bounded — never exceeds the initial overshoot (no runaway) assert max(spans) <= spans[0] # converged near the FILL_SAFETY target (~0.9 window), not 300-500% assert window * 0.7 <= spans[-1] <= window * 1.1 @pytest.mark.db class TestFanoutSignalsExcludesRunning: """Regression: a running fanout has the full dispatch in fans_total but only partial drain elapsed; counting it would inflate fans/span (≈ overlap factor) and the controller would never see the true full-drain rate → persistent overrun. Only FINISHED collect tasks must count.""" CLIENT = DSPClientName.spotify_songwhip @freeze_time("2026-06-02T12:00:00Z") def test_running_only_yields_no_signal(self) -> None: # A fresh fanout, 5m in, full 18k dispatch, nothing finished yet. fanout = create_model( FanoutTask, source=FanoutSource.scheduled, status=TaskStatus.running, started_at=datetime(2026, 6, 2, 11, 55, tzinfo=UTC), ) create_model( CollectTask, fanout_task_id=fanout.id, dsp_client_name=self.CLIENT, started_at=datetime(2026, 6, 2, 11, 55, tzinfo=UTC), finished_at=None, status=TaskStatus.running, fans_total=18_000, requests=5_000, ) assert CollectTask.query.fanout_signals(self.CLIENT) is None @freeze_time("2026-06-02T12:00:00Z") def test_uses_finished_rate_ignoring_running_sibling(self) -> None: # Older fanout actually drained: 60 fans in 1200s → 0.05 fans/s (truth). done = create_model( FanoutTask, source=FanoutSource.scheduled, status=TaskStatus.done, started_at=datetime(2026, 6, 2, 11, 30, tzinfo=UTC), ) create_model( CollectTask, fanout_task_id=done.id, dsp_client_name=self.CLIENT, started_at=datetime(2026, 6, 2, 11, 30, tzinfo=UTC), finished_at=datetime(2026, 6, 2, 11, 50, tzinfo=UTC), # 1200s status=TaskStatus.done, fans_total=60, requests=540, ) # Newer fanout still running with a huge dispatch — must NOT inflate. running = create_model( FanoutTask, source=FanoutSource.scheduled, status=TaskStatus.running, started_at=datetime(2026, 6, 2, 11, 55, tzinfo=UTC), ) create_model( CollectTask, fanout_task_id=running.id, dsp_client_name=self.CLIENT, started_at=datetime(2026, 6, 2, 11, 55, tzinfo=UTC), finished_at=None, status=TaskStatus.running, fans_total=99_999, requests=99_999, ) signals = CollectTask.query.fanout_signals(self.CLIENT) assert signals is not None # true rate from the finished fanout only — the running 99,999 is ignored assert signals.fans_per_s == 60 / 1200 # 0.05 assert signals.duration_s == 1200 assert signals.last_fans == 60 @pytest.mark.db class TestDrainBudget: def test_scheduled_uses_full_window(self) -> None: with override_settings(fan_fanout_window_s=600): assert drain_budget_s(FanoutSource.scheduled) == 600 def test_off_cadence_without_scheduled_history_uses_full_window(self) -> None: with override_settings(fan_fanout_window_s=600): assert drain_budget_s(FanoutSource.manual) == 600 @freeze_time("2026-01-01T00:05:00Z") def test_off_cadence_clamps_to_time_until_next_tick(self) -> None: create_model( FanoutTask, source=FanoutSource.scheduled, started_at=datetime(2026, 1, 1, tzinfo=UTC), status=TaskStatus.done, ) with override_settings(fan_fanout_window_s=600): assert drain_budget_s(FanoutSource.manual) == 300 @freeze_time("2026-01-01T00:23:00Z") def test_off_cadence_clamp_survives_missed_ticks(self) -> None: create_model( FanoutTask, source=FanoutSource.scheduled, started_at=datetime(2026, 1, 1, tzinfo=UTC), status=TaskStatus.done, ) with override_settings(fan_fanout_window_s=600): assert drain_budget_s(FanoutSource.triggered) == 420 @pytest.mark.db class TestPlanFanout: @pytest.fixture(autouse=True) def gateway_mock(self, mocker: MockerFixture) -> mock.MagicMock: m = mocker.patch( "resonance_engine.tasks.planner.dsp_gateway", new_callable=mock.MagicMock, ) m.is_configured.return_value = True m.stats.return_value = RequestStats() return m def test_skips_clients_without_nominal_rps(self) -> None: create_model(DSPClient, name=DSPClientName.spotify_songwhip, nominal_rps=None) results = plan_fanout() assert results == [] def test_skips_paused_clients(self) -> None: create_model( DSPClient, name=DSPClientName.spotify_songwhip, nominal_rps=11, status=DSPClientStatus.paused, ) results = plan_fanout() assert results == [] def test_skips_clients_without_gateway(self, gateway_mock: mock.MagicMock) -> None: gateway_mock.is_configured.return_value = False create_model(DSPClient, name=DSPClientName.spotify_songwhip, nominal_rps=11) results = plan_fanout() assert results == [] def test_returns_plan_and_fans_for_active_client(self) -> None: client = create_model( DSPClient, name=DSPClientName.spotify_songwhip, nominal_rps=11 ) create_model( FanConnection, fan_id="u1", dsp_client_id=client.id, status=FanConnectionStatus.active, ) results = plan_fanout() assert len(results) == 1 assert results[0].client_name == DSPClientName.spotify_songwhip assert results[0].plan.fanout_count > 0 assert len(results[0].fans) == 1 assert results[0].fans[0].fan_id == "u1" def test_orders_by_connect_recency_newest_first(self) -> None: client = create_model( DSPClient, name=DSPClientName.spotify_songwhip, nominal_rps=11 ) # connect_rank = -epoch(last connect); more negative = newer connect. # Never-connected (NULL) sorts last. create_model( FanConnection, fan_id="newest", dsp_client_id=client.id, status=FanConnectionStatus.active, connect_rank=-2000, ) create_model( FanConnection, fan_id="older", dsp_client_id=client.id, status=FanConnectionStatus.active, connect_rank=-1000, ) create_model( FanConnection, fan_id="never_connected", dsp_client_id=client.id, status=FanConnectionStatus.active, connect_rank=None, ) fans = [f.fan_id for f in plan_fanout()[0].fans] assert fans == ["newest", "older", "never_connected"] @freeze_time("2026-06-25T12:00:00Z") def test_reconnect_requeues_recently_collected_fan(self) -> None: client = create_model( DSPClient, name=DSPClientName.spotify_songwhip, nominal_rps=11 ) # Collected an hour ago (within interval → not due) and its last dispatch # completed — normally skipped this tick. connect() re-queues it. create_model( FanConnection, fan_id="u1", dsp_id=DSPId.spotify, dsp_client_id=client.id, status=FanConnectionStatus.active, last_collected_at=datetime(2026, 6, 25, 11, 0, tzinfo=UTC), last_dispatched_at=datetime(2026, 6, 25, 10, 0, tzinfo=UTC), ) FanConnection.query.connect( [ { "fan_id": "u1", "dsp_id": DSPId.spotify, "dsp_client_id": client.id, "token_encrypted": "new", } ] ) results = plan_fanout() assert [f.fan_id for f in results[0].fans] == ["u1"] def test_returns_empty_fans_when_count_zero(self) -> None: create_model(DSPClient, name=DSPClientName.spotify_songwhip, nominal_rps=0) results = plan_fanout() assert len(results) == 1 assert results[0].plan.fanout_count == 0 assert results[0].fans == [] def test_sizes_from_collected_when_fans_drained_last_window(self) -> None: client = create_model( DSPClient, name=DSPClientName.spotify_songwhip, nominal_rps=11 ) create_model( FanCollectionState, fan_id="u1", dsp_id=DSPId.spotify, last_dsp_client_id=client.id, last_collected_at=timezone.now(), last_collection_error=None, ) fanout = create_model(FanoutTask, source=FanoutSource.scheduled) create_model( CollectTask, fanout_task_id=fanout.id, dsp_client_name=DSPClientName.spotify_songwhip, started_at=datetime(2026, 1, 1, 0, 0, tzinfo=UTC), finished_at=datetime(2026, 1, 1, 0, 1, tzinfo=UTC), # 60s status=TaskStatus.done, fans_total=1, requests=9, ) results = plan_fanout() assert results[0].plan.fanout_count_source == "collected" assert results[0].plan.rps == 9 / 60 # 9 requests / 60s span assert results[0].plan.reqs_per_fan == 9 def test_falls_back_to_seed_when_no_signals(self) -> None: create_model(DSPClient, name=DSPClientName.spotify_songwhip, nominal_rps=11) results = plan_fanout() assert results[0].plan.fanout_count_source == "seed" assert results[0].plan.rps is None @freeze_time("2026-06-02T12:00:00Z") def test_overrun_run_dispatches_fewer_fans_and_messages(self) -> None: client = create_model( DSPClient, name=DSPClientName.spotify_songwhip, nominal_rps=11 ) # 80 fans eligible for dispatch. create_model_batch( FanConnection, size=80, dsp_client_id=client.id, status=FanConnectionStatus.active, ) # Previous run finished but took 20m — overran the 10m window. fanout = create_model( FanoutTask, source=FanoutSource.scheduled, status=TaskStatus.done, started_at=datetime(2026, 6, 2, 11, 40, tzinfo=UTC), ) create_model( CollectTask, fanout_task_id=fanout.id, dsp_client_name=DSPClientName.spotify_songwhip, started_at=datetime(2026, 6, 2, 11, 40, tzinfo=UTC), finished_at=datetime(2026, 6, 2, 12, 0, tzinfo=UTC), # 1200s drain status=TaskStatus.done, fans_total=60, requests=540, ) with override_settings( fan_fanout_window_s=600, fan_collect_worker_timeout_s=540, fan_collect_batch_margin_s=120, ): results = plan_fanout() plan = results[0].plan fans = results[0].fans messages = math.ceil(len(fans) / results[0].batch_size) assert plan.rps == 0.45 # 540 requests / 1200s span assert plan.reqs_per_fan == 9 # 540 / 60 fans assert plan.last_run_duration_s == 1200 assert plan.overran is True assert plan.fanout_count == 27 # floor(0.45 * 0.9 * 600 / 9) assert len(fans) == 27 assert results[0].batch_size == 21 # budget 420 / (20s/fan) = 21 assert messages == 2 # ceil(27 / 21) @freeze_time("2026-06-02T12:00:00Z") def test_within_window_run_dispatches_more_fans_and_messages(self) -> None: client = create_model( DSPClient, name=DSPClientName.spotify_songwhip, nominal_rps=11 ) create_model_batch( FanConnection, size=250, dsp_client_id=client.id, status=FanConnectionStatus.active, ) # Previous run drained in 180s of the 600s window — fill the idle. fanout = create_model( FanoutTask, source=FanoutSource.scheduled, status=TaskStatus.done, started_at=datetime(2026, 6, 2, 11, 50, tzinfo=UTC), finished_at=datetime(2026, 6, 2, 11, 53, tzinfo=UTC), ) create_model( CollectTask, fanout_task_id=fanout.id, dsp_client_name=DSPClientName.spotify_songwhip, started_at=datetime(2026, 6, 2, 11, 50, tzinfo=UTC), finished_at=datetime(2026, 6, 2, 11, 53, tzinfo=UTC), # 180s status=TaskStatus.done, fans_total=60, requests=540, ) with override_settings( fan_fanout_window_s=600, fan_collect_worker_timeout_s=540, fan_collect_batch_margin_s=120, ): results = plan_fanout() plan = results[0].plan fans = results[0].fans messages = math.ceil(len(fans) / results[0].batch_size) assert plan.overran is False assert plan.fanout_count == 180 # floor(3 * 0.9 * 600 / 9), fills the idle assert len(fans) == 180 assert messages == 2 # batch=140 (budget 420 / 3s/fan) → ceil(180 / 140)