import logging from typing import cast from fansifter_common.adapters.sendgrid.models import FanData from app.adapters.redis import redis_client from app.config import settings from app.utils import hash_email logger = logging.getLogger(__name__) def make_key(automated_email_id: str, email_hash: bytes) -> str: return f"dedup:automated_email:{automated_email_id}:{email_hash.hex()}" def filter_unsent( automated_email_id: str, fan_data_list: list[FanData] ) -> list[FanData]: if not fan_data_list: return [] hashes = [hash_email(fan_data.email) for fan_data in fan_data_list] keys = [make_key(automated_email_id, h) for h in hashes] results = cast(list[bytes | None], redis_client.mget(keys)) unsent: list[FanData] = [ fan_data for fan_data, value in zip(fan_data_list, results, strict=True) if value is None ] if skipped := len(fan_data_list) - len(unsent): logger.info( "Skipping %d already sent fans.", skipped, extra={"automated_email_id": automated_email_id}, ) return unsent def mark_sent(automated_email_id: str, fan_data_list: list[FanData]) -> None: if not fan_data_list: return hashes = [hash_email(fan_data.email) for fan_data in fan_data_list] with redis_client.pipeline() as pipe: for h in hashes: pipe.set(make_key(automated_email_id, h), 1, ex=settings.redis_dedup_ttl) pipe.execute()