macula_rust/pool.rs
1//! A multi-station connection pool: dials several [`Session`]s concurrently
2//! (bootstrap seeds, optionally grown by discovering more via
3//! `hecate_stations.list_stations`) and gives [`Pool::call`]/
4//! [`Pool::publish`] a choice of which connected one to use, instead of a
5//! caller managing a single [`Session`] by hand.
6//!
7//! **Call/Publish only — no pooled Subscribe.** [`Session::run_subscriber`]'s
8//! own doc (and [`Session::call`]'s) already say a control stream can't
9//! safely serve an in-flight Call's response-wait and an ongoing Subscribe
10//! EVENT loop at once — each discards frames it doesn't recognize, so they'd
11//! steal each other's frames. Building a pool-wide Subscribe fan-out
12//! properly would need either a second `Session` per link dedicated to it,
13//! or a real frame-demultiplexing layer on top of `Session` (dispatch by
14//! `call_id`/frame-type to whichever waiter wants it) — genuinely new
15//! infrastructure, out of scope here. A caller that needs Subscribe still
16//! uses a bare [`Session::subscribe`]/[`Session::run_subscriber`] directly,
17//! unpooled, exactly as before this module existed.
18//!
19//! Ported from the same station-discovery/link-rotation design already
20//! shipped in `macula-go` (`pool/discovery.go`, v0.7.0) and `macula-dotnet`
21//! (`StationDiscovery.cs`, v0.4.0) — see those crates' own doc comments for
22//! the fuller cross-language history. This is a from-scratch build, not a
23//! port of an existing pool: this crate had no `Seed`/multi-link concept at
24//! all before this module.
25//!
26//! **Per-link trust, not one fixed mode for the whole pool** — this is the
27//! one deliberate design difference from the go/dotnet ports, made possible
28//! by building from scratch rather than extending an existing single-Trust
29//! pool: a discovered station whose directory row has NO `hostname` (only a
30//! bare-IP `host_advertised`) but DOES carry a `node_id` dials under
31//! `Trust::Pinned(node_id)` instead of being skipped outright the way
32//! go/dotnet's ports skip every hostname-less row under `Trust::WebPki`
33//! (which can never validate a bare IP with no IP SANs). The underlying
34//! MECHANISM mirrors a shipped, live-verified precedent in a completely
35//! different codebase: `macula-apps/macula-cam2me`'s Android client
36//! (`reachability/StationDiscovery.kt`), confirmed against the real fleet
37//! by 34 (`macula`'s own reference-implementation session) independently
38//! reaching the identical conclusion this module reaches, including a
39//! live TLS-layer `verify=none` warning from `macula_quic` when dialing a
40//! hostname'd station by IP+Pinned instead of its usual WebPki path.
41//!
42//! **This module's PRIORITY ORDER deliberately differs from cam2me's own,
43//! in the safer direction** — cam2me picks Pinned(node_id) whenever a row
44//! carries a node_id at all (true of essentially every row), falling back
45//! to hostname/WebPki only when `host_advertised` itself is missing; that
46//! trades away TLS-layer MITM resistance for the common case, not just the
47//! no-DNS one. Here, [`dial_target_from_station_row`] prefers `hostname`
48//! unconditionally — `Trust::Pinned` is chosen only when a row has NO
49//! usable hostname at all, so a normal Let's-Encrypt-backed station still
50//! dials WebPki exactly as it always has; only a genuine no-DNS station
51//! (`stations-linode-toronto` — see `tests/live_station.rs`'s own
52//! `pinned_trust_full_handshake_succeeds_against_toronto`) ever falls to
53//! Pinned. Bootstrap seeds are entirely unaffected either way — they
54//! always dial under the pool's own configured `Trust`.
55
56use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
57use std::sync::Arc;
58use std::time::Duration;
59
60use tokio::sync::{watch, Mutex, RwLock};
61use tokio::task::JoinSet;
62
63use crate::connection::{self, CallError, Session};
64use crate::dht;
65use crate::frame::{CallResponse, PublishSpec};
66use crate::identity::KeyPair;
67use crate::transport::Trust;
68
69/// A dial target: host+port only, no identity attached (every link in a
70/// pool shares the pool's one identity). Mirrors macula-go's
71/// `connection.Seed` / macula-dotnet's `Seed` record — this crate had no
72/// equivalent type before this pool.
73#[derive(Debug, Clone, PartialEq, Eq, Hash)]
74pub struct Seed {
75 pub host: String,
76 pub port: u16,
77}
78
79impl Seed {
80 pub fn new(host: impl Into<String>, port: u16) -> Self {
81 Self {
82 host: host.into(),
83 port,
84 }
85 }
86}
87
88/// How [`Pool::call`]/[`Pool::publish`] order the pool's currently-connected
89/// links before applying their own existing first-match/`replication_factor`
90/// logic — changes ORDER only, never how many links get used. Matches
91/// macula_client.erl's own `link_selection` option (`first_success`/
92/// `random`) and macula-go/macula-dotnet's identically-named enum, so
93/// config ported from any of those doesn't need re-learning.
94#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
95pub enum LinkSelection {
96 /// (The default.) Derives the actual policy from
97 /// [`StationDiscoveryOptions::enabled`]: [`FirstSuccess`](Self::FirstSuccess)
98 /// if discovery is off (this pool's original, only behavior, unchanged),
99 /// [`Random`](Self::Random) if it's on.
100 #[default]
101 Auto,
102 /// Tries links in Seed-list order (bootstrap seeds in the order the
103 /// caller gave them, then discovered links in discovery order) — this
104 /// pool's baseline behavior, since `links` is a plain append-only
105 /// `Vec` (see [`Pool`]'s own doc on why that, not a `HashMap`, is the
106 /// backing store — insertion order is the ordering, nothing extra to
107 /// track).
108 FirstSuccess,
109 /// Uniformly shuffles the connected-links list before the same
110 /// first-match (Call) or take-first-N (Publish) logic runs. Composes
111 /// safely with a small `replication_factor`: shuffling ahead of a
112 /// 1-element slice is a no-op.
113 Random,
114}
115
116fn resolve_link_selection(
117 configured: LinkSelection,
118 station_discovery_enabled: bool,
119) -> LinkSelection {
120 match configured {
121 LinkSelection::Auto if station_discovery_enabled => LinkSelection::Random,
122 LinkSelection::Auto => LinkSelection::FirstSuccess,
123 other => other,
124 }
125}
126
127fn select_links(
128 mut connected: Vec<Arc<PooledLink>>,
129 resolved: LinkSelection,
130) -> Vec<Arc<PooledLink>> {
131 if resolved != LinkSelection::Random || connected.len() <= 1 {
132 return connected;
133 }
134 use rand::seq::SliceRandom;
135 connected.shuffle(&mut rand::rng());
136 connected
137}
138
139/// Configures opt-in discovery of additional stations via
140/// `hecate_stations.list_stations`, layered on top of the caller-supplied
141/// bootstrap [`Seed`]s. Default (`enabled == false`) is a complete no-op.
142///
143/// Bootstrap seeds keep their exact meaning: dialed first, permanent
144/// fallback if discovery never succeeds, retried forever on failure, never
145/// replaced. Discovery only ADDS links — a station missing from a later
146/// refresh does NOT tear down an existing link. A DISCOVERY-added link that
147/// fails to dial [`DISCOVERY_LINK_MAX_RESPAWN_ATTEMPTS`] times in a row
148/// gives up and frees its slot (so a future refresh can try a different
149/// station instead) — unlike a bootstrap seed, which never gives up. Go's
150/// and dotnet's ports of this same feature have no such give-up mechanism
151/// (a permanently-unreachable discovered station wastes a slot forever,
152/// redialing at the flat respawn delay indefinitely); building this pool
153/// from scratch, with 34's parallel design work on the Erlang reference
154/// specifically adding this exception, was reason enough to include it here
155/// too rather than carry the narrower behavior forward by default.
156#[derive(Debug, Clone)]
157pub struct StationDiscoveryOptions {
158 pub enabled: bool,
159 /// Interval between discovery attempts once at least one bootstrap
160 /// link is up. Default 30 minutes.
161 pub refresh_interval: Duration,
162 /// Bounds discovery's OWN adds only, not the pool's total link count —
163 /// see macula-dotnet's identical `StationDiscoveryOptions.MaxLinks` doc
164 /// for the exact accounting rules this mirrors.
165 pub max_links: usize,
166}
167
168impl Default for StationDiscoveryOptions {
169 fn default() -> Self {
170 Self {
171 enabled: false,
172 refresh_interval: Duration::from_secs(30 * 60),
173 max_links: 5,
174 }
175 }
176}
177
178/// Tunables for [`Pool`]. Defaults match every other port of this feature.
179#[derive(Debug, Clone)]
180pub struct PoolOptions {
181 pub link_selection: LinkSelection,
182 pub station_discovery: StationDiscoveryOptions,
183 /// Flat delay before redialing a link after it dies (or fails to dial
184 /// in the first place). Default 1s — flat, not exponential, matching
185 /// the reference and every other port in this SDK family except ts.
186 pub respawn_delay: Duration,
187 /// Per-CALL timeout, passed through to [`Session::call`]. Default
188 /// matches [`connection::DEFAULT_CALL_TIMEOUT`].
189 pub call_timeout: Duration,
190 /// How many currently-connected links a single publish fans out to.
191 /// Partial success counts as success. Default 1, matching macula-ts
192 /// and macula-dotnet's own `ReplicationFactor` default.
193 pub replication_factor: usize,
194}
195
196impl Default for PoolOptions {
197 fn default() -> Self {
198 Self {
199 link_selection: LinkSelection::default(),
200 station_discovery: StationDiscoveryOptions::default(),
201 respawn_delay: DEFAULT_RESPAWN_DELAY,
202 call_timeout: connection::DEFAULT_CALL_TIMEOUT,
203 replication_factor: 1,
204 }
205 }
206}
207
208pub const DEFAULT_RESPAWN_DELAY: Duration = Duration::from_secs(1);
209
210/// Discovery-added links only (never bootstrap seeds) give up redialing
211/// after this many consecutive failures — see
212/// [`StationDiscoveryOptions`]'s own doc for why.
213const DISCOVERY_LINK_MAX_RESPAWN_ATTEMPTS: u32 = 5;
214
215#[derive(Debug, Clone, Copy, PartialEq, Eq)]
216enum LinkOrigin {
217 Bootstrap,
218 Discovered,
219}
220
221/// A bootstrap link never gives up — matches every other port of this
222/// pool shape, and macula_client.erl's own reference behavior for
223/// caller-supplied seeds. Only a discovery-added link, after
224/// [`DISCOVERY_LINK_MAX_RESPAWN_ATTEMPTS`] straight failures, gives up —
225/// see [`StationDiscoveryOptions`]'s own doc for why this exception
226/// exists at all.
227fn should_give_up(origin: LinkOrigin, consecutive_failures: u32) -> bool {
228 origin == LinkOrigin::Discovered && consecutive_failures >= DISCOVERY_LINK_MAX_RESPAWN_ATTEMPTS
229}
230
231struct LinkState {
232 session: Option<Session>,
233 peer_node_id: Option<[u8; 32]>,
234}
235
236/// One configured link's current state. Never removed from [`Pool`]'s own
237/// `links` list once added (see that field's own doc) — `connected` simply
238/// goes false when the link is down or still dialing.
239pub struct PooledLink {
240 pub seed: Seed,
241 origin: LinkOrigin,
242 /// This link's OWN trust — usually the pool's configured `Trust`, but
243 /// see [`Pool`]'s module-level doc for the one case (a discovered,
244 /// hostname-less, node_id-bearing row) where it differs per link.
245 trust: Trust,
246 /// **Correctness depends on `tokio::sync::Mutex`'s documented FIFO
247 /// ordering** (queued lockers acquire in the exact order they queued)
248 /// — verified explicitly by adversarial review, 2026-09-05, tracing
249 /// the concurrent-`Pool::call`-vs-`Pool::call` race `redialing`
250 /// guards against: two callers racing the same dead session both
251 /// hold this lock for their own full network round-trip before
252 /// either calls `mark_disconnected`, and FIFO ordering is what
253 /// guarantees a STALE caller's own `mark_disconnected` can never run
254 /// after — and clobber — a session a freshly-spawned redial task
255 /// installed in the meantime (the redial task can only begin trying
256 /// to acquire this same lock once a real dial completes, which is
257 /// strictly later than any caller that was already queued before it
258 /// started). If this field's type ever changes to a non-fair lock,
259 /// that invariant needs re-deriving from scratch, not assumed to
260 /// still hold.
261 state: Mutex<LinkState>,
262 connected: AtomicBool,
263 /// Discovery-added links only: consecutive failed dial attempts, reset
264 /// on a successful connect. Bootstrap links never read this — they
265 /// retry forever regardless.
266 consecutive_failures: AtomicU32,
267 gave_up: AtomicBool,
268 /// Guards against two concurrent lifecycle tasks racing for the SAME
269 /// link — found by adversarial review, 2026-09-05: `Pool::call`/
270 /// `Pool::publish` can each independently observe the same dead
271 /// session (both hold+release `state`'s lock separately, one after
272 /// the other, before either calls `mark_disconnected`), so both can
273 /// call `mark_disconnected` for one link. Without this guard, both
274 /// would spawn their own respawn task, causing two concurrent dials
275 /// for one seed — whichever succeeds second silently drops (not
276 /// closes) the first's live `Session`, and a `Discovered` link's
277 /// `consecutive_failures` advances roughly 2x per real redial
278 /// interval, making [`should_give_up`] trigger far sooner than its
279 /// documented threshold implies. Only the task that wins the
280 /// false->true transition (see [`try_claim_redial`]) actually spawns;
281 /// the loser is a harmless no-op. Reset back to `false` exactly once,
282 /// when that task's own lifecycle future finishes (see
283 /// [`link_lifecycle_future`]).
284 redialing: AtomicBool,
285}
286
287impl PooledLink {
288 fn new(seed: Seed, origin: LinkOrigin, trust: Trust) -> Arc<Self> {
289 Arc::new(Self {
290 seed,
291 origin,
292 trust,
293 state: Mutex::new(LinkState {
294 session: None,
295 peer_node_id: None,
296 }),
297 connected: AtomicBool::new(false),
298 consecutive_failures: AtomicU32::new(0),
299 gave_up: AtomicBool::new(false),
300 redialing: AtomicBool::new(false),
301 })
302 }
303
304 pub fn is_connected(&self) -> bool {
305 self.connected.load(Ordering::Acquire)
306 }
307
308 /// A discovery-added link that has permanently given up redialing — see
309 /// [`StationDiscoveryOptions`]'s own doc. Always `false` for a
310 /// bootstrap link, which never gives up.
311 fn has_given_up(&self) -> bool {
312 self.gave_up.load(Ordering::Acquire)
313 }
314
315 async fn peer_node_id(&self) -> Option<[u8; 32]> {
316 self.state.lock().await.peer_node_id
317 }
318}
319
320/// Snapshot of one link, for health/introspection.
321#[derive(Debug, Clone)]
322pub struct LinkInfo {
323 pub seed: Seed,
324 pub connected: bool,
325 pub node_id: Option<[u8; 32]>,
326}
327
328/// Aggregate health snapshot. Lock-free best-effort read.
329#[derive(Debug, Clone, Copy)]
330pub struct PoolStatus {
331 pub healthy_links: usize,
332 pub total_links: usize,
333}
334
335impl PoolStatus {
336 /// At least one link has completed its CONNECT/HELLO handshake.
337 pub fn is_healthy(&self) -> bool {
338 self.healthy_links > 0
339 }
340}
341
342#[derive(Debug)]
343pub enum PoolCallError {
344 /// No link in the pool has completed its CONNECT/HELLO handshake.
345 NoHealthyStation,
346 /// Every currently-connected link's `call` failed — carries the LAST
347 /// one's error.
348 AllFailed(CallError),
349}
350
351impl std::fmt::Display for PoolCallError {
352 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
353 match self {
354 PoolCallError::NoHealthyStation => {
355 write!(f, "pool: no link has completed its CONNECT/HELLO handshake")
356 }
357 PoolCallError::AllFailed(e) => {
358 write!(f, "pool: every connected link's call failed: {e}")
359 }
360 }
361 }
362}
363
364impl std::error::Error for PoolCallError {}
365
366#[derive(Debug)]
367pub enum PoolPublishError {
368 NoHealthyStation,
369 /// Every link the publish was routed to failed — carries the count.
370 AllFailed(usize),
371}
372
373impl std::fmt::Display for PoolPublishError {
374 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
375 match self {
376 PoolPublishError::NoHealthyStation => {
377 write!(f, "pool: no link has completed its CONNECT/HELLO handshake")
378 }
379 PoolPublishError::AllFailed(n) => {
380 write!(f, "pool: publish failed on all {n} targeted link(s)")
381 }
382 }
383 }
384}
385
386impl std::error::Error for PoolPublishError {}
387
388/// A multi-station connection pool — see this module's own doc for the
389/// full design.
390///
391/// `links` is a plain append-only `Vec<Arc<PooledLink>>` behind an async
392/// `RwLock`, deliberately NOT a `HashMap`/`HashSet` — this is the one
393/// change made in direct response to the SAME class of bug found (twice,
394/// independently) porting this feature to macula-go (`map[string]*link`
395/// randomizing iteration order) and macula-dotnet (migrating to
396/// `ConcurrentDictionary` for concurrent-add safety silently broke
397/// `FirstSuccess`'s reliance on insertion order). A `Vec`, appended to
398/// under the same lock that's ALSO taken to read it, has no separate
399/// enumeration-order concept to accidentally break — insertion order IS
400/// the order, by construction, with nothing to track alongside it the way
401/// go's fix (tracking iteration order separately) or dotnet's fix (an
402/// explicit `Ordinal` field) both had to.
403pub struct Pool {
404 identity: Arc<KeyPair>,
405 trust: Trust,
406 options: PoolOptions,
407 links: RwLock<Vec<Arc<PooledLink>>>,
408 stop_tx: watch::Sender<bool>,
409 /// Every background task this pool has spawned (a link's respawn
410 /// lifecycle, the discovery loop). `Pool::close` hard-aborts and
411 /// awaits every one of these BEFORE draining/closing `links` — found
412 /// necessary by adversarial review, 2026-09-05: without this, a task
413 /// mid-dial (`connection::connect` has no internal cooperation point
414 /// with the pool's own stop signal) or mid-discovery-push could
415 /// complete AFTER `close` returns and resurrect a link, or a live,
416 /// connected `Session`, that nothing will ever close again. Sequencing
417 /// `shutdown()` strictly before the drain guarantees any task that
418 /// manages to finish its own critical section either finishes before
419 /// `close` starts draining (so the drain sees and closes it) or gets
420 /// aborted before it can (so nothing is left to see).
421 tasks: Mutex<JoinSet<()>>,
422}
423
424impl Pool {
425 /// Spawn a pool with one link per seed. Returns as soon as every link's
426 /// dial has STARTED, not once any is connected — handshakes complete
427 /// asynchronously, matching macula_client:connect/2 and every other
428 /// port of this pool shape.
429 pub fn connect(
430 seeds: Vec<Seed>,
431 trust: Trust,
432 identity: KeyPair,
433 options: PoolOptions,
434 ) -> Arc<Pool> {
435 assert!(!seeds.is_empty(), "at least one seed is required");
436 let (stop_tx, _stop_rx) = watch::channel(false);
437 let identity = Arc::new(identity);
438 let pool = Arc::new(Pool {
439 identity: identity.clone(),
440 trust,
441 options,
442 links: RwLock::new(Vec::new()),
443 stop_tx,
444 tasks: Mutex::new(JoinSet::new()),
445 });
446
447 let bootstrap_links: Vec<Arc<PooledLink>> = seeds
448 .into_iter()
449 .map(|seed| PooledLink::new(seed, LinkOrigin::Bootstrap, trust))
450 .collect();
451
452 {
453 // Bootstrap links are known synchronously at construction, so
454 // this can be a blocking write via try_write rather than
455 // spawning a task just to populate the initial Vec -- no
456 // other task holds the lock yet.
457 let mut links = pool
458 .links
459 .try_write()
460 .expect("no other task can hold this lock before Pool::connect returns");
461 links.extend(bootstrap_links.iter().cloned());
462 }
463
464 {
465 // Same reasoning as the links lock above: nothing else can
466 // hold this lock yet either.
467 let mut tasks = pool
468 .tasks
469 .try_lock()
470 .expect("no other task can hold this lock before Pool::connect returns");
471 for link in bootstrap_links {
472 if try_claim_redial(&link) {
473 tasks.spawn(link_lifecycle_future(pool.clone(), link));
474 }
475 }
476 if pool.options.station_discovery.enabled {
477 tasks.spawn(discover_stations_loop(pool.clone()));
478 }
479 }
480
481 pool
482 }
483
484 /// Send a signed CALL, choosing among currently-connected links per
485 /// [`PoolOptions::link_selection`], trying each in order until one
486 /// answers (a transport-level failure marks that link disconnected and
487 /// triggers its respawn, then moves to the next candidate — a BOLT#4
488 /// ERROR response is still a successful `call` as far as this pool is
489 /// concerned, exactly like a bare [`Session::call`]).
490 pub async fn call(
491 self: &Arc<Self>,
492 procedure: &str,
493 realm: [u8; 32],
494 payload: crate::cbor::Value,
495 deadline_ms: i128,
496 ) -> Result<CallResponse, PoolCallError> {
497 let candidates = self.select_connected_links().await;
498 if candidates.is_empty() {
499 return Err(PoolCallError::NoHealthyStation);
500 }
501 let mut last_err = None;
502 for link in candidates {
503 let mut state = link.state.lock().await;
504 let Some(session) = state.session.as_mut() else {
505 continue; // raced with a disconnect between selection and lock
506 };
507 match session
508 .call(
509 procedure,
510 realm,
511 payload.clone(),
512 deadline_ms,
513 &self.identity,
514 self.options.call_timeout,
515 )
516 .await
517 {
518 Ok(resp) => return Ok(resp),
519 Err(e) => {
520 drop(state);
521 self.mark_disconnected(&link).await;
522 last_err = Some(e);
523 }
524 }
525 }
526 Err(last_err
527 .map(PoolCallError::AllFailed)
528 .unwrap_or(PoolCallError::NoHealthyStation))
529 }
530
531 /// Send a signed PUBLISH, fanning out to up to
532 /// [`PoolOptions::replication_factor`] currently-connected links
533 /// (ordered by [`PoolOptions::link_selection`]). Partial success counts
534 /// as success, matching macula-ts/macula-dotnet's own publish-fanout
535 /// contract.
536 pub async fn publish(self: &Arc<Self>, spec: &PublishSpec) -> Result<(), PoolPublishError> {
537 let candidates = self.select_connected_links().await;
538 if candidates.is_empty() {
539 return Err(PoolPublishError::NoHealthyStation);
540 }
541 let targets: Vec<_> = candidates
542 .into_iter()
543 .take(self.options.replication_factor.max(1))
544 .collect();
545 let attempted = targets.len();
546 let mut successes = 0usize;
547 for link in targets {
548 let mut state = link.state.lock().await;
549 let Some(session) = state.session.as_mut() else {
550 continue;
551 };
552 match session.publish(spec, &self.identity).await {
553 Ok(()) => successes += 1,
554 Err(_) => {
555 drop(state);
556 self.mark_disconnected(&link).await;
557 }
558 }
559 }
560 if successes > 0 {
561 Ok(())
562 } else {
563 Err(PoolPublishError::AllFailed(attempted))
564 }
565 }
566
567 /// Aggregate health snapshot.
568 pub async fn status(&self) -> PoolStatus {
569 let links = self.links.read().await;
570 let healthy_links = links.iter().filter(|l| l.is_connected()).count();
571 PoolStatus {
572 healthy_links,
573 total_links: links.len(),
574 }
575 }
576
577 /// Per-link snapshot, in seed-list/discovery order — see [`Pool`]'s own
578 /// doc on why a plain `Vec` already guarantees this without any extra
579 /// bookkeeping.
580 pub async fn links(&self) -> Vec<LinkInfo> {
581 let links = self.links.read().await;
582 let mut out = Vec::with_capacity(links.len());
583 for link in links.iter() {
584 out.push(LinkInfo {
585 seed: link.seed.clone(),
586 connected: link.is_connected(),
587 node_id: if link.is_connected() {
588 link.peer_node_id().await
589 } else {
590 None
591 },
592 });
593 }
594 out
595 }
596
597 /// Sends GOODBYE on every currently-connected link and stops all
598 /// background dial/discovery tasks. Waits for every background task
599 /// (respawn lifecycles, the discovery loop) to actually be gone
600 /// BEFORE draining/closing `links` — see [`Pool`]'s own field doc on
601 /// `tasks` for why this ordering, specifically, is load-bearing.
602 /// Does not wait for the GOODBYE writes themselves to finish being
603 /// scheduled beyond [`Session::close`]'s own bounded drain.
604 pub async fn close(&self, reason: &str, detail: Option<&str>) {
605 let _ = self.stop_tx.send(true);
606 self.tasks.lock().await.shutdown().await;
607 let mut links = self.links.write().await;
608 for link in links.drain(..) {
609 let mut state = link.state.lock().await;
610 if let Some(session) = state.session.take() {
611 session.close(reason, detail, &self.identity).await;
612 }
613 link.connected.store(false, Ordering::Release);
614 }
615 }
616
617 async fn select_connected_links(&self) -> Vec<Arc<PooledLink>> {
618 let links = self.links.read().await;
619 let connected: Vec<Arc<PooledLink>> =
620 links.iter().filter(|l| l.is_connected()).cloned().collect();
621 drop(links);
622 let resolved = resolve_link_selection(
623 self.options.link_selection,
624 self.options.station_discovery.enabled,
625 );
626 select_links(connected, resolved)
627 }
628
629 async fn mark_disconnected(self: &Arc<Self>, link: &Arc<PooledLink>) {
630 let mut state = link.state.lock().await;
631 state.session = None;
632 link.connected.store(false, Ordering::Release);
633 drop(state);
634 // Guarded: Pool::call/Pool::publish can each independently observe
635 // the same dead session and both reach this point for the SAME
636 // link (see PooledLink::redialing's own doc) -- only the task that
637 // wins try_claim_redial's CAS actually spawns a respawn task.
638 if try_claim_redial(link) {
639 self.tasks
640 .lock()
641 .await
642 .spawn(link_lifecycle_future(self.clone(), link.clone()));
643 }
644 }
645}
646
647/// Attempts to claim the right to run a lifecycle task for `link`, via a
648/// false->true compare-exchange on [`PooledLink::redialing`]. Returns
649/// `true` only for the ONE caller that wins the race — see that field's
650/// own doc for why this exists. A caller that loses must NOT spawn
651/// anything; the winner's own task is already responsible for this link.
652fn try_claim_redial(link: &PooledLink) -> bool {
653 link.redialing
654 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
655 .is_ok()
656}
657
658/// Builds the future [`try_claim_redial`]'s winner spawns onto
659/// [`Pool`]'s own `tasks` [`JoinSet`] — wraps [`run_link_lifecycle`] with
660/// the one thing every exit path (normal completion, give-up, or a hard
661/// abort from [`Pool::close`]) must do exactly once: release the
662/// `redialing` claim, so a LATER disconnect of this same link can spawn
663/// a fresh lifecycle task again. An abort drops this future without
664/// running the line after `.await`, which is fine — [`Pool::close`]
665/// never spawns anything after calling `shutdown`, so a claim left
666/// `true` forever after an abort has no observable effect.
667async fn link_lifecycle_future(pool: Arc<Pool>, link: Arc<PooledLink>) {
668 run_link_lifecycle(&pool, &link).await;
669 link.redialing.store(false, Ordering::Release);
670}
671
672/// Dial (or redial) `link` until it connects, then return — this task's
673/// job ends at a successful handshake; it does not keep running
674/// afterward (no persistent reader is needed for a Call/Publish-only
675/// pool, see this module's own doc). [`Pool::mark_disconnected`] spawns a
676/// fresh instance of this same function whenever a link goes down, so
677/// respawn is just "run this again", not a separate mechanism.
678async fn run_link_lifecycle(pool: &Arc<Pool>, link: &Arc<PooledLink>) {
679 let mut stop_rx = pool.stop_tx.subscribe();
680 loop {
681 if *stop_rx.borrow() {
682 return;
683 }
684 match connection::connect(&link.seed.host, link.seed.port, link.trust, &pool.identity).await
685 {
686 Ok(session) => {
687 let node_id = session.station.node_id;
688 let mut state = link.state.lock().await;
689 state.session = Some(session);
690 state.peer_node_id = Some(node_id);
691 drop(state);
692 link.connected.store(true, Ordering::Release);
693 link.consecutive_failures.store(0, Ordering::Release);
694 return;
695 }
696 Err(_e) => {
697 let failures = link.consecutive_failures.fetch_add(1, Ordering::AcqRel) + 1;
698 if should_give_up(link.origin, failures) {
699 link.gave_up.store(true, Ordering::Release);
700 return;
701 }
702 tokio::select! {
703 _ = tokio::time::sleep(pool.options.respawn_delay) => {}
704 _ = stop_rx.changed() => {
705 if *stop_rx.borrow() {
706 return;
707 }
708 }
709 }
710 }
711 }
712 }
713}
714
715// ---------------------------------------------------------------------
716// Station discovery — resolves hecate_stations.list_stations' realm via
717// the DHT, calls it through the pool's own `call` (never a raw Session),
718// and additively adds links for whatever it finds.
719// ---------------------------------------------------------------------
720
721const LIST_STATIONS_PROCEDURE: &str = "hecate_stations.list_stations";
722/// `hecate_stations.list_stations` itself is called under its OWN resolved
723/// realm (whatever [`resolve_list_stations_realm`] found), NOT the DHT's
724/// own all-zero realm — that's `dht::find_records_by_type`'s realm
725/// (`dht.rs`'s own private `DHT_REALM` constant), a separate, unrelated
726/// realm this module never needs to name directly since it only ever
727/// reaches the DHT through `dht::find_records_by_type` itself.
728const DISCOVERY_CALL_DEADLINE: Duration = Duration::from_secs(5);
729
730async fn discover_stations_loop(pool: Arc<Pool>) {
731 // Trust::Pinned can never validate a SECOND station's identity --
732 // running discovery at all under it is unconditionally pointless, not
733 // just risky (it would otherwise dial-storm forever against stations
734 // it can never actually trust). Matches the identical guard in
735 // macula-go's and macula-dotnet's own ports of this feature.
736 if matches!(pool.trust, Trust::Pinned(_)) {
737 return;
738 }
739
740 let mut stop_rx = pool.stop_tx.subscribe();
741 if !wait_for_any_healthy_link(&pool, &mut stop_rx).await {
742 return;
743 }
744
745 loop {
746 if *stop_rx.borrow() {
747 return;
748 }
749 discover_once(&pool).await;
750 tokio::select! {
751 _ = tokio::time::sleep(pool.options.station_discovery.refresh_interval) => {}
752 _ = stop_rx.changed() => {
753 if *stop_rx.borrow() {
754 return;
755 }
756 }
757 }
758 }
759}
760
761async fn wait_for_any_healthy_link(pool: &Arc<Pool>, stop_rx: &mut watch::Receiver<bool>) -> bool {
762 loop {
763 if *stop_rx.borrow() {
764 return false;
765 }
766 if pool.status().await.healthy_links > 0 {
767 return true;
768 }
769 tokio::select! {
770 _ = tokio::time::sleep(Duration::from_millis(200)) => {}
771 _ = stop_rx.changed() => {
772 if *stop_rx.borrow() {
773 return false;
774 }
775 }
776 }
777 }
778}
779
780async fn discover_once(pool: &Arc<Pool>) {
781 let Some(realm) = resolve_list_stations_realm(pool).await else {
782 return;
783 };
784
785 let deadline = now_ms() + DISCOVERY_CALL_DEADLINE.as_millis() as i128;
786 let Ok(CallResponse::Result { payload, .. }) = pool
787 .call(
788 LIST_STATIONS_PROCEDURE,
789 realm,
790 crate::cbor::Value::Map(vec![]),
791 deadline,
792 )
793 .await
794 else {
795 return;
796 };
797 let crate::cbor::Value::Map(fields) = &payload else {
798 return;
799 };
800 let Some(crate::cbor::Value::List(stations)) = fields
801 .iter()
802 .find(|(k, _)| matches!(k, crate::cbor::Value::Text(t) if t == "stations"))
803 .map(|(_, v)| v.clone())
804 else {
805 return;
806 };
807
808 add_discovered_links(pool, &stations).await;
809}
810
811/// Resolves `hecate_stations.list_stations`' own realm by scanning every
812/// `procedure_advertisement` DHT record visible from the pool's current
813/// bootstrap connection, matching a `procedure_uri` of the shape
814/// `hex(realm) + "/hecate_stations.list_stations"` — mirrors
815/// `dht::discovery_uri`'s own format and macula-go/macula-dotnet's
816/// identical resolution step. Returns `None` on any failure (no
817/// advertisement found yet, none verify, no healthy link to ask with) —
818/// discovery just tries again at the next refresh tick, same as a bare DHT
819/// lookup miss anywhere else in this crate.
820async fn resolve_list_stations_realm(pool: &Arc<Pool>) -> Option<[u8; 32]> {
821 let links = pool.links.read().await;
822 let link = links.iter().find(|l| l.is_connected())?.clone();
823 drop(links);
824
825 let mut state = link.state.lock().await;
826 let session = state.session.as_mut()?;
827 let records =
828 dht::find_records_by_type(session, &pool.identity, dht::TYPE_PROCEDURE_ADVERTISEMENT)
829 .await
830 .ok()?;
831 drop(state);
832
833 for record in records {
834 if dht::verify(&record).is_err() {
835 continue;
836 }
837 let Ok(advertisement) = dht::read_procedure_advertisement(&record) else {
838 continue;
839 };
840 if let Some(realm) = try_match_list_stations_realm(&advertisement.procedure_uri) {
841 return Some(realm);
842 }
843 }
844 None
845}
846
847fn try_match_list_stations_realm(procedure_uri: &str) -> Option<[u8; 32]> {
848 let suffix = format!("/{LIST_STATIONS_PROCEDURE}");
849 let hex_realm = procedure_uri.strip_suffix(&suffix)?;
850 if hex_realm.len() != 64 {
851 return None;
852 }
853 let bytes = hex_decode(hex_realm)?;
854 bytes.try_into().ok()
855}
856
857fn hex_decode(s: &str) -> Option<Vec<u8>> {
858 if s.len() % 2 != 0 {
859 return None;
860 }
861 (0..s.len())
862 .step_by(2)
863 .map(|i| u8::from_str_radix(&s[i..i + 2], 16).ok())
864 .collect()
865}
866
867/// How many of `links` currently count against
868/// [`StationDiscoveryOptions::max_links`] — MaxLinks bounds discovery's OWN
869/// adds only, not the pool's total link count (see that field's own doc),
870/// so a bootstrap seed must NEVER count against this budget, however many
871/// were configured, or a pool started with >= max_links bootstrap seeds (a
872/// realistic deployment) would have discovery silently do nothing, forever,
873/// from its very first refresh — caught by adversarial review, 2026-09-05,
874/// after an earlier version of this count included every link regardless
875/// of origin. A given-up discovery link ALSO stays in the pool's `links`
876/// forever (no removal path — see [`Pool`]'s own doc on why the backing
877/// store is a plain `Vec`) but must NOT keep occupying its own slot either,
878/// or the whole point of giving up (freeing room for a different
879/// candidate) is defeated.
880fn count_occupied_discovery_slots(links: &[Arc<PooledLink>]) -> usize {
881 links
882 .iter()
883 .filter(|l| l.origin == LinkOrigin::Discovered && !l.has_given_up())
884 .count()
885}
886
887async fn add_discovered_links(pool: &Arc<Pool>, stations: &[crate::cbor::Value]) {
888 let is_web_pki = matches!(pool.trust, Trust::WebPki);
889 for station in stations {
890 let links = pool.links.read().await;
891 let occupied_slots = count_occupied_discovery_slots(&links);
892 drop(links);
893 if occupied_slots >= pool.options.station_discovery.max_links {
894 break;
895 }
896 let Some((host, port, node_id)) = dial_target_from_station_row(station) else {
897 continue;
898 };
899 let is_bare_ip = host.parse::<std::net::IpAddr>().is_ok();
900
901 // Prefer a real hostname under the pool's own configured trust;
902 // fall back to Trust::Pinned(node_id) for a bare-IP-only row when
903 // a node_id is available -- see this module's own doc for why,
904 // and cam2me's StationDiscovery.kt for the shipped precedent.
905 let link_trust = if is_bare_ip && is_web_pki {
906 match node_id {
907 Some(id) => Trust::Pinned(id),
908 None => continue, // no way to validate this station at all
909 }
910 } else {
911 pool.trust
912 };
913
914 if let Some(id) = node_id {
915 if has_link_for_node_id(pool, id).await {
916 continue;
917 }
918 }
919
920 let seed = Seed::new(host, port);
921 spawn_seed_link_if_absent(pool, seed, link_trust).await;
922 }
923}
924
925/// Extracts a dialable `(host, port, node_id)` from one
926/// `hecate_stations.list_stations` response row. Prefers `hostname` over
927/// `host_advertised[0]` when both are present — every station on the real
928/// fleet advertises `host_advertised` as a bare IP literal, never a DNS
929/// name (confirmed live, same finding independently hit by macula-go's and
930/// macula-dotnet's own ports of this feature), so `hostname` is the only
931/// field a WebPki dial can ever succeed against. `host_advertised` is
932/// still read out (and returned) even when a `hostname` is present, since
933/// callers that end up needing `Trust::Pinned` (see
934/// [`add_discovered_links`]) dial by IP regardless of whether a hostname
935/// also exists.
936pub(crate) fn dial_target_from_station_row(
937 row: &crate::cbor::Value,
938) -> Option<(String, u16, Option<[u8; 32]>)> {
939 let crate::cbor::Value::Map(fields) = row else {
940 return None;
941 };
942 let get = |name: &str| {
943 fields
944 .iter()
945 .find(|(k, _)| matches!(k, crate::cbor::Value::Text(t) if t == name))
946 .map(|(_, v)| v)
947 };
948
949 let port = match get("quic_port") {
950 Some(crate::cbor::Value::Int(n)) if (1..=65535).contains(n) => *n as u16,
951 _ => return None,
952 };
953
954 let hostname = match get("hostname") {
955 Some(crate::cbor::Value::Text(t)) if !t.is_empty() => Some(t.clone()),
956 Some(crate::cbor::Value::Bytes(b)) if !b.is_empty() => String::from_utf8(b.clone()).ok(),
957 _ => None,
958 };
959 let host_advertised = match get("host_advertised") {
960 Some(crate::cbor::Value::List(items)) => items.iter().find_map(|item| match item {
961 crate::cbor::Value::Bytes(b) => String::from_utf8(b.clone()).ok(),
962 crate::cbor::Value::Text(t) => Some(t.clone()),
963 _ => None,
964 }),
965 _ => None,
966 };
967
968 let host = hostname.or(host_advertised)?;
969 let node_id = match get("node_id") {
970 Some(crate::cbor::Value::Bytes(b)) => b.as_slice().try_into().ok(),
971 _ => None,
972 };
973 Some((host, port, node_id))
974}
975
976async fn has_link_for_node_id(pool: &Arc<Pool>, node_id: [u8; 32]) -> bool {
977 let links = pool.links.read().await;
978 for link in links.iter() {
979 if link.is_connected() {
980 if let Some(known) = link.peer_node_id().await {
981 if known == node_id {
982 return true;
983 }
984 }
985 }
986 }
987 false
988}
989
990async fn spawn_seed_link_if_absent(pool: &Arc<Pool>, seed: Seed, trust: Trust) {
991 let mut links = pool.links.write().await;
992 if links.iter().any(|l| l.seed == seed) {
993 return;
994 }
995 let link = PooledLink::new(seed, LinkOrigin::Discovered, trust);
996 links.push(link.clone());
997 drop(links);
998 if try_claim_redial(&link) {
999 pool.tasks
1000 .lock()
1001 .await
1002 .spawn(link_lifecycle_future(pool.clone(), link));
1003 }
1004}
1005
1006fn now_ms() -> i128 {
1007 use std::time::{SystemTime, UNIX_EPOCH};
1008 SystemTime::now()
1009 .duration_since(UNIX_EPOCH)
1010 .expect("system clock before 1970")
1011 .as_millis() as i128
1012}
1013
1014#[cfg(test)]
1015mod tests {
1016 use super::*;
1017
1018 fn row(fields: Vec<(&str, crate::cbor::Value)>) -> crate::cbor::Value {
1019 crate::cbor::Value::Map(
1020 fields
1021 .into_iter()
1022 .map(|(k, v)| (crate::cbor::Value::Text(k.to_string()), v))
1023 .collect(),
1024 )
1025 }
1026
1027 #[test]
1028 fn resolve_link_selection_auto_pairs_with_discovery() {
1029 assert_eq!(
1030 resolve_link_selection(LinkSelection::Auto, false),
1031 LinkSelection::FirstSuccess
1032 );
1033 assert_eq!(
1034 resolve_link_selection(LinkSelection::Auto, true),
1035 LinkSelection::Random
1036 );
1037 }
1038
1039 #[test]
1040 fn resolve_link_selection_explicit_survives_either_way() {
1041 assert_eq!(
1042 resolve_link_selection(LinkSelection::FirstSuccess, true),
1043 LinkSelection::FirstSuccess
1044 );
1045 assert_eq!(
1046 resolve_link_selection(LinkSelection::Random, false),
1047 LinkSelection::Random
1048 );
1049 }
1050
1051 #[test]
1052 fn dial_target_prefers_hostname_over_bare_ip() {
1053 let r = row(vec![
1054 (
1055 "hostname",
1056 crate::cbor::Value::Text("station-de-frankfurt.macula.io".into()),
1057 ),
1058 (
1059 "host_advertised",
1060 crate::cbor::Value::List(vec![crate::cbor::Value::Bytes(
1061 b"2a01:7e01::f03c:94ff:fe22:719e".to_vec(),
1062 )]),
1063 ),
1064 ("quic_port", crate::cbor::Value::Int(4433)),
1065 ]);
1066 let (host, port, _) = dial_target_from_station_row(&r).expect("should parse");
1067 assert_eq!(host, "station-de-frankfurt.macula.io");
1068 assert_eq!(port, 4433);
1069 }
1070
1071 #[test]
1072 fn dial_target_falls_back_to_host_advertised_when_hostname_absent() {
1073 let r = row(vec![
1074 (
1075 "host_advertised",
1076 crate::cbor::Value::List(vec![crate::cbor::Value::Bytes(
1077 b"2600:3c0b::2000:1fff:fe35:416b".to_vec(),
1078 )]),
1079 ),
1080 ("quic_port", crate::cbor::Value::Int(4433)),
1081 ]);
1082 let (host, _, _) = dial_target_from_station_row(&r).expect("should parse");
1083 assert_eq!(host, "2600:3c0b::2000:1fff:fe35:416b");
1084 }
1085
1086 #[test]
1087 fn dial_target_rejects_missing_port() {
1088 let r = row(vec![("hostname", crate::cbor::Value::Text("x".into()))]);
1089 assert!(dial_target_from_station_row(&r).is_none());
1090 }
1091
1092 #[test]
1093 fn dial_target_extracts_node_id() {
1094 let node_id = [0xABu8; 32];
1095 let r = row(vec![
1096 ("hostname", crate::cbor::Value::Text("x".into())),
1097 ("quic_port", crate::cbor::Value::Int(4433)),
1098 ("node_id", crate::cbor::Value::Bytes(node_id.to_vec())),
1099 ]);
1100 let (_, _, id) = dial_target_from_station_row(&r).expect("should parse");
1101 assert_eq!(id, Some(node_id));
1102 }
1103
1104 #[test]
1105 fn should_give_up_never_applies_to_a_bootstrap_link() {
1106 assert!(!should_give_up(
1107 LinkOrigin::Bootstrap,
1108 DISCOVERY_LINK_MAX_RESPAWN_ATTEMPTS
1109 ));
1110 assert!(!should_give_up(LinkOrigin::Bootstrap, 1_000_000));
1111 }
1112
1113 #[test]
1114 fn should_give_up_applies_to_a_discovered_link_at_the_threshold() {
1115 assert!(!should_give_up(
1116 LinkOrigin::Discovered,
1117 DISCOVERY_LINK_MAX_RESPAWN_ATTEMPTS - 1
1118 ));
1119 assert!(should_give_up(
1120 LinkOrigin::Discovered,
1121 DISCOVERY_LINK_MAX_RESPAWN_ATTEMPTS
1122 ));
1123 }
1124
1125 fn synthetic_link(origin: LinkOrigin, gave_up: bool) -> Arc<PooledLink> {
1126 let link = PooledLink::new(Seed::new("x.example", 4433), origin, Trust::WebPki);
1127 link.gave_up.store(gave_up, Ordering::Relaxed);
1128 link
1129 }
1130
1131 /// Regression test for the exact bug adversarial review caught: an
1132 /// earlier version of this count included bootstrap links, which can
1133 /// never give up, so a pool started with `bootstrap.len() >=
1134 /// max_links` would have discovery silently do nothing forever.
1135 #[test]
1136 fn discovery_slot_count_excludes_bootstrap_links_entirely() {
1137 let links = vec![
1138 synthetic_link(LinkOrigin::Bootstrap, false),
1139 synthetic_link(LinkOrigin::Bootstrap, false),
1140 synthetic_link(LinkOrigin::Bootstrap, false),
1141 ];
1142 assert_eq!(count_occupied_discovery_slots(&links), 0);
1143 }
1144
1145 #[test]
1146 fn discovery_slot_count_excludes_a_given_up_discovered_link() {
1147 let links = vec![
1148 synthetic_link(LinkOrigin::Discovered, false),
1149 synthetic_link(LinkOrigin::Discovered, true), // gave up -- slot freed
1150 ];
1151 assert_eq!(count_occupied_discovery_slots(&links), 1);
1152 }
1153
1154 #[test]
1155 fn discovery_slot_count_mixed_origins() {
1156 let links = vec![
1157 synthetic_link(LinkOrigin::Bootstrap, false),
1158 synthetic_link(LinkOrigin::Bootstrap, false),
1159 synthetic_link(LinkOrigin::Discovered, false),
1160 synthetic_link(LinkOrigin::Discovered, true),
1161 ];
1162 assert_eq!(count_occupied_discovery_slots(&links), 1);
1163 }
1164
1165 #[test]
1166 fn try_claim_redial_only_lets_one_caller_win() {
1167 let link = synthetic_link(LinkOrigin::Discovered, false);
1168 assert!(try_claim_redial(&link));
1169 assert!(
1170 !try_claim_redial(&link),
1171 "a second claim must fail while the first is outstanding"
1172 );
1173 link.redialing.store(false, Ordering::Release);
1174 assert!(
1175 try_claim_redial(&link),
1176 "releasing the claim must allow a fresh one"
1177 );
1178 }
1179
1180 #[test]
1181 fn try_match_list_stations_realm_matches_expected_format() {
1182 let hex_realm = "0".repeat(64);
1183 let uri = format!("{hex_realm}/hecate_stations.list_stations");
1184 assert_eq!(try_match_list_stations_realm(&uri), Some([0u8; 32]));
1185 }
1186
1187 #[test]
1188 fn try_match_list_stations_realm_rejects_a_different_procedure() {
1189 let hex_realm = "0".repeat(64);
1190 let uri = format!("{hex_realm}/some.other_procedure");
1191 assert_eq!(try_match_list_stations_realm(&uri), None);
1192 }
1193
1194 #[test]
1195 fn select_links_first_success_returns_input_unshuffled() {
1196 // Empty/zero-length input is the simplest observable proof this
1197 // is a passthrough -- Vec equality on non-empty synthetic
1198 // PooledLinks would need constructing real Arc<PooledLink>s with
1199 // no real Session, which is exactly the "no fake-dialer seam"
1200 // situation macula-dotnet's own tests document; a live pool test
1201 // covers the real end-to-end ordering instead.
1202 let empty: Vec<Arc<PooledLink>> = Vec::new();
1203 assert_eq!(select_links(empty, LinkSelection::FirstSuccess).len(), 0);
1204 }
1205}