Skip to main content

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}