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, ReplyType, RequestSpec, VerifiedRequest};
14
15use super::{now_ms, Inner, Link, LinkError};
16
17pub const DEFAULT_CALL_TIMEOUT: Duration = Duration::from_secs(5);
19pub 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#[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
55pub(super) struct Pending {
57 pub(super) request: VerifiedRequest,
58 pub(super) outcome: oneshot::Sender<Result<Value, LinkError>>,
59}
60
61impl Link {
62 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
129pub(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
176pub(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}