import logging from datetime import UTC, datetime, timedelta from unittest import mock import pytest from fansifter_common.utils import timezone from pydantic import SecretStr from pytest_mock import MockerFixture from resonance_engine.dsp.enums import DSPClientName, DSPId, DSPResource from resonance_engine.dsp.exceptions import ( DSPForbiddenError, RateLimitError, StreamingAPIError, TokenRefreshError, TokenRevokedError, ) from resonance_engine.dsp.models import DSPClient from resonance_engine.fandata.collector import ( _resource_last_at, _resource_needed, collect_fan, collect_fans, ) from resonance_engine.fandata.enums import ( FanCollectionError, FanConnectionStatus, ) from resonance_engine.fandata.models import FanCollectionState, FanConnection from resonance_engine.fandata.timings import CollectTimings from resonance_engine.fandata.types import FanRecord from tests.unit.helpers import create_model, override_settings @pytest.mark.db class TestFanCollector: @pytest.fixture def gateway_mock(self, mocker: MockerFixture) -> mock.MagicMock: m = mocker.patch( "resonance_engine.fandata.collector.dsp_gateway", new_callable=mock.MagicMock, ) m.refresh_token.return_value = {"access_token": SecretStr("acc-tok")} m.get_profile.return_value = {} m.get_top_artists.return_value = [{"id": "art-1"}] return m @pytest.fixture def data_sink_mock(self, mocker: MockerFixture) -> mock.MagicMock: return mocker.patch( "resonance_engine.fandata.collector.data_sink", new_callable=mock.MagicMock, ) # --- helpers -------------------------------------------------------------- def test_resource_needed_true_when_never_collected(self) -> None: assert _resource_needed(None, datetime(2026, 5, 1, tzinfo=UTC), 0) def test_resource_needed_false_when_already_done_this_run(self) -> None: started = datetime(2026, 5, 1, tzinfo=UTC) assert not _resource_needed(started, started, 0) def test_resource_needed_true_with_zero_interval(self) -> None: assert _resource_needed( datetime(2026, 4, 1, tzinfo=UTC), datetime(2026, 5, 1, tzinfo=UTC), 0, ) def test_resource_needed_false_when_within_interval(self) -> None: now = timezone.now() assert not _resource_needed( now - timedelta(seconds=30), now - timedelta(seconds=1), 3600 ) def test_resource_needed_true_when_interval_elapsed(self) -> None: now = timezone.now() assert _resource_needed( now - timedelta(hours=2), now - timedelta(seconds=1), 3600 ) def test_resource_last_at_returns_none_for_missing_state(self) -> None: assert _resource_last_at(None, DSPResource.profile) is None def test_resource_last_at_returns_collected_at_for_resource(self) -> None: state = mock.MagicMock( profile_collected_at=datetime(2026, 5, 1, tzinfo=UTC), top_artists_collected_at=None, ) assert _resource_last_at(state, DSPResource.profile) == datetime( 2026, 5, 1, tzinfo=UTC ) assert _resource_last_at(state, DSPResource.top_artists) is None # --- collect_fans (batch budget) ----------------------------------------- def test_collect_fans_processes_whole_batch_within_budget( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) for uid in ("user-1", "user-2"): create_model( FanConnection, fan_id=uid, dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) outcome = collect_fans( fans=[ FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), FanRecord(fan_id="user-2", token_encrypted="refresh-tok"), ], client=dsp_client, ) assert outcome.fans_processed == 2 def test_collect_fans_defers_fans_when_budget_exhausted( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) with override_settings( fan_collect_worker_timeout_s=540, fan_collect_batch_margin_s=600 ): outcome = collect_fans( fans=[FanRecord(fan_id="user-1", token_encrypted="refresh-tok")], client=dsp_client, ) assert outcome.fans_processed == 0 gateway_mock.refresh_token.assert_not_called() def test_collect_fans_isolates_failing_fan( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) for fan_id in ("user-1", "user-2"): create_model( FanConnection, fan_id=fan_id, dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) # user-1 raises an unexpected (non-DSP) error; the batch must continue. gateway_mock.refresh_token.side_effect = [ RuntimeError("boom"), {"access_token": SecretStr("acc-tok")}, ] outcome = collect_fans( fans=[ FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), FanRecord(fan_id="user-2", token_encrypted="refresh-tok"), ], client=dsp_client, ) assert outcome.fans_processed == 1 # user-2 still collected assert outcome.fans_errors == 1 # user-1 counted, not fatal # --- collect_fan ---------------------------------------------------------- def test_success( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) outcome = collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, ) assert outcome.processed == 1 assert outcome.errors == 0 assert outcome.requests == 9 # token + 8 resources assert outcome.requests_skipped == 0 gateway_mock.refresh_token.assert_called_once() gateway_mock.get_profile.assert_called_once() gateway_mock.get_top_artists.assert_called_once() gateway_mock.get_top_tracks.assert_called_once() gateway_mock.get_recently_played.assert_called_once() gateway_mock.get_playlists.assert_called_once() gateway_mock.get_saved_albums.assert_called_once() gateway_mock.get_saved_tracks.assert_called_once() gateway_mock.get_followed_artists.assert_called_once() data_sink_mock.flush.assert_called_once() state = FanCollectionState.query.where( FanCollectionState.fan_id == "user-1", FanCollectionState.dsp_id == DSPId.spotify, ).one() assert state.last_collection_error is None assert state.consecutive_failures == 0 assert state.profile_collected_at is not None assert state.top_artists_collected_at is not None def test_proactive_scope_skip_counts_as_skipped_not_request( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, caplog: pytest.LogCaptureFixture, ) -> None: gateway_mock.refresh_token.return_value = { "access_token": SecretStr("acc-tok"), "scope": "user-read-email user-read-private", } # Gateway pre-skipped the call (known-missing scope) — no request made. gateway_mock.get_top_artists.side_effect = DSPForbiddenError( "token missing scope 'user-top-read' for top_artists", request_made=False ) dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) with caplog.at_level(logging.DEBUG): outcome = collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, ) assert outcome.processed == 1 # fan still succeeds assert outcome.errors == 0 assert outcome.requests_skipped == 1 # the pre-skipped resource — no request assert outcome.requests == len(DSPResource) # token + 7 made; top_artists not state = FanCollectionState.query.where( FanCollectionState.fan_id == "user-1", FanCollectionState.dsp_id == DSPId.spotify, ).one() assert state.top_artists_collected_at is not None assert state.profile_collected_at is not None # other resources still collected assert "missing scope" in caplog.text assert "user-read-email user-read-private" in caplog.text def test_real_403_counts_as_request_not_skipped( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: gateway_mock.refresh_token.return_value = { "access_token": SecretStr("acc-tok"), "scope": "user-read-email user-read-private", } # The API actually returned 403 — a request was made (request_made defaults True). gateway_mock.get_top_artists.side_effect = DSPForbiddenError( "403: Insufficient client scope" ) dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) outcome = collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, ) assert outcome.processed == 1 # fan still succeeds assert outcome.requests_skipped == 0 # the 403 was a real request, not a skip assert outcome.requests == 1 + len(DSPResource) # token + every resource called state = FanCollectionState.query.where( FanCollectionState.fan_id == "user-1", FanCollectionState.dsp_id == DSPId.spotify, ).one() # Stamped on a real 403 too — don't re-attempt (and re-spend a request) # every tick; retry once the interval elapses. assert state.top_artists_collected_at is not None def test_skips_token_refresh_when_no_resource_due( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) # All resource intervals default to 24h or 0; mark everything as collected # 1 minute ago so nothing is due yet (except recently_played which has # interval_s=0 and is always due — set FAN_COLLECT_RECENTLY_PLAYED to # match below). recent = timezone.now() - timedelta(minutes=1) create_model( FanCollectionState, fan_id="user-1", dsp_id=DSPId.spotify, last_dsp_client_id=dsp_client.id, last_collected_at=recent, profile_collected_at=recent, top_artists_collected_at=recent, top_tracks_collected_at=recent, recently_played_collected_at=recent, playlists_collected_at=recent, saved_albums_collected_at=recent, saved_tracks_collected_at=recent, followed_artists_collected_at=recent, ) with mock.patch( "resonance_engine.fandata.collector.settings.fan_collect_recently_played_interval_s", 3600, ): outcome = collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, ) assert outcome.processed == 1 assert outcome.errors == 0 assert outcome.requests == 0 assert outcome.requests_skipped == len(DSPResource) gateway_mock.refresh_token.assert_not_called() def test_processes_when_at_least_one_resource_due( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) # Only profile is stale; the rest were just collected. recent = timezone.now() - timedelta(minutes=1) stale = timezone.now() - timedelta(days=8) create_model( FanCollectionState, fan_id="user-1", dsp_id=DSPId.spotify, last_dsp_client_id=dsp_client.id, last_collected_at=recent, profile_collected_at=stale, top_artists_collected_at=recent, top_tracks_collected_at=recent, recently_played_collected_at=recent, playlists_collected_at=recent, saved_albums_collected_at=recent, saved_tracks_collected_at=recent, followed_artists_collected_at=recent, ) with mock.patch( "resonance_engine.fandata.collector.settings.fan_collect_recently_played_interval_s", 3600, ): outcome = collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, ) assert outcome.processed == 1 assert outcome.requests > 0 gateway_mock.refresh_token.assert_called_once() def test_no_token_skips_without_calling_gateway( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) outcome = collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted=""), client=dsp_client, ) assert outcome.errors == 1 assert outcome.requests == 0 gateway_mock.refresh_token.assert_not_called() data_sink_mock.write_fans.assert_not_called() def test_token_revoked_marks_connection_revoked( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) gateway_mock.refresh_token.side_effect = TokenRevokedError("revoked") outcome = collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, ) assert outcome.stale_tokens == 1 assert outcome.processed == 0 connection = FanConnection.query.where( FanConnection.fan_id == "user-1", FanConnection.dsp_id == DSPId.spotify, ).one() assert connection.status == FanConnectionStatus.revoked state = FanCollectionState.query.where( FanCollectionState.fan_id == "user-1", FanCollectionState.dsp_id == DSPId.spotify, ).one() assert state.last_collection_error == FanCollectionError.token_error def test_token_refresh_failed_keeps_connection_active( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) gateway_mock.refresh_token.side_effect = TokenRefreshError("transient") outcome = collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, ) assert outcome.errors == 1 assert outcome.processed == 0 connection = FanConnection.query.where( FanConnection.fan_id == "user-1", FanConnection.dsp_id == DSPId.spotify, ).one() assert connection.status == FanConnectionStatus.active state = FanCollectionState.query.where( FanCollectionState.fan_id == "user-1", FanCollectionState.dsp_id == DSPId.spotify, ).one() assert state.last_collection_error == FanCollectionError.token_error def test_rate_limit_during_token_refresh( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) gateway_mock.refresh_token.side_effect = RateLimitError("rl", retry_after=1) outcome = collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, ) assert outcome.rate_limited == 1 assert outcome.reset_consecutive is False state = FanCollectionState.query.where( FanCollectionState.fan_id == "user-1", FanCollectionState.dsp_id == DSPId.spotify, ).one() assert state.last_collection_error == FanCollectionError.rate_limited def test_api_error_in_resource_records_api_error( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) gateway_mock.get_profile.side_effect = StreamingAPIError("500", status_code=500) outcome = collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, ) assert outcome.errors == 1 assert outcome.processed == 0 state = FanCollectionState.query.where( FanCollectionState.fan_id == "user-1", FanCollectionState.dsp_id == DSPId.spotify, ).one() assert state.last_collection_error == FanCollectionError.api_error assert state.profile_collected_at is None def test_rate_limit_mid_resources_preserves_completed_resources( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) gateway_mock.get_top_artists.side_effect = RateLimitError("rl") outcome = collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, ) assert outcome.rate_limited == 1 state = FanCollectionState.query.where( FanCollectionState.fan_id == "user-1", FanCollectionState.dsp_id == DSPId.spotify, ).one() assert state.last_collection_error == FanCollectionError.rate_limited # profile was completed before the rate limit hit on top_artists assert state.profile_collected_at is not None assert state.top_artists_collected_at is None def test_rotates_refresh_token_when_returned( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) gateway_mock.refresh_token.return_value = { "access_token": SecretStr("acc"), "refresh_token": SecretStr("new-refresh-tok"), } collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, ) connection = FanConnection.query.where( FanConnection.fan_id == "user-1", FanConnection.dsp_id == DSPId.spotify, ).one() assert connection.token_encrypted == "new-refresh-tok" def test_skips_resource_already_done_this_run( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) future = datetime(2030, 1, 1, tzinfo=UTC) create_model( FanCollectionState, fan_id="user-1", dsp_id=DSPId.spotify, last_dsp_client_id=dsp_client.id, last_collected_at=future, profile_collected_at=future, ) outcome = collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, started_at=datetime(2020, 1, 1, tzinfo=UTC), ) assert outcome.requests_skipped == 1 gateway_mock.get_profile.assert_not_called() gateway_mock.get_top_artists.assert_called_once() def test_populates_timings( self, gateway_mock: mock.MagicMock, data_sink_mock: mock.MagicMock, ) -> None: dsp_client = create_model( DSPClient, dsp_id=DSPId.spotify, name=DSPClientName.spotify_songwhip ) create_model( FanConnection, fan_id="user-1", dsp_id=DSPId.spotify, dsp_client_id=dsp_client.id, token_encrypted="refresh-tok", ) timings = CollectTimings() collect_fan( fan=FanRecord(fan_id="user-1", token_encrypted="refresh-tok"), client=dsp_client, timings=timings, ) assert len(timings.token_ms) == 1 assert len(timings.profile_ms) == 1 assert len(timings.top_artists_ms) == 1 assert len(timings.top_tracks_ms) == 1 assert len(timings.recently_played_ms) == 1 assert len(timings.playlists_ms) == 1 assert len(timings.saved_albums_ms) == 1 assert len(timings.saved_tracks_ms) == 1 assert len(timings.followed_artists_ms) == 1