"""MySQL Outbox Iterable.""" from __future__ import annotations from typing import Iterator from lambdacommon.common_config import logger from src.connectors.repository import Repository from src.schemas import AbacusOutbox class DBOutboxIterable: """DB outbox iterable. Iterates over pending outbox events from the database using pessimistic row-level locking (FOR UPDATE SKIP LOCKED). The iterable acquires row-level locks when querying events but does NOT release them. The consumer is responsible for releasing the lock by calling commit() or rollback() on the connection (e.g., repository). Lock lifecycle: 1. Iterable: SELECT FOR UPDATE SKIP LOCKED (lock acquired) 2. Iterable: yield event 3. Consumer: process event (lock still held) 4. Consumer: commit() or rollback() (lock released) """ def __init__( self, repository: Repository, capacity: int = 10, ): """Initialize MySQL outbox iterable. Args: repository: Repository instance for database operations capacity: Maximum number of events to iterate over """ self._capacity = capacity self._repository = repository def __iter__(self) -> Iterator[AbacusOutbox]: """Iterate over pending outbox events from the database. Each event is fetched and yielded one by one. Important: - `autocommit` must be DISABLED on the DB connection for SKIP LOCKED to work properly. - The iterator acquires database locks but does NOT release them. The consumer MUST commit or rollback to release locks after processing each event. Yields: AbacusOutbox: Individual outbox events, locked until consumer commits Raises: RuntimeError: If `autocommit` is enabled on the database connection. Note: Fetches events one at a time to maintain proper lock granularity and prevent holding multiple row locks simultaneously. """ if self._repository.is_autocommit_enabled(): raise RuntimeError('Requires autocommit to be disabled') for _ in range(self._capacity): logger.info('Getting pending event') events = self._repository.get_pending_events(limit=1) if not events: break yield events[0]