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};
17
18/// Schema label for the first opaque content-envelope transport contract.
19pub const OPAQUE_CONTENT_ENVELOPE_SCHEMA_V1: &str = "appcore.opaque-content.v1";
20
21/// Gateway/sync transport metadata for an opaque encrypted content envelope.
22#[derive(Clone, PartialEq, Eq, Serialize, Deserialize)]
23pub struct OpaqueContentEnvelopeV1 {
24    /// Stable schema label.
25    pub schema: String,
26    /// Tenant routing boundary.
27    pub tenant_id: TenantId,
28    /// Runtime instance that produced this envelope.
29    pub sender_instance_id: InstanceId,
30    /// Idempotency and deduplication identifier.
31    pub message_id: String,
32    /// Optional correlation identifier.
33    pub correlation_id: Option<String>,
34    /// Capability required to consume this envelope.
35    pub capability: CapabilityName,
36    /// Opaque encrypted content envelope, usually DNT bytes.
37    pub content_envelope: Vec<u8>,
38    /// Content-envelope format version, such as DNT envelope version.
39    pub envelope_version: u16,
40    /// Creation timestamp in Unix milliseconds.
41    pub created_at_ms: u64,
42    /// Expiration timestamp in Unix milliseconds.
43    pub expires_at_ms: u64,
44    /// Optional transport priority. Lower values are higher priority.
45    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    /// Creates a versioned opaque content envelope.
69    #[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    /// Validates transport metadata without opening the opaque payload.
98    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/// Transport policy for opaque content envelopes.
132#[derive(Debug, Clone, PartialEq, Eq)]
133pub struct OpaqueEnvelopePolicy {
134    /// Accepted opaque content-envelope versions.
135    pub accepted_envelope_versions: Vec<u16>,
136    /// Maximum opaque payload bytes.
137    pub max_payload_bytes: u64,
138    /// Capabilities this consumer can process.
139    pub accepted_capabilities: Vec<CapabilityName>,
140}
141
142/// Result of opaque transport validation.
143#[derive(Debug, Clone, Copy, PartialEq, Eq)]
144pub enum OpaqueEnvelopeDecision {
145    /// The envelope may be routed or consumed.
146    Accepted,
147    /// The transport schema is unknown.
148    UnsupportedSchema,
149    /// The opaque content-envelope version is not accepted.
150    UnsupportedEnvelopeVersion,
151    /// The envelope is expired.
152    Expired,
153    /// The opaque bytes exceed policy.
154    PayloadTooLarge,
155    /// No compatible consumer capability is available.
156    UnsupportedCapability,
157    /// The message ID was already observed.
158    Duplicate,
159    /// Required transport metadata is malformed.
160    InvalidEnvelope,
161}
162
163/// Bounded in-memory deduplicator keyed by opaque message ID.
164#[derive(Debug, Clone)]
165pub struct OpaqueEnvelopeDeduplicator {
166    max_entries: usize,
167    seen: HashSet<String>,
168    order: VecDeque<String>,
169}
170
171impl OpaqueEnvelopeDeduplicator {
172    /// Creates a bounded deduplicator.
173    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    /// Records a message ID and reports whether it was new.
182    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}