Skip to main content

lean_ctx/core/
capsule_transport.rs

1//! Local reference transport for authenticated OCLA A2A capsule delivery (P11).
2//!
3//! In-process delivery queue that accepts signed, payload-free capsule manifests.
4//! Bytes are measured exactly after serialization; the token count is a local
5//! tokenizer proxy. This is delivery evidence only — never convertible into
6//! compression savings or provider billing.
7
8use std::collections::{BTreeMap, VecDeque};
9use std::sync::Mutex;
10
11use serde::{Deserialize, Serialize};
12
13use crate::core::context_capsule::SignedContextCapsuleV1;
14
15const MAX_INBOX_SIZE: usize = 128;
16
17#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
18#[serde(rename_all = "snake_case")]
19pub enum DeliveredTokenCountKindV1 {
20    LocalTokenizerProxy,
21}
22
23/// Delivery receipt measured after a concrete local queue write.
24#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
25pub struct AgentRelayDeliveryV1 {
26    pub relay_id: String,
27    pub capsule_ref: String,
28    pub recipient_agent_id: String,
29    pub delivered_bytes: u64,
30    pub delivered_tokens: u64,
31    pub token_count_kind: DeliveredTokenCountKindV1,
32}
33
34#[derive(Clone)]
35struct QueuedCapsule {
36    signed_json: String,
37    delivery: AgentRelayDeliveryV1,
38}
39
40/// Bounded local transport. Deployments must place a durable, authenticated
41/// transport behind the same receipt semantics before remote use.
42pub struct LocalSignedCapsuleTransport {
43    inboxes: Mutex<BTreeMap<String, VecDeque<QueuedCapsule>>>,
44}
45
46impl Default for LocalSignedCapsuleTransport {
47    fn default() -> Self {
48        Self {
49            inboxes: Mutex::new(BTreeMap::new()),
50        }
51    }
52}
53
54impl LocalSignedCapsuleTransport {
55    /// Deliver a signed capsule to a recipient's inbox.
56    pub fn deliver(
57        &self,
58        signed: &SignedContextCapsuleV1,
59        recipient_agent_id: &str,
60    ) -> Result<AgentRelayDeliveryV1, TransportError> {
61        signed
62            .validate_structure()
63            .map_err(|e| TransportError::Validation(e.to_string()))?;
64
65        let envelope = signed
66            .capsule
67            .agent_envelope(recipient_agent_id)
68            .map_err(|e| TransportError::Validation(e.to_string()))?;
69
70        let signed_json = serde_json::to_string(signed)
71            .map_err(|e| TransportError::Serialization(e.to_string()))?;
72
73        let delivered_bytes = u64::try_from(signed_json.len()).unwrap_or(u64::MAX);
74        let delivered_tokens = estimate_tokens(&signed_json);
75
76        let delivery = AgentRelayDeliveryV1 {
77            relay_id: envelope.relay_id.clone(),
78            capsule_ref: envelope.capsule_ref,
79            recipient_agent_id: recipient_agent_id.to_string(),
80            delivered_bytes,
81            delivered_tokens,
82            token_count_kind: DeliveredTokenCountKindV1::LocalTokenizerProxy,
83        };
84
85        let mut inboxes = self
86            .inboxes
87            .lock()
88            .unwrap_or_else(std::sync::PoisonError::into_inner);
89        let inbox = inboxes.entry(recipient_agent_id.to_string()).or_default();
90
91        if inbox.len() >= MAX_INBOX_SIZE {
92            inbox.pop_front();
93        }
94
95        inbox.push_back(QueuedCapsule {
96            signed_json,
97            delivery: delivery.clone(),
98        });
99
100        Ok(delivery)
101    }
102
103    /// Receive the next capsule from an agent's inbox.
104    pub fn receive(
105        &self,
106        agent_id: &str,
107    ) -> Option<(SignedContextCapsuleV1, AgentRelayDeliveryV1)> {
108        let mut inboxes = self
109            .inboxes
110            .lock()
111            .unwrap_or_else(std::sync::PoisonError::into_inner);
112        let inbox = inboxes.get_mut(agent_id)?;
113        let queued = inbox.pop_front()?;
114        let signed: SignedContextCapsuleV1 = serde_json::from_str(&queued.signed_json).ok()?;
115        Some((signed, queued.delivery))
116    }
117
118    /// Peek at inbox depth for a given agent.
119    pub fn inbox_depth(&self, agent_id: &str) -> usize {
120        let inboxes = self
121            .inboxes
122            .lock()
123            .unwrap_or_else(std::sync::PoisonError::into_inner);
124        inboxes.get(agent_id).map_or(0, VecDeque::len)
125    }
126}
127
128/// Simple token estimation: ~4 chars per token (conservative for JSON).
129fn estimate_tokens(text: &str) -> u64 {
130    u64::try_from(text.len() / 4).unwrap_or(u64::MAX).max(1)
131}
132
133#[derive(Debug, thiserror::Error)]
134pub enum TransportError {
135    #[error("validation failed: {0}")]
136    Validation(String),
137    #[error("serialization failed: {0}")]
138    Serialization(String),
139}
140
141// ─── Tests ───────────────────────────────────────────────────────────────────
142
143#[cfg(test)]
144mod tests {
145    use super::*;
146    use crate::core::context_capsule::*;
147
148    fn test_signed_capsule() -> SignedContextCapsuleV1 {
149        let mut capsule = ContextCapsuleV1 {
150            schema_version: CONTEXT_CAPSULE_SCHEMA_VERSION,
151            capsule_id: "capsule:pending".to_string(),
152            request_id: "request:1".to_string(),
153            session_id: "session:1".to_string(),
154            agent_id: "sender-agent".to_string(),
155            intent_ref: "intent:test".to_string(),
156            task_ref: "task:1".to_string(),
157            expected_outcome_ref: "outcome:pass".to_string(),
158            acceptance_criteria_refs: vec!["criteria:ci-green".to_string()],
159            references: vec![ContextCapsuleReferenceV1 {
160                kind: CapsuleReferenceKindV1::File,
161                content_ref: "blake3:file1".to_string(),
162                freshness_ref: "freshness:1".to_string(),
163                recovery_ref: None,
164            }],
165            finding_refs: vec![],
166            decision_refs: vec![],
167            uncertainty_refs: vec![],
168            negative_result_refs: vec![],
169            source_ref: "source:test".to_string(),
170            policy_ref: "policy:default".to_string(),
171            contract_ref: "contract:v1".to_string(),
172            freshness_ref: "freshness:now".to_string(),
173            sensitivity: CapsuleSensitivityV1::Internal,
174            allowed_agent_ids: vec!["recipient-agent".to_string()],
175            budget: ContextCapsuleBudgetV1 {
176                tokens_used: 50,
177                tokens_remaining: 950,
178                cost_micros_used: 5,
179                cost_micros_remaining: 95,
180                latency_ms_used: 10,
181                latency_ms_remaining: 90,
182            },
183            chain: ContextCapsuleChainV1 {
184                chain_id: "chain:test".to_string(),
185                parent_capsule_ref: None,
186                owner_agent_id: "sender-agent".to_string(),
187                attribution_ref: "attribution:1".to_string(),
188                hop: 0,
189            },
190            quality_signal_refs: vec![],
191            recovery_refs: vec![],
192            delta_from: None,
193        };
194        capsule.assign_capsule_id().unwrap();
195        let keypair = ed25519_dalek::SigningKey::from_bytes(&[42u8; 32]);
196        SignedContextCapsuleV1::sign(&capsule, &keypair).unwrap()
197    }
198
199    #[test]
200    fn deliver_and_receive_roundtrip() {
201        let transport = LocalSignedCapsuleTransport::default();
202        let signed = test_signed_capsule();
203
204        let delivery = transport.deliver(&signed, "recipient-agent").unwrap();
205        assert!(delivery.delivered_bytes > 0);
206        assert!(delivery.delivered_tokens > 0);
207        assert_eq!(delivery.recipient_agent_id, "recipient-agent");
208        assert_eq!(transport.inbox_depth("recipient-agent"), 1);
209
210        let (received, receipt) = transport.receive("recipient-agent").unwrap();
211        assert_eq!(received.capsule.capsule_id, signed.capsule.capsule_id);
212        assert_eq!(receipt.relay_id, delivery.relay_id);
213        assert_eq!(transport.inbox_depth("recipient-agent"), 0);
214    }
215
216    #[test]
217    fn receive_from_empty_inbox_returns_none() {
218        let transport = LocalSignedCapsuleTransport::default();
219        assert!(transport.receive("nobody").is_none());
220    }
221
222    #[test]
223    fn rejects_unauthorized_recipient() {
224        let transport = LocalSignedCapsuleTransport::default();
225        let signed = test_signed_capsule();
226        assert!(transport.deliver(&signed, "unauthorized-agent").is_err());
227    }
228}