49 lines
2.3 KiB
Python
49 lines
2.3 KiB
Python
|
|
"""Bounded Audit Core transport for the acknowledgment outbox."""
|
||
|
|
import json
|
||
|
|
from urllib.error import HTTPError, URLError
|
||
|
|
from urllib.parse import urlsplit
|
||
|
|
from urllib.request import Request, HTTPRedirectHandler, ProxyHandler, build_opener
|
||
|
|
|
||
|
|
|
||
|
|
class NoRedirect(HTTPRedirectHandler):
|
||
|
|
def redirect_request(self, req, fp, code, msg, headers, newurl):
|
||
|
|
return None
|
||
|
|
|
||
|
|
|
||
|
|
class AuditTransport:
|
||
|
|
def __init__(self, origin, token_provider, *, allow_internal_http=False):
|
||
|
|
parsed = urlsplit(origin)
|
||
|
|
internal = parsed.hostname in ('127.0.0.1', '::1') or (parsed.hostname or '').endswith(('.svc', '.svc.cluster.local'))
|
||
|
|
if (parsed.scheme != 'https' and not (allow_internal_http and internal and parsed.scheme == 'http')
|
||
|
|
or not parsed.hostname or parsed.path or parsed.query or parsed.fragment
|
||
|
|
or parsed.username or parsed.password):
|
||
|
|
raise ValueError('fixed audit origin required')
|
||
|
|
self.url = origin + '/v1/events'
|
||
|
|
self.token_provider = token_provider
|
||
|
|
self.opener = build_opener(ProxyHandler({}), NoRedirect())
|
||
|
|
|
||
|
|
def __call__(self, event):
|
||
|
|
token = self.token_provider()
|
||
|
|
if (not isinstance(token, str) or not 1 <= len(token) <= 32768
|
||
|
|
or any(ord(c) < 33 or ord(c) > 126 for c in token)):
|
||
|
|
raise OSError('audit credential unavailable')
|
||
|
|
request = Request(self.url, data=json.dumps(event, sort_keys=True, separators=(',', ':')).encode(),
|
||
|
|
headers={'Content-Type': 'application/json', 'Authorization': 'Bearer ' + token,
|
||
|
|
'Idempotency-Key': event['id']}, method='POST')
|
||
|
|
try:
|
||
|
|
with self.opener.open(request, timeout=5) as response:
|
||
|
|
raw = response.read(65537)
|
||
|
|
if len(raw) > 65536:
|
||
|
|
raise OSError('audit response invalid')
|
||
|
|
body = json.loads(raw)
|
||
|
|
if not isinstance(body, dict):
|
||
|
|
raise ValueError()
|
||
|
|
return response.status, body
|
||
|
|
except HTTPError as error:
|
||
|
|
# Keep no upstream diagnostic body, URL, credential or response text.
|
||
|
|
status = error.code
|
||
|
|
error.close()
|
||
|
|
return status, {}
|
||
|
|
except (URLError, ValueError, UnicodeError):
|
||
|
|
raise OSError('audit unavailable or invalid response') from None
|