repositories / dotfiles
dotfiles
bugabingas dorkfiles
owned by admin
quickshell/nuguland/vault-native/src/session.rs
Rawuse 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");
}
}
}