Luigit
repositories / dotfiles

dotfiles

bugabingas dorkfiles

owned by admin

quickshell/nuguland/vault-native/src/session.rs

Raw
use super::{Profile, ProtectedUnlockData, read_secret, settings, store_secret};
use bitwarden_core::{UserId, auth::JwtToken};
use serde::{Deserialize, Serialize};
use std::{
    io,
    time::{Duration, SystemTime, UNIX_EPOCH},
};
use tokio::{sync::oneshot, task::JoinSet};

#[derive(Default)]
pub(super) struct SessionTasks {
    tasks: JoinSet<()>,
}

impl SessionTasks {
    // Credential transactions outlive request cancellation, but never retain vault keys.
    pub async fn settle(&mut self) -> io::Result<()> {
        while let Some(result) = self.tasks.join_next().await {
            result.map_err(io::Error::other)?;
        }
        Ok(())
    }

    pub async fn save_login(&mut self, profile: &Profile, session: SavedSession) -> io::Result<()> {
        self.settle().await?;
        let profile = profile.clone();
        let (sender, receiver) = oneshot::channel();
        self.tasks.spawn(async move {
            let result = store_session(&profile, &session).await.and_then(|()| {
                match std::fs::remove_file(revocation_path(&profile)) {
                    Ok(()) => Ok(()),
                    Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()),
                    Err(error) => Err(error),
                }
            });
            drop(sender.send(result));
        });
        receiver.await.map_err(io::Error::other)?
    }

    pub async fn update_unlock(
        &mut self,
        profile: &Profile,
        unlock: &ProtectedUnlockData,
    ) -> io::Result<()> {
        self.settle().await?;
        let (profile, unlock) = (profile.clone(), unlock.clone());
        let (sender, receiver) = oneshot::channel();
        self.tasks.spawn(async move {
            drop(sender.send(update_unlock(&profile, &unlock).await));
        });
        receiver.await.map_err(io::Error::other)?
    }

    pub async fn refresh(&mut self, profile: &Profile, user_id: UserId) -> RefreshResult {
        if self.settle().await.is_err() {
            return RefreshResult::Unavailable;
        }
        let (sender, receiver) = oneshot::channel();
        let profile = profile.clone();
        self.tasks.spawn(refresh_session(profile, user_id, sender));
        receiver.await.unwrap_or(RefreshResult::Unavailable)
    }
}

// Non-secret, durable state: a failed keyring cleanup must not revive a revoked login.
fn revocation_path(profile: &Profile) -> std::path::PathBuf {
    std::env::var_os("XDG_STATE_HOME")
        .map_or_else(
            || {
                std::path::PathBuf::from(std::env::var_os("HOME").unwrap_or_default())
                    .join(".local/state")
            },
            std::path::PathBuf::from,
        )
        .join("nuguland/vault")
        .join(format!("{}.reauthenticate", profile.storage_key()))
}
use zeroize::{Zeroize as _, Zeroizing};

#[derive(Serialize, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub(super) struct SavedSession {
    pub refresh_token: String,
    pub unlock: ProtectedUnlockData,
}

impl Drop for SavedSession {
    fn drop(&mut self) {
        self.refresh_token.zeroize();
    }
}

pub(super) enum RefreshResult {
    Ready(Zeroizing<String>),
    Reauthenticate,
    RevocationFailed,
    Unavailable,
}

pub(super) fn access_valid(token: &str, user_id: UserId) -> bool {
    let now = SystemTime::now()
        .duration_since(UNIX_EPOCH)
        .unwrap_or_default()
        .as_secs();
    token.parse::<JwtToken>().is_ok_and(|jwt| {
        jwt.sub.parse::<UserId>().ok() == Some(user_id) && jwt.exp > now.saturating_add(30)
    })
}

pub(super) async fn read_session(profile: &Profile) -> io::Result<Option<SavedSession>> {
    match std::fs::symlink_metadata(revocation_path(profile)) {
        Ok(_) => return Ok(None),
        Err(error) if error.kind() == io::ErrorKind::NotFound => {}
        Err(error) => return Err(error),
    }
    let Some(bytes) = tokio::time::timeout(
        Duration::from_secs(5),
        read_secret(&profile.storage_key(), "kind", "session"),
    )
    .await
    .map_err(io::Error::other)??
    else {
        return Ok(None);
    };
    let session: SavedSession = serde_json::from_slice(&bytes).map_err(io::Error::other)?;
    if session.refresh_token.is_empty()
        || session.unlock.server_url != profile.server_url
        || session.unlock.account != profile.account
    {
        return Err(io::Error::other(
            "saved login belongs to another profile or is invalid",
        ));
    }
    Ok(Some(session))
}

