use appcore_types::{CapabilityName, InstanceId, TenantId};
use serde::{Deserialize, Serialize};
use std::collections::{HashSet, VecDeque};
use std::fmt::{Debug, Formatter};
use std::sync::Arc;
pub const OPAQUE_CONTENT_ENVELOPE_SCHEMA_V1: &str = "appcore.opaque-content.v1";
pub const MAX_OPAQUE_MESSAGE_ID_BYTES: usize = 1_024;
fn is_valid_message_id(message_id: &str) -> bool {
!message_id.trim().is_empty()
&& message_id.len() <= MAX_OPAQUE_MESSAGE_ID_BYTES
&& !message_id.chars().any(char::is_control)
}
#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OpaqueContentEnvelopeV1 {
pub schema: String,
pub tenant_id: TenantId,
pub sender_instance_id: InstanceId,
pub message_id: String,
pub correlation_id: Option<String>,
pub capability: CapabilityName,
pub content_envelope: Vec<u8>,
pub envelope_version: u16,
pub created_at_ms: u64,
pub expires_at_ms: u64,
pub priority: Option<u8>,
}
impl Debug for OpaqueContentEnvelopeV1 {
fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("OpaqueContentEnvelopeV1")
.field("schema", &self.schema)
.field("tenant_id", &self.tenant_id)
.field("sender_instance_id", &self.sender_instance_id)
.field("message_id", &self.message_id)
.field("correlation_id", &self.correlation_id)
.field("capability", &self.capability)
.field("content_envelope_bytes", &self.content_envelope.len())
.field("envelope_version", &self.envelope_version)
.field("created_at_ms", &self.created_at_ms)
.field("expires_at_ms", &self.expires_at_ms)
.field("priority", &self.priority)
.finish()
}
}
impl OpaqueContentEnvelopeV1 {
#[allow(clippy::too_many_arguments)]
pub fn new(
tenant_id: TenantId,
sender_instance_id: InstanceId,
message_id: impl Into<String>,
correlation_id: Option<String>,
capability: CapabilityName,
content_envelope: Vec<u8>,
envelope_version: u16,
created_at_ms: u64,
expires_at_ms: u64,
priority: Option<u8>,
) -> Self {
Self {
schema: OPAQUE_CONTENT_ENVELOPE_SCHEMA_V1.to_string(),
tenant_id,
sender_instance_id,
message_id: message_id.into(),
correlation_id,
capability,
content_envelope,
envelope_version,
created_at_ms,
expires_at_ms,
priority,
}
}
pub fn validate_transport(
&self,
policy: &OpaqueEnvelopePolicy,
now_ms: u64,
) -> OpaqueEnvelopeDecision {
if self.schema != OPAQUE_CONTENT_ENVELOPE_SCHEMA_V1 {
return OpaqueEnvelopeDecision::UnsupportedSchema;
}
if !is_valid_message_id(&self.message_id) || self.content_envelope.is_empty() {
return OpaqueEnvelopeDecision::InvalidEnvelope;
}
if self.created_at_ms >= self.expires_at_ms {
return OpaqueEnvelopeDecision::Expired;
}
if self.expires_at_ms <= now_ms {
return OpaqueEnvelopeDecision::Expired;
}
if !policy
.accepted_envelope_versions
.contains(&self.envelope_version)
{
return OpaqueEnvelopeDecision::UnsupportedEnvelopeVersion;
}
if self.content_envelope.len() as u64 > policy.max_payload_bytes {
return OpaqueEnvelopeDecision::PayloadTooLarge;
}
if !policy.accepted_capabilities.contains(&self.capability) {
return OpaqueEnvelopeDecision::UnsupportedCapability;
}
OpaqueEnvelopeDecision::Accepted
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct OpaqueEnvelopePolicy {
pub accepted_envelope_versions: Vec<u16>,
pub max_payload_bytes: u64,
pub accepted_capabilities: Vec<CapabilityName>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OpaqueEnvelopeDecision {
Accepted,
UnsupportedSchema,
UnsupportedEnvelopeVersion,
Expired,
PayloadTooLarge,
UnsupportedCapability,
Duplicate,
InvalidEnvelope,
}
#[derive(Debug, Clone)]
pub struct OpaqueEnvelopeDeduplicator {
max_entries: usize,
seen: HashSet<Arc<str>>,
order: VecDeque<Arc<str>>,
}
impl OpaqueEnvelopeDeduplicator {
pub fn new(max_entries: usize) -> Self {
Self {
max_entries: max_entries.max(1),
seen: HashSet::new(),
order: VecDeque::new(),
}
}
pub fn accept(&mut self, message_id: &str) -> OpaqueEnvelopeDecision {
if !is_valid_message_id(message_id) {
return OpaqueEnvelopeDecision::InvalidEnvelope;
}
if self.seen.contains(message_id) {
return OpaqueEnvelopeDecision::Duplicate;
}
let shared: Arc<str> = Arc::from(message_id);
self.seen.insert(Arc::clone(&shared));
self.order.push_back(shared);
while self.seen.len() > self.max_entries {
if let Some(oldest) = self.order.pop_front() {
self.seen.remove(oldest.as_ref());
}
}
OpaqueEnvelopeDecision::Accepted
}
}
#[cfg(test)]
mod tests {
use super::*;
fn envelope(capability: &str, version: u16) -> OpaqueContentEnvelopeV1 {
OpaqueContentEnvelopeV1::new(
TenantId::new("tenant-a").unwrap(),
InstanceId::new("instance-a").unwrap(),
"message-a",
Some("corr-a".to_string()),
CapabilityName::new(capability).unwrap(),
vec![1, 2, 3],
version,
100,
200,
Some(5),
)
}
fn policy() -> OpaqueEnvelopePolicy {
OpaqueEnvelopePolicy {
accepted_envelope_versions: vec![1],
max_payload_bytes: 8,
accepted_capabilities: vec![CapabilityName::new("sync.consumer").unwrap()],
}
}
#[test]
fn opaque_envelope_validates_transport_only() {
assert_eq!(
envelope("sync.consumer", 1).validate_transport(&policy(), 150),
OpaqueEnvelopeDecision::Accepted
);
}
#[test]
fn opaque_envelope_rejects_old_consumer_version_and_capability() {
assert_eq!(
envelope("sync.consumer", 2).validate_transport(&policy(), 150),
OpaqueEnvelopeDecision::UnsupportedEnvelopeVersion
);
assert_eq!(
envelope("sync.publisher", 1).validate_transport(&policy(), 150),
OpaqueEnvelopeDecision::UnsupportedCapability
);
}
#[test]
fn opaque_envelope_deduplicates_message_id() {
let mut dedupe = OpaqueEnvelopeDeduplicator::new(16);
assert_eq!(dedupe.accept("message-a"), OpaqueEnvelopeDecision::Accepted);
assert_eq!(
dedupe.accept("message-a"),
OpaqueEnvelopeDecision::Duplicate
);
}
#[test]
fn opaque_envelope_deduplicator_evicts_in_acceptance_order() {
let mut dedupe = OpaqueEnvelopeDeduplicator::new(2);
assert_eq!(dedupe.accept("message-a"), OpaqueEnvelopeDecision::Accepted);
assert_eq!(dedupe.accept("message-b"), OpaqueEnvelopeDecision::Accepted);
assert_eq!(dedupe.accept("message-c"), OpaqueEnvelopeDecision::Accepted);
assert_eq!(
dedupe.accept("message-b"),
OpaqueEnvelopeDecision::Duplicate
);
assert_eq!(dedupe.accept("message-a"), OpaqueEnvelopeDecision::Accepted);
}
#[test]
fn opaque_envelope_deduplicator_bounds_ids_without_evicting_valid_state() {
let mut dedupe = OpaqueEnvelopeDeduplicator::new(1);
assert_eq!(dedupe.accept("message-a"), OpaqueEnvelopeDecision::Accepted);
for invalid in [
String::new(),
" \t ".to_string(),
"message\ncontrol".to_string(),
"é".repeat(513),
] {
assert_eq!(
dedupe.accept(&invalid),
OpaqueEnvelopeDecision::InvalidEnvelope
);
}
assert_eq!(
dedupe.accept("message-a"),
OpaqueEnvelopeDecision::Duplicate
);
let maximum = "é".repeat(512);
assert_eq!(dedupe.accept(&maximum), OpaqueEnvelopeDecision::Accepted);
assert_eq!(dedupe.accept(&maximum), OpaqueEnvelopeDecision::Duplicate);
}
#[test]
fn opaque_envelope_rejects_malformed_transport_metadata() {
let mut malformed = envelope("sync.consumer", 1);
malformed.message_id.clear();
assert_eq!(
malformed.validate_transport(&policy(), 150),
OpaqueEnvelopeDecision::InvalidEnvelope
);
for message_id in ["message\ncontrol".to_string(), "é".repeat(513)] {
let mut malformed = envelope("sync.consumer", 1);
malformed.message_id = message_id;
assert_eq!(
malformed.validate_transport(&policy(), 150),
OpaqueEnvelopeDecision::InvalidEnvelope
);
}
let mut maximum = envelope("sync.consumer", 1);
maximum.message_id = "é".repeat(512);
assert_eq!(
maximum.validate_transport(&policy(), 150),
OpaqueEnvelopeDecision::Accepted
);
let mut expired_before_creation = envelope("sync.consumer", 1);
expired_before_creation.expires_at_ms = expired_before_creation.created_at_ms;
assert_eq!(
expired_before_creation.validate_transport(&policy(), 150),
OpaqueEnvelopeDecision::Expired
);
}
#[test]
fn opaque_envelope_debug_never_prints_content_bytes() {
let marker = b"secret-marker-must-not-appear";
let mut envelope = envelope("sync.consumer", 1);
envelope.content_envelope = marker.to_vec();
let output = format!("{envelope:?}");
assert!(!output.contains(std::str::from_utf8(marker).unwrap()));
assert!(output.contains("content_envelope_bytes"));
}
}