From 71b01013fae66023c552a3c2a9383436bf0082a5 Mon Sep 17 00:00:00 2001 From: lr90 Date: Tue, 29 Sep 2026 11:12:21 +0800 Subject: [PATCH 1/3] fix(auth): distinguish pending refresh during CLI startup --- .../astra-cli/src/cli/cli_config/cli_utils.rs | 7 +- crates/astra-cli/src/cli/native_auth.rs | 191 +++++-- .../src/cli/session/session_startup.rs | 12 +- .../tests/native_auth_environment.rs | 74 +++ crates/astra-credentials/src/native.rs | 475 ++++++++++++++++-- .../src/native/diagnostics.rs | 192 +++++++ docs/guides/moi-native-login.md | 29 ++ 7 files changed, 899 insertions(+), 81 deletions(-) create mode 100644 crates/astra-credentials/src/native/diagnostics.rs diff --git a/crates/astra-cli/src/cli/cli_config/cli_utils.rs b/crates/astra-cli/src/cli/cli_config/cli_utils.rs index de92b18168..6ec652e03d 100644 --- a/crates/astra-cli/src/cli/cli_config/cli_utils.rs +++ b/crates/astra-cli/src/cli/cli_config/cli_utils.rs @@ -175,8 +175,8 @@ pub(crate) fn cli_owner_auth_snapshot() -> CliOwnerAuthSnapshot { if let Some(binding) = &native_binding { let matches_owner = binding.profile_name() == identity.profile_name && binding - .snapshot() - .is_ok_and(|session| Some(session.astra_user_id) == identity.account_id); + .account_id() + .is_ok_and(|account_id| Some(account_id) == identity.account_id); return CliOwnerAuthSnapshot { owner_scope, access_token: None, @@ -235,8 +235,7 @@ pub(crate) fn configure_cli_profile_identity( admission: CliProfileIdentityAdmission, ) -> Result<(), String> { if let Some(binding) = crate::cli::native_auth::active() { - let session = binding.snapshot()?; - return install_cli_profile_identity(binding.profile_name(), Some(session.astra_user_id)); + return install_cli_profile_identity(binding.profile_name(), Some(binding.account_id()?)); } let creds = credential_store() .load() diff --git a/crates/astra-cli/src/cli/native_auth.rs b/crates/astra-cli/src/cli/native_auth.rs index f40a501fe7..8ca60b1314 100644 --- a/crates/astra-cli/src/cli/native_auth.rs +++ b/crates/astra-cli/src/cli/native_auth.rs @@ -55,14 +55,26 @@ impl Binding { ) } - pub(crate) fn snapshot(&self) -> Result { - let current = self.store.current()?; + fn check_identity( + &self, + current: native::NativeSession, + ) -> Result { if current.generation != self.session.generation || current.environment != self.session.environment || current.subject != self.session.subject { return Err("MOI account or environment changed; restart Astra".into()); } + Ok(current) + } + + /// Local identity binding does not require a usable access token. + pub(crate) fn account_id(&self) -> Result { + Ok(self.check_identity(self.store.current()?)?.astra_user_id) + } + + pub(crate) fn snapshot(&self) -> Result { + let current = self.check_identity(self.store.current()?)?; if current.refresh_pending { return Err("MOI token rotation was interrupted; run astra login".into()); } @@ -72,14 +84,7 @@ impl Binding { /// Same identity check as [`snapshot`], without failing when a rotation is pending. /// The file lock runs on the blocking pool. pub(crate) async fn snapshot_off_runtime(&self) -> Result { - let current = self.store.current_off_runtime().await?; - if current.generation != self.session.generation - || current.environment != self.session.environment - || current.subject != self.session.subject - { - return Err("MOI account or environment changed; restart Astra".into()); - } - Ok(current) + self.check_identity(self.store.current_off_runtime().await?) } pub(crate) async fn access_token(&self) -> Result { @@ -97,6 +102,22 @@ impl Binding { } } +/// Resolve native auth before startup launches cloud work. The progress timer +/// never cancels the credential future or starts another refresh. +pub(crate) async fn startup_access_token(binding: &Binding) -> Result { + let pending = binding.access_token(); + tokio::pin!(pending); + let result = tokio::select! { + biased; + result = &mut pending => result, + _ = tokio::time::sleep(std::time::Duration::from_secs(1)) => { + eprintln!("Refreshing your sign-in…"); + pending.await + } + }; + result.map_err(|error| error.access_miss().startup_warning().to_owned()) +} + impl astra_thin_client::client::BearerProvider for Binding { fn token( &self, @@ -163,10 +184,9 @@ pub(crate) fn bind_after_login( let mut base = api.api_origin(); let binding = bind_process(&mut base, true, None, false)? .ok_or("UC login did not publish credentials")?; - let session = binding.snapshot()?; crate::cli::cli_config::cli_utils::install_cli_profile_identity( binding.profile_name(), - Some(session.astra_user_id), + Some(binding.account_id()?), )?; Ok(api.clone().with_bearer_provider(binding)) } @@ -176,7 +196,7 @@ pub(crate) fn projected_credentials() -> Result Result Result<(), String> { let status = if !configured { serde_json::json!({"version": 1, "state": "not_configured"}) } else { - match store.current() { - Ok(session) => { - serde_json::json!({"version": 1, "state": if session.refresh_pending { "reauthentication_required" } else { "signed_in" }, + match store.status().await { + Ok((session, status)) => { + serde_json::json!({"version": 1, "state": status, "environment": session.environment.key(), "issuer": session.environment.issuer, "astra_url": session.environment.astra_url, "moi_url": session.environment.moi_url, "subject": session.subject, "generation": session.generation, "expires_at": session.expires_at, @@ -269,11 +291,13 @@ fn render_status(status: &serde_json::Value, json: bool) -> String { return status.to_string(); } match status["state"].as_str() { - Some("signed_in" | "reauthentication_required") => { - let heading = if status["state"] == "signed_in" { - "Signed in to MOI." - } else { - "Sign-in needs to be renewed. Run astra login." + Some("signed_in" | "refresh_in_progress" | "reauthentication_required") => { + let heading = match status["state"].as_str() { + Some("signed_in") => "Signed in to MOI.", + Some("refresh_in_progress") => { + "Sign-in refresh is in progress. Please wait; no login is needed." + } + _ => "Sign-in needs to be renewed. Run astra login.", }; format!( "{heading}\nAccount: {}\nAstra: {}\nMOI: {}\nWorkspace: {}", @@ -335,6 +359,7 @@ mod tests { ("not_configured", "not configured"), ("signed_out", "Signed out"), ("signed_in", "Signed in"), + ("refresh_in_progress", "no login is needed"), ("reauthentication_required", "needs to be renewed"), ] { let status = serde_json::json!({ @@ -349,19 +374,19 @@ mod tests { let human = render_status(&status, false); assert!(human.contains(expected), "{human}"); assert!(!human.starts_with('{')); - if matches!(state, "signed_in" | "reauthentication_required") { + if matches!( + state, + "signed_in" | "refresh_in_progress" | "reauthentication_required" + ) { assert!(human.contains("Workspace: not selected")); assert!(human.contains("Account: account-a")); } } } - #[tokio::test] - async fn login_rebinding_does_not_retarget_existing_bearer_provider() { - let root = tempfile::tempdir().unwrap(); - let store = NativeStore::with_directory(root.path().join("auth")); + fn test_session() -> native::NativeSession { let issuer = "https://uc.example.test/realms/moi"; - let original = native::NativeSession { + native::NativeSession { environment: native::Environment { issuer: issuer.into(), astra_url: "https://astra.example.test".into(), @@ -383,7 +408,14 @@ mod tests { workspace_id: None, role_id: None, refresh_pending: false, - }; + } + } + + #[tokio::test] + async fn login_rebinding_does_not_retarget_existing_bearer_provider() { + let root = tempfile::tempdir().unwrap(); + let store = NativeStore::with_directory(root.path().join("auth")); + let original = test_session(); let (session, _) = store.publish(original.clone()).unwrap(); let old = Binding { store: store.clone(), @@ -397,11 +429,112 @@ mod tests { let (session, _) = store.publish(next).unwrap(); let new = Binding { store, session }; assert!(old.snapshot().is_err()); + assert!(old.account_id().is_err()); assert!(old.access_token().await.is_err()); assert_eq!(new.access_token().await.unwrap(), "test-access-b"); assert_ne!(old.profile_name(), new.profile_name()); } + #[tokio::test] + async fn startup_identity_accepts_pending_but_credentials_wait_for_settlement() { + const CHILD: &str = "ASTRA_PENDING_STARTUP_TEST"; + if std::env::var_os(CHILD).is_none() { + let result = std::process::Command::new(std::env::current_exe().unwrap()) + .args(["--exact", "cli::native_auth::tests::startup_identity_accepts_pending_but_credentials_wait_for_settlement", "--nocapture"]) + .env(CHILD, "1").output().unwrap(); + assert!( + result.status.success(), + "{}\n{}", + String::from_utf8_lossy(&result.stdout), + String::from_utf8_lossy(&result.stderr) + ); + return; + } + let _home = crate::test_utils::HomeGuard::temp(); + let _credentials = crate::test_utils::isolate_credentials(); + use fs2::FileExt; + let root = tempfile::tempdir().unwrap(); + let store = NativeStore::with_directory(root.path().join("auth")); + let (session, _) = store.publish(test_session()).unwrap(); + let path = root.path().join("auth/auth.json"); + let lock_path = root + .path() + .join(format!("auth/refresh-{}.lock", session.environment.key())); + use std::os::unix::fs::OpenOptionsExt; + let lock = std::fs::OpenOptions::new() + .create(true) + .truncate(false) + .write(true) + .mode(0o600) + .open(lock_path) + .unwrap(); + lock.lock_exclusive().unwrap(); + let mut state: serde_json::Value = + serde_json::from_slice(&std::fs::read(&path).unwrap()).unwrap(); + state["sessions"][session.environment.key()]["refresh_pending"] = true.into(); + std::fs::write(&path, serde_json::to_vec(&state).unwrap()).unwrap(); + let binding = Arc::new(Binding { store, session }); + let _active = install_active_for_test(binding.clone()); + use crate::cli::cli_config::cli_utils::{self, CliProfileIdentityAdmission}; + cli_utils::configure_cli_profile_identity( + None, + CliProfileIdentityAdmission::RequireBoundAccount, + ) + .unwrap(); + assert!( + cli_utils::cli_owner_auth_snapshot() + .native_binding + .is_some() + ); + let last_session = uuid::Uuid::new_v4().to_string(); + let mut metadata = astra_credentials::CredentialsFile::default(); + metadata.profiles.insert( + binding.profile_name(), + astra_credentials::Profile { + last_session_id: Some(last_session.clone()), + ..Default::default() + }, + ); + cli_utils::save_credentials(&metadata).unwrap(); + let projection = projected_credentials().unwrap().unwrap(); + assert!( + projection.profiles[&binding.profile_name()] + .access_token + .is_none() + ); + assert_eq!(cli_utils::stored_last_session_id(None), Some(last_session)); + assert_eq!(binding.account_id().unwrap(), "astra-a"); + assert!( + binding.snapshot().is_err(), + "pending must not expose an access token" + ); + let pending = startup_access_token(&binding); + tokio::pin!(pending); + // The progress threshold must not cancel/restart the credential future. + assert!( + tokio::time::timeout(std::time::Duration::from_millis(1100), &mut pending) + .await + .is_err() + ); + state["sessions"][binding.session.environment.key()]["refresh_pending"] = false.into(); + state["sessions"][binding.session.environment.key()]["access_token"] = + "settled-access".into(); + std::fs::write(&path, serde_json::to_vec(&state).unwrap()).unwrap(); + FileExt::unlock(&lock).unwrap(); + assert_eq!(pending.await.unwrap(), "settled-access"); + + state["sessions"][binding.session.environment.key()]["refresh_pending"] = true.into(); + std::fs::write(&path, serde_json::to_vec(&state).unwrap()).unwrap(); + assert_eq!( + binding.access_token().await.unwrap_err(), + native::CredentialFailure::RotationInterrupted + ); + assert_eq!( + startup_access_token(&binding).await.unwrap_err(), + native::AccessMiss::ReauthenticationRequired.startup_warning() + ); + } + #[tokio::test] async fn cloud_sync_native_credentials_remain_bound_to_the_selected_origin() { // Native process authority is global. Exercise it in a fresh process so diff --git a/crates/astra-cli/src/cli/session/session_startup.rs b/crates/astra-cli/src/cli/session/session_startup.rs index 35a167a172..bbaf211efa 100644 --- a/crates/astra-cli/src/cli/session/session_startup.rs +++ b/crates/astra-cli/src/cli/session/session_startup.rs @@ -677,6 +677,13 @@ pub(crate) async fn complete_session_startup( no_instructions: bool, cli_context: &crate::cli::cli_config::cli_context::CliContext, ) -> Result { + // Resolve the native credential once, before cloud sync and memory startup + // can swallow its error or independently spend another lock-wait budget. + let native_startup_token = if let Some(binding) = crate::cli::native_auth::active() { + Some(crate::cli::native_auth::startup_access_token(&binding).await?) + } else { + None + }; // Install panic hook to write session_end on unexpected crashes. install_session_panic_hook(); // Install signal handlers so SIGTERM/SIGHUP can drain through normal REPL shutdown. @@ -791,7 +798,10 @@ pub(crate) async fn complete_session_startup( &pref_keys_after_pull, ); - let startup_token = session_runtime::fresh_access_token(api, profile).await; + let startup_token = match native_startup_token { + Some(token) => Some(token), + None => session_runtime::fresh_access_token(api, profile).await, + }; // Keep startup on the same model-selection state machine used by turns and // account commands. The local flag carries the one UI-specific outcome diff --git a/crates/astra-cli/tests/native_auth_environment.rs b/crates/astra-cli/tests/native_auth_environment.rs index a609d47962..2c0d69cd72 100644 --- a/crates/astra-cli/tests/native_auth_environment.rs +++ b/crates/astra-cli/tests/native_auth_environment.rs @@ -1,6 +1,80 @@ //! Native authority selection and legacy compatibility at the binary boundary. use std::process::Command; +#[tokio::test] +#[cfg(unix)] +async fn status_distinguishes_pending_owner_without_refresh_or_credential_disclosure() { + use astra_credentials::native::{Environment, NativeSession, NativeStore}; + use fs2::FileExt; + use std::os::unix::fs::OpenOptionsExt; + let root = tempfile::tempdir().unwrap(); + let store = NativeStore::with_directory(root.path().join(".moi")); + let issuer = "https://uc.example.test/realms/moi"; + let (session, _) = store + .publish(NativeSession { + environment: Environment { + issuer: issuer.into(), + astra_url: "https://astra.example.test".into(), + moi_url: "https://moi.example.test".into(), + authorization_endpoint: format!("{issuer}/protocol/openid-connect/auth"), + token_endpoint: format!("{issuer}/protocol/openid-connect/token"), + revocation_endpoint: format!("{issuer}/protocol/openid-connect/revoke"), + jwks_uri: format!("{issuer}/protocol/openid-connect/certs"), + }, + generation: String::new(), + subject: "user".into(), + session_id: "session".into(), + astra_user_id: "astra-user".into(), + moi_principal_id: "moi-user".into(), + catalog_user_id: "catalog-user".into(), + access_token: "private-access".into(), + refresh_token: "private-refresh".into(), + expires_at: 0, + workspace_id: None, + role_id: None, + refresh_pending: false, + }) + .unwrap(); + let path = root.path().join(".moi/auth.json"); + let mut state: serde_json::Value = + serde_json::from_slice(&std::fs::read(&path).unwrap()).unwrap(); + state["sessions"][session.environment.key()]["refresh_pending"] = true.into(); + let original = serde_json::to_vec(&state).unwrap(); + std::fs::write(&path, &original).unwrap(); + let lock = std::fs::OpenOptions::new() + .create(true) + .truncate(false) + .write(true) + .mode(0o600) + .open( + root.path() + .join(format!(".moi/refresh-{}.lock", session.environment.key())), + ) + .unwrap(); + lock.lock_exclusive().unwrap(); + for expected in ["refresh_in_progress", "reauthentication_required"] { + if expected == "reauthentication_required" { + FileExt::unlock(&lock).unwrap(); + } + let output = + client_output(isolated_client(root.path()).args(["auth", "status", "--json"])).await; + assert!( + output.status.success(), + "{}", + String::from_utf8_lossy(&output.stderr) + ); + let status: serde_json::Value = serde_json::from_slice(&output.stdout).unwrap(); + assert_eq!(status["state"], expected); + assert_eq!(std::fs::read(&path).unwrap(), original); + let public = format!( + "{}{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + assert!(!public.contains("private-access") && !public.contains("private-refresh")); + } +} + #[cfg(unix)] fn isolated_client(root: &std::path::Path) -> tokio::process::Command { let mut command = tokio::process::Command::new(env!("CARGO_BIN_EXE_astra")); diff --git a/crates/astra-credentials/src/native.rs b/crates/astra-credentials/src/native.rs index a2a7f44d90..60a74ed1be 100644 --- a/crates/astra-credentials/src/native.rs +++ b/crates/astra-credentials/src/native.rs @@ -2,6 +2,8 @@ //! profiles because rotation intent, account CAS and logout tombstones are //! part of this protocol, not optional legacy profile attributes. use fs2::FileExt; +mod diagnostics; +use diagnostics::{RotationDiagnostic, response_request_id, transport_cause}; use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; use std::{ @@ -125,6 +127,15 @@ pub struct NativeStore { root: PathBuf, } +/// Read-only local status; observing it never sends a refresh request. +#[derive(Debug, Clone, Copy, Serialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum SignInStatus { + SignedIn, + RefreshInProgress, + ReauthenticationRequired, +} + impl NativeStore { pub fn new() -> Result { let root = match std::env::var_os("MOI_AUTH_DIR") { @@ -543,6 +554,24 @@ pub enum AccessMiss { } impl AccessMiss { + /// Before the workbench exists, recovery commands must be shell commands. + pub fn startup_warning(self) -> &'static str { + match self { + Self::NotLoggedIn => "You are not signed in. Run astra login.", + Self::NetworkUnavailable => { + "Cannot connect to the sign-in service. Your saved sign-in is unchanged; retry when the network is back." + } + Self::ReauthenticationRequired => { + "Your sign-in could not be renewed. Run astra login to continue." + } + Self::RefreshInProgress => { + "Your sign-in is still being refreshed. Try starting Astra again shortly; you do not need to log in again." + } + Self::AccountChanged => "Your sign-in changed. Restart Astra.", + Self::Unavailable => "Cannot use your saved sign-in. Try starting Astra again.", + } + } + pub fn user_warning(self) -> &'static str { match self { Self::NotLoggedIn => " Not logged in. Use /login to authenticate.", @@ -613,6 +642,64 @@ impl NativeStore { self.blocking(|store| store.current()).await } + fn rotation_lock( + &self, + key: &str, + wait: std::time::Duration, + ) -> Result, CredentialFailure> { + let rotation = private_open(&self.root.join(format!("refresh-{key}.lock")), true) + .map_err(CredentialFailure::Unavailable)?; + let deadline = std::time::Instant::now() + wait; + loop { + match FileExt::try_lock_exclusive(&rotation) { + Ok(()) => return Ok(std::sync::Arc::new(rotation)), + Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => { + if std::time::Instant::now() >= deadline { + return Err(CredentialFailure::RefreshInProgress); + } + std::thread::sleep(std::time::Duration::from_millis(25)); + } + Err(_) => { + return Err(CredentialFailure::Unavailable( + "cannot lock native credential refresh".into(), + )); + } + } + } + } + + pub async fn status(&self) -> Result<(NativeSession, SignInStatus), String> { + self.blocking(|store| { + let previous = store.current()?; + if !previous.refresh_pending { + return Ok((previous, SignInStatus::SignedIn)); + } + // Hold the same lock used by credential() while re-reading pending. + // A free lock is not evidence that an old refresh token is reusable. + let rotation = + match store.rotation_lock(&previous.environment.key(), std::time::Duration::ZERO) { + Ok(lock) => Some(lock), + Err(CredentialFailure::RefreshInProgress) => None, + Err(error) => return Err(error.into()), + }; + let current = store.current()?; + if current.generation != previous.generation + || current.environment != previous.environment + { + return Err(CredentialFailure::AccountChanged.into()); + } + let status = if !current.refresh_pending { + SignInStatus::SignedIn + } else if rotation.is_some() { + SignInStatus::ReauthenticationRequired + } else { + SignInStatus::RefreshInProgress + }; + Ok((current, status)) + }) + .await + } + pub async fn credential( &self, target: &str, @@ -652,38 +739,38 @@ impl NativeStore { // logout/account changes on the short global state transaction. // Fresh credentials need only a shared state read, not the rotation // lock. Pending intent still has to wait for its in-flight owner. - let rotation = if frozen.expires_at + let diagnostic = if frozen.expires_at <= unix_now().map_err(CredentialFailure::Unavailable)? + 60 || frozen.refresh_pending { + Some(RotationDiagnostic::new(self, &frozen)) + } else { + None + }; + let rotation = if let Some(diagnostic) = &diagnostic { let key = frozen.environment.key(); - Some( - self.blocking(move |store| { - let rotation = - private_open(&store.root.join(format!("refresh-{key}.lock")), true)?; - let deadline = std::time::Instant::now() + lock_wait; - loop { - match FileExt::try_lock_exclusive(&rotation) { - Ok(()) => return Ok(std::sync::Arc::new(rotation)), - Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => { - if std::time::Instant::now() >= deadline { - return Err(CredentialFailure::RefreshInProgress.to_string()); - } - std::thread::sleep(std::time::Duration::from_millis(25)); - } - Err(_) => return Err("cannot lock native credential refresh".into()), - } - } - }) + diagnostic.record("lock", "waiting", None, None, None).await; + let result = self + .blocking(move |store| Ok(store.rotation_lock(&key, lock_wait))) .await - .map_err(|error| { - if error == CredentialFailure::RefreshInProgress.to_string() { - CredentialFailure::RefreshInProgress + .map_err(CredentialFailure::Unavailable)?; + match result { + Ok(rotation) => { + diagnostic + .record("lock", "acquired", None, None, None) + .await; + Some(rotation) + } + Err(error) => { + let reason = if error == CredentialFailure::RefreshInProgress { + "lock_wait_timeout" } else { - CredentialFailure::Unavailable(error) - } - })?, - ) + "lock_failed" + }; + diagnostic.record("lock", reason, None, None, None).await; + return Err(error); + } + } } else { None }; @@ -698,10 +785,18 @@ impl NativeStore { return Err(CredentialFailure::AccountChanged); } if current.refresh_pending { + if let Some(diagnostic) = &diagnostic { + diagnostic + .record("pending", "previous_rotation_unsettled", None, None, None) + .await; + } return Err(CredentialFailure::RotationInterrupted); } let now = unix_now().map_err(CredentialFailure::Unavailable)?; - if let Some(rotation) = rotation.filter(|_| current.expires_at <= now + 60) { + if let Some((rotation, diagnostic)) = rotation + .zip(diagnostic) + .filter(|_| current.expires_at <= now + 60) + { let client = http_client().map_err(CredentialFailure::Unavailable)?; let request = client .post(¤t.environment.token_endpoint) @@ -721,7 +816,7 @@ impl NativeStore { // not abandon an in-flight rotation or clear a token that may // already have been consumed. let settlement = tokio::spawn(async move { - settle_token_rotation(store, expected, client, request, rotation).await + settle_token_rotation(store, expected, client, request, rotation, diagnostic).await }); current = settlement.await.map_err(|_| { tracing::warn!("native token rotation task panicked before settlement"); @@ -756,10 +851,14 @@ async fn settle_token_rotation( client: reqwest::Client, request: reqwest::Request, rotation: std::sync::Arc, + diagnostic: RotationDiagnostic, ) -> Result { + diagnostic + .record("pending_write", "started", None, None, None) + .await; let write_guard = rotation.clone(); let pending_expected = expected.clone(); - store + let pending = store .blocking(move |store| { // spawn_blocking outlives cancellation of its awaiting task. // Retain rotation ownership until the durable write finishes. @@ -769,17 +868,34 @@ async fn settle_token_rotation( Ok(()) }) }) - .await - .map_err(CredentialFailure::Unavailable)?; + .await; + if let Err(error) = pending { + diagnostic + .record("pending_write", "storage_failed", None, None, Some(&error)) + .await; + return Err(CredentialFailure::Unavailable(error)); + } + diagnostic + .record("http", "request_started", None, None, None) + .await; // After sending, any failure may hide a successful rotation. Keep the // intent; never automatically replay the previous refresh token. Only a // connect error proves the request never reached the issuer. let response = match client.execute(request).await { Ok(response) => response, Err(error) if error.is_connect() => { + diagnostic + .record( + "http", + "connect_failed", + None, + None, + Some(&transport_cause(error)), + ) + .await; let expected = expected.clone(); let write_guard = rotation.clone(); - store + let cleared = store .blocking(move |store| { let _write_guard = write_guard; store.update(&expected, |session| { @@ -787,17 +903,61 @@ async fn settle_token_rotation( Ok(()) }) }) - .await - .map_err(CredentialFailure::Unavailable)?; + .await; + if let Err(error) = cleared { + diagnostic + .record("pending_clear", "storage_failed", None, None, Some(&error)) + .await; + return Err(CredentialFailure::Unavailable(error)); + } + diagnostic + .record( + "settled", + "connection_failure_preserved_session", + None, + None, + None, + ) + .await; return Err(log_rotation_failure(CredentialFailure::NetworkUnavailable)); } - Err(_) => return Err(log_rotation_failure(CredentialFailure::RotationUnconfirmed)), + Err(error) => { + let reason = if error.is_timeout() { + "request_timeout" + } else { + "transport_failed" + }; + diagnostic + .record("http", reason, None, None, Some(&transport_cause(error))) + .await; + return Err(log_rotation_failure(CredentialFailure::RotationUnconfirmed)); + } }; + let status = response.status().as_u16(); + let request_id = response_request_id(&response); if !response.status().is_success() { // Even a 5xx can be generated by a proxy or after the issuer // committed rotation. HTTP status is not proof of non-consumption. + diagnostic + .record( + "http", + "http_rejected", + Some(status), + request_id.as_deref(), + None, + ) + .await; return Err(log_rotation_failure(CredentialFailure::RotationRejected)); } + diagnostic + .record( + "response", + "received", + Some(status), + request_id.as_deref(), + None, + ) + .await; let token: TokenResponse = match bounded_json(response).await { Ok(token) => token, Err(error) => { @@ -805,6 +965,15 @@ async fn settle_token_rotation( // body is not proof the issuer rejected the refresh, and the old // refresh token must not be replayed. tracing::warn!(error = %error, "native token rotation response was unreadable"); + diagnostic + .record( + "response", + "invalid_response", + Some(status), + request_id.as_deref(), + Some(&error), + ) + .await; return Err(log_rotation_failure(CredentialFailure::RotationUnconfirmed)); } }; @@ -814,9 +983,36 @@ async fn settle_token_rotation( || token.expires_in > 86400 || !token.token_type.eq_ignore_ascii_case("bearer") { + diagnostic + .record( + "response", + "invalid_token_fields", + Some(status), + request_id.as_deref(), + None, + ) + .await; return Err(log_rotation_failure(CredentialFailure::RotationUnconfirmed)); } - if let Err(error) = verify_rotated_identity(&expected, &token.access_token).await { + diagnostic + .record( + "identity", + "verification_started", + Some(status), + request_id.as_deref(), + None, + ) + .await; + if let Err(error) = verify_rotated_identity(&expected, &token.access_token, &diagnostic).await { + diagnostic + .record( + "identity", + "verification_failed", + Some(status), + request_id.as_deref(), + Some(&error), + ) + .await; let _ = revoke(&expected.environment, &token.refresh_token).await; tracing::warn!(error = %error, "rotated token identity check failed"); return Err(log_rotation_failure(CredentialFailure::RotationUnconfirmed)); @@ -826,6 +1022,9 @@ async fn settle_token_rotation( let expires_in = token.expires_in; let write_guard = rotation; let saved = expected.clone(); + diagnostic + .record("save", "started", Some(status), request_id.as_deref(), None) + .await; if let Err(error) = store .blocking(move |store| { let _write_guard = write_guard; @@ -839,10 +1038,28 @@ async fn settle_token_rotation( }) .await { + diagnostic + .record( + "save", + "storage_failed", + Some(status), + request_id.as_deref(), + Some(&error), + ) + .await; let _ = revoke(&expected.environment, &token.refresh_token).await; tracing::warn!(cause = %error, "rotated token could not be saved"); return Err(log_rotation_failure(classify_post_rotation_save(error))); } + diagnostic + .record( + "settled", + "success", + Some(status), + request_id.as_deref(), + None, + ) + .await; store .blocking(|store| store.current()) .await @@ -864,7 +1081,11 @@ fn log_rotation_failure(error: CredentialFailure) -> CredentialFailure { error } -async fn verify_rotated_identity(current: &NativeSession, token: &str) -> Result<(), String> { +async fn verify_rotated_identity( + current: &NativeSession, + token: &str, + diagnostic: &RotationDiagnostic, +) -> Result<(), String> { use jsonwebtoken::{Algorithm, DecodingKey, Validation, decode, decode_header, jwk::JwkSet}; if token.len() > 16 * 1024 { return Err("rotated UC token exceeds size limit".into()); @@ -877,8 +1098,36 @@ async fn verify_rotated_identity(current: &NativeSession, token: &str) -> Result let response = http_client()? .get(¤t.environment.jwks_uri) .send() - .await - .map_err(|_| "UC signing keys unavailable")?; + .await; + let response = match response { + Ok(response) => response, + Err(error) => { + let reason = if error.is_timeout() { + "request_timeout" + } else { + "transport_failed" + }; + diagnostic + .record( + "signing_keys", + reason, + None, + None, + Some(&transport_cause(error)), + ) + .await; + return Err("UC signing keys unavailable".into()); + } + }; + diagnostic + .record( + "signing_keys", + "received", + Some(response.status().as_u16()), + response_request_id(&response).as_deref(), + None, + ) + .await; if !response.status().is_success() { return Err("UC signing keys unavailable".into()); } @@ -1051,6 +1300,131 @@ mod tests { } } + #[tokio::test] + async fn account_switch_while_waiting_for_refresh_does_not_return_new_accounts_token() { + let (_directory, store) = store(); + let (published, _) = store.publish(session("A")).unwrap(); + let lock = store + .rotation_lock(&published.environment.key(), std::time::Duration::ZERO) + .unwrap(); + store + .update(&published, |current| { + current.refresh_pending = true; + Ok(()) + }) + .unwrap(); + let pending = store.credential("astra", Some(&published.generation)); + tokio::pin!(pending); + assert!( + tokio::time::timeout(std::time::Duration::from_millis(80), &mut pending) + .await + .is_err() + ); + let (replacement, _) = store.publish(session("B")).unwrap(); + drop(lock); + assert_eq!( + pending.await.unwrap_err(), + CredentialFailure::AccountChanged + ); + assert_eq!(store.current().unwrap().generation, replacement.generation); + assert!(!store.current().unwrap().refresh_pending); + } + + #[tokio::test] + async fn status_observes_live_rotation_without_refreshing_or_requiring_login() { + use std::sync::atomic::Ordering; + let fixture = rotation_fixture("A").await; + let (_directory, store) = store(); + let mut expiring = session("A"); + expiring.environment = fixture.environment.clone(); + expiring.expires_at = unix_now().unwrap(); + store.publish(expiring).unwrap(); + assert_eq!(store.status().await.unwrap().1, SignInStatus::SignedIn); + assert_eq!(fixture.rotations.load(Ordering::SeqCst), 0); + let clone = store.clone(); + let helper = tokio::spawn(async move { clone.credential("astra", None).await }); + fixture.entered.notified().await; + assert_eq!( + store.status().await.unwrap().1, + SignInStatus::RefreshInProgress + ); + assert_eq!(fixture.rotations.load(Ordering::SeqCst), 1); + fixture.release.notify_one(); + helper.await.unwrap().unwrap(); + assert_eq!(store.status().await.unwrap().1, SignInStatus::SignedIn); + assert_eq!(fixture.rotations.load(Ordering::SeqCst), 1); + let log = std::fs::read_to_string(store.root.join("auth-refresh.jsonl")).unwrap(); + let events: Vec = log + .lines() + .map(|line| serde_json::from_str(line).unwrap()) + .collect(); + assert_eq!(events.last().unwrap()["reason"], "success"); + assert!(events.iter().any(|event| event["stage"] == "signing_keys")); + assert!(!log.contains("synthetic-refresh")); + assert!(!log.contains("synthetic-rotated")); + assert!(!log.contains(&store.current().unwrap().access_token)); + } + + #[test] + fn killed_refresh_owner_releases_lock_but_does_not_authorize_replay() { + const CHILD: &str = "ASTRA_REFRESH_CRASH_FIXTURE_DIR"; + if let Some(root) = std::env::var_os(CHILD) { + let store = NativeStore::with_directory(root.into()); + let session = store.current().unwrap(); + let _rotation = store + .rotation_lock(&session.environment.key(), std::time::Duration::ZERO) + .unwrap(); + store + .update(&session, |current| { + current.refresh_pending = true; + Ok(()) + }) + .unwrap(); + std::fs::write(store.root.join("child-ready"), b"ready").unwrap(); + loop { + std::thread::park(); + } + } + let (_directory, store) = store(); + store.publish(session("A")).unwrap(); + let mut child = std::process::Command::new(std::env::current_exe().unwrap()) + .args([ + "--exact", + "native::tests::killed_refresh_owner_releases_lock_but_does_not_authorize_replay", + ]) + .env(CHILD, &store.root) + .stdout(std::process::Stdio::null()) + .spawn() + .unwrap(); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + while !store.root.join("child-ready").exists() { + if std::time::Instant::now() > deadline { + let _ = child.kill(); + let _ = child.wait(); + panic!("child did not acquire refresh lock"); + } + std::thread::sleep(std::time::Duration::from_millis(10)); + } + child.kill().unwrap(); + child.wait().unwrap(); + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async { + assert_eq!( + store.status().await.unwrap().1, + SignInStatus::ReauthenticationRequired + ); + let result = tokio::time::timeout( + std::time::Duration::from_secs(1), + store.credential("astra", None), + ) + .await + .unwrap(); + assert_eq!(result.unwrap_err(), CredentialFailure::RotationInterrupted); + }); + assert!(store.current().unwrap().refresh_pending); + assert_eq!(store.current().unwrap().refresh_token, "synthetic-refresh"); + } + #[tokio::test] async fn concurrent_helpers_rotate_exactly_once_and_keep_the_command_generation() { use std::sync::atomic::Ordering; @@ -1582,7 +1956,7 @@ mod tests { String::from_utf8_lossy(&input[..len]) .starts_with("POST /realms/moi/protocol/openid-connect/token ") ); - stream.write_all(format!("HTTP/1.1 {status}\r\nContent-Length: 25\r\nConnection: close\r\n\r\n{{\"error\":\"invalid_grant\"}}").as_bytes()).await.unwrap(); + stream.write_all(format!("HTTP/1.1 {status}\r\nX-Request-ID: test-request-1\r\nContent-Length: 25\r\nConnection: close\r\n\r\n{{\"error\":\"invalid_grant\"}}").as_bytes()).await.unwrap(); }); let (_directory, store) = store(); let mut expiring = session("A"); @@ -1596,15 +1970,22 @@ mod tests { }; expiring.expires_at = unix_now().unwrap(); store.publish(expiring).unwrap(); - assert!( - store - .credential("moi", None) - .await - .unwrap_err() - .to_string() - .contains("rotation rejected") + assert_eq!( + store.credential("moi", None).await.unwrap_err(), + CredentialFailure::RotationRejected ); request.await.unwrap(); + let log = fs::read_to_string(store.root.join("auth-refresh.jsonl")).unwrap(); + let last: serde_json::Value = serde_json::from_str(log.lines().last().unwrap()).unwrap(); + assert_eq!(last["reason"], "http_rejected"); + assert_eq!(last["http_status"], status[..3].parse::().unwrap()); + assert_eq!(last["request_id"], "test-request-1"); + assert!( + !log.contains("invalid_grant"), + "response bodies are not diagnostics" + ); + assert!(!log.contains("synthetic-refresh")); + assert!(!log.contains("synthetic-access")); assert!(store.current().unwrap().refresh_pending); // No server remains; this must fail from persisted intent, not retry. assert!( diff --git a/crates/astra-credentials/src/native/diagnostics.rs b/crates/astra-credentials/src/native/diagnostics.rs new file mode 100644 index 0000000000..1037f5d292 --- /dev/null +++ b/crates/astra-credentials/src/native/diagnostics.rs @@ -0,0 +1,192 @@ +//! Bounded, private refresh diagnostics. Never used to decide credential state. +use super::{NativeSession, NativeStore, private_open, unix_now}; +use fs2::FileExt; +use serde::Serialize; +use std::{io::Write, time::Instant}; + +const MAX_BYTES: u64 = 64 * 1024; + +pub(super) struct RotationDiagnostic { + store: NativeStore, + operation_id: String, + environment: String, + generation: String, + started: Instant, +} + +#[derive(Serialize)] +struct Event<'a> { + component: &'static str, + operation: &'static str, + operation_id: &'a str, + environment: &'a str, + generation: &'a str, + pid: u32, + timestamp: i64, + elapsed_ms: u128, + stage: &'static str, + reason: &'static str, + http_status: Option, + request_id: Option<&'a str>, + cause: Option<&'a str>, +} + +impl RotationDiagnostic { + pub(super) fn new(store: &NativeStore, session: &NativeSession) -> Self { + Self { + store: store.clone(), + operation_id: uuid::Uuid::new_v4().to_string(), + environment: session.environment.key(), + generation: session.generation.chars().take(128).collect(), + started: Instant::now(), + } + } + + // Callers supply only local classifications/sanitized errors, never token + // responses, headers other than the bounded request ID, or credentials. + pub(super) async fn record( + &self, + stage: &'static str, + reason: &'static str, + http_status: Option, + request_id: Option<&str>, + cause: Option<&str>, + ) { + let cause = cause.map(|value| value.chars().take(1024).collect::()); + let event = Event { + component: "astra-credentials", + operation: "token_refresh", + operation_id: &self.operation_id, + environment: &self.environment, + generation: &self.generation, + pid: std::process::id(), + timestamp: unix_now().unwrap_or_default(), + elapsed_ms: self.started.elapsed().as_millis(), + stage, + reason, + http_status, + request_id, + cause: cause.as_deref(), + }; + tracing::info!(component = event.component, operation = event.operation, + operation_id = %event.operation_id, environment = %event.environment, + generation = %event.generation, pid = event.pid, stage, reason, + elapsed_ms = %event.elapsed_ms, http_status, request_id, cause = event.cause, + "native sign-in refresh"); + let Ok(mut line) = serde_json::to_vec(&event) else { + return; + }; + line.push(b'\n'); + // Logging is best effort: inability to record diagnostics must not + // prevent settlement. File I/O stays off the async runtime, and a busy + // diagnostic writer never delays authentication on a separate lock. + if let Err(error) = self + .store + .blocking(move |store| { + store.check_dir()?; + let mut file = private_open(&store.root.join("auth-refresh.jsonl"), true)?; + match FileExt::try_lock_exclusive(&file) { + Ok(()) => (), + Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => return Ok(()), + Err(_) => return Err("cannot lock sign-in diagnostic file".into()), + } + if file + .metadata() + .map_err(|_| "cannot inspect sign-in diagnostics")? + .len() + + line.len() as u64 + > MAX_BYTES + { + file.set_len(0) + .map_err(|_| "cannot truncate sign-in diagnostics")?; + } + use std::io::{Seek, SeekFrom}; + file.seek(SeekFrom::End(0)) + .map_err(|_| "cannot seek sign-in diagnostics")?; + file.write_all(&line) + .map_err(|_| "cannot write sign-in diagnostics".into()) + }) + .await + { + tracing::warn!(component = "astra-credentials", operation = "token_refresh", + stage = "diagnostics", reason = "diagnostic_write_failed", %error, + "could not record sign-in diagnostics"); + } + } +} + +pub(super) fn response_request_id(response: &reqwest::Response) -> Option { + response + .headers() + .get("x-request-id")? + .to_str() + .ok() + .filter(|value| { + !value.is_empty() + && value.len() <= 128 + && value + .bytes() + .all(|b| b.is_ascii_alphanumeric() || b"-_.:".contains(&b)) + }) + .map(str::to_owned) +} + +pub(super) fn transport_cause(error: reqwest::Error) -> String { + use std::error::Error; + let error = error.without_url(); + let mut cause = error.to_string(); + let mut source = error.source(); + while let Some(next) = source { + if cause.len() >= 1024 { + break; + } + cause.push_str(": "); + cause.extend(next.to_string().chars().take(1024)); + source = next.source(); + } + cause +} + +#[cfg(all(test, unix))] +mod tests { + use super::*; + use std::os::unix::fs::{MetadataExt, symlink}; + + #[tokio::test] + async fn default_diagnostics_are_private_bounded_and_do_not_serialize_credentials() { + let directory = tempfile::tempdir().unwrap(); + let store = NativeStore::with_directory(directory.path().join("auth")); + store.prepare_for_login().unwrap(); + let diagnostic = RotationDiagnostic { + store: store.clone(), + operation_id: "test-operation".into(), + environment: "environment-digest".into(), + generation: "generation".into(), + started: Instant::now(), + }; + let path = store.root.join("auth-refresh.jsonl"); + let mut file = private_open(&path, true).unwrap(); + file.write_all(&vec![b' '; MAX_BYTES as usize]).unwrap(); + diagnostic + .record("http", "http_rejected", Some(400), Some("request-1"), None) + .await; + let contents = std::fs::read(&path).unwrap(); + let event: serde_json::Value = serde_json::from_slice(&contents).unwrap(); + assert_eq!(event["http_status"], 400); + assert_eq!(event["request_id"], "request-1"); + assert_eq!(event["reason"], "http_rejected"); + assert!(contents.len() < MAX_BYTES as usize); + assert_eq!(std::fs::metadata(&path).unwrap().mode() & 0o777, 0o600); + assert!(event.get("access_token").is_none()); + assert!(event.get("refresh_token").is_none()); + + std::fs::remove_file(&path).unwrap(); + let target = directory.path().join("untouched"); + std::fs::write(&target, "unchanged").unwrap(); + symlink(&target, &path).unwrap(); + diagnostic + .record("http", "http_rejected", Some(400), None, None) + .await; + assert_eq!(std::fs::read_to_string(&target).unwrap(), "unchanged"); + } +} diff --git a/docs/guides/moi-native-login.md b/docs/guides/moi-native-login.md index cba9868676..f7dd3415d9 100644 --- a/docs/guides/moi-native-login.md +++ b/docs/guides/moi-native-login.md @@ -55,6 +55,35 @@ set `MOI_ASTRA_BINARY` to its absolute path. Commands freeze account, environmen and login generation; an account switch cannot silently redirect an in-flight operation. An uncertain refresh rotation requires a new login, not refresh replay. +Startup binds the saved identity independently of token availability. When a +refresh is already in progress, credential acquisition waits on the existing +per-environment lock and re-reads the saved session after acquiring it. It does +not submit a second rotation. Interactive startup shows `Refreshing your +sign-in…` after one second of waiting. The progress timer does not cancel the +refresh. The existing 20-second lock wait and per-request HTTP timeouts still +apply; 20 seconds is not an end-to-end startup deadline. A lock-wait timeout +asks the user to retry starting Astra, not to sign in again. + +`astra auth status` never refreshes tokens. While a pending refresh still owns +the lock, its JSON state is `refresh_in_progress`. If the lock can be acquired +and the re-read session still has pending intent, it reports +`reauthentication_required`. Process exit releases the OS lock, but does not +erase pending intent: the issuer may already have consumed the old token. +Account identity and local resume metadata remain readable while pending; +ordinary token snapshots do not expose that session's access token. + +Refresh attempts record bounded diagnostics by default in +`~/.moi/auth-refresh.jsonl` (or the selected `MOI_AUTH_DIR`). The file is private +(0600), limited to 64 KiB, and restarts from empty when the next record would +exceed that limit. Records contain a local operation ID, process ID, environment +digest, login generation, timestamp, elapsed time, stage, error classification, +HTTP status and a bounded request ID when supplied. They never include tokens +or response bodies. Diagnostic write failures do not change refresh outcomes; +the file is evidence only and is never used to authorize recovery. A killed +process may leave an incomplete attempt. Rejected/uncertain rotations still +require a new login; these startup changes do not extend the issuer's session +lifetime or make an uncertain refresh token safe to replay. + Cloud preference, outbox and resume traffic use the same selected native endpoint and generation-bound refreshing credential as chat. A conflicting `ASTRA_API_URL` cannot redirect that credential, and omitting the environment From c2861fc0eceb16f27822de503cf30b4738f35497 Mon Sep 17 00:00:00 2001 From: lr90 Date: Tue, 29 Sep 2026 11:51:59 +0800 Subject: [PATCH 2/3] fix(auth): keep CLI startup available during refresh failures --- .../src/cli/session/session_startup.rs | 64 ++++++++--- crates/astra-credentials/src/native.rs | 45 ++++++-- .../src/native/diagnostics.rs | 101 ++++++++++++++---- docs/guides/moi-native-login.md | 9 +- 4 files changed, 173 insertions(+), 46 deletions(-) diff --git a/crates/astra-cli/src/cli/session/session_startup.rs b/crates/astra-cli/src/cli/session/session_startup.rs index bbaf211efa..ebb33228a4 100644 --- a/crates/astra-cli/src/cli/session/session_startup.rs +++ b/crates/astra-cli/src/cli/session/session_startup.rs @@ -6,7 +6,7 @@ use crate::cli::cli_config::cli_utils::{ preflight_remote_resume_session, }; use crate::cli::cloud_sync::{ - append_cloud_pull_sync_journal, try_cloud_pull, try_cloud_pull_preferences, + CloudPullResult, append_cloud_pull_sync_journal, try_cloud_pull, try_cloud_pull_preferences, }; use crate::cli::edge_lifecycle::register_and_start_heartbeat; use crate::cli::permission_manager; @@ -677,18 +677,26 @@ pub(crate) async fn complete_session_startup( no_instructions: bool, cli_context: &crate::cli::cli_config::cli_context::CliContext, ) -> Result { - // Resolve the native credential once, before cloud sync and memory startup - // can swallow its error or independently spend another lock-wait budget. - let native_startup_token = if let Some(binding) = crate::cli::native_auth::active() { - Some(crate::cli::native_auth::startup_access_token(&binding).await?) - } else { - None - }; // Install panic hook to write session_end on unexpected crashes. install_session_panic_hook(); // Install signal handlers so SIGTERM/SIGHUP can drain through normal REPL shutdown. install_sigterm_handler(); let shutdown_signal_rx = subscribe_shutdown_signal(); + // Resolve native auth before optional cloud work. An unavailable issuer + // should leave the local workbench usable and give the user a clear hint. + let native_binding = crate::cli::native_auth::active(); + let native_startup_token = if let Some(binding) = native_binding.as_ref() { + match crate::cli::native_auth::startup_access_token(binding).await { + Ok(token) => Some(token), + Err(message) => { + eprintln!("warning: {message}"); + None + } + } + } else { + None + }; + let native_auth_unavailable = native_binding.is_some() && native_startup_token.is_none(); // --session-id: override with explicit session UUID if let Some(sid) = cli_context.session_id.as_deref() { @@ -769,8 +777,18 @@ pub(crate) async fn complete_session_startup( astra_turn_core::tool_health_persistence::load_tool_health(profile_name); state.synced_tool_health_entries = astra_turn_core::tool_health_persistence::load_synced_tool_health(profile_name); - let cloud_pull_result = try_cloud_pull(profile_name).await; - let pref_keys = try_cloud_pull_preferences(state).await; + let cloud_pull_result = if native_auth_unavailable { + CloudPullResult { + cloud_reachable: false, + } + } else { + try_cloud_pull(profile_name).await + }; + let pref_keys = if native_auth_unavailable { + Vec::new() + } else { + try_cloud_pull_preferences(state).await + }; (cross_session_health_entries, cloud_pull_result, pref_keys) }; tracer.phase("learning_state"); @@ -780,7 +798,11 @@ pub(crate) async fn complete_session_startup( state.synced_tool_health_entries = cross_session_health_entries; } - state.session_memory_extractor = build_cli_session_memory_extractor(api, profile).await; + state.session_memory_extractor = if native_auth_unavailable { + None + } else { + build_cli_session_memory_extractor(api, profile).await + }; state.team_store = std::sync::Arc::new(crate::cli::http_team_store::HttpTeamStore::new( api.api_origin(), profile, @@ -798,9 +820,19 @@ pub(crate) async fn complete_session_startup( &pref_keys_after_pull, ); - let startup_token = match native_startup_token { - Some(token) => Some(token), - None => session_runtime::fresh_access_token(api, profile).await, + let startup_token = match (native_binding.as_ref(), native_startup_token) { + (Some(binding), Some(_)) => { + // The early token may be near expiry after other startup work. + match crate::cli::native_auth::startup_access_token(binding).await { + Ok(token) => Some(token), + Err(message) => { + eprintln!("warning: {message}"); + None + } + } + } + (Some(_), None) => None, + (None, _) => session_runtime::fresh_access_token(api, profile).await, }; // Keep startup on the same model-selection state machine used by turns and @@ -819,7 +851,9 @@ pub(crate) async fn complete_session_startup( false }; tracer.phase("model_check"); - prune_stale_pending_recovery(api, profile, state).await; + if native_binding.is_none() || startup_token.is_some() { + prune_stale_pending_recovery(api, profile, state).await; + } if state.session_id.is_none() && let Some(sid) = resume_session_id diff --git a/crates/astra-credentials/src/native.rs b/crates/astra-credentials/src/native.rs index 60a74ed1be..a4f335e36b 100644 --- a/crates/astra-credentials/src/native.rs +++ b/crates/astra-credentials/src/native.rs @@ -1216,6 +1216,30 @@ mod tests { use super::*; use std::os::unix::fs::{PermissionsExt, symlink}; + async fn wait_for_diagnostics(store: &NativeStore, required: &[(&str, &str)]) -> String { + let path = store.root.join("auth-refresh.jsonl"); + tokio::time::timeout(std::time::Duration::from_secs(2), async { + loop { + let log = fs::read_to_string(&path).unwrap_or_default(); + if let Ok(events) = log + .lines() + .map(serde_json::from_str::) + .collect::, _>>() + && required.iter().all(|(stage, reason)| { + events + .iter() + .any(|event| event["stage"] == *stage && event["reason"] == *reason) + }) + { + return log; + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + }) + .await + .expect("refresh diagnostic was not written") + } + struct RotationFixture { environment: Environment, rotations: std::sync::Arc, @@ -1353,12 +1377,16 @@ mod tests { helper.await.unwrap().unwrap(); assert_eq!(store.status().await.unwrap().1, SignInStatus::SignedIn); assert_eq!(fixture.rotations.load(Ordering::SeqCst), 1); - let log = std::fs::read_to_string(store.root.join("auth-refresh.jsonl")).unwrap(); + let log = wait_for_diagnostics( + &store, + &[("signing_keys", "received"), ("settled", "success")], + ) + .await; let events: Vec = log .lines() .map(|line| serde_json::from_str(line).unwrap()) .collect(); - assert_eq!(events.last().unwrap()["reason"], "success"); + assert!(events.iter().any(|event| event["reason"] == "success")); assert!(events.iter().any(|event| event["stage"] == "signing_keys")); assert!(!log.contains("synthetic-refresh")); assert!(!log.contains("synthetic-rotated")); @@ -1975,11 +2003,14 @@ mod tests { CredentialFailure::RotationRejected ); request.await.unwrap(); - let log = fs::read_to_string(store.root.join("auth-refresh.jsonl")).unwrap(); - let last: serde_json::Value = serde_json::from_str(log.lines().last().unwrap()).unwrap(); - assert_eq!(last["reason"], "http_rejected"); - assert_eq!(last["http_status"], status[..3].parse::().unwrap()); - assert_eq!(last["request_id"], "test-request-1"); + let log = wait_for_diagnostics(&store, &[("http", "http_rejected")]).await; + let rejected: serde_json::Value = log + .lines() + .map(|line| serde_json::from_str(line).unwrap()) + .find(|event: &serde_json::Value| event["reason"] == "http_rejected") + .unwrap(); + assert_eq!(rejected["http_status"], status[..3].parse::().unwrap()); + assert_eq!(rejected["request_id"], "test-request-1"); assert!( !log.contains("invalid_grant"), "response bodies are not diagnostics" diff --git a/crates/astra-credentials/src/native/diagnostics.rs b/crates/astra-credentials/src/native/diagnostics.rs index 1037f5d292..f5e5a5992d 100644 --- a/crates/astra-credentials/src/native/diagnostics.rs +++ b/crates/astra-credentials/src/native/diagnostics.rs @@ -77,18 +77,31 @@ impl RotationDiagnostic { return; }; line.push(b'\n'); - // Logging is best effort: inability to record diagnostics must not - // prevent settlement. File I/O stays off the async runtime, and a busy - // diagnostic writer never delays authentication on a separate lock. - if let Err(error) = self - .store - .blocking(move |store| { + // The future completes without yielding. A stalled diagnostic file + // must not delay the request or the durable publication of new tokens. + let store = self.store.clone(); + tokio::task::spawn_blocking(move || { + let write = || -> Result<(), String> { store.check_dir()?; let mut file = private_open(&store.root.join("auth-refresh.jsonl"), true)?; - match FileExt::try_lock_exclusive(&file) { - Ok(()) => (), - Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => return Ok(()), - Err(_) => return Err("cannot lock sign-in diagnostic file".into()), + // Detached records from the same rotation can briefly race. + // Give them a bounded chance to serialize without waiting in + // the credential path or hanging on another process's lock. + let lock_deadline = Instant::now() + std::time::Duration::from_millis(250); + loop { + match FileExt::try_lock_exclusive(&file) { + Ok(()) => break, + Err(error) + if error.kind() == std::io::ErrorKind::WouldBlock + && Instant::now() < lock_deadline => + { + std::thread::sleep(std::time::Duration::from_millis(2)); + } + Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => { + return Ok(()); + } + Err(_) => return Err("cannot lock sign-in diagnostic file".into()), + } } if file .metadata() @@ -105,13 +118,13 @@ impl RotationDiagnostic { .map_err(|_| "cannot seek sign-in diagnostics")?; file.write_all(&line) .map_err(|_| "cannot write sign-in diagnostics".into()) - }) - .await - { - tracing::warn!(component = "astra-credentials", operation = "token_refresh", - stage = "diagnostics", reason = "diagnostic_write_failed", %error, - "could not record sign-in diagnostics"); - } + }; + if let Err(error) = write() { + tracing::warn!(component = "astra-credentials", operation = "token_refresh", + stage = "diagnostics", reason = "diagnostic_write_failed", %error, + "could not record sign-in diagnostics"); + } + }); } } @@ -150,7 +163,43 @@ pub(super) fn transport_cause(error: reqwest::Error) -> String { #[cfg(all(test, unix))] mod tests { use super::*; - use std::os::unix::fs::{MetadataExt, symlink}; + use std::{ + future::Future, + os::unix::fs::{MetadataExt, symlink}, + }; + + #[tokio::test] + async fn a_locked_diagnostic_file_does_not_delay_record() { + let directory = tempfile::tempdir().unwrap(); + let store = NativeStore::with_directory(directory.path().join("auth")); + store.prepare_for_login().unwrap(); + let diagnostic = RotationDiagnostic { + store: store.clone(), + operation_id: "test-operation".into(), + environment: "environment-digest".into(), + generation: "generation".into(), + started: Instant::now(), + }; + let path = store.root.join("auth-refresh.jsonl"); + let file = private_open(&path, true).unwrap(); + FileExt::lock_exclusive(&file).unwrap(); + let mut record = Box::pin(diagnostic.record("http", "request_started", None, None, None)); + let mut context = std::task::Context::from_waker(std::task::Waker::noop()); + assert!(record.as_mut().poll(&mut context).is_ready()); + FileExt::unlock(&file).unwrap(); + tokio::time::timeout(std::time::Duration::from_secs(2), async { + loop { + if std::fs::read_to_string(&path) + .is_ok_and(|contents| contents.contains("request_started")) + { + break; + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + }) + .await + .unwrap(); + } #[tokio::test] async fn default_diagnostics_are_private_bounded_and_do_not_serialize_credentials() { @@ -170,7 +219,17 @@ mod tests { diagnostic .record("http", "http_rejected", Some(400), Some("request-1"), None) .await; - let contents = std::fs::read(&path).unwrap(); + let contents = tokio::time::timeout(std::time::Duration::from_secs(2), async { + loop { + let contents = std::fs::read(&path).unwrap(); + if serde_json::from_slice::(&contents).is_ok() { + break contents; + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + }) + .await + .unwrap(); let event: serde_json::Value = serde_json::from_slice(&contents).unwrap(); assert_eq!(event["http_status"], 400); assert_eq!(event["request_id"], "request-1"); @@ -184,9 +243,7 @@ mod tests { let target = directory.path().join("untouched"); std::fs::write(&target, "unchanged").unwrap(); symlink(&target, &path).unwrap(); - diagnostic - .record("http", "http_rejected", Some(400), None, None) - .await; + assert!(private_open(&path, true).is_err()); assert_eq!(std::fs::read_to_string(&target).unwrap(), "unchanged"); } } diff --git a/docs/guides/moi-native-login.md b/docs/guides/moi-native-login.md index f7dd3415d9..ccc982a94c 100644 --- a/docs/guides/moi-native-login.md +++ b/docs/guides/moi-native-login.md @@ -62,7 +62,10 @@ not submit a second rotation. Interactive startup shows `Refreshing your sign-in…` after one second of waiting. The progress timer does not cancel the refresh. The existing 20-second lock wait and per-request HTTP timeouts still apply; 20 seconds is not an end-to-end startup deadline. A lock-wait timeout -asks the user to retry starting Astra, not to sign in again. +asks the user to retry starting Astra, not to sign in again. If credential +acquisition fails, interactive Astra still opens with a warning and skips +native cloud initialization for that startup. MOI requests still require a +valid credential; this does not bypass authentication. `astra auth status` never refreshes tokens. While a pending refresh still owns the lock, its JSON state is `refresh_in_progress`. If the lock can be acquired @@ -80,7 +83,9 @@ digest, login generation, timestamp, elapsed time, stage, error classification, HTTP status and a bounded request ID when supplied. They never include tokens or response bodies. Diagnostic write failures do not change refresh outcomes; the file is evidence only and is never used to authorize recovery. A killed -process may leave an incomplete attempt. Rejected/uncertain rotations still +process may leave an incomplete attempt or lose a queued diagnostic record; +diagnostic file writes do not delay token settlement, and records from the same +attempt may arrive out of order. Rejected/uncertain rotations still require a new login; these startup changes do not extend the issuer's session lifetime or make an uncertain refresh token safe to replay. From 77c7b0378464396bfd7a6cddd8bec305184cd388 Mon Sep 17 00:00:00 2001 From: lr90 Date: Tue, 29 Sep 2026 14:54:18 +0800 Subject: [PATCH 3/3] fix(auth): preserve refresh state in startup banner --- crates/astra-cli/src/cli/native_auth.rs | 12 ++- .../src/cli/session/session_runtime.rs | 91 ++++++++++++++++--- .../src/cli/session/session_startup.rs | 41 +++++---- docs/guides/moi-native-login.md | 5 +- 4 files changed, 114 insertions(+), 35 deletions(-) diff --git a/crates/astra-cli/src/cli/native_auth.rs b/crates/astra-cli/src/cli/native_auth.rs index 8ca60b1314..a018f5a51c 100644 --- a/crates/astra-cli/src/cli/native_auth.rs +++ b/crates/astra-cli/src/cli/native_auth.rs @@ -104,7 +104,7 @@ impl Binding { /// Resolve native auth before startup launches cloud work. The progress timer /// never cancels the credential future or starts another refresh. -pub(crate) async fn startup_access_token(binding: &Binding) -> Result { +pub(crate) async fn startup_access_token(binding: &Binding) -> Result { let pending = binding.access_token(); tokio::pin!(pending); let result = tokio::select! { @@ -115,7 +115,7 @@ pub(crate) async fn startup_access_token(binding: &Binding) -> Result String { result } -pub(crate) fn print_session_banner(profile: Option<&str>, state: &SessionState) { +pub(crate) fn print_session_banner( + profile: Option<&str>, + state: &SessionState, + native_auth: Option>, +) { let creds = load_credentials(); let pname = profile_name(profile, &creds); let p = creds.profiles.get(&pname); - let logged_in = p.and_then(|p| p.access_token.as_ref()).is_some() + let legacy_logged_in = p.and_then(|p| p.access_token.as_ref()).is_some() || active_env_access_token(chrono::Utc::now().timestamp()).is_some(); + let (auth_status, welcome_back) = match native_auth { + Some(Ok(())) => ("logged in", true), + Some(Err(AccessMiss::RefreshInProgress)) => ("sign-in refreshing", p.is_some()), + Some(Err(AccessMiss::ReauthenticationRequired)) => ("sign-in needs renewal", p.is_some()), + Some(Err(AccessMiss::NetworkUnavailable | AccessMiss::Unavailable)) => { + ("sign-in unavailable", p.is_some()) + } + Some(Err(AccessMiss::AccountChanged)) => ("sign-in changed", p.is_some()), + Some(Err(AccessMiss::NotLoggedIn)) => ("not logged in", false), + None if legacy_logged_in => ("logged in", true), + None => ("not logged in", false), + }; let model_display = state.model.as_deref().unwrap_or("auto"); let version = env!("CARGO_PKG_VERSION"); let skills_count = state.unified_skill_registry.len(); @@ -1848,14 +1864,7 @@ pub(crate) fn print_session_banner(profile: Option<&str>, state: &SessionState) )); right.push(style_banner_text( trunc_vis( - &format!( - "{skills_count} skills · {}", - if logged_in { - "logged in" - } else { - "not logged in" - } - ), + &format!("{skills_count} skills · {auth_status}"), right_col_w, ), BannerTextStyle::Body, @@ -2055,7 +2064,7 @@ pub(crate) fn print_session_banner(profile: Option<&str>, state: &SessionState) } eprintln!(); - let welcome = banner_welcome_text(p, logged_in); + let welcome = banner_welcome_text(p, welcome_back); let model_hint = if model_display == "auto" { format!( "{} {}", @@ -2173,15 +2182,73 @@ mod tests { if std::env::var_os("ASTRA_TEST_BANNER_CHILD").is_none() { return; } + let native_auth = match std::env::var("ASTRA_TEST_BANNER_AUTH").as_deref() { + Ok("refreshing") => Some(Err(super::AccessMiss::RefreshInProgress)), + Ok("renewal") => Some(Err(super::AccessMiss::ReauthenticationRequired)), + Ok("signed_in") => Some(Ok(())), + _ => None, + }; + if native_auth.is_some() { + let mut credentials = astra_credentials::CredentialsFile::default(); + credentials.profiles.insert( + "default".to_string(), + astra_credentials::Profile { + username: Some("test-user".to_string()), + ..Default::default() + }, + ); + crate::cli::cli_config::cli_utils::save_credentials(&credentials).unwrap(); + } + let long_profile = "moi-internal-profile-".repeat(10); + let profile = if native_auth.is_some() { + "default" + } else { + &long_profile + }; super::print_session_banner( - Some(&"moi-internal-profile-".repeat(10)), + Some(profile), &super::SessionState { model: Some("模型-deepseek-".repeat(10)), ..Default::default() }, + native_auth, ); } + #[cfg(unix)] + #[test] + fn banner_distinguishes_refreshing_from_reauthentication() { + use std::process::Command; + + for (auth, expected) in [ + ("refreshing", "sign-in refreshing"), + ("renewal", "sign-in needs renewal"), + ] { + let home = tempfile::tempdir().unwrap(); + let output = Command::new(std::env::current_exe().unwrap()) + .args([ + "--exact", + "cli::session::session_runtime::tests::banner_pty_child", + "--nocapture", + ]) + .env("ASTRA_TEST_BANNER_CHILD", "1") + .env("ASTRA_TEST_BANNER_AUTH", auth) + .env("ASTRA_CLI_CREDENTIALS_DIR", home.path()) + .env("HOME", home.path()) + .env("NO_COLOR", "1") + .output() + .unwrap(); + assert!(output.status.success(), "{output:?}"); + let visible = String::from_utf8_lossy(&output.stderr); + assert!(visible.contains(expected), "{auth}: {visible}"); + assert!(!visible.contains("not logged in"), "{auth}: {visible}"); + assert!( + visible.contains("Welcome back, test-user"), + "{auth}: {visible}" + ); + } + } + #[cfg(unix)] #[test] fn banner_animation_leaves_one_card_in_real_terminal() { diff --git a/crates/astra-cli/src/cli/session/session_startup.rs b/crates/astra-cli/src/cli/session/session_startup.rs index ebb33228a4..a812890826 100644 --- a/crates/astra-cli/src/cli/session/session_startup.rs +++ b/crates/astra-cli/src/cli/session/session_startup.rs @@ -685,18 +685,16 @@ pub(crate) async fn complete_session_startup( // Resolve native auth before optional cloud work. An unavailable issuer // should leave the local workbench usable and give the user a clear hint. let native_binding = crate::cli::native_auth::active(); - let native_startup_token = if let Some(binding) = native_binding.as_ref() { - match crate::cli::native_auth::startup_access_token(binding).await { - Ok(token) => Some(token), - Err(message) => { - eprintln!("warning: {message}"); - None - } + let native_startup_auth = if let Some(binding) = native_binding.as_ref() { + let result = crate::cli::native_auth::startup_access_token(binding).await; + if let Err(miss) = &result { + eprintln!("warning: {}", miss.startup_warning()); } + Some(result) } else { None }; - let native_auth_unavailable = native_binding.is_some() && native_startup_token.is_none(); + let native_auth_unavailable = native_startup_auth.as_ref().is_some_and(Result::is_err); // --session-id: override with explicit session UUID if let Some(sid) = cli_context.session_id.as_deref() { @@ -820,19 +818,24 @@ pub(crate) async fn complete_session_startup( &pref_keys_after_pull, ); - let startup_token = match (native_binding.as_ref(), native_startup_token) { - (Some(binding), Some(_)) => { + let native_final_auth = match (native_binding.as_ref(), native_startup_auth) { + (Some(binding), Some(Ok(_))) => { // The early token may be near expiry after other startup work. - match crate::cli::native_auth::startup_access_token(binding).await { - Ok(token) => Some(token), - Err(message) => { - eprintln!("warning: {message}"); - None - } + let result = crate::cli::native_auth::startup_access_token(binding).await; + if let Err(miss) = &result { + eprintln!("warning: {}", miss.startup_warning()); } + Some(result) } - (Some(_), None) => None, - (None, _) => session_runtime::fresh_access_token(api, profile).await, + (_, result) => result, + }; + let banner_native_auth = native_final_auth + .as_ref() + .map(|result| result.as_ref().map(|_| ()).map_err(|miss| *miss)); + let startup_token = match native_final_auth { + Some(Ok(token)) => Some(token), + Some(Err(_)) => None, + None => session_runtime::fresh_access_token(api, profile).await, }; // Keep startup on the same model-selection state machine used by turns and @@ -861,7 +864,7 @@ pub(crate) async fn complete_session_startup( slash_session::restore_session_into_state(sid, profile, api, state).await?; } - print_session_banner(profile, state); + print_session_banner(profile, state, banner_native_auth); tracer.phase("banner"); // Pending recovery is silently retained in state for /resume; no startup banner. diff --git a/docs/guides/moi-native-login.md b/docs/guides/moi-native-login.md index ccc982a94c..d128414c90 100644 --- a/docs/guides/moi-native-login.md +++ b/docs/guides/moi-native-login.md @@ -65,7 +65,10 @@ apply; 20 seconds is not an end-to-end startup deadline. A lock-wait timeout asks the user to retry starting Astra, not to sign in again. If credential acquisition fails, interactive Astra still opens with a warning and skips native cloud initialization for that startup. MOI requests still require a -valid credential; this does not bypass authentication. +valid credential; this does not bypass authentication. The startup card +distinguishes a refresh still owned by another process from a rotation that +needs a new login; neither is shown as a logout merely because its unsettled +access token is withheld. `astra auth status` never refreshes tokens. While a pending refresh still owns the lock, its JSON state is `refresh_in_progress`. If the lock can be acquired