appcore_distributed_contracts/
opaque.rs1use appcore_types::{CapabilityName, InstanceId, TenantId};
14use serde::{Deserialize, Serialize};
15use std::collections::{HashSet, VecDeque};
16use std::fmt::{Debug, Formatter};
17
18pub const OPAQUE_CONTENT_ENVELOPE_SCHEMA_V1: &str = "appcore.opaque-content.v1";
20
21#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
23pub struct OpaqueContentEnvelopeV1 {
24 pub schema: String,
26 pub tenant_id: TenantId,
28 pub sender_instance_id: InstanceId,
30 pub message_id: String,
32 pub correlation_id: Option<String>,
34 pub capability: CapabilityName,
36 pub content_envelope: Vec<u8>,
38 pub envelope_version: u16,
40 pub created_at_ms: u64,
42 pub expires_at_ms: u64,
44 pub priority: Option<u8>,
46}
47
48impl Debug for OpaqueContentEnvelopeV1 {
49 fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result {
50 formatter
51 .debug_struct("OpaqueContentEnvelopeV1")
52 .field("schema", &self.schema)
53 .field("tenant_id", &self.tenant_id)
54 .field("sender_instance_id", &self.sender_instance_id)
55 .field("message_id", &self.message_id)
56 .field("correlation_id", &self.correlation_id)
57 .field("capability", &self.capability)
58 .field("content_envelope_bytes", &self.content_envelope.len())
59 .field("envelope_version", &self.envelope_version)
60 .field("created_at_ms", &self.created_at_ms)
61 .field("expires_at_ms", &self.expires_at_ms)
62 .field("priority", &self.priority)
63 .finish()
64 }
65}
66
67impl OpaqueContentEnvelopeV1 {
68 #[allow(clippy::too_many_arguments)]
70 pub fn new(
71 tenant_id: TenantId,
72 sender_instance_id: InstanceId,
73 message_id: impl Into<String>,
74 correlation_id: Option<String>,
75 capability: CapabilityName,
76 content_envelope: Vec<u8>,
77 envelope_version: u16,
78 created_at_ms: u64,
79 expires_at_ms: u64,
80 priority: Option<u8>,
81 ) -> Self {
82 Self {
83 schema: OPAQUE_CONTENT_ENVELOPE_SCHEMA_V1.to_string(),
84 tenant_id,
85 sender_instance_id,
86 message_id: message_id.into(),
87 correlation_id,
88 capability,
89 content_envelope,
90 envelope_version,
91 created_at_ms,
92 expires_at_ms,
93 priority,
94 }
95 }
96
97 pub fn validate_transport(
99 &self,
100 policy: &OpaqueEnvelopePolicy,
101 now_ms: u64,
102 ) -> OpaqueEnvelopeDecision {
103 if self.schema != OPAQUE_CONTENT_ENVELOPE_SCHEMA_V1 {
104 return OpaqueEnvelopeDecision::UnsupportedSchema;
105 }
106 if self.message_id.trim().is_empty() || self.content_envelope.is_empty() {
107 return OpaqueEnvelopeDecision::InvalidEnvelope;
108 }
109 if self.created_at_ms >= self.expires_at_ms {
110 return OpaqueEnvelopeDecision::Expired;
111 }
112 if self.expires_at_ms <= now_ms {
113 return OpaqueEnvelopeDecision::Expired;
114 }
115 if !policy
116 .accepted_envelope_versions
117 .contains(&self.envelope_version)
118 {
119 return OpaqueEnvelopeDecision::UnsupportedEnvelopeVersion;
120 }
121 if self.content_envelope.len() as u64 > policy.max_payload_bytes {
122 return OpaqueEnvelopeDecision::PayloadTooLarge;
123 }
124 if !policy.accepted_capabilities.contains(&self.capability) {
125 return OpaqueEnvelopeDecision::UnsupportedCapability;
126 }
127 OpaqueEnvelopeDecision::Accepted
128 }
129}
130
131#[derive(Debug, Clone, PartialEq, Eq)]
133pub struct OpaqueEnvelopePolicy {
134 pub accepted_envelope_versions: Vec<u16>,
136 pub max_payload_bytes: u64,
138 pub accepted_capabilities: Vec<CapabilityName>,
140}
141
142#[derive(Debug, Clone, Copy, PartialEq, Eq)]
144pub enum OpaqueEnvelopeDecision {
145 Accepted,
147 UnsupportedSchema,
149 UnsupportedEnvelopeVersion,
151 Expired,
153 PayloadTooLarge,
155 UnsupportedCapability,
157 Duplicate,
159 InvalidEnvelope,
161}
162
163#[derive(Debug, Clone)]
165pub struct OpaqueEnvelopeDeduplicator {
166 max_entries: usize,
167 seen: HashSet<String>,
168 order: VecDeque<String>,
169}
170
171impl OpaqueEnvelopeDeduplicator {
172 pub fn new(max_entries: usize) -> Self {
174 Self {
175 max_entries: max_entries.max(1),
176 seen: HashSet::new(),
177 order: VecDeque::new(),
178 }
179 }
180
181 pub fn accept(&mut self, message_id: &str) -> OpaqueEnvelopeDecision {
183 if self.seen.contains(message_id) {
184 return OpaqueEnvelopeDecision::Duplicate;
185 }
186 self.seen.insert(message_id.to_string());
187 self.order.push_back(message_id.to_string());
188 while self.seen.len() > self.max_entries {
189 if let Some(oldest) = self.order.pop_front() {
190 self.seen.remove(&oldest);
191 }
192 }
193 OpaqueEnvelopeDecision::Accepted
194 }
195}
196
197#[cfg(test)]
198mod tests {
199 use super::*;
200
201 fn envelope(capability: &str, version: u16) -> OpaqueContentEnvelopeV1 {
202 OpaqueContentEnvelopeV1::new(
203 TenantId::new("tenant-a").unwrap(),
204 InstanceId::new("instance-a").unwrap(),
205 "message-a",
206 Some("corr-a".to_string()),
207 CapabilityName::new(capability).unwrap(),
208 vec![1, 2, 3],
209 version,
210 100,
211 200,
212 Some(5),
213 )
214 }
215
216 fn policy() -> OpaqueEnvelopePolicy {
217 OpaqueEnvelopePolicy {
218 accepted_envelope_versions: vec![1],
219 max_payload_bytes: 8,
220 accepted_capabilities: vec![CapabilityName::new("sync.consumer").unwrap()],
221 }
222 }
223
224 #[test]
225 fn opaque_envelope_validates_transport_only() {
226 assert_eq!(
227 envelope("sync.consumer", 1).validate_transport(&policy(), 150),
228 OpaqueEnvelopeDecision::Accepted
229 );
230 }
231
232 #[test]
233 fn opaque_envelope_rejects_old_consumer_version_and_capability() {
234 assert_eq!(
235 envelope("sync.consumer", 2).validate_transport(&policy(), 150),
236 OpaqueEnvelopeDecision::UnsupportedEnvelopeVersion
237 );
238 assert_eq!(
239 envelope("sync.publisher", 1).validate_transport(&policy(), 150),
240 OpaqueEnvelopeDecision::UnsupportedCapability
241 );
242 }
243
244 #[test]
245 fn opaque_envelope_deduplicates_message_id() {
246 let mut dedupe = OpaqueEnvelopeDeduplicator::new(16);
247 assert_eq!(dedupe.accept("message-a"), OpaqueEnvelopeDecision::Accepted);
248 assert_eq!(
249 dedupe.accept("message-a"),
250 OpaqueEnvelopeDecision::Duplicate
251 );
252 }
253
254 #[test]
255 fn opaque_envelope_rejects_malformed_transport_metadata() {
256 let mut malformed = envelope("sync.consumer", 1);
257 malformed.message_id.clear();
258 assert_eq!(
259 malformed.validate_transport(&policy(), 150),
260 OpaqueEnvelopeDecision::InvalidEnvelope
261 );
262
263 let mut expired_before_creation = envelope("sync.consumer", 1);
264 expired_before_creation.expires_at_ms = expired_before_creation.created_at_ms;
265 assert_eq!(
266 expired_before_creation.validate_transport(&policy(), 150),
267 OpaqueEnvelopeDecision::Expired
268 );
269 }
270
271 #[test]
272 fn opaque_envelope_debug_never_prints_content_bytes() {
273 let marker = b"secret-marker-must-not-appear";
274 let mut envelope = envelope("sync.consumer", 1);
275 envelope.content_envelope = marker.to_vec();
276
277 let output = format!("{envelope:?}");
278 assert!(!output.contains(std::str::from_utf8(marker).unwrap()));
279 assert!(output.contains("content_envelope_bytes"));
280 }
281}