Skip to main content

macula_rust/pool/
call.rs

1//! Calls and streams that reach a provider at its own station, and the DHT
2//! through the pool's links.
3//!
4//! A call resolves the procedure's advertisements from the DHT, keeps those
5//! the realm's pinned key authorizes (or, in a node's own namespace, those
6//! that node signed), and tries the freshest first: it dials the serving
7//! station the advertisement names, pinned by its node_id from the station's
8//! own station_endpoint record, and calls the provider there. It moves on to
9//! the next candidate only when a station cannot be reached within that
10//! candidate's share of the deadline, before anything is sent: once the CALL
11//! or STREAM_OPEN has gone out, under the whole deadline, its outcome is
12//! returned as it is, a timeout included, so one call reaches a provider at
13//! most once (macula's call_work and failure_scope/1). A candidate that
14//! answered is remembered until its advertisement expires.
15
16use std::collections::HashSet;
17use std::future::Future;
18use std::sync::{Arc, Weak};
19use std::time::Duration;
20
21use tokio::time::Instant;
22
23use crate::cbor::Value;
24use crate::frame::StreamMode;
25use crate::record::{self, RecordType, Trust, Verified};
26use crate::seal::KEY_ID_SIZE;
27use crate::station_link::{
28    self, Confidentiality, ConfidentialityError, ConfidentialityReason, Link, LinkError, Report,
29    Reseal, Seal, Stream, DEFAULT_CALL_TIMEOUT,
30};
31use crate::transport::Target;
32
33use super::member::Member;
34use super::{Pool, PoolError, PoolInner};
35
36/// No candidate gets less than a second of a call's time.
37const MIN_CANDIDATE_SHARE: Duration = Duration::from_secs(1);
38
39/// A call to a procedure: its realm and name, the provider to call (any
40/// trusted one when zero), the payload, how long to wait
41/// ([`DEFAULT_CALL_TIMEOUT`] when zero), a UCAN and its proofs for a gated
42/// procedure, and how it must be kept.
43#[derive(Debug, Clone, PartialEq)]
44pub struct Call {
45    pub realm: [u8; 32],
46    pub procedure: String,
47    pub provider: [u8; 32],
48    pub payload: Value,
49    pub timeout: Duration,
50    pub token: Option<Vec<u8>>,
51    pub proofs: Vec<Vec<u8>>,
52    pub confidential: Confidentiality,
53}
54
55impl Default for Call {
56    fn default() -> Self {
57        Call {
58            realm: [0; 32],
59            procedure: String::new(),
60            provider: [0; 32],
61            payload: Value::Map(Vec::new()),
62            timeout: Duration::ZERO,
63            token: None,
64            proofs: Vec::new(),
65            confidential: Confidentiality::Preferred,
66        }
67    }
68}
69
70/// A streaming session to open: its realm and name, the provider (any
71/// trusted one when zero), the mode, the open's payload, its deadline (the
72/// link's default when zero), a UCAN and its proofs for a gated procedure,
73/// and how it must be kept. Its default mode is server_stream.
74#[derive(Debug, Clone, PartialEq)]
75pub struct StreamCall {
76    pub realm: [u8; 32],
77    pub procedure: String,
78    pub provider: [u8; 32],
79    pub mode: StreamMode,
80    pub payload: Value,
81    pub deadline: Duration,
82    pub token: Option<Vec<u8>>,
83    pub proofs: Vec<Vec<u8>>,
84    pub confidential: Confidentiality,
85}
86
87impl Default for StreamCall {
88    fn default() -> Self {
89        StreamCall {
90            realm: [0; 32],
91            procedure: String::new(),
92            provider: [0; 32],
93            mode: StreamMode::ServerStream,
94            payload: Value::Map(Vec::new()),
95            deadline: Duration::ZERO,
96            token: None,
97            proofs: Vec::new(),
98            confidential: Confidentiality::Preferred,
99        }
100    }
101}
102
103/// A node serving a procedure, and the station it serves from.
104#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
105pub struct Provider {
106    pub node: [u8; 32],
107    pub station: [u8; 32],
108}
109
110/// A trusted advertisement: its provider, serving station, times, and the
111/// KEM key it names, as carried.
112#[derive(Debug, Clone, PartialEq, Eq)]
113pub(super) struct Candidate {
114    provider: Provider,
115    expires_at: u64,
116    created_at: u64,
117    kem_key: Option<Vec<u8>>,
118}
119
120impl Candidate {
121    /// How a call to this candidate is kept: sealed to the key its
122    /// advertisement names, or clear when it names none.
123    fn seal(&self) -> Seal {
124        match &self.kem_key {
125            Some(key) => Seal::To(key.clone()),
126            None => Seal::Clear,
127        }
128    }
129
130    fn kem_key_id(&self) -> Option<[u8; KEY_ID_SIZE]> {
131        self.kem_key.as_deref().map(crate::seal::key_id)
132    }
133}
134
135#[derive(Debug, Clone, PartialEq, Eq, Hash)]
136pub(super) struct ResolvedKey {
137    realm: [u8; 32],
138    procedure: String,
139    provider: [u8; 32],
140}
141
142impl Pool {
143    /// Calls a procedure at a provider that serves it, as macula 12 calls
144    /// one. A provider's ERROR comes back as
145    /// `PoolError::Link(LinkError::Provider { .. })`; when no candidate
146    /// answers, [`PoolError::NoProvider`] names each one tried.
147    pub async fn call(&self, c: Call) -> Result<Value, PoolError> {
148        self.call_report(c).await.map(|(result, _)| result)
149    }
150
151    /// [`Pool::call`], returning with its result the call's seal report
152    /// (macula's DESIGN_E2E_SEAL_REPORT): `sealed` 1 with the id of the key
153    /// the request was sealed to, which is the key its answer opened under,
154    /// or 0 and no key for a clear call; `provider` is the node the call was
155    /// addressed to. An error comes with no report. It states that sealing
156    /// ran on this exchange, nothing more.
157    pub async fn call_report(&self, c: Call) -> Result<(Value, Report), PoolError> {
158        let inner = &self.inner;
159        let realm_key = inner.realm_key_for(&c.realm, &c.procedure)?;
160        let timeout = if c.timeout.is_zero() {
161            DEFAULT_CALL_TIMEOUT
162        } else {
163            c.timeout
164        };
165        let deadline = Instant::now() + timeout;
166        let key = ResolvedKey {
167            realm: c.realm,
168            procedure: c.procedure.clone(),
169            provider: c.provider,
170        };
171        let candidates = bounded(deadline, inner.candidates(&key, realm_key.clone())).await?;
172        let candidates = callable(candidates, c.confidential)?;
173        let (link, cand) = first_reached(inner, &key, candidates, deadline).await?;
174        let outcome = bounded(deadline, inner.call_at(&link, &cand, &c, deadline)).await;
175        let sent = Sent {
176            link: &link,
177            call: &c,
178            realm_key,
179            deadline,
180        };
181        let (cand, outcome) = inner.resealed(&key, cand, &sent, outcome).await;
182        inner.settled(key, cand, outcome)
183    }
184
185    /// Every provider whose advertisement of `procedure` in `realm` the
186    /// realm's pinned key authorizes, with the station each serves from,
187    /// freshest first.
188    pub async fn providers(
189        &self,
190        realm: &[u8; 32],
191        procedure: &str,
192    ) -> Result<Vec<Provider>, PoolError> {
193        let realm_key = self.inner.realm_key_for(realm, procedure)?;
194        let key = ResolvedKey {
195            realm: *realm,
196            procedure: procedure.to_string(),
197            provider: [0; 32],
198        };
199        let found = self.inner.resolve(&key, realm_key).await?;
200        Ok(found.into_iter().map(|c| c.provider).collect())
201    }
202
203    /// Opens a streaming session at a provider of the procedure, reached as
204    /// [`Pool::call`] reaches one: the next candidate only when a station
205    /// cannot be reached, and the link's own outcome is final. The stream is
206    /// open once its STREAM_OPEN is sent; a provider's or station's refusal
207    /// arrives on its first recv. A sealed stream refused sealed_refused
208    /// after its provider's key rotated, before it has sent anything,
209    /// reopens once behind the same handle, sealed to the key the provider
210    /// names, as a call reseals ([`Pool::call`]); never in the clear.
211    pub async fn open_stream(&self, c: StreamCall) -> Result<Stream, PoolError> {
212        let inner = &self.inner;
213        let realm_key = inner.realm_key_for(&c.realm, &c.procedure)?;
214        let deadline = Instant::now() + DEFAULT_CALL_TIMEOUT;
215        let key = ResolvedKey {
216            realm: c.realm,
217            procedure: c.procedure.clone(),
218            provider: c.provider,
219        };
220        let candidates = bounded(deadline, inner.candidates(&key, realm_key.clone())).await?;
221        let candidates = callable(candidates, c.confidential)?;
222        let (link, cand) = first_reached(inner, &key, candidates, deadline).await?;
223        let reseal = inner.stream_reseal(&key, &cand, &c, realm_key, &link);
224        let outcome = bounded(deadline, inner.open_at(&link, &cand, &c, reseal)).await;
225        inner.settled(key, cand, outcome)
226    }
227
228    /// Where `station` is dialed, from the station_endpoint record the
229    /// station signed itself; one another key signed is refused.
230    pub async fn station_target(&self, station: &[u8; 32]) -> Result<Target, PoolError> {
231        self.inner.station_target(station).await
232    }
233
234    /// A link to `station`: one the pool holds, or a direct link it dials to
235    /// the address in the station's own endpoint record, pinned by its
236    /// node_id, within the default call timeout. A direct link that does not
237    /// come up on its first dial is not kept.
238    pub async fn link_to(&self, station: &[u8; 32]) -> Result<Link, PoolError> {
239        let deadline = Instant::now() + DEFAULT_CALL_TIMEOUT;
240        self.inner.link_to(station, deadline).await
241    }
242
243    /// The record under `key`, verified, from the first link that answers.
244    pub async fn find_record(&self, key: &[u8; 32]) -> Result<Verified, PoolError> {
245        self.inner
246            .first_answer(|l| async move { l.find_record(key).await })
247            .await
248    }
249
250    /// The records under `key` that verify, and how many did not, from the
251    /// first link that answers.
252    pub async fn find_records(&self, key: &[u8; 32]) -> Result<(Vec<Verified>, usize), PoolError> {
253        self.inner
254            .first_answer(|l| async move { l.find_records(key).await })
255            .await
256    }
257
258    /// The records of type `t` that verify, and how many did not, from the
259    /// first link that answers.
260    pub async fn find_records_by_type(
261        &self,
262        t: RecordType,
263    ) -> Result<(Vec<Verified>, usize), PoolError> {
264        self.inner
265            .first_answer(|l| async move { l.find_records_by_type(t).await })
266            .await
267    }
268
269    /// Puts a signed record through the first link whose station takes it.
270    pub async fn put_record(&self, wire: &[u8]) -> Result<(), PoolError> {
271        self.inner
272            .first_answer(|l| async move { l.put_record(wire).await })
273            .await
274    }
275}
276
277impl PoolInner {
278    fn remember(&self, key: ResolvedKey, cand: Candidate) {
279        self.lock().remember.insert(key, cand);
280    }
281
282    fn forget(&self, key: &ResolvedKey) {
283        self.lock().remember.remove(key);
284    }
285
286    /// The remembered candidate while it lives and its station is linked,
287    /// else the procedure's trusted advertisements from the DHT.
288    async fn candidates(
289        &self,
290        key: &ResolvedKey,
291        realm_key: Option<Vec<u8>>,
292    ) -> Result<Vec<Candidate>, PoolError> {
293        let remembered = self.lock().remember.get(key).cloned();
294        let live = remembered.filter(|cand| {
295            cand.expires_at as i64 > now_ms() && self.linked_to(&cand.provider.station).is_some()
296        });
297        if let Some(cand) = live {
298            return Ok(vec![cand]);
299        }
300        self.resolve(key, realm_key).await
301    }
302
303    /// The procedure's advertisements the realm key authorizes, by
304    /// `key.provider` when it is set, freshest first.
305    async fn resolve(
306        &self,
307        key: &ResolvedKey,
308        realm_key: Option<Vec<u8>>,
309    ) -> Result<Vec<Candidate>, PoolError> {
310        let slot = record::procedure_key(&key.realm, &key.procedure);
311        let (found, _) = self
312            .first_answer(|l| async move { l.find_records(&slot).await })
313            .await?;
314        let now = now_ms();
315        let trust = Trust {
316            profile: self.opts.identity.profile(),
317            realm_key,
318        };
319        let mut out: Vec<Candidate> = found
320            .iter()
321            .filter(|v| v.record().record_type == RecordType::PROCEDURE_ADVERTISEMENT)
322            .filter_map(|v| trusted_candidate(v, key, &trust, now))
323            .collect();
324        if out.is_empty() {
325            return Err(PoolError::NoProvider(Vec::new()));
326        }
327        out.sort_by_key(|c| std::cmp::Reverse(c.created_at));
328        Ok(out)
329    }
330
331    /// A link to the candidate's serving station, reached by `share`.
332    /// Nothing is sent.
333    async fn reach(self: &Arc<Self>, cand: &Candidate, share: Instant) -> Result<Link, PoolError> {
334        bounded(share, self.link_to(&cand.provider.station, share)).await
335    }
336
337    /// The outcome of a call or stream sent to `cand`, final whatever it is:
338    /// a CALL that went out is never sent again elsewhere. A candidate whose
339    /// provider answered is remembered, one that did not is forgotten.
340    fn settled<T>(
341        &self,
342        key: ResolvedKey,
343        cand: Candidate,
344        outcome: Result<T, PoolError>,
345    ) -> Result<T, PoolError> {
346        match &outcome {
347            Ok(_) | Err(PoolError::Link(LinkError::Provider { .. })) => self.remember(key, cand),
348            Err(_) => self.forget(&key),
349        }
350        outcome
351    }
352
353    /// Calls the candidate's provider on `link`, under what is left before
354    /// `deadline`, and reports the call: a sealed call's result is only ever
355    /// one its answer opened under the key it was sealed to (the link
356    /// refuses any other), so `sealed` 1 names that key.
357    async fn call_at(
358        &self,
359        link: &Link,
360        cand: &Candidate,
361        c: &Call,
362        deadline: Instant,
363    ) -> Result<(Value, Report), PoolError> {
364        let left = deadline.saturating_duration_since(Instant::now());
365        let result = link
366            .call(station_link::Call {
367                realm: c.realm,
368                procedure: c.procedure.clone(),
369                target: cand.provider.node,
370                payload: c.payload.clone(),
371                timeout: left.max(Duration::from_millis(1)),
372                token: c.token.clone(),
373                proofs: c.proofs.clone(),
374                seal: Some(cand.seal()),
375            })
376            .await?;
377        Ok((result, Report::of(cand.provider.node, cand.kem_key_id())))
378    }
379
380    /// A sealed call refused sealed_refused after its provider's key rotated,
381    /// sealed again once, as macula's `resealed/7` does (E2E design §5.1,
382    /// Amendment A1; macula-rust#20): ONE fresh lookup of the provider's own
383    /// trusted advertisements, then a new request sealed to the one naming
384    /// exactly the key the refusal named, or, from a provider that named
385    /// none, to the first key it advertises. A second refusal is the result.
386    /// Any other key fails key_mismatch naming both, none fails no_kem_key;
387    /// never the clear. Any other outcome is returned as it is. The
388    /// candidate the outcome came from, and the outcome.
389    async fn resealed(
390        &self,
391        key: &ResolvedKey,
392        cand: Candidate,
393        sent: &Sent<'_>,
394        outcome: Result<(Value, Report), PoolError>,
395    ) -> (Candidate, Result<(Value, Report), PoolError>) {
396        let named = match &outcome {
397            Err(PoolError::Link(LinkError::SealedRefused { named })) if cand.kem_key.is_some() => {
398                *named
399            }
400            _ => return (cand, outcome),
401        };
402        let own = ResolvedKey {
403            provider: cand.provider.node,
404            ..key.clone()
405        };
406        let fresh = bounded(sent.deadline, self.resolve(&own, sent.realm_key.clone())).await;
407        let next = match reseal_to(fresh, named) {
408            Ok(next) => next,
409            Err(e) => return (cand, Err(e)),
410        };
411        let resent = self.call_at(sent.link, &next, sent.call, sent.deadline);
412        let outcome = bounded(sent.deadline, resent).await;
413        (next, outcome)
414    }
415
416    /// Opens the stream at the candidate's provider on `link`, holding
417    /// `reseal` for a sealed_refused of its open.
418    async fn open_at(
419        &self,
420        link: &Link,
421        cand: &Candidate,
422        c: &StreamCall,
423        reseal: Option<Reseal>,
424    ) -> Result<Stream, PoolError> {
425        Ok(link
426            .open_stream_resealing(
427                station_link::StreamCall {
428                    realm: c.realm,
429                    procedure: c.procedure.clone(),
430                    target: cand.provider.node,
431                    mode: c.mode,
432                    payload: c.payload.clone(),
433                    deadline: c.deadline,
434                    token: c.token.clone(),
435                    proofs: c.proofs.clone(),
436                    seal: Some(cand.seal()),
437                },
438                reseal,
439            )
440            .await?)
441    }
442
443    /// A sealed stream's one reseal, on a call's terms (see
444    /// [`PoolInner::resealed`]): one fresh lookup of the provider's own
445    /// trusted advertisements, then the stream reopened on `link` sealed
446    /// only to the key the refusal named (another fails naming both, none
447    /// fails closed), or, after a keyless refusal, to the first key the
448    /// provider now names. The reopened stream holds no reseal: a second
449    /// refusal ends it. None for a clear candidate.
450    fn stream_reseal(
451        self: &Arc<Self>,
452        key: &ResolvedKey,
453        cand: &Candidate,
454        c: &StreamCall,
455        realm_key: Option<Vec<u8>>,
456        link: &Link,
457    ) -> Option<Reseal> {
458        cand.kem_key.as_ref()?;
459        let pool = Arc::downgrade(self);
460        let r = Resealing {
461            own: ResolvedKey {
462                provider: cand.provider.node,
463                ..key.clone()
464            },
465            key: key.clone(),
466            c: c.clone(),
467            link: link.clone(),
468            realm_key,
469        };
470        Some(Box::new(move |named| {
471            Box::pin(resealed_stream(pool, r, named))
472        }))
473    }
474
475    /// The link up now to `station`, if any.
476    fn linked_to(&self, station: &[u8; 32]) -> Option<Link> {
477        self.links()
478            .into_iter()
479            .find(|l| l.station_node_id() == *station)
480    }
481
482    pub(super) async fn link_to(
483        self: &Arc<Self>,
484        station: &[u8; 32],
485        deadline: Instant,
486    ) -> Result<Link, PoolError> {
487        if let Some(link) = self.linked_to(station) {
488            return Ok(link);
489        }
490        let (existing, direct) = self.member_for(station)?;
491        let fresh = existing.is_none();
492        let member = match existing {
493            Some(m) => m,
494            None => self.start_direct(station, direct, deadline).await?,
495        };
496        if let Some(link) = member.await_up(deadline, fresh).await {
497            return Ok(link);
498        }
499        if fresh {
500            self.drop_member(&member);
501        }
502        Err(PoolError::StationNotReached {
503            station: *station,
504            cause: member.last_error(),
505        })
506    }
507
508    /// The member already linking to `station`, if any, and how many direct
509    /// links the pool holds; [`PoolError::Closed`] on a closed pool.
510    fn member_for(&self, station: &[u8; 32]) -> Result<(Option<Arc<Member>>, usize), PoolError> {
511        let state = self.lock();
512        if state.closed {
513            return Err(PoolError::Closed);
514        }
515        let existing = state
516            .members
517            .iter()
518            .find(|m| m.target.expected_node_id == *station)
519            .cloned();
520        Ok((existing, state.members.iter().filter(|m| m.direct).count()))
521    }
522
523    /// Starts a direct link to `station`, at the address its own endpoint
524    /// record names, when fewer than `max_direct_links` are held.
525    async fn start_direct(
526        self: &Arc<Self>,
527        station: &[u8; 32],
528        direct: usize,
529        deadline: Instant,
530    ) -> Result<Arc<Member>, PoolError> {
531        if direct >= self.opts.max_direct_links {
532            return Err(PoolError::DirectLinksFull);
533        }
534        let target = bounded(deadline, self.station_target(station)).await?;
535        Ok(self.start_member(target, true))
536    }
537
538    async fn station_target(&self, station: &[u8; 32]) -> Result<Target, PoolError> {
539        let slot = record::station_endpoint_key(station);
540        let verified = self
541            .first_answer(|l| async move { l.find_record(&slot).await })
542            .await
543            .map_err(|e| match e {
544                PoolError::Link(link) => PoolError::NoStationEndpoint(Some(link)),
545                other => other,
546            })?;
547        let r = verified.record();
548        let signer = r.signed.as_ref().map(|s| s.key_id);
549        let endpoint = record::read_station_endpoint(r)
550            .map_err(|e| PoolError::NoStationEndpoint(Some(e.into())))?;
551        match (signer, endpoint.host_advertised.first()) {
552            (Some(signer), Some(host)) if signer == *station && endpoint.quic_port != 0 => {
553                Ok(Target {
554                    host: host.clone(),
555                    port: endpoint.quic_port,
556                    profile: self.opts.identity.profile(),
557                    expected_node_id: *station,
558                })
559            }
560            _ => Err(PoolError::NoStationEndpoint(None)),
561        }
562    }
563
564    /// Runs `ask` on the links in selection order and returns the first
565    /// answer, moving on only when a link could not carry the request: a
566    /// station's own answer, not_found included, is final.
567    pub(super) async fn first_answer<'a, T, F, Fut>(&self, ask: F) -> Result<T, PoolError>
568    where
569        F: Fn(Link) -> Fut,
570        Fut: Future<Output = Result<T, LinkError>> + 'a,
571    {
572        let links = self.links();
573        if links.is_empty() {
574            return Err(PoolError::NoLink(Vec::new()));
575        }
576        let mut errors = Vec::new();
577        for link in links {
578            match ask(link).await {
579                Err(e) if unreachable(&e) => errors.push(e),
580                answered => return answered.map_err(PoolError::Link),
581            }
582        }
583        Err(PoolError::NoLink(errors))
584    }
585}
586
587/// A call that went out: the link it went on, the call, the realm key its
588/// candidates were trusted under, and its deadline.
589struct Sent<'a> {
590    link: &'a Link,
591    call: &'a Call,
592    realm_key: Option<Vec<u8>>,
593    deadline: Instant,
594}
595
596/// The advertisement a refused sealed call is sealed to again, from the
597/// provider's fresh advertisements (`found`): the one naming exactly the key
598/// the refusal `named`, searched for since the DHT may still serve an older
599/// one, or, when the refusal named none, the first that names a key. When
600/// the named key is not among them, key_mismatch naming both; when none
601/// names a key, no_kem_key.
602/// What a stream's reseal reopens: the provider's own key, the key the pool
603/// remembers it under, the open, and its link and realm key.
604struct Resealing {
605    own: ResolvedKey,
606    key: ResolvedKey,
607    c: StreamCall,
608    link: Link,
609    realm_key: Option<Vec<u8>>,
610}
611
612/// The stream reopened as [`PoolInner::stream_reseal`] says.
613async fn resealed_stream(
614    pool: Weak<PoolInner>,
615    r: Resealing,
616    named: Option<[u8; KEY_ID_SIZE]>,
617) -> Result<Stream, LinkError> {
618    let pool = pool.upgrade().ok_or(LinkError::Closed)?;
619    let deadline = Instant::now() + DEFAULT_CALL_TIMEOUT;
620    let fresh = bounded(deadline, pool.resolve(&r.own, r.realm_key)).await;
621    let reopened = match reseal_to(fresh, named) {
622        Ok(next) => bounded(deadline, pool.open_at(&r.link, &next, &r.c, None))
623            .await
624            .map(|stream| (stream, next)),
625        Err(e) => Err(e),
626    };
627    // As a call settles: the provider's new key is remembered, and a stale
628    // one is forgotten, so the next open looks it up afresh.
629    match reopened {
630        Ok((stream, next)) => {
631            pool.remember(r.key, next);
632            Ok(stream)
633        }
634        Err(e) => {
635            pool.forget(&r.key);
636            Err(as_link_error(e))
637        }
638    }
639}
640
641/// A pool error as the stream it ends carries it.
642fn as_link_error(e: PoolError) -> LinkError {
643    match e {
644        PoolError::Link(e) => e,
645        PoolError::Confidentiality(e) => LinkError::Confidentiality(e),
646        other => LinkError::Io(other.to_string()),
647    }
648}
649
650fn reseal_to(
651    found: Result<Vec<Candidate>, PoolError>,
652    named: Option<[u8; KEY_ID_SIZE]>,
653) -> Result<Candidate, PoolError> {
654    let keyed: Vec<Candidate> = found
655        .unwrap_or_default()
656        .into_iter()
657        .filter(|c| c.kem_key.is_some())
658        .collect();
659    let advertised: Vec<[u8; KEY_ID_SIZE]> =
660        keyed.iter().filter_map(Candidate::kem_key_id).collect();
661    let chosen = match named {
662        Some(id) => keyed.into_iter().find(|c| c.kem_key_id() == Some(id)),
663        None => keyed.into_iter().next(),
664    };
665    let reason = match advertised.is_empty() {
666        true => ConfidentialityReason::NoKemKey,
667        false => ConfidentialityReason::KeyMismatch,
668    };
669    chosen.ok_or(PoolError::Confidentiality(ConfidentialityError {
670        reason,
671        advertised,
672        named,
673    }))
674}
675
676/// The first of `candidates` whose serving station is reached within its
677/// share of `deadline`, and the link to it; nothing is sent. A candidate not
678/// reached is forgotten and named in [`PoolError::NoProvider`], and none is
679/// tried once `deadline` has passed.
680async fn first_reached(
681    inner: &Arc<PoolInner>,
682    key: &ResolvedKey,
683    candidates: Vec<Candidate>,
684    deadline: Instant,
685) -> Result<(Link, Candidate), PoolError> {
686    let mut tried = Vec::new();
687    let count = candidates.len();
688    for (i, cand) in candidates.into_iter().enumerate() {
689        match inner
690            .reach(&cand, candidate_share(deadline, count - i))
691            .await
692        {
693            Ok(link) => return Ok((link, cand)),
694            Err(PoolError::Closed) => return Err(PoolError::Closed),
695            Err(e) => {
696                inner.forget(key);
697                tried.push((cand.provider, e));
698            }
699        }
700        if Instant::now() >= deadline {
701            break;
702        }
703    }
704    Err(PoolError::NoProvider(tried))
705}
706
707/// The candidate an advertisement under `key`'s slot makes, when it is one
708/// of `key`'s procedure (by `key.provider` when it is set) and `trust`
709/// authorizes it at `now`.
710fn trusted_candidate(
711    v: &Verified,
712    key: &ResolvedKey,
713    trust: &Trust,
714    now: i64,
715) -> Option<Candidate> {
716    let ad = record::read_procedure_advertisement(v.record()).ok()?;
717    let wanted = ad.realm_id == key.realm
718        && ad.procedure == key.procedure
719        && (key.provider == [0; 32] || ad.advertiser_node == key.provider);
720    if !wanted || record::verify_authorization(v, trust, now).is_err() {
721        return None;
722    }
723    Some(Candidate {
724        provider: Provider {
725            node: ad.advertiser_node,
726            station: ad.serving_station,
727        },
728        expires_at: v.record().expires_at,
729        created_at: v.record().created_at,
730        kem_key: ad.kem_key.map(|(key, _)| key),
731    })
732}
733
734/// A failure of the link to carry a request, as opposed to the station's
735/// answer: no reply in time, or the link ended.
736fn unreachable(e: &LinkError) -> bool {
737    matches!(
738        e,
739        LinkError::CallTimeout
740            | LinkError::Closed
741            | LinkError::LivenessLost
742            | LinkError::V5DowngradeRefused
743            | LinkError::Io(_)
744            | LinkError::Goodbye(_)
745            | LinkError::StatusExpired
746            | LinkError::BindingExpired
747    )
748}
749
750/// One candidate's part of what is left before `deadline` with `left`
751/// candidates to try, at least a second, as macula shares it, and never past
752/// `deadline` itself.
753fn candidate_share(deadline: Instant, left: usize) -> Instant {
754    let now = Instant::now();
755    let remaining = deadline.saturating_duration_since(now);
756    (now + (remaining / left.max(1) as u32).max(MIN_CANDIDATE_SHARE)).min(deadline)
757}
758
759/// `work` bounded by `deadline`: past it, the call timed out.
760async fn bounded<T>(
761    deadline: Instant,
762    work: impl Future<Output = Result<T, PoolError>>,
763) -> Result<T, PoolError> {
764    tokio::time::timeout_at(deadline, work)
765        .await
766        .unwrap_or(Err(PoolError::Link(LinkError::CallTimeout)))
767}
768
769fn now_ms() -> i64 {
770    crate::uuid_v7::now_ms() as i64
771}
772
773/// The candidates a call or an open under `confidential` may reach, in their
774/// order, before anything is sent. A candidate whose advertisement names a
775/// KEM key is called sealed to it. Under `Preferred` a keyless one is called
776/// in the clear, unless its provider names a key in another advertisement
777/// the DHT still serves (it rotated onto `kem_advertise`): a provider that
778/// names a key is never called in the clear. Under `Required` only keyed
779/// candidates are called. None left is a [`ConfidentialityError`] naming the
780/// advertised key ids. `Off` is refused: a clear call is an explicit
781/// target's, made on a link with [`Seal::Clear`].
782fn callable(
783    candidates: Vec<Candidate>,
784    confidential: Confidentiality,
785) -> Result<Vec<Candidate>, PoolError> {
786    if confidential == Confidentiality::Off {
787        return Err(PoolError::InvalidOpts(
788            "confidential off is refused for a pool call or open: it is preferred or required"
789                .into(),
790        ));
791    }
792    let advertised: Vec<[u8; KEY_ID_SIZE]> = candidates
793        .iter()
794        .filter_map(Candidate::kem_key_id)
795        .collect();
796    let keyed: HashSet<[u8; 32]> = candidates
797        .iter()
798        .filter(|c| c.kem_key.is_some())
799        .map(|c| c.provider.node)
800        .collect();
801    let kept: Vec<Candidate> = candidates
802        .into_iter()
803        .filter(|c| match (confidential, &c.kem_key) {
804            (_, Some(_)) => true,
805            (Confidentiality::Preferred, None) => !keyed.contains(&c.provider.node),
806            (_, None) => false,
807        })
808        .collect();
809    if kept.is_empty() {
810        return Err(PoolError::Confidentiality(ConfidentialityError {
811            reason: ConfidentialityReason::NoKemKey,
812            advertised,
813            named: None,
814        }));
815    }
816    Ok(kept)
817}
818
819#[cfg(test)]
820mod tests {
821    use super::*;
822
823    fn cand(node: u8, kem_key: Option<u8>) -> Candidate {
824        Candidate {
825            provider: Provider {
826                node: [node; 32],
827                station: [9; 32],
828            },
829            expires_at: 0,
830            created_at: 0,
831            kem_key: kem_key.map(|k| vec![k; 1568]),
832        }
833    }
834
835    fn refused(r: Result<Vec<Candidate>, PoolError>) -> ConfidentialityError {
836        match r {
837            Err(PoolError::Confidentiality(e)) => e,
838            other => panic!("not a confidentiality refusal: {other:?}"),
839        }
840    }
841
842    fn key_id_of(k: u8) -> [u8; KEY_ID_SIZE] {
843        crate::seal::key_id(&[k; 1568])
844    }
845
846    #[test]
847    fn a_refused_call_is_sealed_again_only_to_the_key_named() {
848        let found = || Ok(vec![cand(1, None), cand(1, Some(7)), cand(1, Some(8))]);
849        let next = reseal_to(found(), Some(key_id_of(8))).unwrap();
850        assert_eq!(next, cand(1, Some(8)));
851        let e = refused(reseal_to(found(), Some(key_id_of(9))).map(|c| vec![c]));
852        assert_eq!(e.reason, ConfidentialityReason::KeyMismatch);
853        assert_eq!(e.named, Some(key_id_of(9)));
854        assert_eq!(e.advertised, vec![key_id_of(7), key_id_of(8)]);
855        // A provider that named no key is sealed to the first it advertises.
856        assert_eq!(reseal_to(found(), None).unwrap(), cand(1, Some(7)));
857        // Nothing keyed, or nothing found: no_kem_key, never the clear.
858        let e = refused(reseal_to(Ok(vec![cand(1, None)]), Some(key_id_of(8))).map(|c| vec![c]));
859        assert_eq!(e.reason, ConfidentialityReason::NoKemKey);
860        let e = refused(reseal_to(Err(PoolError::NoProvider(Vec::new())), None).map(|c| vec![c]));
861        assert_eq!(e.reason, ConfidentialityReason::NoKemKey);
862    }
863
864    #[test]
865    fn preferred_seals_to_a_keyed_provider_and_calls_a_keyless_one_in_the_clear() {
866        let kept = callable(
867            vec![cand(1, Some(7)), cand(2, None), cand(3, None)],
868            Confidentiality::Preferred,
869        )
870        .unwrap();
871        assert_eq!(kept, vec![cand(1, Some(7)), cand(2, None), cand(3, None)]);
872        assert_eq!(kept[0].seal(), Seal::To(vec![7; 1568]));
873        assert_eq!(kept[1].seal(), Seal::Clear);
874    }
875
876    #[test]
877    fn a_provider_that_names_a_key_is_not_called_through_its_older_keyless_ad() {
878        let kept = callable(
879            vec![cand(1, Some(7)), cand(1, None), cand(2, None)],
880            Confidentiality::Preferred,
881        )
882        .unwrap();
883        assert_eq!(kept, vec![cand(1, Some(7)), cand(2, None)]);
884    }
885
886    #[test]
887    fn required_calls_only_keyed_providers_and_refuses_when_there_are_none() {
888        let kept = callable(
889            vec![cand(1, None), cand(2, Some(7))],
890            Confidentiality::Required,
891        )
892        .unwrap();
893        assert_eq!(kept, vec![cand(2, Some(7))]);
894        let e = refused(callable(vec![cand(1, None)], Confidentiality::Required));
895        assert_eq!(e.reason, ConfidentialityReason::NoKemKey);
896        assert!(e.advertised.is_empty());
897        assert_eq!(e.to_string(), "confidentiality: no_kem_key");
898    }
899
900    #[test]
901    fn confidential_parses_and_off_is_refused_for_a_pool_call() {
902        assert_eq!(Confidentiality::default(), Confidentiality::Preferred);
903        for (text, parsed) in [
904            ("", Confidentiality::Preferred),
905            ("preferred", Confidentiality::Preferred),
906            ("required", Confidentiality::Required),
907            ("off", Confidentiality::Off),
908        ] {
909            assert_eq!(text.parse::<Confidentiality>(), Ok(parsed), "{text}");
910        }
911        for refused in ["Required", "none", "optional"] {
912            assert!(refused.parse::<Confidentiality>().is_err(), "{refused}");
913        }
914        assert!(matches!(
915            callable(vec![cand(1, Some(7))], Confidentiality::Off),
916            Err(PoolError::InvalidOpts(_))
917        ));
918    }
919}