pub(super) async fn store_session(profile: &Profile, session: &SavedSession) -> io::Result<()> {
    if session.refresh_token.is_empty()
        || session.unlock.server_url != profile.server_url
        || session.unlock.account != profile.account
    {
        return Err(io::Error::other("saved login identity mismatch"));
    }
    let bytes = Zeroizing::new(serde_json::to_vec(session).map_err(io::Error::other)?);
    tokio::time::timeout(
        Duration::from_secs(5),
        store_secret(&profile.storage_key(), "kind", "session", &bytes),
    )
    .await
    .map_err(io::Error::other)?
}

pub(super) async fn update_unlock(
    profile: &Profile,
    unlock: &ProtectedUnlockData,
) -> io::Result<()> {
    let mut session = read_session(profile)
        .await?
        .ok_or_else(|| io::Error::other("saved login unavailable"))?;
    if session.unlock.user_id != unlock.user_id {
        return Err(io::Error::other("saved login identity mismatch"));
    }
    session.unlock = unlock.clone();
    store_session(profile, &session).await
}

#[derive(Deserialize)]
struct RefreshedTokens {
    access_token: String,
    refresh_token: Option<String>,
}

impl Drop for RefreshedTokens {
    fn drop(&mut self) {
        self.access_token.zeroize();
        if let Some(token) = &mut self.refresh_token {
            token.zeroize();
        }
    }
}

// The client-managed SDK handler accepts access tokens but does not renew them.
// Only a definitive invalid_grant response requests interactive sign-in.
async fn refresh_session(
    profile: Profile,
    user_id: UserId,
    mut sender: oneshot::Sender<RefreshResult>,
) {
    let mut session = match read_session(&profile).await {
        Ok(Some(session)) if session.unlock.user_id == user_id => session,
        result => {
            drop(sender.send(if matches!(result, Ok(None)) {
                RefreshResult::Reauthenticate
            } else {
                RefreshResult::Unavailable
            }));
            return;
        }
    };
    let result = match request_refresh(&settings(&profile).identity_url, &mut session).await {
        RefreshResult::Ready(mut access) => {
            let store = store_session(&profile, &session);
            tokio::pin!(store);
            let saved = tokio::select! {
                result = &mut store => result,
                () = sender.closed() => { access.zeroize(); store.await },
            };
            if saved.is_ok() {
                RefreshResult::Ready(access)
            } else {
                RefreshResult::Unavailable
            }
        }
        RefreshResult::Reauthenticate => {
            // Publish synchronously before yielding or notifying a cancellable caller.
            if super::write_private_json(
                &revocation_path(&profile),
                &serde_json::json!({"reauthenticate":true}),
            )
            .is_err()
            {
                RefreshResult::RevocationFailed
            } else {
                RefreshResult::Reauthenticate
            }
        }
        result => result,
    };
    drop(sender.send(result));
}

