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