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