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