pub(super) async fn request_refresh(
    identity_url: &str,
    session: &mut SavedSession,
) -> RefreshResult {
    let Ok(client) = reqwest::Client::builder()
        .timeout(Duration::from_secs(15))
        .redirect(reqwest::redirect::Policy::none())
        .build()
    else {
        return RefreshResult::Unavailable;
    };
    let response = client
        .post(format!(
            "{}/connect/token",
            identity_url.trim_end_matches('/')
        ))
        .form(&[
            ("grant_type", "refresh_token"),
            ("client_id", "connector"),
            ("refresh_token", session.refresh_token.as_str()),
        ])
        .send()
        .await;
    let Ok(mut response) = response else {
        return RefreshResult::Unavailable;
    };
    let status = response.status();
    if response
        .content_length()
        .is_some_and(|length| length > 1_048_576)
    {
        return RefreshResult::Unavailable;
    }
    let mut bytes = Zeroizing::new(Vec::new());
    loop {
        match response.chunk().await {
            Ok(Some(chunk)) if bytes.len() + chunk.len() <= 1_048_576 => {
                bytes.extend_from_slice(&chunk);
            }
            Ok(None) => break,
            _ => return RefreshResult::Unavailable,
        }
    }
    if !status.is_success() {
        return if status == reqwest::StatusCode::BAD_REQUEST
            && serde_json::from_slice::<serde_json::Value>(&bytes)
                .is_ok_and(|error| error["error"] == "invalid_grant")
        {
            RefreshResult::Reauthenticate
        } else {
            RefreshResult::Unavailable
        };
    }
    let Ok(mut tokens) = serde_json::from_slice::<RefreshedTokens>(&bytes) else {
        return RefreshResult::Unavailable;
    };
    if !access_valid(&tokens.access_token, session.unlock.user_id)
        || tokens.refresh_token.as_ref().is_some_and(String::is_empty)
    {
        return RefreshResult::Unavailable;
    }
    if let Some(rotated) = tokens.refresh_token.take() {
        session.refresh_token.zeroize();
        session.refresh_token = rotated;
    }
    RefreshResult::Ready(Zeroizing::new(std::mem::take(&mut tokens.access_token)))
}

#[cfg(test)]
pub(super) mod tests {
    use super::*;
    use crate::{App, OpenResult, PendingLogin, TokenSuccess};
    use bitwarden_crypto::{HashPurpose, Kdf, MasterKey};
    use serde_json::json;
    use std::{
        collections::VecDeque,
        fs,
        path::Path,
        sync::{
            Arc, Mutex,
            atomic::{AtomicUsize, Ordering},
        },
    };
    use tokio::io::{AsyncBufReadExt as _, AsyncReadExt as _, AsyncWriteExt as _, BufReader};

