Skip to main content

kanade_shared/
nats_client.rs

1//! Shared NATS client constructor.
2//!
3//! Every binary names the [`NatsRole`] it connects as, and the token is
4//! resolved per role (first match wins):
5//!
6//!   1. Windows registry — `HKLM\SOFTWARE\kanade\<role>\NatsToken`
7//!      (`REG_SZ`). The role-specific credential. Hardened ACL (SYSTEM +
8//!      Admin only) keeps the token out of low-privilege users' reach,
9//!      which Machine-scope env vars cannot do.
10//!   2. Windows registry — `HKLM\SOFTWARE\kanade\agent\NatsToken`. The
11//!      **shared** credential every role used before roles existed. Kept as
12//!      a fallback so an existing deployment keeps working untouched; see
13//!      "Staged migration" below.
14//!   3. `$KANADE_NATS_TOKEN` environment variable. Dev / fallback path. The
15//!      agent service runs as LocalSystem so user-session env vars never
16//!      reach it; this branch only fires for `cargo run` / interactive
17//!      shells.
18//!   4. No token — connect unauthenticated. Works against a broker started
19//!      without `authorization { … }`.
20//!
21//! # Per-role users, with the token as fallback
22//!
23//! A role may also hold a NATS *user*: `NatsUser` and `NatsPassword`
24//! (`REG_SZ`) under its **own** registry key, or `$KANADE_NATS_USER` /
25//! `$KANADE_NATS_PASSWORD`. There is deliberately no shared user — the shared
26//! credential is the token. A pair with only one half is a configuration
27//! error that fails the connect, never a silent fallback.
28//!
29//! With no user provisioned nothing below applies: the connection presents
30//! the token (or nothing) exactly as before and no probe is ever made.
31//!
32//! With a user provisioned the broker may be on either side of the
33//! token → `users` switch, and the switch is atomic on the broker, so the
34//! credential is chosen per connection attempt (initial and every reconnect)
35//! from an auth callback. Whether a rejected attempt is retried is up to the
36//! async-nats version (0.50 keeps retrying, and re-runs the callback, while
37//! another version may end the connection task for good, leaving a `Client`
38//! that never talks again), so "present the user, and on rejection retry with
39//! the token" is not left to its retry loop. Instead each attempt first opens
40//! a short-lived *probe* connection with the user, which does not retry:
41//!
42//!   * accepted → present the user;
43//!   * rejected (authorization violation) → present the token, or fail
44//!     naming the missing token if there is none;
45//!   * anything else (unreachable, timeout) → defer the attempt with an
46//!     authentication callback error. A cached answer describes an old broker
47//!     mode and must not become a refused handshake after a restart.
48//!
49//! The process keeps one `Client`: a deferred attempt is retried and re-probed.
50//! That a refusal is retried and never ends the connection task is pinned
51//! against a real broker for the pinned async-nats, so a version change that
52//! breaks it fails a test instead of leaving a silent client.
53//!
54//! What a switch or restart guarantees is stated as [`RESUME_BOUND`] and
55//! [`EXIT_BOUND`].
56//! A client whose connection task has terminated for any reason (a panic in
57//! it, a drain, a version that treats a violation as terminal) is detected by
58//! [`wait_until_dead`] so its process can exit for a supervised restart. The
59//! same call also reports a broker that refuses the credential on every
60//! attempt for a sustained period, since async-nats would otherwise retry it
61//! forever and leave a process that looks alive but can never talk.
62//!
63//! # Why roles exist here (#1155)
64//!
65//! The broker authorises a *connection*, and a connection is only as
66//! specific as the credential that opened it. While every binary presented
67//! the same token, the broker could not tell an agent from the backend from
68//! the CLI, so no `permissions` block could say "only the backend may
69//! subscribe `remote.frame.>`" — there was nothing to hang the rule on.
70//! That is why a shared token means a token holder can execute code on any
71//! endpoint and silently watch any remote-assistance session (#1140).
72//!
73//! Distinct credentials do not fix that by themselves; the broker config has
74//! to grow the matching `authorization { users: [...] }` entries. This
75//! module is the half that makes those entries *expressible*.
76//!
77//! # Staged migration
78//!
79//! Step 2 above is the whole migration strategy. A fleet running today has
80//! one token, provisioned at `…\kanade\agent\NatsToken` on every host
81//! regardless of role. After this change it keeps working: no role key
82//! exists, so every role falls through to the shared one and presents
83//! exactly what it presented before.
84//!
85//! Rolling out per-role credentials is then per-host and reversible — write
86//! `…\kanade\backend\NatsToken` on the backend host and it starts using it;
87//! delete it and it falls back. The broker only needs to start
88//! *distinguishing* the roles once every host has its own, so the config
89//! change lands last, when it can no longer lock anyone out.
90//!
91//! No deploy script writes a role key yet: `deploy-backend.ps1` still
92//! provisions the shared path, so today the role key is a manual registry
93//! write. That is deliberate — the scripted path should start writing role
94//! keys in the same change that teaches the broker to tell the roles apart,
95//! because until then a role key changes nothing and a script that writes
96//! only the role key (dropping the shared one) would strand the CLI on a
97//! backend-only host.
98//!
99//! The order matters and is deliberate: role key first, shared key second.
100//! The reverse would make the shared token permanent — a host that still has
101//! it (all of them, today) would never notice its role key.
102//!
103//! # What the broker will and will not accept (measured, #1270)
104//!
105//! Two nats-server behaviours constrain every plan built on this module, so
106//! they are recorded here rather than rediscovered:
107//!
108//! * A config may not carry **both** a `token` and a `users` array —
109//!   nats-server refuses to start: *"Can not have a token and a users
110//!   array"*. And once `users` are defined, a client presenting a token is
111//!   rejected with an Authorization Violation, even when the token equals a
112//!   user's password. So the shared token and a per-role `users` split
113//!   cannot coexist for a transition window: the flip is atomic, and every
114//!   host must already hold a credential of the new shape before it happens.
115//!   Resolving a *token* per role, which is all this module does today, is
116//!   therefore not sufficient for that split — the client has to learn to
117//!   present a user as well.
118//! * `/connz?auth=1` reports `authorized_user` per connection. Under
119//!   `users` that is the username — the per-host answer #1270 wants. Under
120//!   `token`, nats-server 2.14.3 reports the literal `[REDACTED]`: it hides
121//!   the credential, so token mode can say *that* a host is on the shared
122//!   token but never anything finer. Whether the value is hidden is the
123//!   broker build's choice, not ours, so a consumer must assume it may be
124//!   handling a secret; see [`CredentialProbe`] for the one question it can
125//!   safely ask about one.
126//!
127//! # Limits worth naming
128//!
129//! A per-role token still cannot express per-*agent* identity. A role
130//! credential permitted to subscribe `commands.pc.*` lets any agent holding
131//! it read another agent's inbox. This narrows a fleet-wide compromise to a
132//! fleet-wide **agent-role** compromise, which is better, not solved. The
133//! end state is per-agent identity (NKeys / NATS-JWT), for which the plan is
134//! to grow `ConnectOptions` here so every binary picks up the upgrade for
135//! free. Same for mTLS.
136
137use std::collections::HashMap;
138use std::sync::atomic::Ordering;
139use std::sync::{Arc, Mutex, OnceLock, Weak};
140use std::time::Duration;
141
142use anyhow::{Context, Result, bail};
143use tracing::{debug, info, warn};
144
145use crate::secrets;
146
147const ENV_TOKEN: &str = "KANADE_NATS_TOKEN";
148const REG_VALUE: &str = "NatsToken";
149const ENV_USER: &str = "KANADE_NATS_USER";
150const ENV_PASSWORD: &str = "KANADE_NATS_PASSWORD";
151const REG_USER: &str = "NatsUser";
152const REG_PASSWORD: &str = "NatsPassword";
153
154/// How long the credential probe may take. Short on purpose: it runs inside
155/// the reconnect path, so a black-holed broker must not stall every attempt.
156const PROBE_TIMEOUT: Duration = Duration::from_secs(3);
157
158/// Name of the probe connection. Deliberately *not* a `kanade-` name: the
159/// backend attributes `kanade-<role>` connections to hosts, and a transient
160/// probe must not show up there as an unknown role.
161const PROBE_NAME: &str = "auth-probe";
162
163/// How often the client sends a protocol PING when idle. Set explicitly (it
164/// equals async-nats' default) because every liveness bound below is derived
165/// from it: a healthy idle connection receives a PONG at least this often.
166const PING_INTERVAL: Duration = Duration::from_secs(60);
167
168/// How often [`wait_until_dead`] runs its local check. The check only reads
169/// counters and queues a flush on the local command channel, so a short
170/// interval costs no network traffic.
171const CHECK_INTERVAL: Duration = Duration::from_secs(5);
172
173/// How long without any received byte before the client is suspected stalled:
174/// three ping periods, so a healthy idle connection (one PONG per
175/// [`PING_INTERVAL`]) is two full periods away from being suspected.
176const STALL_BOUND: Duration = Duration::from_secs(3 * PING_INTERVAL.as_secs());
177
178const HEALTH_TIMEOUT: Duration = Duration::from_secs(5);
179
180/// The longest async-nats waits between reconnect attempts (its backoff cap).
181const RECONNECT_DELAY_MAX_SECS: u64 = 4;
182
183/// async-nats' default timeout for one connection handshake.
184const CONNECT_TIMEOUT_SECS: u64 = 5;
185
186/// The longest a single reconnect attempt can take: the backoff, then the
187/// credential probe (its own timeout plus the guard around it), then the
188/// handshake the probe's answer decides.
189const ATTEMPT_MAX_SECS: u64 =
190    RECONNECT_DELAY_MAX_SECS + PROBE_TIMEOUT.as_secs() + 1 + CONNECT_TIMEOUT_SECS;
191
192/// How long after the broker starts answering in its new authentication mode
193/// (a restart, or a switch between `token` and `users`) a client may take to
194/// talk again: two attempts, allowing a probe interrupted by the switch,
195/// plus slack. Both bounds are
196/// counted from the moment the broker answers, not from the moment it went
197/// away; a broker that is down is waited for indefinitely.
198pub const RESUME_BOUND: Duration = Duration::from_secs(2 * ATTEMPT_MAX_SECS + 4);
199
200/// How long after the broker starts answering in its new mode a client that
201/// has *not* resumed may stay alive before its process exits non-zero for the
202/// service manager to restart: the stall threshold, then the next local
203/// check, a bounded flush, and a bounded credential witness. The service
204/// manager's own restart delay comes on top. A broker outage seen by the
205/// witness only moves the threshold to [`RESUME_BOUND`] after it ends, which
206/// is never later than this.
207pub const EXIT_BOUND: Duration = Liveness::PRODUCTION.exit_bound();
208
209const _: () = {
210    assert!(STALL_BOUND.as_secs() >= PING_INTERVAL.as_secs() + 60);
211    assert!(RESUME_BOUND.as_secs() <= STALL_BOUND.as_secs());
212    assert!(EXIT_BOUND.as_secs() < 300);
213};
214
215/// The timing of [`wait_until_dead_with`]. Only [`Liveness::PRODUCTION`] is
216/// meant for shipped code; other values exist so tests can shorten the
217/// bounds, and a `stall_bound` at or below the ping interval would judge a
218/// healthy idle connection stalled.
219#[derive(Debug, Clone, Copy)]
220pub struct Liveness {
221    /// How often the local check runs.
222    pub check_interval: Duration,
223    /// How long without received bytes before a credential witness is tried.
224    pub stall_bound: Duration,
225}
226
227impl Liveness {
228    pub const PRODUCTION: Liveness = Liveness {
229        check_interval: CHECK_INTERVAL,
230        stall_bound: STALL_BOUND,
231    };
232
233    /// The longest a stalled client stays alive after the broker answers.
234    pub const fn exit_bound(&self) -> Duration {
235        Duration::from_secs(
236            self.stall_bound.as_secs()
237                + self.check_interval.as_secs()
238                + 2 * HEALTH_TIMEOUT.as_secs(),
239        )
240    }
241
242    /// The grace given after a witness finds the broker reachable again.
243    fn resume_grace(&self) -> Duration {
244        RESUME_BOUND.min(self.stall_bound)
245    }
246}
247
248/// What the supervisor should do after a credential witness.
249#[derive(Debug, Clone, Copy, PartialEq, Eq)]
250enum Verdict {
251    /// Keep watching.
252    Wait,
253    /// The broker answered and the client is still silent: exit, unless a
254    /// recheck shows progress.
255    Stalled,
256}
257
258/// The stall decision as a pure state machine over caller-supplied times, so
259/// the long-window behaviour is testable without waiting.
260struct StallJudge {
261    liveness: Liveness,
262    deadline: tokio::time::Instant,
263    /// The last witness found the broker unreachable. The silence then
264    /// describes the outage, not the client.
265    outage: bool,
266    next_witness: Option<tokio::time::Instant>,
267}
268
269impl StallJudge {
270    fn new(liveness: Liveness, now: tokio::time::Instant) -> Self {
271        Self {
272            liveness,
273            deadline: now + liveness.stall_bound,
274            outage: false,
275            next_witness: None,
276        }
277    }
278
279    /// Receive progress was observed: the clock restarts and any outage
280    /// verdict no longer applies.
281    fn progressed(&mut self, now: tokio::time::Instant) {
282        self.deadline = now + self.liveness.stall_bound;
283        self.outage = false;
284        self.next_witness = None;
285    }
286
287    /// Whether a witness connection is due. Never true before the deadline.
288    fn witness_due(&self, now: tokio::time::Instant) -> bool {
289        now >= self.deadline && self.next_witness.is_none_or(|at| now >= at)
290    }
291
292    fn witnessed(&mut self, now: tokio::time::Instant, outcome: ProbeOutcome) -> Verdict {
293        if outcome == ProbeOutcome::Unreachable {
294            self.outage = true;
295            self.next_witness = Some(now + self.liveness.resume_grace());
296            return Verdict::Wait;
297        }
298        if std::mem::take(&mut self.outage) {
299            // First sight of the broker after an outage: the client may still
300            // be in its reconnect backoff, so give it the resume grace.
301            self.deadline = self.deadline.max(now + self.liveness.resume_grace());
302            self.next_witness = None;
303            return Verdict::Wait;
304        }
305        Verdict::Stalled
306    }
307}
308
309/// A connection that has been refused this many times in a row, over at
310/// least [`AUTH_REJECTION_WINDOW`], with no successful connect in between, is
311/// treated as failed. async-nats retries a refused credential forever, which
312/// leaves a process that can never talk but looks alive; the window keeps a
313/// broker that is merely mid-reload (a few quick refusals) from counting.
314const AUTH_REJECTION_LIMIT: usize = 5;
315const AUTH_REJECTION_WINDOW: Duration = Duration::from_secs(10);
316
317/// Prefix every kanade connection announces in its `name`.
318const NAME_PREFIX: &str = "kanade-";
319
320/// Separator between the role and the host identity in a connection name.
321/// `/` is safe as a delimiter because the identity is a Windows computer
322/// name, which cannot contain one.
323const NAME_SEP: char = '/';
324
325/// Registry subkey holding the pre-#1155 shared credential. Also the agent's
326/// role key, which is not a coincidence — the shared token was provisioned
327/// under the agent's path because agents were the first thing to need it.
328const REG_SHARED_SUBKEY: &str = r"SOFTWARE\kanade\agent";
329
330/// Which kanade binary is opening the connection.
331///
332/// Named on every call rather than inferred, because the broker's view of a
333/// connection comes entirely from the credential it presents: a caller that
334/// picks the wrong role does not get a warning, it gets someone else's
335/// permissions.
336#[derive(Debug, Clone, Copy, PartialEq, Eq)]
337pub enum NatsRole {
338    /// The endpoint agent. The most numerous and least trusted role — one
339    /// compromised endpoint holds this credential.
340    Agent,
341    /// The backend. The only role that needs to see the whole fleet.
342    Backend,
343    /// The operator CLI, including the backend-down recovery path that
344    /// drives agents over NATS directly.
345    Cli,
346}
347
348impl NatsRole {
349    pub fn as_str(self) -> &'static str {
350        match self {
351            NatsRole::Agent => "agent",
352            NatsRole::Backend => "backend",
353            NatsRole::Cli => "cli",
354        }
355    }
356
357    /// Registry subkey holding this role's credential.
358    fn reg_subkey(self) -> String {
359        format!(r"SOFTWARE\kanade\{}", self.as_str())
360    }
361}
362
363/// Resolve a role's token, given a registry reader and the environment
364/// fallback.
365///
366/// Split from [`resolve_token`] so the *ordering* — the part that decides
367/// whether a migration is reversible — is testable without a Windows
368/// registry to write to.
369fn resolve_token_with(
370    role: NatsRole,
371    read_reg: impl Fn(&str, &str) -> Option<String>,
372    env: Option<String>,
373) -> Option<String> {
374    if let Some(t) = read_reg(&role.reg_subkey(), REG_VALUE) {
375        return Some(t);
376    }
377    if let Some(t) = read_reg(REG_SHARED_SUBKEY, REG_VALUE) {
378        return Some(t);
379    }
380    env.filter(|t| !t.is_empty())
381}
382
383fn resolve_token(role: NatsRole) -> Option<String> {
384    resolve_token_with(
385        role,
386        secrets::read_hklm_value,
387        std::env::var(ENV_TOKEN).ok(),
388    )
389}
390
391/// A named NATS user. The password is never shown by `Debug`.
392#[derive(Clone, PartialEq, Eq)]
393struct UserCredential {
394    name: String,
395    password: String,
396}
397
398impl std::fmt::Debug for UserCredential {
399    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
400        f.debug_struct("UserCredential")
401            .field("name", &self.name)
402            .field("password", &"<redacted>")
403            .finish()
404    }
405}
406
407/// Turn the two halves of a user pair into a credential, or a configuration
408/// error naming what is missing and where (never a value).
409fn pair(
410    name: Option<String>,
411    password: Option<String>,
412    source: &str,
413) -> Result<Option<UserCredential>> {
414    match (name, password) {
415        (Some(name), Some(password)) => Ok(Some(UserCredential { name, password })),
416        (None, None) => Ok(None),
417        (Some(_), None) => bail!("NATS user is set but its password is missing ({source})"),
418        (None, Some(_)) => bail!("NATS password is set but its user name is missing ({source})"),
419    }
420}
421
422/// Resolve a role's user pair: its own registry key first, then the
423/// environment. Unlike the token there is no shared-key step.
424///
425/// The registry is judged as a whole before the environment is consulted: a
426/// half pair there is an error even when the environment is complete,
427/// because quietly mixing sources would hide a half-provisioned host.
428fn resolve_user_with(
429    role: NatsRole,
430    read_reg: impl Fn(&str, &str) -> Option<String>,
431    env_user: Option<String>,
432    env_password: Option<String>,
433) -> Result<Option<UserCredential>> {
434    let subkey = role.reg_subkey();
435    let reg_user = read_reg(&subkey, REG_USER);
436    let reg_password = read_reg(&subkey, REG_PASSWORD);
437    if reg_user.is_some() || reg_password.is_some() {
438        let source = format!(r"HKLM\{subkey}: {REG_USER} / {REG_PASSWORD}");
439        return pair(reg_user, reg_password, &source);
440    }
441    let source = format!("${ENV_USER} / ${ENV_PASSWORD}");
442    pair(
443        env_user.filter(|v| !v.is_empty()),
444        env_password.filter(|v| !v.is_empty()),
445        &source,
446    )
447}
448
449fn resolve_user(role: NatsRole) -> Result<Option<UserCredential>> {
450    resolve_user_with(
451        role,
452        secrets::read_hklm_value,
453        std::env::var(ENV_USER).ok(),
454        std::env::var(ENV_PASSWORD).ok(),
455    )
456}
457
458/// The credentials a connection is built from.
459///
460/// [`connect`] resolves them from the registry / environment; this type
461/// exists so a caller (tests, chiefly) can supply them directly without
462/// touching process-global state. `Debug` never prints a secret.
463#[derive(Clone, Default)]
464pub struct NatsCredentials {
465    token: Option<String>,
466    user: Option<UserCredential>,
467}
468
469impl NatsCredentials {
470    pub fn new(token: Option<String>, user: Option<(String, String)>) -> Self {
471        Self {
472            token,
473            user: user.map(|(name, password)| UserCredential { name, password }),
474        }
475    }
476
477    fn resolve(role: NatsRole) -> Result<Self> {
478        Ok(Self {
479            token: resolve_token(role),
480            user: resolve_user(role)?,
481        })
482    }
483}
484
485impl std::fmt::Debug for NatsCredentials {
486    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
487        f.debug_struct("NatsCredentials")
488            .field("token", &self.token.as_ref().map(|_| "<redacted>"))
489            .field("user", &self.user)
490            .finish()
491    }
492}
493
494/// Which credential a connection attempt presents.
495#[derive(Debug, Clone, Copy, PartialEq, Eq)]
496enum Choice {
497    User,
498    Token,
499}
500
501impl Choice {
502    fn label(self) -> &'static str {
503        match self {
504            Choice::User => "user",
505            Choice::Token => "token",
506        }
507    }
508}
509
510/// What the probe connection learned about the broker.
511#[derive(Debug, Clone, Copy, PartialEq, Eq)]
512enum ProbeOutcome {
513    Accepted,
514    Rejected,
515    Unreachable,
516}
517
518impl ProbeOutcome {
519    fn label(self) -> &'static str {
520        match self {
521            ProbeOutcome::Accepted => "accepted",
522            ProbeOutcome::Rejected => "rejected",
523            ProbeOutcome::Unreachable => "unreachable",
524        }
525    }
526}
527
528/// The credential the broker last *proved* it accepts or rejects for a role,
529/// shared between the connection (which writes it on every attempt) and the
530/// [`CredentialProbe`] (which reads it). Only a probe answer is recorded — a
531/// guess made while the broker was unreachable is not evidence.
532#[derive(Default)]
533struct Live {
534    decision: Mutex<Option<Choice>>,
535    /// Consecutive authentication refusals: when the first one happened and
536    /// how many there have been. Cleared by a successful connect.
537    refusals: Mutex<Option<(std::time::Instant, usize)>>,
538}
539
540impl Live {
541    /// Track the connection's events for [`Live::auth_failed`].
542    fn observe(&self, ev: &async_nats::Event) {
543        let mut refusals = self.refusals.lock().unwrap_or_else(|e| e.into_inner());
544        match ev {
545            // A successful connect, a drop, or any other kind of failed
546            // attempt (broker unreachable, timeout) ends a run of refusals:
547            // only an unbroken run is evidence the credential itself is bad,
548            // and an offline broker is something to wait out.
549            async_nats::Event::Connected | async_nats::Event::Disconnected => *refusals = None,
550            async_nats::Event::ClientError(async_nats::ClientError::Other(kind))
551                if *kind == async_nats::ConnectErrorKind::AuthorizationViolation.to_string() =>
552            {
553                let (first, n) = refusals.unwrap_or((std::time::Instant::now(), 0));
554                *refusals = Some((first, n + 1));
555            }
556            // Callback errors include inconclusive probes, not broker refusals.
557            async_nats::Event::ClientError(async_nats::ClientError::Other(kind))
558                if *kind == async_nats::ConnectErrorKind::Authentication.to_string() => {}
559            async_nats::Event::ClientError(_) => *refusals = None,
560            _ => {}
561        }
562    }
563
564    /// Whether the broker has refused every attempt for long enough that
565    /// waiting longer is not going to help.
566    fn auth_failed(&self) -> bool {
567        match *self.refusals.lock().unwrap_or_else(|e| e.into_inner()) {
568            Some((first, n)) => {
569                n >= AUTH_REJECTION_LIMIT && first.elapsed() >= AUTH_REJECTION_WINDOW
570            }
571            None => false,
572        }
573    }
574
575    fn get(&self) -> Option<Choice> {
576        *self.decision.lock().unwrap_or_else(|e| e.into_inner())
577    }
578
579    fn set(&self, c: Choice) {
580        *self.decision.lock().unwrap_or_else(|e| e.into_inner()) = Some(c);
581    }
582}
583
584/// Process-wide per-role state, so [`CredentialProbe::for_role`] sees the
585/// decision the already-open connection made without the callers having to
586/// thread anything through.
587fn role_lives() -> &'static Mutex<HashMap<&'static str, Arc<Live>>> {
588    static LIVE: OnceLock<Mutex<HashMap<&'static str, Arc<Live>>>> = OnceLock::new();
589    LIVE.get_or_init(Default::default)
590}
591
592fn publish_role_live(role: NatsRole, live: &Live) {
593    // Preserve the reporting handle held by CredentialProbe::for_role.
594    let reported = live_for(role);
595    *reported.decision.lock().unwrap_or_else(|e| e.into_inner()) = live.get();
596}
597
598fn live_for(role: NatsRole) -> Arc<Live> {
599    role_lives()
600        .lock()
601        .unwrap_or_else(|e| e.into_inner())
602        .entry(role.as_str())
603        .or_default()
604        .clone()
605}
606
607struct ConnectionHealth {
608    live: Arc<Live>,
609    url: String,
610    creds: NatsCredentials,
611}
612
613struct HealthEntry {
614    stats: Weak<async_nats::client::Statistics>,
615    health: Arc<ConnectionHealth>,
616}
617
618fn health_entries() -> &'static Mutex<Vec<HealthEntry>> {
619    static HEALTH: OnceLock<Mutex<Vec<HealthEntry>>> = OnceLock::new();
620    HEALTH.get_or_init(Default::default)
621}
622
623fn register_health(client: &async_nats::Client, health: Arc<ConnectionHealth>) {
624    let mut entries = health_entries().lock().unwrap_or_else(|e| e.into_inner());
625    entries.retain(|entry| entry.stats.strong_count() > 0);
626    entries.push(HealthEntry {
627        stats: Arc::downgrade(&client.statistics()),
628        health,
629    });
630}
631
632fn connection_health(client: &async_nats::Client) -> Option<Arc<ConnectionHealth>> {
633    let stats = client.statistics();
634    health_entries()
635        .lock()
636        .unwrap_or_else(|e| e.into_inner())
637        .iter()
638        .find(|entry| entry.stats.ptr_eq(&Arc::downgrade(&stats)))
639        .map(|entry| entry.health.clone())
640}
641
642/// A fresh handshake proves the broker can accept this client's credentials.
643/// It never retries, and neither its traffic nor its errors advance the
644/// original client's progress clock or refusal accounting.
645async fn witness(health: &ConnectionHealth) -> ProbeOutcome {
646    let attempt = async {
647        let mut opts = async_nats::ConnectOptions::new()
648            .name("auth-witness")
649            .connection_timeout(HEALTH_TIMEOUT);
650        match &health.creds.user {
651            Some(user) => match probe_user(&health.url, user).await {
652                ProbeOutcome::Accepted => {
653                    opts = opts.user_and_password(user.name.clone(), user.password.clone());
654                }
655                ProbeOutcome::Rejected => match &health.creds.token {
656                    Some(token) => opts = opts.token(token.clone()),
657                    None => return ProbeOutcome::Rejected,
658                },
659                ProbeOutcome::Unreachable => return ProbeOutcome::Unreachable,
660            },
661            None => {
662                if let Some(token) = &health.creds.token {
663                    opts = opts.token(token.clone());
664                }
665            }
666        }
667        match opts.connect(&health.url).await {
668            Ok(_) => ProbeOutcome::Accepted,
669            Err(error) if error.kind() == async_nats::ConnectErrorKind::AuthorizationViolation => {
670                ProbeOutcome::Rejected
671            }
672            Err(_) => ProbeOutcome::Unreachable,
673        }
674    };
675    tokio::time::timeout(HEALTH_TIMEOUT, attempt)
676        .await
677        .unwrap_or(ProbeOutcome::Unreachable)
678}
679
680/// Pick the credential for one attempt from the probe's answer.
681///
682/// Pure so the whole decision table is testable without a broker. The cached
683/// decision is intentionally ignored: it describes a previous broker mode.
684fn select(
685    outcome: ProbeOutcome,
686    _last_worked: Option<Choice>,
687    have_token: bool,
688) -> std::result::Result<Choice, &'static str> {
689    match outcome {
690        ProbeOutcome::Accepted => Ok(Choice::User),
691        ProbeOutcome::Rejected if have_token => Ok(Choice::Token),
692        ProbeOutcome::Rejected => {
693            Err("the broker rejected the NATS user and no NatsToken is provisioned")
694        }
695        ProbeOutcome::Unreachable => {
696            Err("NATS credential probe inconclusive; deferring this attempt")
697        }
698    }
699}
700
701/// Try the user on a connection of its own, which does not retry, so a
702/// rejection costs this throwaway connection and nothing else.
703///
704/// Runs on a spawned task because the auth callback's future must be `Sync`
705/// and async-nats' connect future is not; a join handle is.
706async fn probe_user(url: &str, user: &UserCredential) -> ProbeOutcome {
707    let (url, name, password) = (url.to_string(), user.name.clone(), user.password.clone());
708    let probe = tokio::spawn(async move {
709        let attempt = async_nats::ConnectOptions::new()
710            .name(PROBE_NAME)
711            .connection_timeout(PROBE_TIMEOUT)
712            .user_and_password(name, password)
713            .connect(url);
714        match tokio::time::timeout(PROBE_TIMEOUT + Duration::from_secs(1), attempt).await {
715            Ok(Ok(_client)) => ProbeOutcome::Accepted,
716            Ok(Err(e)) if e.kind() == async_nats::ConnectErrorKind::AuthorizationViolation => {
717                ProbeOutcome::Rejected
718            }
719            Ok(Err(_)) | Err(_) => ProbeOutcome::Unreachable,
720        }
721    });
722    probe.await.unwrap_or(ProbeOutcome::Unreachable)
723}
724
725/// Probe, decide, record. `probe` is injected so the unit tests can drive
726/// every outcome.
727async fn choose<P, Fut>(
728    role: NatsRole,
729    have_token: bool,
730    live: &Live,
731    probe: P,
732) -> std::result::Result<Choice, &'static str>
733where
734    P: FnOnce() -> Fut,
735    Fut: std::future::Future<Output = ProbeOutcome>,
736{
737    let outcome = probe().await;
738    let previous = live.get();
739    let choice = select(outcome, previous, have_token).inspect_err(|_| {
740        if outcome == ProbeOutcome::Rejected {
741            live.observe(&async_nats::Event::ClientError(
742                async_nats::ClientError::Other(
743                    async_nats::ConnectErrorKind::AuthorizationViolation.to_string(),
744                ),
745            ));
746        } else {
747            *live.refusals.lock().unwrap_or_else(|e| e.into_inner()) = None;
748            debug!(
749                role = role.as_str(),
750                "NATS credential probe inconclusive; deferring this attempt"
751            );
752        }
753    })?;
754    if outcome != ProbeOutcome::Unreachable {
755        live.set(choice);
756        live_for(role).set(choice);
757    }
758    // Log what was selected and the outcome class, never a value; quiet when
759    // nothing changed so a flapping broker does not fill the log.
760    if previous == Some(choice) {
761        debug!(
762            role = role.as_str(),
763            credential = choice.label(),
764            probe = outcome.label(),
765            "NATS credential selected"
766        );
767    } else {
768        info!(
769            role = role.as_str(),
770            credential = choice.label(),
771            probe = outcome.label(),
772            "NATS credential selected"
773        );
774    }
775    Ok(choice)
776}
777
778/// Whether a connection needs per-attempt selection at all.
779enum AuthPlan {
780    /// No user: present the token (or nothing), exactly as before roles had
781    /// users. No callback, no probe.
782    Static(Option<String>),
783    /// A user is provisioned: select per attempt.
784    Select {
785        token: Option<String>,
786        user: UserCredential,
787    },
788}
789
790fn auth_plan(creds: NatsCredentials) -> AuthPlan {
791    match creds.user {
792        None => AuthPlan::Static(creds.token),
793        Some(user) => AuthPlan::Select {
794            token: creds.token,
795            user,
796        },
797    }
798}
799
800/// The `name` a kanade process announces on its NATS connection.
801///
802/// Without an identity this is `kanade-<role>`; with one it is
803/// `kanade-<role>/<identity>`. The broker echoes it back verbatim in
804/// `/connz`, which is what lets the backend attribute a connection — and
805/// therefore the credential the broker authenticated it with — to a pc_id
806/// (#1270). Nothing else on a connection carries the pc_id: the CID is
807/// assigned by the server and the IP is not a stable identifier on a fleet
808/// of laptops.
809///
810/// The name is client-supplied and therefore claimed, not proved. What
811/// `/connz` makes unforgeable is the *credential* half of the pair; a host
812/// can still lie about which pc_id it is. Under one fleet-wide token that
813/// changes nothing (every host can already impersonate every other), and
814/// closing it for good is per-agent identity, not a naming convention.
815pub fn client_name(role: NatsRole, identity: Option<&str>) -> String {
816    match identity.map(str::trim).filter(|s| !s.is_empty()) {
817        Some(id) => format!("{NAME_PREFIX}{}{NAME_SEP}{id}", role.as_str()),
818        None => format!("{NAME_PREFIX}{}", role.as_str()),
819    }
820}
821
822/// A connection name split back into its parts — see [`client_name`].
823#[derive(Debug, Clone, Copy, PartialEq, Eq)]
824pub struct ClientName<'a> {
825    /// The role segment as announced. A `&str` rather than a [`NatsRole`]
826    /// on purpose: a connection from a future (or foreign) build may name a
827    /// role this binary does not know, and dropping it on the floor would
828    /// hide exactly the host worth looking at.
829    pub role: &'a str,
830    /// The host identity, when the connection carried one. `None` for the
831    /// backend / CLI (which are not per-host) and for agents predating
832    /// #1270 — those simply cannot be attributed.
833    pub identity: Option<&'a str>,
834}
835
836/// Parse a connection name produced by [`client_name`]. `None` for any name
837/// that is not a kanade connection at all (a `nats` CLI session, a
838/// monitoring tool), which the caller should ignore rather than guess about.
839pub fn parse_client_name(name: &str) -> Option<ClientName<'_>> {
840    let rest = name.strip_prefix(NAME_PREFIX)?;
841    Some(match rest.split_once(NAME_SEP) {
842        // An empty identity (`kanade-agent/`) is not an identity.
843        Some((role, id)) if !role.is_empty() && !id.is_empty() => ClientName {
844            role,
845            identity: Some(id),
846        },
847        Some((role, _)) if !role.is_empty() => ClientName {
848            role,
849            identity: None,
850        },
851        Some(_) => return None,
852        None if !rest.is_empty() => ClientName {
853            role: rest,
854            identity: None,
855        },
856        None => return None,
857    })
858}
859
860/// Answers "is this credential the one *we* present?" without handing the
861/// credential itself to the caller.
862///
863/// #1270: the NATS monitoring endpoint reports `authorized_user` per
864/// connection. A current nats-server hides that field for
865/// token-authenticated connections, but that is the broker build's
866/// behaviour, not a guarantee this side can lean on — a consumer of
867/// `/connz` has to treat the value as possibly being the fleet-wide secret.
868/// The one question it may safely answer about it is whether it equals the
869/// credential this process already holds, and that answer is enough to
870/// label a connection ("still on the shared token") without ever storing or
871/// serving the value.
872///
873/// Constructed once and reused: [`resolve_token`] hits the Windows registry,
874/// and the caller compares against every connection on the broker.
875pub struct CredentialProbe {
876    token: Option<String>,
877    user: Option<String>,
878    /// The connection's own decision, when this probe belongs to one. `None`
879    /// for the explicit constructors, which describe a fixed credential.
880    live: Option<Arc<Live>>,
881}
882
883/// Which shape of credential a [`CredentialProbe`] holds.
884#[derive(Debug, Clone, Copy, PartialEq, Eq)]
885pub enum CredentialKind {
886    None,
887    Token,
888    /// A named user — and, when the connection presenting it is **live**,
889    /// positive proof that the broker is running `users` rather than a
890    /// token: nats-server refuses to load a config carrying both a `token`
891    /// and a `users` array ("Can not have a token and a users array") and
892    /// rejects token authentication outright once `users` are defined.
893    ///
894    /// That proof is the only thing that makes a reported `authorized_user`
895    /// safe to record verbatim. Note what is *not* proof: holding no
896    /// credential locally. A process that never authenticated at all can
897    /// still read a monitoring endpoint, and inferring the broker's mode
898    /// from a local absence would let a misconfigured host store the very
899    /// secret the rest of this type exists to protect.
900    User,
901}
902
903/// Hand-written so a stray `{:?}` in a log line cannot print the credential.
904/// A username is not a secret and is shown; a token never is.
905impl std::fmt::Debug for CredentialProbe {
906    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
907        let mut held = Vec::new();
908        if self.token.is_some() {
909            held.push("<redacted token>".to_string());
910        }
911        if let Some(name) = &self.user {
912            held.push(format!("user {name}"));
913        }
914        let rendered = if held.is_empty() {
915            "<none>".to_string()
916        } else {
917            held.join(" + ")
918        };
919        f.debug_struct("CredentialProbe")
920            .field("presented", &rendered)
921            .field("kind", &self.kind())
922            .finish()
923    }
924}
925
926impl CredentialProbe {
927    /// Resolve the credential `role` would present, exactly as [`connect`]
928    /// does, and follow the decision that role's open connection makes.
929    ///
930    /// A half-provisioned user pair is ignored here — [`connect`] is where
931    /// it fails loudly.
932    pub fn for_role(role: NatsRole) -> Self {
933        Self {
934            token: resolve_token(role),
935            user: resolve_user(role).ok().flatten().map(|u| u.name),
936            live: Some(live_for(role)),
937        }
938    }
939
940    /// Build a probe around an explicitly-supplied token.
941    ///
942    /// The production path is [`Self::for_role`]; this exists so callers can
943    /// be tested against a known credential without a Windows registry to
944    /// write to — and, on a developer's machine, without accidentally
945    /// probing the real fleet token that `for_role` would find there.
946    pub fn from_token(token: Option<String>) -> Self {
947        Self {
948            token,
949            user: None,
950            live: None,
951        }
952    }
953
954    /// Build a probe for a process that authenticates as a named user.
955    ///
956    /// Describes a fixed user credential, for testing the consumers of
957    /// [`CredentialKind::User`] without a connection.
958    pub fn from_user(name: impl Into<String>) -> Self {
959        Self {
960            token: None,
961            user: Some(name.into()),
962            live: None,
963        }
964    }
965
966    /// Which shape of credential this process is presenting **right now**.
967    ///
968    /// A role that holds a user but is on the token fallback reports
969    /// [`CredentialKind::Token`]; the user shape is reported only once the
970    /// broker has accepted it. Before any decision a role holding a token
971    /// reports the token, and one holding only a user reports `None`: the
972    /// user shape is proof of the broker's mode and is not claimed on a
973    /// guess.
974    pub fn kind(&self) -> CredentialKind {
975        let decided = self.live.as_ref().and_then(|l| l.get());
976        match decided {
977            Some(Choice::User) if self.user.is_some() => return CredentialKind::User,
978            Some(Choice::Token) if self.token.is_some() => return CredentialKind::Token,
979            _ => {}
980        }
981        if self.live.is_none() && self.user.is_some() {
982            CredentialKind::User
983        } else if self.token.is_some() {
984            CredentialKind::Token
985        } else {
986            CredentialKind::None
987        }
988    }
989
990    /// Whether `candidate` is the **secret** this process presents.
991    ///
992    /// Only ever true for a token. A user's secret is its password, which
993    /// `authorized_user` never carries — matching a username here would
994    /// mean "this connection is on the same account", a different and much
995    /// weaker statement than the one callers use this for.
996    ///
997    /// A plain comparison: `candidate` comes from the broker's own report of
998    /// connections it already authenticated, not from an attacker-chosen
999    /// input, so there is no oracle to time.
1000    pub fn is_ours(&self, candidate: &str) -> bool {
1001        self.token.as_deref() == Some(candidate)
1002    }
1003}
1004
1005/// Connect to NATS at `url` as `role`. Resolves the credential from the
1006/// registry (Windows) or the environment: a user pair when one is
1007/// provisioned (with the token as fallback, see the module docs), otherwise
1008/// the bearer token; connects unauthenticated when neither is set.
1009///
1010/// The connection is announced as `kanade-<role>` with no host identity —
1011/// right for the backend and the CLI, which are not per-host. A role that
1012/// has to be attributable to a specific machine (the agent) must use
1013/// [`connect_with_event_callback`] and pass one; see [`client_name`].
1014pub async fn connect(role: NatsRole, url: &str) -> Result<async_nats::Client> {
1015    connect_inner(
1016        role,
1017        url,
1018        None,
1019        None::<fn(async_nats::Event) -> std::future::Ready<()>>,
1020        NatsCredentials::resolve(role)?,
1021        Arc::new(Live::default()),
1022    )
1023    .await
1024}
1025
1026/// Same as [`connect`] but also wires an `event_callback` that fires
1027/// whenever async-nats publishes a `ConnectEvent` (Connected,
1028/// Disconnected, ServerError, etc.). The callback's `Future` runs on
1029/// the async-nats internal task — keep it cheap and non-blocking
1030/// (set a flag, send on a channel, that kind of thing) so the
1031/// connection state machine isn't held up.
1032///
1033/// Used by the agent's v0.26 Layer 2 staleness tracker: the callback
1034/// stamps a shared `Mutex<Option<Instant>>` on every Connected event,
1035/// so `decide()` at fire time can answer "how long ago were we last
1036/// definitely-talking-to-the-broker" without a polling loop.
1037///
1038/// `identity` names the host this connection belongs to (the agent's
1039/// pc_id). It becomes part of the connection name the broker echoes in
1040/// `/connz`, which is the only thing tying a connection — and the
1041/// credential that opened it — back to a machine (#1270).
1042pub async fn connect_with_event_callback<F, Fut>(
1043    role: NatsRole,
1044    url: &str,
1045    identity: Option<&str>,
1046    cb: F,
1047) -> Result<async_nats::Client>
1048where
1049    F: Fn(async_nats::Event) -> Fut + Send + Sync + 'static,
1050    Fut: std::future::Future<Output = ()> + Send + Sync + 'static,
1051{
1052    connect_inner(
1053        role,
1054        url,
1055        identity,
1056        Some(cb),
1057        NatsCredentials::resolve(role)?,
1058        Arc::new(Live::default()),
1059    )
1060    .await
1061}
1062
1063/// Connect with explicitly supplied credentials instead of resolving them
1064/// from the registry / environment, with refusal state private to this
1065/// connection. Credential reporting follows the latest decision for the role.
1066pub async fn connect_with_credentials(
1067    role: NatsRole,
1068    url: &str,
1069    creds: NatsCredentials,
1070) -> Result<async_nats::Client> {
1071    connect_inner(
1072        role,
1073        url,
1074        None,
1075        None::<fn(async_nats::Event) -> std::future::Ready<()>>,
1076        creds,
1077        Arc::new(Live::default()),
1078    )
1079    .await
1080}
1081
1082/// [`connect_with_credentials`] with an event callback, so a test can observe
1083/// what the connection reports while it cannot authenticate.
1084pub async fn connect_with_credentials_and_event_callback<F, Fut>(
1085    role: NatsRole,
1086    url: &str,
1087    creds: NatsCredentials,
1088    cb: F,
1089) -> Result<async_nats::Client>
1090where
1091    F: Fn(async_nats::Event) -> Fut + Send + Sync + 'static,
1092    Fut: std::future::Future<Output = ()> + Send + Sync + 'static,
1093{
1094    connect_inner(role, url, None, Some(cb), creds, Arc::new(Live::default())).await
1095}
1096
1097async fn connect_inner<F, Fut>(
1098    role: NatsRole,
1099    url: &str,
1100    identity: Option<&str>,
1101    cb: Option<F>,
1102    creds: NatsCredentials,
1103    live: Arc<Live>,
1104) -> Result<async_nats::Client>
1105where
1106    F: Fn(async_nats::Event) -> Fut + Send + Sync + 'static,
1107    Fut: std::future::Future<Output = ()> + Send + Sync + 'static,
1108{
1109    // #1187: a workspace build unifies rustls features, so any binary built
1110    // alongside a reqwest user (the CLI, the backend) links BOTH aws-lc-rs and
1111    // ring. rustls 0.23 then cannot auto-pick a process-level provider, and the
1112    // first TLS handshake — a `wss://`/`tls://` broker — panics inside
1113    // async-nats' connection task. That panic does not surface here: the
1114    // client just goes dead. Every production connect funnels through this
1115    // function, so installing ring here covers every binary. `install_default`
1116    // returns Err when a provider is already installed — ignore it.
1117    let _ = rustls::crypto::ring::default_provider().install_default();
1118
1119    let health = Arc::new(ConnectionHealth {
1120        live: live.clone(),
1121        url: url.to_string(),
1122        creds: creds.clone(),
1123    });
1124    publish_role_live(role, &live);
1125    let tracker = live.clone();
1126    // Only a provisioned user changes how the credential is presented; every
1127    // other host takes the token path untouched.
1128    let opts = match auth_plan(creds) {
1129        AuthPlan::Static(token) => {
1130            let opts = async_nats::ConnectOptions::new();
1131            match token {
1132                Some(token) => opts.token(token),
1133                None => opts,
1134            }
1135        }
1136        AuthPlan::Select { token, user } => {
1137            let url = url.to_string();
1138            async_nats::ConnectOptions::with_auth_callback(move |_nonce| {
1139                let (url, token, user, live) =
1140                    (url.clone(), token.clone(), user.clone(), live.clone());
1141                async move {
1142                    let choice = choose(role, token.is_some(), &live, || probe_user(&url, &user))
1143                        .await
1144                        .map_err(async_nats::AuthError::new)?;
1145                    let mut auth = async_nats::Auth::new();
1146                    match choice {
1147                        Choice::User => {
1148                            auth.username = Some(user.name);
1149                            auth.password = Some(user.password);
1150                        }
1151                        Choice::Token => auth.token = token,
1152                    }
1153                    Ok(auth)
1154                }
1155            })
1156        }
1157    };
1158
1159    // v0.38 / #137: offline-tolerant boot. Without
1160    // `retry_on_initial_connect`, `opts.connect(url).await` blocks-then-
1161    // errors when the broker is unreachable at startup — the agent
1162    // process dies, SCM ticks its restart counter, and the offline-
1163    // tolerant subsystems (local_scheduler, outbox drain) never spawn.
1164    // With this flag, connect() returns `Ok(Client)` immediately and
1165    // async-nats does the reconnect in the background; subscribe()
1166    // calls queue the SUB frame until the link is up.
1167    let opts = opts
1168        .retry_on_initial_connect()
1169        .connection_timeout(HEALTH_TIMEOUT)
1170        .ping_interval(PING_INTERVAL)
1171        // Names the connection in `nats server report connections`, in the
1172        // broker's own logs, and in `/connz`. Free observability while the
1173        // fleet is mid-migration: it shows which roles are connecting even
1174        // before their credentials differ, which is exactly the window in
1175        // which a wrongly-provisioned host is otherwise invisible. With an
1176        // identity it also carries the pc_id, so #1270 can join the broker's
1177        // per-connection `authorized_user` back onto the agents row.
1178        .name(client_name(role, identity));
1179    // Always installed, even without a caller callback: it is how a refused
1180    // credential is noticed, since async-nats only retries it.
1181    let opts = opts.event_callback(move |ev| {
1182        tracker.observe(&ev);
1183        let forwarded = cb.as_ref().map(|cb| cb(ev));
1184        async move {
1185            if let Some(forwarded) = forwarded {
1186                forwarded.await;
1187            }
1188        }
1189    });
1190    let client = opts
1191        .connect(url)
1192        .await
1193        .with_context(|| format!("connect to NATS at {url}"))?;
1194    register_health(&client, health);
1195    Ok(client)
1196}
1197
1198/// Whether the client's connection task has terminated.
1199///
1200/// When that task ends (a panic inside it, a drain, exhausted reconnects) the
1201/// `Client` stays alive and never talks again — a silent zombie. The
1202/// observable difference from an ordinary disconnect is the command channel:
1203/// a task that is alive but reconnecting simply does not answer a flush, so
1204/// the flush waits and the timeout elapses; a task that is gone has dropped
1205/// the receiving end, so queueing the flush fails with a send error. The
1206/// check publishes to no subject, so no role needs any publish right for it.
1207///
1208/// This tells whether the task still accepts commands, not whether the broker
1209/// is reachable. A flush that fails after it was queued (its observer was
1210/// dropped) is not treated as death: a task that really is gone fails the
1211/// next check at the send step, and mistaking a plain disconnect for death
1212/// would make a process exit needlessly.
1213pub async fn is_dead(client: &async_nats::Client) -> bool {
1214    match tokio::time::timeout(Duration::from_secs(5), client.flush()).await {
1215        Ok(Err(e)) => match e.kind() {
1216            async_nats::client::FlushErrorKind::SendError => true,
1217            async_nats::client::FlushErrorKind::FlushError => false,
1218        },
1219        Ok(Ok(())) | Err(_) => false,
1220    }
1221}
1222
1223/// Resolve once `client` can no longer be expected to work: its connection
1224/// task has terminated (see [`is_dead`]), its credential is persistently
1225/// refused, or receive progress stalls while a fresh handshake reaches the
1226/// broker (accepted or explicitly refused).
1227/// Protocol PING/PONG traffic provides progress even for idle roles. A flush
1228/// alone only proves a local write, not that the broker answered. The local
1229/// check runs every `liveness.check_interval`; a witness connection is made
1230/// only after `liveness.stall_bound` without received bytes, never on a
1231/// healthy connection, and an unreachable broker never ends the wait. A
1232/// process that depends on the client should treat this as fatal and exit
1233/// non-zero so its service manager restarts it.
1234///
1235/// `role` names the connection the client was opened with.
1236pub async fn wait_until_dead_with(role: NatsRole, client: &async_nats::Client, liveness: Liveness) {
1237    let health = connection_health(client);
1238    let mut tick = tokio::time::interval(liveness.check_interval);
1239    tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
1240    let stats = client.statistics();
1241    let mut received = stats.in_bytes.load(Ordering::Relaxed);
1242    let mut judge = StallJudge::new(liveness, tokio::time::Instant::now());
1243    loop {
1244        tick.tick().await;
1245        let now_received = stats.in_bytes.load(Ordering::Relaxed);
1246        if client.connection_state() == async_nats::connection::State::Connected
1247            && now_received != received
1248        {
1249            judge.progressed(tokio::time::Instant::now());
1250        }
1251        received = now_received;
1252        if health.as_ref().is_some_and(|h| h.live.auth_failed()) {
1253            warn!(
1254                role = role.as_str(),
1255                "the NATS broker keeps refusing this role's credential; the client cannot connect"
1256            );
1257            return;
1258        }
1259        if is_dead(client).await {
1260            warn!(
1261                role = role.as_str(),
1262                "NATS connection task has terminated; the client can no longer talk to the broker"
1263            );
1264            return;
1265        }
1266        let Some(health) = &health else { continue };
1267        if !judge.witness_due(tokio::time::Instant::now()) {
1268            continue;
1269        }
1270        let outcome = witness(health).await;
1271        if judge.witnessed(tokio::time::Instant::now(), outcome) != Verdict::Stalled {
1272            continue;
1273        }
1274        // Recheck after the witness: a reconnect racing it is healthy.
1275        if client.connection_state() == async_nats::connection::State::Connected
1276            && stats.in_bytes.load(Ordering::Relaxed) != received
1277        {
1278            judge.progressed(tokio::time::Instant::now());
1279            continue;
1280        }
1281        warn!(
1282            role = role.as_str(),
1283            witness = outcome.label(),
1284            "NATS client made no receive progress while a fresh credential witness reached the broker; exiting for a supervised restart"
1285        );
1286        return;
1287    }
1288}
1289
1290/// [`wait_until_dead_with`] at the production stall bound, checking every
1291/// `interval`.
1292pub async fn wait_until_dead_every(
1293    role: NatsRole,
1294    client: &async_nats::Client,
1295    interval: Duration,
1296) {
1297    wait_until_dead_with(
1298        role,
1299        client,
1300        Liveness {
1301            check_interval: interval,
1302            ..Liveness::PRODUCTION
1303        },
1304    )
1305    .await
1306}
1307
1308/// [`wait_until_dead_with`] at the production timing.
1309pub async fn wait_until_dead(role: NatsRole, client: &async_nats::Client) {
1310    wait_until_dead_with(role, client, Liveness::PRODUCTION).await
1311}
1312
1313/// Start supervision at `liveness` and exit the process when it resolves.
1314pub fn exit_on_dead_with(role: NatsRole, client: &async_nats::Client, liveness: Liveness) {
1315    let client = client.clone();
1316    tokio::spawn(async move {
1317        wait_until_dead_with(role, &client, liveness).await;
1318        std::process::exit(1);
1319    });
1320}
1321
1322/// Start supervision before subscriptions or bootstrap can block. Fatal
1323/// liveness failure bypasses application shutdown so a blocked subsystem
1324/// cannot keep a silent service alive beyond the exit deadline.
1325pub fn exit_on_dead(role: NatsRole, client: &async_nats::Client) {
1326    exit_on_dead_with(role, client, Liveness::PRODUCTION);
1327}
1328
1329#[cfg(test)]
1330mod tests {
1331    use super::*;
1332
1333    fn secs(n: u64) -> Duration {
1334        Duration::from_secs(n)
1335    }
1336
1337    /// Drive `judge` on the local check cadence until `until`, with receive
1338    /// progress every `progress_every` (none: silent). Returns the first time
1339    /// a witness was due and, if it ends in a stall verdict, when.
1340    fn simulate(
1341        judge: &mut StallJudge,
1342        start: tokio::time::Instant,
1343        until: Duration,
1344        progress_every: Option<Duration>,
1345        outcome: ProbeOutcome,
1346        witnesses: &mut Vec<Duration>,
1347    ) -> Option<Duration> {
1348        let step = judge.liveness.check_interval;
1349        let mut elapsed = Duration::ZERO;
1350        let mut last_progress = Duration::ZERO;
1351        while elapsed < until {
1352            elapsed += step;
1353            let now = start + elapsed;
1354            if progress_every.is_some_and(|p| elapsed - last_progress >= p) {
1355                last_progress = elapsed;
1356                judge.progressed(now);
1357            }
1358            if judge.witness_due(now) {
1359                witnesses.push(elapsed);
1360                if judge.witnessed(now, outcome) == Verdict::Stalled {
1361                    return Some(elapsed);
1362                }
1363            }
1364        }
1365        None
1366    }
1367
1368    #[test]
1369    fn the_production_bounds_leave_a_healthy_idle_connection_a_wide_margin() {
1370        assert!(STALL_BOUND >= PING_INTERVAL + secs(60));
1371        assert!(EXIT_BOUND < secs(300));
1372        assert_eq!(PING_INTERVAL, secs(60));
1373        assert_eq!(STALL_BOUND, secs(180));
1374        assert_eq!(EXIT_BOUND, secs(195));
1375    }
1376
1377    #[test]
1378    fn an_idle_healthy_client_is_never_judged_stalled() {
1379        let start = tokio::time::Instant::now();
1380        let mut judge = StallJudge::new(Liveness::PRODUCTION, start);
1381        let mut witnesses = Vec::new();
1382        // One PONG per ping interval, for a day.
1383        let verdict = simulate(
1384            &mut judge,
1385            start,
1386            secs(24 * 3600),
1387            Some(PING_INTERVAL),
1388            ProbeOutcome::Accepted,
1389            &mut witnesses,
1390        );
1391        assert_eq!(verdict, None);
1392        assert!(
1393            witnesses.is_empty(),
1394            "a witness connected to a healthy broker"
1395        );
1396    }
1397
1398    #[test]
1399    fn an_offline_broker_never_exits_and_is_polled_at_the_resume_interval() {
1400        let start = tokio::time::Instant::now();
1401        let mut judge = StallJudge::new(Liveness::PRODUCTION, start);
1402        let mut witnesses = Vec::new();
1403        let verdict = simulate(
1404            &mut judge,
1405            start,
1406            secs(3 * 3600),
1407            None,
1408            ProbeOutcome::Unreachable,
1409            &mut witnesses,
1410        );
1411        assert_eq!(verdict, None);
1412        assert_eq!(witnesses[0], STALL_BOUND);
1413        assert!(witnesses.windows(2).all(|w| w[1] - w[0] >= RESUME_BOUND));
1414    }
1415
1416    #[test]
1417    fn a_stall_after_an_outage_is_judged_against_the_normal_bound() {
1418        let start = tokio::time::Instant::now();
1419        let mut judge = StallJudge::new(Liveness::PRODUCTION, start);
1420        let mut witnesses = Vec::new();
1421        assert_eq!(
1422            simulate(
1423                &mut judge,
1424                start,
1425                secs(3600),
1426                None,
1427                ProbeOutcome::Unreachable,
1428                &mut witnesses
1429            ),
1430            None
1431        );
1432        assert!(judge.outage);
1433        // The broker answers but the client has not resumed: only a grace.
1434        let back = start + secs(3600);
1435        witnesses.clear();
1436        let exit = simulate(
1437            &mut judge,
1438            back,
1439            EXIT_BOUND,
1440            None,
1441            ProbeOutcome::Accepted,
1442            &mut witnesses,
1443        );
1444        assert!(!judge.outage);
1445        let exit = exit.expect("a silent client with a reachable broker must exit");
1446        assert!(exit <= EXIT_BOUND, "{exit:?}");
1447
1448        // Progress after the outage clears it; a later genuine stall is then
1449        // judged against the full stall bound from the last progress.
1450        let mut judge = StallJudge::new(Liveness::PRODUCTION, start);
1451        let mut witnesses = Vec::new();
1452        simulate(
1453            &mut judge,
1454            start,
1455            secs(3600),
1456            None,
1457            ProbeOutcome::Unreachable,
1458            &mut witnesses,
1459        );
1460        assert!(judge.outage);
1461        let resumed = start + secs(3600);
1462        judge.progressed(resumed);
1463        assert!(!judge.outage);
1464        assert!(!judge.witness_due(resumed + STALL_BOUND - secs(1)));
1465        witnesses.clear();
1466        let exit = simulate(
1467            &mut judge,
1468            resumed,
1469            STALL_BOUND + secs(60),
1470            None,
1471            ProbeOutcome::Rejected,
1472            &mut witnesses,
1473        );
1474        assert_eq!(exit, Some(STALL_BOUND));
1475    }
1476    use std::collections::HashMap;
1477
1478    /// A stand-in registry. Keys are `subkey\value`.
1479    fn reg(entries: &[(&str, &str)]) -> impl Fn(&str, &str) -> Option<String> {
1480        let map: HashMap<String, String> = entries
1481            .iter()
1482            .map(|(k, v)| ((*k).to_string(), (*v).to_string()))
1483            .collect();
1484        move |subkey: &str, value: &str| map.get(&format!(r"{subkey}\{value}")).cloned()
1485    }
1486
1487    #[test]
1488    fn role_subkeys_are_distinct_and_agent_matches_the_shared_path() {
1489        assert_eq!(NatsRole::Backend.reg_subkey(), r"SOFTWARE\kanade\backend");
1490        assert_eq!(NatsRole::Cli.reg_subkey(), r"SOFTWARE\kanade\cli");
1491        // The agent's role key IS the historical shared key, so an agent
1492        // never sees a migration at all.
1493        assert_eq!(NatsRole::Agent.reg_subkey(), REG_SHARED_SUBKEY);
1494    }
1495
1496    #[test]
1497    fn an_unmigrated_fleet_keeps_presenting_the_shared_token() {
1498        // The state of every host today: one token, under the agent path.
1499        let registry = reg(&[(r"SOFTWARE\kanade\agent\NatsToken", "shared")]);
1500        for role in [NatsRole::Agent, NatsRole::Backend, NatsRole::Cli] {
1501            assert_eq!(
1502                resolve_token_with(role, &registry, None).as_deref(),
1503                Some("shared"),
1504                "{role:?} must keep working before its own key is provisioned"
1505            );
1506        }
1507    }
1508
1509    #[test]
1510    fn a_role_key_wins_over_the_shared_one() {
1511        let registry = reg(&[
1512            (r"SOFTWARE\kanade\agent\NatsToken", "shared"),
1513            (r"SOFTWARE\kanade\backend\NatsToken", "backend-only"),
1514        ]);
1515        // The migrated role uses its own credential...
1516        assert_eq!(
1517            resolve_token_with(NatsRole::Backend, &registry, None).as_deref(),
1518            Some("backend-only")
1519        );
1520        // ...while a role that has not been migrated yet is unaffected.
1521        assert_eq!(
1522            resolve_token_with(NatsRole::Cli, &registry, None).as_deref(),
1523            Some("shared")
1524        );
1525    }
1526
1527    #[test]
1528    fn removing_a_role_key_falls_back_rather_than_failing() {
1529        // Rollback of a per-host migration step: the role key is gone, and
1530        // the host must return to the shared credential instead of
1531        // connecting unauthenticated (which a broker with `authorization`
1532        // would refuse — turning a rollback into an outage).
1533        let registry = reg(&[(r"SOFTWARE\kanade\agent\NatsToken", "shared")]);
1534        assert_eq!(
1535            resolve_token_with(NatsRole::Backend, &registry, None).as_deref(),
1536            Some("shared")
1537        );
1538    }
1539
1540    #[test]
1541    fn the_registry_outranks_the_environment() {
1542        // Unchanged from before roles existed: a dev shell's env var must
1543        // not quietly override a provisioned production credential.
1544        let registry = reg(&[(r"SOFTWARE\kanade\agent\NatsToken", "shared")]);
1545        assert_eq!(
1546            resolve_token_with(NatsRole::Agent, &registry, Some("from-env".into())).as_deref(),
1547            Some("shared")
1548        );
1549    }
1550
1551    #[test]
1552    fn the_environment_serves_hosts_with_no_registry_at_all() {
1553        let empty = reg(&[]);
1554        assert_eq!(
1555            resolve_token_with(NatsRole::Cli, &empty, Some("from-env".into())).as_deref(),
1556            Some("from-env")
1557        );
1558        // An empty env var is not a credential — it must fall through to
1559        // "no token" so a dev broker without `authorization` still works,
1560        // rather than presenting the empty string and being rejected.
1561        assert_eq!(
1562            resolve_token_with(NatsRole::Cli, &empty, Some(String::new())),
1563            None
1564        );
1565        assert_eq!(resolve_token_with(NatsRole::Cli, &empty, None), None);
1566    }
1567
1568    // ── #1270: connection naming ─────────────────────────────────────
1569
1570    #[test]
1571    fn an_identity_round_trips_through_the_connection_name() {
1572        // The pc_id is the join key between `/connz` and the agents table,
1573        // so the name has to survive the trip unchanged — including the
1574        // casing, which is NOT uniform across the fleet and which NATS
1575        // subjects treat as significant.
1576        for pc in ["PC001", "minipc", "Web%01", "ws-9"] {
1577            let name = client_name(NatsRole::Agent, Some(pc));
1578            let parsed = parse_client_name(&name).expect("our own name must parse");
1579            assert_eq!(parsed.role, "agent");
1580            assert_eq!(parsed.identity, Some(pc));
1581        }
1582    }
1583
1584    #[test]
1585    fn a_role_without_an_identity_keeps_the_pre_1270_name() {
1586        // The backend and the CLI are not per-host, and an agent that
1587        // predates #1270 announces this shape too. Both must parse as "a
1588        // kanade connection we cannot attribute" rather than as an error or
1589        // as an empty pc_id.
1590        assert_eq!(client_name(NatsRole::Backend, None), "kanade-backend");
1591        let parsed = parse_client_name("kanade-agent").unwrap();
1592        assert_eq!(parsed.role, "agent");
1593        assert_eq!(parsed.identity, None);
1594        // Whitespace-only is not an identity either — it would otherwise
1595        // produce a name that parses back into a pc_id no row can match.
1596        assert_eq!(client_name(NatsRole::Agent, Some("  ")), "kanade-agent");
1597    }
1598
1599    #[test]
1600    fn foreign_connections_do_not_parse_as_kanade_ones() {
1601        // A `nats` CLI session or a monitoring tool shares the broker. The
1602        // projector must skip those rather than attribute them to a host.
1603        assert!(parse_client_name("NATS CLI Version 0.1.5").is_none());
1604        assert!(parse_client_name("").is_none());
1605        assert!(parse_client_name("kanade-").is_none());
1606        assert!(parse_client_name("kanade-/PC001").is_none());
1607        // A trailing separator with no identity is a role, not a pc_id.
1608        assert_eq!(parse_client_name("kanade-agent/").unwrap().identity, None);
1609    }
1610
1611    #[test]
1612    fn an_unknown_role_is_preserved_rather_than_dropped() {
1613        // A future build (or something impersonating one) naming a role this
1614        // binary has never heard of is precisely the connection an operator
1615        // wants to see.
1616        let parsed = parse_client_name("kanade-relay/PC001").unwrap();
1617        assert_eq!(parsed.role, "relay");
1618        assert_eq!(parsed.identity, Some("PC001"));
1619    }
1620
1621    // ── #1270: credential probe ──────────────────────────────────────
1622
1623    #[test]
1624    fn the_probe_recognises_only_the_credential_we_present() {
1625        let probe = CredentialProbe::from_token(Some("shared".into()));
1626        assert_eq!(probe.kind(), CredentialKind::Token);
1627        assert!(probe.is_ours("shared"));
1628        assert!(!probe.is_ours("something-else"));
1629        // The empty string is what the broker reports for a connection it
1630        // did not authenticate at all. It must never read as "ours".
1631        assert!(!probe.is_ours(""));
1632    }
1633
1634    #[test]
1635    fn a_probe_with_no_credential_matches_nothing() {
1636        // Dev broker with no `authorization` block. We hold nothing, so we
1637        // can prove nothing about anyone else's credential — including that
1638        // it is safe to store.
1639        let probe = CredentialProbe::from_token(None);
1640        assert_eq!(probe.kind(), CredentialKind::None);
1641        assert!(!probe.is_ours(""));
1642        assert!(!probe.is_ours("anything"));
1643    }
1644
1645    #[test]
1646    fn a_username_is_not_a_secret_we_can_recognise() {
1647        // `is_ours` answers "is this MY secret", and a user's secret is its
1648        // password. Matching on the username instead would answer a much
1649        // weaker question while reading like the strong one.
1650        let probe = CredentialProbe::from_user("kanade-backend");
1651        assert_eq!(probe.kind(), CredentialKind::User);
1652        assert!(!probe.is_ours("kanade-backend"));
1653    }
1654
1655    #[test]
1656    fn the_probe_never_prints_the_credential() {
1657        // `/connz` handling logs liberally; one `{:?}` on the probe must not
1658        // be the thing that puts the fleet's token in a log file.
1659        let probe = CredentialProbe::from_token(Some("super-secret-token".into()));
1660        let rendered = format!("{probe:?}");
1661        assert!(!rendered.contains("super-secret-token"), "{rendered}");
1662        assert!(rendered.contains("redacted"), "{rendered}");
1663        assert!(
1664            format!("{:?}", CredentialProbe::from_token(None)).contains("none"),
1665            "the no-credential case should be visible, just not the value"
1666        );
1667        // A username is not a secret — hiding it would cost diagnosability
1668        // for nothing.
1669        assert!(
1670            format!("{:?}", CredentialProbe::from_user("kanade-backend"))
1671                .contains("kanade-backend"),
1672        );
1673    }
1674
1675    // ── per-role users ───────────────────────────────────────────────
1676
1677    fn user(r: Result<Option<UserCredential>>) -> Option<(String, String)> {
1678        r.unwrap().map(|u| (u.name, u.password))
1679    }
1680
1681    #[test]
1682    fn a_users_registry_pair_comes_from_the_roles_own_key() {
1683        let registry = reg(&[
1684            (r"SOFTWARE\kanade\backend\NatsUser", "kanade-backend"),
1685            (r"SOFTWARE\kanade\backend\NatsPassword", "pw"),
1686        ]);
1687        assert_eq!(
1688            user(resolve_user_with(NatsRole::Backend, &registry, None, None)),
1689            Some(("kanade-backend".into(), "pw".into()))
1690        );
1691        // No shared user: another role does not inherit it, not even the
1692        // agent's key (which is where the shared *token* lives).
1693        for role in [NatsRole::Agent, NatsRole::Cli] {
1694            assert_eq!(user(resolve_user_with(role, &registry, None, None)), None);
1695        }
1696        let agent = reg(&[
1697            (r"SOFTWARE\kanade\agent\NatsUser", "a"),
1698            (r"SOFTWARE\kanade\agent\NatsPassword", "b"),
1699        ]);
1700        assert_eq!(
1701            user(resolve_user_with(NatsRole::Cli, &agent, None, None)),
1702            None
1703        );
1704    }
1705
1706    #[test]
1707    fn the_registry_pair_outranks_the_environment_pair() {
1708        let registry = reg(&[
1709            (r"SOFTWARE\kanade\cli\NatsUser", "reg-user"),
1710            (r"SOFTWARE\kanade\cli\NatsPassword", "reg-pw"),
1711        ]);
1712        assert_eq!(
1713            user(resolve_user_with(
1714                NatsRole::Cli,
1715                &registry,
1716                Some("env-user".into()),
1717                Some("env-pw".into())
1718            )),
1719            Some(("reg-user".into(), "reg-pw".into()))
1720        );
1721    }
1722
1723    #[test]
1724    fn the_environment_serves_a_user_pair_when_the_registry_has_none() {
1725        assert_eq!(
1726            user(resolve_user_with(
1727                NatsRole::Cli,
1728                reg(&[]),
1729                Some("u".into()),
1730                Some("p".into())
1731            )),
1732            Some(("u".into(), "p".into()))
1733        );
1734        // Empty values are not a credential.
1735        assert_eq!(
1736            user(resolve_user_with(
1737                NatsRole::Cli,
1738                reg(&[]),
1739                Some(String::new()),
1740                Some(String::new())
1741            )),
1742            None
1743        );
1744        assert_eq!(
1745            user(resolve_user_with(NatsRole::Cli, reg(&[]), None, None)),
1746            None
1747        );
1748    }
1749
1750    #[test]
1751    fn half_a_pair_is_an_error_not_a_fallback() {
1752        let err = |r: Result<Option<UserCredential>>| format!("{:#}", r.unwrap_err());
1753        // Registry halves.
1754        let only_user = reg(&[(r"SOFTWARE\kanade\cli\NatsUser", "u")]);
1755        assert!(
1756            err(resolve_user_with(NatsRole::Cli, &only_user, None, None))
1757                .contains("password is missing")
1758        );
1759        let only_pw = reg(&[(r"SOFTWARE\kanade\cli\NatsPassword", "secret-pw")]);
1760        let e = err(resolve_user_with(NatsRole::Cli, &only_pw, None, None));
1761        assert!(e.contains("user name is missing"), "{e}");
1762        assert!(!e.contains("secret-pw"), "{e}");
1763        // A half registry pair is not rescued by a complete environment pair.
1764        assert!(
1765            resolve_user_with(
1766                NatsRole::Cli,
1767                &only_user,
1768                Some("u".into()),
1769                Some("p".into())
1770            )
1771            .is_err()
1772        );
1773        // Environment halves.
1774        assert!(resolve_user_with(NatsRole::Cli, reg(&[]), Some("u".into()), None).is_err());
1775        assert!(resolve_user_with(NatsRole::Cli, reg(&[]), None, Some("p".into())).is_err());
1776    }
1777
1778    #[test]
1779    fn credentials_never_print_a_secret() {
1780        let creds = NatsCredentials::new(
1781            Some("tok-secret".into()),
1782            Some(("kanade-agent".into(), "pw-secret".into())),
1783        );
1784        let rendered = format!("{creds:?}");
1785        assert!(!rendered.contains("tok-secret"), "{rendered}");
1786        assert!(!rendered.contains("pw-secret"), "{rendered}");
1787        assert!(rendered.contains("kanade-agent"), "{rendered}");
1788    }
1789
1790    #[test]
1791    fn no_user_means_the_static_token_path_with_no_selection() {
1792        match auth_plan(NatsCredentials::new(Some("t".into()), None)) {
1793            AuthPlan::Static(Some(t)) => assert_eq!(t, "t"),
1794            _ => panic!("a token-only host must take the static path"),
1795        }
1796        assert!(matches!(
1797            auth_plan(NatsCredentials::new(None, None)),
1798            AuthPlan::Static(None)
1799        ));
1800        assert!(matches!(
1801            auth_plan(NatsCredentials::new(
1802                Some("t".into()),
1803                Some(("u".into(), "p".into()))
1804            )),
1805            AuthPlan::Select { .. }
1806        ));
1807    }
1808
1809    #[test]
1810    fn selection_follows_the_probe() {
1811        use Choice::{Token, User};
1812        use ProbeOutcome::*;
1813        assert_eq!(select(Accepted, None, true), Ok(User));
1814        assert_eq!(select(Accepted, Some(Token), false), Ok(User));
1815        assert_eq!(select(Rejected, Some(User), true), Ok(Token));
1816        let missing = select(Rejected, None, false).unwrap_err();
1817        assert!(missing.contains("NatsToken"), "{missing}");
1818        // An inconclusive probe defers even when a previous mode is cached.
1819        assert!(select(Unreachable, None, true).is_err());
1820        assert!(select(Unreachable, Some(Token), true).is_err());
1821        assert!(select(Unreachable, Some(User), true).is_err());
1822    }
1823
1824    #[tokio::test]
1825    async fn only_a_probe_answer_is_remembered() {
1826        let live = Live::default();
1827        let role = NatsRole::Cli;
1828        let run = |o| choose(role, true, &live, move || async move { o });
1829        // An unreachable broker cannot justify a handshake.
1830        assert!(run(ProbeOutcome::Unreachable).await.is_err());
1831        assert_eq!(live.get(), None);
1832        assert_eq!(run(ProbeOutcome::Rejected).await, Ok(Choice::Token));
1833        assert_eq!(live.get(), Some(Choice::Token));
1834        // The broker goes away: its previous mode is no longer evidence.
1835        assert!(run(ProbeOutcome::Unreachable).await.is_err());
1836        // ...and the flip back is followed.
1837        assert_eq!(run(ProbeOutcome::Accepted).await, Ok(Choice::User));
1838        assert_eq!(live.get(), Some(Choice::User));
1839        // A rejection with no token fails and leaves the cache alone.
1840        let err = choose(role, false, &live, || async { ProbeOutcome::Rejected })
1841            .await
1842            .unwrap_err();
1843        assert!(err.contains("NatsToken"));
1844        assert_eq!(live.get(), Some(Choice::User));
1845    }
1846
1847    #[test]
1848    fn the_probe_reports_the_shape_in_use_right_now() {
1849        let live = Arc::new(Live::default());
1850        let probe = CredentialProbe {
1851            token: Some("tok".into()),
1852            user: Some("kanade-agent".into()),
1853            live: Some(live.clone()),
1854        };
1855        // Nothing proven yet: not the user shape.
1856        assert_eq!(probe.kind(), CredentialKind::Token);
1857        live.set(Choice::User);
1858        assert_eq!(probe.kind(), CredentialKind::User);
1859        // Holding a user while on the token fallback is the token shape.
1860        live.set(Choice::Token);
1861        assert_eq!(probe.kind(), CredentialKind::Token);
1862        assert!(probe.is_ours("tok"));
1863
1864        let user_only = CredentialProbe {
1865            token: None,
1866            user: Some("u".into()),
1867            live: Some(Arc::new(Live::default())),
1868        };
1869        assert_eq!(user_only.kind(), CredentialKind::None);
1870        user_only.live.as_ref().unwrap().set(Choice::User);
1871        assert_eq!(user_only.kind(), CredentialKind::User);
1872    }
1873
1874    #[test]
1875    fn a_probe_holding_both_credentials_prints_neither_secret() {
1876        let probe = CredentialProbe {
1877            token: Some("tok-secret".into()),
1878            user: Some("kanade-agent".into()),
1879            live: None,
1880        };
1881        let rendered = format!("{probe:?}");
1882        assert!(!rendered.contains("tok-secret"), "{rendered}");
1883        assert!(rendered.contains("kanade-agent"), "{rendered}");
1884    }
1885
1886    #[test]
1887    fn sustained_refusals_fail_the_client_and_a_connect_clears_them() {
1888        let refused = || {
1889            async_nats::Event::ClientError(async_nats::ClientError::Other(
1890                async_nats::ConnectErrorKind::AuthorizationViolation.to_string(),
1891            ))
1892        };
1893        let live = Live::default();
1894        for _ in 0..AUTH_REJECTION_LIMIT {
1895            live.observe(&refused());
1896        }
1897        // Enough refusals, but not over a long enough window: a broker
1898        // mid-reload must not count.
1899        assert!(!live.auth_failed());
1900        // Age the first refusal past the window.
1901        let aged = std::time::Instant::now() - AUTH_REJECTION_WINDOW - Duration::from_secs(1);
1902        *live.refusals.lock().unwrap() = Some((aged, AUTH_REJECTION_LIMIT));
1903        assert!(live.auth_failed());
1904        live.observe(&async_nats::Event::Connected);
1905        assert!(!live.auth_failed());
1906        // A different kind of failure (broker gone) ends the run, so refusals
1907        // from before an outage cannot add up to a failure after it.
1908        *live.refusals.lock().unwrap() = Some((aged, AUTH_REJECTION_LIMIT));
1909        live.observe(&async_nats::Event::ClientError(
1910            async_nats::ClientError::Other("io".into()),
1911        ));
1912        assert!(!live.auth_failed());
1913        live.observe(&refused());
1914        assert!(!live.auth_failed());
1915    }
1916
1917    #[tokio::test]
1918    async fn deferred_probes_do_not_count_as_credential_refusals() {
1919        let live = Live::default();
1920        for _ in 0..AUTH_REJECTION_LIMIT {
1921            assert!(
1922                choose(NatsRole::Cli, true, &live, || async {
1923                    ProbeOutcome::Unreachable
1924                })
1925                .await
1926                .is_err()
1927            );
1928            live.observe(&async_nats::Event::ClientError(
1929                async_nats::ClientError::Other(
1930                    async_nats::ConnectErrorKind::Authentication.to_string(),
1931                ),
1932            ));
1933        }
1934        assert!(live.refusals.lock().unwrap().is_none());
1935        for _ in 0..AUTH_REJECTION_LIMIT {
1936            assert!(
1937                choose(NatsRole::Cli, false, &live, || async {
1938                    ProbeOutcome::Rejected
1939                })
1940                .await
1941                .is_err()
1942            );
1943            live.observe(&async_nats::Event::ClientError(
1944                async_nats::ClientError::Other(
1945                    async_nats::ConnectErrorKind::Authentication.to_string(),
1946                ),
1947            ));
1948        }
1949        assert_eq!(
1950            live.refusals.lock().unwrap().as_ref().unwrap().1,
1951            AUTH_REJECTION_LIMIT
1952        );
1953        assert!(!live.auth_failed(), "quick refusals must still allow retry");
1954        assert!(
1955            choose(NatsRole::Cli, true, &live, || async {
1956                ProbeOutcome::Unreachable
1957            })
1958            .await
1959            .is_err()
1960        );
1961        assert!(
1962            live.refusals.lock().unwrap().is_none(),
1963            "an outage interrupts the refusal run"
1964        );
1965    }
1966
1967    #[tokio::test]
1968    async fn same_role_connections_have_independent_health_state() {
1969        let first = connect_with_credentials(
1970            NatsRole::Agent,
1971            "nats://127.0.0.1:1",
1972            NatsCredentials::default(),
1973        )
1974        .await
1975        .unwrap();
1976        let second = connect_with_credentials(
1977            NatsRole::Agent,
1978            "nats://127.0.0.1:1",
1979            NatsCredentials::default(),
1980        )
1981        .await
1982        .unwrap();
1983        let first_health = connection_health(&first).unwrap();
1984        let second_health = connection_health(&second).unwrap();
1985        assert!(!Arc::ptr_eq(&first_health.live, &second_health.live));
1986        assert!(Arc::ptr_eq(
1987            &first_health,
1988            &connection_health(&first.clone()).unwrap()
1989        ));
1990        *first_health.live.refusals.lock().unwrap() = Some((
1991            std::time::Instant::now() - AUTH_REJECTION_WINDOW,
1992            AUTH_REJECTION_LIMIT,
1993        ));
1994        assert!(first_health.live.auth_failed());
1995        assert!(!second_health.live.auth_failed());
1996    }
1997
1998    #[test]
1999    fn connect_leaves_a_process_crypto_provider_installed() {
2000        // #1187: a binary that links both aws-lc-rs and ring (any workspace
2001        // build) panics on its first wss handshake unless a process-level
2002        // provider was installed first. That panic happens inside async-nats'
2003        // background task, so the only thing a caller ever sees is a client
2004        // that silently never works — pin the precondition here instead.
2005        // `retry_on_initial_connect` makes connect return without a broker.
2006        let rt = tokio::runtime::Builder::new_current_thread()
2007            .enable_all()
2008            .build()
2009            .unwrap();
2010        rt.block_on(async {
2011            connect(NatsRole::Agent, "nats://127.0.0.1:1")
2012                .await
2013                .unwrap();
2014        });
2015        assert!(rustls::crypto::CryptoProvider::get_default().is_some());
2016    }
2017}