import json import logging from typing import Any from resonance_engine.adapters.kafka import kafka_producer from resonance_engine.config import settings from resonance_engine.dsp.enums import DSPId from resonance_engine.fandata.types import FanRecord from .base import DataSink logger = logging.getLogger(__name__) class KafkaSink(DataSink): def write_fans( self, *, fan: FanRecord, dsp_id: DSPId, data: dict[str, Any], ) -> None: self._publish( topic=settings.kafka_topic_fan_profile, key=f"{dsp_id}:{fan.fan_id}", payload={ "fan_id": fan.fan_id, "dsp_id": dsp_id, "data": data, }, ) def write_fan_top_artists( self, *, fan: FanRecord, dsp_id: DSPId, items: list[dict[str, Any]], ) -> None: for index, data in enumerate(items, start=1): self._publish( topic=settings.kafka_topic_fan_top_artists, key=f"{dsp_id}:{fan.fan_id}", payload={ "dsp_id": dsp_id, "fan_id": fan.fan_id, "data": data, "position": index, }, ) def write_fan_top_tracks( self, *, fan: FanRecord, dsp_id: DSPId, items: list[dict[str, Any]], ) -> None: for index, data in enumerate(items, start=1): self._publish( topic=settings.kafka_topic_fan_top_tracks, key=f"{dsp_id}:{fan.fan_id}", payload={ "dsp_id": dsp_id, "fan_id": fan.fan_id, "data": data, "position": index, }, ) def write_fan_recently_played( self, *, fan: FanRecord, dsp_id: DSPId, items: list[dict[str, Any]], ) -> None: for data in items: self._publish( topic=settings.kafka_topic_fan_recently_played, key=f"{dsp_id}:{fan.fan_id}", payload={ "dsp_id": dsp_id, "fan_id": fan.fan_id, "data": data, }, ) def write_fan_playlists( self, *, fan: FanRecord, dsp_id: DSPId, items: list[dict[str, Any]], ) -> None: for data in items: self._publish( topic=settings.kafka_topic_fan_playlists, key=f"{dsp_id}:{fan.fan_id}", payload={ "dsp_id": dsp_id, "fan_id": fan.fan_id, "data": data, }, ) def write_fan_saved_albums( self, *, fan: FanRecord, dsp_id: DSPId, items: list[dict[str, Any]], ) -> None: for data in items: self._publish( topic=settings.kafka_topic_fan_saved_albums, key=f"{dsp_id}:{fan.fan_id}", payload={ "dsp_id": dsp_id, "fan_id": fan.fan_id, "data": data, }, ) def write_fan_saved_tracks( self, *, fan: FanRecord, dsp_id: DSPId, items: list[dict[str, Any]], ) -> None: for data in items: self._publish( topic=settings.kafka_topic_fan_saved_tracks, key=f"{dsp_id}:{fan.fan_id}", payload={ "dsp_id": dsp_id, "fan_id": fan.fan_id, "data": data, }, ) def write_fan_followed_artists( self, *, fan: FanRecord, dsp_id: DSPId, items: list[dict[str, Any]], ) -> None: for data in items: self._publish( topic=settings.kafka_topic_fan_followed_artists, key=f"{dsp_id}:{fan.fan_id}", payload={ "dsp_id": dsp_id, "fan_id": fan.fan_id, "data": data, }, ) def flush(self) -> None: undelivered = kafka_producer.flush(timeout=30.0) if undelivered > 0: raise RuntimeError( f"Kafka flush: {undelivered} messages undelivered — failing batch for retry" ) @staticmethod def _publish(*, topic: str, key: str, payload: dict[str, Any]) -> None: kafka_producer.publish( topic=topic, key=key, value=json.dumps(payload, default=str), )