diff --git a/plan/credential-store.md b/plan/credential-store.md index 944c9d4..5262b9e 100644 --- a/plan/credential-store.md +++ b/plan/credential-store.md @@ -77,6 +77,28 @@ re-proposed: the OS keyring, and encrypting the DPoP key at rest. ## Done +- [x] The session store's own read-modify-write was not under the lock — + only its write was. `SessionStore::put` and `remove` read the whole map, + compared, and *then* called a `write` that took the lock, which is one + step too late: two processes both read map `M`, both take the lock in + turn, and the second writes `M` plus its own change over the first. Both + writes are whole and atomic, so rename repairs nothing and the log shows + two ordinary `StoreWrite`s. This is the same defect this epic already + records as closed for atgc's own paths, still open on jacquard's, which + is about ninety writes in every ninety-one. + What it costs is not a cached field. A refresh token is single-use, so + restoring the copy another process has already spent leaves + `invalid_grant` on next use — classified permanent, which deletes the + session. The account is logged out by a write that reported success. + The read is inside the hold now, in one `edit`, and `write` takes a + `&Hold` so the pairing is a compile-time obligation the way + `sessions::write_store`'s `&Guard` already was. The unlocked pre-check + stays: nearly every call here writes a session back exactly as it was + read, and putting those through the lock would queue every atgc on the + machine to change nothing. Losing that race costs one lock acquisition + and no write. Found while asking what breaks when many agents share one + `$HOME`, which is a deployment this was not designed against + - [x] A read-modify-write no longer starts from a file it could not read. `accounts.json` is read whole, edited and written back under the lock, and `load` answers an unparseable file with an *empty* registry — the diff --git a/src/clients/atproto/oauth/store.rs b/src/clients/atproto/oauth/store.rs index 4bffb16..db1e6ae 100644 --- a/src/clients/atproto/oauth/store.rs +++ b/src/clients/atproto/oauth/store.rs @@ -72,26 +72,33 @@ impl SessionStore { super::sessions::read_store_or_empty_at(&self.path).map_err(into_store_error) } - /// Write the whole map, under the lock. + /// Write the whole map. The hold is a proof obligation, not data. /// - /// 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. + /// `_hold` is here for the reason `sessions::write_store` takes a + /// `&Guard`: this file is read whole, edited and written back, so a write + /// is only correct if the read it was computed from happened under the + /// same hold. Requiring the token makes that pairing a compile-time + /// matter, and [`SessionStore::edit`] is the one thing that can satisfy + /// it. /// + /// The lock used to live *here*, which is one step too late: two + /// processes could both read map `M`, then take the lock in turn, and + /// the second would write `M` plus its own change over the first — a + /// lost update that atomic rename cannot repair, because both writes + /// landed whole. + /// + /// The lock [`SessionStore::edit`] takes is /// [`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( + fn write( &self, + _hold: &crate::config::lock::Hold, 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::sessions::write_store_at(&self.path, map, reason).map_err(into_store_error) } @@ -117,25 +124,71 @@ impl SessionStore { value: Value, reason: Reason, ) -> Result<(), SessionStoreError> { - let mut map = self.read()?; - if map.get(&key) == Some(&value) { + // Unlocked, and only to decide whether there is anything to do. Most + // calls here are a session being written back exactly as it was + // read, and taking the lock for those would put every atgc on the + // machine through one queue to change nothing. + if self.read()?.get(&key) == Some(&value) { crate::logging::debug::log(format!( "session store already holds {}; not rewriting it", crate::logging::oauth::redact_key(&key) )); return Ok(()); } - map.insert(key, value); - self.write(&map, reason).await + self.edit(reason, |map| { + if map.get(&key) == Some(&value) { + return false; + } + map.insert(key.clone(), value.clone()); + true + }) + .await } /// Drop `key`, unless it is already absent. async fn remove(&self, key: &str, reason: Reason) -> Result<(), SessionStoreError> { + // The same unlocked pre-check as `put`, for the same reason. + if !self.read()?.contains_key(key) { + return Ok(()); + } + self.edit(reason, |map| map.remove(key).is_some()).await + } + + /// Read the store, change it, and write it back — all under one hold. + /// + /// **The read has to be inside the lock, not merely before the write.** + /// This file is read whole, edited and written back, so a write is only + /// correct if it was computed from a read nothing has landed on top of. + /// The lock used to sit inside `write`, which is a write-modify order + /// that atomic rename does not repair: two processes both read map `M`, + /// both take the lock in turn, and the second writes `M` plus its own + /// change over the first's — losing it. + /// + /// That is not a cache here. `sessions.json` holds refresh tokens, and a + /// refresh token is single-use: restoring the copy another process has + /// already spent leaves an `invalid_grant` on next use, which is + /// classified permanent and deletes the session. The account is logged + /// out, by a write that reported success. + /// + /// `change` returns whether it changed anything, so a race lost between + /// the pre-check and the lock costs a lock acquisition and no write. + async fn edit( + &self, + reason: Reason, + change: impl FnOnce(&mut Map) -> bool, + ) -> Result<(), SessionStoreError> { + let hold = crate::config::lock::take_or_inherit(crate::logging::oauth::Purpose::Store) + .await + .map_err(into_store_error)?; let mut map = self.read()?; - if map.remove(key).is_none() { + if !change(&mut map) { + crate::logging::debug::log( + "session store already had this change when the lock came free; \ + nothing written", + ); return Ok(()); } - self.write(&map, reason).await + self.write(&hold, &map, reason) } } @@ -370,4 +423,111 @@ mod tests { std::fs::remove_dir_all(&dir).ok(); } + + /// Two writers landing at once both survive. + /// + /// The shape the lock exists for, and the one it did not cover until the + /// read moved inside it: `put` read the whole map, *then* took the lock + /// to write it back. Two processes reading the same map and taking the + /// lock in turn both write whole, valid files, and the second one's file + /// is the first one's map with only the second one's change in it — so + /// the first write is lost with nothing to show it happened. Atomic + /// rename does not help; both writes were atomic. + /// + /// What makes that worth a test rather than a comment is what the file + /// holds. `sessions.json` carries refresh tokens, and a refresh token is + /// single-use: putting back the copy another process has already spent + /// leaves an `invalid_grant` on next use, which is classified permanent + /// and deletes the session. The account is logged out by a write that + /// reported success. + /// + /// The two writers are tasks rather than processes, which is the honest + /// limit of this test: `flock` is what separates real processes, and + /// what is exercised here is that the read and the write are one + /// critical section. The slow `change` is what forces the overlap — + /// without it the two finish in sequence and nothing is proved. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn a_concurrent_write_is_not_lost() { + let dir = store_dir("concurrent-put"); + let path = dir.path().join("sessions.json"); + let store = SessionStore::new(path.clone()); + + let alice = serde_json::to_value(StoredSession::ClientSession( + crate::testutil::oauth_session("alice-token"), + )) + .unwrap(); + let bob = serde_json::to_value(StoredSession::ClientSession( + crate::testutil::oauth_session("bob-token"), + )) + .unwrap(); + + let slow = store.edit(Reason::Upsert, |map| { + map.insert("oauth:did:plc:alice/s".to_string(), alice.clone()); + // Inside the critical section, so the other writer is still + // waiting for the lock when this decides what to write. + std::thread::sleep(std::time::Duration::from_millis(150)); + true + }); + let other = async { + tokio::time::sleep(std::time::Duration::from_millis(20)).await; + store + .put( + "oauth:did:plc:bob/s".to_string(), + bob.clone(), + Reason::Upsert, + ) + .await + }; + let (first, second) = tokio::join!(slow, other); + first.expect("the slow write"); + second.expect("the write that had to wait"); + + let map = store.read().expect("read the store back"); + assert!( + map.contains_key("oauth:did:plc:alice/s"), + "the first writer's session was overwritten: {:?}", + map.keys().collect::>() + ); + assert!( + map.contains_key("oauth:did:plc:bob/s"), + "the second writer's session is missing: {:?}", + map.keys().collect::>() + ); + } + + /// A change another writer already made costs a lock and no write. + /// + /// `put`'s unlocked pre-check is what keeps the common case — a session + /// written back exactly as it was read, which is nearly all of them — + /// from queueing every atgc on the machine behind one lock. Losing that + /// race is not an error: the second look under the lock finds the value + /// already there and writes nothing. + #[tokio::test] + async fn a_put_that_is_already_true_writes_nothing() { + let dir = store_dir("put-noop"); + let path = dir.path().join("sessions.json"); + let store = SessionStore::new(path.clone()); + + let value = serde_json::to_value(StoredSession::ClientSession( + crate::testutil::oauth_session("token"), + )) + .unwrap(); + let key = "oauth:did:plc:alice/s".to_string(); + store + .put(key.clone(), value.clone(), Reason::Upsert) + .await + .expect("the first put"); + let before = std::fs::metadata(&path).expect("stat").len(); + + store + .put(key, value, Reason::Upsert) + .await + .expect("the second put"); + + assert_eq!( + std::fs::metadata(&path).expect("stat").len(), + before, + "the store was rewritten for a change that was already there" + ); + } }