    const ACCESS: &str = "e30.eyJleHAiOjQxMDI0NDQ4MDAsInN1YiI6IjAwMDAwMDAwLTAwMDAtNDAwMC04MDAwLTAwMDAwMDAwMDAwMSIsInNjb3BlIjpbImFwaSIsIm9mZmxpbmVfYWNjZXNzIl19.synthetic-signature";
    type Reply = Arc<Mutex<(&'static str, String)>>;
    type ApiStatuses = Arc<Mutex<VecDeque<u16>>>;

    pub(crate) async fn verify_saved_login(directory: &Path, original: &ProtectedUnlockData) {
        let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
            .await
            .expect("session server");
        // Configure validates HTTPS in production; this isolated fixture uses loopback HTTP.
        let profile = Profile {
            server_url: format!("http://{}", listener.local_addr().expect("address")),
            account: original.account.clone(),
            ..Profile::default()
        };
        let mut unlock = original.clone();
        unlock.server_url = profile.server_url.clone();
        let reply: Reply = Arc::new(Mutex::new((
            "200 OK",
            json!({"access_token":ACCESS,"refresh_token":"refresh-rotated"}).to_string(),
        )));
        let grants = Arc::new(AtomicUsize::new(0));
        let passwords = Arc::new(AtomicUsize::new(0));
        let api_statuses = ApiStatuses::default();
        let server = tokio::spawn(serve_tokens(
            listener,
            reply.clone(),
            grants.clone(),
            passwords.clone(),
            api_statuses.clone(),
        ));
        assert!(
            read_session(&profile)
                .await
                .expect("lookup empty keyring")
                .is_none()
        );
        let mut app = App::new(profile.clone());
        first_login(&mut app, &unlock).await;
        assert_eq!(
            read_session(&profile)
                .await
                .expect("saved login")
                .expect("session")
                .refresh_token,
            "refresh-initial"
        );
        app.lock().await;
        assert!(app.runtime.is_none());
        assert!(
            read_session(&profile)
                .await
                .expect("session survives lock")
                .is_some()
        );
        drop(app);

        let mut app = App::new(profile.clone());
        Box::pin(app.unlock("synthetic-password".to_owned())).await;
        assert!(app.runtime.is_some());
        assert_eq!(app.status, "unlocked");
        assert_eq!(grants.load(Ordering::SeqCst), 0, "restart unlock is local");
        assert_eq!(passwords.load(Ordering::SeqCst), 0);
        assert_eq!(app.public_state()["syncNeeded"], true);
        assert!(app.ensure_access().await);
        assert_eq!(grants.load(Ordering::SeqCst), 1);
        assert_eq!(
            read_session(&profile)
                .await
                .expect("rotated login")
                .expect("session")
                .refresh_token,
            "refresh-rotated"
        );
        assert!(app.ensure_access().await);
        assert_eq!(grants.load(Ordering::SeqCst), 1, "valid access is reused");
        verify_api_renewal(&mut app, &api_statuses, &grants).await;
        Box::pin(verify_cancelled_rotation(&mut app, directory, &reply)).await;
        verify_unreleased_rotation(directory, &profile).await;
        Box::pin(verify_network_and_cache_loss(
            &mut app, directory, &profile, unlock, &reply, &passwords,
        ))
        .await;
        server.abort();
        assert!(server.await.expect_err("fixture stopped").is_cancelled());
    }

    async fn verify_network_and_cache_loss(
        app: &mut App,
        directory: &Path,
        profile: &Profile,
        unlock: ProtectedUnlockData,
        reply: &Reply,
        passwords: &AtomicUsize,
    ) {
        app.runtime
            .as_ref()
            .expect("runtime")
            .token
            .0
            .write()
            .await
            .take();
        *reply.lock().expect("response") = ("timeout", String::new());
        assert!(!app.ensure_access().await);
        assert!(
            app.runtime.is_some(),
            "network loss keeps local vault usable"
        );
        assert!(
            read_session(profile)
                .await
                .expect("offline login preserved")
                .is_some()
        );
        verify_refresh_failures(directory, profile, reply).await;
        *reply.lock().expect("response") = (
            "400 Bad Request",
            json!({"error":"invalid_grant"}).to_string(),
        );
        assert!(!app.ensure_access().await);
        assert!(app.runtime.is_none(), "revocation drops decrypted keys");
        assert!(
            read_session(profile)
                .await
                .expect("revoked session lookup")
                .is_none()
        );
        Box::pin(app.unlock("synthetic-password".to_owned())).await;
        assert!(
            app.runtime.is_none(),
            "known revocation cannot fall back to offline login"
        );
        assert_eq!(
            passwords.load(Ordering::SeqCst),
            1,
            "only revoked login invokes prelogin"
        );

        SessionTasks::default()
            .save_login(
                profile,
                SavedSession {
                    refresh_token: "refresh-initial".to_owned(),
                    unlock,
                },
            )
            .await
            .expect("restore fixture session");
        fs::remove_file(crate::cache_path(profile)).expect("evict disposable cache");
        let mut empty = App::new(profile.clone());
        Box::pin(empty.unlock("synthetic-password".to_owned())).await;
        assert_eq!(empty.status, "loading");
        assert!(
            empty.runtime.is_some(),
            "wrapped key restores without cache"
        );
        *reply.lock().expect("response") = (
            "200 OK",
            json!({"access_token":ACCESS,"refresh_token":"refresh-rotated"}).to_string(),
        );
        Box::pin(empty.sync()).await;
        assert_eq!(empty.status, "unlocked");
        assert!(
            crate::read_private_json(&crate::cache_path(profile)).is_ok(),
            "sync rebuilt cache without email verification"
        );
        assert_eq!(passwords.load(Ordering::SeqCst), 1);
        empty.lock().await;
    }

    async fn first_login(app: &mut App, unlock: &ProtectedUnlockData) {
        let Kdf::PBKDF2 { iterations } = &unlock.master_password_unlock.kdf else {
            panic!("PBKDF2 fixture");
        };
        let crate::WrappedAccountCryptographicState::V1 { private_key } =
            &unlock.account_cryptographic_state
        else {
            panic!("V1 fixture");
        };
        let success: TokenSuccess = serde_json::from_value(json!({
            "access_token": ACCESS,
            "refresh_token": "refresh-initial",
            "PrivateKey": private_key.to_string(),
            "UserDecryptionOptions": {"MasterPasswordUnlock": {
                "kdf": {"kdfType":0,"iterations":iterations.get()},
                "masterKeyEncryptedUserKey": unlock.master_password_unlock.master_key_wrapped_user_key.to_string(),
                "salt": unlock.master_password_unlock.salt,
            }},
        })).expect("sign-in response fixture");
        let master_key = MasterKey::derive(
            "synthetic-password",
            &unlock.master_password_unlock.salt,
            &unlock.master_password_unlock.kdf,
        )
        .expect("fixture master key");
        let pending = PendingLogin {
            local_password_hash: Zeroizing::new(
                master_key
                    .derive_master_key_hash(b"synthetic-password", HashPurpose::LocalAuthorization)
                    .to_string(),
            ),
            server_password_hash: Zeroizing::new("synthetic-derived-hash".to_owned()),
            master_key,
            providers: Vec::new(),
        };
        let result = Box::pin(app.finish_login(success, pending)).await;
        assert!(
            matches!(&result, Ok(OpenResult::Opened(_, true))),
            "first sign-in opens and syncs"
        );
        app.apply_open_result(result);
    }

    async fn refresh_session(profile: &Profile, user_id: UserId) -> RefreshResult {
        let mut tasks = SessionTasks::default();
        let result = tasks.refresh(profile, user_id).await;
        tasks
            .settle()
            .await
            .expect("credential transaction completed");
        result
    }

    async fn verify_refresh_failures(directory: &Path, profile: &Profile, reply: &Reply) {
        let saved = read_session(profile)
            .await
            .expect("saved session")
            .expect("session");
        let other = "00000000-0000-4000-8000-000000000002"
            .parse()
            .expect("other identity");
        assert!(!access_valid(ACCESS, other));
        assert!(!access_valid("invalid", saved.unlock.user_id));
        assert!(!access_valid(
            "e30.eyJleHAiOjAsInN1YiI6IjAwMDAwMDAwLTAwMDAtNDAwMC04MDAwLTAwMDAwMDAwMDAwMSIsInNjb3BlIjpbXX0.synthetic-signature",
            saved.unlock.user_id
        ));
        for (status, body) in [
            (
                "503 Service Unavailable",
                json!({"error":"invalid_grant"}).to_string(),
            ),
            ("200 OK", "not json".to_owned()),
            (
                "200 OK",
                json!({"access_token":"invalid","refresh_token":"must-not-store"}).to_string(),
            ),
            (
                "200 OK",
                json!({"access_token":ACCESS,"refresh_token":""}).to_string(),
            ),
            ("302 Found", "{}".to_owned()),
            ("large", "{}".to_owned()),
        ] {
            *reply.lock().expect("response") = (status, body);
            assert!(matches!(
                refresh_session(profile, saved.unlock.user_id).await,
                RefreshResult::Unavailable
            ));
            assert_eq!(
                read_session(profile)
                    .await
                    .expect("preserved login")
                    .expect("session")
                    .refresh_token,
                saved.refresh_token
            );
        }
        *reply.lock().expect("response") = (
            "200 OK",
            json!({"access_token":ACCESS,"refresh_token":"must-not-publish"}).to_string(),
        );
        let blocked = directory.join("fail-store");
        fs::write(&blocked, b"").expect("fail protected writes");
        assert!(
            matches!(
                refresh_session(profile, saved.unlock.user_id).await,
                RefreshResult::Unavailable
            ),
            "do not use access credentials before refresh rotation is stored"
        );
        fs::remove_file(blocked).expect("restore protected writes");
        assert_eq!(
            read_session(profile)
                .await
                .expect("old login")
                .expect("session")
                .refresh_token,
            saved.refresh_token
        );
        *reply.lock().expect("response") = ("200 OK", json!({"access_token":ACCESS}).to_string());
        assert!(matches!(
            refresh_session(profile, saved.unlock.user_id).await,
            RefreshResult::Ready(_)
        ));
        assert_eq!(
            read_session(profile)
                .await
                .expect("unrotated login")
                .expect("session")
                .refresh_token,
            saved.refresh_token
        );
    }

    async fn verify_api_renewal(app: &mut App, statuses: &ApiStatuses, grants: &AtomicUsize) {
        let id = "00000000-0000-4000-8000-000000000099";
        for favorite in [false, true] {
            if favorite {
                app.runtime.as_mut().expect("runtime").encrypted.push(serde_json::from_value(json!({
                    "id":id, "collectionIds":[], "type":1, "favorite":false, "reprompt":0,
                    "organizationUseTotp":false, "edit":true, "viewPassword":true,
                    "creationDate":"2026-01-01T00:00:00Z", "revisionDate":"2026-01-01T00:00:00Z",
                })).expect("favorite fixture"));
            }
            for (responses, refreshes, cleared) in [
                (vec![200, 503], 0, false),
                (vec![503, 200], 0, false),
                (vec![401, 200, 503], 1, false),
                (vec![401, 401, 200], 1, true),
            ] {
                *app.runtime.as_ref().expect("runtime").token.0.write().await =
                    Some(ACCESS.to_owned());
                *statuses.lock().expect("API script") = responses.into();
                let before = grants.load(Ordering::SeqCst);
                let request = if favorite {
                    crate::Request::Favorite {
                        id: id.to_owned(),
                        favorite: true,
                    }
                } else {
                    crate::Request::Sync
                };
                Box::pin(crate::handle_request(app, request)).await;
                assert_eq!(grants.load(Ordering::SeqCst) - before, refreshes);
                assert_eq!(
                    statuses.lock().expect("API script").len(),
                    1,
                    "retry is bounded to one"
                );
                assert_eq!(
                    app.runtime
                        .as_ref()
                        .expect("runtime")
                        .token
                        .0
                        .read()
                        .await
                        .is_none(),
                    cleared
                );
                statuses.lock().expect("API script").clear();
            }
        }
    }

    async fn verify_cancelled_rotation(app: &mut App, directory: &Path, reply: &Reply) {
        let release = directory.join("session-release");
        let block = directory.join("block-session-store");
        let started = directory.join("session-store-started");
        for eof in [false, true] {
            if app.runtime.is_none() {
                Box::pin(app.unlock("synthetic-password".to_owned())).await;
            }
            app.runtime.as_ref().expect("runtime").token.clear().await;
            *reply.lock().expect("response") = (
                "200 OK",
                json!({"access_token":ACCESS,"refresh_token":"refresh-after-cancel"}).to_string(),
            );
            fs::write(&block, b"").expect("block session storage");
            let (released, locked) = oneshot::channel::<()>();
            app.runtime
                .as_mut()
                .expect("runtime")
                .favicon_tasks
                .spawn(async move {
                    let _released = released;
                    std::future::pending::<()>().await;
                });
            let (mut input, reader) = tokio::io::duplex(512);
            let (_suspend, events) = tokio::sync::mpsc::unbounded_channel();
            let replacement = App::new(app.profile.clone());
            let mut owned = std::mem::replace(app, replacement);
            let service = tokio::spawn(async move {
                crate::serve(&mut owned, BufReader::new(reader), events).await;
                owned
            });
            input
                .write_all(b"{\"op\":\"sync\"}\n")
                .await
                .expect("sync request");
            tokio::time::timeout(Duration::from_secs(2), async {
                while !started.exists() {
                    tokio::time::sleep(Duration::from_millis(5)).await;
                }
            })
            .await
            .expect("rotation reached persistence barrier");
            if eof {
                input.shutdown().await.expect("EOF");
            } else {
                input
                    .write_all(b"{\"op\":\"lock\"}\n")
                    .await
                    .expect("lock request");
            }
            assert!(
                tokio::time::timeout(Duration::from_secs(2), locked)
                    .await
                    .expect("keys cleared without waiting for storage")
                    .is_err()
            );
            // A cancelled caller cannot abort the credential commit. A file marker cannot
            // strand a blocking FIFO writer if the helper has already exited.
            fs::write(&release, b"release").expect("release commit");
            input.shutdown().await.expect("close request stream");
            *app = tokio::time::timeout(Duration::from_secs(2), service)
                .await
                .expect("credential drain completes")
                .expect("service");
            assert!(app.runtime.is_none());
            assert_eq!(
                read_session(&app.profile)
                    .await
                    .expect("persisted session")
                    .expect("login")
                    .refresh_token,
                "refresh-after-cancel"
            );
            fs::remove_file(&block).expect("unblock session storage");
            fs::remove_file(&started).expect("remove barrier");
            fs::remove_file(&release).expect("remove release marker");
        }
        Box::pin(app.unlock("synthetic-password".to_owned())).await;
    }

    async fn verify_unreleased_rotation(directory: &Path, profile: &Profile) {
        let mut saved = read_session(profile)
            .await
            .expect("saved session")
            .expect("login");
        let previous = saved.refresh_token.clone();
        saved.refresh_token = "must-not-publish-without-release".to_owned();
        let block = directory.join("block-session-store");
        let started = directory.join("session-store-started");
        fs::write(&block, b"").expect("withhold commit release");
        // Test the helper's own deadline, not a race with store_session's five-second timeout.
        let bytes = serde_json::to_vec(&saved).expect("synthetic session");
        let error = tokio::time::timeout(
            Duration::from_secs(15),
            store_secret(&profile.storage_key(), "kind", "session", &bytes),
        )
        .await
        .expect("unreleased helper exits within the test deadline")
        .expect_err("unreleased helper must exit");
        assert_eq!(error.to_string(), "protected storage write failed");
        let pid: u32 = fs::read_to_string(&started)
            .expect("helper reached barrier")
            .parse()
            .expect("helper pid");
        assert!(
            !Path::new(&format!("/proc/{pid}")).exists(),
            "helper reaped"
        );
        assert_eq!(
            read_session(profile)
                .await
                .expect("previous session remains readable")
                .expect("login")
                .refresh_token,
            previous
        );
        fs::remove_file(block).expect("restore storage");
        fs::remove_file(started).expect("remove expired barrier");
    }

    async fn serve_tokens(
        listener: tokio::net::TcpListener,
        reply: Reply,
        grants: Arc<AtomicUsize>,
        passwords: Arc<AtomicUsize>,
        api_statuses: ApiStatuses,
    ) {
        loop {
            let (socket, _) = listener.accept().await.expect("fixture connection");
            let mut reader = BufReader::new(socket);
            let mut first = String::new();
            reader.read_line(&mut first).await.expect("request line");
            let mut length = 0;
            let mut authorized = false;
            loop {
                let mut line = String::new();
                reader.read_line(&mut line).await.expect("request header");
                if line == "\r\n" {
                    break;
                }
                if let Some(value) = line.to_ascii_lowercase().strip_prefix("content-length:") {
                    length = value.trim().parse::<usize>().expect("request length");
                }
                if line
                    .strip_prefix("authorization:")
                    .is_some_and(|value| value.trim().strip_prefix("Bearer ") == Some(ACCESS))
                {
                    authorized = true;
                }
            }
            assert!(length < 4096);
            let mut bytes = vec![0; length];
            reader.read_exact(&mut bytes).await.expect("request body");
            let (status, body) =
                if first.starts_with("GET /api/sync") || first.starts_with("PUT /api/ciphers/") {
                    assert!(authorized, "API request uses renewed access");
                    let status = match api_statuses
                        .lock()
                        .expect("API status")
                        .pop_front()
                        .unwrap_or(200)
                    {
                        200 => "200 OK",
                        401 => "401 Unauthorized",
                        503 => "503 Service Unavailable",
                        _ => panic!("unsupported fixture status"),
                    };
                    (status, "{}".to_owned())
                } else if first.starts_with("POST /identity/accounts/prelogin/password") {
                    passwords.fetch_add(1, Ordering::SeqCst);
                    ("503 Service Unavailable", "{}".to_owned())
                } else {
                    assert!(first.starts_with("POST /identity/connect/token"));
                    let form: Vec<(String, String)> =
                        serde_urlencoded::from_bytes(&bytes).expect("refresh form");
                    assert_eq!(form.len(), 3, "no master password or email code is sent");
                    assert_eq!(
                        form[0],
                        ("grant_type".to_owned(), "refresh_token".to_owned())
                    );
                    assert_eq!(form[1], ("client_id".to_owned(), "connector".to_owned()));
                    assert_eq!(form[2].0, "refresh_token");
                    grants.fetch_add(1, Ordering::SeqCst);
                    reply.lock().expect("response").clone()
                };
            if status == "timeout" {
                tokio::time::sleep(Duration::from_secs(16)).await;
                continue;
            }
            let length = if status == "large" {
                1_048_577
            } else {
                body.len()
            };
            let status = if status == "large" { "200 OK" } else { status };
            reader.get_mut().write_all(format!("HTTP/1.1 {status}\r\nContent-Type: application/json\r\nContent-Length: {length}\r\nConnection: close\r\n\r\n{body}").as_bytes()).await.expect("fixture response");
        }
    }
}