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. 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, ReplyType, RequestSpec, VerifiedRequest};
14
15use super::{now_ms, Inner, Link, LinkError};
16
17/// macula's default timeout for a call.
18pub const DEFAULT_CALL_TIMEOUT: Duration = Duration::from_secs(5);
19/// The longest a call waits: the far edge of a provider's deadline window.
20pub const MAX_CALL_TIMEOUT: Duration = Duration::from_secs(10 * 60);
21
22const LIVENESS_EVERY: Duration = Duration::from_secs(30);
23const LIVENESS_TIMEOUT: Duration = Duration::from_secs(30);
24const LIVENESS_PROCEDURE: &str = "_macula.ping";
25
26/// One request: the realm and procedure, the node it targets (the station
27/// itself when zero, as for `_dht.*`; the provider's node_id otherwise), its
28/// payload, how long to wait (the default when zero), and a UCAN token and
29/// its delegation chain's proofs for a gated procedure.
30#[derive(Debug, Clone, PartialEq)]
31pub struct Call {
32    pub realm: [u8; 32],
33    pub procedure: String,
34    pub target: [u8; 32],
35    pub payload: Value,
36    pub timeout: Duration,
37    pub token: Option<Vec<u8>>,
38    pub proofs: Vec<Vec<u8>>,
39}
40
41impl Default for Call {
42    fn default() -> Self {
43        Call {
44            realm: [0; 32],
45            procedure: String::new(),
46            target: [0; 32],
47            payload: Value::Map(Vec::new()),
48            timeout: Duration::ZERO,
49            token: None,
50            proofs: Vec::new(),
51        }
52    }
53}
54
55/// A call waiting for its reply.
56pub(super) struct Pending {
57    pub(super) request: VerifiedRequest,
58    pub(super) outcome: oneshot::Sender<Result<Value, LinkError>>,
59}
60
61impl Link {
62    /// Signs `c` as a CALL, sends it, and waits for the RESULT or ERROR that
63    /// verifies for it.
64    pub async fn call(&self, c: Call) -> Result<Value, LinkError> {
65        call(&self.inner, c).await
66    }
67}
68
69pub(super) async fn call(inner: &Arc<Inner>, c: Call) -> Result<Value, LinkError> {
70    let timeout = if c.timeout.is_zero() {
71        DEFAULT_CALL_TIMEOUT
72    } else {
73        c.timeout.min(MAX_CALL_TIMEOUT)
74    };
75    let target = if c.target == [0; 32] {
76        inner.station.node_id
77    } else {
78        c.target
79    };
80    let mut request_id = [0u8; 16];
81    aws_lc_rs::rand::fill(&mut request_id).map_err(|_| LinkError::Io("no randomness".into()))?;
82    let signed = frame::sign_call(
83        &RequestSpec {
84            request_id,
85            realm: c.realm,
86            procedure: c.procedure,
87            target,
88            deadline: (now_ms() + timeout.as_millis() as i64) as u64,
89            payload: c.payload,
90            mode: None,
91            token: c.token,
92            proofs: c.proofs,
93            source_route: None,
94            retry_budget: None,
95        },
96        &inner.key,
97    )?;
98    let request = frame::verify_request(&signed, inner.profile)?;
99    let (outcome_tx, outcome) = oneshot::channel();
100    {
101        let mut state = inner.lock();
102        if let Some(e) = &state.ended {
103            return Err(e.clone());
104        }
105        state.pending.insert(
106            request_id,
107            Pending {
108                request,
109                outcome: outcome_tx,
110            },
111        );
112    }
113    let forget = || {
114        inner.lock().pending.remove(&request_id);
115    };
116    if let Err(e) = inner.write_control(&signed).await {
117        forget();
118        return Err(e);
119    }
120    let answered = tokio::time::timeout(timeout, outcome).await;
121    forget();
122    match answered {
123        Ok(Ok(outcome)) => outcome,
124        Ok(Err(_)) => Err(inner.lock().ended.clone().unwrap_or(LinkError::Closed)),
125        Err(_) => Err(LinkError::CallTimeout),
126    }
127}
128
129/// A RESULT or ERROR matched to its pending call by the ids it claims, and
130/// handed on only once it verifies for that call's request.
131pub(super) fn replied(inner: &Arc<Inner>, v: &Value) {
132    let Ok((request_id, _)) = frame::claimed_reply_ids(v) else {
133        inner.count("malformed_reply");
134        return;
135    };
136    let request = match inner.lock().pending.get(&request_id) {
137        Some(p) => p.request.clone(),
138        None => {
139            inner.count("unmatched_reply");
140            return;
141        }
142    };
143    let Some(outcome) = verified_outcome(inner, v, &request) else {
144        inner.count("unverified_reply");
145        return;
146    };
147    if let Some(p) = inner.lock().pending.remove(&request_id) {
148        let _ = p.outcome.send(outcome);
149    }
150}
151
152fn verified_outcome(
153    inner: &Inner,
154    v: &Value,
155    request: &VerifiedRequest,
156) -> Option<Result<Value, LinkError>> {
157    if v.get("reply").is_some() {
158        let reply = frame::verify_reply(v, request, inner.profile).ok()?;
159        return Some(match reply.frame_type {
160            ReplyType::Result => Ok(reply.payload.unwrap_or(Value::Null)),
161            ReplyType::Error => Err(LinkError::Provider {
162                responded_by: reply.responded_by,
163                code: reply.code.unwrap_or_default(),
164                detail: reply.detail,
165            }),
166        });
167    }
168    let relayed =
169        frame::verify_relay_error(v, request, inner.profile, &inner.station.node_id).ok()?;
170    Some(Err(LinkError::Relay {
171        reported_by: relayed.reported_by,
172        code: relayed.code,
173    }))
174}
175
176/// A liveness probe every 30 seconds; two misses in a row end the link. Any
177/// verified answer counts: it proves the station is there.
178pub(super) async fn probe(link: Weak<Inner>) {
179    let mut misses = 0;
180    loop {
181        let Some(mut done) = link.upgrade().map(|l| l.done_rx.clone()) else {
182            return;
183        };
184        tokio::select! {
185            _ = done.wait_for(|ended| *ended) => return,
186            _ = tokio::time::sleep(LIVENESS_EVERY) => {}
187        }
188        let Some(inner) = link.upgrade() else { return };
189        let outcome = call(
190            &inner,
191            Call {
192                procedure: LIVENESS_PROCEDURE.to_string(),
193                timeout: LIVENESS_TIMEOUT,
194                ..Call::default()
195            },
196        )
197        .await;
198        misses = if matches!(outcome, Err(LinkError::CallTimeout)) {
199            misses + 1
200        } else {
201            0
202        };
203        if misses >= 2 {
204            inner.end(LinkError::LivenessLost);
205            return;
206        }
207    }
208}