From 918680aa9edecad5014b355756cd37123803e2a5 Mon Sep 17 00:00:00 2001 From: "@permadeath.com" Date: Wed, 26 Aug 2026 14:31:06 -0400 Subject: [PATCH] fix(config): wait for the account lock without parking a worker thread MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `registry::update` took the lock with the blocking `take`, justified by a doc comment claiming its callers were sync. Every one of them is reached from an `async fn`, so a contended registry write parked a tokio worker for up to thirty seconds — very often the thread that would run the peer's token refresh. `take` had no other caller and is gone, leaving `take_async` as the one way in. Change-Id: I3ebd16fb59e839ebb46bf22af1a453fe4198e141 --- src/cmd/auth.rs | 13 ++--- src/config/account/registry.rs | 30 ++++++----- src/config/account/selection.rs | 2 +- src/config/lock.rs | 88 ++++++++++++++++++++++----------- 4 files changed, 87 insertions(+), 46 deletions(-) diff --git a/src/cmd/auth.rs b/src/cmd/auth.rs index 54afd87..011d73d 100644 --- a/src/cmd/auth.rs +++ b/src/cmd/auth.rs @@ -528,12 +528,12 @@ pub(crate) async fn login(input: &str) -> Result<()> { // as everywhere, including in scripts that had not been touched. // Adding a credential and choosing an identity are different // decisions, and only the first one was asked for here. - crate::config::account::remember(&did, who.handle.as_deref())?; + crate::config::account::remember(&did, who.handle.as_deref()).await?; // The one durable fact this login produces that nothing else can // reconstruct: the port in here is closed as soon as this // function returns. Without it the first refresh is refused and // the session is deleted. See `client_metadata`. - crate::config::account::remember_client_id(&did, &client_id)?; + crate::config::account::remember_client_id(&did, &client_id).await?; println!( "logged in as {}", @@ -582,7 +582,7 @@ pub(crate) async fn login(input: &str) -> Result<()> { // is the state selection refuses outright. Writing it now is // what makes "the first account stays active" durable. if standing.pointer.is_none() { - crate::config::account::set_active(&did)?; + crate::config::account::set_active(&did).await?; standing.pointer = Some(did.clone()); // Only worth saying when it resolves something. On a fresh // install with one account there was nothing to resolve. @@ -741,7 +741,7 @@ async fn logout(account: Option, all: bool, json: bool) -> Result<()> { } for (did, handle) in &targets { - crate::config::account::forget(did)?; + crate::config::account::forget(did).await?; if !json { println!( "logged out {}", @@ -868,7 +868,8 @@ async fn status(json: bool) -> Result<()> { // their accounts is wrong. A registry another atgc holds the // lock on, or one this refuses to overwrite because it could // not be read, must not take the answer down with it. - if let Err(e) = crate::config::account::remember(&account.did, Some(&handle)) { + if let Err(e) = crate::config::account::remember(&account.did, Some(&handle)).await + { crate::term::say::warning!( Config, "could not record {handle}'s handle in the registry: {e}" @@ -1066,7 +1067,7 @@ pub(crate) struct LogoutJson { async fn switch(spec: &str, json: bool) -> Result<()> { crate::term::jsonout::init(json); let selection = crate::config::account::lookup(spec).await?; - crate::config::account::set_active(&selection.did)?; + crate::config::account::set_active(&selection.did).await?; if !json { println!("active account is now {}", selection.display()); } diff --git a/src/config/account/registry.rs b/src/config/account/registry.rs index 543a76d..3df5b25 100644 --- a/src/config/account/registry.rs +++ b/src/config/account/registry.rs @@ -169,12 +169,17 @@ fn save(_guard: &crate::config::lock::Guard, registry: &Registry) -> Result<()> /// every later refresh presents the wrong client and is refused. A lost /// update there does not cost a cached handle; it costs the session. /// -/// Takes the lock blocking rather than async because these callers are sync, -/// and because the work inside is two small file operations with no I/O wait -/// of its own. The wait can still be a peer's refresh, which is why it is -/// bounded — see [`crate::config::lock`]. -fn update(change: impl FnOnce(&mut Registry)) -> Result<()> { - let guard = crate::config::lock::take(crate::logging::oauth::Purpose::Registry)?; +/// Async, and the lock is taken with `take_async`, because every caller is +/// reached from an `async fn` — `cmd::auth`'s login, logout, status and +/// switch, and `selection`'s handle repair. The doc comment here used to say +/// the opposite ("these callers are sync"), and the blocking wait it +/// justified parked a tokio worker thread for up to the whole of the lock's +/// timeout. The wait is very often on a peer doing a token refresh, so the +/// thread being parked is one of the few that could run the I/O ending the +/// wait; on the one- or two-core machines agents run on there may be no +/// other. +async fn update(change: impl FnOnce(&mut Registry)) -> Result<()> { + let guard = crate::config::lock::take_async(crate::logging::oauth::Purpose::Registry).await?; let reading = read(); // **A read-modify-write may not start from a file it could not read.** // [`load`] answers an unparseable file with an empty registry, which is @@ -243,7 +248,7 @@ pub fn known() -> Result> { /// Record an account, or update its cached handle. Called after a login and /// after any successful handle resolution. -pub fn remember(did: &str, handle: Option<&str>) -> Result<()> { +pub async fn remember(did: &str, handle: Option<&str>) -> Result<()> { update(|registry| { let entry = registry.accounts.entry(did.to_string()).or_default(); if entry.added_at.is_none() { @@ -262,6 +267,7 @@ pub fn remember(did: &str, handle: Option<&str>) -> Result<()> { ); } }) + .await } /// Record the `client_id` a fresh grant was issued to. Called once, by @@ -270,7 +276,7 @@ pub fn remember(did: &str, handle: Option<&str>) -> Result<()> { /// Overwrites unconditionally: a second login for the same DID supersedes the /// first grant (login prunes the old session), so the newest `client_id` is /// the only one with a session behind it. -pub fn remember_client_id(did: &str, client_id: &str) -> Result<()> { +pub async fn remember_client_id(did: &str, client_id: &str) -> Result<()> { update(|registry| { registry .accounts @@ -278,6 +284,7 @@ pub fn remember_client_id(did: &str, client_id: &str) -> Result<()> { .or_default() .client_id = Some(client_id.to_string()); }) + .await } /// The `client_id` this account's grant was issued to, if it was recorded. @@ -292,18 +299,19 @@ pub fn cached_handle(did: &str) -> Option { load().accounts.get(did).and_then(|a| a.handle.clone()) } -pub fn set_active(did: &str) -> Result<()> { - update(|registry| registry.active = Some(did.to_string())) +pub async fn set_active(did: &str) -> Result<()> { + update(|registry| registry.active = Some(did.to_string())).await } /// Drop an account from the registry (logout). Deliberately separate from /// dropping its session: an expired session keeps the entry, an explicit /// logout removes it. -pub fn forget(did: &str) -> Result<()> { +pub async fn forget(did: &str) -> Result<()> { update(|registry| { registry.accounts.remove(did); if registry.active.as_deref() == Some(did) { registry.active = None; } }) + .await } diff --git a/src/config/account/selection.rs b/src/config/account/selection.rs index a9fceaf..5fd317b 100644 --- a/src/config/account/selection.rs +++ b/src/config/account/selection.rs @@ -737,7 +737,7 @@ async fn resolve_spec(spec: &str, known: &[Known], source: Source) -> Result Result { - let started = Instant::now(); - let file = open()?; - loop { - match file.try_lock() { - Ok(()) => return Ok(acquired(file, purpose, started)), - Err(TryLockError::WouldBlock) if started.elapsed() < WAIT => {} - Err(e) => return Err(gave_up(purpose, started, LockFailure::from(e))), - } - std::thread::sleep(POLL); - } -} - /// Take the lock, yielding to the executor while waiting. /// /// The wait can be as long as another process's token refresh, which is a /// network round trip; blocking the thread for that would stall every other /// task in this process, including the timeouts meant to bound it. /// +/// This is the only way in. There used to be a blocking `take` beside it, +/// for "callers that are not sync" — but every caller it had was reached +/// from an `async fn`, so what it actually did was park a tokio worker on +/// top of the peer whose I/O would end the wait. On the one- or two-core +/// machines agents run on, that is a deadlock in all but name. +/// /// Fails rather than proceeding unlocked; never nest inside another live /// [`Guard`]. Both are the module docs' subject. pub async fn take_async(purpose: Purpose) -> Result { + let path = lock_path_in(&crate::config::dir::config_dir()?); + take_async_at(&path, purpose).await +} + +/// The wait loop, over a lock file named outright. +/// +/// Split from [`take_async`] for the reason [`lock_path_in`] is split from +/// `HOME`: `HOME` cannot be redirected inside a test in this crate (see +/// [`crate::docs::testing`]), so the part that waits takes the path it waits +/// on as an argument and the one line that resolves the directory is the +/// untested wrapper above it. +async fn take_async_at(path: &Path, purpose: Purpose) -> Result { let started = Instant::now(); - let file = open()?; + let file = open_at(path)?; loop { match file.try_lock() { Ok(()) => return Ok(acquired(file, purpose, started)), @@ -214,11 +211,6 @@ impl From for LockFailure { /// 0600 from creation, like everything else in this directory: the file /// holds nothing, but a world-writable lock file is a way for another local /// user to hold atgc still. -fn open() -> Result { - let path = lock_path_in(&crate::config::dir::config_dir()?); - open_at(&path) -} - fn open_at(path: &Path) -> Result { let mut options = std::fs::OpenOptions::new(); // Write access is what `flock` needs from this handle; the file's @@ -494,6 +486,46 @@ mod tests { drop(held); } + /// The invariant this module states and `registry::update` used to + /// violate: waiting for the lock from async code must not park the + /// thread the peer's progress depends on. + /// + /// Driven on a single-threaded runtime, which is the shape of the + /// machines this matters on: if the waiter blocks its thread, the + /// releaser below never gets to run and the waiter can only fail after + /// the whole of `WAIT`. Because `take_async` yields, the release runs + /// and the waiter acquires. The tick count is asserted so that + /// "acquired" cannot pass by having waited for a release that happened + /// on some other thread. + #[tokio::test(flavor = "current_thread")] + async fn waiting_for_the_lock_lets_the_rest_of_this_process_run() { + let dir = temp_dir("async-wait"); + let path = dir.path().join(".lock"); + + let held = open_at(&path).expect("open the lock file"); + held.try_lock().expect("take it"); + + let ticks = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let counter = std::sync::Arc::clone(&ticks); + let releaser = tokio::spawn(async move { + for _ in 0..3 { + counter.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + tokio::time::sleep(POLL).await; + } + held.unlock().unwrap(); + }); + + let guard = take_async_at(&path, Purpose::Registry) + .await + .expect("the lock was released while we waited"); + releaser.await.expect("the releaser ran"); + assert!( + ticks.load(std::sync::atomic::Ordering::SeqCst) > 0, + "the waiter blocked its thread: nothing else on this runtime ran" + ); + drop(guard); + } + #[test] fn lock_file_sits_beside_the_files_it_protects() { let dir = Path::new("/home/someone/.config/atgc"); -- 2.51.2