import json from collections.abc import Iterator from unittest import mock import pytest from resonance_engine.config import settings from resonance_engine.dsp.enums import DSPId from resonance_engine.fandata.sink.kafka import KafkaSink from resonance_engine.fandata.types import FanRecord FAN = FanRecord(fan_id="user-1", token_encrypted="tok") KEY = f"{DSPId.spotify}:user-1" @pytest.fixture def producer_mock() -> Iterator[mock.MagicMock]: with mock.patch( "resonance_engine.fandata.sink.kafka.kafka_producer", new_callable=mock.MagicMock, ) as m: yield m def published(producer_mock: mock.MagicMock) -> list[tuple[str, str, dict]]: return [ (call.kwargs["topic"], call.kwargs["key"], json.loads(call.kwargs["value"])) for call in producer_mock.publish.call_args_list ] class TestKafkaSink: def test_write_fans(self, producer_mock: mock.MagicMock) -> None: KafkaSink().write_fans( fan=FAN, dsp_id=DSPId.spotify, data={ "email": "alice@example.com", "display_name": "Alice", "country": "US", "product": "premium", }, ) assert published(producer_mock) == [ ( settings.kafka_topic_fan_profile, KEY, { "fan_id": "user-1", "dsp_id": "spotify", "data": { "email": "alice@example.com", "country": "US", "product": "premium", "display_name": "Alice", }, }, ), ] def test_write_fan_top_artists(self, producer_mock: mock.MagicMock) -> None: KafkaSink().write_fan_top_artists( fan=FAN, dsp_id=DSPId.spotify, items=[ {"id": "a1", "name": "Artist One"}, {"id": "a2", "name": "Artist Two"}, ], ) assert published(producer_mock) == [ ( settings.kafka_topic_fan_top_artists, KEY, { "dsp_id": "spotify", "fan_id": "user-1", "data": {"id": "a1", "name": "Artist One"}, "position": 1, }, ), ( settings.kafka_topic_fan_top_artists, KEY, { "dsp_id": "spotify", "fan_id": "user-1", "data": {"id": "a2", "name": "Artist Two"}, "position": 2, }, ), ] def test_write_fan_top_tracks(self, producer_mock: mock.MagicMock) -> None: KafkaSink().write_fan_top_tracks( fan=FAN, dsp_id=DSPId.spotify, items=[{"id": "t1", "name": "Song"}], ) assert published(producer_mock) == [ ( settings.kafka_topic_fan_top_tracks, KEY, { "dsp_id": "spotify", "fan_id": "user-1", "data": {"id": "t1", "name": "Song"}, "position": 1, }, ), ] def test_write_fan_recently_played(self, producer_mock: mock.MagicMock) -> None: KafkaSink().write_fan_recently_played( fan=FAN, dsp_id=DSPId.spotify, items=[ { "track_id": "t1", "played_at": "2026-04-30T10:00:00Z", "context": "playlist", } ], ) assert published(producer_mock) == [ ( settings.kafka_topic_fan_recently_played, KEY, { "dsp_id": "spotify", "fan_id": "user-1", "data": { "track_id": "t1", "played_at": "2026-04-30T10:00:00Z", "context": "playlist", }, }, ), ] def test_write_fan_playlists(self, producer_mock: mock.MagicMock) -> None: KafkaSink().write_fan_playlists( fan=FAN, dsp_id=DSPId.spotify, items=[{"playlist_id": "p1", "name": "Favs"}] ) assert published(producer_mock) == [ ( settings.kafka_topic_fan_playlists, KEY, { "dsp_id": "spotify", "fan_id": "user-1", "data": { "playlist_id": "p1", "name": "Favs", }, }, ), ] def test_write_fan_saved_albums(self, producer_mock: mock.MagicMock) -> None: KafkaSink().write_fan_saved_albums( fan=FAN, dsp_id=DSPId.spotify, items=[ { "album_id": "al1", "added_at": "2026-04-29T08:00:00Z", "name": "Album", } ], ) assert published(producer_mock) == [ ( settings.kafka_topic_fan_saved_albums, KEY, { "dsp_id": "spotify", "fan_id": "user-1", "data": { "album_id": "al1", "name": "Album", "added_at": "2026-04-29T08:00:00Z", }, }, ), ] def test_write_fan_saved_tracks(self, producer_mock: mock.MagicMock) -> None: KafkaSink().write_fan_saved_tracks( fan=FAN, dsp_id=DSPId.spotify, items=[ { "track_id": "t1", "added_at": "2026-04-28T09:00:00Z", "name": "Song", } ], ) assert published(producer_mock) == [ ( settings.kafka_topic_fan_saved_tracks, KEY, { "dsp_id": "spotify", "fan_id": "user-1", "data": { "track_id": "t1", "name": "Song", "added_at": "2026-04-28T09:00:00Z", }, }, ), ] def test_write_fan_followed_artists(self, producer_mock: mock.MagicMock) -> None: KafkaSink().write_fan_followed_artists( fan=FAN, dsp_id=DSPId.spotify, items=[{"artist_id": "a1", "name": "Artist"}], ) assert published(producer_mock) == [ ( settings.kafka_topic_fan_followed_artists, KEY, { "dsp_id": "spotify", "fan_id": "user-1", "data": {"artist_id": "a1", "name": "Artist"}, }, ), ] def test_flush_succeeds_when_all_delivered( self, producer_mock: mock.MagicMock ) -> None: producer_mock.flush.return_value = 0 KafkaSink().flush() producer_mock.flush.assert_called_once_with(timeout=30.0) def test_flush_raises_when_undelivered(self, producer_mock: mock.MagicMock) -> None: producer_mock.flush.return_value = 3 with pytest.raises(RuntimeError, match="3 messages undelivered"): KafkaSink().flush()