1mod 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
46pub 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#[derive(Debug, Clone, PartialEq, Eq)]
56pub enum PoolError {
57 NoSeeds,
59 SeedNotPinned(String),
62 TooManySeeds { given: usize, max: usize },
64 RealmTrustInvalid([u8; 32]),
67 InvalidOpts(String),
69 NoLink(Vec<LinkError>),
72 Closed,
74 NoRealmKey,
76 NoProvider(Vec<(Provider, PoolError)>),
79 NoStationEndpoint(Option<LinkError>),
81 DirectLinksFull,
84 StationNotReached {
86 station: [u8; 32],
87 cause: Option<LinkError>,
88 },
89 NotServed(Vec<LinkError>),
91 Link(LinkError),
93 NotShared,
96 ContentUnavailable(Vec<([u8; 32], PoolError)>),
99 ContentMismatch(String),
102 ContentTooLarge(String),
104 ContentReply(String),
106 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#[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
160pub enum LinkSelection {
161 #[default]
163 FirstSuccess,
164 Random,
166}
167
168#[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#[derive(Clone)]
180pub struct Opts {
181 pub identity: Arc<NodeKey>,
183 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 pub connect_timeout: Duration,
193 pub admission: Option<AdmissionLimits>,
196 pub link_selection: LinkSelection,
197 pub kem_advertise: bool,
203 pub on_link_event: Option<Arc<dyn Fn(LinkEvent) + Send + Sync>>,
205 pub on_issuer_error: Option<Arc<dyn Fn(IssuerError) + Send + Sync>>,
209}
210
211impl Opts {
212 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#[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 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#[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 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 pub fn node_id(&self) -> [u8; 32] {
336 self.inner.self_id
337 }
338
339 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 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 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 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 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 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
459fn 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
519fn 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}