use crate::models::ZeroMqFormat;
use crate::traits::PublisherError;
use crate::CanonicalMessage;
use anyhow::anyhow;
use bytes::Bytes;
use tracing::{debug, warn};
pub(crate) fn encode_frames(
message: &mut CanonicalMessage,
format: &ZeroMqFormat,
) -> Result<Vec<Bytes>, PublisherError> {
if matches!(format, ZeroMqFormat::RawFramed) {
message.strip_source_metadata();
let meta = serde_json::to_vec(&message.metadata)
.map_err(|e| PublisherError::NonRetryable(anyhow!(e)))?;
Ok(vec![Bytes::from(meta), message.payload.clone()])
} else {
Ok(vec![message.payload.clone()])
}
}
pub(crate) fn decode_frames(
frames: Vec<Bytes>,
is_sub: bool,
format: &ZeroMqFormat,
) -> anyhow::Result<Vec<CanonicalMessage>> {
let payload = frames.last().cloned().unwrap_or_default();
if payload.is_empty() && frames.len() <= 1 {
return Ok(vec![]);
}
let topic = if is_sub && frames.len() > 1 {
Some(String::from_utf8_lossy(frames[0].as_ref()).into_owned())
} else {
None
};
let mut messages = match format {
ZeroMqFormat::Raw => {
let payload_frames = if is_sub && frames.len() > 1 {
&frames[1..]
} else {
&frames[..]
};
payload_frames
.iter()
.map(|f| CanonicalMessage::new_bytes(f.clone(), None))
.collect()
}
ZeroMqFormat::RawFramed => {
let framed: &[Bytes] = if is_sub && frames.len() > 1 {
&frames[1..]
} else {
&frames[..]
};
let payload = framed.last().cloned().unwrap_or_default();
let mut msg = CanonicalMessage::new_bytes(payload, None);
if framed.len() > 1 {
match serde_json::from_slice::<std::collections::HashMap<String, String>>(
framed[0].as_ref(),
) {
Ok(meta) => msg.metadata = meta,
Err(e) => warn!(
is_sub,
frames = frames.len(),
error = %e,
"zeromq raw_framed: leading frame is not valid JSON metadata; treating the message as payload-only. Check the producer's framing."
),
}
if framed.len() > 2 {
debug!(
is_sub,
discarded = framed.len() - 2,
"zeromq raw_framed: discarding frames beyond [metadata, payload]; a raw_framed message carries at most two payload frames. Check the producer's framing."
);
}
}
vec![msg]
}
ZeroMqFormat::Json => {
if let Ok(messages) = serde_json::from_slice::<Vec<CanonicalMessage>>(&payload) {
messages
} else if let Ok(message) = serde_json::from_slice::<CanonicalMessage>(&payload) {
vec![message]
} else {
vec![CanonicalMessage::new(payload.to_vec(), None)]
}
}
};
for message in &mut messages {
message.strip_source_metadata();
}
if crate::canonical_message::source_metadata_enabled() {
if let Some(topic) = topic {
for message in &mut messages {
message
.metadata
.insert("mqb.src.zeromq_topic".to_string(), topic.clone());
}
}
}
Ok(messages)
}