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) → reuse whichever credential
46//!     last worked, defaulting to the user.
47//!
48//! A wrong guess is thus survivable and the process keeps one `Client`.
49//! A client whose connection task has terminated for any reason (a panic in
50//! it, a drain, a version that treats a violation as terminal) is detected by
51//! [`wait_until_dead`] so its process can exit for a supervised restart. The
52//! same call also reports a broker that refuses the credential on every
53//! attempt for a sustained period, since async-nats would otherwise retry it
54//! forever and leave a process that looks alive but can never talk.
55//!
56//! # Why roles exist here (#1155)
57//!
58//! The broker authorises a *connection*, and a connection is only as
59//! specific as the credential that opened it. While every binary presented
60//! the same token, the broker could not tell an agent from the backend from
61//! the CLI, so no `permissions` block could say "only the backend may
62//! subscribe `remote.frame.>`" — there was nothing to hang the rule on.
63//! That is why a shared token means a token holder can execute code on any
64//! endpoint and silently watch any remote-assistance session (#1140).
65//!
66//! Distinct credentials do not fix that by themselves; the broker config has
67//! to grow the matching `authorization { users: [...] }` entries. This
68//! module is the half that makes those entries *expressible*.
69//!
70//! # Staged migration
71//!
72//! Step 2 above is the whole migration strategy. A fleet running today has
73//! one token, provisioned at `…\kanade\agent\NatsToken` on every host
74//! regardless of role. After this change it keeps working: no role key
75//! exists, so every role falls through to the shared one and presents
76//! exactly what it presented before.
77//!
78//! Rolling out per-role credentials is then per-host and reversible — write
79//! `…\kanade\backend\NatsToken` on the backend host and it starts using it;
80//! delete it and it falls back. The broker only needs to start
81//! *distinguishing* the roles once every host has its own, so the config
82//! change lands last, when it can no longer lock anyone out.
83//!
84//! No deploy script writes a role key yet: `deploy-backend.ps1` still
85//! provisions the shared path, so today the role key is a manual registry
86//! write. That is deliberate — the scripted path should start writing role
87//! keys in the same change that teaches the broker to tell the roles apart,
88//! because until then a role key changes nothing and a script that writes
89//! only the role key (dropping the shared one) would strand the CLI on a
90//! backend-only host.
91//!
92//! The order matters and is deliberate: role key first, shared key second.
93//! The reverse would make the shared token permanent — a host that still has
94//! it (all of them, today) would never notice its role key.
95//!
96//! # What the broker will and will not accept (measured, #1270)
97//!
98//! Two nats-server behaviours constrain every plan built on this module, so
99//! they are recorded here rather than rediscovered:
100//!
101//! * A config may not carry **both** a `token` and a `users` array —
102//!   nats-server refuses to start: *"Can not have a token and a users
103//!   array"*. And once `users` are defined, a client presenting a token is
104//!   rejected with an Authorization Violation, even when the token equals a
105//!   user's password. So the shared token and a per-role `users` split
106//!   cannot coexist for a transition window: the flip is atomic, and every
107//!   host must already hold a credential of the new shape before it happens.
108//!   Resolving a *token* per role, which is all this module does today, is
109//!   therefore not sufficient for that split — the client has to learn to
110//!   present a user as well.
111//! * `/connz?auth=1` reports `authorized_user` per connection. Under
112//!   `users` that is the username — the per-host answer #1270 wants. Under
113//!   `token`, nats-server 2.14.3 reports the literal `[REDACTED]`: it hides
114//!   the credential, so token mode can say *that* a host is on the shared
115//!   token but never anything finer. Whether the value is hidden is the
116//!   broker build's choice, not ours, so a consumer must assume it may be
117//!   handling a secret; see [`CredentialProbe`] for the one question it can
118//!   safely ask about one.
119//!
120//! # Limits worth naming
121//!
122//! A per-role token still cannot express per-*agent* identity. A role
123//! credential permitted to subscribe `commands.pc.*` lets any agent holding
124//! it read another agent's inbox. This narrows a fleet-wide compromise to a
125//! fleet-wide **agent-role** compromise, which is better, not solved. The
126//! end state is per-agent identity (NKeys / NATS-JWT), for which the plan is
127//! to grow `ConnectOptions` here so every binary picks up the upgrade for
128//! free. Same for mTLS.
129
130use std::collections::HashMap;
131use std::sync::{Arc, Mutex, OnceLock};
132use std::time::Duration;
133
134use anyhow::{Context, Result, bail};
135use tracing::{debug, info, warn};
136
137use crate::secrets;
138
139const ENV_TOKEN: &str = "KANADE_NATS_TOKEN";
140const REG_VALUE: &str = "NatsToken";
141const ENV_USER: &str = "KANADE_NATS_USER";
142const ENV_PASSWORD: &str = "KANADE_NATS_PASSWORD";
143const REG_USER: &str = "NatsUser";
144const REG_PASSWORD: &str = "NatsPassword";
145
146/// How long the credential probe may take. Short on purpose: it runs inside
147/// the reconnect path, so a black-holed broker must not stall every attempt.
148const PROBE_TIMEOUT: Duration = Duration::from_secs(3);
149
150/// Name of the probe connection. Deliberately *not* a `kanade-` name: the
151/// backend attributes `kanade-<role>` connections to hosts, and a transient
152/// probe must not show up there as an unknown role.
153const PROBE_NAME: &str = "auth-probe";
154
155/// How often [`wait_until_dead`] checks the connection task is still there.
156const DEAD_CHECK_INTERVAL: Duration = Duration::from_secs(15);
157
158/// A connection that has been refused this many times in a row, over at
159/// least [`AUTH_REJECTION_WINDOW`], with no successful connect in between, is
160/// treated as failed. async-nats retries a refused credential forever, which
161/// leaves a process that can never talk but looks alive; the window keeps a
162/// broker that is merely mid-reload (a few quick refusals) from counting.
163const AUTH_REJECTION_LIMIT: usize = 5;
164const AUTH_REJECTION_WINDOW: Duration = Duration::from_secs(10);
165
166/// Prefix every kanade connection announces in its `name`.
167const NAME_PREFIX: &str = "kanade-";
168
169/// Separator between the role and the host identity in a connection name.
170/// `/` is safe as a delimiter because the identity is a Windows computer
171/// name, which cannot contain one.
172const NAME_SEP: char = '/';
173
174/// Registry subkey holding the pre-#1155 shared credential. Also the agent's
175/// role key, which is not a coincidence — the shared token was provisioned
176/// under the agent's path because agents were the first thing to need it.
177const REG_SHARED_SUBKEY: &str = r"SOFTWARE\kanade\agent";
178
179/// Which kanade binary is opening the connection.
180///
181/// Named on every call rather than inferred, because the broker's view of a
182/// connection comes entirely from the credential it presents: a caller that
183/// picks the wrong role does not get a warning, it gets someone else's
184/// permissions.
185#[derive(Debug, Clone, Copy, PartialEq, Eq)]
186pub enum NatsRole {
187    /// The endpoint agent. The most numerous and least trusted role — one
188    /// compromised endpoint holds this credential.
189    Agent,
190    /// The backend. The only role that needs to see the whole fleet.
191    Backend,
192    /// The operator CLI, including the backend-down recovery path that
193    /// drives agents over NATS directly.
194    Cli,
195}
196
197impl NatsRole {
198    pub fn as_str(self) -> &'static str {
199        match self {
200            NatsRole::Agent => "agent",
201            NatsRole::Backend => "backend",
202            NatsRole::Cli => "cli",
203        }
204    }
205
206    /// Registry subkey holding this role's credential.
207    fn reg_subkey(self) -> String {
208        format!(r"SOFTWARE\kanade\{}", self.as_str())
209    }
210}
211
212/// Resolve a role's token, given a registry reader and the environment
213/// fallback.
214///
215/// Split from [`resolve_token`] so the *ordering* — the part that decides
216/// whether a migration is reversible — is testable without a Windows
217/// registry to write to.
218fn resolve_token_with(
219    role: NatsRole,
220    read_reg: impl Fn(&str, &str) -> Option<String>,
221    env: Option<String>,
222) -> Option<String> {
223    if let Some(t) = read_reg(&role.reg_subkey(), REG_VALUE) {
224        return Some(t);
225    }
226    if let Some(t) = read_reg(REG_SHARED_SUBKEY, REG_VALUE) {
227        return Some(t);
228    }
229    env.filter(|t| !t.is_empty())
230}
231
232fn resolve_token(role: NatsRole) -> Option<String> {
233    resolve_token_with(
234        role,
235        secrets::read_hklm_value,
236        std::env::var(ENV_TOKEN).ok(),
237    )
238}
239
240/// A named NATS user. The password is never shown by `Debug`.
241#[derive(Clone, PartialEq, Eq)]
242struct UserCredential {
243    name: String,
244    password: String,
245}
246
247impl std::fmt::Debug for UserCredential {
248    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
249        f.debug_struct("UserCredential")
250            .field("name", &self.name)
251            .field("password", &"<redacted>")
252            .finish()
253    }
254}
255
256/// Turn the two halves of a user pair into a credential, or a configuration
257/// error naming what is missing and where (never a value).
258fn pair(
259    name: Option<String>,
260    password: Option<String>,
261    source: &str,
262) -> Result<Option<UserCredential>> {
263    match (name, password) {
264        (Some(name), Some(password)) => Ok(Some(UserCredential { name, password })),
265        (None, None) => Ok(None),
266        (Some(_), None) => bail!("NATS user is set but its password is missing ({source})"),
267        (None, Some(_)) => bail!("NATS password is set but its user name is missing ({source})"),
268    }
269}
270
271/// Resolve a role's user pair: its own registry key first, then the
272/// environment. Unlike the token there is no shared-key step.
273///
274/// The registry is judged as a whole before the environment is consulted: a
275/// half pair there is an error even when the environment is complete,
276/// because quietly mixing sources would hide a half-provisioned host.
277fn resolve_user_with(
278    role: NatsRole,
279    read_reg: impl Fn(&str, &str) -> Option<String>,
280    env_user: Option<String>,
281    env_password: Option<String>,
282) -> Result<Option<UserCredential>> {
283    let subkey = role.reg_subkey();
284    let reg_user = read_reg(&subkey, REG_USER);
285    let reg_password = read_reg(&subkey, REG_PASSWORD);
286    if reg_user.is_some() || reg_password.is_some() {
287        let source = format!(r"HKLM\{subkey}: {REG_USER} / {REG_PASSWORD}");
288        return pair(reg_user, reg_password, &source);
289    }
290    let source = format!("${ENV_USER} / ${ENV_PASSWORD}");
291    pair(
292        env_user.filter(|v| !v.is_empty()),
293        env_password.filter(|v| !v.is_empty()),
294        &source,
295    )
296}
297
298fn resolve_user(role: NatsRole) -> Result<Option<UserCredential>> {
299    resolve_user_with(
300        role,
301        secrets::read_hklm_value,
302        std::env::var(ENV_USER).ok(),
303        std::env::var(ENV_PASSWORD).ok(),
304    )
305}
306
307/// The credentials a connection is built from.
308///
309/// [`connect`] resolves them from the registry / environment; this type
310/// exists so a caller (tests, chiefly) can supply them directly without
311/// touching process-global state. `Debug` never prints a secret.
312#[derive(Clone, Default)]
313pub struct NatsCredentials {
314    token: Option<String>,
315    user: Option<UserCredential>,
316}
317
318impl NatsCredentials {
319    pub fn new(token: Option<String>, user: Option<(String, String)>) -> Self {
320        Self {
321            token,
322            user: user.map(|(name, password)| UserCredential { name, password }),
323        }
324    }
325
326    fn resolve(role: NatsRole) -> Result<Self> {
327        Ok(Self {
328            token: resolve_token(role),
329            user: resolve_user(role)?,
330        })
331    }
332}
333
334impl std::fmt::Debug for NatsCredentials {
335    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
336        f.debug_struct("NatsCredentials")
337            .field("token", &self.token.as_ref().map(|_| "<redacted>"))
338            .field("user", &self.user)
339            .finish()
340    }
341}
342
343/// Which credential a connection attempt presents.
344#[derive(Debug, Clone, Copy, PartialEq, Eq)]
345enum Choice {
346    User,
347    Token,
348}
349
350impl Choice {
351    fn label(self) -> &'static str {
352        match self {
353            Choice::User => "user",
354            Choice::Token => "token",
355        }
356    }
357}
358
359/// What the probe connection learned about the broker.
360#[derive(Debug, Clone, Copy, PartialEq, Eq)]
361enum ProbeOutcome {
362    Accepted,
363    Rejected,
364    Unreachable,
365}
366
367impl ProbeOutcome {
368    fn label(self) -> &'static str {
369        match self {
370            ProbeOutcome::Accepted => "accepted",
371            ProbeOutcome::Rejected => "rejected",
372            ProbeOutcome::Unreachable => "unreachable",
373        }
374    }
375}
376
377/// The credential the broker last *proved* it accepts or rejects for a role,
378/// shared between the connection (which writes it on every attempt) and the
379/// [`CredentialProbe`] (which reads it). Only a probe answer is recorded — a
380/// guess made while the broker was unreachable is not evidence.
381#[derive(Default)]
382struct Live {
383    decision: Mutex<Option<Choice>>,
384    /// Consecutive authentication refusals: when the first one happened and
385    /// how many there have been. Cleared by a successful connect.
386    refusals: Mutex<Option<(std::time::Instant, usize)>>,
387}
388
389impl Live {
390    /// Track the connection's events for [`Live::auth_failed`].
391    fn observe(&self, ev: &async_nats::Event) {
392        let mut refusals = self.refusals.lock().unwrap_or_else(|e| e.into_inner());
393        match ev {
394            // A successful connect, a drop, or any other kind of failed
395            // attempt (broker unreachable, timeout) ends a run of refusals:
396            // only an unbroken run is evidence the credential itself is bad,
397            // and an offline broker is something to wait out.
398            async_nats::Event::Connected | async_nats::Event::Disconnected => *refusals = None,
399            async_nats::Event::ClientError(async_nats::ClientError::Other(kind))
400                if *kind == async_nats::ConnectErrorKind::AuthorizationViolation.to_string()
401                    || *kind == async_nats::ConnectErrorKind::Authentication.to_string() =>
402            {
403                let (first, n) = refusals.unwrap_or((std::time::Instant::now(), 0));
404                *refusals = Some((first, n + 1));
405            }
406            async_nats::Event::ClientError(_) => *refusals = None,
407            _ => {}
408        }
409    }
410
411    /// Whether the broker has refused every attempt for long enough that
412    /// waiting longer is not going to help.
413    fn auth_failed(&self) -> bool {
414        match *self.refusals.lock().unwrap_or_else(|e| e.into_inner()) {
415            Some((first, n)) => {
416                n >= AUTH_REJECTION_LIMIT && first.elapsed() >= AUTH_REJECTION_WINDOW
417            }
418            None => false,
419        }
420    }
421
422    fn get(&self) -> Option<Choice> {
423        *self.decision.lock().unwrap_or_else(|e| e.into_inner())
424    }
425
426    fn set(&self, c: Choice) {
427        *self.decision.lock().unwrap_or_else(|e| e.into_inner()) = Some(c);
428    }
429}
430
431/// Process-wide per-role state, so [`CredentialProbe::for_role`] sees the
432/// decision the already-open connection made without the callers having to
433/// thread anything through.
434fn live_for(role: NatsRole) -> Arc<Live> {
435    static LIVE: OnceLock<Mutex<HashMap<&'static str, Arc<Live>>>> = OnceLock::new();
436    let mut map = LIVE
437        .get_or_init(Default::default)
438        .lock()
439        .unwrap_or_else(|e| e.into_inner());
440    map.entry(role.as_str()).or_default().clone()
441}
442
443/// Pick the credential for one attempt from the probe's answer.
444///
445/// Pure so the whole decision table is testable without a broker. `Err`
446/// carries the credential that is missing.
447fn select(
448    outcome: ProbeOutcome,
449    last_worked: Option<Choice>,
450    have_token: bool,
451) -> std::result::Result<Choice, &'static str> {
452    match outcome {
453        ProbeOutcome::Accepted => Ok(Choice::User),
454        ProbeOutcome::Rejected if have_token => Ok(Choice::Token),
455        ProbeOutcome::Rejected => {
456            Err("the broker rejected the NATS user and no NatsToken is provisioned")
457        }
458        ProbeOutcome::Unreachable => Ok(last_worked.unwrap_or(Choice::User)),
459    }
460}
461
462/// Try the user on a connection of its own, which does not retry, so a
463/// rejection costs this throwaway connection and nothing else.
464///
465/// Runs on a spawned task because the auth callback's future must be `Sync`
466/// and async-nats' connect future is not; a join handle is.
467async fn probe_user(url: &str, user: &UserCredential) -> ProbeOutcome {
468    let (url, name, password) = (url.to_string(), user.name.clone(), user.password.clone());
469    let probe = tokio::spawn(async move {
470        let attempt = async_nats::ConnectOptions::new()
471            .name(PROBE_NAME)
472            .connection_timeout(PROBE_TIMEOUT)
473            .user_and_password(name, password)
474            .connect(url);
475        match tokio::time::timeout(PROBE_TIMEOUT + Duration::from_secs(1), attempt).await {
476            Ok(Ok(_client)) => ProbeOutcome::Accepted,
477            Ok(Err(e)) if e.kind() == async_nats::ConnectErrorKind::AuthorizationViolation => {
478                ProbeOutcome::Rejected
479            }
480            Ok(Err(_)) | Err(_) => ProbeOutcome::Unreachable,
481        }
482    });
483    probe.await.unwrap_or(ProbeOutcome::Unreachable)
484}
485
486/// Probe, decide, record. `probe` is injected so the unit tests can drive
487/// every outcome.
488async fn choose<P, Fut>(
489    role: NatsRole,
490    have_token: bool,
491    live: &Live,
492    probe: P,
493) -> std::result::Result<Choice, &'static str>
494where
495    P: FnOnce() -> Fut,
496    Fut: std::future::Future<Output = ProbeOutcome>,
497{
498    let outcome = probe().await;
499    let previous = live.get();
500    let choice = select(outcome, previous, have_token)?;
501    if outcome != ProbeOutcome::Unreachable {
502        live.set(choice);
503    }
504    // Log what was selected and the outcome class, never a value; quiet when
505    // nothing changed so a flapping broker does not fill the log.
506    if previous == Some(choice) {
507        debug!(
508            role = role.as_str(),
509            credential = choice.label(),
510            probe = outcome.label(),
511            "NATS credential selected"
512        );
513    } else {
514        info!(
515            role = role.as_str(),
516            credential = choice.label(),
517            probe = outcome.label(),
518            "NATS credential selected"
519        );
520    }
521    Ok(choice)
522}
523
524/// Whether a connection needs per-attempt selection at all.
525enum AuthPlan {
526    /// No user: present the token (or nothing), exactly as before roles had
527    /// users. No callback, no probe.
528    Static(Option<String>),
529    /// A user is provisioned: select per attempt.
530    Select {
531        token: Option<String>,
532        user: UserCredential,
533    },
534}
535
536fn auth_plan(creds: NatsCredentials) -> AuthPlan {
537    match creds.user {
538        None => AuthPlan::Static(creds.token),
539        Some(user) => AuthPlan::Select {
540            token: creds.token,
541            user,
542        },
543    }
544}
545
546/// The `name` a kanade process announces on its NATS connection.
547///
548/// Without an identity this is `kanade-<role>`; with one it is
549/// `kanade-<role>/<identity>`. The broker echoes it back verbatim in
550/// `/connz`, which is what lets the backend attribute a connection — and
551/// therefore the credential the broker authenticated it with — to a pc_id
552/// (#1270). Nothing else on a connection carries the pc_id: the CID is
553/// assigned by the server and the IP is not a stable identifier on a fleet
554/// of laptops.
555///
556/// The name is client-supplied and therefore claimed, not proved. What
557/// `/connz` makes unforgeable is the *credential* half of the pair; a host
558/// can still lie about which pc_id it is. Under one fleet-wide token that
559/// changes nothing (every host can already impersonate every other), and
560/// closing it for good is per-agent identity, not a naming convention.
561pub fn client_name(role: NatsRole, identity: Option<&str>) -> String {
562    match identity.map(str::trim).filter(|s| !s.is_empty()) {
563        Some(id) => format!("{NAME_PREFIX}{}{NAME_SEP}{id}", role.as_str()),
564        None => format!("{NAME_PREFIX}{}", role.as_str()),
565    }
566}
567
568/// A connection name split back into its parts — see [`client_name`].
569#[derive(Debug, Clone, Copy, PartialEq, Eq)]
570pub struct ClientName<'a> {
571    /// The role segment as announced. A `&str` rather than a [`NatsRole`]
572    /// on purpose: a connection from a future (or foreign) build may name a
573    /// role this binary does not know, and dropping it on the floor would
574    /// hide exactly the host worth looking at.
575    pub role: &'a str,
576    /// The host identity, when the connection carried one. `None` for the
577    /// backend / CLI (which are not per-host) and for agents predating
578    /// #1270 — those simply cannot be attributed.
579    pub identity: Option<&'a str>,
580}
581
582/// Parse a connection name produced by [`client_name`]. `None` for any name
583/// that is not a kanade connection at all (a `nats` CLI session, a
584/// monitoring tool), which the caller should ignore rather than guess about.
585pub fn parse_client_name(name: &str) -> Option<ClientName<'_>> {
586    let rest = name.strip_prefix(NAME_PREFIX)?;
587    Some(match rest.split_once(NAME_SEP) {
588        // An empty identity (`kanade-agent/`) is not an identity.
589        Some((role, id)) if !role.is_empty() && !id.is_empty() => ClientName {
590            role,
591            identity: Some(id),
592        },
593        Some((role, _)) if !role.is_empty() => ClientName {
594            role,
595            identity: None,
596        },
597        Some(_) => return None,
598        None if !rest.is_empty() => ClientName {
599            role: rest,
600            identity: None,
601        },
602        None => return None,
603    })
604}
605
606/// Answers "is this credential the one *we* present?" without handing the
607/// credential itself to the caller.
608///
609/// #1270: the NATS monitoring endpoint reports `authorized_user` per
610/// connection. A current nats-server hides that field for
611/// token-authenticated connections, but that is the broker build's
612/// behaviour, not a guarantee this side can lean on — a consumer of
613/// `/connz` has to treat the value as possibly being the fleet-wide secret.
614/// The one question it may safely answer about it is whether it equals the
615/// credential this process already holds, and that answer is enough to
616/// label a connection ("still on the shared token") without ever storing or
617/// serving the value.
618///
619/// Constructed once and reused: [`resolve_token`] hits the Windows registry,
620/// and the caller compares against every connection on the broker.
621pub struct CredentialProbe {
622    token: Option<String>,
623    user: Option<String>,
624    /// The connection's own decision, when this probe belongs to one. `None`
625    /// for the explicit constructors, which describe a fixed credential.
626    live: Option<Arc<Live>>,
627}
628
629/// Which shape of credential a [`CredentialProbe`] holds.
630#[derive(Debug, Clone, Copy, PartialEq, Eq)]
631pub enum CredentialKind {
632    None,
633    Token,
634    /// A named user — and, when the connection presenting it is **live**,
635    /// positive proof that the broker is running `users` rather than a
636    /// token: nats-server refuses to load a config carrying both a `token`
637    /// and a `users` array ("Can not have a token and a users array") and
638    /// rejects token authentication outright once `users` are defined.
639    ///
640    /// That proof is the only thing that makes a reported `authorized_user`
641    /// safe to record verbatim. Note what is *not* proof: holding no
642    /// credential locally. A process that never authenticated at all can
643    /// still read a monitoring endpoint, and inferring the broker's mode
644    /// from a local absence would let a misconfigured host store the very
645    /// secret the rest of this type exists to protect.
646    User,
647}
648
649/// Hand-written so a stray `{:?}` in a log line cannot print the credential.
650/// A username is not a secret and is shown; a token never is.
651impl std::fmt::Debug for CredentialProbe {
652    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
653        let mut held = Vec::new();
654        if self.token.is_some() {
655            held.push("<redacted token>".to_string());
656        }
657        if let Some(name) = &self.user {
658            held.push(format!("user {name}"));
659        }
660        let rendered = if held.is_empty() {
661            "<none>".to_string()
662        } else {
663            held.join(" + ")
664        };
665        f.debug_struct("CredentialProbe")
666            .field("presented", &rendered)
667            .field("kind", &self.kind())
668            .finish()
669    }
670}
671
672impl CredentialProbe {
673    /// Resolve the credential `role` would present, exactly as [`connect`]
674    /// does, and follow the decision that role's open connection makes.
675    ///
676    /// A half-provisioned user pair is ignored here — [`connect`] is where
677    /// it fails loudly.
678    pub fn for_role(role: NatsRole) -> Self {
679        Self {
680            token: resolve_token(role),
681            user: resolve_user(role).ok().flatten().map(|u| u.name),
682            live: Some(live_for(role)),
683        }
684    }
685
686    /// Build a probe around an explicitly-supplied token.
687    ///
688    /// The production path is [`Self::for_role`]; this exists so callers can
689    /// be tested against a known credential without a Windows registry to
690    /// write to — and, on a developer's machine, without accidentally
691    /// probing the real fleet token that `for_role` would find there.
692    pub fn from_token(token: Option<String>) -> Self {
693        Self {
694            token,
695            user: None,
696            live: None,
697        }
698    }
699
700    /// Build a probe for a process that authenticates as a named user.
701    ///
702    /// Describes a fixed user credential, for testing the consumers of
703    /// [`CredentialKind::User`] without a connection.
704    pub fn from_user(name: impl Into<String>) -> Self {
705        Self {
706            token: None,
707            user: Some(name.into()),
708            live: None,
709        }
710    }
711
712    /// Which shape of credential this process is presenting **right now**.
713    ///
714    /// A role that holds a user but is on the token fallback reports
715    /// [`CredentialKind::Token`]; the user shape is reported only once the
716    /// broker has accepted it. Before any decision a role holding a token
717    /// reports the token, and one holding only a user reports `None`: the
718    /// user shape is proof of the broker's mode and is not claimed on a
719    /// guess.
720    pub fn kind(&self) -> CredentialKind {
721        let decided = self.live.as_ref().and_then(|l| l.get());
722        match decided {
723            Some(Choice::User) if self.user.is_some() => return CredentialKind::User,
724            Some(Choice::Token) if self.token.is_some() => return CredentialKind::Token,
725            _ => {}
726        }
727        if self.live.is_none() && self.user.is_some() {
728            CredentialKind::User
729        } else if self.token.is_some() {
730            CredentialKind::Token
731        } else {
732            CredentialKind::None
733        }
734    }
735
736    /// Whether `candidate` is the **secret** this process presents.
737    ///
738    /// Only ever true for a token. A user's secret is its password, which
739    /// `authorized_user` never carries — matching a username here would
740    /// mean "this connection is on the same account", a different and much
741    /// weaker statement than the one callers use this for.
742    ///
743    /// A plain comparison: `candidate` comes from the broker's own report of
744    /// connections it already authenticated, not from an attacker-chosen
745    /// input, so there is no oracle to time.
746    pub fn is_ours(&self, candidate: &str) -> bool {
747        self.token.as_deref() == Some(candidate)
748    }
749}
750
751/// Connect to NATS at `url` as `role`. Resolves the credential from the
752/// registry (Windows) or the environment: a user pair when one is
753/// provisioned (with the token as fallback, see the module docs), otherwise
754/// the bearer token; connects unauthenticated when neither is set.
755///
756/// The connection is announced as `kanade-<role>` with no host identity —
757/// right for the backend and the CLI, which are not per-host. A role that
758/// has to be attributable to a specific machine (the agent) must use
759/// [`connect_with_event_callback`] and pass one; see [`client_name`].
760pub async fn connect(role: NatsRole, url: &str) -> Result<async_nats::Client> {
761    connect_inner(
762        role,
763        url,
764        None,
765        None::<fn(async_nats::Event) -> std::future::Ready<()>>,
766        NatsCredentials::resolve(role)?,
767        live_for(role),
768    )
769    .await
770}
771
772/// Same as [`connect`] but also wires an `event_callback` that fires
773/// whenever async-nats publishes a `ConnectEvent` (Connected,
774/// Disconnected, ServerError, etc.). The callback's `Future` runs on
775/// the async-nats internal task — keep it cheap and non-blocking
776/// (set a flag, send on a channel, that kind of thing) so the
777/// connection state machine isn't held up.
778///
779/// Used by the agent's v0.26 Layer 2 staleness tracker: the callback
780/// stamps a shared `Mutex<Option<Instant>>` on every Connected event,
781/// so `decide()` at fire time can answer "how long ago were we last
782/// definitely-talking-to-the-broker" without a polling loop.
783///
784/// `identity` names the host this connection belongs to (the agent's
785/// pc_id). It becomes part of the connection name the broker echoes in
786/// `/connz`, which is the only thing tying a connection — and the
787/// credential that opened it — back to a machine (#1270).
788pub async fn connect_with_event_callback<F, Fut>(
789    role: NatsRole,
790    url: &str,
791    identity: Option<&str>,
792    cb: F,
793) -> Result<async_nats::Client>
794where
795    F: Fn(async_nats::Event) -> Fut + Send + Sync + 'static,
796    Fut: std::future::Future<Output = ()> + Send + Sync + 'static,
797{
798    connect_inner(
799        role,
800        url,
801        identity,
802        Some(cb),
803        NatsCredentials::resolve(role)?,
804        live_for(role),
805    )
806    .await
807}
808
809/// Connect with explicitly supplied credentials instead of resolving them
810/// from the registry / environment, and with decision state private to this
811/// connection. For tests, which must not depend on process-global state
812/// (the per-role decision and refusal state is still shared, so concurrent
813/// tests should use different roles).
814pub async fn connect_with_credentials(
815    role: NatsRole,
816    url: &str,
817    creds: NatsCredentials,
818) -> Result<async_nats::Client> {
819    connect_inner(
820        role,
821        url,
822        None,
823        None::<fn(async_nats::Event) -> std::future::Ready<()>>,
824        creds,
825        live_for(role),
826    )
827    .await
828}
829
830/// [`connect_with_credentials`] with an event callback, so a test can observe
831/// what the connection reports while it cannot authenticate.
832pub async fn connect_with_credentials_and_event_callback<F, Fut>(
833    role: NatsRole,
834    url: &str,
835    creds: NatsCredentials,
836    cb: F,
837) -> Result<async_nats::Client>
838where
839    F: Fn(async_nats::Event) -> Fut + Send + Sync + 'static,
840    Fut: std::future::Future<Output = ()> + Send + Sync + 'static,
841{
842    connect_inner(role, url, None, Some(cb), creds, live_for(role)).await
843}
844
845async fn connect_inner<F, Fut>(
846    role: NatsRole,
847    url: &str,
848    identity: Option<&str>,
849    cb: Option<F>,
850    creds: NatsCredentials,
851    live: Arc<Live>,
852) -> Result<async_nats::Client>
853where
854    F: Fn(async_nats::Event) -> Fut + Send + Sync + 'static,
855    Fut: std::future::Future<Output = ()> + Send + Sync + 'static,
856{
857    // #1187: a workspace build unifies rustls features, so any binary built
858    // alongside a reqwest user (the CLI, the backend) links BOTH aws-lc-rs and
859    // ring. rustls 0.23 then cannot auto-pick a process-level provider, and the
860    // first TLS handshake — a `wss://`/`tls://` broker — panics inside
861    // async-nats' connection task. That panic does not surface here: the
862    // client just goes dead. Every production connect funnels through this
863    // function, so installing ring here covers every binary. `install_default`
864    // returns Err when a provider is already installed — ignore it.
865    let _ = rustls::crypto::ring::default_provider().install_default();
866
867    let tracker = live.clone();
868    // Only a provisioned user changes how the credential is presented; every
869    // other host takes the token path untouched.
870    let opts = match auth_plan(creds) {
871        AuthPlan::Static(token) => {
872            let opts = async_nats::ConnectOptions::new();
873            match token {
874                Some(token) => opts.token(token),
875                None => opts,
876            }
877        }
878        AuthPlan::Select { token, user } => {
879            let url = url.to_string();
880            async_nats::ConnectOptions::with_auth_callback(move |_nonce| {
881                let (url, token, user, live) =
882                    (url.clone(), token.clone(), user.clone(), live.clone());
883                async move {
884                    let choice = choose(role, token.is_some(), &live, || probe_user(&url, &user))
885                        .await
886                        .map_err(async_nats::AuthError::new)?;
887                    let mut auth = async_nats::Auth::new();
888                    match choice {
889                        Choice::User => {
890                            auth.username = Some(user.name);
891                            auth.password = Some(user.password);
892                        }
893                        Choice::Token => auth.token = token,
894                    }
895                    Ok(auth)
896                }
897            })
898        }
899    };
900
901    // v0.38 / #137: offline-tolerant boot. Without
902    // `retry_on_initial_connect`, `opts.connect(url).await` blocks-then-
903    // errors when the broker is unreachable at startup — the agent
904    // process dies, SCM ticks its restart counter, and the offline-
905    // tolerant subsystems (local_scheduler, outbox drain) never spawn.
906    // With this flag, connect() returns `Ok(Client)` immediately and
907    // async-nats does the reconnect in the background; subscribe()
908    // calls queue the SUB frame until the link is up.
909    let opts = opts
910        .retry_on_initial_connect()
911        // Names the connection in `nats server report connections`, in the
912        // broker's own logs, and in `/connz`. Free observability while the
913        // fleet is mid-migration: it shows which roles are connecting even
914        // before their credentials differ, which is exactly the window in
915        // which a wrongly-provisioned host is otherwise invisible. With an
916        // identity it also carries the pc_id, so #1270 can join the broker's
917        // per-connection `authorized_user` back onto the agents row.
918        .name(client_name(role, identity));
919    // Always installed, even without a caller callback: it is how a refused
920    // credential is noticed, since async-nats only retries it.
921    let opts = opts.event_callback(move |ev| {
922        tracker.observe(&ev);
923        let forwarded = cb.as_ref().map(|cb| cb(ev));
924        async move {
925            if let Some(forwarded) = forwarded {
926                forwarded.await;
927            }
928        }
929    });
930    opts.connect(url)
931        .await
932        .with_context(|| format!("connect to NATS at {url}"))
933}
934
935/// Whether the client's connection task has terminated.
936///
937/// When that task ends (a panic inside it, a drain, exhausted reconnects) the
938/// `Client` stays alive and never talks again — a silent zombie. The
939/// observable difference from an ordinary disconnect is the command channel:
940/// a task that is alive but reconnecting simply does not answer a flush, so
941/// the flush waits and the timeout elapses; a task that is gone has dropped
942/// the receiving end, so queueing the flush fails with a send error. The
943/// check publishes to no subject, so no role needs any publish right for it.
944///
945/// This tells whether the task still accepts commands, not whether the broker
946/// is reachable. A flush that fails after it was queued (its observer was
947/// dropped) is not treated as death: a task that really is gone fails the
948/// next check at the send step, and mistaking a plain disconnect for death
949/// would make a process exit needlessly.
950pub async fn is_dead(client: &async_nats::Client) -> bool {
951    match tokio::time::timeout(Duration::from_secs(5), client.flush()).await {
952        Ok(Err(e)) => match e.kind() {
953            async_nats::client::FlushErrorKind::SendError => true,
954            async_nats::client::FlushErrorKind::FlushError => false,
955        },
956        Ok(Ok(())) | Err(_) => false,
957    }
958}
959
960/// Resolve once `client` can no longer be expected to work: its connection
961/// task has terminated (see [`is_dead`]), or the broker has refused its
962/// credential on every attempt for [`AUTH_REJECTION_WINDOW`]. Checked every
963/// `interval`. A process that depends on the client should treat this as
964/// fatal and exit non-zero so its service manager restarts it.
965///
966/// `role` names the connection the client was opened with.
967pub async fn wait_until_dead_every(
968    role: NatsRole,
969    client: &async_nats::Client,
970    interval: Duration,
971) {
972    let live = live_for(role);
973    let mut tick = tokio::time::interval(interval);
974    tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
975    loop {
976        tick.tick().await;
977        if live.auth_failed() {
978            warn!(
979                role = role.as_str(),
980                "the NATS broker keeps refusing this role's credential; the client cannot connect"
981            );
982            return;
983        }
984        if is_dead(client).await {
985            warn!(
986                role = role.as_str(),
987                "NATS connection task has terminated; the client can no longer talk to the broker"
988            );
989            return;
990        }
991    }
992}
993
994/// [`wait_until_dead_every`] at the production cadence.
995pub async fn wait_until_dead(role: NatsRole, client: &async_nats::Client) {
996    wait_until_dead_every(role, client, DEAD_CHECK_INTERVAL).await
997}
998
999#[cfg(test)]
1000mod tests {
1001    use super::*;
1002    use std::collections::HashMap;
1003
1004    /// A stand-in registry. Keys are `subkey\value`.
1005    fn reg(entries: &[(&str, &str)]) -> impl Fn(&str, &str) -> Option<String> {
1006        let map: HashMap<String, String> = entries
1007            .iter()
1008            .map(|(k, v)| ((*k).to_string(), (*v).to_string()))
1009            .collect();
1010        move |subkey: &str, value: &str| map.get(&format!(r"{subkey}\{value}")).cloned()
1011    }
1012
1013    #[test]
1014    fn role_subkeys_are_distinct_and_agent_matches_the_shared_path() {
1015        assert_eq!(NatsRole::Backend.reg_subkey(), r"SOFTWARE\kanade\backend");
1016        assert_eq!(NatsRole::Cli.reg_subkey(), r"SOFTWARE\kanade\cli");
1017        // The agent's role key IS the historical shared key, so an agent
1018        // never sees a migration at all.
1019        assert_eq!(NatsRole::Agent.reg_subkey(), REG_SHARED_SUBKEY);
1020    }
1021
1022    #[test]
1023    fn an_unmigrated_fleet_keeps_presenting_the_shared_token() {
1024        // The state of every host today: one token, under the agent path.
1025        let registry = reg(&[(r"SOFTWARE\kanade\agent\NatsToken", "shared")]);
1026        for role in [NatsRole::Agent, NatsRole::Backend, NatsRole::Cli] {
1027            assert_eq!(
1028                resolve_token_with(role, &registry, None).as_deref(),
1029                Some("shared"),
1030                "{role:?} must keep working before its own key is provisioned"
1031            );
1032        }
1033    }
1034
1035    #[test]
1036    fn a_role_key_wins_over_the_shared_one() {
1037        let registry = reg(&[
1038            (r"SOFTWARE\kanade\agent\NatsToken", "shared"),
1039            (r"SOFTWARE\kanade\backend\NatsToken", "backend-only"),
1040        ]);
1041        // The migrated role uses its own credential...
1042        assert_eq!(
1043            resolve_token_with(NatsRole::Backend, &registry, None).as_deref(),
1044            Some("backend-only")
1045        );
1046        // ...while a role that has not been migrated yet is unaffected.
1047        assert_eq!(
1048            resolve_token_with(NatsRole::Cli, &registry, None).as_deref(),
1049            Some("shared")
1050        );
1051    }
1052
1053    #[test]
1054    fn removing_a_role_key_falls_back_rather_than_failing() {
1055        // Rollback of a per-host migration step: the role key is gone, and
1056        // the host must return to the shared credential instead of
1057        // connecting unauthenticated (which a broker with `authorization`
1058        // would refuse — turning a rollback into an outage).
1059        let registry = reg(&[(r"SOFTWARE\kanade\agent\NatsToken", "shared")]);
1060        assert_eq!(
1061            resolve_token_with(NatsRole::Backend, &registry, None).as_deref(),
1062            Some("shared")
1063        );
1064    }
1065
1066    #[test]
1067    fn the_registry_outranks_the_environment() {
1068        // Unchanged from before roles existed: a dev shell's env var must
1069        // not quietly override a provisioned production credential.
1070        let registry = reg(&[(r"SOFTWARE\kanade\agent\NatsToken", "shared")]);
1071        assert_eq!(
1072            resolve_token_with(NatsRole::Agent, &registry, Some("from-env".into())).as_deref(),
1073            Some("shared")
1074        );
1075    }
1076
1077    #[test]
1078    fn the_environment_serves_hosts_with_no_registry_at_all() {
1079        let empty = reg(&[]);
1080        assert_eq!(
1081            resolve_token_with(NatsRole::Cli, &empty, Some("from-env".into())).as_deref(),
1082            Some("from-env")
1083        );
1084        // An empty env var is not a credential — it must fall through to
1085        // "no token" so a dev broker without `authorization` still works,
1086        // rather than presenting the empty string and being rejected.
1087        assert_eq!(
1088            resolve_token_with(NatsRole::Cli, &empty, Some(String::new())),
1089            None
1090        );
1091        assert_eq!(resolve_token_with(NatsRole::Cli, &empty, None), None);
1092    }
1093
1094    // ── #1270: connection naming ─────────────────────────────────────
1095
1096    #[test]
1097    fn an_identity_round_trips_through_the_connection_name() {
1098        // The pc_id is the join key between `/connz` and the agents table,
1099        // so the name has to survive the trip unchanged — including the
1100        // casing, which is NOT uniform across the fleet and which NATS
1101        // subjects treat as significant.
1102        for pc in ["PC001", "minipc", "Web%01", "ws-9"] {
1103            let name = client_name(NatsRole::Agent, Some(pc));
1104            let parsed = parse_client_name(&name).expect("our own name must parse");
1105            assert_eq!(parsed.role, "agent");
1106            assert_eq!(parsed.identity, Some(pc));
1107        }
1108    }
1109
1110    #[test]
1111    fn a_role_without_an_identity_keeps_the_pre_1270_name() {
1112        // The backend and the CLI are not per-host, and an agent that
1113        // predates #1270 announces this shape too. Both must parse as "a
1114        // kanade connection we cannot attribute" rather than as an error or
1115        // as an empty pc_id.
1116        assert_eq!(client_name(NatsRole::Backend, None), "kanade-backend");
1117        let parsed = parse_client_name("kanade-agent").unwrap();
1118        assert_eq!(parsed.role, "agent");
1119        assert_eq!(parsed.identity, None);
1120        // Whitespace-only is not an identity either — it would otherwise
1121        // produce a name that parses back into a pc_id no row can match.
1122        assert_eq!(client_name(NatsRole::Agent, Some("  ")), "kanade-agent");
1123    }
1124
1125    #[test]
1126    fn foreign_connections_do_not_parse_as_kanade_ones() {
1127        // A `nats` CLI session or a monitoring tool shares the broker. The
1128        // projector must skip those rather than attribute them to a host.
1129        assert!(parse_client_name("NATS CLI Version 0.1.5").is_none());
1130        assert!(parse_client_name("").is_none());
1131        assert!(parse_client_name("kanade-").is_none());
1132        assert!(parse_client_name("kanade-/PC001").is_none());
1133        // A trailing separator with no identity is a role, not a pc_id.
1134        assert_eq!(parse_client_name("kanade-agent/").unwrap().identity, None);
1135    }
1136
1137    #[test]
1138    fn an_unknown_role_is_preserved_rather_than_dropped() {
1139        // A future build (or something impersonating one) naming a role this
1140        // binary has never heard of is precisely the connection an operator
1141        // wants to see.
1142        let parsed = parse_client_name("kanade-relay/PC001").unwrap();
1143        assert_eq!(parsed.role, "relay");
1144        assert_eq!(parsed.identity, Some("PC001"));
1145    }
1146
1147    // ── #1270: credential probe ──────────────────────────────────────
1148
1149    #[test]
1150    fn the_probe_recognises_only_the_credential_we_present() {
1151        let probe = CredentialProbe::from_token(Some("shared".into()));
1152        assert_eq!(probe.kind(), CredentialKind::Token);
1153        assert!(probe.is_ours("shared"));
1154        assert!(!probe.is_ours("something-else"));
1155        // The empty string is what the broker reports for a connection it
1156        // did not authenticate at all. It must never read as "ours".
1157        assert!(!probe.is_ours(""));
1158    }
1159
1160    #[test]
1161    fn a_probe_with_no_credential_matches_nothing() {
1162        // Dev broker with no `authorization` block. We hold nothing, so we
1163        // can prove nothing about anyone else's credential — including that
1164        // it is safe to store.
1165        let probe = CredentialProbe::from_token(None);
1166        assert_eq!(probe.kind(), CredentialKind::None);
1167        assert!(!probe.is_ours(""));
1168        assert!(!probe.is_ours("anything"));
1169    }
1170
1171    #[test]
1172    fn a_username_is_not_a_secret_we_can_recognise() {
1173        // `is_ours` answers "is this MY secret", and a user's secret is its
1174        // password. Matching on the username instead would answer a much
1175        // weaker question while reading like the strong one.
1176        let probe = CredentialProbe::from_user("kanade-backend");
1177        assert_eq!(probe.kind(), CredentialKind::User);
1178        assert!(!probe.is_ours("kanade-backend"));
1179    }
1180
1181    #[test]
1182    fn the_probe_never_prints_the_credential() {
1183        // `/connz` handling logs liberally; one `{:?}` on the probe must not
1184        // be the thing that puts the fleet's token in a log file.
1185        let probe = CredentialProbe::from_token(Some("super-secret-token".into()));
1186        let rendered = format!("{probe:?}");
1187        assert!(!rendered.contains("super-secret-token"), "{rendered}");
1188        assert!(rendered.contains("redacted"), "{rendered}");
1189        assert!(
1190            format!("{:?}", CredentialProbe::from_token(None)).contains("none"),
1191            "the no-credential case should be visible, just not the value"
1192        );
1193        // A username is not a secret — hiding it would cost diagnosability
1194        // for nothing.
1195        assert!(
1196            format!("{:?}", CredentialProbe::from_user("kanade-backend"))
1197                .contains("kanade-backend"),
1198        );
1199    }
1200
1201    // ── per-role users ───────────────────────────────────────────────
1202
1203    fn user(r: Result<Option<UserCredential>>) -> Option<(String, String)> {
1204        r.unwrap().map(|u| (u.name, u.password))
1205    }
1206
1207    #[test]
1208    fn a_users_registry_pair_comes_from_the_roles_own_key() {
1209        let registry = reg(&[
1210            (r"SOFTWARE\kanade\backend\NatsUser", "kanade-backend"),
1211            (r"SOFTWARE\kanade\backend\NatsPassword", "pw"),
1212        ]);
1213        assert_eq!(
1214            user(resolve_user_with(NatsRole::Backend, &registry, None, None)),
1215            Some(("kanade-backend".into(), "pw".into()))
1216        );
1217        // No shared user: another role does not inherit it, not even the
1218        // agent's key (which is where the shared *token* lives).
1219        for role in [NatsRole::Agent, NatsRole::Cli] {
1220            assert_eq!(user(resolve_user_with(role, &registry, None, None)), None);
1221        }
1222        let agent = reg(&[
1223            (r"SOFTWARE\kanade\agent\NatsUser", "a"),
1224            (r"SOFTWARE\kanade\agent\NatsPassword", "b"),
1225        ]);
1226        assert_eq!(
1227            user(resolve_user_with(NatsRole::Cli, &agent, None, None)),
1228            None
1229        );
1230    }
1231
1232    #[test]
1233    fn the_registry_pair_outranks_the_environment_pair() {
1234        let registry = reg(&[
1235            (r"SOFTWARE\kanade\cli\NatsUser", "reg-user"),
1236            (r"SOFTWARE\kanade\cli\NatsPassword", "reg-pw"),
1237        ]);
1238        assert_eq!(
1239            user(resolve_user_with(
1240                NatsRole::Cli,
1241                &registry,
1242                Some("env-user".into()),
1243                Some("env-pw".into())
1244            )),
1245            Some(("reg-user".into(), "reg-pw".into()))
1246        );
1247    }
1248
1249    #[test]
1250    fn the_environment_serves_a_user_pair_when_the_registry_has_none() {
1251        assert_eq!(
1252            user(resolve_user_with(
1253                NatsRole::Cli,
1254                reg(&[]),
1255                Some("u".into()),
1256                Some("p".into())
1257            )),
1258            Some(("u".into(), "p".into()))
1259        );
1260        // Empty values are not a credential.
1261        assert_eq!(
1262            user(resolve_user_with(
1263                NatsRole::Cli,
1264                reg(&[]),
1265                Some(String::new()),
1266                Some(String::new())
1267            )),
1268            None
1269        );
1270        assert_eq!(
1271            user(resolve_user_with(NatsRole::Cli, reg(&[]), None, None)),
1272            None
1273        );
1274    }
1275
1276    #[test]
1277    fn half_a_pair_is_an_error_not_a_fallback() {
1278        let err = |r: Result<Option<UserCredential>>| format!("{:#}", r.unwrap_err());
1279        // Registry halves.
1280        let only_user = reg(&[(r"SOFTWARE\kanade\cli\NatsUser", "u")]);
1281        assert!(
1282            err(resolve_user_with(NatsRole::Cli, &only_user, None, None))
1283                .contains("password is missing")
1284        );
1285        let only_pw = reg(&[(r"SOFTWARE\kanade\cli\NatsPassword", "secret-pw")]);
1286        let e = err(resolve_user_with(NatsRole::Cli, &only_pw, None, None));
1287        assert!(e.contains("user name is missing"), "{e}");
1288        assert!(!e.contains("secret-pw"), "{e}");
1289        // A half registry pair is not rescued by a complete environment pair.
1290        assert!(
1291            resolve_user_with(
1292                NatsRole::Cli,
1293                &only_user,
1294                Some("u".into()),
1295                Some("p".into())
1296            )
1297            .is_err()
1298        );
1299        // Environment halves.
1300        assert!(resolve_user_with(NatsRole::Cli, reg(&[]), Some("u".into()), None).is_err());
1301        assert!(resolve_user_with(NatsRole::Cli, reg(&[]), None, Some("p".into())).is_err());
1302    }
1303
1304    #[test]
1305    fn credentials_never_print_a_secret() {
1306        let creds = NatsCredentials::new(
1307            Some("tok-secret".into()),
1308            Some(("kanade-agent".into(), "pw-secret".into())),
1309        );
1310        let rendered = format!("{creds:?}");
1311        assert!(!rendered.contains("tok-secret"), "{rendered}");
1312        assert!(!rendered.contains("pw-secret"), "{rendered}");
1313        assert!(rendered.contains("kanade-agent"), "{rendered}");
1314    }
1315
1316    #[test]
1317    fn no_user_means_the_static_token_path_with_no_selection() {
1318        match auth_plan(NatsCredentials::new(Some("t".into()), None)) {
1319            AuthPlan::Static(Some(t)) => assert_eq!(t, "t"),
1320            _ => panic!("a token-only host must take the static path"),
1321        }
1322        assert!(matches!(
1323            auth_plan(NatsCredentials::new(None, None)),
1324            AuthPlan::Static(None)
1325        ));
1326        assert!(matches!(
1327            auth_plan(NatsCredentials::new(
1328                Some("t".into()),
1329                Some(("u".into(), "p".into()))
1330            )),
1331            AuthPlan::Select { .. }
1332        ));
1333    }
1334
1335    #[test]
1336    fn selection_follows_the_probe() {
1337        use Choice::{Token, User};
1338        use ProbeOutcome::*;
1339        assert_eq!(select(Accepted, None, true), Ok(User));
1340        assert_eq!(select(Accepted, Some(Token), false), Ok(User));
1341        assert_eq!(select(Rejected, Some(User), true), Ok(Token));
1342        let missing = select(Rejected, None, false).unwrap_err();
1343        assert!(missing.contains("NatsToken"), "{missing}");
1344        // Unreachable keeps what last worked, defaulting to the user.
1345        assert_eq!(select(Unreachable, None, true), Ok(User));
1346        assert_eq!(select(Unreachable, Some(Token), true), Ok(Token));
1347        assert_eq!(select(Unreachable, Some(User), true), Ok(User));
1348    }
1349
1350    #[tokio::test]
1351    async fn only_a_probe_answer_is_remembered() {
1352        let live = Live::default();
1353        let role = NatsRole::Cli;
1354        let run = |o| choose(role, true, &live, move || async move { o });
1355        // An unreachable broker is a guess, not evidence.
1356        assert_eq!(run(ProbeOutcome::Unreachable).await, Ok(Choice::User));
1357        assert_eq!(live.get(), None);
1358        assert_eq!(run(ProbeOutcome::Rejected).await, Ok(Choice::Token));
1359        assert_eq!(live.get(), Some(Choice::Token));
1360        // The broker goes away: the token that last worked is reused.
1361        assert_eq!(run(ProbeOutcome::Unreachable).await, Ok(Choice::Token));
1362        // ...and the flip back is followed.
1363        assert_eq!(run(ProbeOutcome::Accepted).await, Ok(Choice::User));
1364        assert_eq!(live.get(), Some(Choice::User));
1365        // A rejection with no token fails and leaves the cache alone.
1366        let err = choose(role, false, &live, || async { ProbeOutcome::Rejected })
1367            .await
1368            .unwrap_err();
1369        assert!(err.contains("NatsToken"));
1370        assert_eq!(live.get(), Some(Choice::User));
1371    }
1372
1373    #[test]
1374    fn the_probe_reports_the_shape_in_use_right_now() {
1375        let live = Arc::new(Live::default());
1376        let probe = CredentialProbe {
1377            token: Some("tok".into()),
1378            user: Some("kanade-agent".into()),
1379            live: Some(live.clone()),
1380        };
1381        // Nothing proven yet: not the user shape.
1382        assert_eq!(probe.kind(), CredentialKind::Token);
1383        live.set(Choice::User);
1384        assert_eq!(probe.kind(), CredentialKind::User);
1385        // Holding a user while on the token fallback is the token shape.
1386        live.set(Choice::Token);
1387        assert_eq!(probe.kind(), CredentialKind::Token);
1388        assert!(probe.is_ours("tok"));
1389
1390        let user_only = CredentialProbe {
1391            token: None,
1392            user: Some("u".into()),
1393            live: Some(Arc::new(Live::default())),
1394        };
1395        assert_eq!(user_only.kind(), CredentialKind::None);
1396        user_only.live.as_ref().unwrap().set(Choice::User);
1397        assert_eq!(user_only.kind(), CredentialKind::User);
1398    }
1399
1400    #[test]
1401    fn a_probe_holding_both_credentials_prints_neither_secret() {
1402        let probe = CredentialProbe {
1403            token: Some("tok-secret".into()),
1404            user: Some("kanade-agent".into()),
1405            live: None,
1406        };
1407        let rendered = format!("{probe:?}");
1408        assert!(!rendered.contains("tok-secret"), "{rendered}");
1409        assert!(rendered.contains("kanade-agent"), "{rendered}");
1410    }
1411
1412    #[test]
1413    fn sustained_refusals_fail_the_client_and_a_connect_clears_them() {
1414        let refused = || {
1415            async_nats::Event::ClientError(async_nats::ClientError::Other(
1416                async_nats::ConnectErrorKind::AuthorizationViolation.to_string(),
1417            ))
1418        };
1419        let live = Live::default();
1420        for _ in 0..AUTH_REJECTION_LIMIT {
1421            live.observe(&refused());
1422        }
1423        // Enough refusals, but not over a long enough window: a broker
1424        // mid-reload must not count.
1425        assert!(!live.auth_failed());
1426        // Age the first refusal past the window.
1427        let aged = std::time::Instant::now() - AUTH_REJECTION_WINDOW - Duration::from_secs(1);
1428        *live.refusals.lock().unwrap() = Some((aged, AUTH_REJECTION_LIMIT));
1429        assert!(live.auth_failed());
1430        live.observe(&async_nats::Event::Connected);
1431        assert!(!live.auth_failed());
1432        // A different kind of failure (broker gone) ends the run, so refusals
1433        // from before an outage cannot add up to a failure after it.
1434        *live.refusals.lock().unwrap() = Some((aged, AUTH_REJECTION_LIMIT));
1435        live.observe(&async_nats::Event::ClientError(
1436            async_nats::ClientError::Other("io".into()),
1437        ));
1438        assert!(!live.auth_failed());
1439        live.observe(&refused());
1440        assert!(!live.auth_failed());
1441    }
1442
1443    #[test]
1444    fn connect_leaves_a_process_crypto_provider_installed() {
1445        // #1187: a binary that links both aws-lc-rs and ring (any workspace
1446        // build) panics on its first wss handshake unless a process-level
1447        // provider was installed first. That panic happens inside async-nats'
1448        // background task, so the only thing a caller ever sees is a client
1449        // that silently never works — pin the precondition here instead.
1450        // `retry_on_initial_connect` makes connect return without a broker.
1451        let rt = tokio::runtime::Builder::new_current_thread()
1452            .enable_all()
1453            .build()
1454            .unwrap();
1455        rt.block_on(async {
1456            connect(NatsRole::Agent, "nats://127.0.0.1:1")
1457                .await
1458                .unwrap();
1459        });
1460        assert!(rustls::crypto::CryptoProvider::get_default().is_some());
1461    }
1462}