import json import logging import uuid from time import time from typing import Any from anydi import singleton from confluent_kafka import KafkaError, Message, Producer from app.exceptions import RequestValidationError from app.models import Event from app.validator import RequestValidator logger = logging.getLogger(__name__) @singleton class TwilioWebhooksOutboundHandler: def __init__( self, kafka_producer: Producer, kafka_event_topic: str, request_validator: RequestValidator, ) -> None: self.kafka_producer = kafka_producer self.kafka_event_topic = kafka_event_topic self.request_validator = request_validator def handle(self, event_dict: dict[str, Any]) -> dict[str, Any]: logger.info("Event received", extra={"event": event_dict}) try: timestamp = round(time()) event = Event.model_validate(event_dict) self.request_validator.validate( event.request_full_uri, event.body, event.headers.request_signature, ) self.kafka_producer.produce( self.kafka_event_topic, key=str(uuid.uuid4()), value=json.dumps( {"body": event.body, "params": event.params, "timestamp": timestamp} ), callback=self.kafka_producer_callback, ) self.kafka_producer.flush() return {"statusCode": 200, "body": "OK"} except RequestValidationError: logger.error("Request validation error", extra={"event": event_dict}) return {"statusCode": 403, "body": "Forbidden"} except Exception: logger.exception("Internal server error", extra={"event": event_dict}) return {"statusCode": 500, "body": "Internal Server Error"} @staticmethod def kafka_producer_callback(error: KafkaError | None, message: Message) -> None: if error: logger.error("Kafka producer error %s", error) else: logger.info("Message %s delivered", message.key())