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, ®istry, 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, ®istry, 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, ®istry, 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, ®istry, 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, ®istry, 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, ®istry, 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, ®istry, 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 ®istry,
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}