Skip to main content

macula_rust/
pool.rs

1//! A macula 12 node's set of station links, as macula's client pool keeps
2//! them: one link to each seed station, every seed pinned by its node_id,
3//! all links of one node sharing its identity key, statement issuer, request
4//! admission, publication seq and event dedup. A link that ends is dialed
5//! again after the respawn delay and given back the node's subscriptions and
6//! served procedures.
7//!
8//! Calls reach providers directly, as macula 12 calls them: the procedure's
9//! advertisements are resolved from the DHT and checked against the realm key
10//! the pool pins for the realm, the serving station an advertisement names is
11//! dialed (pinned by its node_id, from its own station_endpoint record) and
12//! the provider called there. Station procedures (`_dht.*`) go to the pool's
13//! links.
14//!
15//! Station discovery beyond the seeds is not here: macula discovers stations
16//! through mcl-stations/list_stations, which this pool does not call.
17
18mod call;
19mod content;
20mod member;
21mod pubsub;
22mod serve;
23
24pub use crate::station_link::Confidentiality;
25pub use crate::station_link::{ConfidentialityError, ConfidentialityReason, Report, ReportError};
26pub use call::{Call, Provider, StreamCall};
27pub use content::{content_procedure_bound, ContentOptions, CONTENT_PROCEDURE};
28pub use pubsub::Subscription;
29pub use serve::{Offer, Served};
30
31use std::collections::HashMap;
32use std::fmt;
33use std::sync::{Arc, Mutex, MutexGuard};
34use std::time::Duration;
35
36use crate::node_key::{carried_key_well_formed, NodeKey, Purpose};
37use crate::seal::Keyring;
38use crate::statement_issuer::{IssuerError, StatementIssuer};
39use crate::station_link::{
40    Admission, AdmissionLimits, EventDedup, Link, LinkError, PublicationSeq,
41};
42use crate::transport::Target;
43
44use member::Member;
45
46/// macula's defaults and caps for a pool's bounds.
47pub const DEFAULT_REPLICATION_FACTOR: usize = 2;
48pub const DEFAULT_RESPAWN_DELAY: Duration = Duration::from_secs(1);
49pub const DEFAULT_MAX_SEEDS: usize = 16;
50pub const DEFAULT_MAX_DIRECT_LINKS: usize = 8;
51pub const DEFAULT_CONNECT_TIMEOUT: Duration = Duration::from_secs(30);
52const MAX_LINK_LIMIT: usize = 64;
53
54/// Why a pool, or one of its operations, failed.
55#[derive(Debug, Clone, PartialEq, Eq)]
56pub enum PoolError {
57    /// A pool given no seed station.
58    NoSeeds,
59    /// A seed without the station's node_id, as host:port: a pool dials
60    /// only stations it can check.
61    SeedNotPinned(String),
62    /// More seeds than `max_seeds`.
63    TooManySeeds { given: usize, max: usize },
64    /// A realm key that is not a well-formed key of the pool's profile, and
65    /// its realm.
66    RealmTrustInvalid([u8; 32]),
67    /// An option out of its range, or a key that is no identity key.
68    InvalidOpts(String),
69    /// No link came up within the connect timeout, or none is up to carry an
70    /// operation; each link's last error.
71    NoLink(Vec<LinkError>),
72    /// An operation on a closed pool.
73    Closed,
74    /// A realm the pool pins no key for: nothing in it is served or trusted.
75    NoRealmKey,
76    /// No provider the pinned realm key authorizes answered: each candidate
77    /// tried and why it failed; none when there was no candidate at all.
78    NoProvider(Vec<(Provider, PoolError)>),
79    /// A serving station with no endpoint record it signed itself.
80    NoStationEndpoint(Option<LinkError>),
81    /// A serving station not yet linked while `max_direct_links` direct
82    /// links are held.
83    DirectLinksFull,
84    /// A station that could not be linked, and the last dial's error.
85    StationNotReached {
86        station: [u8; 32],
87        cause: Option<LinkError>,
88    },
89    /// A procedure no link would serve, and each link's error.
90    NotServed(Vec<LinkError>),
91    /// A link's own failure, or a provider's or station's answer.
92    Link(LinkError),
93    /// Content no node announces in the realm, or a sharer that does not
94    /// hold it.
95    NotShared,
96    /// Content every announcing sharer failed to give: each sharer's node
97    /// and why.
98    ContentUnavailable(Vec<([u8; 32], PoolError)>),
99    /// A block, manifest or whole that does not match the content id it was
100    /// asked for by, and which.
101    ContentMismatch(String),
102    /// Content over the fetch's bounds, and how large.
103    ContentTooLarge(String),
104    /// An answer a sharer never gives, and what it was.
105    ContentReply(String),
106    /// A call or an open that could not be kept confidential, so nothing
107    /// was sent.
108    Confidentiality(ConfidentialityError),
109}
110
111impl fmt::Display for PoolError {
112    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
113        match self {
114            PoolError::Link(e) => write!(f, "{e}"),
115            PoolError::Confidentiality(e) => write!(f, "{e}"),
116            PoolError::NoProvider(tried) if tried.is_empty() => {
117                f.write_str("no trusted provider advertises the procedure")
118            }
119            PoolError::NoProvider(tried) => {
120                f.write_str("no trusted provider answered:")?;
121                for (p, e) in tried {
122                    write!(f, " [{} at {}: {e}]", short(&p.node), short(&p.station))?;
123                }
124                Ok(())
125            }
126            PoolError::ContentUnavailable(tried) => {
127                f.write_str("no sharer gave the content:")?;
128                for (node, e) in tried {
129                    write!(f, " [{}: {e}]", short(node))?;
130                }
131                Ok(())
132            }
133            other => write!(f, "{other:?}"),
134        }
135    }
136}
137
138impl std::error::Error for PoolError {}
139
140impl From<LinkError> for PoolError {
141    fn from(e: LinkError) -> Self {
142        PoolError::Link(e)
143    }
144}
145
146fn short(id: &[u8; 32]) -> String {
147    id[..4].iter().map(|b| format!("{b:02x}")).collect()
148}
149
150/// A station to link to: where it is dialed and the node_id it must prove.
151#[derive(Debug, Clone, PartialEq, Eq)]
152pub struct Seed {
153    pub host: String,
154    pub port: u16,
155    pub node_id: [u8; 32],
156}
157
158/// The order calls, publications and DHT operations try the pool's links in.
159#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
160pub enum LinkSelection {
161    /// The links in seed order.
162    #[default]
163    FirstSuccess,
164    /// A fresh random order each time.
165    Random,
166}
167
168/// A link coming up, or ending or failing to dial with its error.
169#[derive(Debug, Clone, PartialEq, Eq)]
170pub struct LinkEvent {
171    pub station: [u8; 32],
172    pub direct: bool,
173    pub up: bool,
174    pub error: Option<LinkError>,
175}
176
177/// A pool's configuration. [`Opts::new`] gives macula's defaults; a zero
178/// bound also means its default.
179#[derive(Clone)]
180pub struct Opts {
181    /// The node's identity key; its profile is the pool's.
182    pub identity: Arc<NodeKey>,
183    /// Each realm's key as carried: an advertisement in a realm is trusted
184    /// only when its authorization verifies against it, and an org
185    /// procedure is served only in a realm it names.
186    pub realm_trust: HashMap<[u8; 32], Vec<u8>>,
187    pub replication_factor: usize,
188    pub respawn_delay: Duration,
189    pub max_seeds: usize,
190    pub max_direct_links: usize,
191    /// How long [`Pool::connect`] waits for a first link.
192    pub connect_timeout: Duration,
193    /// Bounds on the requests the node's served procedures take; `None` is
194    /// macula's defaults, with the cap one share per link the pool may hold.
195    pub admission: Option<AdmissionLimits>,
196    pub link_selection: LinkSelection,
197    /// Gives the node a KEM keyring and names its current key in the
198    /// advertisements of procedures served confidentially (macula 13, E2E
199    /// design amendment A1), so callers seal to it. Off by default: switch it
200    /// on only once every station runs macula 12.11 or later and every
201    /// caller runs 13. Without it the node opens no sealed request.
202    pub kem_advertise: bool,
203    /// Hears every link coming up and going down, on a task of its own.
204    pub on_link_event: Option<Arc<dyn Fn(LinkEvent) + Send + Sync>>,
205    /// Hears each failure to reissue the node's status statements or rotate
206    /// its CONNECT key; `None` writes it to stderr. Left failing, the links
207    /// end when their statements lapse.
208    pub on_issuer_error: Option<Arc<dyn Fn(IssuerError) + Send + Sync>>,
209}
210
211impl Opts {
212    /// macula's defaults for `identity`, trusting no realm.
213    pub fn new(identity: Arc<NodeKey>) -> Opts {
214        Opts {
215            identity,
216            realm_trust: HashMap::new(),
217            replication_factor: DEFAULT_REPLICATION_FACTOR,
218            respawn_delay: DEFAULT_RESPAWN_DELAY,
219            max_seeds: DEFAULT_MAX_SEEDS,
220            max_direct_links: DEFAULT_MAX_DIRECT_LINKS,
221            connect_timeout: DEFAULT_CONNECT_TIMEOUT,
222            admission: None,
223            link_selection: LinkSelection::FirstSuccess,
224            kem_advertise: false,
225            on_link_event: None,
226            on_issuer_error: None,
227        }
228    }
229}
230
231/// A node's station links. Cloning it shares the pool; the pool closes when
232/// [`Pool::close`] is called or its last handle is dropped.
233#[derive(Clone)]
234pub struct Pool {
235    inner: Arc<PoolInner>,
236}
237
238pub(crate) struct PoolInner {
239    opts: Opts,
240    self_id: [u8; 32],
241    issuer: StatementIssuer,
242    publication_seq: Arc<PublicationSeq>,
243    admission: Arc<Admission>,
244    dedup: Arc<EventDedup>,
245    /// One node identity, one keyring, shared by every link.
246    keyring: Option<Arc<Keyring>>,
247    state: Mutex<State>,
248    ticks: tokio::task::JoinHandle<()>,
249    content: content::Sharer,
250}
251
252struct State {
253    members: Vec<Arc<Member>>,
254    subs: HashMap<u64, Arc<pubsub::SubInner>>,
255    served: HashMap<u64, Arc<serve::ServedInner>>,
256    remember: HashMap<call::ResolvedKey, call::Candidate>,
257    closed: bool,
258}
259
260/// One of the pool's links.
261#[derive(Debug, Clone, PartialEq, Eq)]
262pub struct LinkStatus {
263    pub station: [u8; 32],
264    pub host: String,
265    pub port: u16,
266    pub direct: bool,
267    pub up: bool,
268}
269
270impl Pool {
271    /// Checks `seeds` and `opts`, dials every seed, and returns once one link
272    /// is up, or [`PoolError::NoLink`] with each link's last error when the
273    /// connect timeout passes first. Links not yet up keep dialing.
274    pub async fn connect(seeds: Vec<Seed>, opts: Opts) -> Result<Pool, PoolError> {
275        let opts = checked(&seeds, opts)?;
276        let self_id = opts
277            .identity
278            .node_id()
279            .map_err(|e| PoolError::InvalidOpts(e.to_string()))?;
280        let issuer = StatementIssuer::with_wall_clock(opts.identity.clone())
281            .map_err(|e| PoolError::InvalidOpts(e.to_string()))?;
282        let on_error = opts.on_issuer_error.clone();
283        let ticks = issuer.spawn_ticks(move |e| match &on_error {
284            Some(f) => f(e),
285            None => eprintln!("macula-rust pool: the statement issuer failed: {e}"),
286        });
287        let admission = opts.admission.expect("checked fills the admission limits");
288        let keyring = match opts.kem_advertise {
289            true => Some(Arc::new(
290                Keyring::system(opts.identity.profile())
291                    .map_err(|e| PoolError::InvalidOpts(e.to_string()))?,
292            )),
293            false => None,
294        };
295        let inner = Arc::new(PoolInner {
296            self_id,
297            issuer,
298            publication_seq: Arc::default(),
299            admission: Arc::new(Admission::new(admission)),
300            dedup: Arc::default(),
301            keyring,
302            state: Mutex::new(State {
303                members: Vec::new(),
304                subs: HashMap::new(),
305                served: HashMap::new(),
306                remember: HashMap::new(),
307                closed: false,
308            }),
309            ticks,
310            content: content::Sharer::default(),
311            opts,
312        });
313        let pool = Pool { inner };
314        for seed in &seeds {
315            pool.inner.start_member(pool.target(seed), false);
316        }
317        let deadline = tokio::time::Instant::now() + pool.inner.opts.connect_timeout;
318        if let Err(e) = pool.inner.await_up(deadline).await {
319            pool.close().await;
320            return Err(e);
321        }
322        Ok(pool)
323    }
324
325    fn target(&self, seed: &Seed) -> Target {
326        Target {
327            host: seed.host.clone(),
328            port: seed.port,
329            profile: self.inner.opts.identity.profile(),
330            expected_node_id: seed.node_id,
331        }
332    }
333
334    /// The node_id the pool links as.
335    pub fn node_id(&self) -> [u8; 32] {
336        self.inner.self_id
337    }
338
339    /// Every link the pool holds, seeds first.
340    pub fn status(&self) -> Vec<LinkStatus> {
341        let members = self.inner.lock().members.clone();
342        members
343            .iter()
344            .map(|m| LinkStatus {
345                station: m.target.expected_node_id,
346                host: m.target.host.clone(),
347                port: m.target.port,
348                direct: m.direct,
349                up: m.current().is_some(),
350            })
351            .collect()
352    }
353
354    /// Ends every link with a GOODBYE and every subscription. It withdraws
355    /// nothing: an advertisement lapses with its link.
356    pub async fn close(&self) {
357        let (members, subs) = {
358            let mut state = self.inner.lock();
359            if state.closed {
360                return;
361            }
362            state.closed = true;
363            state.served.clear();
364            (
365                std::mem::take(&mut state.members),
366                std::mem::take(&mut state.subs),
367            )
368        };
369        self.inner.ticks.abort();
370        for m in &members {
371            m.retire();
372        }
373        for m in &members {
374            m.stopped().await;
375        }
376        for sub in subs.into_values() {
377            let _ = sub.end().await;
378        }
379    }
380}
381
382impl fmt::Debug for Pool {
383    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
384        f.debug_struct("Pool")
385            .field("node_id", &short(&self.inner.self_id))
386            .field("links", &self.status())
387            .finish()
388    }
389}
390
391impl PoolInner {
392    fn lock(&self) -> MutexGuard<'_, State> {
393        self.state.lock().unwrap_or_else(|p| p.into_inner())
394    }
395
396    /// The links up now, in the pool's selection order.
397    fn links(&self) -> Vec<Link> {
398        let members = self.lock().members.clone();
399        let mut up: Vec<Link> = members.iter().filter_map(|m| m.current()).collect();
400        if self.opts.link_selection == LinkSelection::Random {
401            shuffle(&mut up);
402        }
403        up
404    }
405
406    /// Waits until a link is up, or `deadline` passes.
407    async fn await_up(self: &Arc<Self>, deadline: tokio::time::Instant) -> Result<(), PoolError> {
408        loop {
409            if !self.links().is_empty() {
410                return Ok(());
411            }
412            if tokio::time::Instant::now() >= deadline {
413                let members = self.lock().members.clone();
414                return Err(PoolError::NoLink(
415                    members.iter().filter_map(|m| m.last_error()).collect(),
416                ));
417            }
418            tokio::time::sleep(Duration::from_millis(10)).await;
419        }
420    }
421
422    fn event(&self, e: LinkEvent) {
423        if let Some(f) = self.opts.on_link_event.clone() {
424            tokio::spawn(async move { f(e) });
425        }
426    }
427
428    /// The key a procedure's authorization is checked against: the realm's
429    /// pinned key, or none for a procedure in a node's own namespace, which
430    /// its advertisement's signature alone authorizes.
431    fn realm_key_for(
432        &self,
433        realm: &[u8; 32],
434        procedure: &str,
435    ) -> Result<Option<Vec<u8>>, PoolError> {
436        if crate::record::in_own_namespace(procedure) {
437            return Ok(None);
438        }
439        self.opts
440            .realm_trust
441            .get(realm)
442            .cloned()
443            .map(Some)
444            .ok_or(PoolError::NoRealmKey)
445    }
446}
447
448impl Drop for PoolInner {
449    /// A pool whose last handle is dropped stops dialing and closes its links.
450    fn drop(&mut self) {
451        self.ticks.abort();
452        let state = self.state.get_mut().unwrap_or_else(|p| p.into_inner());
453        for m in &state.members {
454            m.retire();
455        }
456    }
457}
458
459/// Refuses, before anything is dialed, what macula's pool refuses, and fills
460/// in the defaults.
461fn checked(seeds: &[Seed], mut opts: Opts) -> Result<Opts, PoolError> {
462    if opts.identity.purpose() != Purpose::Identity {
463        return Err(PoolError::InvalidOpts("an identity key is required".into()));
464    }
465    let profile = opts.identity.profile();
466    for (realm, key) in &opts.realm_trust {
467        if !carried_key_well_formed(key, profile) {
468            return Err(PoolError::RealmTrustInvalid(*realm));
469        }
470    }
471    for (name, limit) in [
472        ("max_seeds", opts.max_seeds),
473        ("max_direct_links", opts.max_direct_links),
474        ("replication_factor", opts.replication_factor),
475    ] {
476        if limit > MAX_LINK_LIMIT {
477            return Err(PoolError::InvalidOpts(format!(
478                "{name} of {limit}, outside 1 to {MAX_LINK_LIMIT}"
479            )));
480        }
481    }
482    let or_default = |v: usize, d: usize| if v == 0 { d } else { v };
483    opts.max_seeds = or_default(opts.max_seeds, DEFAULT_MAX_SEEDS);
484    opts.max_direct_links = or_default(opts.max_direct_links, DEFAULT_MAX_DIRECT_LINKS);
485    opts.replication_factor = or_default(opts.replication_factor, DEFAULT_REPLICATION_FACTOR);
486    if opts.respawn_delay.is_zero() {
487        opts.respawn_delay = DEFAULT_RESPAWN_DELAY;
488    }
489    if opts.connect_timeout.is_zero() {
490        opts.connect_timeout = DEFAULT_CONNECT_TIMEOUT;
491    }
492    let admission = opts.admission.unwrap_or_else(|| {
493        let mut limits = AdmissionLimits::default();
494        limits.cap = limits.share * (opts.max_seeds + opts.max_direct_links);
495        limits
496    });
497    admission
498        .validate()
499        .map_err(|e| PoolError::InvalidOpts(e.to_string()))?;
500    opts.admission = Some(admission);
501    if seeds.is_empty() {
502        return Err(PoolError::NoSeeds);
503    }
504    if seeds.len() > opts.max_seeds {
505        return Err(PoolError::TooManySeeds {
506            given: seeds.len(),
507            max: opts.max_seeds,
508        });
509    }
510    if let Some(unpinned) = seeds.iter().find(|s| s.node_id == [0; 32]) {
511        return Err(PoolError::SeedNotPinned(format!(
512            "{}:{}",
513            unpinned.host, unpinned.port
514        )));
515    }
516    Ok(opts)
517}
518
519/// Fisher-Yates over the system's randomness.
520fn shuffle<T>(items: &mut [T]) {
521    for i in (1..items.len()).rev() {
522        let mut r = [0u8; 8];
523        if aws_lc_rs::rand::fill(&mut r).is_err() {
524            return;
525        }
526        let j = (u64::from_le_bytes(r) % (i as u64 + 1)) as usize;
527        items.swap(i, j);
528    }
529}