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