repositories / dotfiles
dotfiles
bugabingas dorkfiles
owned by admin
quickshell/nuguland/vault-native/src/main.rs
Rawmod session;
use session::{RefreshResult, SavedSession, SessionTasks, read_session};
use std::{
collections::{HashMap, HashSet},
env,
fs::{self, OpenOptions},
io::{self, Write as _},
net::{IpAddr, Ipv4Addr, Ipv6Addr},
os::unix::fs::{OpenOptionsExt as _, PermissionsExt as _},
path::{Path, PathBuf},
sync::{
Arc, Mutex,
atomic::{AtomicU64, Ordering},
},
time::{Duration, SystemTime},
};
use async_trait::async_trait;
use atspi::{
AccessibilityConnection, State,
proxy::{accessible::ObjectRefExt, proxy_ext::ProxyExt},
};
use bitwarden_api_api::models::{
CipherDetailsResponseModel, MasterPasswordUnlockResponseModel, PrivateKeysResponseModel,
SyncResponseModel,
};
use bitwarden_core::{
ClientBuilder, ClientSettings, DeviceType, OrganizationId, UserId,
auth::{ClientManagedTokenHandler, ClientManagedTokens, JwtToken},
key_management::{
MasterPasswordUnlockData,
account_cryptographic_state::WrappedAccountCryptographicState,
crypto::{InitOrgCryptoRequest, InitUserCryptoMethod, InitUserCryptoRequest},
},
};
use bitwarden_crypto::{HashPurpose, MasterKey, SymmetricCryptoKey, UnsignedSharedKey};
use bitwarden_crypto_sync_handler::{CryptoSyncData, CryptoSyncHandler};
use bitwarden_pm::PasswordManagerClient;
use bitwarden_sync::{SyncHandler, SyncHandlerError, SyncRequest};
use bitwarden_vault::{Cipher, CipherRepromptType, CipherType, CipherView, FieldType};
use chrono::{Duration as ChronoDuration, Utc};
use futures_util::StreamExt as _;
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use sha2::{Digest as _, Sha256};
use zeroize::{Zeroize as _, Zeroizing};
use tokio::{
io::{AsyncBufReadExt as _, AsyncWriteExt as _, BufReader},
process::Command as TokioCommand,
sync::{RwLock, Semaphore, mpsc, oneshot, watch},
task::JoinSet,
};
#[derive(Clone, Deserialize, Serialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct Profile {
name: String,
server_url: String,
account: String,
unlock_timeout_seconds: u64,
clipboard_timeout_seconds: u64,
}
impl Default for Profile {
fn default() -> Self {
Self {
name: "personal".to_owned(),
server_url: String::new(),
account: String::new(),
unlock_timeout_seconds: 900,
clipboard_timeout_seconds: 30,
}
}
}
impl Profile {
fn configured(&self) -> bool {
!self.server_url.is_empty() && !self.account.is_empty()
}
fn validate(self) -> Result<Self, &'static str> {
let name = self.name.trim();
let account = self.account.trim();
let server = self.server_url.trim();
if name.is_empty() || account.is_empty() || account.chars().any(char::is_control) {
return Err("server url and account are required");
}
let url = if server.contains("://") {
reqwest::Url::parse(server)
} else {
reqwest::Url::parse(&format!("https://{server}"))
}
.map_err(|_| "valid https server url required")?;
if url.scheme() != "https"
|| url.host_str().is_none()
|| !url.username().is_empty()
|| url.password().is_some()
|| url.query().is_some()
|| url.fragment().is_some()
|| server.chars().any(char::is_control)
{
return Err("use an https server url without credentials, query or fragment");
}
if !(1..=86_400).contains(&self.unlock_timeout_seconds)
|| !(1..=300).contains(&self.clipboard_timeout_seconds)
{
return Err("invalid vault timeout");
}
Ok(Self {
name: name.to_owned(),
account: account.to_owned(),
server_url: url.as_str().trim_end_matches('/').to_owned(),
..self
})
}
fn storage_key(&self) -> String {
format!(
"{:x}",
Sha256::digest(format!("{}\0{}", self.server_url, self.account).as_bytes())
)
}
}
#[derive(Deserialize)]
#[serde(tag = "op", rename_all = "kebab-case")]
enum Request {
Hello,
Configure {
profile: Profile,
},
TestConnection {
profile: Profile,
},
ConfigureLock {
id: u64,
},
Context {
#[serde(rename = "appId")]
app_id: Option<String>,
},
Unlock {
password: String,
},
SecondFactor {
token: String,
provider: u8,
},
CancelUnlock,
Reprompt {
id: String,
password: String,
},
State,
Detail {
id: String,
},
Copy {
id: String,
target: String,
slot: Option<i64>,
},
Favorite {
id: String,
favorite: bool,
},
Open {
id: String,
field: String,
},
Sync,
Lock,
CommitCopy,
}
#[derive(Debug)]
struct TokenStore(RwLock<Option<String>>);
#[async_trait]
impl ClientManagedTokens for TokenStore {
async fn get_access_token(&self) -> Option<String> {
self.0.read().await.clone()
}
}
impl TokenStore {
async fn clear(&self) {
if let Some(mut token) = self.0.write().await.take() {
token.zeroize();
}
}
}
impl Drop for TokenStore {
fn drop(&mut self) {
if let Some(token) = self.0.get_mut() {
token.zeroize();
}
}
}
struct Capture {
result: Arc<Mutex<Option<CapturedSync>>>,
cache_path: PathBuf,
profile: Profile,
fallback_unlock: ProtectedUnlockData,
}
#[async_trait]
impl SyncHandler for Capture {
async fn on_sync(&self, response: &SyncResponseModel) -> Result<(), SyncHandlerError> {
let mut payload = EncryptedVaultCache {
ciphers: response.ciphers.clone().unwrap_or_default(),
organization_keys: encrypted_organization_keys(response),
unlock_generation: String::new(),
};
parse_organization_keys(&payload.organization_keys).map_err(io::Error::other)?;
let protected_unlock = protected_unlock_from_sync(response, &self.profile)
.unwrap_or_else(|| self.fallback_unlock.clone());
// Store immutable key material first; the atomic cache rename publishes its exact generation.
// Before publication, interrupted writes leave the previous pair intact; both generations remain readable.
payload.unlock_generation =
store_protected_unlock(&self.profile.storage_key(), &protected_unlock).await?;
write_private_json(&self.cache_path, &payload)?;
*self.result.lock().map_err(|_| "sync capture poisoned")? = Some(CapturedSync {
payload,
protected_unlock,
});
Ok(())
}
}
#[derive(Clone)]
struct CapturedSync {
payload: EncryptedVaultCache,
protected_unlock: ProtectedUnlockData,
}
#[derive(Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct EncryptedVaultCache {
ciphers: Vec<CipherDetailsResponseModel>,
organization_keys: Vec<EncryptedOrganizationKey>,
unlock_generation: String,
}
#[derive(Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct EncryptedOrganizationKey {
id: String,
key: String,
}
#[derive(Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct ProtectedUnlockData {
user_id: UserId,
server_url: String,
account: String,
master_password_unlock: MasterPasswordUnlockData,
account_cryptographic_state: WrappedAccountCryptographicState,
}
struct RuntimeVault {
client: PasswordManagerClient,
token: Arc<TokenStore>,
encrypted: Vec<Cipher>,
summaries: Vec<Summary>,
needs_sync: bool,
favicon_tasks: JoinSet<()>,
api_url: String,
cache_path: PathBuf,
local_password_hash: Zeroizing<String>,
protected_unlock: ProtectedUnlockData,
}
struct PendingLogin {
master_key: MasterKey,
server_password_hash: Zeroizing<String>,
local_password_hash: Zeroizing<String>,
providers: Vec<SecondFactorProvider>,
}
#[derive(Clone, Serialize)]
#[serde(rename_all = "camelCase")]
struct SecondFactorProvider {
id: u8,
label: &'static str,
code_supported: bool,
}
#[derive(Deserialize)]
struct TokenSuccess {
access_token: String,
refresh_token: Option<String>,
#[serde(rename = "PrivateKey", alias = "privateKey")]
private_key: Option<String>,
#[serde(rename = "AccountKeys", alias = "accountKeys")]
account_keys: Option<PrivateKeysResponseModel>,
#[serde(rename = "UserDecryptionOptions", alias = "userDecryptionOptions")]
user_decryption_options: Option<UserDecryptionOptions>,
}
impl TokenSuccess {
fn account_state(&self) -> Result<WrappedAccountCryptographicState, &'static str> {
if let Some(keys) = self.account_keys.as_ref() {
WrappedAccountCryptographicState::try_from(keys).map_err(|_| "invalid account keys")
} else {
let private_key = self
.private_key
.as_deref()
.ok_or("missing account keys")?
.parse()
.map_err(|_| "invalid account keys")?;
Ok(WrappedAccountCryptographicState::V1 { private_key })
}
}
fn user_id(&self) -> Result<UserId, &'static str> {
// Decode only tokens returned directly by the configured HTTPS identity endpoint.
self.access_token
.parse::<JwtToken>()
.map_err(|_| "invalid sign-in token")?
.sub
.parse()
.map_err(|_| "invalid account identity")
}
}
#[derive(Deserialize)]
struct UserDecryptionOptions {
#[serde(rename = "MasterPasswordUnlock", alias = "masterPasswordUnlock")]
master_password_unlock: Option<MasterPasswordUnlockResponseModel>,
}
#[derive(Deserialize)]
struct TokenError {
#[serde(default, rename = "TwoFactorProviders", alias = "twoFactorProviders")]
two_factor_providers: Vec<String>,
}
enum TokenResult {
Authenticated(TokenSuccess),
Challenge(Vec<SecondFactorProvider>),
AuthenticationFailure,
Unavailable,
}
enum OpenResult {
Opened(Box<RuntimeVault>, bool),
Challenge(PendingLogin),
}
struct PendingCopy {
id: u64,
value: Zeroizing<String>,
field: String,
}
struct CopyOwner {
cancel: Option<oneshot::Sender<()>>,
task: tokio::task::JoinHandle<()>,
}
#[derive(Clone, Copy)]
enum SuspendEvent {
Ready,
Sleep,
Failed,
}
struct App {
profile: Profile,
runtime: Option<RuntimeVault>,
sessions: SessionTasks,
last_sync: Option<SystemTime>,
status: &'static str,
message: &'static str,
browser_origin: String,
pending_copy: Option<PendingCopy>,
copy_owner: Option<CopyOwner>,
pending_login: Option<PendingLogin>,
reprompt_grants: HashSet<String>,
next_copy_id: u64,
suspend_ready: bool,
lock_deadline: Option<tokio::time::Instant>,
}
#[derive(Clone, Serialize)]
#[serde(rename_all = "camelCase")]
struct Summary {
id: String,
title: String,
item_type: &'static str,
subtitle: String,
favorite: bool,
search_text: String,
favicon: String,
origins: Vec<String>,
has_totp: bool,
}
impl App {
fn new(profile: Profile) -> Self {
let cache_path = cache_path(&profile);
let last_sync = if profile.configured() {
fs::metadata(&cache_path).and_then(|m| m.modified()).ok()
} else {
None
};
let status = if profile.configured() {
"locked"
} else {
"needs-setup"
};
Self {
profile,
runtime: None,
sessions: SessionTasks::default(),
last_sync,
status,
message: "",
browser_origin: String::new(),
pending_copy: None,
copy_owner: None,
pending_login: None,
reprompt_grants: HashSet::new(),
next_copy_id: 1,
suspend_ready: false,
lock_deadline: None,
}
}
fn public_state(&self) -> Value {
let cache_age_seconds = self
.last_sync
.and_then(|time| time.elapsed().ok())
.map_or(-1, |age| i64::try_from(age.as_secs()).unwrap_or(i64::MAX));
let items = self
.runtime
.as_ref()
.map_or_else(Vec::new, |runtime| runtime.summaries.clone());
json!({
"event": "state",
"status": self.status,
"message": self.message,
"browserOrigin": if self.runtime.is_some() { self.browser_origin.as_str() } else { "" },
"challenge": self.pending_login.as_ref().map(|pending| json!({
"providers": pending.providers,
})),
"profileKey": self.profile.storage_key(),
"profile": {
"name": self.profile.name,
"serverUrl": self.profile.server_url,
"account": self.profile.account,
"unlockTimeoutSeconds": self.profile.unlock_timeout_seconds,
"clipboardTimeoutSeconds": self.profile.clipboard_timeout_seconds,
},
"cache": {
"available": self.last_sync.is_some(),
"ageSeconds": cache_age_seconds,
"stale": cache_age_seconds > 86_400,
},
"items": items,
"syncNeeded": self.runtime.as_ref().is_some_and(|runtime| runtime.needs_sync),
})
}
async fn configure(&mut self, profile: Profile) {
self.lock().await;
match profile.validate() {
Ok(profile) => {
self.profile = profile;
self.last_sync = fs::metadata(cache_path(&self.profile))
.and_then(|m| m.modified())
.ok();
self.status = "locked";
emit(&json!({"event": "configured", "profile": self.profile}));
}
Err(message) => emit(&json!({"event": "configure-failed", "message": message})),
}
emit(&self.public_state());
}
async fn unlock(&mut self, password: String) {
let password = Zeroizing::new(password);
let browser_origin = std::mem::take(&mut self.browser_origin);
self.lock().await;
self.browser_origin = browser_origin;
self.status = "loading";
emit(&self.public_state());
if self.sessions.settle().await.is_err() {
self.apply_open_result(Err("saved login unavailable"));
return;
}
let result = match read_session(&self.profile).await {
Ok(Some(session)) => self.resume_login(&password, &session).await,
Ok(None) => Box::pin(self.begin_login(&password)).await,
Err(_) => Err("saved login unavailable"),
};
self.apply_open_result(result);
}
async fn begin_login(&mut self, password: &str) -> Result<OpenResult, &'static str> {
let settings = settings(&self.profile);
let login_host = PasswordManagerClient::new(Some(settings.clone()));
let Ok(prelogin) = login_host
.auth()
.login(settings)
.get_password_prelogin(self.profile.account.clone())
.await
else {
return Err("vault unavailable; sign-in requires a connection");
};
let master_key = MasterKey::derive(password, &prelogin.salt, &prelogin.kdf)
.map_err(|_| "invalid password derivation parameters")?;
let pending = PendingLogin {
server_password_hash: Zeroizing::new(
master_key
.derive_master_key_hash(password.as_bytes(), HashPurpose::ServerAuthorization)
.to_string(),
),
local_password_hash: Zeroizing::new(
master_key
.derive_master_key_hash(password.as_bytes(), HashPurpose::LocalAuthorization)
.to_string(),
),
master_key,
providers: Vec::new(),
};
match request_tokens(&self.profile, &pending, None).await {
TokenResult::Authenticated(success) => self.finish_login(success, pending).await,
TokenResult::Challenge(providers) => Ok(OpenResult::Challenge(PendingLogin {
providers,
..pending
})),
TokenResult::AuthenticationFailure => Err("authentication-failure"),
TokenResult::Unavailable => Err("vault unavailable; sign-in requires a connection"),
}
}
async fn offline_login(
&self,
password: &str,
user_id: UserId,
) -> Result<OpenResult, &'static str> {
let payload = read_private_json(&cache_path(&self.profile))
.map_err(|_| "offline unlock unavailable")?;
let protected =
read_protected_unlock(&self.profile.storage_key(), &payload.unlock_generation)
.await
.map_err(|_| "offline unlock unavailable")?;
if protected.server_url != self.profile.server_url
|| protected.account != self.profile.account
|| protected.user_id != user_id
{
return Err("offline cache belongs to another profile");
}
self.open_cache(password, protected, payload).await
}
async fn resume_login(
&self,
password: &str,
session: &SavedSession,
) -> Result<OpenResult, &'static str> {
match read_private_json(&cache_path(&self.profile)) {
Ok(_) => self.offline_login(password, session.unlock.user_id).await,
Err(error) if error.kind() == io::ErrorKind::NotFound => {
self.open_cache(
password,
session.unlock.clone(),
EncryptedVaultCache {
ciphers: Vec::new(),
organization_keys: Vec::new(),
unlock_generation: String::new(),
},
)
.await
}
Err(_) => Err("cached vault unavailable"),
}
}
async fn open_cache(
&self,
password: &str,
protected: ProtectedUnlockData,
payload: EncryptedVaultCache,
) -> Result<OpenResult, &'static str> {
let master_key = MasterKey::derive(
password,
&protected.master_password_unlock.salt,
&protected.master_password_unlock.kdf,
)
.map_err(|_| "authentication-failure")?;
let local_password_hash = Zeroizing::new(
master_key
.derive_master_key_hash(password.as_bytes(), HashPurpose::LocalAuthorization)
.to_string(),
);
let decrypted_user_key = master_key
.decrypt_user_key(
protected
.master_password_unlock
.master_key_wrapped_user_key
.clone(),
)
.map_err(|_| "authentication-failure")?;
let settings = settings(&self.profile);
let token = Arc::new(TokenStore(RwLock::new(None)));
let client = vault_client(settings.clone(), token.clone()).await;
initialize_vault_crypto(&client, &self.profile, &protected, &decrypted_user_key).await?;
initialize_organization_keys(&client, &payload.organization_keys).await?;
let (encrypted, mut failures) = convert_models(&payload.ciphers);
let (mut summaries, favicon_tasks) =
build_summaries(&client, &encrypted, &self.profile.storage_key()).await;
summaries.append(&mut failures);
Ok(OpenResult::Opened(
Box::new(RuntimeVault {
client,
token,
encrypted,
summaries,
needs_sync: true,
favicon_tasks,
api_url: settings.api_url,
cache_path: cache_path(&self.profile),
local_password_hash,
protected_unlock: protected,
}),
false,
))
}
async fn second_factor(&mut self, token: String, provider: u8) {
let token = Zeroizing::new(token);
let Some(pending) = self.pending_login.as_ref() else {
return;
};
if !pending
.providers
.iter()
.any(|candidate| candidate.id == provider && candidate.code_supported)
{
self.status = "second-factor";
self.message = "unsupported verification method";
emit(&self.public_state());
return;
}
self.message = "verifying";
emit(&self.public_state());
let Some(mut pending) = self.pending_login.take() else {
return;
};
let result = match request_tokens(&self.profile, &pending, Some((&token, provider))).await {
TokenResult::Authenticated(success) => self.finish_login(success, pending).await,
TokenResult::Challenge(providers) => {
pending.providers = providers;
self.pending_login = Some(pending);
self.status = "second-factor";
self.message = "verification failed";
emit(&self.public_state());
return;
}
TokenResult::AuthenticationFailure => {
self.pending_login = Some(pending);
self.status = "second-factor";
self.message = "verification failed";
emit(&self.public_state());
return;
}
TokenResult::Unavailable => Err("vault unavailable"),
};
self.apply_open_result(result);
}
async fn finish_login(
&mut self,
success: TokenSuccess,
pending: PendingLogin,
) -> Result<OpenResult, &'static str> {
let user_id = success.user_id()?;
let account_cryptographic_state = success.account_state()?;
let refresh_token = success
.refresh_token
.filter(|token| !token.is_empty())
.ok_or("server did not provide refresh credentials")?;
let unlock = success
.user_decryption_options
.and_then(|options| options.master_password_unlock)
.ok_or("unsupported unlock method")?;
let master_password_unlock =
MasterPasswordUnlockData::try_from(&unlock).map_err(|_| "invalid unlock data")?;
let decrypted_user_key = pending
.master_key
.decrypt_user_key(master_password_unlock.master_key_wrapped_user_key.clone())
.map_err(|_| "vault key decryption failed")?;
let settings = settings(&self.profile);
let token = Arc::new(TokenStore(RwLock::new(Some(success.access_token))));
let client = vault_client(settings.clone(), token.clone()).await;
let protected_unlock = ProtectedUnlockData {
user_id,
server_url: self.profile.server_url.clone(),
account: self.profile.account.clone(),
master_password_unlock,
account_cryptographic_state,
};
initialize_vault_crypto(
&client,
&self.profile,
&protected_unlock,
&decrypted_user_key,
)
.await?;
let session = SavedSession {
refresh_token,
unlock: protected_unlock.clone(),
};
self.sessions
.save_login(&self.profile, session)
.await
.map_err(|_| "could not save login in protected storage")?;
let cache_path = cache_path(&self.profile);
let captured = Arc::new(Mutex::new(None));
let sync = client.sync();
sync.register_sync_handler(Arc::new(CryptoSyncHandler::new(client.0.clone())));
sync.register_sync_handler(Arc::new(Capture {
result: captured.clone(),
cache_path: cache_path.clone(),
profile: self.profile.clone(),
fallback_unlock: protected_unlock.clone(),
}));
let synced = sync
.sync(SyncRequest {
force: true,
exclude_subdomains: None,
})
.await
.is_ok();
let (payload, protected_unlock) = if synced {
let captured = captured
.lock()
.map_err(|_| "sync failed")?
.clone()
.ok_or("sync failed")?;
(captured.payload, captured.protected_unlock)
} else {
(
read_private_json(&cache_path).map_err(|_| "sync failed")?,
protected_unlock,
)
};
self.sessions
.update_unlock(&self.profile, &protected_unlock)
.await
.map_err(|_| "could not save login in protected storage")?;
initialize_organization_keys(&client, &payload.organization_keys).await?;
let (encrypted, mut failures) = convert_models(&payload.ciphers);
let (mut summaries, favicon_tasks) =
build_summaries(&client, &encrypted, &self.profile.storage_key()).await;
summaries.append(&mut failures);
Ok(OpenResult::Opened(
Box::new(RuntimeVault {
client,
token,
encrypted,
summaries,
needs_sync: false,
favicon_tasks,
api_url: settings.api_url,
cache_path,
local_password_hash: pending.local_password_hash,
protected_unlock,
}),
synced,
))
}
fn apply_open_result(&mut self, result: Result<OpenResult, &'static str>) {
match result {
Ok(OpenResult::Opened(runtime, synced)) => {
let local = runtime.needs_sync;
self.runtime = Some(*runtime);
self.pending_login = None;
self.lock_deadline = Some(
tokio::time::Instant::now()
+ Duration::from_secs(self.profile.unlock_timeout_seconds),
);
if synced {
self.last_sync = Some(SystemTime::now());
self.status = "unlocked";
self.message = "";
} else {
self.status = if self
.last_sync
.and_then(|time| time.elapsed().ok())
.is_some_and(|age| age > Duration::from_hours(24))
{
"stale-cache"
} else if local {
if self.last_sync.is_some() {
"unlocked"
} else {
"loading"
}
} else {
"sync-failure"
};
self.message = if local {
"restored login; syncing vault"
} else {
"sync failed; encrypted cache is in use"
};
}
}
Ok(OpenResult::Challenge(pending)) => {
self.pending_login = Some(pending);
self.lock_deadline = Some(
tokio::time::Instant::now()
+ Duration::from_secs(self.profile.unlock_timeout_seconds),
);
self.status = "second-factor";
self.message = "verification required";
}
Err(code) => {
self.pending_login = None;
self.status = if code == "authentication-failure" {
"authentication-failure"
} else {
"unavailable-vault"
};
self.message = code;
}
}
emit(&self.public_state());
}
async fn stop_copy(&mut self) {
if let Some(owner) = &mut self.copy_owner {
owner.cancel.take();
let _ = (&mut owner.task).await;
}
self.copy_owner = None;
}
async fn lock(&mut self) {
if let Some(mut runtime) = self.runtime.take() {
runtime.favicon_tasks.shutdown().await;
}
self.browser_origin.clear();
let pending_id = self.pending_copy.take().map(|pending| pending.id);
self.pending_login = None;
self.reprompt_grants.clear();
self.lock_deadline = None;
self.status = if self.profile.configured() {
"locked"
} else {
"needs-setup"
};
self.message = "";
self.stop_copy().await;
if let Some(id) = pending_id {
emit(&json!({"event": "copy-ended", "id": id}));
}
}
async fn ensure_access(&mut self) -> bool {
let Some(runtime) = &self.runtime else {
return false;
};
let user_id = runtime.protected_unlock.user_id;
if runtime
.token
.0
.read()
.await
.as_ref()
.is_some_and(|token| session::access_valid(token, user_id))
{
return true;
}
match self.sessions.refresh(&self.profile, user_id).await {
RefreshResult::Ready(token) => {
let Some(runtime) = &self.runtime else {
return false;
};
let mut access = runtime.token.0.write().await;
if let Some(old) = access.as_mut() {
old.zeroize();
}
*access = Some(token.to_string());
true
}
result @ (RefreshResult::Reauthenticate | RefreshResult::RevocationFailed) => {
self.lock().await;
self.message = if matches!(result, RefreshResult::RevocationFailed) {
"could not persist revoked login; local storage unavailable"
} else {
"session expired; sign in again"
};
emit(&self.public_state());
false
}
RefreshResult::Unavailable => {
self.status = if self.last_sync.is_some() {
"sync-failure"
} else {
"unavailable-vault"
};
self.message = "sync unavailable; saved login retained";
emit(&self.public_state());
false
}
}
}
async fn sync(&mut self) {
for attempt in 0..2 {
if !self.ensure_access().await {
return;
}
let Some(runtime) = &mut self.runtime else {
return;
};
let captured = Arc::new(Mutex::new(None));
let sync = runtime.client.sync();
sync.register_sync_handler(Arc::new(CryptoSyncHandler::new(runtime.client.0.clone())));
sync.register_sync_handler(Arc::new(Capture {
result: captured.clone(),
cache_path: runtime.cache_path.clone(),
profile: self.profile.clone(),
fallback_unlock: runtime.protected_unlock.clone(),
}));
let result = sync
.sync(SyncRequest {
force: true,
exclude_subdomains: None,
})
.await;
if matches!(&result, Err(bitwarden_sync::SyncError::Api(bitwarden_core::ApiError::Response(response))) if response.status == reqwest::StatusCode::UNAUTHORIZED)
{
runtime.token.clear().await;
if attempt == 0 {
continue;
}
}
if result.is_ok() {
let captured = captured.lock().ok().and_then(|guard| guard.clone());
if let Some(captured) = captured
&& initialize_organization_keys(
&runtime.client,
&captured.payload.organization_keys,
)
.await
.is_ok()
{
let (encrypted, mut failures) = convert_models(&captured.payload.ciphers);
runtime.encrypted = encrypted;
runtime.protected_unlock = captured.protected_unlock;
if self
.sessions
.update_unlock(&self.profile, &runtime.protected_unlock)
.await
.is_err()
{
self.status = "sync-failure";
self.message = "could not update saved login";
emit(&self.public_state());
return;
}
runtime.favicon_tasks.shutdown().await;
(runtime.summaries, runtime.favicon_tasks) = build_summaries(
&runtime.client,
&runtime.encrypted,
&self.profile.storage_key(),
)
.await;
runtime.summaries.append(&mut failures);
self.last_sync = Some(SystemTime::now());
self.status = "unlocked";
self.message = "";
} else {
self.status = "sync-failure";
self.message = "sync response could not be opened";
}
} else {
self.status = "sync-failure";
self.message = "sync failed; cached items remain available";
}
emit(&self.public_state());
return;
}
}
async fn detail(&self, id: &str) -> Option<Value> {
let runtime = self.runtime.as_ref()?;
let cipher = find_cipher(&runtime.encrypted, id)?.clone();
let view = runtime
.client
.vault()
.ciphers()
.decrypt(cipher)
.await
.ok()?;
detail_value(
&runtime.client,
&view,
self.reprompt_grants.contains(id),
&self.profile.storage_key(),
)
}
async fn reprompt(&mut self, id: String, password: String) {
let password = Zeroizing::new(password);
let valid = {
let Some(runtime) = self.runtime.as_ref() else {
return;
};
let Some(cipher) = find_cipher(&runtime.encrypted, &id).cloned() else {
emit(&json!({"event": "reprompt-failed", "id": id}));
return;
};
let Ok(view) = runtime.client.vault().ciphers().decrypt(cipher).await else {
emit(&json!({"event": "reprompt-failed", "id": id}));
return;
};
if !view.view_password {
emit(&json!({"event": "action-failed", "action": "item access"}));
return;
}
if view.reprompt != CipherRepromptType::Password {
emit(&json!({"event": "reprompt-failed", "id": id}));
return;
}
let Ok(expected) = runtime.local_password_hash.parse() else {
emit(&json!({"event": "reprompt-failed", "id": id}));
return;
};
runtime
.client
.0
.auth()
.validate_password(password.to_string(), expected)
.await
.unwrap_or(false)
};
if valid {
self.reprompt_grants.insert(id.clone());
emit(&json!({"event": "reprompt-succeeded", "id": id}));
if let Some(detail) = self.detail(&id).await {
emit(&detail);
}
} else {
emit(&json!({"event": "reprompt-failed", "id": id}));
}
}
async fn totp_for_slot(&self, id: &str, slot: i64) -> Option<String> {
let runtime = self.runtime.as_ref()?;
let cipher = find_cipher(&runtime.encrypted, id)?.clone();
let view = runtime
.client
.vault()
.ciphers()
.decrypt(cipher)
.await
.ok()?;
if !view.view_password
|| (view.reprompt == CipherRepromptType::Password && !self.reprompt_grants.contains(id))
{
return None;
}
let key = serde_json::to_value(view)
.ok()?
.get("login")?
.get("totp")?
.as_str()?
.to_owned();
totp_code_for_slot(&runtime.client, &key, slot)
}
async fn prepare_copy(&mut self, id: &str, target: &str, slot: Option<i64>) {
self.pending_copy = None;
let Some(detail) = self.detail(id).await else {
emit(&json!({"event": "copy-failed", "field": target}));
return;
};
let value = if target == "totp-current" || target == "totp-next" {
if let Some(slot) = slot {
self.totp_for_slot(id, slot).await
} else {
None
}
} else {
detail
.get("fields")
.and_then(Value::as_array)
.and_then(|fields| {
fields.iter().find(|field| {
field.get("id").and_then(Value::as_str) == Some(target)
&& matches!(
field.get("kind").and_then(Value::as_str),
Some("text" | "url")
)
})
})
.and_then(|field| field.get("value"))
.and_then(Value::as_str)
.filter(|value| !value.is_empty())
.map(str::to_owned)
};
let Some(value) = value else {
emit(&json!({"event": "copy-failed", "field": target}));
return;
};
let copy_id = self.next_copy_id;
self.next_copy_id = self.next_copy_id.wrapping_add(1).max(1);
self.pending_copy = Some(PendingCopy {
id: copy_id,
value: Zeroizing::new(value),
field: target.to_owned(),
});
emit(&json!({"event": "copy-ready", "id": copy_id}));
}
async fn commit_copy(&mut self) {
self.stop_copy().await;
let Some(pending) = self.pending_copy.take() else {
emit(&json!({"event": "copy-failed"}));
emit(&json!({"event": "copy-ended"}));
return;
};
let Ok((owner, ready)) = start_copy(
pending.id,
pending.value,
self.profile.clipboard_timeout_seconds,
) else {
emit(&json!({"event": "copy-failed", "id": pending.id, "field": pending.field}));
emit(&json!({"event": "copy-ended", "id": pending.id}));
return;
};
self.copy_owner = Some(owner);
let event = if ready.await.unwrap_or(false) {
"copied"
} else {
"copy-failed"
};
emit(&json!({"event": event, "id": pending.id, "field": pending.field}));
}
async fn favorite(&mut self, id: &str, favorite: bool) {
for attempt in 0..2 {
if !self.ensure_access().await {
return;
}
let Some(runtime) = &mut self.runtime else {
return;
};
let Some(cipher) = find_cipher(&runtime.encrypted, id) else {
return;
};
let folder_id = cipher.folder_id.as_ref().map(ToString::to_string);
let Some(token) = runtime.token.get_access_token().await else {
return;
};
let endpoint = format!(
"{}/ciphers/{id}/partial",
runtime.api_url.trim_end_matches('/')
);
let result = reqwest::Client::new()
.put(endpoint)
.bearer_auth(token)
.json(&json!({"folderId": folder_id, "favorite": favorite}))
.send()
.await;
if result
.as_ref()
.is_ok_and(|response| response.status() == reqwest::StatusCode::UNAUTHORIZED)
{
runtime.token.clear().await;
if attempt == 0 {
continue;
}
}
if result.is_ok_and(|response| response.status().is_success()) {
if let Some(cipher) = find_cipher_mut(&mut runtime.encrypted, id) {
cipher.favorite = favorite;
}
if let Some(summary) = runtime.summaries.iter_mut().find(|item| item.id == id) {
summary.favorite = favorite;
}
self.status = "unlocked";
self.message = "";
emit(&self.public_state());
} else {
emit(&json!({"event": "action-failed", "action": "favorite"}));
}
return;
}
}
async fn open_url(&self, id: &str, field_id: &str) {
let Some(detail) = self.detail(id).await else {
return;
};
let Some(url) = detail
.get("fields")
.and_then(Value::as_array)
.and_then(|fields| {
fields.iter().find(|field| {
field.get("id").and_then(Value::as_str) == Some(field_id)
&& field.get("kind").and_then(Value::as_str) == Some("url")
})
})
.and_then(|field| field.get("value"))
.and_then(Value::as_str)
else {
return;
};
let status = TokioCommand::new("gdbus")
.args([
"call",
"--session",
"--dest",
"org.freedesktop.portal.Desktop",
"--object-path",
"/org/freedesktop/portal/desktop",
"--method",
"org.freedesktop.portal.OpenURI.OpenURI",
"",
url,
"{}",
])
.status()
.await;
emit(&json!({
"event": if status.is_ok_and(|s| s.success()) { "opened" } else { "action-failed" },
"action": "open",
}));
}
}
#[derive(Serialize)]
struct TokenRequest<'a> {
client_id: &'static str,
grant_type: &'static str,
scope: &'static str,
#[serde(rename = "deviceType")]
device_type: DeviceType,
#[serde(rename = "deviceIdentifier")]
device_identifier: &'static str,
#[serde(rename = "deviceName")]
device_name: &'static str,
username: &'a str,
password: &'a str,
#[serde(rename = "twoFactorToken", skip_serializing_if = "Option::is_none")]
two_factor_token: Option<&'a str>,
#[serde(rename = "twoFactorProvider", skip_serializing_if = "Option::is_none")]
two_factor_provider: Option<u8>,
#[serde(rename = "twoFactorRemember", skip_serializing_if = "Option::is_none")]
two_factor_remember: Option<bool>,
}
async fn request_tokens(
profile: &Profile,
pending: &PendingLogin,
second_factor: Option<(&str, u8)>,
) -> TokenResult {
let settings = settings(profile);
let request = TokenRequest {
client_id: "connector",
grant_type: "password",
scope: "api offline_access",
device_type: DeviceType::LinuxDesktop,
device_identifier: "0fd397b8-c02f-4a9d-8b3e-70a46687ef89",
device_name: "nuguland vault",
username: &profile.account,
password: &pending.server_password_hash,
two_factor_token: second_factor.map(|(token, _)| token),
two_factor_provider: second_factor.map(|(_, provider)| provider),
two_factor_remember: second_factor.map(|_| false),
};
let Ok(http) = reqwest::Client::builder()
.timeout(Duration::from_secs(15))
.build()
else {
return TokenResult::Unavailable;
};
let Ok(response) = http
.post(format!(
"{}/connect/token",
settings.identity_url.trim_end_matches('/')
))
.header("Device-Identifier", request.device_identifier)
.header("Device-Type", (request.device_type as u8).to_string())
.header("Bitwarden-Client-Name", "desktop")
.header("User-Agent", settings.user_agent)
.form(&request)
.send()
.await
else {
return TokenResult::Unavailable;
};
let status = response.status();
let Some(bytes) = response_bytes_limited(response, 1_048_576).await else {
return TokenResult::Unavailable;
};
classify_token_response(status, &bytes)
}
fn classify_token_response(status: reqwest::StatusCode, bytes: &[u8]) -> TokenResult {
if status.is_success() {
return serde_json::from_slice(bytes)
.map_or(TokenResult::Unavailable, TokenResult::Authenticated);
}
let Ok(error) = serde_json::from_slice::<TokenError>(bytes) else {
return if status.is_server_error() {
TokenResult::Unavailable
} else {
TokenResult::AuthenticationFailure
};
};
let providers: Vec<_> = error
.two_factor_providers
.iter()
.filter_map(|provider| provider.parse().ok())
.filter_map(second_factor_provider)
.collect();
if providers.is_empty() {
TokenResult::AuthenticationFailure
} else {
TokenResult::Challenge(providers)
}
}
async fn response_bytes_limited(response: reqwest::Response, limit: usize) -> Option<Vec<u8>> {
if response
.content_length()
.is_some_and(|length| length > limit as u64)
{
return None;
}
let mut bytes = Vec::new();
let mut stream = response.bytes_stream();
while let Some(chunk) = stream.next().await {
let chunk = chunk.ok()?;
if bytes.len().saturating_add(chunk.len()) > limit {
return None;
}
bytes.extend_from_slice(&chunk);
}
Some(bytes)
}
fn second_factor_provider(id: u8) -> Option<SecondFactorProvider> {
let (label, code_supported) = match id {
0 => ("authenticator app", true),
1 => ("email code", true),
2 => ("duo", false),
3 => ("yubikey otp", true),
4 => ("u2f", false),
5 => ("remembered device", false),
6 => ("organization duo", false),
7 => ("webauthn", false),
_ => return None,
};
Some(SecondFactorProvider {
id,
label,
code_supported,
})
}
async fn initialize_vault_crypto(
client: &PasswordManagerClient,
profile: &Profile,
protected: &ProtectedUnlockData,
user_key: &SymmetricCryptoKey,
) -> Result<(), &'static str> {
client
.crypto()
.initialize_user_crypto(InitUserCryptoRequest {
user_id: Some(protected.user_id),
kdf_params: protected.master_password_unlock.kdf.clone(),
email: profile.account.clone(),
account_cryptographic_state: protected.account_cryptographic_state.clone(),
method: InitUserCryptoMethod::DecryptedKey {
decrypted_user_key: user_key.to_base64().to_string(),
},
upgrade_token: None,
})
.await
.map_err(|_| "vault key initialization failed")
}
async fn test_connection(settings: ClientSettings, account: String) -> Result<(), &'static str> {
let client = PasswordManagerClient::new(Some(settings.clone()));
let login = client.auth().login(settings);
match tokio::time::timeout(Duration::from_secs(3), login.get_password_prelogin(account)).await {
Ok(Ok(_)) => Ok(()),
Ok(Err(_)) => Err("could not reach a valid bitwarden prelogin endpoint"),
Err(_) => Err("connection test timed out"),
}
}
fn settings(profile: &Profile) -> ClientSettings {
let base = profile.server_url.trim_end_matches('/');
let official = base == "https://vault.bitwarden.com";
ClientSettings {
identity_url: if official {
"https://identity.bitwarden.com".to_owned()
} else {
format!("{base}/identity")
},
api_url: if official {
"https://api.bitwarden.com".to_owned()
} else {
format!("{base}/api")
},
user_agent: "Nuguland Vault/0.1".to_owned(),
device_type: DeviceType::LinuxDesktop,
device_identifier: Some("0fd397b8-c02f-4a9d-8b3e-70a46687ef89".to_owned()),
bitwarden_client_version: None,
bitwarden_package_type: None,
}
}
fn protected_unlock_from_sync(
response: &SyncResponseModel,
profile: &Profile,
) -> Option<ProtectedUnlockData> {
let crypto = CryptoSyncData::try_from(response).ok()?;
Some(ProtectedUnlockData {
user_id: UserId::new(response.profile.as_ref()?.id?),
server_url: profile.server_url.clone(),
account: profile.account.clone(),
master_password_unlock: crypto.user_decryption?.master_password_unlock?,
account_cryptographic_state: crypto.account_cryptographic_state?,
})
}
fn encrypted_organization_keys(response: &SyncResponseModel) -> Vec<EncryptedOrganizationKey> {
let Some(profile) = response.profile.as_ref() else {
return Vec::new();
};
profile
.organizations_new
.as_ref()
.or(profile.organizations.as_ref())
.into_iter()
.flatten()
.filter_map(|organization| {
Some(EncryptedOrganizationKey {
id: organization.id?.to_string(),
key: organization.key.clone()?,
})
})
.collect()
}
fn parse_organization_keys(
keys: &[EncryptedOrganizationKey],
) -> Result<HashMap<OrganizationId, UnsignedSharedKey>, &'static str> {
keys.iter()
.map(|entry| {
let id: OrganizationId = entry.id.parse().map_err(|_| "invalid organization id")?;
let key: UnsignedSharedKey =
entry.key.parse().map_err(|_| "invalid organization key")?;
Ok((id, key))
})
.collect()
}
async fn initialize_organization_keys(
client: &PasswordManagerClient,
keys: &[EncryptedOrganizationKey],
) -> Result<(), &'static str> {
let organization_keys = parse_organization_keys(keys)?;
client
.crypto()
.initialize_org_crypto(InitOrgCryptoRequest { organization_keys })
.await
.map_err(|_| "organization keys unavailable")
}
async fn vault_client(settings: ClientSettings, token: Arc<TokenStore>) -> PasswordManagerClient {
let core = ClientBuilder::new()
.with_settings(settings)
.with_token_handler(ClientManagedTokenHandler::new(token))
.build();
core.flags()
.load(HashMap::from([(
"pm-34500-strict-cipher-decryption".to_owned(),
true,
)]))
.await;
PasswordManagerClient(core)
}
async fn store_protected_unlock(profile: &str, value: &ProtectedUnlockData) -> io::Result<String> {
let bytes = Zeroizing::new(serde_json::to_vec(value).map_err(io::Error::other)?);
let generation = format!("{:x}", Sha256::digest(&bytes));
store_secret(profile, "generation", &generation, &bytes).await?;
Ok(generation)
}
async fn store_secret(profile: &str, field: &str, key: &str, bytes: &[u8]) -> io::Result<()> {
let mut child = TokioCommand::new("secret-tool")
.args([
"store",
"--label=nuguland vault",
"application",
"nuguland-vault",
"profile",
profile,
field,
key,
])
.kill_on_drop(true)
.stdin(std::process::Stdio::piped())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()?;
let Some(mut stdin) = child.stdin.take() else {
return Err(io::Error::other("protected storage input unavailable"));
};
stdin.write_all(bytes).await?;
drop(stdin);
if child.wait().await?.success() {
Ok(())
} else {
Err(io::Error::other("protected storage write failed"))
}
}
async fn read_protected_unlock(profile: &str, generation: &str) -> io::Result<ProtectedUnlockData> {
if generation.len() != 64 || !generation.bytes().all(|byte| byte.is_ascii_hexdigit()) {
return Err(io::Error::other("invalid cache credential generation"));
}
let bytes = read_secret(profile, "generation", generation)
.await?
.ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "protected unlock not found"))?;
// secret-tool lookup appends a newline to the stored JSON.
let serialized = bytes.trim_ascii_end();
if format!("{:x}", Sha256::digest(serialized)) != generation {
return Err(io::Error::other("cache credential generation mismatch"));
}
serde_json::from_slice(serialized).map_err(io::Error::other)
}
async fn read_secret(
profile: &str,
field: &str,
key: &str,
) -> io::Result<Option<Zeroizing<Vec<u8>>>> {
let output = TokioCommand::new("secret-tool")
.args([
"lookup",
"application",
"nuguland-vault",
"profile",
profile,
field,
key,
])
.kill_on_drop(true)
.stdin(std::process::Stdio::null())
.stderr(std::process::Stdio::piped())
.output()
.await?;
let bytes = Zeroizing::new(output.stdout);
if output.status.code() == Some(1) && bytes.is_empty() && output.stderr.is_empty() {
return Ok(None);
}
if !output.status.success() || bytes.is_empty() || bytes.len() > 1_048_576 {
return Err(io::Error::other("protected storage read failed"));
}
Ok(Some(bytes))
}
fn find_cipher<'a>(ciphers: &'a [Cipher], id: &str) -> Option<&'a Cipher> {
ciphers.iter().find(|cipher| {
cipher
.id
.as_ref()
.is_some_and(|value| value.to_string() == id)
})
}
fn find_cipher_mut<'a>(ciphers: &'a mut [Cipher], id: &str) -> Option<&'a mut Cipher> {
ciphers.iter_mut().find(|cipher| {
cipher
.id
.as_ref()
.is_some_and(|value| value.to_string() == id)
})
}
fn unavailable_summary(id: String, item_type: &'static str) -> Summary {
Summary {
id,
title: "unavailable item".to_owned(),
item_type,
subtitle: "decryption failed".to_owned(),
favorite: false,
search_text: format!("unavailable item {item_type}"),
favicon: String::new(),
origins: Vec::new(),
has_totp: false,
}
}
fn convert_models(models: &[CipherDetailsResponseModel]) -> (Vec<Cipher>, Vec<Summary>) {
let mut encrypted = Vec::with_capacity(models.len());
let mut failures = Vec::new();
for (index, model) in models.iter().enumerate() {
match Cipher::try_from(model.clone()) {
Ok(cipher) => encrypted.push(cipher),
Err(_) => failures.push(unavailable_summary(
model
.id
.as_ref()
.map_or_else(|| format!("unavailable-{index}"), ToString::to_string),
"unsupported",
)),
}
}
(encrypted, failures)
}
async fn build_summaries(
client: &PasswordManagerClient,
ciphers: &[Cipher],
profile_key: &str,
) -> (Vec<Summary>, JoinSet<()>) {
let mut favicon_tasks = JoinSet::new();
let mut summaries = Vec::with_capacity(ciphers.len());
let mut favicon_jobs = Vec::new();
for cipher in ciphers {
if let Ok(view) = client.vault().ciphers().decrypt(cipher.clone()).await
&& let Some(summary) = summary_value(&view, profile_key)
{
if summary.favicon.is_empty()
&& let Some(uri) = login_uris(&view).into_iter().next()
{
favicon_jobs.push((summary.id.clone(), uri));
}
summaries.push(summary);
} else {
summaries.push(unavailable_summary(
cipher
.id
.as_ref()
.map_or_else(|| "unavailable".to_owned(), ToString::to_string),
item_type(cipher.r#type),
));
}
}
let permits = Arc::new(Semaphore::new(8));
for (id, uri) in favicon_jobs {
let permits = permits.clone();
let profile_key = profile_key.to_owned();
favicon_tasks.spawn(async move {
let Ok(_permit) = permits.acquire_owned().await else {
return;
};
if let Some(path) = download_favicon(&uri, &profile_key).await {
emit(
&json!({"event": "favicon", "profileKey": profile_key, "id": id, "path": path}),
);
}
});
}
(summaries, favicon_tasks)
}
fn summary_value(view: &CipherView, profile_key: &str) -> Option<Summary> {
let data = serde_json::to_value(view).unwrap_or_default();
let item_type = item_type(view.r#type);
let payload = data
.get(type_key(view.r#type))
.cloned()
.unwrap_or(Value::Null);
let username = payload
.get("username")
.and_then(Value::as_str)
.unwrap_or("");
let uris = login_uris(view);
let first_uri = uris.first().map_or("", String::as_str);
let subtitle = if username.is_empty() {
first_uri.to_owned()
} else {
username.to_owned()
};
let mut searchable = vec![
view.name.clone(),
item_type.to_owned(),
username.to_owned(),
first_uri.to_owned(),
];
collect_searchable(&payload, "", &mut searchable);
if let Some(fields) = &view.fields {
for field in fields {
if let Some(name) = &field.name {
searchable.push(name.clone());
}
if field.r#type == FieldType::Text
&& let Some(value) = &field.value
{
searchable.push(value.clone());
}
}
}
let favicon = favicon_path(first_uri, profile_key)
.filter(|path| path.exists())
.map_or_else(String::new, |path| path.to_string_lossy().into_owned());
Some(Summary {
id: view.id.as_ref()?.to_string(),
title: view.name.clone(),
item_type,
subtitle,
favorite: view.favorite,
search_text: searchable.join(" ").to_lowercase(),
favicon,
origins: uris
.iter()
.filter_map(|uri| normalized_origin(uri))
.collect(),
has_totp: payload.get("totp").and_then(Value::as_str).is_some(),
})
}
fn detail_value(
client: &PasswordManagerClient,
view: &CipherView,
reprompt_granted: bool,
profile_key: &str,
) -> Option<Value> {
let summary = summary_value(view, profile_key)?;
let data = serde_json::to_value(view).unwrap_or_default();
let payload = data
.get(type_key(view.r#type))
.cloned()
.unwrap_or(Value::Null);
let mut fields = Vec::new();
flatten_fields(&payload, type_key(view.r#type), &mut fields);
if let Some(notes) = &view.notes
&& !notes.is_empty()
{
fields.push(field_value("notes", "notes", notes, true, "text"));
}
if let Some(custom) = &view.fields {
for (index, field) in custom.iter().enumerate() {
let name = field.name.as_deref().unwrap_or("custom field");
let value = field.value.as_deref().unwrap_or("");
let (sensitive, kind) = match field.r#type {
FieldType::Text | FieldType::Boolean => (false, "text"),
FieldType::Hidden => (true, "text"),
FieldType::Linked => (true, "unsupported"),
};
fields.push(json!({
"id": format!("custom.{index}"),
"label": name,
"value": if field.r#type == FieldType::Linked { "linked field" } else { value },
"sensitive": sensitive,
"kind": kind,
"fieldType": field.r#type as u8,
"linkedId": field.linked_id.map(u32::from),
}));
}
}
append_auxiliary_fields(view, &mut fields);
let requires_reprompt =
view.view_password && view.reprompt == CipherRepromptType::Password && !reprompt_granted;
for field in &mut fields {
let Some(object) = field.as_object_mut() else {
continue;
};
let sensitive = object
.get("sensitive")
.and_then(Value::as_bool)
.unwrap_or(false);
if field_access_restricted(sensitive, requires_reprompt, view.view_password) {
object.insert("value".to_owned(), Value::String(String::new()));
object.insert("kind".to_owned(), Value::String("restricted".to_owned()));
object.insert("sensitive".to_owned(), Value::Bool(true));
}
}
let totp = (!requires_reprompt && view.view_password)
.then(|| payload.get("totp").and_then(Value::as_str))
.flatten()
.and_then(|key| totp_value(client, key));
Some(json!({
"event": "detail",
"id": summary.id,
"title": summary.title,
"itemType": summary.item_type,
"subtitle": summary.subtitle,
"favorite": summary.favorite,
"favicon": summary.favicon,
"requiresReprompt": requires_reprompt,
"fields": fields,
"totp": totp,
}))
}
fn field_access_restricted(
sensitive: bool,
requires_reprompt: bool,
can_view_password: bool,
) -> bool {
sensitive && (requires_reprompt || !can_view_password)
}
fn append_auxiliary_fields(view: &CipherView, fields: &mut Vec<Value>) {
if let Some(attachments) = &view.attachments {
for (index, attachment) in attachments.iter().enumerate() {
fields.push(field_value(
&format!("attachment.{index}"),
"attachment",
attachment.file_name.as_deref().unwrap_or("attachment"),
false,
"unsupported",
));
}
}
if let Some(attachments) = &view.attachment_decryption_failures {
for (index, attachment) in attachments.iter().enumerate() {
fields.push(field_value(
&format!("attachmentFailure.{index}"),
"attachment decryption failed",
attachment.file_name.as_deref().unwrap_or("attachment"),
false,
"unsupported",
));
}
}
if let Some(history) = &view.password_history {
for (index, entry) in history.iter().enumerate() {
fields.push(field_value(
&format!("passwordHistory.{index}.password"),
&format!(
"previous password · {}",
entry.last_used_date.format("%Y-%m-%d")
),
&entry.password,
true,
"text",
));
}
}
}
fn item_type(kind: CipherType) -> &'static str {
match kind {
CipherType::Login => "login",
CipherType::SecureNote => "secure note",
CipherType::Card => "card",
CipherType::Identity => "identity",
CipherType::SshKey => "ssh key",
CipherType::BankAccount => "bank account",
CipherType::DriversLicense => "driver license",
CipherType::Passport => "passport",
}
}
fn type_key(kind: CipherType) -> &'static str {
match kind {
CipherType::Login => "login",
CipherType::SecureNote => "secureNote",
CipherType::Card => "card",
CipherType::Identity => "identity",
CipherType::SshKey => "sshKey",
CipherType::BankAccount => "bankAccount",
CipherType::DriversLicense => "driversLicense",
CipherType::Passport => "passport",
}
}
fn collect_searchable(value: &Value, key: &str, result: &mut Vec<String>) {
match value {
Value::Object(map) => {
for (name, child) in map {
if name != "totp" && name != "uriChecksum" {
collect_searchable(child, name, result);
}
}
}
Value::Array(values) => {
for child in values {
collect_searchable(child, key, result);
}
}
Value::String(text) if !text.is_empty() && is_searchable(key) => result.push(text.clone()),
_ => {}
}
}
fn flatten_fields(value: &Value, path: &str, result: &mut Vec<Value>) {
match value {
Value::Object(map) => {
for (name, child) in map {
if name == "totp" || name == "uriChecksum" || name == "match" {
continue;
}
let next = format!("{path}.{name}");
flatten_fields(child, &next, result);
}
}
Value::Array(values) => {
for (index, child) in values.iter().enumerate() {
flatten_fields(child, &format!("{path}.{index}"), result);
}
}
Value::String(text) if !text.is_empty() => {
let key = path.rsplit('.').next().unwrap_or(path);
let kind =
if key == "uri" && (text.starts_with("https://") || text.starts_with("http://")) {
"url"
} else if key == "keyValue" {
"unsupported"
} else {
"text"
};
result.push(field_value(
path,
&label(key),
text,
is_sensitive(path),
kind,
));
}
Value::Bool(boolean) => result.push(field_value(
path,
&label(path),
&boolean.to_string(),
false,
"text",
)),
Value::Number(number) => result.push(field_value(
path,
&label(path),
&number.to_string(),
is_sensitive(path),
"text",
)),
_ => {}
}
}
fn field_value(id: &str, label: &str, value: &str, sensitive: bool, kind: &str) -> Value {
json!({
"id": id,
"label": label,
"value": value,
"sensitive": sensitive,
"kind": kind,
})
}
fn is_searchable(name: &str) -> bool {
matches!(
name.to_ascii_lowercase().as_str(),
"uri"
| "username"
| "firstname"
| "middlename"
| "lastname"
| "email"
| "company"
| "brand"
| "cardholdername"
| "country"
)
}
fn is_sensitive(path: &str) -> bool {
let path = path.to_ascii_lowercase();
let leaf = path.rsplit('.').next().unwrap_or(&path);
matches!(
leaf,
"password"
| "code"
| "ssn"
| "socialsecuritynumber"
| "privatekey"
| "securitycode"
| "secret"
| "keyvalue"
| "credentialid"
| "userhandle"
| "totp"
| "pin"
| "number"
| "accountnumber"
| "routingnumber"
| "branchnumber"
| "swiftcode"
| "iban"
| "passportnumber"
| "licensenumber"
| "nationalidentificationnumber"
)
}
fn label(name: &str) -> String {
let leaf = name.rsplit('.').next().unwrap_or(name);
let mut result = String::new();
for character in leaf.chars() {
if character.is_ascii_uppercase() {
result.push(' ');
result.push(character.to_ascii_lowercase());
} else {
result.push(character);
}
}
result
}
fn totp_code_for_slot(client: &PasswordManagerClient, key: &str, slot: i64) -> Option<String> {
let now = Utc::now();
let current = client
.vault()
.totp()
.generate_totp(key.to_owned(), Some(now))
.ok()?;
let period = i64::from(current.period);
let current_slot = now.timestamp().div_euclid(period) * period;
if slot.rem_euclid(period) != 0 || (slot - current_slot).abs() > period {
return None;
}
let time = chrono::DateTime::from_timestamp(slot, 0)?;
client
.vault()
.totp()
.generate_totp(key.to_owned(), Some(time))
.ok()
.map(|totp| totp.code)
}
fn totp_value(client: &PasswordManagerClient, key: &str) -> Option<Value> {
let now = Utc::now();
let current = client
.vault()
.totp()
.generate_totp(key.to_owned(), Some(now))
.ok()?;
let period = i64::from(current.period);
let current_slot = now.timestamp().div_euclid(period) * period;
let next_slot = current_slot + period;
let next = client
.vault()
.totp()
.generate_totp(key.to_owned(), Some(now + ChronoDuration::seconds(period)))
.ok()?;
Some(json!({
"current": current.code,
"next": next.code,
"period": period,
"remaining": period - now.timestamp().rem_euclid(period),
"currentSlot": current_slot,
"nextSlot": next_slot,
}))
}
async fn browser_origin(app_id: Option<&str>) -> Option<String> {
tokio::time::timeout(
Duration::from_millis(750),
browser_origin_from_accessibility(app_id),
)
.await
.ok()
.flatten()
}
async fn browser_origin_from_accessibility(app_id: Option<&str>) -> Option<String> {
let accessibility = AccessibilityConnection::new().await.ok()?;
let connection = accessibility.connection();
let root = accessibility.root_accessible_on_registry().await.ok()?;
let applications = root.get_children().await.ok()?;
for application in applications {
let Ok(application) = application.into_accessible_proxy(connection).await else {
continue;
};
let app_matches = if let Some(id) = app_id {
application.name().await.is_ok_and(|name| {
let name = name.to_lowercase();
let id = id.to_lowercase();
name.contains(&id) || id.contains(&name)
})
} else {
false
};
if app_id.is_some() && !app_matches {
continue;
}
let Ok(frames) = application.get_children().await else {
continue;
};
for frame in frames {
let Ok(frame) = frame.into_accessible_proxy(connection).await else {
continue;
};
if !frame
.get_state()
.await
.is_ok_and(|state| state.contains(State::Active))
{
continue;
}
let mut stack = frame.get_children().await.unwrap_or_default();
let mut remaining = 4096;
while let Some(object) = stack.pop() {
if remaining == 0 {
return None;
}
remaining -= 1;
let Ok(object) = object.into_accessible_proxy(connection).await else {
continue;
};
if let Ok(proxies) = object.proxies().await
&& let Ok(document) = proxies.document().await
&& let Ok(attributes) = document.get_attributes().await
&& let Some(uri) = attributes.get("URI")
&& let Some(origin) = normalized_origin(uri)
{
return Some(origin);
}
if let Ok(mut children) = object.get_children().await {
stack.append(&mut children);
}
}
}
}
None
}
fn normalized_origin(uri: &str) -> Option<String> {
let url = reqwest::Url::parse(uri).ok()?;
matches!(url.scheme(), "http" | "https").then(|| url.origin().ascii_serialization())
}
fn login_uris(view: &CipherView) -> Vec<String> {
serde_json::to_value(view)
.ok()
.and_then(|value| value.get("login").cloned())
.and_then(|login| login.get("uris").cloned())
.and_then(|uris| uris.as_array().cloned())
.map_or_else(Vec::new, |uris| {
uris.into_iter()
.filter_map(|uri| uri.get("uri").and_then(Value::as_str).map(str::to_owned))
.collect()
})
}
async fn download_favicon(uri: &str, profile_key: &str) -> Option<String> {
let mut url = reqwest::Url::parse(uri).ok()?;
if !matches!(url.scheme(), "http" | "https") {
return None;
}
url.set_path("/favicon.ico");
url.set_query(None);
url.set_fragment(None);
let path = favicon_path(url.as_str(), profile_key)?;
if path.exists() {
return Some(path.to_string_lossy().into_owned());
}
let host = url.host_str()?.to_owned();
let port = url.port_or_known_default()?;
let addresses: Vec<_> = tokio::net::lookup_host((host.as_str(), port))
.await
.ok()?
.collect();
if addresses.is_empty() || addresses.iter().any(|address| !is_public_ip(address.ip())) {
return None;
}
let mut client = reqwest::Client::builder()
.timeout(Duration::from_secs(3))
.redirect(reqwest::redirect::Policy::none());
if host.parse::<IpAddr>().is_err() {
client = client.resolve(&host, addresses[0]);
}
let response = client.build().ok()?.get(url).send().await.ok()?;
if !response.status().is_success()
|| response
.content_length()
.is_none_or(|length| length > 1_048_576)
{
return None;
}
let bytes = response.bytes().await.ok()?;
if bytes.is_empty() || bytes.len() > 1_048_576 {
return None;
}
write_private_bytes(&path, &bytes).ok()?;
Some(path.to_string_lossy().into_owned())
}
fn is_public_ip(ip: IpAddr) -> bool {
match ip {
IpAddr::V4(ip) => is_public_ipv4(ip),
IpAddr::V6(ip) => is_public_ipv6(ip),
}
}
fn is_public_ipv4(ip: Ipv4Addr) -> bool {
let [a, b, _, _] = ip.octets();
!(a == 0
|| a == 10
|| a == 127
|| (a == 100 && (64..=127).contains(&b))
|| (a == 169 && b == 254)
|| (a == 172 && (16..=31).contains(&b))
|| (a == 192 && b == 168)
|| (a == 192 && b == 0)
|| (a == 198 && (b == 18 || b == 19))
|| ip.is_broadcast()
|| ip.is_documentation()
|| ip.is_multicast()
|| ip.is_unspecified()
|| a >= 240)
}
fn is_public_ipv6(ip: Ipv6Addr) -> bool {
if let Some(ipv4) = ip.to_ipv4_mapped() {
return is_public_ipv4(ipv4);
}
let segments = ip.segments();
(0x2000..=0x3fff).contains(&segments[0]) && !(segments[0] == 0x2001 && segments[1] == 0x0db8)
}
fn cache_root() -> PathBuf {
env::var_os("XDG_CACHE_HOME").map_or_else(
|| PathBuf::from(env::var_os("HOME").unwrap_or_default()).join(".cache"),
PathBuf::from,
)
}
fn cache_path(profile: &Profile) -> PathBuf {
cache_root()
.join("nuguland")
.join("vault")
.join(format!("{}.encrypted.json", profile.storage_key()))
}
fn favicon_path(uri: &str, profile_key: &str) -> Option<PathBuf> {
let origin = uri.split('/').take(3).collect::<Vec<_>>().join("/");
if !origin.starts_with("https://") && !origin.starts_with("http://") {
return None;
}
let hash = format!("{:x}", Sha256::digest(origin.as_bytes()));
Some(
cache_root()
.join("nuguland")
.join("vault")
.join(profile_key)
.join("favicons")
.join(hash),
)
}
fn write_private_json(path: &Path, value: &impl Serialize) -> io::Result<()> {
let bytes = serde_json::to_vec(value).map_err(io::Error::other)?;
write_private_bytes(path, &bytes)
}
fn write_private_bytes(path: &Path, bytes: &[u8]) -> io::Result<()> {
static NEXT_TEMPORARY: AtomicU64 = AtomicU64::new(0);
let Some(parent) = path.parent() else {
return Err(io::Error::other("cache path has no parent"));
};
fs::create_dir_all(parent)?;
fs::set_permissions(parent, fs::Permissions::from_mode(0o700))?;
let (temporary, mut file) = loop {
let sequence = NEXT_TEMPORARY.fetch_add(1, Ordering::Relaxed);
let temporary = parent.join(format!(
".vault-write-{}-{sequence}.tmp",
std::process::id()
));
match OpenOptions::new()
.create_new(true)
.write(true)
.mode(0o600)
.open(&temporary)
{
Ok(file) => break (temporary, file),
Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {}
Err(error) => return Err(error),
}
};
let result = (|| {
file.write_all(bytes)?;
file.sync_all()?;
fs::rename(&temporary, path)?;
fs::File::open(parent)
.and_then(|directory| directory.sync_all())
.map_err(|error| {
io::Error::new(
error.kind(),
format!("cache published; directory durability unconfirmed: {error}"),
)
})
})();
if result.is_err() {
let _ = fs::remove_file(&temporary);
}
result
}
fn read_private_json(path: &Path) -> io::Result<EncryptedVaultCache> {
let bytes = fs::read(path)?;
serde_json::from_slice(&bytes).map_err(io::Error::other)
}
fn start_copy(
copy_id: u64,
secret: Zeroizing<String>,
timeout_seconds: u64,
) -> io::Result<(CopyOwner, oneshot::Receiver<bool>)> {
let copy = TokioCommand::new("wl-copy")
.args(["--foreground", "--type", "text/plain"])
.kill_on_drop(true)
.stdin(std::process::Stdio::piped())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()?;
Ok(track_copy(
copy,
copy_id,
secret,
Duration::from_secs(timeout_seconds),
))
}
fn track_copy(
mut copy: tokio::process::Child,
copy_id: u64,
secret: Zeroizing<String>,
timeout: Duration,
) -> (CopyOwner, oneshot::Receiver<bool>) {
let (cancel, mut cancelled) = oneshot::channel();
let (ready, receiver) = oneshot::channel();
let task = tokio::spawn(async move {
let initialized = tokio::select! {
biased;
_ = &mut cancelled => false,
result = initialize_copy(&mut copy, secret) => result.is_ok(),
};
let _ = ready.send(initialized);
if initialized {
supervise_copy(copy, cancelled, timeout).await;
} else {
let _ = copy.kill().await;
let _ = copy.wait().await;
}
emit(&json!({"event": "copy-ended", "id": copy_id}));
});
(
CopyOwner {
cancel: Some(cancel),
task,
},
receiver,
)
}
async fn initialize_copy(
copy: &mut tokio::process::Child,
secret: Zeroizing<String>,
) -> io::Result<()> {
if let Some(mut stdin) = copy.stdin.take() {
stdin.write_all(secret.as_bytes()).await?;
}
drop(secret);
tokio::time::sleep(Duration::from_millis(25)).await;
if copy.try_wait()?.is_some_and(|status| !status.success()) {
return Err(io::Error::other("wl-copy failed"));
}
Ok(())
}
async fn supervise_copy(
mut copy: tokio::process::Child,
cancelled: oneshot::Receiver<()>,
timeout: Duration,
) {
tokio::select! {
_ = copy.wait() => return,
_ = cancelled => {}
() = tokio::time::sleep(timeout) => {}
}
let _ = copy.kill().await;
let _ = copy.wait().await;
}
fn emit(value: &Value) {
let mut stdout = io::stdout().lock();
if serde_json::to_writer(&mut stdout, value).is_ok() {
let _ = stdout.write_all(b"\n");
let _ = stdout.flush();
}
}
async fn watch_suspend(sender: &mpsc::UnboundedSender<SuspendEvent>) -> Option<()> {
let connection = atspi::zbus::Connection::system().await.ok()?;
let proxy = atspi::zbus::Proxy::new(
&connection,
"org.freedesktop.login1",
"/org/freedesktop/login1",
"org.freedesktop.login1.Manager",
)
.await
.ok()?;
let mut signals = proxy.receive_signal("PrepareForSleep").await.ok()?;
let _ = sender.send(SuspendEvent::Ready);
while let Some(message) = signals.next().await {
if message
.body()
.deserialize::<bool>()
.is_ok_and(|sleeping| sleeping)
{
let _ = sender.send(SuspendEvent::Sleep);
}
}
None
}
async fn handle_suspend_event(app: &mut App, event: SuspendEvent) {
match event {
SuspendEvent::Ready => {
app.suspend_ready = true;
if app.message == "suspend monitor unavailable" {
app.status = if app.profile.configured() {
"locked"
} else {
"needs-setup"
};
app.message = "";
emit(&app.public_state());
}
}
SuspendEvent::Sleep => {
app.lock().await;
emit(&app.public_state());
}
SuspendEvent::Failed => {
app.lock().await;
app.suspend_ready = false;
app.status = "unavailable-vault";
app.message = "suspend monitor unavailable";
emit(&app.public_state());
}
}
}
async fn handle_request(app: &mut App, request: Request) {
match request {
Request::Hello => {
if let Some(origin) = browser_origin(None).await {
app.browser_origin = origin;
}
emit(&app.public_state());
}
Request::Context { app_id } => {
if let Some(origin) = browser_origin(app_id.as_deref()).await {
app.browser_origin = origin;
}
emit(&app.public_state());
}
Request::State => emit(&app.public_state()),
Request::Configure { profile } => app.configure(profile).await,
Request::TestConnection { profile } => {
let result = match profile.validate() {
Ok(profile) => test_connection(settings(&profile), profile.account).await,
Err(message) => Err(message),
};
emit(&json!({
"event": "connection-result",
"success": result.is_ok(),
"message": result.err().unwrap_or("connection successful; account not verified"),
}));
}
Request::Unlock { password } => {
if !app.profile.configured() {
app.status = "needs-setup";
emit(&app.public_state());
} else if app.suspend_ready {
Box::pin(app.unlock(password)).await;
} else {
app.status = "unavailable-vault";
app.message = "suspend monitor unavailable";
emit(&app.public_state());
}
}
Request::SecondFactor { token, provider } => app.second_factor(token, provider).await,
Request::Reprompt { id, password } => app.reprompt(id, password).await,
Request::Detail { id } => {
if let Some(detail) = app.detail(&id).await {
emit(&detail);
} else {
emit(&json!({"event": "action-failed", "action": "detail"}));
}
}
Request::Copy { id, target, slot } => app.prepare_copy(&id, &target, slot).await,
Request::CommitCopy => app.commit_copy().await,
Request::Favorite { id, favorite } => app.favorite(&id, favorite).await,
Request::Open { id, field } => app.open_url(&id, &field).await,
Request::Sync => {
if let Some(runtime) = &mut app.runtime {
runtime.needs_sync = false;
app.status = "syncing";
app.message = "syncing";
emit(&app.public_state());
}
app.sync().await;
}
Request::CancelUnlock | Request::Lock | Request::ConfigureLock { .. } => {
app.lock().await;
emit(&app.public_state());
}
}
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum InputControl {
Running,
Lock,
ConfigureLock(u64),
Closed,
}
enum OperationStop {
Complete,
Lock,
ConfigureLock(u64),
Closed,
Suspend(SuspendEvent),
}
async fn read_requests(
reader: impl tokio::io::AsyncBufRead + Unpin,
requests: mpsc::Sender<Request>,
control: watch::Sender<InputControl>,
) {
let mut lines = reader.lines();
while let Ok(Some(line)) = lines.next_line().await {
let line = Zeroizing::new(line);
match serde_json::from_str::<Request>(&line) {
Ok(Request::Lock | Request::CancelUnlock) => {
if control.send(InputControl::Lock).is_err() {
return;
}
}
Ok(Request::ConfigureLock { id }) => {
if control.send(InputControl::ConfigureLock(id)).is_err() {
return;
}
}
Ok(request) => match requests.try_send(request) {
Ok(()) => {}
Err(mpsc::error::TrySendError::Full(_)) => {
emit(&json!({"event": "protocol-error", "message": "request queue full"}));
}
Err(mpsc::error::TrySendError::Closed(_)) => return,
},
Err(_) => emit(&json!({"event": "protocol-error"})),
}
}
let _ = control.send(InputControl::Closed);
}
async fn wait_deadline(deadline: Option<tokio::time::Instant>) {
match deadline {
Some(deadline) => tokio::time::sleep_until(deadline).await,
None => std::future::pending().await,
}
}
async fn interruptible(
operation: impl std::future::Future<Output = ()>,
control: &mut watch::Receiver<InputControl>,
suspend: &mut mpsc::UnboundedReceiver<SuspendEvent>,
deadline: Option<tokio::time::Instant>,
) -> (OperationStop, bool) {
tokio::pin!(operation);
let mut monitor_ready = false;
let stop = loop {
tokio::select! {
biased;
changed = control.changed() => {
break if changed.is_err() {
OperationStop::Closed
} else {
match *control.borrow_and_update() {
InputControl::Closed => OperationStop::Closed,
InputControl::ConfigureLock(id) => OperationStop::ConfigureLock(id),
InputControl::Running | InputControl::Lock => OperationStop::Lock,
}
};
}
() = wait_deadline(deadline) => break OperationStop::Lock,
Some(event) = suspend.recv() => match event {
SuspendEvent::Ready => monitor_ready = true,
event => break OperationStop::Suspend(event),
},
() = &mut operation => break OperationStop::Complete,
}
};
(stop, monitor_ready)
}
async fn serve(
app: &mut App,
reader: impl tokio::io::AsyncBufRead + Unpin + Send + 'static,
mut suspend: mpsc::UnboundedReceiver<SuspendEvent>,
) {
let (sender, mut requests) = mpsc::channel(32);
let (control_sender, mut control) = watch::channel(InputControl::Running);
let reader_task = tokio::spawn(read_requests(reader, sender, control_sender));
loop {
let mut request = None;
let (stop, ready) = interruptible(
async { request = requests.recv().await },
&mut control,
&mut suspend,
app.lock_deadline,
)
.await;
if ready {
handle_suspend_event(app, SuspendEvent::Ready).await;
}
let stop = if matches!(stop, OperationStop::Complete) {
let Some(request) = request else { break };
let deadline = if matches!(request, Request::Unlock { .. }) {
Some(
tokio::time::Instant::now()
+ Duration::from_secs(app.profile.unlock_timeout_seconds),
)
} else {
app.lock_deadline
};
let (stop, ready) = interruptible(
Box::pin(handle_request(app, request)),
&mut control,
&mut suspend,
deadline,
)
.await;
if ready {
handle_suspend_event(app, SuspendEvent::Ready).await;
}
stop
} else {
stop
};
let configure_id = match &stop {
OperationStop::ConfigureLock(id) => Some(*id),
_ => None,
};
let acknowledge_lock =
matches!(&stop, OperationStop::Lock | OperationStop::ConfigureLock(_));
match stop {
OperationStop::Complete => continue,
OperationStop::Closed => break,
OperationStop::Lock | OperationStop::ConfigureLock(_) => app.lock().await,
OperationStop::Suspend(event) => handle_suspend_event(app, event).await,
}
// Requests queued before the lock must never reopen or access the vault afterward.
while requests.try_recv().is_ok() {}
if acknowledge_lock {
emit(&app.public_state());
}
if let Some(id) = configure_id {
emit(&json!({"event": "configure-ready", "id": id}));
}
}
app.lock().await;
if app.sessions.settle().await.is_err() {
emit(&json!({"event":"action-failed", "action":"save-login"}));
}
reader_task.abort();
}
#[tokio::main]
async fn main() {
let mut app = App::new(Profile::default());
emit(&app.public_state());
let (suspend_sender, suspend_receiver) = mpsc::unbounded_channel();
tokio::spawn(async move {
if watch_suspend(&suspend_sender).await.is_none() {
let _ = suspend_sender.send(SuspendEvent::Failed);
}
});
serve(
&mut app,
BufReader::new(tokio::io::stdin()),
suspend_receiver,
)
.await;
}
#[cfg(test)]
mod tests {
use std::num::NonZeroU32;
use bitwarden_crypto::Kdf;
use super::*;
#[test]
fn authenticated_token_supplies_a_valid_user_identity() {
let user_id = "00000000-0000-4000-8000-000000000001"
.parse()
.expect("fixture user id");
let token: TokenSuccess = serde_json::from_value(json!({
"access_token": "e30.eyJleHAiOjQxMDI0NDQ4MDAsInN1YiI6IjAwMDAwMDAwLTAwMDAtNDAwMC04MDAwLTAwMDAwMDAwMDAwMSIsInNjb3BlIjpbImFwaSIsIm9mZmxpbmVfYWNjZXNzIl19.synthetic-signature"
})).expect("synthetic server response");
assert_eq!(token.user_id(), Ok(user_id));
for (token, expected) in [
("invalid", "invalid sign-in token"),
(
"e30.eyJleHAiOjQxMDI0NDQ4MDAsInN1YiI6Im5vdC1hLXVzZXItaWQiLCJzY29wZSI6W119.synthetic-signature",
"invalid account identity",
),
] {
let token: TokenSuccess =
serde_json::from_value(json!({"access_token": token})).expect("response");
assert_eq!(token.user_id(), Err(expected));
}
}
fn unlock_fixture() -> (Profile, ProtectedUnlockData, SymmetricCryptoKey) {
use bitwarden_core::key_management::KeySlotIds;
use bitwarden_crypto::{
KeyStore, PublicKeyEncryptionAlgorithm, SymmetricCryptoKey, SymmetricKeyAlgorithm,
};
let user_key = SymmetricCryptoKey::make(SymmetricKeyAlgorithm::Aes256CbcHmac);
let store = KeyStore::<KeySlotIds>::default();
let account_state = {
let mut context = store.context_mut();
let key_id = context.add_local_symmetric_key(user_key.clone());
let private_id = context.make_private_key(PublicKeyEncryptionAlgorithm::RsaOaepSha1);
WrappedAccountCryptographicState::V1 {
private_key: context
.wrap_private_key(key_id, private_id)
.expect("wrap fixture private key"),
}
};
let user_id = "00000000-0000-4000-8000-000000000001"
.parse()
.expect("fixture user id");
let profile = Profile {
server_url: "https://vault.example.com".to_owned(),
account: "test@example.com".to_owned(),
..Profile::default()
};
// Exercise real crypto at the SDK's valid minimum, not production KDF cost on every unlock.
let kdf = Kdf::PBKDF2 {
iterations: NonZeroU32::new(5_000).expect("nonzero fixture iterations"),
};
let master_key = MasterKey::derive("synthetic-password", &profile.account, &kdf)
.expect("fixture master key");
let protected = ProtectedUnlockData {
user_id,
server_url: profile.server_url.clone(),
account: profile.account.clone(),
master_password_unlock: MasterPasswordUnlockData {
kdf,
master_key_wrapped_user_key: master_key
.encrypt_user_key(&user_key)
.expect("wrapped fixture user key"),
salt: profile.account.clone(),
contained_key_id: None,
},
account_cryptographic_state: account_state,
};
(profile, protected, user_key)
}
#[tokio::test]
async fn sdk_unlock_requires_authenticated_user_identity() {
let (profile, protected, user_key) = unlock_fixture();
for id in [None, Some(protected.user_id)] {
let client = PasswordManagerClient(ClientBuilder::new().build());
let result = client
.crypto()
.initialize_user_crypto(InitUserCryptoRequest {
user_id: id,
kdf_params: Kdf::default_pbkdf2(),
email: profile.account.clone(),
account_cryptographic_state: protected.account_cryptographic_state.clone(),
method: InitUserCryptoMethod::DecryptedKey {
decrypted_user_key: user_key.to_base64().to_string(),
},
upgrade_token: None,
})
.await;
assert_eq!(
result.is_ok(),
id.is_some(),
"valid key material requires a user identity"
);
}
let serialized = serde_json::to_value(&protected).expect("offline unlock serialization");
let restored: ProtectedUnlockData =
serde_json::from_value(serialized.clone()).expect("offline unlock round trip");
assert_eq!(restored.user_id, protected.user_id);
for unlock in [&protected, &restored] {
let tokens = Arc::new(TokenStore(RwLock::new(None)));
let client = PasswordManagerClient(
ClientBuilder::new()
.with_token_handler(ClientManagedTokenHandler::new(tokens))
.build(),
);
assert!(
initialize_vault_crypto(&client, &profile, unlock, &user_key)
.await
.is_ok(),
"online and offline use the same identity-aware initialization"
);
}
assert_favicon_jobs_cancelled(&profile, &protected).await;
let mut without_identity = serialized;
without_identity
.as_object_mut()
.expect("object")
.remove("userId");
assert!(serde_json::from_value::<ProtectedUnlockData>(without_identity).is_err());
let mut app = App::new(profile);
app.apply_open_result(Err("vault key initialization failed"));
assert_eq!(app.status, "unavailable-vault");
assert_eq!(app.message, "vault key initialization failed");
}
#[tokio::test]
async fn cache_publication_keeps_matching_credential_generations() {
let Ok(directory) = env::var("NUGULAND_CACHE_PAIR_TEST") else {
let directory =
env::temp_dir().join(format!("nuguland-cache-pair-{}", std::process::id()));
fs::create_dir(&directory).expect("isolated credential fixture");
let helper = directory.join("secret-tool");
fs::write(
&helper,
r#"#!/bin/sh
set -eu
operation=$1
case "$operation" in
store) shift 2 ;;
lookup) shift ;;
*) exit 2 ;;
esac
[ "$#" = 6 ]
[ "$1" = application ] && [ "$2" = nuguland-vault ]
[ "$3" = profile ]
case "$4" in *[!0-9a-f]*|'') exit 2 ;; esac
case "$5" in
generation) case "$6" in *[!0-9a-f]*|'') exit 2 ;; esac ;;
kind) [ "$6" = session ] ;;
*) exit 2 ;;
esac
directory=$(dirname "$0")
file="$directory/$4.$6"
case "$operation" in
store)
if [ -e "$directory/fail-store" ]; then exit 2; fi
if [ "$5" = kind ] && [ -e "$directory/block-session-store" ]; then
cat > "$file.pending"
printf '%s' "$$" > "$directory/session-store-started"
attempts=0
until [ -e "$directory/session-release" ]; do
[ "$attempts" -lt 300 ] || exit 2
attempts=$((attempts + 1))
sleep 0.01
done
mv -- "$file.pending" "$file"
else
cat > "$file"
fi
if [ -e "$directory/block-store" ]; then
printf '%s' "$$" > "$directory/store-started"
attempts=0
while [ -e "$directory/block-store" ]; do
[ "$attempts" -lt 300 ] || exit 2
attempts=$((attempts + 1))
sleep 0.01
done
fi
;;
lookup) [ -f "$file" ] || exit 1; cat "$file"; printf '\n' ;;
esac
"#,
)
.expect("synthetic credential helper");
fs::set_permissions(&helper, fs::Permissions::from_mode(0o700))
.expect("fixture executable");
let paths = std::iter::once(directory.clone())
.chain(env::split_paths(&env::var_os("PATH").expect("PATH")).collect::<Vec<_>>());
// A process-level deadline also covers blocked runtime threads and pipe drains.
let output = TokioCommand::new("timeout")
.args(["--kill-after=5s", "120s"])
.arg(env::current_exe().expect("test executable"))
.args([
"--exact",
"tests::cache_publication_keeps_matching_credential_generations",
"--nocapture",
])
.env("NUGULAND_CACHE_PAIR_TEST", &directory)
.env("XDG_CACHE_HOME", directory.join("cache"))
.env("XDG_STATE_HOME", directory.join("state"))
.env("PATH", env::join_paths(paths).expect("isolated PATH"))
.env(
"DBUS_SESSION_BUS_ADDRESS",
"unix:path=/nonexistent-nuguland-test-bus",
)
.kill_on_drop(true)
.output()
.await;
fs::remove_dir_all(&directory).expect("clean credential fixture");
let output = output.expect("isolated test process");
assert!(
output.status.success(),
"isolated test exited with {}\n{}\n{}",
output.status,
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
return;
};
Box::pin(assert_cache_publication(Path::new(&directory))).await;
}
async fn assert_cache_publication(directory: &Path) {
let (profile, protected, user_key) = unlock_fixture();
let profile_key = profile.storage_key();
let mut capture = Capture {
result: Arc::new(Mutex::new(None)),
cache_path: cache_path(&profile),
profile: profile.clone(),
fallback_unlock: protected.clone(),
};
let response = SyncResponseModel::default();
capture
.on_sync(&response)
.await
.expect("initial publication");
let first = read_private_json(&capture.cache_path).expect("initial cache");
let app = App::new(profile);
assert_offline_open(&app, "synthetic-password").await;
let rotated = MasterKey::derive(
"rotated-password",
&protected.master_password_unlock.salt,
&protected.master_password_unlock.kdf,
)
.expect("rotated master key");
capture
.fallback_unlock
.master_password_unlock
.master_key_wrapped_user_key = rotated
.encrypt_user_key(&user_key)
.expect("rotated wrapping");
assert_cancelled_publication(&capture, directory).await;
assert_eq!(
read_private_json(&capture.cache_path)
.expect("cache after cancellation")
.unlock_generation,
first.unlock_generation
);
let cache = capture.cache_path.clone();
capture.cache_path = cache.join("cannot-publish-below-a-file");
*capture.result.lock().expect("capture lock") = None;
assert!(
capture.on_sync(&response).await.is_err(),
"credential store succeeds before cache publication fails"
);
assert!(capture.result.lock().expect("capture lock").is_none());
assert_eq!(
read_private_json(&cache)
.expect("previous cache")
.unlock_generation,
first.unlock_generation
);
assert_offline_open(&app, "synthetic-password").await;
capture.cache_path = cache;
capture
.on_sync(&response)
.await
.expect("rotated publication");
let current = read_private_json(&capture.cache_path).expect("rotated cache");
assert_ne!(current.unlock_generation, first.unlock_generation);
assert_offline_open(&app, "rotated-password").await;
assert!(matches!(
app.offline_login("synthetic-password", protected.user_id)
.await,
Err("authentication-failure")
));
assert!(
read_protected_unlock(&profile_key, &first.unlock_generation)
.await
.is_ok(),
"readers of the old cache retain matching credentials"
);
assert!(
read_protected_unlock("other-profile", ¤t.unlock_generation)
.await
.is_err()
);
for invalid in ["", "../other", "00"] {
assert!(read_protected_unlock(&profile_key, invalid).await.is_err());
}
capture
.on_sync(&response)
.await
.expect("same generation publication");
assert_eq!(
read_private_json(&capture.cache_path)
.expect("same-generation cache")
.unlock_generation,
current.unlock_generation
);
assert_unpaired_metadata_rejected(
directory,
&profile_key,
¤t.unlock_generation,
&protected,
)
.await;
Box::pin(session::tests::verify_saved_login(directory, &protected)).await;
}
async fn assert_unpaired_metadata_rejected(
directory: &Path,
profile_key: &str,
generation: &str,
protected: &ProtectedUnlockData,
) {
let stored = directory.join(format!("{profile_key}.{generation}"));
fs::write(
stored,
serde_json::to_vec(protected).expect("mismatched fixture"),
)
.expect("replace fixture credential");
assert!(
read_protected_unlock(profile_key, generation)
.await
.is_err(),
"wrong metadata under matching keyring attributes is rejected"
);
let unpaired = directory.join("unpaired-cache.json");
fs::write(&unpaired, br#"{"ciphers":[],"organizationKeys":[]}"#)
.expect("unpaired cache fixture");
assert!(
read_private_json(&unpaired).is_err(),
"unpaired cache requires an online resync"
);
}
async fn assert_offline_open(app: &App, password: &str) {
let user_id = "00000000-0000-4000-8000-000000000001"
.parse()
.expect("fixture identity");
assert!(matches!(
app.offline_login(password, user_id).await,
Ok(OpenResult::Opened(_, false))
));
}
async fn assert_cancelled_publication(capture: &Capture, directory: &Path) {
let block = directory.join("block-store");
let started = directory.join("store-started");
fs::write(&block, b"").expect("block credential helper");
let pending = Capture {
result: capture.result.clone(),
cache_path: capture.cache_path.clone(),
profile: capture.profile.clone(),
fallback_unlock: capture.fallback_unlock.clone(),
};
let task =
tokio::spawn(async move { pending.on_sync(&SyncResponseModel::default()).await });
tokio::time::timeout(Duration::from_secs(5), async {
while fs::read_to_string(&started)
.ok()
.and_then(|value| value.parse::<u32>().ok())
.is_none()
{
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("credential write reached cancellation barrier");
let pid = fs::read_to_string(&started).expect("owned helper pid");
task.abort();
assert!(task.await.expect_err("sync was cancelled").is_cancelled());
tokio::time::timeout(Duration::from_secs(5), async {
while Path::new(&format!("/proc/{pid}")).exists() {
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("cancelled credential helper is reaped");
fs::remove_file(block).expect("unblock future writes");
fs::remove_file(started).expect("remove barrier");
}
async fn assert_favicon_jobs_cancelled(profile: &Profile, protected: &ProtectedUnlockData) {
for explicit_lock in [false, true] {
let mut app = App::new(profile.clone());
let mut favicon_tasks = JoinSet::new();
let (started, ready) = oneshot::channel();
let (guard, dropped) = oneshot::channel::<()>();
favicon_tasks.spawn(async move {
let _guard = guard;
started.send(()).expect("favicon started");
std::future::pending::<()>().await;
});
ready.await.expect("favicon reached its await");
app.runtime = Some(RuntimeVault {
client: PasswordManagerClient(ClientBuilder::new().build()),
token: Arc::new(TokenStore(RwLock::new(None))),
encrypted: Vec::new(),
summaries: Vec::new(),
needs_sync: false,
favicon_tasks,
api_url: String::new(),
cache_path: PathBuf::new(),
local_password_hash: Zeroizing::new(String::new()),
protected_unlock: protected.clone(),
});
if explicit_lock {
app.lock().await;
assert!(app.runtime.is_none());
} else {
drop(app.runtime.take());
}
assert!(
tokio::time::timeout(Duration::from_secs(2), dropped)
.await
.expect("favicon cancellation completes")
.is_err(),
"locking or dropping a runtime releases its favicon tasks"
);
}
}
#[tokio::test]
async fn connection_test_requires_a_valid_prelogin_response_and_is_bounded() {
use tokio::io::AsyncReadExt as _;
for (status, body, expected) in [
(
"200 OK",
r#"{"kdfSettings":{"kdfType":0,"iterations":600000},"salt":"test@example.com"}"#,
Ok(()),
),
(
"200 OK",
"{}",
Err("could not reach a valid bitwarden prelogin endpoint"),
),
(
"200 OK",
"not json",
Err("could not reach a valid bitwarden prelogin endpoint"),
),
(
"404 Not Found",
"{}",
Err("could not reach a valid bitwarden prelogin endpoint"),
),
("timeout", "", Err("connection test timed out")),
] {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("test listener");
let address = listener.local_addr().expect("test address");
let server = tokio::spawn(async move {
let (socket, _) = listener.accept().await.expect("prelogin connection");
let mut request = BufReader::new(socket);
let mut line = String::new();
request.read_line(&mut line).await.expect("request line");
assert_eq!(
line,
"POST /identity/accounts/prelogin/password HTTP/1.1\r\n"
);
let mut length = 0;
loop {
line.clear();
request.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("body length");
}
}
assert!(length < 1024);
let mut bytes = vec![0; length];
request.read_exact(&mut bytes).await.expect("prelogin body");
assert_eq!(
serde_json::from_slice::<Value>(&bytes).expect("json"),
json!({"email":"test@example.com"}),
"prelogin sends no password or token"
);
if status == "timeout" {
std::future::pending::<()>().await;
}
let response = format!(
"HTTP/1.1 {status}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
body.len()
);
request
.get_mut()
.write_all(response.as_bytes())
.await
.expect("response");
});
// The production request handler validates HTTPS before constructing these settings.
let settings = ClientSettings {
identity_url: format!("http://{address}/identity"),
api_url: format!("http://{address}/api"),
..settings(&Profile::default())
};
assert_eq!(
test_connection(settings, "test@example.com".to_owned()).await,
expected
);
if status == "timeout" {
server.abort();
assert!(server.await.expect_err("cancel fixture").is_cancelled());
} else {
server.await.expect("fixture served prelogin");
}
}
}
#[test]
fn private_writes_are_atomic_concurrent_and_leave_no_temporary_files() {
let directory = env::temp_dir().join(format!(
"nuguland-private-write-test-{}-{}",
std::process::id(),
SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.expect("clock")
.as_nanos()
));
fs::create_dir(&directory).expect("fixture directory");
let path = directory.join("cache");
let barrier = std::sync::Barrier::new(16);
std::thread::scope(|scope| {
let mut writers = Vec::new();
for value in 0..16_u8 {
let path = &path;
let barrier = &barrier;
writers.push(scope.spawn(move || {
barrier.wait();
write_private_bytes(path, &vec![value; 65_536]).expect("concurrent write");
let observed = fs::read(path).expect("read complete published value");
assert_eq!(observed.len(), 65_536);
assert!(observed.iter().all(|byte| *byte == observed[0]));
}));
}
for writer in writers {
writer.join().expect("writer assertions");
}
});
assert_eq!(
fs::metadata(&path)
.expect("cache metadata")
.permissions()
.mode()
& 0o777,
0o600
);
assert_eq!(
fs::metadata(&directory)
.expect("directory metadata")
.permissions()
.mode()
& 0o777,
0o700
);
assert_eq!(fs::read_dir(&directory).expect("entries").count(), 1);
// A failed rename must remove only its own temporary file.
let blocked = directory.join("blocked");
fs::create_dir(&blocked).expect("blocked destination");
assert!(write_private_bytes(&blocked, b"replacement").is_err());
assert!(blocked.is_dir());
assert_eq!(
fs::read_dir(&directory)
.expect("entries after failure")
.count(),
2
);
fs::remove_dir_all(directory).expect("clean fixture");
}
#[test]
fn profile_validation_and_storage_are_identity_scoped() {
let profile = Profile {
server_url: " https://VAULT.example.com:443/ ".to_owned(),
account: " test@example.com ".to_owned(),
..Profile::default()
}
.validate()
.expect("valid profile");
assert_eq!(profile.server_url, "https://vault.example.com");
assert_eq!(profile.account, "test@example.com");
assert!(!Profile::default().configured());
for (input, expected) in [
("vault.example.com", "https://vault.example.com"),
(" VAULT.example.com:443/ ", "https://vault.example.com"),
(
"vault.example.com:8443/vault/",
"https://vault.example.com:8443/vault",
),
] {
let normalized = Profile {
server_url: input.to_owned(),
..profile.clone()
}
.validate()
.expect("hostname defaults to https");
assert_eq!(normalized.server_url, expected);
let explicit = Profile {
server_url: expected.to_owned(),
..profile.clone()
};
assert_eq!(normalized.storage_key(), explicit.storage_key());
}
for url in [
"",
"http://vault.example.com",
"HTTP://vault.example.com",
"ftp://vault.example.com",
"https://",
"not a url",
"https://user:pass@vault.example.com",
"https://vault.example.com/?token=secret",
"https://vault.example.com/#fragment",
"https://vault.exa\nmple.com",
] {
let mut invalid = profile.clone();
url.clone_into(&mut invalid.server_url);
assert!(invalid.validate().is_err(), "invalid server: {url:?}");
}
for account in ["", " ", "test\0@example.com"] {
let mut invalid = profile.clone();
account.clone_into(&mut invalid.account);
assert!(invalid.validate().is_err());
}
for (unlock, clipboard) in [(0, 30), (86_401, 30), (900, 0), (900, 301)] {
let mut invalid = profile.clone();
invalid.unlock_timeout_seconds = unlock;
invalid.clipboard_timeout_seconds = clipboard;
assert!(invalid.validate().is_err());
}
let mut renamed = profile.clone();
"../different-name".clone_into(&mut renamed.name);
assert_eq!(cache_path(&profile), cache_path(&renamed));
for other in [
Profile {
server_url: "https://other.example.com".to_owned(),
..profile.clone()
},
Profile {
account: "other@example.com".to_owned(),
..profile.clone()
},
] {
assert_ne!(
profile.storage_key(),
other.storage_key(),
"keyring entries are identity scoped"
);
assert_ne!(cache_path(&profile), cache_path(&other));
assert_ne!(
favicon_path("https://example.com", &profile.storage_key()),
favicon_path("https://example.com", &other.storage_key())
);
}
assert_eq!(profile.storage_key().len(), 64);
assert!(profile.storage_key().bytes().all(|b| b.is_ascii_hexdigit()));
let mut data = serde_json::to_value(&profile).expect("profile json");
data["password"] = json!("synthetic");
assert!(
serde_json::from_value::<Profile>(data).is_err(),
"settings cannot carry passwords"
);
}
#[tokio::test]
async fn configuration_locks_and_clears_previous_session() {
let mut app = App::new(Profile::default());
assert_eq!(app.status, "needs-setup");
app.suspend_ready = true;
Box::pin(handle_request(
&mut app,
Request::Unlock {
password: "synthetic".to_owned(),
},
))
.await;
assert_eq!(app.status, "needs-setup");
assert!(app.pending_login.is_none());
let profile = Profile {
server_url: "https://vault.example.com".to_owned(),
account: "test@example.com".to_owned(),
..Profile::default()
};
Box::pin(handle_request(
&mut app,
Request::Configure {
profile: profile.clone(),
},
))
.await;
assert_eq!(app.status, "locked");
assert_eq!(app.profile.storage_key(), profile.storage_key());
for valid in [false, true] {
app.lock_deadline = Some(tokio::time::Instant::now() + Duration::from_mins(15));
app.reprompt_grants.insert("old-item".to_owned());
"https://old.example.com".clone_into(&mut app.browser_origin);
app.pending_copy = Some(PendingCopy {
id: 1,
value: Zeroizing::new("synthetic".to_owned()),
field: "password".to_owned(),
});
let (cancel, cancelled) = oneshot::channel();
app.copy_owner = Some(CopyOwner {
cancel: Some(cancel),
task: tokio::spawn(async move {
assert!(cancelled.await.is_err());
}),
});
let mut next = profile.clone();
(if valid { "other@example.com" } else { "" }).clone_into(&mut next.account);
Box::pin(handle_request(
&mut app,
Request::Configure {
profile: next.clone(),
},
))
.await;
assert_eq!(app.status, "locked");
assert!(app.runtime.is_none());
assert!(app.pending_login.is_none());
assert!(app.pending_copy.is_none());
assert!(app.copy_owner.is_none());
assert!(app.lock_deadline.is_none());
assert!(app.reprompt_grants.is_empty());
assert!(app.browser_origin.is_empty());
assert_eq!(app.public_state()["items"], json!([]));
assert_eq!(
&app.profile.account,
if valid {
&next.account
} else {
&profile.account
}
);
}
}
#[tokio::test]
async fn lock_and_eof_cancel_in_flight_operations() {
for command in [
InputControl::Lock,
InputControl::ConfigureLock(42),
InputControl::Closed,
] {
let (sender, mut control) = watch::channel(InputControl::Running);
let (_suspend_sender, mut suspend) = mpsc::unbounded_channel();
let (started, ready) = oneshot::channel();
let (guard, dropped) = oneshot::channel::<()>();
let operation = async move {
let _guard = guard;
started.send(()).expect("operation started");
std::future::pending::<()>().await;
};
let cancel = async {
ready.await.expect("operation reached its await");
sender.send(command).expect("control receiver alive");
};
let ((stop, _), ()) = tokio::join!(
interruptible(operation, &mut control, &mut suspend, None),
cancel,
);
assert!(matches!(
(command, stop),
(InputControl::Lock, OperationStop::Lock)
| (InputControl::Closed, OperationStop::Closed)
| (
InputControl::ConfigureLock(42),
OperationStop::ConfigureLock(42)
)
));
assert!(
dropped.await.is_err(),
"cancelled future must release its state"
);
}
}
#[tokio::test]
async fn monitor_readiness_does_not_discard_a_request() {
let (_sender, mut control) = watch::channel(InputControl::Running);
let (sender, mut suspend) = mpsc::unbounded_channel();
sender.send(SuspendEvent::Ready).expect("monitor ready");
let mut completed = false;
let (stop, ready) =
interruptible(async { completed = true }, &mut control, &mut suspend, None).await;
assert!(ready);
assert!(completed);
assert!(matches!(stop, OperationStop::Complete));
}
#[tokio::test]
async fn suspend_and_deadline_interrupt_operations() {
let (_sender, mut control) = watch::channel(InputControl::Running);
let (sender, mut suspend) = mpsc::unbounded_channel();
for event in [SuspendEvent::Sleep, SuspendEvent::Failed] {
sender.send(event).expect("suspend event");
let (stop, _) =
interruptible(std::future::pending(), &mut control, &mut suspend, None).await;
assert!(matches!(stop, OperationStop::Suspend(_)));
}
let (stop, _) = interruptible(
std::future::pending(),
&mut control,
&mut suspend,
Some(tokio::time::Instant::now()),
)
.await;
assert!(matches!(stop, OperationStop::Lock));
}
#[tokio::test]
async fn cancellation_reaps_only_the_owned_clipboard_process() {
for expire in [false, true] {
let child = TokioCommand::new("sleep")
.arg("60")
.kill_on_drop(true)
.spawn()
.expect("start owned test process");
let pid = child.id().expect("owned process id");
let (sender, receiver) = oneshot::channel();
if !expire {
sender.send(()).expect("cancel clipboard owner");
}
tokio::time::timeout(
Duration::from_secs(2),
supervise_copy(
child,
receiver,
if expire {
Duration::ZERO
} else {
Duration::from_mins(1)
},
),
)
.await
.expect("clipboard owner is reaped promptly");
assert!(!Path::new(&format!("/proc/{pid}")).exists());
}
}
#[tokio::test]
async fn cancelling_copy_setup_reaps_the_child_before_completion() {
let child = TokioCommand::new("sleep")
.arg("60")
.stdin(std::process::Stdio::piped())
.kill_on_drop(true)
.spawn()
.expect("start blocked copy fixture");
let pid = child.id().expect("copy fixture id");
let (mut owner, ready) = track_copy(
child,
42,
Zeroizing::new("x".repeat(1_048_576)),
Duration::from_mins(1),
);
tokio::task::yield_now().await;
owner.cancel.take();
tokio::time::timeout(Duration::from_secs(2), owner.task)
.await
.expect("setup cancellation completes")
.expect("copy task exits");
assert!(!ready.await.expect("setup outcome"));
assert!(!Path::new(&format!("/proc/{pid}")).exists());
}
#[tokio::test]
async fn lock_bypasses_a_full_request_queue() {
let (mut input, output) = tokio::io::duplex(4096);
let (sender, _requests) = mpsc::channel(1);
let (control_sender, mut control) = watch::channel(InputControl::Running);
let task = tokio::spawn(read_requests(
BufReader::new(output),
sender,
control_sender,
));
input
.write_all(b"{\"op\":\"state\"}\n{\"op\":\"state\"}\n{\"op\":\"cancel-unlock\"}\n")
.await
.expect("queue and cancel requests");
tokio::time::timeout(Duration::from_secs(2), control.changed())
.await
.expect("lock bypasses queue")
.expect("control alive");
assert!(*control.borrow_and_update() == InputControl::Lock);
drop(input);
task.await.expect("reader exits on EOF");
assert!(*control.borrow() == InputControl::Closed);
}
#[tokio::test]
async fn lock_clears_pending_copy_and_monitor_failure_blocks_unlock() {
let profile = Profile {
name: "test".to_owned(),
server_url: "https://example.com".to_owned(),
account: "test@example.com".to_owned(),
unlock_timeout_seconds: 900,
clipboard_timeout_seconds: 30,
};
let mut app = App::new(profile);
let (sender, receiver) = oneshot::channel();
let task = tokio::spawn(async move {
assert!(receiver.await.is_err(), "lock drops cancellation sender");
});
app.copy_owner = Some(CopyOwner {
cancel: Some(sender),
task,
});
app.pending_copy = Some(PendingCopy {
id: 1,
value: Zeroizing::new("synthetic".to_owned()),
field: "password".to_owned(),
});
app.reprompt_grants.insert("test-item".to_owned());
app.suspend_ready = true;
handle_suspend_event(&mut app, SuspendEvent::Failed).await;
assert!(app.copy_owner.is_none());
assert!(app.pending_copy.is_none());
assert!(app.reprompt_grants.is_empty());
assert!(!app.suspend_ready);
assert!(
app.public_state()["items"]
.as_array()
.expect("items")
.is_empty()
);
}
#[test]
fn cached_unlock_material_accepts_only_the_master_password() {
let kdf = Kdf::PBKDF2 {
iterations: NonZeroU32::new(600_000).expect("nonzero iterations"),
};
let correct = MasterKey::derive("test_password", "test@example.com", &kdf)
.expect("derive correct master key");
let (_, wrapped_user_key) = correct.make_user_key().expect("make wrapped user key");
let unlock = MasterPasswordUnlockData {
kdf,
master_key_wrapped_user_key: wrapped_user_key,
salt: "test@example.com".to_owned(),
contained_key_id: None,
};
assert!(
correct
.decrypt_user_key(unlock.master_key_wrapped_user_key.clone())
.is_ok()
);
let wrong =
MasterKey::derive("wrong", &unlock.salt, &unlock.kdf).expect("derive wrong master key");
assert!(
wrong
.decrypt_user_key(unlock.master_key_wrapped_user_key)
.is_err()
);
}
#[test]
fn sensitivity_and_labels_cover_secret_fields() {
assert!(is_sensitive("privateKey"));
assert!(is_sensitive("card.number"));
assert!(is_sensitive("card.code"));
assert!(is_sensitive("identity.ssn"));
assert!(is_sensitive("bankAccount.routingNumber"));
assert!(is_sensitive("login.fido2Credentials.0.keyValue"));
assert!(is_sensitive("login.fido2Credentials.0.credentialId"));
assert!(is_sensitive("login.fido2Credentials.0.userHandle"));
assert!(!is_sensitive("username"));
assert_eq!(
label("login.passwordRevisionDate"),
"password revision date"
);
}
#[test]
fn access_controls_restrict_every_sensitive_field() {
assert!(field_access_restricted(true, true, true));
assert!(field_access_restricted(true, false, false));
assert!(!field_access_restricted(true, false, true));
assert!(!field_access_restricted(false, true, false));
}
#[test]
fn protocol_accepts_challenge_and_reprompt_requests() {
assert!(matches!(
serde_json::from_str::<Request>(
r#"{"op":"second-factor","token":"123456","provider":0}"#
),
Ok(Request::SecondFactor { provider: 0, .. })
));
assert!(matches!(
serde_json::from_str::<Request>(
r#"{"op":"reprompt","id":"item","password":"candidate"}"#
),
Ok(Request::Reprompt { .. })
));
assert!(matches!(
serde_json::from_str::<Request>(r#"{"op":"cancel-unlock"}"#),
Ok(Request::CancelUnlock)
));
}
#[test]
fn identity_token_request_uses_bitwarden_form_fields() {
let request = TokenRequest {
client_id: "connector",
grant_type: "password",
scope: "api offline_access",
device_type: DeviceType::LinuxDesktop,
device_identifier: "device-id",
device_name: "nuguland vault",
username: "person@example.com",
password: "derived-hash",
two_factor_token: Some("123456"),
two_factor_provider: Some(0),
two_factor_remember: Some(false),
};
let form = serde_urlencoded::to_string(request).expect("serialize token request");
assert!(form.contains("deviceType=LinuxDesktop"));
assert!(form.contains("username=person%40example.com"));
assert!(form.contains("password=derived-hash"));
assert!(form.contains("twoFactorToken=123456"));
assert!(form.contains("twoFactorProvider=0"));
assert!(form.contains("twoFactorRemember=false"));
}
#[test]
fn identity_challenge_response_preserves_available_providers() {
let response = br#"{
"error": "invalid_grant",
"error_description": "Two factor required.",
"TwoFactorProviders": ["0", "1", "7"]
}"#;
let TokenResult::Challenge(providers) =
classify_token_response(reqwest::StatusCode::BAD_REQUEST, response)
else {
panic!("expected second-factor challenge");
};
assert_eq!(providers.len(), 3);
assert_eq!(providers[0].label, "authenticator app");
assert!(providers[0].code_supported);
assert!(!providers[2].code_supported);
}
#[test]
fn second_factor_providers_expose_only_code_capable_flows() {
assert!(second_factor_provider(0).is_some_and(|provider| provider.code_supported));
assert!(second_factor_provider(1).is_some_and(|provider| provider.code_supported));
assert!(second_factor_provider(3).is_some_and(|provider| provider.code_supported));
assert!(second_factor_provider(7).is_some_and(|provider| !provider.code_supported));
assert!(second_factor_provider(255).is_none());
}
#[test]
fn flatten_preserves_nested_fields_and_actions() {
let mut fields = Vec::new();
flatten_fields(
&json!({
"username": "oli",
"password": "secret",
"uris": [{"uri": "https://example.com", "match": 0}],
"totp": "ignored",
}),
"login",
&mut fields,
);
assert_eq!(fields.len(), 3);
assert_eq!(fields[1]["sensitive"], true);
assert_eq!(fields[2]["kind"], "url");
assert!(!fields.iter().any(|field| field["value"] == "ignored"));
}
#[test]
fn searchable_metadata_excludes_secrets() {
let mut values = Vec::new();
collect_searchable(
&json!({
"username": "oli",
"password": "secret",
"ssn": "private",
"code": "123",
"uri": "https://example.com"
}),
"",
&mut values,
);
assert!(values.iter().any(|value| value == "oli"));
assert!(!values.iter().any(|value| value == "secret"));
assert!(!values.iter().any(|value| value == "private"));
assert!(!values.iter().any(|value| value == "123"));
}
#[test]
fn favicon_destinations_reject_local_networks() {
assert!(is_public_ip("8.8.8.8".parse().expect("public address")));
assert!(is_public_ip(
"2606:4700:4700::1111".parse().expect("public address")
));
for address in [
"127.0.0.1",
"10.0.0.1",
"100.64.0.1",
"169.254.1.1",
"172.16.0.1",
"192.168.0.1",
"::1",
"fc00::1",
"fe80::1",
"2001:db8::1",
"::ffff:127.0.0.1",
] {
assert!(!is_public_ip(address.parse().expect("local address")));
}
}
#[test]
fn browser_context_keeps_only_web_origin() {
assert_eq!(
normalized_origin("https://example.com/private/path?query=secret"),
Some("https://example.com".to_owned())
);
assert_eq!(normalized_origin("file:///tmp/private"), None);
assert_eq!(normalized_origin("not a URL"), None);
}
#[test]
fn official_profile_uses_official_api_hosts() {
let profile = Profile {
name: "official".to_owned(),
server_url: "https://vault.bitwarden.com".to_owned(),
account: "person@example.com".to_owned(),
unlock_timeout_seconds: 900,
clipboard_timeout_seconds: 30,
};
let settings = settings(&profile);
assert_eq!(settings.identity_url, "https://identity.bitwarden.com");
assert_eq!(settings.api_url, "https://api.bitwarden.com");
}
#[test]
fn sdk_totp_matches_rfc_6238_vector() {
let client = PasswordManagerClient::new(None);
let time = chrono::DateTime::from_timestamp(59, 0).expect("valid RFC timestamp");
let public_rfc_secret = ["GEZDGNBV", "GY3TQOJQ", "GEZDGNBV", "GY3TQOJQ"].concat();
let uri =
"otpauth://totp/test?secret=".to_owned() + &public_rfc_secret + "&digits=8&period=30";
let result = client
.vault()
.totp()
.generate_totp(uri, Some(time))
.expect("valid RFC TOTP");
assert_eq!(result.code, "94287082");
assert_eq!(result.period, 30);
}
}