macula_rust/station_link/
call.rs1use std::sync::{Arc, Weak};
8use std::time::Duration;
9
10use tokio::sync::oneshot;
11
12use crate::cbor::Value;
13use crate::frame::{self, RequestSpec, VerifiedRequest};
14use crate::handshake::VERSION_5;
15
16use super::confidential::{reply_outcome, sealed_request, stated, CallSeal, Seal};
17use super::{now_ms, Inner, Link, LinkError};
18use crate::seal;
19
20pub const DEFAULT_CALL_TIMEOUT: Duration = Duration::from_secs(5);
22pub const MAX_CALL_TIMEOUT: Duration = Duration::from_secs(10 * 60);
24
25const LIVENESS_EVERY: Duration = Duration::from_secs(30);
26const LIVENESS_TIMEOUT: Duration = Duration::from_secs(30);
27const LIVENESS_PROCEDURE: &str = "_macula.ping";
28
29#[derive(Debug, Clone, PartialEq)]
38pub struct Call {
39 pub realm: [u8; 32],
40 pub procedure: String,
41 pub target: [u8; 32],
42 pub payload: Value,
43 pub timeout: Duration,
44 pub token: Option<Vec<u8>>,
45 pub proofs: Vec<Vec<u8>>,
46 pub seal: Option<Seal>,
47}
48
49impl Default for Call {
50 fn default() -> Self {
51 Call {
52 realm: [0; 32],
53 procedure: String::new(),
54 target: [0; 32],
55 payload: Value::Map(Vec::new()),
56 timeout: Duration::ZERO,
57 token: None,
58 proofs: Vec::new(),
59 seal: None,
60 }
61 }
62}
63
64pub(super) struct Pending {
66 pub(super) request: VerifiedRequest,
67 pub(super) seal: Option<CallSeal>,
68 pub(super) outcome: oneshot::Sender<Result<Value, LinkError>>,
69}
70
71impl Link {
72 pub async fn call(&self, c: Call) -> Result<Value, LinkError> {
75 call(&self.inner, c).await
76 }
77}
78
79pub(super) async fn call(inner: &Arc<Inner>, c: Call) -> Result<Value, LinkError> {
80 let timeout = if c.timeout.is_zero() {
81 DEFAULT_CALL_TIMEOUT
82 } else {
83 c.timeout.min(MAX_CALL_TIMEOUT)
84 };
85 stated(&c.target, &inner.station.node_id, &c.seal)?;
86 let target = if c.target == [0; 32] {
87 inner.station.node_id
88 } else {
89 c.target
90 };
91 let mut request_id = [0u8; 16];
92 aws_lc_rs::rand::fill(&mut request_id).map_err(|_| LinkError::Io("no randomness".into()))?;
93 let deadline = (now_ms() + timeout.as_millis() as i64) as u64;
94 let (sealed, call_seal) = sealed_call(inner, &c, target, request_id, deadline)?;
95 let signed = frame::sign_call(
96 &RequestSpec {
97 request_id,
98 realm: c.realm,
99 procedure: c.procedure,
100 target,
101 deadline,
102 payload: c.payload,
103 sealed,
104 mode: None,
105 token: c.token,
106 proofs: c.proofs,
107 source_route: None,
108 retry_budget: None,
109 },
110 &inner.key,
111 )?;
112 let request = frame::verify_request(&signed, inner.profile)?;
113 let (outcome_tx, outcome) = oneshot::channel();
114 await_reply(
115 inner,
116 request_id,
117 Pending {
118 request,
119 seal: call_seal,
120 outcome: outcome_tx,
121 },
122 )?;
123 let forget = || {
124 inner.lock().pending.remove(&request_id);
125 };
126 if let Err(e) = inner.write_control(&signed).await {
127 forget();
128 return Err(e);
129 }
130 let answered = tokio::time::timeout(timeout, outcome).await;
131 forget();
132 match answered {
133 Ok(Ok(outcome)) => outcome,
134 Ok(Err(_)) => Err(inner.lock().ended.clone().unwrap_or(LinkError::Closed)),
135 Err(_) => Err(LinkError::CallTimeout),
136 }
137}
138
139fn sealed_call(
142 inner: &Inner,
143 c: &Call,
144 target: [u8; 32],
145 request_id: [u8; 16],
146 deadline: u64,
147) -> Result<(Option<frame::Sealed>, Option<CallSeal>), LinkError> {
148 match &c.seal {
149 Some(Seal::To(key)) => {
150 let (sealed, s) = sealed_request(
151 inner.profile,
152 key,
153 seal::FRAME_CALL,
154 c.realm,
155 &c.procedure,
156 inner.self_id,
157 target,
158 request_id,
159 deadline,
160 &c.payload,
161 )?;
162 Ok((Some(sealed), Some(s)))
163 }
164 _ => Ok((None, None)),
165 }
166}
167
168fn await_reply(inner: &Inner, request_id: [u8; 16], pending: Pending) -> Result<(), LinkError> {
171 let mut state = inner.lock();
172 if let Some(e) = &state.ended {
173 return Err(e.clone());
174 }
175 state.pending.insert(request_id, pending);
176 Ok(())
177}
178
179pub(super) fn replied(inner: &Arc<Inner>, v: &Value) {
182 let Ok((request_id, _)) = frame::claimed_reply_ids(v) else {
183 inner.count("malformed_reply");
184 return;
185 };
186 let (request, seal) = match inner.lock().pending.get(&request_id) {
187 Some(p) => (p.request.clone(), p.seal.clone()),
188 None => {
189 inner.count("unmatched_reply");
190 return;
191 }
192 };
193 let Some(outcome) = verified_outcome(inner, v, &request, seal.as_ref()) else {
194 inner.count("unverified_reply");
195 return;
196 };
197 if let Some(p) = inner.lock().pending.remove(&request_id) {
198 let _ = p.outcome.send(outcome);
199 }
200}
201
202fn verified_outcome(
203 inner: &Inner,
204 v: &Value,
205 request: &VerifiedRequest,
206 seal: Option<&CallSeal>,
207) -> Option<Result<Value, LinkError>> {
208 if v.get("reply").is_some() {
209 let reply = frame::verify_reply(v, request, inner.profile).ok()?;
210 return Some(reply_outcome(reply, request, seal));
211 }
212 let relayed =
213 frame::verify_relay_error(v, request, inner.profile, &inner.station.node_id).ok()?;
214 Some(Err(LinkError::Relay {
215 reported_by: relayed.reported_by,
216 code: relayed.code,
217 }))
218}
219
220pub(super) async fn probe(link: Weak<Inner>) {
226 let mut misses = 0;
227 loop {
228 let Some(mut done) = link.upgrade().map(|l| l.done_rx.clone()) else {
229 return;
230 };
231 tokio::select! {
232 _ = done.wait_for(|ended| *ended) => return,
233 _ = tokio::time::sleep(LIVENESS_EVERY) => {}
234 }
235 let Some(inner) = link.upgrade() else { return };
236 let outcome = match inner.version {
237 VERSION_5 => probe_v5(&inner).await,
238 _ => call(
239 &inner,
240 Call {
241 procedure: LIVENESS_PROCEDURE.to_string(),
242 timeout: LIVENESS_TIMEOUT,
243 ..Call::default()
244 },
245 )
246 .await
247 .map(|_| ()),
248 };
249 misses = if matches!(outcome, Err(LinkError::CallTimeout)) {
250 misses + 1
251 } else {
252 0
253 };
254 if misses >= 2 {
255 inner.end(LinkError::LivenessLost);
256 return;
257 }
258 }
259}
260
261async fn probe_v5(inner: &Inner) -> Result<(), LinkError> {
264 let mut nonce = [0u8; frame::LIVENESS_NONCE_SIZE];
265 aws_lc_rs::rand::fill(&mut nonce).map_err(|_| LinkError::Io("no randomness".into()))?;
266 let mut pongs = inner.pongs_rx.lock().await;
267 while pongs.try_recv().is_ok() {}
270 inner
271 .send_control(&frame::liveness_ping_frame(&nonce))
272 .await?;
273 let mut done = inner.done_rx.clone();
274 let answered = pong_of(&mut pongs, nonce);
275 tokio::select! {
276 _ = done.wait_for(|ended| *ended) => Err(LinkError::Closed),
277 outcome = tokio::time::timeout(LIVENESS_TIMEOUT, answered) => {
278 outcome.unwrap_or(Err(LinkError::CallTimeout))
279 }
280 }
281}
282
283async fn pong_of(
286 pongs: &mut tokio::sync::mpsc::Receiver<[u8; frame::LIVENESS_NONCE_SIZE]>,
287 nonce: [u8; frame::LIVENESS_NONCE_SIZE],
288) -> Result<(), LinkError> {
289 while let Some(pong) = pongs.recv().await {
290 if pong == nonce {
291 return Ok(());
292 }
293 }
294 Err(LinkError::Closed)
295}