1use 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#[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
40pub 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 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 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 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
128fn 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#[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}