"""EventBridge connector.""" from __future__ import annotations from collections.abc import Callable from functools import wraps from typing import Any, Optional, ParamSpec, TypeVar import boto3 from botocore.client import BaseClient, Config from botocore.exceptions import ClientError from lambdacommon.common_config import logger from src.errors import TransientError from src.schemas.event_bridge import ( EventBridgeEvent, EventBridgeEventDetail, EventBridgeEventMetadata, ) P = ParamSpec('P') T = TypeVar('T') def handle_eb_errors(func: Callable[P, T]) -> Callable[P, T]: """Decorate to handle transient EventBridge errors.""" @wraps(func) def wrapper(*args: P.args, **kwargs: P.kwargs) -> T: try: return func(*args, **kwargs) except ClientError as e: logger.error(f'EventBridge client error: {e}') raise TransientError(f'EventBridge error: {e}') from e return wrapper class EventBridgeConnector: """EventBridge connector.""" def __init__(self, source: str, events_client: BaseClient) -> None: """Initialize EventBridge connector. Args: source: Event source (e.g., 'abacus.outbox') events_client: Boto3 Events client instance """ self.source = source self.events_client = events_client @handle_eb_errors def put_event( self, detail_type: str, metadata: EventBridgeEventMetadata, data: Any = None, ) -> None: """Put events to EventBridge. Args: detail_type: The type of event (e.g., 'file_upload.completed'). metadata: The event metadata. data: The event business data. Raises: TransientError: If EventBridge call fails (retriable). """ event = EventBridgeEvent( source=self.source, detail_type=detail_type, detail=EventBridgeEventDetail(metadata=metadata, data=data), ) entry = event.model_dump(by_alias=True, exclude_none=True) response = self.events_client.put_events(Entries=[entry]) failed_count = response.get('FailedEntryCount', 0) if failed_count > 0: error_code = response['Entries'][0].get('ErrorCode') error_msg = response['Entries'][0].get('ErrorMessage') raise TransientError( f'Failed to publish event with error {error_code}: {error_msg}' ) def get_events_client(config: Optional[Config] = None) -> BaseClient: """Get EventBridge client. Args: config: Optional botocore Config for the client Returns: Configured Boto3 Events client instance """ if config is None: config = Config() return boto3.client('events', config=config)