Source code for palaestrai.core.serialisation

from __future__ import annotations

from typing import Any, Dict, List, Optional, Tuple
import logging

import pickle
from pyzstd import compress, decompress
from palaestrai.util.otel import get_tracer

LOG = logging.getLogger(__name__)


[docs] def serialize(request: Any) -> bytes: """Serialize a message, injecting OTel trace context.""" tracer = get_tracer() msg_type = type(request).__name__ with tracer.start_as_current_span( "serialisation.serialize", attributes={ "palaestrai.message_type": msg_type, }, ) as span: trace_context = getattr(request, "_otel_trace_context", None) trace_context_source = "attached" if not isinstance(trace_context, dict): trace_context_source = "current" trace_context = {} try: from opentelemetry import propagate propagate.inject(trace_context) except ImportError: pass LOG.debug( "OTEL_CTX serialize msg=%s src=%s sender=%s receiver=%s traceparent=%s", msg_type, trace_context_source, getattr(request, "sender", None), getattr(request, "receiver", None), trace_context.get("traceparent"), ) # Wrap the original object with trace context in a tuple envelope envelope = (request, trace_context) pick = pickle.dumps(envelope) payload = compress(pick) span.set_attribute( "palaestrai.trace_context_source", trace_context_source ) span.set_attribute("palaestrai.serialized_size", len(payload)) return payload
[docs] def deserialize(response: List[bytes]) -> Tuple[Any, Optional[Dict[str, str]]]: """Deserialize a message, extracting OTel trace context. Returns (message, trace_context_dict). """ tracer = get_tracer() with tracer.start_as_current_span("serialisation.deserialize") as span: if len(response) == 0: span.set_attribute("palaestrai.empty_response", True) return None, None span.set_attribute("palaestrai.empty_response", False) span.set_attribute("palaestrai.serialized_size", len(response[0])) decompressed = decompress(response[0]) data = pickle.loads(decompressed) # Handle both old-format (plain object) and new-format (envelope tuple) if ( isinstance(data, tuple) and len(data) == 2 and isinstance(data[1], dict) ): message, trace_context = data span.set_attribute( "palaestrai.message_type", type(message).__name__ ) span.set_attribute("palaestrai.trace_context_present", True) LOG.debug( "OTEL_CTX deserialize msg=%s sender=%s receiver=%s traceparent=%s", type(message).__name__, getattr(message, "sender", None), getattr(message, "receiver", None), trace_context.get("traceparent"), ) return message, trace_context # Backward compatibility: no trace context span.set_attribute("palaestrai.message_type", type(data).__name__) span.set_attribute("palaestrai.trace_context_present", False) LOG.debug( "OTEL_CTX deserialize-old msg=%s sender=%s receiver=%s traceparent=None", type(data).__name__, getattr(data, "sender", None), getattr(data, "receiver", None), ) return data, None