"""Kafka sink — produces to separate MSK topics per data type. Each method produces one message per row, fields matching the Postgres table with JSONB `data` columns spread inline. """ import hashlib import json import logging from datetime import datetime from typing import Any from confluent_kafka import Producer from app.dsp.enums import DSPId from app.dsp.types import ( FollowedArtistsResult, PlaylistsResult, ProfileResult, RecentlyPlayedResult, SavedAlbumsResult, SavedTracksResult, TopArtistsResult, TopTracksResult, ) from app.fandata.sinks.base import DataSink from app.fandata.types import FanRecord logger = logging.getLogger(__name__) class KafkaSink(DataSink): def __init__( self, producer: Producer, topic_fans: str, topic_top_artists: str, topic_top_tracks: str, topic_recently_played: str, topic_playlists: str, topic_saved_albums: str, topic_saved_tracks: str, topic_followed_artists: str, ) -> None: self._producer = producer self._topic_fans = topic_fans self._topic_top_artists = topic_top_artists self._topic_top_tracks = topic_top_tracks self._topic_recently_played = topic_recently_played self._topic_playlists = topic_playlists self._topic_saved_albums = topic_saved_albums self._topic_saved_tracks = topic_saved_tracks self._topic_followed_artists = topic_followed_artists def write_fans( self, *, fan: FanRecord, dsp_id: DSPId, profile: ProfileResult, collected_at: datetime, ) -> None: email = profile.email or "" fan_id = hashlib.sha256(email.encode()).hexdigest() self._produce( self._topic_fans, fan.dsp_user_id, { "dsp_user_id": fan.dsp_user_id, "dsp_id": dsp_id, "email": email, "fan_id": fan_id, "country": profile.country, "product": profile.product, "display_name": profile.display_name, "collected_at": collected_at.isoformat(), }, ) def write_fan_top_artists( self, *, fan: FanRecord, dsp_id: DSPId, result: TopArtistsResult, collected_at: datetime, ) -> None: for item in result.items: self._produce( self._topic_top_artists, fan.dsp_user_id, { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "artist_id": item.id, "collected_at": collected_at.isoformat(), **item.data, }, ) def write_fan_top_tracks( self, *, fan: FanRecord, dsp_id: DSPId, result: TopTracksResult, collected_at: datetime, ) -> None: for item in result.items: self._produce( self._topic_top_tracks, fan.dsp_user_id, { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "track_id": item.id, "collected_at": collected_at.isoformat(), **item.data, }, ) def write_fan_recently_played( self, *, fan: FanRecord, dsp_id: DSPId, result: RecentlyPlayedResult, collected_at: datetime, ) -> None: for item in result.items: self._produce( self._topic_recently_played, fan.dsp_user_id, { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "track_id": item.track_id, "played_at": item.played_at, "collected_at": collected_at.isoformat(), **item.data, }, ) def write_fan_playlists( self, *, fan: FanRecord, dsp_id: DSPId, result: PlaylistsResult, collected_at: datetime, ) -> None: for item in result.items: self._produce( self._topic_playlists, fan.dsp_user_id, { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "playlist_id": item.id, "collected_at": collected_at.isoformat(), **item.data, }, ) def write_fan_saved_albums( self, *, fan: FanRecord, dsp_id: DSPId, result: SavedAlbumsResult, collected_at: datetime, ) -> None: for item in result.items: self._produce( self._topic_saved_albums, fan.dsp_user_id, { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "album_id": item.album_id, "added_at": item.added_at, "collected_at": collected_at.isoformat(), **item.data, }, ) def write_fan_saved_tracks( self, *, fan: FanRecord, dsp_id: DSPId, result: SavedTracksResult, collected_at: datetime, ) -> None: for item in result.items: self._produce( self._topic_saved_tracks, fan.dsp_user_id, { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "track_id": item.track_id, "added_at": item.added_at, "collected_at": collected_at.isoformat(), **item.data, }, ) def write_fan_followed_artists( self, *, fan: FanRecord, dsp_id: DSPId, result: FollowedArtistsResult, collected_at: datetime, ) -> None: for item in result.items: self._produce( self._topic_followed_artists, fan.dsp_user_id, { "dsp_id": dsp_id, "dsp_user_id": fan.dsp_user_id, "artist_id": item.id, "collected_at": collected_at.isoformat(), **item.data, }, ) def flush(self) -> None: undelivered = self._producer.flush(timeout=30) if undelivered > 0: raise RuntimeError( f"Kafka flush: {undelivered} messages undelivered — failing batch for retry" ) def _produce(self, topic: str, key: str, message: dict[str, Any]) -> None: self._producer.produce( topic, value=json.dumps(message).encode("utf-8"), key=key.encode("utf-8"), on_delivery=_delivery_report, ) self._producer.poll(0) def _delivery_report(err: object, msg: object) -> None: # noqa: ARG001 if err is not None: logger.error("Kafka delivery failed: %s", err)