Skip to main content

appcore_distributed_contracts/
opaque.rs

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