Skip to main content

appcore_distributed_contracts/
opaque.rs

1// =============================================================================
2//        #######
3//     ###       ###     F: opaque.rs
4//    ##   ## ##   ##    P: AppCore-Runtime
5//         ## ##
6//                       C: 2026/08/02 00:04:12 by dnettoRaw
7//    ##   ## ##   ##    U: 2026/08/02 12:48:56 by dnettoRaw
8//      ###########      S: 1.0.1-rc.8
9// =============================================================================
10
11//! Opaque content-envelope transport contracts.
12
13use appcore_types::{CapabilityName, InstanceId, TenantId};
14use serde::{Deserialize, Serialize};
15use std::collections::{HashSet, VecDeque};
16use std::fmt::{Debug, Formatter};
17use std::sync::Arc;
18
19/// Schema label for the first opaque content-envelope transport contract.
20pub const OPAQUE_CONTENT_ENVELOPE_SCHEMA_V1: &str = "appcore.opaque-content.v1";
21/// Maximum UTF-8 byte length accepted for an opaque message identifier.
22pub 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/// Gateway/sync transport metadata for an opaque encrypted content envelope.
31#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
32pub struct OpaqueContentEnvelopeV1 {
33    /// Stable schema label.
34    pub schema: String,
35    /// Tenant routing boundary.
36    pub tenant_id: TenantId,
37    /// Runtime instance that produced this envelope.
38    pub sender_instance_id: InstanceId,
39    /// Idempotency and deduplication identifier, bounded by
40    /// [`MAX_OPAQUE_MESSAGE_ID_BYTES`] at validation boundaries.
41    pub message_id: String,
42    /// Optional correlation identifier.
43    pub correlation_id: Option<String>,
44    /// Capability required to consume this envelope.
45    pub capability: CapabilityName,
46    /// Opaque encrypted content envelope, usually DNT bytes.
47    pub content_envelope: Vec<u8>,
48    /// Content-envelope format version, such as DNT envelope version.
49    pub envelope_version: u16,
50    /// Creation timestamp in Unix milliseconds.
51    pub created_at_ms: u64,
52    /// Expiration timestamp in Unix milliseconds.
53    pub expires_at_ms: u64,
54    /// Optional transport priority. Lower values are higher priority.
55    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    /// Creates a versioned opaque content envelope.
79    #[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    /// Validates transport metadata without opening the opaque payload.
108    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/// Transport policy for opaque content envelopes.
142#[derive(Debug, Clone, PartialEq, Eq)]
143pub struct OpaqueEnvelopePolicy {
144    /// Accepted opaque content-envelope versions.
145    pub accepted_envelope_versions: Vec<u16>,
146    /// Maximum opaque payload bytes.
147    pub max_payload_bytes: u64,
148    /// Capabilities this consumer can process.
149    pub accepted_capabilities: Vec<CapabilityName>,
150}
151
152/// Result of opaque transport validation.
153#[derive(Debug, Clone, Copy, PartialEq, Eq)]
154pub enum OpaqueEnvelopeDecision {
155    /// The envelope may be routed or consumed.
156    Accepted,
157    /// The transport schema is unknown.
158    UnsupportedSchema,
159    /// The opaque content-envelope version is not accepted.
160    UnsupportedEnvelopeVersion,
161    /// The envelope is expired.
162    Expired,
163    /// The opaque bytes exceed policy.
164    PayloadTooLarge,
165    /// No compatible consumer capability is available.
166    UnsupportedCapability,
167    /// The message ID was already observed.
168    Duplicate,
169    /// Required transport metadata is malformed.
170    InvalidEnvelope,
171}
172
173/// Bounded in-memory deduplicator keyed by opaque message ID.
174///
175/// Each accepted identifier has one shared allocation across membership and
176/// acceptance-order indexes.
177#[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    /// Creates a bounded deduplicator.
186    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    /// Records a valid bounded message ID and reports whether it was new.
195    ///
196    /// Invalid identifiers return [`OpaqueEnvelopeDecision::InvalidEnvelope`]
197    /// without changing the retained window.
198    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}