appcore_distributed_contracts/
opaque.rs1use appcore_types::{CapabilityName, InstanceId, TenantId};
14use serde::{Deserialize, Serialize};
15use std::collections::{HashSet, VecDeque};
16use std::fmt::{Debug, Formatter};
17use std::sync::Arc;
18
19pub const OPAQUE_CONTENT_ENVELOPE_SCHEMA_V1: &str = "appcore.opaque-content.v1";
21pub const MAX_OPAQUE_MESSAGE_ID_BYTES: usize = 1_024;
23
24fn is_valid_message_id(message_id: &str) -> bool {
25 !message_id.trim().is_empty()
26 && message_id.len() <= MAX_OPAQUE_MESSAGE_ID_BYTES
27 && !message_id.chars().any(char::is_control)
28}
29
30#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
32pub struct OpaqueContentEnvelopeV1 {
33 pub schema: String,
35 pub tenant_id: TenantId,
37 pub sender_instance_id: InstanceId,
39 pub message_id: String,
42 pub correlation_id: Option<String>,
44 pub capability: CapabilityName,
46 pub content_envelope: Vec<u8>,
48 pub envelope_version: u16,
50 pub created_at_ms: u64,
52 pub expires_at_ms: u64,
54 pub priority: Option<u8>,
56}
57
58impl Debug for OpaqueContentEnvelopeV1 {
59 fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result {
60 formatter
61 .debug_struct("OpaqueContentEnvelopeV1")
62 .field("schema", &self.schema)
63 .field("tenant_id", &self.tenant_id)
64 .field("sender_instance_id", &self.sender_instance_id)
65 .field("message_id", &self.message_id)
66 .field("correlation_id", &self.correlation_id)
67 .field("capability", &self.capability)
68 .field("content_envelope_bytes", &self.content_envelope.len())
69 .field("envelope_version", &self.envelope_version)
70 .field("created_at_ms", &self.created_at_ms)
71 .field("expires_at_ms", &self.expires_at_ms)
72 .field("priority", &self.priority)
73 .finish()
74 }
75}
76
77impl OpaqueContentEnvelopeV1 {
78 #[allow(clippy::too_many_arguments)]
80 pub fn new(
81 tenant_id: TenantId,
82 sender_instance_id: InstanceId,
83 message_id: impl Into<String>,
84 correlation_id: Option<String>,
85 capability: CapabilityName,
86 content_envelope: Vec<u8>,
87 envelope_version: u16,
88 created_at_ms: u64,
89 expires_at_ms: u64,
90 priority: Option<u8>,
91 ) -> Self {
92 Self {
93 schema: OPAQUE_CONTENT_ENVELOPE_SCHEMA_V1.to_string(),
94 tenant_id,
95 sender_instance_id,
96 message_id: message_id.into(),
97 correlation_id,
98 capability,
99 content_envelope,
100 envelope_version,
101 created_at_ms,
102 expires_at_ms,
103 priority,
104 }
105 }
106
107 pub fn validate_transport(
109 &self,
110 policy: &OpaqueEnvelopePolicy,
111 now_ms: u64,
112 ) -> OpaqueEnvelopeDecision {
113 if self.schema != OPAQUE_CONTENT_ENVELOPE_SCHEMA_V1 {
114 return OpaqueEnvelopeDecision::UnsupportedSchema;
115 }
116 if !is_valid_message_id(&self.message_id) || self.content_envelope.is_empty() {
117 return OpaqueEnvelopeDecision::InvalidEnvelope;
118 }
119 if self.created_at_ms >= self.expires_at_ms {
120 return OpaqueEnvelopeDecision::Expired;
121 }
122 if self.expires_at_ms <= now_ms {
123 return OpaqueEnvelopeDecision::Expired;
124 }
125 if !policy
126 .accepted_envelope_versions
127 .contains(&self.envelope_version)
128 {
129 return OpaqueEnvelopeDecision::UnsupportedEnvelopeVersion;
130 }
131 if self.content_envelope.len() as u64 > policy.max_payload_bytes {
132 return OpaqueEnvelopeDecision::PayloadTooLarge;
133 }
134 if !policy.accepted_capabilities.contains(&self.capability) {
135 return OpaqueEnvelopeDecision::UnsupportedCapability;
136 }
137 OpaqueEnvelopeDecision::Accepted
138 }
139}
140
141#[derive(Debug, Clone, PartialEq, Eq)]
143pub struct OpaqueEnvelopePolicy {
144 pub accepted_envelope_versions: Vec<u16>,
146 pub max_payload_bytes: u64,
148 pub accepted_capabilities: Vec<CapabilityName>,
150}
151
152#[derive(Debug, Clone, Copy, PartialEq, Eq)]
154pub enum OpaqueEnvelopeDecision {
155 Accepted,
157 UnsupportedSchema,
159 UnsupportedEnvelopeVersion,
161 Expired,
163 PayloadTooLarge,
165 UnsupportedCapability,
167 Duplicate,
169 InvalidEnvelope,
171}
172
173#[derive(Debug, Clone)]
178pub struct OpaqueEnvelopeDeduplicator {
179 max_entries: usize,
180 seen: HashSet<Arc<str>>,
181 order: VecDeque<Arc<str>>,
182}
183
184impl OpaqueEnvelopeDeduplicator {
185 pub fn new(max_entries: usize) -> Self {
187 Self {
188 max_entries: max_entries.max(1),
189 seen: HashSet::new(),
190 order: VecDeque::new(),
191 }
192 }
193
194 pub fn accept(&mut self, message_id: &str) -> OpaqueEnvelopeDecision {
199 if !is_valid_message_id(message_id) {
200 return OpaqueEnvelopeDecision::InvalidEnvelope;
201 }
202 if self.seen.contains(message_id) {
203 return OpaqueEnvelopeDecision::Duplicate;
204 }
205 let shared: Arc<str> = Arc::from(message_id);
206 self.seen.insert(Arc::clone(&shared));
207 self.order.push_back(shared);
208 while self.seen.len() > self.max_entries {
209 if let Some(oldest) = self.order.pop_front() {
210 self.seen.remove(oldest.as_ref());
211 }
212 }
213 OpaqueEnvelopeDecision::Accepted
214 }
215}
216
217#[cfg(test)]
218mod tests {
219 use super::*;
220
221 fn envelope(capability: &str, version: u16) -> OpaqueContentEnvelopeV1 {
222 OpaqueContentEnvelopeV1::new(
223 TenantId::new("tenant-a").unwrap(),
224 InstanceId::new("instance-a").unwrap(),
225 "message-a",
226 Some("corr-a".to_string()),
227 CapabilityName::new(capability).unwrap(),
228 vec![1, 2, 3],
229 version,
230 100,
231 200,
232 Some(5),
233 )
234 }
235
236 fn policy() -> OpaqueEnvelopePolicy {
237 OpaqueEnvelopePolicy {
238 accepted_envelope_versions: vec![1],
239 max_payload_bytes: 8,
240 accepted_capabilities: vec![CapabilityName::new("sync.consumer").unwrap()],
241 }
242 }
243
244 #[test]
245 fn opaque_envelope_validates_transport_only() {
246 assert_eq!(
247 envelope("sync.consumer", 1).validate_transport(&policy(), 150),
248 OpaqueEnvelopeDecision::Accepted
249 );
250 }
251
252 #[test]
253 fn opaque_envelope_rejects_old_consumer_version_and_capability() {
254 assert_eq!(
255 envelope("sync.consumer", 2).validate_transport(&policy(), 150),
256 OpaqueEnvelopeDecision::UnsupportedEnvelopeVersion
257 );
258 assert_eq!(
259 envelope("sync.publisher", 1).validate_transport(&policy(), 150),
260 OpaqueEnvelopeDecision::UnsupportedCapability
261 );
262 }
263
264 #[test]
265 fn opaque_envelope_deduplicates_message_id() {
266 let mut dedupe = OpaqueEnvelopeDeduplicator::new(16);
267 assert_eq!(dedupe.accept("message-a"), OpaqueEnvelopeDecision::Accepted);
268 assert_eq!(
269 dedupe.accept("message-a"),
270 OpaqueEnvelopeDecision::Duplicate
271 );
272 }
273
274 #[test]
275 fn opaque_envelope_deduplicator_evicts_in_acceptance_order() {
276 let mut dedupe = OpaqueEnvelopeDeduplicator::new(2);
277 assert_eq!(dedupe.accept("message-a"), OpaqueEnvelopeDecision::Accepted);
278 assert_eq!(dedupe.accept("message-b"), OpaqueEnvelopeDecision::Accepted);
279 assert_eq!(dedupe.accept("message-c"), OpaqueEnvelopeDecision::Accepted);
280 assert_eq!(
281 dedupe.accept("message-b"),
282 OpaqueEnvelopeDecision::Duplicate
283 );
284 assert_eq!(dedupe.accept("message-a"), OpaqueEnvelopeDecision::Accepted);
285 }
286
287 #[test]
288 fn opaque_envelope_deduplicator_bounds_ids_without_evicting_valid_state() {
289 let mut dedupe = OpaqueEnvelopeDeduplicator::new(1);
290 assert_eq!(dedupe.accept("message-a"), OpaqueEnvelopeDecision::Accepted);
291
292 for invalid in [
293 String::new(),
294 " \t ".to_string(),
295 "message\ncontrol".to_string(),
296 "é".repeat(513),
297 ] {
298 assert_eq!(
299 dedupe.accept(&invalid),
300 OpaqueEnvelopeDecision::InvalidEnvelope
301 );
302 }
303 assert_eq!(
304 dedupe.accept("message-a"),
305 OpaqueEnvelopeDecision::Duplicate
306 );
307 let maximum = "é".repeat(512);
308 assert_eq!(dedupe.accept(&maximum), OpaqueEnvelopeDecision::Accepted);
309 assert_eq!(dedupe.accept(&maximum), OpaqueEnvelopeDecision::Duplicate);
310 }
311
312 #[test]
313 fn opaque_envelope_rejects_malformed_transport_metadata() {
314 let mut malformed = envelope("sync.consumer", 1);
315 malformed.message_id.clear();
316 assert_eq!(
317 malformed.validate_transport(&policy(), 150),
318 OpaqueEnvelopeDecision::InvalidEnvelope
319 );
320
321 for message_id in ["message\ncontrol".to_string(), "é".repeat(513)] {
322 let mut malformed = envelope("sync.consumer", 1);
323 malformed.message_id = message_id;
324 assert_eq!(
325 malformed.validate_transport(&policy(), 150),
326 OpaqueEnvelopeDecision::InvalidEnvelope
327 );
328 }
329
330 let mut maximum = envelope("sync.consumer", 1);
331 maximum.message_id = "é".repeat(512);
332 assert_eq!(
333 maximum.validate_transport(&policy(), 150),
334 OpaqueEnvelopeDecision::Accepted
335 );
336
337 let mut expired_before_creation = envelope("sync.consumer", 1);
338 expired_before_creation.expires_at_ms = expired_before_creation.created_at_ms;
339 assert_eq!(
340 expired_before_creation.validate_transport(&policy(), 150),
341 OpaqueEnvelopeDecision::Expired
342 );
343 }
344
345 #[test]
346 fn opaque_envelope_debug_never_prints_content_bytes() {
347 let marker = b"secret-marker-must-not-appear";
348 let mut envelope = envelope("sync.consumer", 1);
349 envelope.content_envelope = marker.to_vec();
350
351 let output = format!("{envelope:?}");
352 assert!(!output.contains(std::str::from_utf8(marker).unwrap()));
353 assert!(output.contains("content_envelope_bytes"));
354 }
355}