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), 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::().is_ok_and(|jwt| { jwt.sub.parse::().ok() == Some(user_id) && jwt.exp > now.saturating_add(30) }) } pub(super) async fn read_session(profile: &Profile) -> io::Result> { 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, } 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, ) { 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::(&bytes) .is_ok_and(|error| error["error"] == "invalid_grant") { RefreshResult::Reauthenticate } else { RefreshResult::Unavailable }; } let Ok(mut tokens) = serde_json::from_slice::(&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>; type ApiStatuses = Arc>>; 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, passwords: Arc, 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::().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"); } } }