"""Restricted outbound transports for security audit evidence.""" from __future__ import annotations import json import socket import ssl from urllib.parse import urlsplit from urllib.request import HTTPRedirectHandler, HTTPSHandler, Request, build_opener class _NoRedirect(HTTPRedirectHandler): def redirect_request(self, req, fp, code, msg, headers, newurl): return None class HttpsWebhookSyslogTransport: """Deliver bounded safe envelopes after service-level endpoint validation.""" def __init__(self, *, timeout: float = 5.0, maximum_bytes: int = 1_048_576): self.timeout = float(timeout) self.maximum_bytes = int(maximum_bytes) def deliver(self, sink, envelope): payload = json.dumps( envelope, ensure_ascii=False, sort_keys=True, separators=(",", ":") ).encode("utf-8") if len(payload) > self.maximum_bytes: return {"status": "failed", "error_code": "payload_too_large"} try: if sink["sink_type"] == "webhook": return self._webhook(sink["endpoint"], payload) return self._syslog_tls(sink["endpoint"], payload) except (OSError, ssl.SSLError, TimeoutError): return {"status": "failed", "error_code": "transport_unavailable"} def _webhook(self, endpoint, payload): request = Request( endpoint, data=payload, headers={"Content-Type": "application/json", "User-Agent": "dataops-security/1"}, method="POST", ) opener = build_opener(HTTPSHandler(), _NoRedirect()) with opener.open(request, timeout=self.timeout) as response: status = int(response.status) if status < 200 or status >= 300: return {"status": "failed", "error_code": f"http_{status}"} remote_ref = str(response.headers.get("X-Request-ID") or "")[:300] or None return {"status": "delivered", "remote_ref": remote_ref} def _syslog_tls(self, endpoint, payload): parsed = urlsplit(endpoint) port = int(parsed.port or 6514) context = ssl.create_default_context() with ( socket.create_connection((parsed.hostname, port), timeout=self.timeout) as raw, context.wrap_socket(raw, server_hostname=parsed.hostname) as secured, ): secured.sendall(payload + b"\n") return {"status": "delivered"}