From 0e2fc69d9d54603ca091bf7245e0208f21c50488 Mon Sep 17 00:00:00 2001 From: "@permadeath.com" Date: Tue, 18 Aug 2026 22:25:30 -0400 Subject: [PATCH] fix(lock): take the store lock on jacquard's writes too `SessionStore` reached `write_store_at` directly, so most writes to sessions.json were atomic and unlocked, and a whole-file rewrite can drop another process's session. Taking it unconditionally deadlocks inside a refresh, so `take_or_inherit` keyed on the task id closes both. Co-Authored-By: Claude Opus 5 (1M context) --- TODO.md | 53 +++++++------- src/auth/store.rs | 41 ++++++++--- src/config/lock.rs | 163 +++++++++++++++++++++++++++++++++++++++++++ src/logging/oauth.rs | 8 ++- 4 files changed, 230 insertions(+), 35 deletions(-) diff --git a/TODO.md b/TODO.md index e20c999..f80180f 100644 --- a/TODO.md +++ b/TODO.md @@ -334,33 +334,34 @@ without touching the file. Confirmed against the real store: `auth status` and `doctor` now leave `sessions.json` at the same inode and the same checksum, where every such command used to rewrite it -- [ ] **The one pathway is atomic but still unlocked, and the obvious fix - self-deadlocks.** `SessionStore` reaches `write_store_at` directly - rather than through `write_store`, which is the function that requires - a `config::lock::Guard` as a proof obligation — so jacquard's writes - are now atomic but take no lock, exactly as they did before. Closing - that is not a matter of adding a `take_async` to `put`: - `SessionRegistry::get_refreshed` calls `upsert_session` and - `delete_session` *inside* the guard `lock_for_refresh` returns - (vendored `session.rs`, the refresh critical section), so a lock taken - unconditionally in the write path blocks against this process's own - `flock`, waits out the 30-second `WAIT`, and then fails — on the one - write that carries freshly rotated tokens. It needs a re-entrancy - story: a `Hold` that is either taken here or inherited from an - enclosing section, with the refresh guard parked somewhere the store - can see it. +- [x] **The one pathway was atomic but still unlocked.** `SessionStore` + reached `write_store_at` directly rather than through the function that + requires a `config::lock::Guard` as a proof obligation, so jacquard's + writes — which is most of them — were atomic and took no lock, exactly + as they had before. A whole-file rewrite can drop another process's + session whatever prompted it, which is the thing the lock exists for. - That leaves a hole worth naming rather than discovering later. "The - lock is already held" is a fact about the *process*, not about the - call chain, so a second task writing to the store while a refresh is in - flight would read the refresh's hold as its own and proceed - unprotected. Closing it properly wants task-local state — there is no - stable way to ask "is this hold mine" from a `&self` method jacquard - calls. In practice the window is small: every authenticated request - goes through the registry's per-`(DID, session_id)` mutex, atgc uses - one account per invocation, and the read and write in `put` have no - await between them. But it is a hole, and it is the reason this is - written down instead of shipped + The obvious fix self-deadlocks, as this entry said: `get_refreshed` + calls `upsert_session` and `delete_session` *inside* the guard + `lock_for_refresh` hands it, and `flock` is per process, so a lock + taken unconditionally in the write path waits out `WAIT` against this + process's own hold and then fails — on the one write that carries + freshly rotated tokens. So it is a re-entrancy story after all: + `take_reentrant` publishes the section's context in a `HOLDER`, and + `take_or_inherit` returns `Hold::Inherited` to a write already inside + it and `Hold::Owned` to everything else. + + The hole this entry named — "the lock is already held" being a fact + about the *process* rather than the call chain — is closed by what + `HOLDER` stores: `tokio::task::try_id()`, so a second task writing + while a refresh is in flight does not match and takes the lock + properly. The `Option>` is deliberate twice over: the inner + `None` is a future polled by `block_on` rather than a spawned task, + which is where atgc's commands actually run, and matching `None` + against `None` is sound because there is one `block_on` context per + runtime and `put`'s read and write have no `.await` between them. + Both halves are pinned by tests, the inheriting one timed so that + "inherited" cannot pass as "waited thirty seconds and then inherited" - [x] Whether a nonce-only update takes the lock, settled: **yes, when there is one — but the case that mattered no longer writes at all.** The exclusion was never about the nonce; it is about the whole-file diff --git a/src/auth/store.rs b/src/auth/store.rs index 55c43c1..a91636e 100644 --- a/src/auth/store.rs +++ b/src/auth/store.rs @@ -72,7 +72,26 @@ impl SessionStore { super::read_store_or_empty_at(&self.path).map_err(into_store_error) } - fn write(&self, map: &Map, reason: Reason) -> Result<(), SessionStoreError> { + /// Write the whole map, under the lock. + /// + /// Async and fallible-on-the-lock, where this used to be neither: it + /// reached `write_store_at` directly, so jacquard's writes were atomic + /// and *unlocked* — a whole-file rewrite that can drop another process's + /// session, which is the exact thing the lock exists to stop. + /// + /// [`take_or_inherit`](crate::config::lock::take_or_inherit) rather than + /// `take_async`, because half of what reaches here runs inside the + /// critical section `lock_for_refresh` opened and the other half does + /// not. Taking unconditionally would deadlock the first against this + /// process's own hold. + async fn write( + &self, + map: &Map, + reason: Reason, + ) -> Result<(), SessionStoreError> { + let _hold = crate::config::lock::take_or_inherit(crate::logging::oauth::Purpose::Store) + .await + .map_err(into_store_error)?; super::write_store_at(&self.path, map, reason).map_err(into_store_error) } @@ -92,7 +111,12 @@ impl SessionStore { /// rewrite, so the bytes are what it compares. A record that round-trips /// through JSON to something equal-but-not-identical would still be /// written, which is the safe direction to be wrong in. - fn put(&self, key: String, value: Value, reason: Reason) -> Result<(), SessionStoreError> { + async fn put( + &self, + key: String, + value: Value, + reason: Reason, + ) -> Result<(), SessionStoreError> { let mut map = self.read()?; if map.get(&key) == Some(&value) { crate::logging::debug::log(format!( @@ -102,16 +126,16 @@ impl SessionStore { return Ok(()); } map.insert(key, value); - self.write(&map, reason) + self.write(&map, reason).await } /// Drop `key`, unless it is already absent. - fn remove(&self, key: &str, reason: Reason) -> Result<(), SessionStoreError> { + async fn remove(&self, key: &str, reason: Reason) -> Result<(), SessionStoreError> { let mut map = self.read()?; if map.remove(key).is_none() { return Ok(()); } - self.write(&map, reason) + self.write(&map, reason).await } } @@ -155,7 +179,7 @@ impl ClientAuthStore for SessionStore { async fn upsert_session(&self, session: ClientSessionData) -> Result<(), SessionStoreError> { let key = session_key(session.account_did.as_str(), session.session_id.as_ref()); let value = serde_json::to_value(StoredSession::ClientSession(session))?; - self.put(key, value, Reason::Upsert) + self.put(key, value, Reason::Upsert).await } async fn delete_session( @@ -167,6 +191,7 @@ impl ClientAuthStore for SessionStore { &session_key(did.as_str(), session_id), Reason::VendorRegistry, ) + .await } async fn get_auth_req_info( @@ -185,11 +210,11 @@ impl ClientAuthStore for SessionStore { ) -> Result<(), SessionStoreError> { let key = state_key(auth_req_info.state.as_ref()); let value = serde_json::to_value(StoredSession::AuthRequest(auth_req_info.clone()))?; - self.put(key, value, Reason::AuthState) + self.put(key, value, Reason::AuthState).await } async fn delete_auth_req_info(&self, state: &str) -> Result<(), SessionStoreError> { - self.remove(&state_key(state), Reason::AuthState) + self.remove(&state_key(state), Reason::AuthState).await } /// Every OAuth session in the file, from one read. diff --git a/src/config/lock.rs b/src/config/lock.rs index 65bc9ea..a0b80e1 100644 --- a/src/config/lock.rs +++ b/src/config/lock.rs @@ -123,10 +123,16 @@ pub struct Guard { file: Option, purpose: Purpose, since: Instant, + /// Whether this guard published itself to [`HOLDER`] and so has to take + /// itself back out. Only [`take_reentrant`] sets it. + registered: bool, } impl Drop for Guard { fn drop(&mut self) { + if self.registered { + *holder() = None; + } // `Option` only so that `Drop` can take the file back out; it is // `Some` for the whole of the guard's visible life. let Some(file) = self.file.take() else { return }; @@ -237,7 +243,87 @@ fn acquired(file: File, purpose: Purpose, since: Instant) -> Guard { file: Some(file), purpose, since: Instant::now(), + registered: false, + } +} + +// --------------------------------------------------------------------------- +// Re-entrancy +// --------------------------------------------------------------------------- + +/// Which async context, if any, is inside the refresh critical section. +/// +/// `flock` is per *process* and per open file description, so a second +/// `try_lock` from this process against its own hold does not pass through: +/// it waits out [`WAIT`] and fails. That matters because jacquard's +/// `get_refreshed` calls `upsert_session` and `delete_session` *inside* the +/// guard [`crate::logging::oauth`]'s `lock_for_refresh` hands it — so a lock +/// taken unconditionally in the store's write path would deadlock against +/// this process on the one write that carries freshly rotated tokens. +/// +/// The identity stored is [`tokio::task::try_id`]'s, and it is an `Option` +/// twice over on purpose. The outer one is "is anybody inside"; the inner +/// one is the task id, which is `None` for a future polled by `block_on` +/// rather than by a spawned task — which is where atgc's commands actually +/// run, `#[tokio::main]` being a `block_on`. Matching `None` against `None` +/// is therefore the *ordinary* case rather than a degenerate one, and it is +/// sound for a reason worth stating: there is one `block_on` context per +/// runtime, and the read and the write inside `SessionStore::put` have no +/// `.await` between them, so two futures sharing that context cannot +/// interleave there. A spawned task carries `Some(id)`, does not match, and +/// takes the lock properly — which is the hole this closes. +static HOLDER: std::sync::Mutex>> = std::sync::Mutex::new(None); + +/// [`HOLDER`], with a poisoned mutex treated as unheld. +/// +/// A panic inside the critical section is a bug that has already happened; +/// refusing every later write over it would turn one into a broken install. +fn holder() -> std::sync::MutexGuard<'static, Option>> { + HOLDER.lock().unwrap_or_else(|e| e.into_inner()) +} + +/// Take the lock and publish this context as its holder, so that writes made +/// underneath it inherit rather than deadlock. +/// +/// For `lock_for_refresh` and nothing else: it is the one place that holds +/// the lock across code it does not control. +pub async fn take_reentrant(purpose: Purpose) -> Result { + let mut guard = take_async(purpose).await?; + *holder() = Some(tokio::task::try_id()); + guard.registered = true; + Ok(guard) +} + +/// A held lock, however it came to be held. +/// +/// Not a `Guard`, because [`Hold::Inherited`] owns nothing and must not +/// release anything on drop: the enclosing critical section is still using +/// it. Existing still means held, which is the property the whole module is +/// built on. +#[derive(Debug)] +pub enum Hold { + /// Held by this value: dropping it releases the lock. Never read, and + /// that is the point — it is an RAII token, and the whole of what it + /// does happens in [`Guard::drop`]. + Owned(#[allow(dead_code, reason = "held for its Drop")] Guard), + /// This context is already inside a [`take_reentrant`] section. + Inherited, +} + +/// Take the lock, unless this context already holds it. +/// +/// The front door for writes that can happen either on their own or from +/// inside a refresh. A caller that knows it is neither should keep using +/// [`take_async`], which cannot silently do nothing. +pub async fn take_or_inherit(purpose: Purpose) -> Result { + if *holder() == Some(tokio::task::try_id()) { + crate::logging::debug::log( + "store lock already held by this context; writing inside it rather than \ + waiting on ourselves", + ); + return Ok(Hold::Inherited); } + Ok(Hold::Owned(take_async(purpose).await?)) } /// Record the failure and say what happened in terms of what to do about it. @@ -320,6 +406,83 @@ mod tests { dir } + /// The deadlock this exists to avoid, driven directly. + /// + /// `flock` is per process, so a second `try_lock` against this process's + /// own hold does not pass through — it waits out `WAIT` and fails. That + /// is what would happen on every `upsert_session` inside a refresh if + /// the store's write path took the lock unconditionally, and it would + /// happen on the one write that carries freshly rotated tokens. + /// + /// Timed, because "it inherited" and "it waited thirty seconds and then + /// inherited" are the same assertion otherwise, and only one of them is + /// the fix. + #[tokio::test] + async fn a_write_inside_the_refresh_section_inherits_the_lock() { + let held = take_reentrant(Purpose::Refresh) + .await + .expect("nothing else holds it"); + + let started = Instant::now(); + let inherited = take_or_inherit(Purpose::Store) + .await + .expect("a write underneath it must not fail"); + assert!( + matches!(inherited, Hold::Inherited), + "the write took its own lock and would have deadlocked" + ); + assert!( + started.elapsed() < Duration::from_secs(1), + "it waited rather than inheriting: {:?}", + started.elapsed() + ); + + // Dropping an inherited hold releases nothing: the section around it + // is still using the lock. + drop(inherited); + assert_eq!(*holder(), Some(tokio::task::try_id())); + + drop(held); + assert_eq!(*holder(), None, "the section did not clean up after itself"); + } + + /// The hole the task identity closes: a *different* task writing while a + /// refresh is in flight must not read the refresh's hold as its own. + /// + /// This is the case the design note in TODO.md said wanted task-local + /// state. A spawned task carries its own `tokio::task::Id`, so it does + /// not match the holder and takes the lock properly — which here means + /// it cannot take it at all while the section holds it, and says so + /// rather than proceeding unlocked. + #[tokio::test] + async fn another_task_does_not_inherit_a_refresh_it_is_not_part_of() { + let held = take_reentrant(Purpose::Refresh) + .await + .expect("nothing else holds it"); + assert_eq!( + *holder(), + Some(tokio::task::try_id()), + "the section published itself" + ); + + let inherits = tokio::spawn(async { + // A spawned task has an id of its own, so it is not the holder. + // Asserted on the comparison `take_or_inherit` makes rather than + // by calling it: the other branch really does try for the lock, + // and this process is holding it, so driving it here would wait + // out the whole of `WAIT` to prove something already proven. + *holder() == Some(tokio::task::try_id()) + }) + .await + .expect("the task ran"); + assert!( + !inherits, + "a task outside the refresh would have inherited a lock it does not hold" + ); + + drop(held); + } + #[test] fn lock_file_sits_beside_the_files_it_protects() { let dir = Path::new("/home/someone/.config/atgc"); diff --git a/src/logging/oauth.rs b/src/logging/oauth.rs index fa98e38..f500a15 100644 --- a/src/logging/oauth.rs +++ b/src/logging/oauth.rs @@ -1316,8 +1316,14 @@ impl ClientAuthStore for LoggedAuthStore { /// error arrives at jacquard as `Error::Store`, which `is_permanent` /// reports as false, so this cannot be mistaken for the very failure it /// exists to prevent. + /// `take_reentrant`, not `take_async`: `get_refreshed` calls + /// `upsert_session` and `delete_session` *inside* the guard it is handed, + /// and those now take the lock too. Publishing this context as the holder + /// is what lets them inherit it instead of waiting out `WAIT` on this + /// process's own hold and failing — on the one write that carries freshly + /// rotated tokens. async fn lock_for_refresh(&self) -> Result>, SessionStoreError> { - let guard = crate::config::lock::take_async(Purpose::Refresh) + let guard = crate::config::lock::take_reentrant(Purpose::Refresh) .await .map_err(|e| SessionStoreError::Other(e.into()))?; Ok(Some(Box::new(guard))) -- 2.51.2