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