import logging from datetime import timedelta import requests from celery import shared_task from django.template import Context, Template from django.utils import timezone from infrasynth.shared.settings_utils import get_setting from .models import OutboundDelivery, OutboundSubscription from .signals import outbound_delivery_failed, outbound_delivery_succeeded from .signature import sign_payload logger = logging.getLogger(__name__) @shared_task( name="infrasynth.webhooks.deliver_webhook", bind=True, max_retries=5, default_retry_delay=60, ) def deliver_webhook(self, subscription_id, event_name, payload, payload_template, tenant_id=None): """Delivers an outbound webhook with HMAC signature and retry/backoff.""" from infrasynth.tenancy.context import tenant_context from infrasynth.tenancy.models import Tenant try: subscription = OutboundSubscription.all_objects.select_related("endpoint").get(pk=subscription_id) except OutboundSubscription.DoesNotExist: logger.warning("Subscription %s not found", subscription_id) return None tenant = Tenant.objects.filter(pk=tenant_id).first() if tenant_id else subscription.tenant with tenant_context(tenant): return _deliver_webhook_inner(self, subscription, event_name, payload, payload_template, tenant) def _deliver_webhook_inner(self, subscription, event_name, payload, payload_template, tenant): endpoint = subscription.endpoint body = _build_payload(payload, payload_template) delivery = OutboundDelivery.all_objects.create( tenant=tenant or subscription.tenant, subscription=subscription, payload=body, status=OutboundDelivery.Status.RETRYING, attempt=self.request.retries + 1, ) headers = {"Content-Type": "application/json"} headers.update(endpoint.headers or {}) headers[get_setting("INFRASYNTH_WEBHOOKS", "SIGNATURE_HEADER", "X-Webhook-Signature")] = sign_payload( endpoint.secret, body ) timeout = endpoint.timeout_seconds or get_setting("INFRASYNTH_WEBHOOKS", "DEFAULT_TIMEOUT_SECONDS", 10) try: response = requests.post( endpoint.url, data=body, headers=headers, timeout=timeout, ) except requests.RequestException as exc: return _handle_failure(self, delivery, subscription, endpoint, str(exc)) delivery.response_status = response.status_code delivery.response_body = response.text[:4000] delivery.completed_at = timezone.now() if 200 <= response.status_code < 300: delivery.status = OutboundDelivery.Status.SUCCESS delivery.save() outbound_delivery_succeeded.send( sender=OutboundDelivery, delivery_id=delivery.id, event_name=event_name, status_code=response.status_code, ) logger.info( "Webhook delivered: %s -> %s (%s)", event_name, endpoint.url, response.status_code, ) return delivery.id delivery.save() return _handle_failure( self, delivery, subscription, endpoint, f"HTTP {response.status_code}: {response.text[:500]}", ) def _build_payload(payload: dict, payload_template: str) -> str: import json if payload_template: try: template = Template(payload_template) return template.render(Context({"event": payload, "payload": payload})) except Exception: # noqa: BLE001 logger.exception("Failed to render webhook payload template") return json.dumps(payload, default=str) def _handle_failure(self, delivery, subscription, endpoint, error: str): max_retries = int( (endpoint.retry_policy or {}).get("max_retries", get_setting("INFRASYNTH_WEBHOOKS", "MAX_RETRIES", 5)) ) backoff = (endpoint.retry_policy or {}).get( "backoff", get_setting("INFRASYNTH_WEBHOOKS", "RETRY_BACKOFF", "exponential") ) initial_delay = int(get_setting("INFRASYNTH_WEBHOOKS", "RETRY_INITIAL_DELAY_SECONDS", 60)) delivery.status = OutboundDelivery.Status.RETRYING delivery.response_body = error[:4000] if not delivery.response_body else delivery.response_body if self.request.retries < max_retries: countdown = initial_delay * (2**self.request.retries) if backoff == "exponential" else initial_delay delivery.next_retry_at = timezone.now() + timedelta(seconds=countdown) delivery.save() logger.warning("Webhook delivery failed (attempt %s): %s", self.request.retries + 1, error) raise self.retry(exc=Exception(error), countdown=countdown) from None delivery.status = OutboundDelivery.Status.FAILED delivery.completed_at = timezone.now() delivery.save() outbound_delivery_failed.send( sender=OutboundDelivery, delivery_id=delivery.id, event_name=delivery.subscription.event_name, error=error, ) logger.error("Webhook delivery gave up: %s", error) return None @shared_task( name="infrasynth.webhooks.process_inbound_event", bind=True, max_retries=3, default_retry_delay=60, ) def process_inbound_event(self, event_id, tenant_id=None): """Runs the endpoint's handler against a verified inbound event.""" from django.utils.module_loading import import_string from infrasynth.tenancy.context import tenant_context from infrasynth.tenancy.models import Tenant from .inbound.handlers import HMACInboundHandler from .models import InboundEvent from .signals import inbound_event_processed try: event = InboundEvent.all_objects.select_related("endpoint").get(pk=event_id) except InboundEvent.DoesNotExist: logger.warning("Inbound event %s not found", event_id) return None tenant = Tenant.objects.filter(pk=tenant_id).first() if tenant_id else event.tenant with tenant_context(tenant): handler = None if event.endpoint.handler: try: handler = import_string(event.endpoint.handler)() except (ImportError, TypeError): handler = None handler = handler or HMACInboundHandler() try: result = handler.process(event.event_type, event.raw_payload) except Exception as exc: # noqa: BLE001 logger.exception("Inbound event %s processing failed", event_id) event.error = str(exc) event.save(update_fields=["error"]) raise self.retry(exc=exc) from exc event.result = result if isinstance(result, dict) else {"result": result} event.is_processed = True event.error = "" event.processed_at = timezone.now() event.save(update_fields=["result", "is_processed", "error", "processed_at"]) inbound_event_processed.send( sender=InboundEvent, tenant_id=str(event.tenant_id) if event.tenant_id else None, event_id=event.id, event_type=event.event_type, result=event.result, ) return event.id