From 8ad56b938d57c6464ec749d0160a2104ef83e364 Mon Sep 17 00:00:00 2001 From: Alex-Emmet Date: Tue, 29 Sep 2026 09:23:00 +0200 Subject: [PATCH] Update OPAQUE credential profile handling --- Cargo.lock | 4 + iota-auth/Cargo.toml | 2 + iota-auth/src/password.rs | 237 ++++- .../src/users/password_credentials.rs | 20 +- iota-storage/src/util/db.rs | 143 ++- mtp-type-maps | 2 +- omikron-connector/Cargo.toml | 4 + omikron-connector/src/lib.rs | 2 +- omikron-connector/src/omikron_connection.rs | 37 +- omikron-connector/src/password/error.rs | 26 + omikron-connector/src/password/management.rs | 209 ++++ omikron-connector/src/password/mod.rs | 496 ++++++++++ omikron-connector/src/password/protected.rs | 168 ++++ .../src/password/provisioning.rs | 468 +++++++++ omikron-connector/src/password/runtime.rs | 285 ++++++ omikron-connector/src/password/wire.rs | 79 ++ .../src/password_provisioning.rs | 904 ------------------ 17 files changed, 2113 insertions(+), 973 deletions(-) create mode 100644 omikron-connector/src/password/error.rs create mode 100644 omikron-connector/src/password/management.rs create mode 100644 omikron-connector/src/password/mod.rs create mode 100644 omikron-connector/src/password/protected.rs create mode 100644 omikron-connector/src/password/provisioning.rs create mode 100644 omikron-connector/src/password/runtime.rs create mode 100644 omikron-connector/src/password/wire.rs delete mode 100644 omikron-connector/src/password_provisioning.rs diff --git a/Cargo.lock b/Cargo.lock index e78b115..2dc0b0b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2232,6 +2232,7 @@ dependencies = [ "argon2", "async-trait", "dashmap", + "generic-array", "iota-identity", "mtp", "opaque-ke", @@ -2240,6 +2241,7 @@ dependencies = [ "thiserror 2.0.20", "tokio", "uuid", + "zeroize", ] [[package]] @@ -3297,6 +3299,7 @@ dependencies = [ "iota-util", "json", "mtp", + "opaque-ke", "rand_core 0.6.4", "reqwest", "serde", @@ -3308,6 +3311,7 @@ dependencies = [ "trust-dns-resolver", "url", "uuid", + "zeroize", ] [[package]] diff --git a/iota-auth/Cargo.toml b/iota-auth/Cargo.toml index 5195580..a99188e 100644 --- a/iota-auth/Cargo.toml +++ b/iota-auth/Cargo.toml @@ -11,6 +11,8 @@ mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735d rand_core = { version = "0.6", features = ["getrandom", "std"] } opaque-ke = { version = "4.0.1", features = ["argon2"] } argon2 = "0.5" +generic-array = "0.14" +zeroize = "1" sha2 = "0.10" thiserror = "2" uuid = { version = "*", features = ["v4"] } diff --git a/iota-auth/src/password.rs b/iota-auth/src/password.rs index 2929ff9..856b517 100644 --- a/iota-auth/src/password.rs +++ b/iota-auth/src/password.rs @@ -1,4 +1,6 @@ +use generic_array::{ArrayLength, GenericArray}; use mtp::crypto::PublicKeyBundle; +use opaque_ke::ksf::Ksf; use opaque_ke::{ CipherSuite, CredentialFinalization, CredentialRequest, Identifiers, RegistrationRequest, RegistrationUpload, Ristretto255, ServerLogin, ServerLoginParameters, ServerRegistration, @@ -6,13 +8,78 @@ use opaque_ke::{ }; use rand_core::OsRng; use sha2::{Digest, Sha256}; +use zeroize::Zeroizing; -pub struct TensaminOpaque; +pub struct OpaqueProfileV1; -impl CipherSuite for TensaminOpaque { +impl OpaqueProfileV1 { + pub const ID: i64 = 1; + pub const SERVER_ID_DOMAIN: &'static [u8] = b"tensamin:iota-password\0"; + pub const LOGIN_CONTEXT_DOMAIN: &'static [u8] = b"tensamin:password-provisioning\0"; + + pub fn server_identifier(iota_id: i64) -> Vec { + let mut result = Self::SERVER_ID_DOMAIN.to_vec(); + result.extend_from_slice(&iota_id.to_be_bytes()); + result + } + + pub fn login_context( + principal: &str, + iota_id: i64, + session_id: uuid::Uuid, + contact_key: &PublicKeyBundle, + ) -> Result, PasswordAuthError> { + let mut result = Self::LOGIN_CONTEXT_DOMAIN.to_vec(); + let principal_len = + u32::try_from(principal.len()).map_err(|_| PasswordAuthError::FieldTooLarge)?; + result.extend_from_slice(&principal_len.to_be_bytes()); + result.extend_from_slice(principal.as_bytes()); + result.extend_from_slice(&iota_id.to_be_bytes()); + result.extend_from_slice(session_id.as_bytes()); + let key = contact_key + .try_as_bytes() + .map_err(|error| PasswordAuthError::ContactKey(error.to_string()))?; + result.extend_from_slice(&Sha256::digest(key)); + Ok(result) + } +} + +pub const CURRENT_OPAQUE_PROFILE: i64 = OpaqueProfileV1::ID; + +/// KSF parameters are persistent semantics of OpaqueProfileV1. Do not change deployed profile-1 credentials. +pub struct PasswordKsfV1(argon2::Argon2<'static>); + +impl Default for PasswordKsfV1 { + fn default() -> Self { + let params = argon2::Params::new(19 * 1024, 2, 1, None) + .expect("PasswordKsfV1 parameters must be valid"); + Self(argon2::Argon2::new( + argon2::Algorithm::Argon2id, + argon2::Version::V0x13, + params, + )) + } +} + +impl Ksf for PasswordKsfV1 { + fn hash>( + &self, + input: GenericArray, + ) -> Result, opaque_ke::errors::InternalError> { + let mut output = GenericArray::::default(); + self.0 + .hash_password_into(&input, &[0; argon2::RECOMMENDED_SALT_LEN], &mut output) + .map_err(|_| opaque_ke::errors::InternalError::KsfError)?; + Ok(output) + } +} + +pub type TensaminOpaque = OpaqueProfileV1; + +impl CipherSuite for OpaqueProfileV1 { type OprfCs = Ristretto255; type KeyExchange = TripleDh; - type Ksf = argon2::Argon2<'static>; + type Ksf = PasswordKsfV1; } pub type PasswordServerSetup = ServerSetup; @@ -39,6 +106,13 @@ pub fn deserialize_server_setup(bytes: &[u8]) -> Result, +) -> Result { + let bytes = Zeroizing::new(bytes); + deserialize_server_setup(&bytes) +} + pub fn registration_start( setup: &PasswordServerSetup, request: &[u8], @@ -59,10 +133,14 @@ pub fn registration_finish(upload: &[u8]) -> Result, PasswordAuthError> .to_vec()) } +pub fn validate_login_request(request: &[u8]) -> Result<(), PasswordAuthError> { + CredentialRequest::::deserialize(request) + .map(|_| ()) + .map_err(|error| PasswordAuthError::Opaque(error.to_string())) +} + pub fn server_identifier(iota_id: i64) -> Vec { - let mut result = b"tensamin:iota-password:v1\0".to_vec(); - result.extend_from_slice(&iota_id.to_be_bytes()); - result + OpaqueProfileV1::server_identifier(iota_id) } pub fn login_context( @@ -71,25 +149,16 @@ pub fn login_context( session_id: uuid::Uuid, contact_key: &PublicKeyBundle, ) -> Result, PasswordAuthError> { - let mut result = b"tensamin:password-provisioning:v1\0".to_vec(); - let principal_len = - u32::try_from(principal.len()).map_err(|_| PasswordAuthError::FieldTooLarge)?; - result.extend_from_slice(&principal_len.to_be_bytes()); - result.extend_from_slice(principal.as_bytes()); - result.extend_from_slice(&iota_id.to_be_bytes()); - result.extend_from_slice(session_id.as_bytes()); - let key = contact_key - .try_as_bytes() - .map_err(|error| PasswordAuthError::ContactKey(error.to_string()))?; - result.extend_from_slice(&Sha256::digest(key)); - Ok(result) + OpaqueProfileV1::login_context(principal, iota_id, session_id, contact_key) } pub struct LoginStartResult { pub response: Vec, - pub state: Vec, + pub state: SerializedLoginState, } +pub type SerializedLoginState = Zeroizing>; + pub fn login_start( setup: &PasswordServerSetup, record: Option<&[u8]>, @@ -121,21 +190,29 @@ pub fn login_start( .map_err(|error| PasswordAuthError::Opaque(error.to_string()))?; Ok(LoginStartResult { response: result.message.serialize().to_vec(), - state: result.state.serialize().to_vec(), + state: Zeroizing::new(result.state.serialize().to_vec()), }) } +#[derive(Debug, thiserror::Error)] +pub enum LoginFinishError { + #[error("invalid OPAQUE finalization message")] + InvalidMessage, + #[error("OPAQUE authentication failed")] + AuthenticationFailed, +} + pub fn login_finish( serialized_state: &[u8], finalization: &[u8], principal: &[u8], server_id: &[u8], context: &[u8], -) -> Result<(), PasswordAuthError> { +) -> Result<(), LoginFinishError> { let state = ServerLogin::::deserialize(serialized_state) - .map_err(|error| PasswordAuthError::Opaque(error.to_string()))?; + .map_err(|_| LoginFinishError::InvalidMessage)?; let message = CredentialFinalization::::deserialize(finalization) - .map_err(|error| PasswordAuthError::Opaque(error.to_string()))?; + .map_err(|_| LoginFinishError::InvalidMessage)?; state .finish( message, @@ -147,7 +224,7 @@ pub fn login_finish( }, }, ) - .map_err(|error| PasswordAuthError::Opaque(error.to_string()))?; + .map_err(|_| LoginFinishError::AuthenticationFailed)?; Ok(()) } @@ -161,6 +238,46 @@ mod tests { }; use uuid::Uuid; + #[test] + fn profile_one_persistent_semantics_are_stable() { + assert_eq!(OpaqueProfileV1::ID, 1); + assert_eq!( + OpaqueProfileV1::SERVER_ID_DOMAIN, + b"tensamin:iota-password\0" + ); + assert_eq!( + OpaqueProfileV1::LOGIN_CONTEXT_DOMAIN, + b"tensamin:password-provisioning\0" + ); + assert_eq!( + OpaqueProfileV1::server_identifier(11), + [ + b"tensamin:iota-password\0".as_slice(), + &11_i64.to_be_bytes() + ] + .concat() + ); + let input = GenericArray::from([7_u8; 64]); + let pinned = PasswordKsfV1::default().hash(input.clone()).unwrap(); + let previous = argon2::Argon2::default().hash(input).unwrap(); + assert_eq!(pinned, previous); + } + + #[test] + fn malformed_finalization_is_not_an_authentication_failure() { + assert!(matches!( + login_finish(b"invalid", b"invalid", b"account", b"server", b"context"), + Err(LoginFinishError::InvalidMessage) + )); + } + + #[test] + fn login_request_validation_rejects_arbitrary_bytes() { + assert!(validate_login_request(b"not a KE1").is_err()); + let request = ClientLogin::::start(&mut OsRng, b"password").unwrap(); + assert!(validate_login_request(&request.message.serialize()).is_ok()); + } + #[test] fn registration_and_login_bind_context_identifiers_and_persisted_setup() { let setup = generate_server_setup(); @@ -186,7 +303,8 @@ mod tests { .unwrap(); let record = registration_finish(&upload.message.serialize()).unwrap(); let key = Keyring::generate().public_key_bundle(); - let context = login_context("omega-key:example#7", 11, Uuid::new_v4(), &key).unwrap(); + let session = Uuid::new_v4(); + let context = login_context("omega-key:example#7", 11, session, &key).unwrap(); let client = ClientLogin::::start(&mut OsRng, b"correct horse").unwrap(); let result = login_start( @@ -218,28 +336,55 @@ mod tests { ) .is_ok() ); - let client = ClientLogin::::start(&mut OsRng, b"correct horse").unwrap(); - let mismatched = login_start( - &restored, - Some(&record), - &client.message.serialize(), - principal, - &server, - b"wrong context", - ) - .unwrap(); - assert!( - client - .state - .finish( - &mut OsRng, - b"correct horse", - CredentialResponse::::deserialize(&mismatched.response) - .unwrap(), - ClientLoginFinishParameters::new(Some(&context), identifiers, None) + for (binding, changed_context) in [ + ( + "session", + login_context("omega-key:example#7", 11, Uuid::new_v4(), &key).unwrap(), + ), + ( + "contact key", + login_context( + "omega-key:example#7", + 11, + session, + &Keyring::generate().public_key_bundle(), ) - .is_err() - ); + .unwrap(), + ), + ( + "account", + login_context("omega-key:example#8", 11, session, &key).unwrap(), + ), + ( + "Iota", + login_context("omega-key:example#7", 12, session, &key).unwrap(), + ), + ] { + let client = + ClientLogin::::start(&mut OsRng, b"correct horse").unwrap(); + let mismatched = login_start( + &restored, + Some(&record), + &client.message.serialize(), + principal, + &server, + &changed_context, + ) + .unwrap(); + assert!( + client + .state + .finish( + &mut OsRng, + b"correct horse", + CredentialResponse::::deserialize(&mismatched.response) + .unwrap(), + ClientLoginFinishParameters::new(Some(&context), identifiers, None), + ) + .is_err(), + "login accepted a different {binding}" + ); + } let wrong_server = server_identifier(12); for (login_principal, login_server, password) in [ diff --git a/iota-storage/src/users/password_credentials.rs b/iota-storage/src/users/password_credentials.rs index baa1443..5188494 100644 --- a/iota-storage/src/users/password_credentials.rs +++ b/iota-storage/src/users/password_credentials.rs @@ -5,8 +5,7 @@ use crate::{storage_error::StorageError, util::db}; #[derive(Clone, Debug, PartialEq, Eq)] pub struct PasswordCredential { pub user_id: i64, - pub protocol_version: i64, - pub credential_format_version: i64, + pub opaque_profile: i64, pub opaque_record: Vec, pub encrypted_tu_credential: Vec, pub account_public_key_sha256: Vec, @@ -17,12 +16,11 @@ pub struct PasswordCredential { pub fn get(user_id: i64) -> Result, StorageError> { db::with_db(|conn| { conn.query_row( - "SELECT user_id, protocol_version, credential_format_version, opaque_record, encrypted_tu_credential, account_public_key_sha256, created_at, updated_at FROM user_password_credentials WHERE user_id = ?1", + "SELECT user_id, opaque_profile, opaque_record, encrypted_tu_credential, account_public_key_sha256, created_at, updated_at FROM user_password_credentials WHERE user_id = ?1", params![user_id], |row| Ok(PasswordCredential { - user_id: row.get(0)?, protocol_version: row.get(1)?, credential_format_version: row.get(2)?, - opaque_record: row.get(3)?, encrypted_tu_credential: row.get(4)?, - account_public_key_sha256: row.get(5)?, created_at: row.get(6)?, updated_at: row.get(7)?, + user_id: row.get(0)?, opaque_profile: row.get(1)?, opaque_record: row.get(2)?, encrypted_tu_credential: row.get(3)?, + account_public_key_sha256: row.get(4)?, created_at: row.get(5)?, updated_at: row.get(6)?, }), ).optional().map_err(Into::into) }) @@ -53,17 +51,15 @@ pub fn any() -> Result { pub fn upsert(credential: &PasswordCredential) -> Result<(), StorageError> { db::with_immediate_transaction(|tx| { tx.execute( - "INSERT INTO user_password_credentials (user_id, protocol_version, credential_format_version, opaque_record, encrypted_tu_credential, account_public_key_sha256, created_at, updated_at) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8) + "INSERT INTO user_password_credentials (user_id, opaque_profile, opaque_record, encrypted_tu_credential, account_public_key_sha256, created_at, updated_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) ON CONFLICT(user_id) DO UPDATE SET - protocol_version = excluded.protocol_version, - credential_format_version = excluded.credential_format_version, + opaque_profile = excluded.opaque_profile, opaque_record = excluded.opaque_record, encrypted_tu_credential = excluded.encrypted_tu_credential, account_public_key_sha256 = excluded.account_public_key_sha256, updated_at = excluded.updated_at", - params![credential.user_id, credential.protocol_version, credential.credential_format_version, - credential.opaque_record, credential.encrypted_tu_credential, + params![credential.user_id, credential.opaque_profile, credential.opaque_record, credential.encrypted_tu_credential, credential.account_public_key_sha256, credential.created_at, credential.updated_at], )?; Ok(()) diff --git a/iota-storage/src/util/db.rs b/iota-storage/src/util/db.rs index a46d0ff..71a17f4 100644 --- a/iota-storage/src/util/db.rs +++ b/iota-storage/src/util/db.rs @@ -2027,8 +2027,6 @@ fn run_migrations_on_connection(conn: &Connection) -> Result<(), StorageError> { r#" CREATE TABLE user_password_credentials ( user_id INTEGER PRIMARY KEY, - protocol_version INTEGER NOT NULL, - credential_format_version INTEGER NOT NULL, opaque_record BLOB NOT NULL, encrypted_tu_credential BLOB NOT NULL, account_public_key_sha256 BLOB NOT NULL, @@ -2048,6 +2046,49 @@ fn run_migrations_on_connection(conn: &Connection) -> Result<(), StorageError> { )?; } + if current_version < 48 { + let has_version_columns: bool = conn.query_row( + "SELECT EXISTS(SELECT 1 FROM pragma_table_info('user_password_credentials') WHERE name = 'protocol_version')", + [], + |row| row.get(0), + )?; + if has_version_columns { + conn.execute_batch( + r#" + BEGIN IMMEDIATE; + CREATE TABLE user_password_credentials_unversioned ( + user_id INTEGER PRIMARY KEY, + opaque_record BLOB NOT NULL, + encrypted_tu_credential BLOB NOT NULL, + account_public_key_sha256 BLOB NOT NULL, + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL + ); + INSERT INTO user_password_credentials_unversioned + (user_id, opaque_record, encrypted_tu_credential, account_public_key_sha256, created_at, updated_at) + SELECT user_id, opaque_record, encrypted_tu_credential, account_public_key_sha256, created_at, updated_at + FROM user_password_credentials; + DROP TABLE user_password_credentials; + ALTER TABLE user_password_credentials_unversioned RENAME TO user_password_credentials; + PRAGMA user_version = 48; + COMMIT; + "#, + )?; + } else { + conn.pragma_update(None, "user_version", 48)?; + } + } + + if current_version < 49 { + conn.execute_batch( + r#" + ALTER TABLE user_password_credentials + ADD COLUMN opaque_profile INTEGER NOT NULL DEFAULT 1; + PRAGMA user_version = 49; + "#, + )?; + } + Ok(()) } @@ -2122,7 +2163,7 @@ mod tests { run_migrations_on_connection(&conn)?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; - assert_eq!(version, 47); + assert_eq!(version, 49); for column in ["height", "reply_to", "edited_count", "deleted_by_external"] { let mut statement = conn.prepare("SELECT 1 FROM pragma_table_info('messages') WHERE name = ?1")?; @@ -2141,7 +2182,7 @@ mod tests { run_migrations_on_connection(&conn)?; run_migrations_on_connection(&conn)?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; - assert_eq!(version, 47); + assert_eq!(version, 49); for table in [ "sync_heads", "sync_events", @@ -2185,6 +2226,94 @@ mod tests { conn.prepare("SELECT 1 FROM pragma_table_info('pending_relays') WHERE name = ?1")?; assert!(statement.exists([column])?); } + let password_columns: Vec = conn + .prepare("SELECT name FROM pragma_table_info('user_password_credentials')")? + .query_map([], |row| row.get(0))? + .collect::>()?; + for required in [ + "user_id", + "opaque_profile", + "opaque_record", + "encrypted_tu_credential", + "account_public_key_sha256", + "created_at", + "updated_at", + ] { + assert!( + password_columns.iter().any(|column| column == required), + "missing password column {required}" + ); + } + Ok(()) + } + + #[test] + fn removes_password_version_columns_without_losing_credentials() -> Result<(), StorageError> { + let conn = Connection::open_in_memory()?; + conn.execute_batch( + r#" + CREATE TABLE user_password_credentials ( + user_id INTEGER PRIMARY KEY, + protocol_version INTEGER NOT NULL, + credential_format_version INTEGER NOT NULL, + opaque_record BLOB NOT NULL, + encrypted_tu_credential BLOB NOT NULL, + account_public_key_sha256 BLOB NOT NULL, + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL + ); + INSERT INTO user_password_credentials VALUES (7, 1, 1, x'0102', x'0304', x'0506', 10, 20); + PRAGMA user_version = 47; + "#, + )?; + + run_migrations_on_connection(&conn)?; + run_migrations_on_connection(&conn)?; + + let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; + assert_eq!(version, 49); + let credential: (i64, Vec, Vec, Vec, i64, i64) = conn.query_row( + "SELECT opaque_profile, opaque_record, encrypted_tu_credential, account_public_key_sha256, created_at, updated_at FROM user_password_credentials WHERE user_id = 7", + [], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?, row.get(4)?, row.get(5)?)), + )?; + assert_eq!(credential, (1, vec![1, 2], vec![3, 4], vec![5, 6], 10, 20)); + let old_columns: i64 = conn.query_row( + "SELECT COUNT(*) FROM pragma_table_info('user_password_credentials') WHERE name IN ('protocol_version', 'credential_format_version')", + [], + |row| row.get(0), + )?; + assert_eq!(old_columns, 0); + Ok(()) + } + + #[test] + fn migration_adds_opaque_profile_to_existing_credentials() -> Result<(), StorageError> { + let conn = Connection::open_in_memory()?; + conn.execute_batch( + r#" + CREATE TABLE user_password_credentials ( + user_id INTEGER PRIMARY KEY, + opaque_record BLOB NOT NULL, + encrypted_tu_credential BLOB NOT NULL, + account_public_key_sha256 BLOB NOT NULL, + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL + ); + INSERT INTO user_password_credentials VALUES (7, x'0102', x'0304', x'0506', 10, 20); + PRAGMA user_version = 48; + "#, + )?; + run_migrations_on_connection(&conn)?; + run_migrations_on_connection(&conn)?; + let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; + assert_eq!(version, 49); + let credential: (i64, Vec, Vec, Vec, i64, i64) = conn.query_row( + "SELECT opaque_profile, opaque_record, encrypted_tu_credential, account_public_key_sha256, created_at, updated_at FROM user_password_credentials WHERE user_id = 7", + [], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?, row.get(4)?, row.get(5)?)), + )?; + assert_eq!(credential, (1, vec![1, 2], vec![3, 4], vec![5, 6], 10, 20)); Ok(()) } @@ -2224,7 +2353,7 @@ mod tests { run_migrations_on_connection(&conn)?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; - assert_eq!(version, 47); + assert_eq!(version, 49); for column in [ "id", "user_id", @@ -2319,7 +2448,7 @@ mod tests { )?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; assert_eq!(preserved, "remote_committed"); - assert_eq!(version, 47); + assert_eq!(version, 49); Ok(()) } @@ -2357,7 +2486,7 @@ mod tests { })?; assert_eq!(count, 0); let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; - assert_eq!(version, 47); + assert_eq!(version, 49); Ok(()) } diff --git a/mtp-type-maps b/mtp-type-maps index 4f18c7a..4f81a27 160000 --- a/mtp-type-maps +++ b/mtp-type-maps @@ -1 +1 @@ -Subproject commit 4f18c7a0d9b04d38a77fbb011c4f0b21c25bf7bf +Subproject commit 4f81a27820982b54da79213d8fcae358728e983a diff --git a/omikron-connector/Cargo.toml b/omikron-connector/Cargo.toml index e09a5de..a8936d2 100644 --- a/omikron-connector/Cargo.toml +++ b/omikron-connector/Cargo.toml @@ -34,3 +34,7 @@ base64 = "0.22.1" rand_core = { version = "0.6", features = ["getrandom", "std"] } trust-dns-resolver = { version = "0.23", features = ["tokio-runtime"] } url = "2" +zeroize = "1" + +[dev-dependencies] +opaque-ke = { version = "4.0.1", features = ["argon2"] } diff --git a/omikron-connector/src/lib.rs b/omikron-connector/src/lib.rs index 902b523..f5238d8 100644 --- a/omikron-connector/src/lib.rs +++ b/omikron-connector/src/lib.rs @@ -2,7 +2,7 @@ pub mod client; pub mod identity; pub mod omega_discovery; pub mod omikron_connection; -mod password_provisioning; +mod password; pub mod router; pub mod tauth; pub mod user_ops; diff --git a/omikron-connector/src/omikron_connection.rs b/omikron-connector/src/omikron_connection.rs index a5d156c..9dbf5be 100644 --- a/omikron-connector/src/omikron_connection.rs +++ b/omikron-connector/src/omikron_connection.rs @@ -24,7 +24,7 @@ use uuid::Uuid; use crate::client::{OmikronClient, OmikronError}; use crate::omega_discovery; -use crate::password_provisioning::PasswordAuthRuntime; +use crate::password::PasswordAuthRuntime; use iota_connection::message_common::*; use iota_connection::message_handlers; @@ -631,6 +631,12 @@ impl OmikronConnection { } } + fn initialize_password_auth(&self, identity: &Path) { + if let Err(error) = self.password_auth.initialize(identity) { + log!("Password authentication unavailable: {}", error); + } + } + async fn connect_once(self: Arc) -> Result { self.set_state(ConnectionState::Connecting).await; log_t!("omikron_connecting"); @@ -640,8 +646,8 @@ impl OmikronConnection { .await .map_err(|error| format!("Iota identity initialization failed: {error}"))?; let keyring = identity.keyring(); - self.password_auth.initialize(identity_path())?; *self.keyring.write().await = Some(keyring.clone()); + self.initialize_password_auth(identity_path()); let existing_iota_id = CONFIG.load().iota_id; @@ -698,6 +704,7 @@ impl OmikronConnection { *self.sender.write().await = Some(sender_arc.clone()); self.set_state(ConnectionState::Connected { identified: true }) .await; + self.password_auth.prune(); // Start read loop let connection = Arc::new(connection); @@ -3672,6 +3679,32 @@ mod tests { use super::*; use std::fs; + #[tokio::test] + async fn damaged_password_setup_does_not_invalidate_iota_identity() { + let directory = test_path("password-isolation"); + fs::create_dir(&directory).unwrap(); + let identity_path = directory.join("identity.keyring"); + let identity = load_or_migrate_keyring_at(&identity_path, None).unwrap(); + fs::write( + identity_path.with_file_name("password-auth.setup"), + b"corrupt setup", + ) + .unwrap(); + let connection = OmikronConnection::new( + Arc::new(DashSet::new()), + Arc::new(std::sync::Mutex::new(AppState::new())), + ); + *connection.keyring.write().await = Some(identity.keyring()); + connection.initialize_password_auth(&identity_path); + assert!(connection.keyring.read().await.is_some()); + assert!(connection.password_auth.setup().is_err()); + assert_eq!( + fs::read(identity_path.with_file_name("password-auth.setup")).unwrap(), + b"corrupt setup" + ); + fs::remove_dir_all(directory).unwrap(); + } + fn test_path(name: &str) -> PathBuf { std::env::temp_dir().join(format!( "iota-identity-{name}-{}-{}", diff --git a/omikron-connector/src/password/error.rs b/omikron-connector/src/password/error.rs new file mode 100644 index 0000000..4733873 --- /dev/null +++ b/omikron-connector/src/password/error.rs @@ -0,0 +1,26 @@ +#[derive(Debug, thiserror::Error)] +pub(super) enum PasswordFinishError { + #[error("authentication failed")] + Authentication, + #[error("invalid password protocol message")] + InvalidClientMessage, + #[error("OPAQUE session binding mismatch")] + BindingMismatch, + #[error("password credential changed during login")] + CredentialChanged, + #[error("credential unavailable")] + CredentialUnavailable, + #[error("password storage unavailable")] + Storage, + #[error("Iota keyring unavailable")] + NodeUnavailable, +} + +impl PasswordFinishError { + pub(super) fn counts_as_guess(&self) -> bool { + matches!( + self, + Self::Authentication | Self::InvalidClientMessage | Self::BindingMismatch + ) + } +} diff --git a/omikron-connector/src/password/management.rs b/omikron-connector/src/password/management.rs new file mode 100644 index 0000000..84cc120 --- /dev/null +++ b/omikron-connector/src/password/management.rs @@ -0,0 +1,209 @@ +use std::{sync::Arc, time::Instant}; + +use iota_auth::password; +use iota_storage::users::{ + password_credentials::{self, PasswordCredential}, + user_manager, +}; +use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue}; +use uuid::Uuid; + +use super::{ + MAX_CREDENTIAL_BYTES, MAX_OPAQUE_BYTES, OmikronConnection, + app_protection::{PASSWORD_RESPONSE_ENCRYPTION, PASSWORD_RESPONSE_SIGNATURE}, + wire::{bytes, key_fingerprint, now_millis, session, typed}, +}; + +impl OmikronConnection { + pub(crate) async fn handle_password_enrollment_start( + self: Arc, + frame: &CommunicationValue, + ) { + match self.enrollment_start(frame).await { + Ok(response) => { + let _ = self.send_message(&response).await; + } + Err(error) => { + iota_logger::log!("Password enrollment start failed: {error}"); + self.password_error(frame, CommunicationType::ErrorInvalidData) + .await; + } + } + } + + async fn enrollment_start( + &self, + frame: &CommunicationValue, + ) -> Result { + let opened = self + .open_password_command(frame, CommunicationType::PasswordEnrollmentStart) + .await?; + let user_id = i64::try_from(opened.signer_id()).map_err(|error| error.to_string())?; + if self.password_auth.enrollments.len() >= self.password_auth.max_pending { + return Err("too many enrollments".into()); + } + let principal = self.password_principal(user_id).await?; + let request = bytes(opened.content(), DataType::OpaqueMessage, MAX_OPAQUE_BYTES)?; + let setup = self + .password_auth + .setup() + .map_err(|error| error.to_string())?; + let response = password::registration_start(&setup, &request, principal.as_bytes()) + .map_err(|error| error.to_string())?; + let enrollment_id = Uuid::new_v4(); + self.password_auth + .enrollments + .insert(enrollment_id, (user_id, Instant::now())); + self.protected_response( + frame, + CommunicationType::PasswordEnrollmentResponse, + user_id, + DataValue::Container(vec![ + typed(DataType::Uuid, DataValue::Str(enrollment_id.to_string()))?, + typed(DataType::OpaqueMessage, DataValue::Bytes(response))?, + ]), + PASSWORD_RESPONSE_SIGNATURE, + PASSWORD_RESPONSE_ENCRYPTION, + ) + .await + } + + pub(crate) async fn handle_password_enrollment_finish( + self: Arc, + frame: &CommunicationValue, + ) { + match self.enrollment_finish(frame).await { + Ok(response) => { + let _ = self.send_message(&response).await; + } + Err(error) => { + iota_logger::log!("Password enrollment finish failed: {error}"); + self.password_error(frame, CommunicationType::ErrorInvalidData) + .await; + } + } + } + + pub(super) async fn enrollment_finish( + &self, + frame: &CommunicationValue, + ) -> Result { + let opened = self + .open_password_command(frame, CommunicationType::PasswordEnrollmentFinish) + .await?; + let user_id = i64::try_from(opened.signer_id()).map_err(|error| error.to_string())?; + let enrollment_id = session(opened.content())?; + let (_, (pending_user, started)) = self + .password_auth + .enrollments + .remove(&enrollment_id) + .ok_or("enrollment expired")?; + if pending_user != user_id || started.elapsed() >= self.password_auth.pending_ttl { + return Err("enrollment mismatch".into()); + } + let upload = bytes(opened.content(), DataType::OpaqueMessage, MAX_OPAQUE_BYTES)?; + let encrypted = bytes( + opened.content(), + DataType::EncryptedCredential, + MAX_CREDENTIAL_BYTES, + )?; + let opaque_record = + password::registration_finish(&upload).map_err(|error| error.to_string())?; + let profile = user_manager::get_user(user_id) + .map_err(|error| error.to_string())? + .ok_or("account not hosted")?; + let fingerprint = key_fingerprint(&profile.public_key)?; + let now = now_millis() as i64; + password_credentials::upsert(&PasswordCredential { + user_id, + opaque_profile: password::CURRENT_OPAQUE_PROFILE, + opaque_record, + encrypted_tu_credential: encrypted, + account_public_key_sha256: fingerprint, + created_at: now, + updated_at: now, + }) + .map_err(|error| error.to_string())?; + self.password_auth.cancel_logins_for_user(user_id); + self.protected_response( + frame, + CommunicationType::PasswordEnrollmentStatus, + user_id, + DataValue::Container(vec![typed(DataType::Enabled, DataValue::Bool(true))?]), + PASSWORD_RESPONSE_SIGNATURE, + PASSWORD_RESPONSE_ENCRYPTION, + ) + .await + } + + pub(crate) async fn handle_password_enrollment_status( + self: Arc, + frame: &CommunicationValue, + ) { + let result = async { + let opened = self + .open_password_command(frame, CommunicationType::PasswordEnrollmentStatus) + .await?; + let user_id = i64::try_from(opened.signer_id()).map_err(|error| error.to_string())?; + let enabled = password_credentials::get(user_id) + .map_err(|error| error.to_string())? + .is_some_and(|credential| { + credential.opaque_profile == password::CURRENT_OPAQUE_PROFILE + }); + self.protected_response( + frame, + CommunicationType::PasswordEnrollmentStatus, + user_id, + DataValue::Container(vec![typed(DataType::Enabled, DataValue::Bool(enabled))?]), + PASSWORD_RESPONSE_SIGNATURE, + PASSWORD_RESPONSE_ENCRYPTION, + ) + .await + } + .await; + match result { + Ok(response) => { + let _ = self.send_message(&response).await; + } + Err(error) => { + iota_logger::log!("Password status failed: {error}"); + self.password_error(frame, CommunicationType::ErrorInvalidData) + .await; + } + } + } + + pub(crate) async fn handle_password_enrollment_disable( + self: Arc, + frame: &CommunicationValue, + ) { + let result = async { + let opened = self + .open_password_command(frame, CommunicationType::PasswordEnrollmentDisable) + .await?; + let user_id = i64::try_from(opened.signer_id()).map_err(|error| error.to_string())?; + password_credentials::delete(user_id).map_err(|error| error.to_string())?; + self.password_auth.cancel_logins_for_user(user_id); + self.protected_response( + frame, + CommunicationType::PasswordEnrollmentStatus, + user_id, + DataValue::Container(vec![typed(DataType::Enabled, DataValue::Bool(false))?]), + PASSWORD_RESPONSE_SIGNATURE, + PASSWORD_RESPONSE_ENCRYPTION, + ) + .await + } + .await; + match result { + Ok(response) => { + let _ = self.send_message(&response).await; + } + Err(error) => { + iota_logger::log!("Password disable failed: {error}"); + self.password_error(frame, CommunicationType::ErrorInvalidData) + .await; + } + } + } +} diff --git a/omikron-connector/src/password/mod.rs b/omikron-connector/src/password/mod.rs new file mode 100644 index 0000000..20016b0 --- /dev/null +++ b/omikron-connector/src/password/mod.rs @@ -0,0 +1,496 @@ +use std::time::Duration; + +use iota_storage::users::user_manager; +use mtp::crypto::PublicKeyBundle; + +#[cfg(test)] +use iota_auth::password; +#[cfg(test)] +use mtp::codec::{ + CommunicationType, CommunicationValue, DataType, DataValue, ProtectedMessageBuilder, + ProtectedOpenOptions, ProtectionPolicy, open_protected_with_checked, +}; +#[cfg(test)] +use mtp::crypto::DualSigner; +#[cfg(test)] +use rand_core::OsRng; +#[cfg(test)] +use sha2::{Digest, Sha256}; +#[cfg(test)] +use std::time::Instant; +#[cfg(test)] +use uuid::Uuid; +#[cfg(test)] +use zeroize::Zeroizing; + +use crate::omikron_connection::OmikronConnection; + +#[allow(dead_code)] +mod app_protection { + include!(concat!( + env!("CARGO_MANIFEST_DIR"), + "/../mtp-type-maps/app_protection.rs" + )); +} +#[cfg(test)] +use app_protection::{PASSWORD_MANAGEMENT_ENCRYPTION, PASSWORD_MANAGEMENT_SIGNATURE, purpose}; + +const MAX_OPAQUE_BYTES: usize = 16 * 1024; +const MAX_CREDENTIAL_BYTES: usize = 256 * 1024; + +mod error; +mod management; +mod protected; +mod provisioning; +mod runtime; +use error::PasswordFinishError; + +pub(crate) use runtime::PasswordAuthRuntime; +#[cfg(test)] +use runtime::PasswordRuntimeError; +use runtime::PendingLogin; +mod wire; +#[cfg(test)] +use wire::{now_millis, typed}; + +impl OmikronConnection { + // Password login uses Omega's key-scoped principal, including for accounts + // whose older hosted-principal record still uses the legacy host locator. + async fn password_principal(&self, user_id: i64) -> Result { + if user_id <= 0 + || user_manager::get_user(user_id) + .map_err(|error| error.to_string())? + .is_none() + { + return Err("account not hosted".into()); + } + let authority = if let Some(cached) = self.password_auth.omega_authority.get() { + cached.clone() + } else { + let url = format!( + "{}/.well-known/tensamin", + crate::omega_discovery::api_base() + ); + let response = + tokio::time::timeout(Duration::from_secs(10), self.http_client.get(url).send()) + .await + .map_err(|_| "Omega identity lookup timed out")? + .map_err(|error| error.to_string())? + .error_for_status() + .map_err(|error| error.to_string())?; + let discovery: serde_json::Value = + response.json().await.map_err(|error| error.to_string())?; + let key = discovery + .get("public_key") + .and_then(serde_json::Value::as_str) + .ok_or("Omega discovery is missing its key")?; + let key = PublicKeyBundle::from_base64(key).map_err(|error| error.to_string())?; + let authority = iota_identity::AuthorityId::for_omega(&key) + .map_err(|error| error.to_string())? + .as_str() + .to_owned(); + if discovery + .get("authority_id") + .and_then(serde_json::Value::as_str) + != Some(authority.as_str()) + { + return Err("Omega discovery identity mismatch".into()); + } + let _ = self.password_auth.omega_authority.set(authority.clone()); + authority + }; + Ok(format!("{authority}#{user_id}")) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use mtp::{codec::InMemoryReplayGuard, crypto::Keyring}; + use opaque_ke::{ + ClientLogin, ClientLoginFinishParameters, ClientRegistration, + ClientRegistrationFinishParameters, CredentialResponse, Identifiers, RegistrationResponse, + }; + + #[test] + fn damaged_setup_stays_unavailable() { + let directory = std::env::temp_dir().join(format!("password-setup-{}", Uuid::new_v4())); + std::fs::create_dir(&directory).unwrap(); + let identity = directory.join("identity.keyring"); + std::fs::write(directory.join("password-auth.setup"), b"corrupt setup").unwrap(); + let runtime = PasswordAuthRuntime::default(); + assert!(runtime.initialize(&identity).is_err()); + assert!(matches!( + runtime.setup(), + Err(PasswordRuntimeError::SetupUnavailable) + )); + assert_eq!( + std::fs::read(directory.join("password-auth.setup")).unwrap(), + b"corrupt setup" + ); + std::fs::remove_dir_all(directory).unwrap(); + } + + fn limited_runtime() -> std::sync::Arc { + std::sync::Arc::new(PasswordAuthRuntime::default()) + } + + fn pending_login( + user_id: i64, + attempt: runtime::AttemptLease, + created: Instant, + ) -> PendingLogin { + PendingLogin { + user_id, + participant_id: 1, + contact_key: Keyring::generate().public_key_bundle(), + context: vec![], + state: Zeroizing::new(vec![]), + real_record: false, + record_hash: None, + created, + attempt: Some(attempt), + } + } + + #[test] + fn issued_ke2_reserves_global_capacity() { + let runtime = limited_runtime(); + let (_attempt, delay) = runtime.begin_attempt(7).unwrap(); + assert_eq!(delay, Duration::ZERO); + assert_eq!(runtime.attempts.get(&7).unwrap().in_flight, 1); + } + + #[test] + fn global_capacity_is_bounded_independently_of_account() { + let runtime = limited_runtime(); + let attempts = (0..runtime.max_pending) + .map(|index| runtime.begin_attempt(index as i64).unwrap().0) + .collect::>(); + assert!(runtime.begin_attempt(9999).is_err()); + drop(attempts); + assert!(runtime.begin_attempt(9999).is_ok()); + } + + #[test] + fn successful_login_resets_pressure() { + let runtime = limited_runtime(); + runtime.begin_attempt(7).unwrap().0.failure(); + let (attempt, delay) = runtime.begin_attempt(7).unwrap(); + assert!(delay > Duration::ZERO); + attempt.success(); + assert_eq!(runtime.attempts.get(&7).unwrap().failures, 0); + assert_eq!(runtime.begin_attempt(7).unwrap().1, Duration::ZERO); + } + + #[test] + fn abandoned_logins_cause_bounded_delay_without_account_lockout() { + let runtime = limited_runtime(); + for _ in 0..30 { + runtime.begin_attempt(7).unwrap().0.failure(); + } + let (attempt, delay) = runtime.begin_attempt(7).unwrap(); + assert_eq!(delay, runtime.throttle.soft_delay_cap); + attempt.cancel(); + } + + #[test] + fn preserved_transport_attempt_fails_after_ttl() { + let runtime = limited_runtime(); + let (attempt, _) = runtime.begin_attempt(7).unwrap(); + let id = Uuid::new_v4(); + runtime.logins.insert( + id, + pending_login( + 7, + attempt, + Instant::now() - runtime.pending_ttl - Duration::from_secs(1), + ), + ); + runtime.prune(); + assert!(!runtime.logins.contains_key(&id)); + assert_eq!(runtime.attempts.get(&7).unwrap().failures, 1); + } + + #[test] + fn password_replacement_cancels_pending_attempts() { + let runtime = limited_runtime(); + let (attempt, _) = runtime.begin_attempt(7).unwrap(); + let id = Uuid::new_v4(); + runtime + .logins + .insert(id, pending_login(7, attempt, Instant::now())); + runtime.cancel_logins_for_user(7); + assert!(!runtime.logins.contains_key(&id)); + assert_eq!(runtime.attempts.get(&7).unwrap().failures, 0); + assert_eq!(runtime.attempts.get(&7).unwrap().in_flight, 0); + } + + #[test] + fn transport_loss_does_not_forgive_pending_attempt() { + let runtime = limited_runtime(); + let (attempt, _) = runtime.begin_attempt(7).unwrap(); + let id = Uuid::new_v4(); + runtime + .logins + .insert(id, pending_login(7, attempt, Instant::now())); + assert!(runtime.logins.contains_key(&id)); + assert_eq!(runtime.attempts.get(&7).unwrap().in_flight, 1); + } + + #[test] + fn cancellation_preserves_previous_failures() { + let runtime = limited_runtime(); + runtime.begin_attempt(7).unwrap().0.failure(); + runtime.begin_attempt(7).unwrap().0.cancel(); + assert_eq!(runtime.attempts.get(&7).unwrap().failures, 1); + } + + #[test] + fn expired_failure_window_resets_pressure() { + let runtime = limited_runtime(); + runtime.begin_attempt(7).unwrap().0.failure(); + runtime.attempts.get_mut(&7).unwrap().window_started -= runtime.throttle.failure_window; + assert_eq!(runtime.begin_attempt(7).unwrap().1, Duration::ZERO); + } + + #[test] + fn finish_errors_distinguish_client_failures_from_server_invalidation() { + for error in [ + PasswordFinishError::Authentication, + PasswordFinishError::InvalidClientMessage, + PasswordFinishError::BindingMismatch, + ] { + assert!(error.counts_as_guess()); + } + for error in [ + PasswordFinishError::CredentialChanged, + PasswordFinishError::CredentialUnavailable, + PasswordFinishError::Storage, + PasswordFinishError::NodeUnavailable, + ] { + assert!(!error.counts_as_guess()); + } + } + + #[test] + fn wrong_password_client_stops_after_ke2_and_expiry_consumes_budget() { + let setup = password::generate_server_setup(); + let principal = b"omega-key:example#7"; + let server = password::server_identifier(11); + let identifiers = Identifiers { + client: Some(principal), + server: Some(&server), + }; + let registration = + ClientRegistration::::start(&mut OsRng, b"correct password") + .unwrap(); + let response = + password::registration_start(&setup, ®istration.message.serialize(), principal) + .unwrap(); + let upload = registration + .state + .finish( + &mut OsRng, + b"correct password", + RegistrationResponse::::deserialize(&response).unwrap(), + ClientRegistrationFinishParameters::new(identifiers, None), + ) + .unwrap(); + let record = password::registration_finish(&upload.message.serialize()).unwrap(); + + let runtime = limited_runtime(); + let id = Uuid::new_v4(); + let contact_key = Keyring::generate().public_key_bundle(); + let context = password::login_context("omega-key:example#7", 11, id, &contact_key).unwrap(); + let client = + ClientLogin::::start(&mut OsRng, b"wrong password").unwrap(); + let ke2 = password::login_start( + &setup, + Some(&record), + &client.message.serialize(), + principal, + &server, + &context, + ) + .unwrap(); + let (attempt, _) = runtime.begin_attempt(7).unwrap(); + assert!( + client + .state + .finish( + &mut OsRng, + b"wrong password", + CredentialResponse::::deserialize(&ke2.response) + .unwrap(), + ClientLoginFinishParameters::new(Some(&context), identifiers, None), + ) + .is_err() + ); + runtime.logins.insert( + id, + PendingLogin { + user_id: 7, + participant_id: 1, + contact_key, + context, + state: ke2.state, + real_record: true, + record_hash: Some(Sha256::digest(&record).to_vec()), + created: Instant::now() - runtime.pending_ttl - Duration::from_secs(1), + attempt: Some(attempt), + }, + ); + runtime.prune(); + assert!(!runtime.logins.contains_key(&id)); + assert_eq!(runtime.attempts.get(&7).unwrap().failures, 1); + assert!(runtime.begin_attempt(7).is_ok()); + } + + #[test] + fn correct_ke3_releases_the_reserved_account_slot() { + let setup = password::generate_server_setup(); + let principal = b"omega-key:example#7"; + let server = password::server_identifier(11); + let identifiers = Identifiers { + client: Some(principal), + server: Some(&server), + }; + let registration = + ClientRegistration::::start(&mut OsRng, b"correct password") + .unwrap(); + let response = + password::registration_start(&setup, ®istration.message.serialize(), principal) + .unwrap(); + let upload = registration + .state + .finish( + &mut OsRng, + b"correct password", + RegistrationResponse::::deserialize(&response).unwrap(), + ClientRegistrationFinishParameters::new(identifiers, None), + ) + .unwrap(); + let record = password::registration_finish(&upload.message.serialize()).unwrap(); + let id = Uuid::new_v4(); + let context = password::login_context( + "omega-key:example#7", + 11, + id, + &Keyring::generate().public_key_bundle(), + ) + .unwrap(); + let client = + ClientLogin::::start(&mut OsRng, b"correct password") + .unwrap(); + let ke2 = password::login_start( + &setup, + Some(&record), + &client.message.serialize(), + principal, + &server, + &context, + ) + .unwrap(); + let runtime = limited_runtime(); + let (attempt, _) = runtime.begin_attempt(7).unwrap(); + let ke3 = client + .state + .finish( + &mut OsRng, + b"correct password", + CredentialResponse::::deserialize(&ke2.response).unwrap(), + ClientLoginFinishParameters::new(Some(&context), identifiers, None), + ) + .unwrap(); + password::login_finish( + &ke2.state, + &ke3.message.serialize(), + principal, + &server, + &context, + ) + .unwrap(); + attempt.success(); + assert_eq!(runtime.begin_attempt(7).unwrap().1, Duration::ZERO); + } + + #[test] + fn password_management_requires_signed_encrypted_single_use_frames() { + let account = Keyring::generate(); + let iota = Keyring::generate(); + let signer = DualSigner::new( + &account.sig_cl_secret_key, + &account.sig_pq_secret_key, + &account.sig_pq_public_key, + ) + .unwrap(); + let frame = ProtectedMessageBuilder::new( + CommunicationType::PasswordEnrollmentFinish, + DataValue::Container(vec![ + typed( + DataType::EncryptedCredential, + DataValue::Bytes(vec![1, 2, 3]), + ) + .unwrap(), + ]), + 7, + 11, + &signer, + purpose(PASSWORD_MANAGEMENT_SIGNATURE), + purpose(PASSWORD_MANAGEMENT_ENCRYPTION), + ) + .message_id(123_u128) + .created_at(now_millis()) + .recipients(vec![iota.public_key_bundle()]) + .frame_id(42) + .build() + .unwrap(); + let options = ProtectedOpenOptions::new( + Some(11), + purpose(PASSWORD_MANAGEMENT_SIGNATURE), + purpose(PASSWORD_MANAGEMENT_ENCRYPTION), + ProtectionPolicy::dual(), + ); + let mut guard = InMemoryReplayGuard::default(); + let resolve = |_| Some(vec![account.public_key_bundle()]); + let opened = + open_protected_with_checked(&frame, &[&iota], Some(7), resolve, options, &mut guard) + .unwrap(); + assert_eq!(opened.signer_id(), 7); + assert!( + open_protected_with_checked(&frame, &[&iota], Some(7), resolve, options, &mut guard) + .is_err() + ); + assert!( + open_protected_with_checked( + &frame, + &[&iota], + Some(7), + resolve, + ProtectedOpenOptions::new( + Some(12), + purpose(PASSWORD_MANAGEMENT_SIGNATURE), + purpose(PASSWORD_MANAGEMENT_ENCRYPTION), + ProtectionPolicy::dual() + ), + &mut InMemoryReplayGuard::default() + ) + .is_err() + ); + let clear = CommunicationValue::new(CommunicationType::PasswordEnrollmentFinish) + .with_receiver(11) + .with_id(42); + assert!( + open_protected_with_checked( + &clear, + &[&iota], + Some(7), + resolve, + options, + &mut InMemoryReplayGuard::default() + ) + .is_err() + ); + } +} diff --git a/omikron-connector/src/password/protected.rs b/omikron-connector/src/password/protected.rs new file mode 100644 index 0000000..b0bb6ff --- /dev/null +++ b/omikron-connector/src/password/protected.rs @@ -0,0 +1,168 @@ +use iota_storage::{users::user_manager, util::protected_replay}; +use mtp::{ + codec::{ + CommunicationType, CommunicationValue, DataValue, ProtectedMessageBuilder, + ProtectedOpenOptions, ProtectionPolicy, ReplayError, ReplayGuard, VerifiedProtectedMessage, + open_protected_with_checked, + }, + crypto::{DualSigner, PublicKeyBundle}, +}; +use rand_core::{OsRng, RngCore}; + +use super::{ + OmikronConnection, + app_protection::{PASSWORD_MANAGEMENT_ENCRYPTION, PASSWORD_MANAGEMENT_SIGNATURE, purpose}, + wire::now_millis, +}; + +struct CommandReplayGuard; + +impl ReplayGuard for CommandReplayGuard { + fn accept( + &mut self, + signer_id: u64, + message_id: mtp::common::MessageId, + created_at: u64, + ) -> Result { + let now = now_millis(); + if created_at < now.saturating_sub(7 * 24 * 60 * 60 * 1000) + || created_at > now.saturating_add(5 * 60 * 1000) + { + return Ok(false); + } + let signer = + i64::try_from(signer_id).map_err(|error| ReplayError::Store(error.to_string()))?; + let time = + i64::try_from(created_at).map_err(|error| ReplayError::Store(error.to_string()))?; + protected_replay::accept(signer, &message_id.to_string(), time) + .map_err(|error| ReplayError::Store(error.to_string())) + } +} + +impl OmikronConnection { + pub(super) async fn open_password_command( + &self, + frame: &CommunicationValue, + kind: CommunicationType, + ) -> Result { + let iota_id = iota_storage::util::config_util::CONFIG + .load() + .iota_id + .ok_or("Iota has no identity")?; + let keyring = self + .keyring + .read() + .await + .clone() + .ok_or("Iota has no keyring")?; + let mut replay = CommandReplayGuard; + let opened = open_protected_with_checked( + frame, + &[keyring.as_ref()], + None, + |signer_id| { + let id = i64::try_from(signer_id).ok()?; + let profile = user_manager::get_user(id).ok()??; + Some(vec![ + PublicKeyBundle::from_base64(&profile.public_key).ok()?, + ]) + }, + ProtectedOpenOptions::new( + Some(iota_id), + purpose(PASSWORD_MANAGEMENT_SIGNATURE), + purpose(PASSWORD_MANAGEMENT_ENCRYPTION), + ProtectionPolicy::dual(), + ), + &mut replay, + ) + .map_err(|error| error.to_string())?; + if opened.message_type() != kind { + return Err("wrong protected command type".into()); + } + Ok(opened) + } + + pub(super) async fn protected_response( + &self, + frame: &CommunicationValue, + kind: CommunicationType, + user_id: i64, + content: DataValue, + signature: u8, + encryption: u8, + ) -> Result { + let iota_id = iota_storage::util::config_util::CONFIG + .load() + .iota_id + .ok_or("Iota has no identity")?; + let keyring = self + .keyring + .read() + .await + .clone() + .ok_or("Iota has no keyring")?; + let recipient = PublicKeyBundle::from_base64( + &user_manager::get_user(user_id) + .map_err(|error| error.to_string())? + .ok_or("account not hosted")? + .public_key, + ) + .map_err(|error| error.to_string())?; + self.build_protected( + frame, + kind, + iota_id, + user_id as u64, + content, + recipient, + signature, + encryption, + &keyring, + ) + } + + pub(super) fn build_protected( + &self, + frame: &CommunicationValue, + kind: CommunicationType, + signer_id: u64, + recipient_id: u64, + content: DataValue, + recipient: PublicKeyBundle, + signature: u8, + encryption: u8, + keyring: &mtp::crypto::Keyring, + ) -> Result { + let signer = DualSigner::new( + &keyring.sig_cl_secret_key, + &keyring.sig_pq_secret_key, + &keyring.sig_pq_public_key, + ) + .map_err(|error| error.to_string())?; + let mut rng = OsRng; + let id = ((rng.next_u64() as u128) << 64) | rng.next_u64() as u128; + ProtectedMessageBuilder::new( + kind, + content, + signer_id, + recipient_id, + &signer, + purpose(signature), + purpose(encryption), + ) + .message_id(id) + .created_at(now_millis()) + .recipients(vec![recipient]) + .frame_id(frame.id().ok_or("missing request ID")?) + .build() + .map_err(|error| error.to_string()) + } + + pub(super) async fn password_error(&self, frame: &CommunicationValue, kind: CommunicationType) { + let mut response = CommunicationValue::new(kind); + if let Some(id) = frame.id() { + response = response.with_id(id); + } + let _ = self.send_message(&response).await; + } +} diff --git a/omikron-connector/src/password/provisioning.rs b/omikron-connector/src/password/provisioning.rs new file mode 100644 index 0000000..c5259f5 --- /dev/null +++ b/omikron-connector/src/password/provisioning.rs @@ -0,0 +1,468 @@ +use std::{ + sync::Arc, + time::{Duration, Instant}, +}; + +use dashmap::mapref::entry::Entry; +use iota_auth::password::{self, PasswordServerSetup}; +use iota_storage::users::{password_credentials, user_manager}; +use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue}; +use mtp::crypto::PublicKeyBundle; +use sha2::{Digest, Sha256}; +use uuid::Uuid; + +use super::{ + MAX_CREDENTIAL_BYTES, MAX_OPAQUE_BYTES, OmikronConnection, PasswordFinishError, PendingLogin, + app_protection::{PROVISIONING_CREDENTIAL_ENCRYPTION, PROVISIONING_CREDENTIAL_SIGNATURE}, + runtime::AttemptLease, + wire::{bytes, key_fingerprint, positive, public_key, session, typed}, +}; + +struct PreparedPasswordStart { + request_id: u32, + session_id: Uuid, + participant_id: u64, + user_id: i64, + iota_id: i64, + contact_key: PublicKeyBundle, + ke1: Vec, + setup: Arc, + attempt: AttemptLease, + delay: Duration, +} + +impl OmikronConnection { + pub(crate) async fn handle_password_provisioning_start( + self: Arc, + frame: &CommunicationValue, + ) { + let prepared = match self.prepare_password_provisioning_start(frame) { + Ok(prepared) => prepared, + Err(error) => { + iota_logger::log!("Password KE1 admission failed: {error}"); + self.password_error(frame, CommunicationType::ErrorNotAuthenticated) + .await; + return; + } + }; + let request_id = prepared.request_id; + let connection = self.clone(); + tokio::spawn(async move { + match connection + .execute_password_provisioning_start(prepared) + .await + { + Ok(response) => { + if let Err(error) = connection.send_message(&response).await { + iota_logger::log!("Password KE2 delivery failed: {error}"); + } + } + Err(error) => { + iota_logger::log!("Password KE1 processing failed: {error}"); + let response = + CommunicationValue::new(CommunicationType::ErrorNotAuthenticated) + .with_id(request_id); + let _ = connection.send_message(&response).await; + } + } + }); + } + + fn prepare_password_provisioning_start( + &self, + frame: &CommunicationValue, + ) -> Result { + let content = frame.payload(); + let request_id = frame.id().ok_or("missing request ID")?; + let session_id = session(content)?; + let participant_id = u64::try_from(positive(content, DataType::ProvisioningParticipantId)?) + .map_err(|error| error.to_string())?; + let user_id = positive(content, DataType::UserId)?; + let iota_id = positive(content, DataType::IotaId)?; + if Some(iota_id as u64) != iota_storage::util::config_util::CONFIG.load().iota_id { + return Err("wrong Iota".into()); + } + let contact_key = public_key(content)?; + let ke1 = bytes(content, DataType::OpaqueMessage, MAX_OPAQUE_BYTES)?; + password::validate_login_request(&ke1).map_err(|error| error.to_string())?; + let setup = self + .password_auth + .setup() + .map_err(|error| error.to_string())?; + let (attempt, delay) = self + .password_auth + .begin_attempt(user_id) + .map_err(|error| error.to_string())?; + Ok(PreparedPasswordStart { + request_id, + session_id, + participant_id, + user_id, + iota_id, + contact_key, + ke1, + setup, + attempt, + delay, + }) + } + + async fn execute_password_provisioning_start( + &self, + prepared: PreparedPasswordStart, + ) -> Result { + let PreparedPasswordStart { + request_id, + session_id, + participant_id, + user_id, + iota_id, + contact_key, + ke1, + setup, + attempt, + delay, + } = prepared; + if !delay.is_zero() { + tokio::time::sleep(delay).await; + } + let principal = self.password_principal(user_id).await?; + let context = password::login_context(&principal, iota_id, session_id, &contact_key) + .map_err(|error| error.to_string())?; + let credential = password_credentials::get(user_id).map_err(|error| error.to_string())?; + let compatible_credential = credential + .as_ref() + .filter(|credential| credential.opaque_profile == password::CURRENT_OPAQUE_PROFILE); + let response = password::login_start( + &setup, + compatible_credential.map(|value| value.opaque_record.as_slice()), + &ke1, + principal.as_bytes(), + &password::server_identifier(iota_id), + &context, + ) + .map_err(|error| error.to_string())?; + let pending = PendingLogin { + user_id, + participant_id, + contact_key, + context, + state: response.state, + real_record: compatible_credential.is_some(), + created: Instant::now(), + attempt: Some(attempt), + record_hash: compatible_credential + .map(|value| Sha256::digest(&value.opaque_record).to_vec()), + }; + match self.password_auth.logins.entry(session_id) { + Entry::Vacant(slot) => { + slot.insert(pending); + } + Entry::Occupied(_) => { + return Err("session already started".into()); + } + } + Ok( + CommunicationValue::new(CommunicationType::PasswordProvisioningResponse) + .with_id(request_id) + .add_typed_default(DataType::Uuid, DataValue::Str(session_id.to_string())) + .add_typed_default(DataType::OpaqueMessage, DataValue::Bytes(response.response)), + ) + } + + pub(crate) async fn handle_password_provisioning_finish( + self: Arc, + frame: &CommunicationValue, + ) { + match self.provisioning_finish(frame).await { + Ok(response) => { + let _ = self.send_message(&response).await; + } + Err(error) => { + iota_logger::log!("Password KE3 failed: {error}"); + self.password_error(frame, CommunicationType::ErrorNotAuthenticated) + .await; + } + } + } + + pub(super) async fn provisioning_finish( + &self, + frame: &CommunicationValue, + ) -> Result { + let id = session(frame.payload())?; + let (_, mut pending) = self + .password_auth + .logins + .remove(&id) + .ok_or("OPAQUE session expired")?; + match self.finish_pending_login(frame, id, &mut pending).await { + Ok(response) => Ok(response), + Err(error) => { + if let Some(attempt) = pending.attempt.take() { + if error.counts_as_guess() { + attempt.failure(); + } else { + attempt.cancel(); + } + } + Err(error.to_string()) + } + } + } + + async fn finish_pending_login( + &self, + frame: &CommunicationValue, + id: Uuid, + pending: &mut PendingLogin, + ) -> Result { + let content = frame.payload(); + if pending.created.elapsed() >= self.password_auth.pending_ttl + || pending.user_id + != positive(content, DataType::UserId) + .map_err(|_| PasswordFinishError::InvalidClientMessage)? + || pending.participant_id + != u64::try_from( + positive(content, DataType::ProvisioningParticipantId) + .map_err(|_| PasswordFinishError::InvalidClientMessage)?, + ) + .map_err(|_| PasswordFinishError::InvalidClientMessage)? + || pending + .contact_key + .try_as_bytes() + .map_err(|_| PasswordFinishError::InvalidClientMessage)? + != public_key(content) + .map_err(|_| PasswordFinishError::InvalidClientMessage)? + .try_as_bytes() + .map_err(|_| PasswordFinishError::InvalidClientMessage)? + { + return Err(PasswordFinishError::BindingMismatch); + } + let iota_id = positive(content, DataType::IotaId) + .map_err(|_| PasswordFinishError::InvalidClientMessage)?; + if Some(iota_id as u64) != iota_storage::util::config_util::CONFIG.load().iota_id { + return Err(PasswordFinishError::BindingMismatch); + } + let ke3 = bytes(content, DataType::OpaqueMessage, MAX_OPAQUE_BYTES) + .map_err(|_| PasswordFinishError::InvalidClientMessage)?; + let principal = self + .password_principal(pending.user_id) + .await + .map_err(|_| PasswordFinishError::CredentialUnavailable)?; + match password::login_finish( + &pending.state, + &ke3, + principal.as_bytes(), + &password::server_identifier(iota_id), + &pending.context, + ) { + Ok(()) => {} + Err(password::LoginFinishError::AuthenticationFailed) => { + return Err(PasswordFinishError::Authentication); + } + Err(password::LoginFinishError::InvalidMessage) => { + return Err(PasswordFinishError::InvalidClientMessage); + } + } + if !pending.real_record { + return Err(PasswordFinishError::Authentication); + } + let credential = password_credentials::get(pending.user_id) + .map_err(|_| PasswordFinishError::Storage)? + .ok_or(PasswordFinishError::CredentialUnavailable)?; + if credential.opaque_profile != password::CURRENT_OPAQUE_PROFILE { + return Err(PasswordFinishError::CredentialUnavailable); + } + if pending.record_hash.as_deref() + != Some(Sha256::digest(&credential.opaque_record).as_slice()) + { + return Err(PasswordFinishError::CredentialChanged); + } + // Only the same registered credential can clear password pressure. + if let Some(attempt) = pending.attempt.take() { + attempt.success(); + } + let profile = user_manager::get_user(pending.user_id) + .map_err(|_| PasswordFinishError::Storage)? + .ok_or(PasswordFinishError::CredentialUnavailable)?; + if credential.account_public_key_sha256 + != key_fingerprint(&profile.public_key) + .map_err(|_| PasswordFinishError::CredentialUnavailable)? + || credential.encrypted_tu_credential.is_empty() + || credential.encrypted_tu_credential.len() > MAX_CREDENTIAL_BYTES + { + return Err(PasswordFinishError::CredentialUnavailable); + } + let wire = + |kind, value| typed(kind, value).map_err(|_| PasswordFinishError::InvalidClientMessage); + let content = DataValue::Container(vec![ + wire(DataType::Uuid, DataValue::Str(id.to_string()))?, + wire( + DataType::ProvisioningMethod, + DataValue::Str("password".into()), + )?, + wire( + DataType::UserId, + DataValue::SignedNumber(pending.user_id.into()), + )?, + wire(DataType::IotaId, DataValue::SignedNumber(iota_id.into()))?, + wire( + DataType::EncryptedCredential, + DataValue::Bytes(credential.encrypted_tu_credential), + )?, + wire( + DataType::Sha256, + DataValue::Bytes(credential.account_public_key_sha256), + )?, + ]); + let keyring = self + .keyring + .read() + .await + .clone() + .ok_or(PasswordFinishError::NodeUnavailable)?; + self.build_protected( + frame, + CommunicationType::ProvisioningCredential, + iota_id as u64, + pending.participant_id, + content, + pending.contact_key.clone(), + PROVISIONING_CREDENTIAL_SIGNATURE, + PROVISIONING_CREDENTIAL_ENCRYPTION, + &keyring, + ) + .map_err(|_| PasswordFinishError::NodeUnavailable) + } +} + +#[cfg(test)] +mod admission_tests { + use super::*; + use dashmap::DashSet; + use iota_storage::util::config_util::CONFIG; + use mtp::crypto::Keyring; + use opaque_ke::ClientLogin; + use rand_core::OsRng; + use tokio::sync::Semaphore; + + #[tokio::test] + async fn delayed_password_start_releases_general_handler_capacity() { + let mut config = (**CONFIG.load()).clone(); + let iota_id = config.iota_id.unwrap_or(11); + config.iota_id = Some(iota_id); + let previous = CONFIG.swap(Arc::new(config)); + + let mut connection = OmikronConnection::new( + Arc::new(DashSet::new()), + Arc::new(std::sync::Mutex::new(iota_state::AppState::new())), + ); + Arc::get_mut(&mut connection.password_auth) + .unwrap() + .throttle + .soft_delay_step = Duration::from_millis(400); + let connection = Arc::new(connection); + connection + .password_auth + .setup + .set(Arc::new(password::generate_server_setup())) + .ok() + .unwrap(); + let user_id = 777; + connection + .password_auth + .begin_attempt(user_id) + .unwrap() + .0 + .failure(); + let ke1 = ClientLogin::::start(&mut OsRng, b"password") + .unwrap() + .message + .serialize() + .to_vec(); + let frame = CommunicationValue::new(CommunicationType::PasswordProvisioningStart) + .with_id(42) + .add_typed_default(DataType::Uuid, DataValue::Str(Uuid::new_v4().to_string())) + .add_typed_default( + DataType::ProvisioningParticipantId, + DataValue::UnsignedNumber(1), + ) + .add_typed_default(DataType::UserId, DataValue::SignedNumber(user_id.into())) + .add_typed_default(DataType::IotaId, DataValue::SignedNumber(iota_id.into())) + .add_typed_default( + DataType::PublicKey, + DataValue::Str( + Keyring::generate() + .public_key_bundle() + .try_to_base64() + .unwrap(), + ), + ) + .add_typed_default(DataType::OpaqueMessage, DataValue::Bytes(ke1)); + let malformed = frame.clone().with_data( + DataType::OpaqueMessage, + DataValue::Bytes(b"not a KE1".to_vec()), + ); + assert!( + connection + .prepare_password_provisioning_start(&malformed) + .is_err() + ); + assert_eq!( + connection + .password_auth + .attempts + .get(&user_id) + .unwrap() + .in_flight, + 0 + ); + let general = Arc::new(Semaphore::new(1)); + let permit = general.clone().acquire_owned().await.unwrap(); + connection + .clone() + .handle_password_provisioning_start(&frame) + .await; + drop(permit); + + assert_eq!( + connection + .password_auth + .attempts + .get(&user_id) + .unwrap() + .in_flight, + 1 + ); + let normal = tokio::time::timeout(Duration::from_millis(100), general.acquire()) + .await + .expect("normal work must acquire the general permit during password delay") + .unwrap(); + drop(normal); + assert_eq!( + connection + .password_auth + .attempts + .get(&user_id) + .unwrap() + .in_flight, + 1 + ); + tokio::time::timeout(Duration::from_secs(2), async { + while connection + .password_auth + .attempts + .get(&user_id) + .unwrap() + .in_flight + != 0 + { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .unwrap(); + CONFIG.store(previous); + } +} diff --git a/omikron-connector/src/password/runtime.rs b/omikron-connector/src/password/runtime.rs new file mode 100644 index 0000000..ae99eaa --- /dev/null +++ b/omikron-connector/src/password/runtime.rs @@ -0,0 +1,285 @@ +use std::{ + path::Path, + sync::{Arc, OnceLock, Weak}, + time::{Duration, Instant}, +}; + +use dashmap::DashMap; +use iota_auth::password::{self, PasswordServerSetup}; +use iota_storage::users::password_credentials; +use mtp::crypto::PublicKeyBundle; +use tokio::sync::{OwnedSemaphorePermit, Semaphore}; +use uuid::Uuid; +use zeroize::Zeroizing; + +pub(super) struct PendingLogin { + pub(super) user_id: i64, + pub(super) participant_id: u64, + pub(super) contact_key: PublicKeyBundle, + pub(super) context: Vec, + pub(super) state: password::SerializedLoginState, + pub(super) real_record: bool, + pub(super) record_hash: Option>, + pub(super) created: Instant, + pub(super) attempt: Option, +} + +#[derive(Debug)] +pub(super) struct AccountAttemptState { + pub(super) window_started: Instant, + pub(super) failures: u32, + pub(super) in_flight: u32, +} + +pub(super) struct PasswordThrottleConfig { + pub(super) failure_window: Duration, + pub(super) soft_delay_step: Duration, + pub(super) soft_delay_cap: Duration, +} + +#[derive(Debug, thiserror::Error)] +pub(super) enum AttemptAdmissionError { + #[error("password authentication service busy")] + ServerBusy, +} + +pub(super) struct AttemptLease { + runtime: Weak, + user_id: i64, + settled: bool, + _global_permit: OwnedSemaphorePermit, +} + +impl AttemptLease { + pub(super) fn success(mut self) { + if let Some(runtime) = self.runtime.upgrade() { + runtime.finish_attempt_success(self.user_id); + } + self.settled = true; + } + + pub(super) fn failure(mut self) { + if let Some(runtime) = self.runtime.upgrade() { + runtime.finish_attempt_failure(self.user_id); + } + self.settled = true; + } + + pub(super) fn cancel(mut self) { + if let Some(runtime) = self.runtime.upgrade() { + runtime.finish_attempt_cancel(self.user_id); + } + self.settled = true; + } +} + +impl Drop for AttemptLease { + fn drop(&mut self) { + if !self.settled { + if let Some(runtime) = self.runtime.upgrade() { + runtime.finish_attempt_cancel(self.user_id); + } + } + } +} + +pub(crate) struct PasswordAuthRuntime { + pub(super) setup: OnceLock>, + pub(super) omega_authority: OnceLock, + pub(super) logins: DashMap, + pub(super) enrollments: DashMap, + pub(super) attempts: DashMap, + pub(super) max_pending: usize, + exchange_slots: Arc, + pub(super) pending_ttl: Duration, + pub(super) throttle: PasswordThrottleConfig, +} + +fn configured_positive(name: &str, default: T) -> T { + std::env::var(name) + .ok() + .and_then(|value| value.parse::().ok()) + .filter(|value| *value > T::default()) + .unwrap_or(default) +} + +impl Default for PasswordAuthRuntime { + fn default() -> Self { + let max_pending = configured_positive("PASSWORD_MAX_PENDING_EXCHANGES", 1024); + Self { + setup: OnceLock::new(), + omega_authority: OnceLock::new(), + logins: DashMap::new(), + enrollments: DashMap::new(), + attempts: DashMap::new(), + max_pending, + exchange_slots: Arc::new(Semaphore::new(max_pending)), + pending_ttl: Duration::from_secs(configured_positive( + "PASSWORD_PENDING_TTL_SECONDS", + 180, + )), + throttle: PasswordThrottleConfig { + failure_window: Duration::from_secs(configured_positive( + "PASSWORD_ATTEMPT_WINDOW_SECONDS", + 300, + )), + soft_delay_step: Duration::from_millis(configured_positive( + "PASSWORD_SOFT_DELAY_STEP_MS", + 250, + )), + soft_delay_cap: Duration::from_millis(configured_positive( + "PASSWORD_SOFT_DELAY_CAP_MS", + 3000, + )), + }, + } + } +} + +impl PasswordAuthRuntime { + pub(crate) fn setup(&self) -> Result, PasswordRuntimeError> { + self.setup + .get() + .cloned() + .ok_or(PasswordRuntimeError::SetupUnavailable) + } + + pub(crate) fn initialize(&self, identity: &Path) -> Result<(), String> { + if self.setup.get().is_some() { + return Ok(()); + } + let path = identity.with_file_name("password-auth.setup"); + let setup = match std::fs::read(&path) { + Ok(bytes) => password::deserialize_server_setup_owned(bytes) + .map_err(|error| error.to_string())?, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => { + if password_credentials::any().map_err(|error| error.to_string())? { + return Err("OPAQUE setup is missing while password credentials exist".into()); + } + let setup = password::generate_server_setup(); + let serialized = Zeroizing::new(password::serialize_server_setup(&setup)); + iota_util::atomic_file::replace_private(&path, &serialized, 0) + .map_err(|error| error.to_string())?; + setup + } + Err(error) => return Err(error.to_string()), + }; + let _ = self.setup.set(Arc::new(setup)); + Ok(()) + } + + pub(crate) fn prune(&self) { + let expired = self + .logins + .iter() + .filter_map(|entry| { + (entry.created.elapsed() >= self.pending_ttl).then_some(*entry.key()) + }) + .collect::>(); + for id in expired { + if let Some((_, pending)) = self.logins.remove(&id) { + if let Some(attempt) = pending.attempt { + attempt.failure(); + } + } + } + self.enrollments + .retain(|_, (_, created)| created.elapsed() < self.pending_ttl); + self.attempts.retain(|_, state| { + state.in_flight > 0 || state.window_started.elapsed() < self.throttle.failure_window + }); + } + + pub(super) fn begin_attempt( + self: &Arc, + user_id: i64, + ) -> Result<(AttemptLease, Duration), AttemptAdmissionError> { + let permit = self + .exchange_slots + .clone() + .try_acquire_owned() + .map_err(|_| AttemptAdmissionError::ServerBusy)?; + let mut state = self + .attempts + .entry(user_id) + .or_insert_with(|| AccountAttemptState { + window_started: Instant::now(), + failures: 0, + in_flight: 0, + }); + if state.window_started.elapsed() >= self.throttle.failure_window { + state.window_started = Instant::now(); + state.failures = 0; + } + state.in_flight = state.in_flight.saturating_add(1); + let pressure = state + .failures + .saturating_add(state.in_flight.saturating_sub(1)); + let delay = self + .throttle + .soft_delay_step + .saturating_mul(pressure) + .min(self.throttle.soft_delay_cap); + Ok(( + AttemptLease { + runtime: Arc::downgrade(self), + user_id, + settled: false, + _global_permit: permit, + }, + delay, + )) + } + + fn finish_attempt_success(&self, user_id: i64) { + if let Some(mut state) = self.attempts.get_mut(&user_id) { + state.in_flight = state.in_flight.saturating_sub(1); + state.failures = 0; + // Keep the idle entry until pruning to avoid racing a new reservation. + } + } + + fn finish_attempt_failure(&self, user_id: i64) { + let mut state = self + .attempts + .entry(user_id) + .or_insert_with(|| AccountAttemptState { + window_started: Instant::now(), + failures: 0, + in_flight: 0, + }); + if state.window_started.elapsed() >= self.throttle.failure_window { + state.window_started = Instant::now(); + state.failures = 0; + } + state.in_flight = state.in_flight.saturating_sub(1); + state.failures = state.failures.saturating_add(1); + } + + fn finish_attempt_cancel(&self, user_id: i64) { + if let Some(mut state) = self.attempts.get_mut(&user_id) { + state.in_flight = state.in_flight.saturating_sub(1); + } + } + + pub(super) fn cancel_logins_for_user(&self, user_id: i64) { + let ids = self + .logins + .iter() + .filter_map(|entry| (entry.user_id == user_id).then_some(*entry.key())) + .collect::>(); + for id in ids { + if let Some((_, pending)) = self.logins.remove(&id) { + if let Some(attempt) = pending.attempt { + attempt.cancel(); + } + } + } + } +} + +#[derive(Debug, thiserror::Error)] +pub(crate) enum PasswordRuntimeError { + #[error("password authentication setup unavailable")] + SetupUnavailable, +} diff --git a/omikron-connector/src/password/wire.rs b/omikron-connector/src/password/wire.rs new file mode 100644 index 0000000..db61e53 --- /dev/null +++ b/omikron-connector/src/password/wire.rs @@ -0,0 +1,79 @@ +use std::time::{SystemTime, UNIX_EPOCH}; + +use mtp::{ + codec::{DataType, DataTypeId, DataValue, TypeMap}, + crypto::PublicKeyBundle, +}; +use sha2::{Digest, Sha256}; +use uuid::Uuid; + +pub(super) fn now_millis() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_millis() + .try_into() + .unwrap_or(u64::MAX) +} + +pub(super) fn field(content: &DataValue, kind: DataType) -> Result<&DataValue, String> { + let id = kind + .try_to_id(&TypeMap::latest()) + .ok_or("unknown data field")?; + let DataValue::Container(fields) = content else { + return Err("content is not a container".into()); + }; + fields + .iter() + .find_map(|(field_id, value)| (*field_id == id).then_some(value)) + .ok_or_else(|| format!("missing {kind:?}")) +} + +pub(super) fn typed(kind: DataType, value: DataValue) -> Result<(DataTypeId, DataValue), String> { + Ok(( + kind.try_to_id(&TypeMap::latest()).ok_or("unknown field")?, + value, + )) +} + +pub(super) fn bytes(content: &DataValue, kind: DataType, max: usize) -> Result, String> { + let value = field(content, kind)? + .as_bytes() + .ok_or("field is not bytes")?; + if value.is_empty() || value.len() > max { + return Err("field exceeds size limits".into()); + } + Ok(value) +} + +pub(super) fn positive(content: &DataValue, kind: DataType) -> Result { + field(content, kind)? + .as_number() + .and_then(|value| i64::try_from(value).ok()) + .filter(|value| *value > 0) + .ok_or_else(|| format!("invalid {kind:?}")) +} + +pub(super) fn session(content: &DataValue) -> Result { + Uuid::parse_str( + field(content, DataType::Uuid)? + .as_str() + .ok_or("invalid session UUID")?, + ) + .map_err(|error| error.to_string()) +} + +pub(super) fn public_key(content: &DataValue) -> Result { + let encoded = field(content, DataType::PublicKey)? + .as_str() + .ok_or("invalid public key")?; + if encoded.len() > 16 * 1024 { + return Err("public key too large".into()); + } + PublicKeyBundle::from_base64(encoded).map_err(|error| error.to_string()) +} + +pub(super) fn key_fingerprint(encoded: &str) -> Result, String> { + let bundle = PublicKeyBundle::from_base64(encoded).map_err(|error| error.to_string())?; + Ok(Sha256::digest(bundle.try_as_bytes().map_err(|error| error.to_string())?).to_vec()) +} diff --git a/omikron-connector/src/password_provisioning.rs b/omikron-connector/src/password_provisioning.rs deleted file mode 100644 index c50b63d..0000000 --- a/omikron-connector/src/password_provisioning.rs +++ /dev/null @@ -1,904 +0,0 @@ -use std::{ - path::Path, - sync::{Arc, OnceLock}, - time::{Duration, Instant, SystemTime, UNIX_EPOCH}, -}; - -use dashmap::{DashMap, mapref::entry::Entry}; -use iota_auth::password::{self, PasswordServerSetup}; -use iota_storage::{ - users::{ - password_credentials::{self, PasswordCredential}, - user_manager, - }, - util::protected_replay, -}; -use mtp::{ - codec::{ - CommunicationType, CommunicationValue, DataType, DataTypeId, - DataValue, ProtectedMessageBuilder, ProtectedOpenOptions, ProtectionPolicy, - ReplayError, ReplayGuard, TypeMap, VerifiedProtectedMessage, - open_protected_with_checked, - }, - crypto::{DualSigner, PublicKeyBundle}, -}; -use rand_core::{OsRng, RngCore}; -use sha2::{Digest, Sha256}; -use uuid::Uuid; - -use crate::omikron_connection::OmikronConnection; - -#[allow(dead_code)] -mod app_protection { - include!(concat!(env!("CARGO_MANIFEST_DIR"), "/../mtp-type-maps/app_protection.rs")); -} -use app_protection::{ - purpose, PASSWORD_MANAGEMENT_SIGNATURE, PASSWORD_MANAGEMENT_ENCRYPTION, - PASSWORD_RESPONSE_SIGNATURE, PASSWORD_RESPONSE_ENCRYPTION, - PROVISIONING_CREDENTIAL_SIGNATURE, PROVISIONING_CREDENTIAL_ENCRYPTION, -}; - -const MAX_OPAQUE_BYTES: usize = 16 * 1024; -const MAX_CREDENTIAL_BYTES: usize = 256 * 1024; - -struct CommandReplayGuard; - -impl ReplayGuard for CommandReplayGuard { - fn accept( - &mut self, - signer_id: u64, - message_id: mtp::common::MessageId, - created_at: u64, - ) -> Result { - let now = now_millis(); - if created_at < now.saturating_sub(7 * 24 * 60 * 60 * 1000) - || created_at > now.saturating_add(5 * 60 * 1000) - { - return Ok(false); - } - let signer = - i64::try_from(signer_id).map_err(|error| ReplayError::Store(error.to_string()))?; - let time = - i64::try_from(created_at).map_err(|error| ReplayError::Store(error.to_string()))?; - protected_replay::accept(signer, &message_id.to_string(), time) - .map_err(|error| ReplayError::Store(error.to_string())) - } -} - -pub(super) struct PendingLogin { - user_id: i64, - participant_id: u64, - contact_key: PublicKeyBundle, - context: Vec, - state: Vec, - real_record: bool, - record_hash: Option>, - created: Instant, -} - -pub(super) struct PasswordAuthRuntime { - setup: OnceLock>, - omega_authority: OnceLock, - logins: DashMap, - enrollments: DashMap, - attempts: DashMap, - max_pending: usize, - pending_ttl: Duration, - attempt_window: Duration, - max_attempts: u32, -} - -fn configured_positive(name: &str, default: T) -> T { - std::env::var(name) - .ok() - .and_then(|value| value.parse::().ok()) - .filter(|value| *value > T::default()) - .unwrap_or(default) -} - -impl Default for PasswordAuthRuntime { - fn default() -> Self { - Self { - setup: OnceLock::new(), - omega_authority: OnceLock::new(), - logins: DashMap::new(), - enrollments: DashMap::new(), - attempts: DashMap::new(), - max_pending: configured_positive("PASSWORD_MAX_PENDING_EXCHANGES", 1024), - pending_ttl: Duration::from_secs(configured_positive( - "PASSWORD_PENDING_TTL_SECONDS", - 180, - )), - attempt_window: Duration::from_secs(configured_positive( - "PASSWORD_ATTEMPT_WINDOW_SECONDS", - 300, - )), - max_attempts: configured_positive("PASSWORD_MAX_ATTEMPTS_PER_ACCOUNT", 10), - } - } -} - -impl PasswordAuthRuntime { - pub(super) fn initialize(&self, identity: &Path) -> Result<(), String> { - if self.setup.get().is_some() { - return Ok(()); - } - let path = identity.with_file_name("password-auth.setup"); - let setup = match std::fs::read(&path) { - Ok(bytes) => { - password::deserialize_server_setup(&bytes).map_err(|error| error.to_string())? - } - Err(error) if error.kind() == std::io::ErrorKind::NotFound => { - if password_credentials::any().map_err(|error| error.to_string())? { - return Err("OPAQUE setup is missing while password credentials exist".into()); - } - let setup = password::generate_server_setup(); - iota_util::atomic_file::replace_private( - &path, - &password::serialize_server_setup(&setup), - 0, - ) - .map_err(|error| error.to_string())?; - setup - } - Err(error) => return Err(error.to_string()), - }; - let _ = self.setup.set(Arc::new(setup)); - Ok(()) - } - - pub(super) fn prune(&self) { - self.logins - .retain(|_, value| value.created.elapsed() < self.pending_ttl); - self.enrollments - .retain(|_, (_, created)| created.elapsed() < self.pending_ttl); - self.attempts - .retain(|_, (started, _)| started.elapsed() < self.attempt_window); - } - - fn allow_attempt(&self, user_id: i64) -> bool { - let mut attempt = self.attempts.entry(user_id).or_insert((Instant::now(), 0)); - if attempt.0.elapsed() >= self.attempt_window { - *attempt = (Instant::now(), 0); - } - if attempt.1 >= self.max_attempts { - return false; - } - attempt.1 += 1; - true - } -} - -fn now_millis() -> u64 { - SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_millis() - .try_into() - .unwrap_or(u64::MAX) -} - -fn field<'a>(content: &'a DataValue, kind: DataType) -> Result<&'a DataValue, String> { - let id = kind - .try_to_id(&TypeMap::latest()) - .ok_or("unknown data field")?; - let DataValue::Container(fields) = content else { - return Err("content is not a container".into()); - }; - fields - .iter() - .find_map(|(field_id, value)| (*field_id == id).then_some(value)) - .ok_or_else(|| format!("missing {kind:?}")) -} - -fn typed(kind: DataType, value: DataValue) -> Result<(DataTypeId, DataValue), String> { - Ok(( - kind.try_to_id(&TypeMap::latest()).ok_or("unknown field")?, - value, - )) -} - -fn bytes(content: &DataValue, kind: DataType, max: usize) -> Result, String> { - let value = field(content, kind)? - .as_bytes() - .ok_or("field is not bytes")?; - if value.is_empty() || value.len() > max { - return Err("field exceeds size limits".into()); - } - Ok(value) -} - -fn positive(content: &DataValue, kind: DataType) -> Result { - field(content, kind)? - .as_number() - .and_then(|value| i64::try_from(value).ok()) - .filter(|value| *value > 0) - .ok_or_else(|| format!("invalid {kind:?}")) -} - -fn session(content: &DataValue) -> Result { - Uuid::parse_str( - field(content, DataType::Uuid)? - .as_str() - .ok_or("invalid session UUID")?, - ) - .map_err(|error| error.to_string()) -} - -fn public_key(content: &DataValue) -> Result { - let encoded = field(content, DataType::PublicKey)? - .as_str() - .ok_or("invalid public key")?; - if encoded.len() > 16 * 1024 { - return Err("public key too large".into()); - } - PublicKeyBundle::from_base64(encoded).map_err(|error| error.to_string()) -} - -fn key_fingerprint(encoded: &str) -> Result, String> { - let bundle = PublicKeyBundle::from_base64(encoded).map_err(|error| error.to_string())?; - Ok(Sha256::digest(bundle.try_as_bytes().map_err(|error| error.to_string())?).to_vec()) -} - -impl OmikronConnection { - // Password login uses Omega's key-scoped principal, including for accounts - // whose older hosted-principal record still uses the legacy host locator. - async fn password_principal(&self, user_id: i64) -> Result { - if user_id <= 0 - || user_manager::get_user(user_id) - .map_err(|error| error.to_string())? - .is_none() - { - return Err("account not hosted".into()); - } - let authority = if let Some(cached) = self.password_auth.omega_authority.get() { - cached.clone() - } else { - let url = format!( - "{}/.well-known/tensamin", - crate::omega_discovery::api_base() - ); - let response = - tokio::time::timeout(Duration::from_secs(10), self.http_client.get(url).send()) - .await - .map_err(|_| "Omega identity lookup timed out")? - .map_err(|error| error.to_string())? - .error_for_status() - .map_err(|error| error.to_string())?; - let discovery: serde_json::Value = - response.json().await.map_err(|error| error.to_string())?; - let key = discovery - .get("public_key") - .and_then(serde_json::Value::as_str) - .ok_or("Omega discovery is missing its key")?; - let key = PublicKeyBundle::from_base64(key).map_err(|error| error.to_string())?; - let authority = iota_identity::AuthorityId::for_omega(&key) - .map_err(|error| error.to_string())? - .as_str() - .to_owned(); - if discovery - .get("authority_id") - .and_then(serde_json::Value::as_str) - != Some(authority.as_str()) - { - return Err("Omega discovery identity mismatch".into()); - } - let _ = self.password_auth.omega_authority.set(authority.clone()); - authority - }; - Ok(format!("{authority}#{user_id}")) - } - - async fn open_password_command( - &self, - frame: &CommunicationValue, - kind: CommunicationType, - ) -> Result { - let iota_id = iota_storage::util::config_util::CONFIG - .load() - .iota_id - .ok_or("Iota has no identity")?; - let keyring = self - .keyring - .read() - .await - .clone() - .ok_or("Iota has no keyring")?; - let mut replay = CommandReplayGuard; - let opened = open_protected_with_checked( - frame, - &[keyring.as_ref()], - None, - |signer_id| { - let id = i64::try_from(signer_id).ok()?; - let profile = user_manager::get_user(id).ok()??; - Some(vec![ - PublicKeyBundle::from_base64(&profile.public_key).ok()?, - ]) - }, - ProtectedOpenOptions::new( - Some(iota_id), - purpose(PASSWORD_MANAGEMENT_SIGNATURE), - purpose(PASSWORD_MANAGEMENT_ENCRYPTION), - ProtectionPolicy::dual(), - ), - &mut replay, - ) - .map_err(|error| error.to_string())?; - if opened.message_type() != kind { - return Err("wrong protected command type".into()); - } - Ok(opened) - } - - async fn protected_response( - &self, - frame: &CommunicationValue, - kind: CommunicationType, - user_id: i64, - content: DataValue, - signature: u8, - encryption: u8, - ) -> Result { - let iota_id = iota_storage::util::config_util::CONFIG - .load() - .iota_id - .ok_or("Iota has no identity")?; - let keyring = self - .keyring - .read() - .await - .clone() - .ok_or("Iota has no keyring")?; - let recipient = PublicKeyBundle::from_base64( - &user_manager::get_user(user_id) - .map_err(|error| error.to_string())? - .ok_or("account not hosted")? - .public_key, - ) - .map_err(|error| error.to_string())?; - self.build_protected( - frame, - kind, - iota_id, - user_id as u64, - content, - recipient, - signature, - encryption, - &keyring, - ) - } - - fn build_protected( - &self, - frame: &CommunicationValue, - kind: CommunicationType, - signer_id: u64, - recipient_id: u64, - content: DataValue, - recipient: PublicKeyBundle, - signature: u8, - encryption: u8, - keyring: &mtp::crypto::Keyring, - ) -> Result { - let signer = DualSigner::new( - &keyring.sig_cl_secret_key, - &keyring.sig_pq_secret_key, - &keyring.sig_pq_public_key, - ) - .map_err(|error| error.to_string())?; - let mut rng = OsRng; - let id = ((rng.next_u64() as u128) << 64) | rng.next_u64() as u128; - ProtectedMessageBuilder::new( - kind, - content, - signer_id, - recipient_id, - &signer, - purpose(signature), - purpose(encryption), - ) - .message_id(id) - .created_at(now_millis()) - .recipients(vec![recipient]) - .frame_id(frame.id().ok_or("missing request ID")?) - .build() - .map_err(|error| error.to_string()) - } - - async fn password_error(&self, frame: &CommunicationValue, kind: CommunicationType) { - let mut response = CommunicationValue::new(kind); - if let Some(id) = frame.id() { - response = response.with_id(id); - } - let _ = self.send_message(&response).await; - } - - pub(super) async fn handle_password_enrollment_start( - self: Arc, - frame: &CommunicationValue, - ) { - let result = self.enrollment_start(frame).await; - match result { - Ok(response) => { - let _ = self.send_message(&response).await; - } - Err(error) => { - iota_logger::log!("Password enrollment start failed: {error}"); - self.password_error(frame, CommunicationType::ErrorInvalidData) - .await; - } - } - } - - async fn enrollment_start( - &self, - frame: &CommunicationValue, - ) -> Result { - let opened = self - .open_password_command(frame, CommunicationType::PasswordEnrollmentStart) - .await?; - let user_id = i64::try_from(opened.signer_id()).map_err(|error| error.to_string())?; - if self.password_auth.enrollments.len() >= self.password_auth.max_pending { - return Err("too many enrollments".into()); - } - let principal = self.password_principal(user_id).await?; - let request = bytes(opened.content(), DataType::OpaqueMessage, MAX_OPAQUE_BYTES)?; - let setup = self - .password_auth - .setup - .get() - .ok_or("OPAQUE setup unavailable")?; - let response = password::registration_start(setup, &request, principal.as_bytes()) - .map_err(|error| error.to_string())?; - let enrollment_id = Uuid::new_v4(); - self.password_auth - .enrollments - .insert(enrollment_id, (user_id, Instant::now())); - self.protected_response( - frame, - CommunicationType::PasswordEnrollmentResponse, - user_id, - DataValue::Container(vec![ - typed(DataType::Uuid, DataValue::Str(enrollment_id.to_string()))?, - typed(DataType::OpaqueMessage, DataValue::Bytes(response))?, - ]), - PASSWORD_RESPONSE_SIGNATURE, - PASSWORD_RESPONSE_ENCRYPTION, - ) - .await - } - - pub(super) async fn handle_password_enrollment_finish( - self: Arc, - frame: &CommunicationValue, - ) { - match self.enrollment_finish(frame).await { - Ok(response) => { - let _ = self.send_message(&response).await; - } - Err(error) => { - iota_logger::log!("Password enrollment finish failed: {error}"); - self.password_error(frame, CommunicationType::ErrorInvalidData) - .await; - } - } - } - - async fn enrollment_finish( - &self, - frame: &CommunicationValue, - ) -> Result { - let opened = self - .open_password_command(frame, CommunicationType::PasswordEnrollmentFinish) - .await?; - let user_id = i64::try_from(opened.signer_id()).map_err(|error| error.to_string())?; - let enrollment_id = session(opened.content())?; - let (_, (pending_user, started)) = self - .password_auth - .enrollments - .remove(&enrollment_id) - .ok_or("enrollment expired")?; - if pending_user != user_id || started.elapsed() >= self.password_auth.pending_ttl { - return Err("enrollment mismatch".into()); - } - let upload = bytes(opened.content(), DataType::OpaqueMessage, MAX_OPAQUE_BYTES)?; - let encrypted = bytes( - opened.content(), - DataType::EncryptedCredential, - MAX_CREDENTIAL_BYTES, - )?; - let version = positive(opened.content(), DataType::CredentialFormatVersion)?; - if version != 1 { - return Err("unsupported credential format".into()); - } - let opaque_record = - password::registration_finish(&upload).map_err(|error| error.to_string())?; - let profile = user_manager::get_user(user_id) - .map_err(|error| error.to_string())? - .ok_or("account not hosted")?; - let fingerprint = key_fingerprint(&profile.public_key)?; - let now = now_millis() as i64; - password_credentials::upsert(&PasswordCredential { - user_id, - protocol_version: 1, - credential_format_version: version, - opaque_record, - encrypted_tu_credential: encrypted, - account_public_key_sha256: fingerprint, - created_at: now, - updated_at: now, - }) - .map_err(|error| error.to_string())?; - self.password_auth - .logins - .retain(|_, pending| pending.user_id != user_id); - self.protected_response( - frame, - CommunicationType::PasswordEnrollmentStatus, - user_id, - DataValue::Container(vec![typed(DataType::Enabled, DataValue::Bool(true))?]), - PASSWORD_RESPONSE_SIGNATURE, - PASSWORD_RESPONSE_ENCRYPTION, - ) - .await - } - - pub(super) async fn handle_password_enrollment_status( - self: Arc, - frame: &CommunicationValue, - ) { - let result = async { - let opened = self - .open_password_command(frame, CommunicationType::PasswordEnrollmentStatus) - .await?; - let user_id = i64::try_from(opened.signer_id()).map_err(|error| error.to_string())?; - let enabled = - password_credentials::exists(user_id).map_err(|error| error.to_string())?; - self.protected_response( - frame, - CommunicationType::PasswordEnrollmentStatus, - user_id, - DataValue::Container(vec![typed(DataType::Enabled, DataValue::Bool(enabled))?]), - PASSWORD_RESPONSE_SIGNATURE, - PASSWORD_RESPONSE_ENCRYPTION, - ) - .await - } - .await; - match result { - Ok(response) => { - let _ = self.send_message(&response).await; - } - Err(error) => { - iota_logger::log!("Password status failed: {error}"); - self.password_error(frame, CommunicationType::ErrorInvalidData) - .await; - } - } - } - - pub(super) async fn handle_password_enrollment_disable( - self: Arc, - frame: &CommunicationValue, - ) { - let result = async { - let opened = self - .open_password_command(frame, CommunicationType::PasswordEnrollmentDisable) - .await?; - let user_id = i64::try_from(opened.signer_id()).map_err(|error| error.to_string())?; - password_credentials::delete(user_id).map_err(|error| error.to_string())?; - self.password_auth - .logins - .retain(|_, pending| pending.user_id != user_id); - self.protected_response( - frame, - CommunicationType::PasswordEnrollmentStatus, - user_id, - DataValue::Container(vec![typed(DataType::Enabled, DataValue::Bool(false))?]), - PASSWORD_RESPONSE_SIGNATURE, - PASSWORD_RESPONSE_ENCRYPTION, - ) - .await - } - .await; - match result { - Ok(response) => { - let _ = self.send_message(&response).await; - } - Err(error) => { - iota_logger::log!("Password disable failed: {error}"); - self.password_error(frame, CommunicationType::ErrorInvalidData) - .await; - } - } - } - - pub(super) async fn handle_password_provisioning_start( - self: Arc, - frame: &CommunicationValue, - ) { - match self.provisioning_start(frame).await { - Ok(response) => { - let _ = self.send_message(&response).await; - } - Err(error) => { - iota_logger::log!("Password KE1 failed: {error}"); - self.password_error(frame, CommunicationType::ErrorNotAuthenticated) - .await; - } - } - } - - async fn provisioning_start( - &self, - frame: &CommunicationValue, - ) -> Result { - let content = frame.payload(); - let id = session(content)?; - let participant_id = u64::try_from(positive(content, DataType::ProvisioningParticipantId)?) - .map_err(|error| error.to_string())?; - let user_id = positive(content, DataType::UserId)?; - let iota_id = positive(content, DataType::IotaId)?; - if Some(iota_id as u64) != iota_storage::util::config_util::CONFIG.load().iota_id { - return Err("wrong Iota".into()); - } - if self.password_auth.logins.len() >= self.password_auth.max_pending - || !self.password_auth.allow_attempt(user_id) - { - return Err("password attempt limit reached".into()); - } - let contact_key = public_key(content)?; - let ke1 = bytes(content, DataType::OpaqueMessage, MAX_OPAQUE_BYTES)?; - let principal = self.password_principal(user_id).await?; - let context = password::login_context(&principal, iota_id, id, &contact_key) - .map_err(|error| error.to_string())?; - let credential = password_credentials::get(user_id).map_err(|error| error.to_string())?; - let setup = self - .password_auth - .setup - .get() - .ok_or("OPAQUE setup unavailable")?; - let response = password::login_start( - setup, - credential - .as_ref() - .map(|value| value.opaque_record.as_slice()), - &ke1, - principal.as_bytes(), - &password::server_identifier(iota_id), - &context, - ) - .map_err(|error| error.to_string())?; - let pending = PendingLogin { - user_id, - participant_id, - contact_key, - context, - state: response.state, - real_record: credential.is_some(), - created: Instant::now(), - record_hash: credential - .as_ref() - .map(|value| Sha256::digest(&value.opaque_record).to_vec()), - }; - match self.password_auth.logins.entry(id) { - Entry::Vacant(slot) => { - slot.insert(pending); - } - Entry::Occupied(_) => return Err("session already started".into()), - } - Ok( - CommunicationValue::new(CommunicationType::PasswordProvisioningResponse) - .with_id(frame.id().ok_or("missing request ID")?) - .add_typed_default(DataType::Uuid, DataValue::Str(id.to_string())) - .add_typed_default(DataType::OpaqueMessage, DataValue::Bytes(response.response)), - ) - } - - pub(super) async fn handle_password_provisioning_finish( - self: Arc, - frame: &CommunicationValue, - ) { - match self.provisioning_finish(frame).await { - Ok(response) => { - let _ = self.send_message(&response).await; - } - Err(error) => { - iota_logger::log!("Password KE3 failed: {error}"); - self.password_error(frame, CommunicationType::ErrorNotAuthenticated) - .await; - } - } - } - - async fn provisioning_finish( - &self, - frame: &CommunicationValue, - ) -> Result { - let content = frame.payload(); - let id = session(content)?; - let (_, pending) = self - .password_auth - .logins - .remove(&id) - .ok_or("OPAQUE session expired")?; - if pending.created.elapsed() >= self.password_auth.pending_ttl - || pending.user_id != positive(content, DataType::UserId)? - || pending.participant_id - != u64::try_from(positive(content, DataType::ProvisioningParticipantId)?) - .map_err(|error| error.to_string())? - || pending - .contact_key - .try_as_bytes() - .map_err(|error| error.to_string())? - != public_key(content)? - .try_as_bytes() - .map_err(|error| error.to_string())? - { - return Err("OPAQUE session binding mismatch".into()); - } - let iota_id = positive(content, DataType::IotaId)?; - if Some(iota_id as u64) != iota_storage::util::config_util::CONFIG.load().iota_id { - return Err("wrong Iota".into()); - } - let ke3 = bytes(content, DataType::OpaqueMessage, MAX_OPAQUE_BYTES)?; - let principal = self.password_principal(pending.user_id).await?; - password::login_finish( - &pending.state, - &ke3, - principal.as_bytes(), - &password::server_identifier(iota_id), - &pending.context, - ) - .map_err(|error| error.to_string())?; - if !pending.real_record { - return Err("authentication failed".into()); - } - let credential = password_credentials::get(pending.user_id) - .map_err(|error| error.to_string())? - .ok_or("credential unavailable")?; - if pending.record_hash.as_deref() - != Some(Sha256::digest(&credential.opaque_record).as_slice()) - { - return Err("password credential changed during login".into()); - } - let profile = user_manager::get_user(pending.user_id) - .map_err(|error| error.to_string())? - .ok_or("account unavailable")?; - if credential.account_public_key_sha256 != key_fingerprint(&profile.public_key)? - || credential.encrypted_tu_credential.is_empty() - || credential.encrypted_tu_credential.len() > MAX_CREDENTIAL_BYTES - || credential.credential_format_version != 1 - { - return Err("credential stale".into()); - } - let content = DataValue::Container(vec![ - typed(DataType::Uuid, DataValue::Str(id.to_string()))?, - typed( - DataType::ProvisioningMethod, - DataValue::Str("password".into()), - )?, - typed( - DataType::UserId, - DataValue::SignedNumber(pending.user_id.into()), - )?, - typed(DataType::IotaId, DataValue::SignedNumber(iota_id.into()))?, - typed( - DataType::EncryptedCredential, - DataValue::Bytes(credential.encrypted_tu_credential), - )?, - typed( - DataType::CredentialFormatVersion, - DataValue::SignedNumber(credential.credential_format_version.into()), - )?, - typed( - DataType::Sha256, - DataValue::Bytes(credential.account_public_key_sha256), - )?, - ]); - let keyring = self - .keyring - .read() - .await - .clone() - .ok_or("Iota keyring unavailable")?; - self.build_protected( - frame, - CommunicationType::ProvisioningCredential, - iota_id as u64, - pending.participant_id, - content, - pending.contact_key, - PROVISIONING_CREDENTIAL_SIGNATURE, - PROVISIONING_CREDENTIAL_ENCRYPTION, - &keyring, - ) - } -} - -#[cfg(test)] -mod tests { - use super::*; - use mtp::{codec::InMemoryReplayGuard, crypto::Keyring}; - - #[test] - fn password_management_requires_signed_encrypted_single_use_frames() { - let account = Keyring::generate(); - let iota = Keyring::generate(); - let signer = DualSigner::new( - &account.sig_cl_secret_key, - &account.sig_pq_secret_key, - &account.sig_pq_public_key, - ) - .unwrap(); - let frame = ProtectedMessageBuilder::new( - CommunicationType::PasswordEnrollmentFinish, - DataValue::Container(vec![ - typed( - DataType::EncryptedCredential, - DataValue::Bytes(vec![1, 2, 3]), - ) - .unwrap(), - ]), - 7, - 11, - &signer, - purpose(PASSWORD_MANAGEMENT_SIGNATURE), - purpose(PASSWORD_MANAGEMENT_ENCRYPTION), - ) - .message_id(123_u128) - .created_at(now_millis()) - .recipients(vec![iota.public_key_bundle()]) - .frame_id(42) - .build() - .unwrap(); - let options = ProtectedOpenOptions::new( - Some(11), - purpose(PASSWORD_MANAGEMENT_SIGNATURE), - purpose(PASSWORD_MANAGEMENT_ENCRYPTION), - ProtectionPolicy::dual(), - ); - let mut guard = InMemoryReplayGuard::default(); - let resolve = |_| Some(vec![account.public_key_bundle()]); - let opened = - open_protected_with_checked(&frame, &[&iota], Some(7), resolve, options, &mut guard) - .unwrap(); - assert_eq!(opened.signer_id(), 7); - assert!( - open_protected_with_checked(&frame, &[&iota], Some(7), resolve, options, &mut guard) - .is_err() - ); - assert!( - open_protected_with_checked( - &frame, - &[&iota], - Some(7), - resolve, - ProtectedOpenOptions::new( - Some(12), - purpose(PASSWORD_MANAGEMENT_SIGNATURE), - purpose(PASSWORD_MANAGEMENT_ENCRYPTION), - ProtectionPolicy::dual() - ), - &mut InMemoryReplayGuard::default() - ) - .is_err() - ); - let clear = CommunicationValue::new(CommunicationType::PasswordEnrollmentFinish) - .with_receiver(11) - .with_id(42); - assert!( - open_protected_with_checked( - &clear, - &[&iota], - Some(7), - resolve, - options, - &mut InMemoryReplayGuard::default() - ) - .is_err() - ); - } -}