use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use crate::ids::{
IdempotencyKey, JobId, MessageId, SessionId, SpanId, StreamId, SubscriptionId, TraceId,
};
use crate::messages::MessageType;
use crate::PROTOCOL_VERSION;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum Priority {
Low,
#[default]
Normal,
High,
Critical,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Envelope {
pub arcp: String,
pub id: MessageId,
pub timestamp: DateTime<Utc>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub target: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub session_id: Option<SessionId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub job_id: Option<JobId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub stream_id: Option<StreamId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub subscription_id: Option<SubscriptionId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub trace_id: Option<TraceId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub span_id: Option<SpanId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent_span_id: Option<SpanId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub correlation_id: Option<MessageId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub causation_id: Option<MessageId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idempotency_key: Option<IdempotencyKey>,
#[serde(default, skip_serializing_if = "is_default_priority")]
pub priority: Priority,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<serde_json::Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub event_seq: Option<u64>,
#[serde(flatten)]
pub payload: MessageType,
}
#[allow(clippy::trivially_copy_pass_by_ref)]
const fn is_default_priority(p: &Priority) -> bool {
matches!(p, Priority::Normal)
}
impl Envelope {
#[must_use]
pub fn new(payload: MessageType) -> Self {
Self {
arcp: PROTOCOL_VERSION.to_owned(),
id: MessageId::new(),
timestamp: Utc::now(),
source: None,
target: None,
session_id: None,
job_id: None,
stream_id: None,
subscription_id: None,
trace_id: None,
span_id: None,
parent_span_id: None,
correlation_id: None,
causation_id: None,
idempotency_key: None,
priority: Priority::default(),
extensions: None,
event_seq: None,
payload,
}
}
pub fn into_raw(self) -> Result<RawEnvelope, crate::error::ARCPError> {
let value = serde_json::to_value(&self)?;
Ok(serde_json::from_value(value)?)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RawEnvelope {
pub arcp: String,
pub id: MessageId,
pub timestamp: DateTime<Utc>,
#[serde(rename = "type")]
pub type_name: String,
#[serde(default)]
pub payload: serde_json::Value,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub target: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub session_id: Option<SessionId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub job_id: Option<JobId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub stream_id: Option<StreamId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub subscription_id: Option<SubscriptionId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub trace_id: Option<TraceId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub span_id: Option<SpanId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent_span_id: Option<SpanId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub correlation_id: Option<MessageId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub causation_id: Option<MessageId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idempotency_key: Option<IdempotencyKey>,
#[serde(default, skip_serializing_if = "is_default_priority")]
pub priority: Priority,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub extensions: Option<serde_json::Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub event_seq: Option<u64>,
}
impl RawEnvelope {
pub fn try_into_typed(self) -> Result<Envelope, crate::error::ARCPError> {
let value = serde_json::to_value(&self)?;
Ok(serde_json::from_value(value)?)
}
}
#[cfg(test)]
#[allow(
clippy::expect_used,
clippy::unwrap_used,
clippy::panic,
clippy::missing_panics_doc
)]
mod tests {
use chrono::TimeZone;
use super::*;
use crate::ids::{MessageId, SessionId};
use crate::messages::{LogLevel, LogPayload, PingPayload};
fn fixed_timestamp() -> DateTime<Utc> {
Utc.with_ymd_and_hms(2026, 5, 7, 21, 30, 0).unwrap()
}
#[test]
fn envelope_round_trips_through_serde() {
let env = Envelope::new(MessageType::Ping(PingPayload {
nonce: Some("n".into()),
}));
let json = serde_json::to_string(&env).expect("serialize");
let back: Envelope = serde_json::from_str(&json).expect("deserialize");
assert_eq!(env, back);
}
#[test]
fn envelope_wire_format_is_flat() {
let mut env = Envelope::new(MessageType::Ping(PingPayload { nonce: None }));
env.id = "msg_01JABC0123456789ABCDEFGHJK".parse().expect("valid id");
env.timestamp = fixed_timestamp();
let value = serde_json::to_value(&env).expect("serialize");
assert_eq!(value["type"], "ping");
assert!(value.get("payload").is_some());
assert_eq!(value["arcp"], "1.1");
assert_eq!(value["id"], "msg_01JABC0123456789ABCDEFGHJK");
}
#[test]
fn envelope_with_optional_fields_round_trips() {
let mut env = Envelope::new(MessageType::Log(LogPayload {
level: LogLevel::Warn,
message: "retrying".into(),
attributes: None,
}));
env.session_id = Some(SessionId::new());
env.trace_id = Some(TraceId::new("trace_789").expect("non-empty"));
env.correlation_id = Some(MessageId::new());
env.priority = Priority::High;
let json = serde_json::to_string(&env).expect("serialize");
let back: Envelope = serde_json::from_str(&json).expect("deserialize");
assert_eq!(env, back);
}
#[test]
fn priority_default_is_omitted_on_serialize() {
let env = Envelope::new(MessageType::Ping(PingPayload::default()));
let value = serde_json::to_value(&env).expect("serialize");
assert!(
value.get("priority").is_none(),
"default priority must be elided"
);
}
#[test]
fn priority_round_trips_each_variant() {
for p in [
Priority::Low,
Priority::Normal,
Priority::High,
Priority::Critical,
] {
let s = serde_json::to_string(&p).expect("serialize");
let back: Priority = serde_json::from_str(&s).expect("deserialize");
assert_eq!(p, back);
}
}
#[test]
fn raw_envelope_round_trips_to_typed() {
let env = Envelope::new(MessageType::Ping(PingPayload::default()));
let raw = env.clone().into_raw().expect("to raw");
assert_eq!(raw.type_name, "ping");
let back = raw.try_into_typed().expect("to typed");
assert_eq!(env, back);
}
#[test]
fn raw_envelope_unknown_type_does_not_fail_decode() {
let wire = serde_json::json!({
"arcp": "1.1",
"id": "msg_01JABC0123456789ABCDEFGHJK",
"timestamp": "2026-05-07T21:30:00Z",
"type": "arcpx.example.v1",
"payload": {"hello": "world"},
});
let raw: RawEnvelope = serde_json::from_value(wire).expect("raw parse");
assert_eq!(raw.type_name, "arcpx.example.v1");
assert_eq!(raw.payload["hello"], "world");
assert!(raw.try_into_typed().is_err());
}
}