from collections.abc import Iterator import sqlalchemy as sa from app.adapters.db import Repository from app.models import BatchRecipient, BatchSendRecord, Campaign, CampaignBatch, Domain from app.types import Recipient class DomainRepository(Repository[Domain]): pass class CampaignRepository(Repository[Campaign]): pass class CampaignBatchRepository(Repository[CampaignBatch]): def get_for_update(self, batch_id: int) -> CampaignBatch | None: return self.db.session.get(self.model, batch_id, with_for_update=True) def get_first_in_campaign(self, campaign_id: str) -> CampaignBatch | None: query = ( sa.select(CampaignBatch) .where(CampaignBatch.campaign_id == campaign_id) .order_by(CampaignBatch.id.asc()) .limit(1) ) result = self.db.session.execute(query) return result.scalars().first() class BatchRecipientRepository(Repository[BatchRecipient]): def iter_next_in_batch( self, batch_id: int, *, limit: int, offset: int ) -> Iterator[Recipient]: query = sa.text(""" SELECT fan_id, fan_email, fan_country_iso2, profile_id, fan_first_name FROM email_campaign_batch_recipient WHERE batch_id = :batch_id ORDER BY id ASC LIMIT :limit OFFSET :offset """) result = self.db.session.execute( query, {"batch_id": batch_id, "limit": limit, "offset": offset}, execution_options={"stream_results": True}, ) for row in result: yield Recipient(*row) class BatchSendRecordRepository(Repository[BatchSendRecord]): pass