//! One advisory lock over `~/.config/atgc`, so several atgc processes can //! share the files in it. //! //! # Why this exists //! //! [`crate::config::dir::write_atomic`] made each individual write to an //! identity's session and account files land whole. That closes the *crash* //! window and nothing else: replacing a file atomically says nothing about //! whether the bytes being written were computed from the file's latest //! contents. Both files are read, edited in memory and written back whole, //! from several code paths and from jacquard's session store underneath //! them, so two invocations that overlap still produce a lost update — B //! reads, A reads, A writes, B writes B's now-stale copy, and A's change is //! gone. `read_store`'s doc comment has named this as the open half of the //! problem since the incident that produced [`crate::logging::oauth`]. //! //! Agents make it routine rather than rare. Several atgc processes running //! at once against one `$HOME` is the normal case now, not a coincidence of //! two terminals. //! //! The worst instance is not a lost edit but a destroyed grant, and it is //! not fixable by locking the file *writes*. ATProto refresh tokens are //! single-use and rotating. Two processes that both find an expiring session //! both spend the same refresh token; the authorization server honors one //! and answers the other `invalid_grant`; jacquard's `SessionRegistry` //! classifies that as permanent and deletes the session outright. The //! critical section that has to be serialized is therefore not "the write" //! but *read → token request → write*, spanning a network round trip. That //! is what [`Purpose::Refresh`] holds this lock across, via //! `ClientAuthStore::lock_for_refresh` in the vendored crate. //! //! # Shape //! //! One lock file, `~/.config/atgc/.lock`, for the whole directory rather //! than one per JSON file. The files are small and the critical sections are //! rare, so there is nothing to gain from finer granularity, and a single //! lock cannot be acquired in two different orders — which is the only way //! two locks could deadlock against each other. //! //! The lock is advisory (`flock(2)`, through `std::fs::File::try_lock`) and //! released by the kernel when the process dies, so a killed atgc cannot //! leave a stale lock behind for the next one to wait out. //! //! # Nesting is not allowed //! //! `flock` is held per open file description, so a *second* [`take_async`] in the //! same process blocks against the first one rather than passing through it. //! No caller may hold a guard while calling code that takes another: the //! session-store guard and the account-registry guard are always taken in //! sequence, never nested (see `cmd::auth::logout`, which scopes its store guard //! to a block precisely so that the `account::forget` calls after it start //! from nothing held). A slip here waits out [`WAIT`] against itself and then //! fails the command, which is loud, bounded, and cannot corrupt anything. //! //! # Both ends are bounded, and the bounds are related //! //! A caller waits up to [`WAIT`] and then **gives up and returns an error**. //! It does not go ahead unlocked. //! //! For that to be a fair answer, a waiter has to outlast the longest hold a //! holder is allowed to take — otherwise "another atgc got there first" is //! reported for work that was proceeding normally. Nothing enforced that. //! [`HOLD_BUDGET`] now caps the refresh critical section, [`WAIT`] is derived //! from it, and a build fails if the two ever cross. //! //! The alternative was tried and rejected. Proceeding without the lock is //! only "the behavior every version so far had" in the sense that the bug //! this module exists to fix was also that behavior: an unlocked //! read-modify-write is how a live session gets deleted, and a session //! deleted is a login the user has to perform again, possibly without //! understanding why. Waiting a few seconds costs an agent nothing it will //! notice. Being silently logged out costs a person a browser round trip and //! a chunk of trust, and the failure is invisible at the moment it is caused. //! Between a command that fails saying exactly what is wrong and one that //! succeeds while destroying a grant, the loud one wins every time. //! //! Nor does the *waiter* retry on its own, which is a separate question and //! was reconsidered for the agent-fleet case: many short-lived atgc processes //! against one `$HOME`. It is left as it is deliberately. A retry loop here //! would be indistinguishable from a longer [`WAIT`] — the wait already *is* //! a retry loop — while hiding how long a command really took and stacking //! more pollers onto a lock that has no queue. The bound belongs in one //! place, and it is [`WAIT`]. A caller that wants to sit longer than that //! should say so by retrying the whole command, which is exactly what //! [`crate::exit::Exit::Conflict`] tells it to do. //! //! One consequence to know about: this makes the lock a hard dependency, so //! a filesystem where `flock` is unavailable rather than merely contended //! (`ENOLCK` on some remote mounts) fails these commands outright instead of //! degrading. That is the intended trade — see plan/credential-store.md, //! which records it as the one way this can bite. use anyhow::{Context, Result}; use std::fs::{File, TryLockError}; use std::path::{Path, PathBuf}; use std::time::{Duration, Instant}; use crate::logging::oauth::Purpose; /// The most time the refresh critical section may hold the lock. /// /// Enforced, by [`under_hold_budget`], and that is the point: every number /// below is derived from this one, so the derivation is only worth anything /// if the holder actually obeys it. Nothing else in the process enforces it — /// [`crate::clients::http`]'s `READ_TIMEOUT` bounds the gap *between* bytes /// rather than the exchange, and `MAX_RETRY_AFTER` allows a further backoff /// on top, so the sum of the client's own deadlines is well past anything a /// waiter is prepared to sit through. /// /// Twenty-five seconds is generous for what is inside it: a cached /// authorization-server metadata lookup (see /// [`crate::clients::atproto::oauth::metadata`]) and one token POST, which /// together are well under a second against a healthy server. It is sized /// for a server having a bad day, not for the ordinary case. const HOLD_BUDGET: Duration = Duration::from_secs(25); /// How many held sections a waiter is willing to queue behind. /// /// One would make [`WAIT`] the bare invariant and nothing more: the process /// that arrives second gets in, the third does not. Two is the smallest /// number that lets a small burst — which is what an agent fleet produces — /// drain rather than fail, without making the ceiling absurd. const QUEUE_DEPTH: u32 = 2; /// Head-room over the queued holds: scheduling, the poll interval, and the /// fact that `try_lock` has no queue, so a waiter can lose a race it was /// first in line for. const SLACK: Duration = Duration::from_secs(10); /// How long a caller waits for the lock before giving up and failing. /// /// **The invariant: this must exceed [`HOLD_BUDGET`].** It did not used to. /// The comment here claimed that thirty seconds was "several" token /// refreshes and so that a caller reaching it was not queued behind honest /// work — which was false in both halves. The critical section covered two /// uncached network round trips, not one, and the client bounds around them /// (5s connect, a 30s-per-read stall allowance, up to 30s of `Retry-After` /// backoff) sum past thirty seconds on their own. One slow authorization /// server therefore failed every other atgc on the machine with /// [`crate::exit::Exit::Conflict`] while it was queued behind entirely /// honest work. /// /// So the two bounds are now derived from each other rather than asserted to /// be compatible, and the holder's is enforced. const WAIT: Duration = wait_bound(HOLD_BUDGET, QUEUE_DEPTH, SLACK); /// How long a waiter is prepared to sit behind `queue_depth` full holds. /// /// A function, and a `const` one, so the relationship is checkable at compile /// time by the assertion below and testable without waiting out a single /// second of it. const fn wait_bound(hold: Duration, queue_depth: u32, slack: Duration) -> Duration { Duration::from_secs(hold.as_secs() * queue_depth as u64 + slack.as_secs()) } /// The invariant, as a build failure rather than a comment. /// /// A waiter that gives up sooner than the holder is allowed to take is the /// defect this module was corrected for: it turns one slow authorization /// server into a failure for every other process on the machine. const _: () = assert!( WAIT.as_secs() > HOLD_BUDGET.as_secs(), "a waiter must outlast the longest hold the holder is allowed to take" ); /// How long to sleep between attempts while waiting, before jitter. /// /// `try_lock` in a poll loop rather than the blocking `lock`, because the /// async callers must not park the executor thread while a peer does network /// I/O, and because a bounded wait needs a wakeup to bound it. 25ms is far /// below the cost of anything this serializes and far above the cost of the /// syscall. const POLL_MS: u64 = 25; /// How far either side of [`POLL_MS`] a sleep may land, in percent. /// /// `try_lock` in a loop is not a queue: there is no fairness and no ordering, /// only whoever happens to call at the moment the lock is free. Processes /// that start polling together stay phase-aligned and keep colliding, and the /// same one can keep losing — a starvation tail with nothing bounding it but /// [`WAIT`]. Spreading each sleep breaks the alignment, which is the cheapest /// approximation of fairness available without a real queue. /// /// Forty percent, so a 25ms interval lands anywhere in 15–35ms. Wide enough /// that two processes drift apart within a few polls, narrow enough that the /// interval is still recognizably 25ms. const POLL_JITTER_PERCENT: u64 = 40; /// Where the lock lives, given a config directory. /// /// Split from `HOME` resolution for the same reason /// [`crate::logging::file::Log::path_in`] is: `HOME` cannot be redirected inside a /// test in this crate, so the part that does not read the environment takes /// its input as an argument. pub fn lock_path_in(config_dir: &Path, scope: &Scope) -> Result { Ok(match scope { Scope::Global => config_dir.join(".lock"), Scope::Identity(did) => { crate::config::dir::identity_path_in(config_dir, did)?.join(".lock") } }) } /// What a lock covers. /// /// The lock used to be one file for the whole configuration directory, /// because the thing it protected was one file: `sessions.json` held every /// account, so any write to it excluded every other. Sessions are per /// identity now, and the lock's scope follows the data rather than the /// directory. /// /// This is not only a speed question. The refresh critical section is held /// **across the token request** — see `lock_for_refresh` — so a global lock /// serialises unrelated accounts behind each other's network round trips. On /// a machine running an identity per agent that is the whole fleet through /// one queue, for an exclusion none of them needed from each other: what must /// not happen twice is two spends of *one* session's single-use refresh /// token, and that is one identity's business. /// /// [`Scope::Global`] remains for what is still global — `config.toml`, and /// the pending authorization requests that have no identity yet. #[derive(Debug, Clone, PartialEq, Eq, Hash)] pub enum Scope { /// Files shared by every account: the settings, and pending requests. Global, /// One identity's directory, named by its DID. Identity(String), } impl Scope { /// The scope a session-store key belongs to: its identity, or global for /// an authorization request that does not have one yet. pub fn of_key(key: &str) -> Self { match key .strip_prefix("oauth:") .and_then(|rest| rest.rsplit_once('/')) { Some((did, _)) => Scope::Identity(did.to_string()), None => Scope::Global, } } } /// A held lock. Released on drop. /// /// Callers keep this alive across the whole read-modify-write it protects and /// otherwise ignore it; the type exists to make "I am holding the lock" a /// value that can be passed to the functions that require it, so a write that /// forgets to take it does not compile. See /// `crate::clients::atproto::oauth::sessions::write_store`. /// /// Existing means held: the only constructor is [`acquired`], reached only /// after `try_lock` succeeded. There is no "guard that missed the lock", /// because there is no caller allowed to proceed on one. #[derive(Debug)] pub struct Guard { file: Option, purpose: Purpose, /// What this guard covers, so `Drop` can unpublish the right entry. scope: Scope, since: Instant, /// Whether this guard published itself to [`HOLDER`] and so has to take /// itself back out. Only [`take_reentrant`] sets it. registered: bool, /// The flow this was taken on, so the debug-only nesting check releases /// against the same key it registered against. #[cfg(debug_assertions)] flow: Flow, } impl Drop for Guard { fn drop(&mut self) { #[cfg(debug_assertions)] forget_held(&self.flow, &self.scope); if self.registered && let Some(map) = holders().as_mut() { map.remove(&self.scope); } // `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 }; // Closing the file would release the lock on its own. Unlocking // explicitly first keeps the release ordered before the line that // reports it, so the log cannot show a waiter acquiring before the // holder let go. let _ = file.unlock(); crate::logging::oauth::emit(crate::logging::oauth::Event::LockReleased { purpose: self.purpose, held_ms: ms_since(self.since), }); } } /// 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, scope: Scope) -> Result { take_async_in(&crate::config::dir::config_dir()?, purpose, scope).await } /// [`take_async`], against a configuration directory named outright. /// /// The lock has to live beside the thing it protects. Deriving it from /// `config_dir()` instead means a store rooted anywhere else -- which is /// every store a test builds -- locks a path in the real configuration /// directory while writing somewhere else entirely: the lock guards nothing, /// and the run leaves files in the home of whoever ran it. pub async fn take_async_in(config_dir: &Path, purpose: Purpose, scope: Scope) -> Result { let path = lock_path_in(config_dir, &scope)?; // An identity's lock lives in its directory, which a first write may be // about to create. The global one's directory always exists. if let Some(parent) = path.parent() { crate::config::dir::create_private_dir(parent) .with_context(|| format!("failed to create {}", parent.display()))?; } take_async_at(&path, purpose, scope).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, scope: Scope) -> Result { let started = Instant::now(); let file = open_at(path)?; loop { match file.try_lock() { Ok(()) => return Ok(acquired(file, purpose, scope, started)), Err(TryLockError::WouldBlock) if started.elapsed() < WAIT => {} Err(e) => return Err(gave_up(purpose, started, LockFailure::from(e))), } tokio::time::sleep(poll_delay()).await; } } /// Bound a refresh so it cannot hold the lock past what a waiter will sit /// through. /// /// Wraps the whole of the critical section — jacquard's /// `SessionRegistry::get_refreshed`, which takes the lock through /// `lock_for_refresh` and does not give it back until it returns. Cancelling /// the future drops the guard with it, which is what makes this a bound on /// the *lock* and not merely on the command. /// /// # Cancelling a token request is not free /// /// The window this can land in badly is between the authorization server /// answering with rotated tokens and jacquard writing them: cancel there and /// the old refresh token has been spent for a new one nobody kept. That /// window is the length of a `serde` parse and a file rename, and /// [`HOLD_BUDGET`] is sized so that reaching it at all means the server has /// already stopped answering. Against that: a holder with no bound fails /// *every* other atgc on the machine, deterministically, for as long as one /// authorization server is slow. A rare cancellation of a request that was /// not going to complete is the cheaper of the two. /// /// A timeout is [`crate::exit::Exit::Unreachable`], not `Conflict`: nothing /// is contended and the answer is to try again unchanged. pub async fn under_hold_budget(what: &str, work: impl Future) -> Result { within(HOLD_BUDGET, what, work).await } /// [`under_hold_budget`] with the budget as an argument. /// /// Split out only so a test can drive the overrun in a millisecond instead of /// twenty-five seconds. Production has exactly one budget, which is why the /// public function does not take one. async fn within(budget: Duration, what: &str, work: impl Future) -> Result { match tokio::time::timeout(budget, work).await { Ok(done) => Ok(done), Err(_) => { crate::logging::debug::log(format!( "gave up on {what} after {budget:?} so other atgc processes are not \ queued behind it", )); Err(crate::exit::fail( crate::exit::Exit::Unreachable, format!( "gave up after {}s trying to {what}\n\ the authorization server did not finish answering, and atgc holds a \ lock every other atgc process on this machine waits on\n\ retry; nothing has been changed", budget.as_secs(), ), )) } } } /// How long to sleep before the next `try_lock`, spread so that processes /// polling together do not stay in step. fn poll_delay() -> Duration { jittered_ms(POLL_MS, POLL_JITTER_PERCENT, entropy()) } /// `base_ms`, moved up or down by up to `spread_percent`, picked by /// `entropy`. /// /// Pure, and separate from [`poll_delay`], because the arithmetic is the part /// worth pinning: a test can assert what this decides for a given draw, where /// asserting on real sleeps would be a test that measures the scheduler. fn jittered_ms(base_ms: u64, spread_percent: u64, entropy: u64) -> Duration { // `100 - spread` .. `100 + spread` inclusive, so a zero spread is the // identity rather than a special case. let span = 2 * spread_percent + 1; let percent = (100 - spread_percent.min(100)) + entropy % span; Duration::from_millis(base_ms * percent / 100) } /// A fresh draw from a per-process xorshift. /// /// Seeded from [`std::hash::RandomState`], which is randomized per process: /// two atgc invocations that start together must not draw the same sequence, /// or the jitter above would move them in lockstep instead of apart. No /// dependency for it, because nothing here is deciding anything a stranger /// gets to predict — it is a tie-breaker between local processes. fn entropy() -> u64 { use std::sync::atomic::{AtomicU64, Ordering}; static STATE: AtomicU64 = AtomicU64::new(0); let mut x = STATE.load(Ordering::Relaxed); if x == 0 { use std::hash::{BuildHasher, RandomState}; // `| 1`: zero is xorshift's fixed point and this module's "unseeded" // marker, and it must be neither. x = RandomState::new().hash_one(std::process::id()) | 1; } x ^= x << 13; x ^= x >> 7; x ^= x << 17; STATE.store(x, Ordering::Relaxed); x } /// Why the lock was not taken, which decides what the user is told. enum LockFailure { /// Someone else held it for the whole of [`WAIT`]. Contended, /// `flock` itself refused — the filesystem does not support it, or the /// handle went bad. Different advice entirely: waiting will not help. Unsupported(std::io::Error), } impl From for LockFailure { fn from(e: TryLockError) -> Self { match e { TryLockError::WouldBlock => LockFailure::Contended, TryLockError::Error(e) => LockFailure::Unsupported(e), } } } /// Open (creating if needed) the lock file. /// /// 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_at(path: &Path) -> Result { let mut options = std::fs::OpenOptions::new(); // Write access is what `flock` needs from this handle; the file's // contents are never read or written, so it stays empty forever. options.create(true).write(true).truncate(false); #[cfg(unix)] { use std::os::unix::fs::OpenOptionsExt; options.mode(0o600); } options .open(path) .with_context(|| format!("failed to open the lock file {}", path.display())) } /// Every lock this process is holding right now, in debug builds. /// /// The module's rule is that nothing holds two guards at once: the locks are /// `flock`, which is per open file description, so a second `take_async` for /// a scope this process already holds blocks on *itself* until [`WAIT`] runs /// out and then reports a conflict naming another process. Two guards on /// different scopes are worse, because two processes taking them in opposite /// orders deadlock until both give up. /// /// The rule was enforced by a comment and by `forget` scoping its guard in a /// block by hand. This makes it enforced by a panic in the build that tests /// run in, at the moment the second one is taken, with both scopes named -- /// rather than sixty seconds later as somebody else's fault. /// One logical flow of work, which is the thing that may not hold two locks. /// /// Not the process: the test binary runs its tests in parallel threads, each /// legitimately holding its own lock over its own directory, and a /// process-wide check calls that a violation. Not the thread either: a guard /// is held across `await`, and a multi-threaded runtime may resume the task /// somewhere else. The task when there is one, the thread when there is not. #[cfg(debug_assertions)] #[derive(Debug, Clone, PartialEq, Eq, Hash)] enum Flow { Task(tokio::task::Id), Thread(std::thread::ThreadId), } #[cfg(debug_assertions)] fn flow() -> Flow { match tokio::task::try_id() { Some(id) => Flow::Task(id), None => Flow::Thread(std::thread::current().id()), } } #[cfg(debug_assertions)] static HELD: std::sync::Mutex>>> = std::sync::Mutex::new(None); #[cfg(debug_assertions)] fn note_held(flow: &Flow, scope: &Scope) { let mut held = HELD.lock().unwrap_or_else(|e| e.into_inner()); let mine = held .get_or_insert_default() .entry(flow.clone()) .or_default(); assert!( mine.is_empty(), "atgc took the lock for {scope:?} while still holding {:?}. Nothing may \ hold two: `flock` is per open file description, so re-taking one this \ flow holds waits out the whole budget and then blames another \ process, and two scopes held at once deadlock against a peer taking \ them the other way round. Release the first -- `forget` scopes its \ guard in a block for exactly this reason -- or use `take_or_inherit` \ if this is the reentrant refresh path.", mine.last() ); mine.push(scope.clone()); } /// Unregister against the flow the guard was *taken* on, which is why the /// guard carries it: a task that moved threads would otherwise leave its /// entry behind and the next lock on that flow would look like a nesting. #[cfg(debug_assertions)] fn forget_held(flow: &Flow, scope: &Scope) { let mut held = HELD.lock().unwrap_or_else(|e| e.into_inner()); let Some(map) = held.as_mut() else { return }; if let Some(mine) = map.get_mut(flow) && let Some(at) = mine.iter().rposition(|s| s == scope) { mine.remove(at); if mine.is_empty() { map.remove(flow); } } } fn acquired(file: File, purpose: Purpose, scope: Scope, since: Instant) -> Guard { crate::logging::oauth::emit(crate::logging::oauth::Event::Lock { purpose, waited_ms: ms_since(since), acquired: true, }); #[cfg(debug_assertions)] let flow = { let flow = flow(); note_held(&flow, &scope); flow }; Guard { file: Some(file), purpose, scope, since: Instant::now(), registered: false, #[cfg(debug_assertions)] flow, } } // --------------------------------------------------------------------------- // 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. /// Keyed by scope, because there is a lock per identity now and a context /// holding one of them holds nothing about the others: a refresh of alice /// must not let a write for bob skip its own lock. 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 holders() -> std::sync::MutexGuard<'static, Option>>> { HOLDER.lock().unwrap_or_else(|e| e.into_inner()) } /// Who, if anyone, holds `scope` in this process. fn holder_of(scope: &Scope) -> Option> { holders().as_ref().and_then(|m| m.get(scope).copied()) } /// 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(config_dir: &Path, purpose: Purpose, scope: Scope) -> Result { let mut guard = take_async_in(config_dir, purpose, scope.clone()).await?; holders() .get_or_insert_with(Default::default) .insert(scope, 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(config_dir: &Path, purpose: Purpose, scope: Scope) -> Result { if holder_of(&scope) == 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_in(config_dir, purpose, scope).await?, )) } /// Record the failure and say what happened in terms of what to do about it. /// /// The contended case is the ordinary one and it is not really an error in /// the user's world — another atgc is mid-write — so it says that, rather /// than naming a lock file most people will never have heard of. fn gave_up(purpose: Purpose, since: Instant, why: LockFailure) -> anyhow::Error { let waited_ms = ms_since(since); crate::logging::oauth::emit(crate::logging::oauth::Event::Lock { purpose, waited_ms, acquired: false, }); let what = match purpose { Purpose::Refresh => "refresh this account's token", Purpose::Store => "update the session store", Purpose::Registry => "update the account registry", }; match why { // Contended is the textbook `Conflict`: nothing is wrong, another // atgc simply got there first, and the whole advice is "retry". LockFailure::Contended => crate::exit::fail( crate::exit::Exit::Conflict, format!( "timed out after {}s waiting for another atgc process\n\ it holds {}, needed to {what} without losing a live session\n\ retry once that command has finished", WAIT.as_secs(), lock_display(), ), ), LockFailure::Unsupported(e) => anyhow::anyhow!( "cannot lock {}: {e}\n\ atgc serializes account state with flock(2), which this filesystem does \ not support; without it a live session can be destroyed\n\ move $HOME to a local filesystem", lock_display(), ), } } /// The lock's path for an error message, or a plain description of it when /// even `HOME` cannot be resolved. fn lock_display() -> String { match crate::config::dir::config_dir_path() { Ok(dir) => match lock_path_in(&dir, &Scope::Global) { Ok(path) => path.display().to_string(), Err(_) => "the atgc config directory lock".to_string(), }, Err(_) => "the atgc config directory lock".to_string(), } } /// Refuse to go on without the lock, for callers whose failure is not /// worth propagating but whose *work* must not happen unlocked. /// /// Only `oauth::login::discard_auth_state` uses this: it runs on paths that /// are already failing, so replacing "the login timed out" with a lock error /// would report the wrong thing — but skipping a cleanup is harmless, while /// doing it unlocked is not. pub fn or_skip(taken: Result, what: &str) -> Option { match taken { Ok(guard) => Some(guard), Err(e) => { crate::logging::debug::log(format!("skipping {what}: {e}")); None } } } fn ms_since(start: Instant) -> u64 { start.elapsed().as_millis() as u64 } #[cfg(test)] mod tests { use super::*; /// A throwaway directory, named for the test so a leftover from a killed /// run says which one made it. fn temp_dir(label: &str) -> tempfile::TempDir { tempfile::Builder::new() .prefix(&format!("atgc-lock-{label}-")) // 0o700, because these stand in for ~/.config/atgc, which // production creates owner-only. `tempfile` defaults a // directory to 0o777 & ~umask, which would quietly make // every mode assertion below weaker than the real thing. .permissions( ::from_mode(0o700), ) .tempdir() .unwrap() } /// 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. /// The rule the module rests on, as a failure rather than a comment. /// /// Debug-only, like the check itself: a release build has no bookkeeping /// to consult, and the cost of the rule being broken there is the sixty /// seconds and the misleading conflict this exists to prevent. #[cfg(debug_assertions)] #[tokio::test] #[should_panic(expected = "while still holding")] async fn holding_two_locks_at_once_is_refused() { let home = tempfile::tempdir().expect("a temp dir"); let _first = take_async_in(home.path(), Purpose::Store, Scope::Global) .await .expect("nothing else holds it"); // A different scope, which is the deadlock-against-a-peer case rather // than the wait-on-yourself one. Both are the same rule. let _second = take_async_in( home.path(), Purpose::Store, Scope::Identity("did:plc:someone".into()), ) .await; } /// Two flows holding their own locks at once is the ordinary case and /// must not trip the check -- the test binary itself does it constantly. #[tokio::test] async fn two_flows_may_each_hold_one() { let home = tempfile::tempdir().expect("a temp dir"); let root = home.path().to_path_buf(); let peer = tokio::spawn(async move { let _held = take_async_in(&root, Purpose::Store, Scope::Identity("did:plc:two".into())) .await .expect("its own scope is free"); tokio::time::sleep(Duration::from_millis(150)).await; }); let _mine = take_async_in( home.path(), Purpose::Store, Scope::Identity("did:plc:one".into()), ) .await .expect("a different identity's lock is free"); peer.await.expect("the peer finished"); } #[tokio::test] async fn a_write_inside_the_refresh_section_inherits_the_lock() { let home = tempfile::tempdir().expect("a temp dir"); let held = take_reentrant(home.path(), Purpose::Refresh, Scope::Global) .await .expect("nothing else holds it"); let started = Instant::now(); let home = tempfile::tempdir().expect("a temp dir"); let inherited = take_or_inherit(home.path(), Purpose::Store, Scope::Global) .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_of(&Scope::Global), Some(tokio::task::try_id())); drop(held); assert_eq!( holder_of(&Scope::Global), 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 plan/credential-store.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 home = tempfile::tempdir().expect("a temp dir"); let held = take_reentrant(home.path(), Purpose::Refresh, Scope::Global) .await .expect("nothing else holds it"); assert_eq!( holder_of(&Scope::Global), 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_of(&Scope::Global) == 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); } /// 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_delay()).await; } held.unlock().unwrap(); }); let guard = take_async_at(&path, Purpose::Registry, Scope::Global) .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); } /// **The defect this change exists for**, as an assertion on the two /// numbers rather than on a stopwatch: a waiter must outlast the longest /// hold a holder is permitted to take. /// /// Before this, `WAIT` was thirty seconds and nothing bounded the hold at /// all — the critical section covered two uncached round trips under /// client deadlines summing past thirty seconds. Every other atgc on the /// machine then failed with `Conflict` while queued behind honest work, /// which is precisely what `Conflict` is supposed to mean is *not* /// happening. The compile-time assertion above catches the same thing; /// this states it where a reader looking for the property will find it. #[test] fn a_waiter_outlasts_the_longest_permitted_hold() { assert!( WAIT > HOLD_BUDGET, "a holder may take {HOLD_BUDGET:?} but a waiter gives up at {WAIT:?}", ); } /// The wait is derived from the hold rather than picked next to it, so /// this pins the derivation: raising what a refresh may take must raise /// what a waiter will sit through, automatically. /// /// Driven through the pure function on invented numbers, because the real /// constants would make this a restatement of their values instead of a /// test of the arithmetic that produced them. #[test] fn the_wait_covers_a_full_queue_of_holds_plus_head_room() { let hold = Duration::from_secs(10); let slack = Duration::from_secs(3); assert_eq!(wait_bound(hold, 1, slack), Duration::from_secs(13)); assert_eq!(wait_bound(hold, 2, slack), Duration::from_secs(23)); assert!( wait_bound(hold, 1, Duration::ZERO) > Duration::ZERO, "even the tightest depth has to leave room for one whole hold" ); assert_eq!( WAIT, wait_bound(HOLD_BUDGET, QUEUE_DEPTH, SLACK), "the shipped wait stopped being a function of the shipped hold" ); } /// Jitter has to stay inside the band it advertises. A draw that could /// land outside it would either make the poll interval long enough to /// waste a meaningful slice of `WAIT` or short enough to turn the wait /// into a spin on `flock`. /// /// Every residue is exercised rather than a sample, since the arithmetic /// is a modulo and the interesting values are its ends. #[test] fn a_jittered_poll_stays_inside_its_band() { let low = Duration::from_millis(POLL_MS * (100 - POLL_JITTER_PERCENT) / 100); let high = Duration::from_millis(POLL_MS * (100 + POLL_JITTER_PERCENT) / 100); for draw in 0..1_000u64 { let delay = jittered_ms(POLL_MS, POLL_JITTER_PERCENT, draw); assert!( delay >= low && delay <= high, "draw {draw} produced {delay:?}, outside {low:?}..={high:?}" ); } assert!(low > Duration::ZERO, "a zero sleep would spin on try_lock"); } /// The band must actually be a band. A jitter function that returned a /// constant would pass the bounds test above while leaving every waiter /// phase-aligned, which is the starvation this exists to break up. #[test] fn jitter_actually_spreads_the_poll_interval() { let drawn: std::collections::BTreeSet<_> = (0..1_000u64) .map(|draw| jittered_ms(POLL_MS, POLL_JITTER_PERCENT, draw)) .collect(); assert!( drawn.len() > 10, "the poll interval took only {} distinct values; pollers stay in step", drawn.len() ); } /// A zero spread has to be the identity rather than a special case, so /// that turning jitter off is a one-line change to a constant and not a /// change to the code path. #[test] fn no_spread_leaves_the_interval_alone() { for draw in 0..100u64 { assert_eq!( jittered_ms(POLL_MS, 0, draw), Duration::from_millis(POLL_MS) ); } } /// The entropy the jitter draws on must move. A generator stuck on its /// seed — xorshift's fixed point is zero, which is also this one's /// "unseeded" marker — would hand every poll the same offset and quietly /// undo the spreading above. #[test] fn the_entropy_source_does_not_stick() { let drawn: std::collections::BTreeSet<_> = (0..64).map(|_| entropy()).collect(); assert!( drawn.len() > 32, "the generator repeated itself: {} distinct draws in 64", drawn.len() ); assert!(!drawn.contains(&0), "the generator reached its fixed point"); } /// Work that finishes inside the budget must come back untouched, and /// nothing about the bound may leak into the ordinary path. #[tokio::test] async fn work_within_the_budget_is_returned_as_it_is() { let out = under_hold_budget("do the thing", async { Ok::<_, ()>(7) }) .await .expect("finishing promptly is not a failure"); assert_eq!(out, Ok(7)); } /// Overrunning the budget has to fail, and has to fail as /// `Unreachable` rather than `Conflict`: nothing is contended, the server /// simply stopped answering, and the two exit codes send a caller looking /// in different places. /// /// The clock is paused rather than slept through, so this asserts on the /// decision at `HOLD_BUDGET` without spending it. #[tokio::test] async fn overrunning_the_budget_fails_as_unreachable() { // A millisecond against a future that never completes: the decision // is what is asserted, and no scheduler is fast enough to make a // never-completing future finish first. let err = within( Duration::from_millis(1), "refresh this account's token", std::future::pending::<()>(), ) .await .expect_err("a future that never completes must not be waited on forever"); assert_eq!( crate::exit::classify(&err), crate::exit::Exit::Unreachable, "a stalled server was reported as contention: {err:#}" ); let text = format!("{err:#}"); assert!( text.contains("refresh this account's token"), "the message does not say what was given up on: {text}" ); } /// Cancellation is the whole mechanism: the timeout drops the future, and /// dropping the future is what releases the lock the refresh was holding. /// A bound that let the work carry on would leave every waiter queued /// behind exactly the hold it was meant to cut short. #[tokio::test] async fn overrunning_the_budget_drops_the_work() { struct Notice(std::sync::Arc); impl Drop for Notice { fn drop(&mut self) { self.0.store(true, std::sync::atomic::Ordering::SeqCst); } } let dropped = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)); let notice = Notice(dropped.clone()); let _ = within(Duration::from_millis(1), "hold something", async move { let _notice = notice; std::future::pending::<()>().await; }) .await; assert!( dropped.load(std::sync::atomic::Ordering::SeqCst), "the overrunning work was left running, so it is still holding the lock" ); } #[test] fn lock_file_sits_beside_the_files_it_protects() { let dir = Path::new("/home/someone/.config/atgc"); assert_eq!( lock_path_in(dir, &Scope::Global).expect("the global lock"), dir.join(".lock") ); } /// The property the whole module is for: while one handle holds the /// lock, a second one cannot take it. Two handles on the same file in /// one process is exactly what two processes look like to `flock`, which /// is what makes this testable without spawning anything — and it is /// also why nesting is banned, since this is the same collision a /// careless caller would hit against itself. #[test] fn a_second_holder_is_excluded_until_the_first_lets_go() { let dir = temp_dir("exclusion"); let path = dir.path().join(".lock"); let first = open_at(&path).expect("open the lock file"); first.try_lock().expect("take it"); let second = open_at(&path).expect("open it again"); assert!( matches!(second.try_lock(), Err(TryLockError::WouldBlock)), "a second holder got the lock while the first held it" ); first.unlock().unwrap(); second.try_lock().expect("the lock was not released"); second.unlock().unwrap(); } /// Dropping the guard has to release the lock, or the first command to /// take it would wedge every later one for 30 seconds apiece. Checked /// through `Guard`'s own `Drop` rather than by calling `unlock`, since /// the release being automatic is the part callers rely on. #[test] fn dropping_the_guard_releases_the_lock() { let dir = temp_dir("release"); let path = dir.path().join(".lock"); let file = open_at(&path).expect("open the lock file"); file.try_lock().expect("take it"); let guard = acquired(file, Purpose::Store, Scope::Global, Instant::now()); drop(guard); let after = open_at(&path).expect("open it again"); after.try_lock().expect("the guard did not release on drop"); after.unlock().unwrap(); std::fs::remove_dir_all(path.parent().unwrap()).ok(); } /// Contention has to end in an error, not in a guard. This is the /// behavior the module was changed to have: proceeding unlocked risks /// deleting a live session, and a command that fails loudly is the /// cheaper of the two outcomes by a wide margin. /// /// `WAIT` is not waited out — that would be a 30-second test. The /// contended path is driven directly, which is the same code `take` runs /// when its deadline passes. #[test] fn losing_the_lock_is_an_error_and_says_which_kind() { let contended = gave_up(Purpose::Refresh, Instant::now(), LockFailure::Contended); let text = format!("{contended:#}"); assert!(text.contains("another atgc process"), "{text}"); assert!( text.contains("refresh this account's token"), "the message does not say what was refused: {text}" ); let unsupported = gave_up( Purpose::Store, Instant::now(), LockFailure::Unsupported(std::io::Error::from_raw_os_error(37)), ); let text = format!("{unsupported:#}"); assert!( text.contains("flock(2)"), "a filesystem that cannot lock needs different advice: {text}" ); } /// **The property the module exists for**, stated as the bug it /// prevents: a read-modify-write cycle run concurrently must not lose an /// update. /// /// Two threads stand in for two atgc processes, which is not an /// approximation — `flock` excludes per *open file description*, so two /// handles in one process contend exactly as two processes do, and that /// is the same fact the ban on nesting rests on. /// /// The counter file is `sessions.json` in miniature: read it whole, /// change it, write it back. Without the lock this is the interleaving /// from the incident — both threads read `n`, both write `n+1`, one /// increment vanishes — so an exact final count is the assertion that /// the exclusion held every single time, not merely most of the time. /// `take` itself is not called because it resolves `$HOME`, which a /// hermetic test cannot redirect (see [`crate::docs::testing`]); this drives the same /// file and the same `try_lock` through a temp path. #[test] fn the_lock_stops_a_concurrent_read_modify_write_losing_an_update() { const THREADS: usize = 4; const EACH: usize = 25; let dir = temp_dir("lost-update"); let lock = dir.path().join(".lock"); let counter = dir.path().join("counter"); std::fs::write(&counter, "0").unwrap(); std::thread::scope(|scope| { for _ in 0..THREADS { scope.spawn(|| { for _ in 0..EACH { let file = open_at(&lock).expect("open the lock file"); while file.try_lock().is_err() { std::thread::sleep(Duration::from_millis(1)); } // The window a lost update needs: a read, a gap, a // write of something derived from that read. The // yield widens it so that a broken lock fails this // test reliably rather than occasionally. let n: usize = std::fs::read_to_string(&counter) .unwrap() .trim() .parse() .unwrap(); std::thread::yield_now(); std::fs::write(&counter, (n + 1).to_string()).unwrap(); file.unlock().unwrap(); } }); } }); let total: usize = std::fs::read_to_string(&counter) .unwrap() .trim() .parse() .unwrap(); assert_eq!( total, THREADS * EACH, "an update was lost: the lock did not exclude every overlap" ); std::fs::remove_dir_all(&dir).ok(); } /// Two identities do not exclude each other; one identity still does. /// /// This is the whole of why the scope moved. The refresh critical section /// is held **across the token request**, so a global lock put every /// account's network round trip in one queue — on a machine running an /// identity per agent, the fleet behind an exclusion none of them needed /// from each other. What still has to serialise is two refreshes of *one* /// session, because a refresh token is single-use and spending it twice /// deletes the session. /// /// Asserted on the lock files themselves rather than on timing: two /// scopes that resolve to one path cannot help but exclude, and two that /// resolve to different paths cannot exclude at all. #[test] fn a_lock_covers_one_identity_and_not_its_neighbours() { let root = tempfile::tempdir().expect("a temp dir"); let alice = Scope::Identity("did:plc:aaaaaaaaaaaaaaaaaaaaaaaa".into()); let bob = Scope::Identity("did:web:bob.example".into()); let a = lock_path_in(root.path(), &alice).expect("alice's lock"); let b = lock_path_in(root.path(), &bob).expect("bob's lock"); let global = lock_path_in(root.path(), &Scope::Global).expect("the global lock"); assert_ne!(a, b, "two identities share a lock"); assert_ne!(a, global, "an identity shares the global lock"); assert_eq!( a, lock_path_in(root.path(), &alice).expect("alice's lock again"), "one identity resolves to two locks" ); assert_eq!( a, root.path() .join("dids/plc/aaaaaaaaaaaaaaaaaaaaaaaa") .join(".lock"), "an identity's lock is not in its own directory" ); } /// A session key names the identity whose lock covers it, and an /// authorization request — which has no identity yet — falls to global. #[test] fn a_store_key_says_which_lock_covers_it() { assert_eq!( Scope::of_key("oauth:did:plc:alice/session-one"), Scope::Identity("did:plc:alice".into()) ); // A did:web carries dots and a percent-encoded port; the split is on // the last slash, which only the session id introduces. assert_eq!( Scope::of_key("oauth:did:web:localhost%3A3000/s"), Scope::Identity("did:web:localhost%3A3000".into()) ); assert_eq!(Scope::of_key("oauth-state:abc123"), Scope::Global); } /// Inheriting is scoped: holding one identity's lock is not holding /// another's. /// /// The failure this guards is silent and is the module's whole subject. /// [`take_or_inherit`] answers "this context already holds it" by looking /// the holder up; if that lookup ignored which scope was asked about, a /// refresh holding alice's lock would let a write for bob report /// `Inherited` and go to disk **with no lock at all** — an unlocked /// read-modify-write on somebody else's sessions, which is how the /// incident in this module's documentation happened. /// /// Asserted on the `Hold` rather than on timing, because the wrong answer /// here is fast, not slow. #[tokio::test] async fn holding_one_identitys_lock_is_not_holding_anothers() { let alice = Scope::Identity("did:plc:aaaaaaaaaaaaaaaaaaaaaaaa".into()); let bob = Scope::Identity("did:web:bob.example".into()); // Published without touching the filesystem: this is about the // bookkeeping `take_or_inherit` consults, and `take_reentrant` would // need a real config directory to open a lock file in. holders() .get_or_insert_with(Default::default) .insert(alice.clone(), tokio::task::try_id()); assert_eq!( holder_of(&alice), Some(tokio::task::try_id()), "the scope that was published is not held" ); assert_eq!( holder_of(&bob), None, "holding one identity reported as holding another" ); assert_eq!( holder_of(&Scope::Global), None, "holding an identity reported as holding the global lock" ); holders().as_mut().expect("the map").remove(&alice); } }