Skip to main content

macula_rust/station_link/
call.rs

1//! Calls on a link: a CALL signed with the link's identity key, answered by
2//! the RESULT or ERROR that verifies for it: a provider reply signed by its
3//! target, or a relay error the connected station signed. A reply that does
4//! not verify is counted and ignored, and the call keeps waiting, as macula's
5//! link does. On a v4 link the liveness probe is a call too.
6
7use 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
20/// macula's default timeout for a call.
21pub const DEFAULT_CALL_TIMEOUT: Duration = Duration::from_secs(5);
22/// The longest a call waits: the far edge of a provider's deadline window.
23pub 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/// One request: the realm and procedure, the node it targets (the station
30/// itself when zero, as for `_dht.*`; the provider's node_id otherwise), its
31/// payload, how long to wait (the default when zero), a UCAN token and its
32/// delegation chain's proofs for a gated procedure, and how it is kept. A
33/// call to a provider must state `seal`: sealed to the key its verified
34/// advertisement names, or clear by the application's decision; one that
35/// states neither is refused no_signed_state before anything is sent. A call
36/// to the station is always clear.
37#[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
64/// A call waiting for its reply, with what its sealed request agreed.
65pub(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    /// Signs `c` as a CALL, sends it, and waits for the RESULT or ERROR that
73    /// verifies for it.
74    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
139/// The call's payload sealed to the key its seal names, with what the sealing
140/// agreed; nothing when the call goes clear.
141fn 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
168/// Files the call as pending under its request id, unless the link has
169/// ended, in which case its end is the call's error.
170fn 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
179/// A RESULT or ERROR matched to its pending call by the ids it claims, and
180/// handed on only once it verifies for that call's request.
181pub(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
220/// A liveness probe every 30 seconds; two misses in a row end the link. On v5
221/// the probe is a liveness_ping, answered by the station's connection with a
222/// liveness_pong of its nonce, and nothing is signed. On v4 it is a
223/// `_macula.ping` call: any verified answer counts, a provider's or a relay
224/// error from the station alike, as it proves the station is there.
225pub(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
261/// One v5 liveness probe: answered when a liveness_pong of its nonce arrives
262/// within [`LIVENESS_TIMEOUT`], [`LinkError::CallTimeout`] when none does.
263async 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    // A pong that arrived after its probe's deadline must not take the place
268    // of this probe's.
269    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
283/// Waits for the liveness_pong of `nonce`, passing over any other;
284/// [`LinkError::Closed`] when the pongs stop.
285async 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}