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