Update OPAQUE credential profile handling

This commit is contained in:
Alex-Emmet 2026-09-29 09:23:00 +02:00
commit 8ad56b938d
17 changed files with 2113 additions and 973 deletions

4
Cargo.lock generated
View file

@ -2232,6 +2232,7 @@ dependencies = [
"argon2", "argon2",
"async-trait", "async-trait",
"dashmap", "dashmap",
"generic-array",
"iota-identity", "iota-identity",
"mtp", "mtp",
"opaque-ke", "opaque-ke",
@ -2240,6 +2241,7 @@ dependencies = [
"thiserror 2.0.20", "thiserror 2.0.20",
"tokio", "tokio",
"uuid", "uuid",
"zeroize",
] ]
[[package]] [[package]]
@ -3297,6 +3299,7 @@ dependencies = [
"iota-util", "iota-util",
"json", "json",
"mtp", "mtp",
"opaque-ke",
"rand_core 0.6.4", "rand_core 0.6.4",
"reqwest", "reqwest",
"serde", "serde",
@ -3308,6 +3311,7 @@ dependencies = [
"trust-dns-resolver", "trust-dns-resolver",
"url", "url",
"uuid", "uuid",
"zeroize",
] ]
[[package]] [[package]]

View file

@ -11,6 +11,8 @@ mtp = { git = "https://git.methanium.net/Methanium/mtp.git", rev = "bb0f682b735d
rand_core = { version = "0.6", features = ["getrandom", "std"] } rand_core = { version = "0.6", features = ["getrandom", "std"] }
opaque-ke = { version = "4.0.1", features = ["argon2"] } opaque-ke = { version = "4.0.1", features = ["argon2"] }
argon2 = "0.5" argon2 = "0.5"
generic-array = "0.14"
zeroize = "1"
sha2 = "0.10" sha2 = "0.10"
thiserror = "2" thiserror = "2"
uuid = { version = "*", features = ["v4"] } uuid = { version = "*", features = ["v4"] }

View file

@ -1,4 +1,6 @@
use generic_array::{ArrayLength, GenericArray};
use mtp::crypto::PublicKeyBundle; use mtp::crypto::PublicKeyBundle;
use opaque_ke::ksf::Ksf;
use opaque_ke::{ use opaque_ke::{
CipherSuite, CredentialFinalization, CredentialRequest, Identifiers, RegistrationRequest, CipherSuite, CredentialFinalization, CredentialRequest, Identifiers, RegistrationRequest,
RegistrationUpload, Ristretto255, ServerLogin, ServerLoginParameters, ServerRegistration, RegistrationUpload, Ristretto255, ServerLogin, ServerLoginParameters, ServerRegistration,
@ -6,13 +8,78 @@ use opaque_ke::{
}; };
use rand_core::OsRng; use rand_core::OsRng;
use sha2::{Digest, Sha256}; 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<u8> {
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<Vec<u8>, 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<L: ArrayLength<u8>>(
&self,
input: GenericArray<u8, L>,
) -> Result<GenericArray<u8, L>, opaque_ke::errors::InternalError> {
let mut output = GenericArray::<u8, L>::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 OprfCs = Ristretto255;
type KeyExchange = TripleDh<Ristretto255, sha2::Sha512>; type KeyExchange = TripleDh<Ristretto255, sha2::Sha512>;
type Ksf = argon2::Argon2<'static>; type Ksf = PasswordKsfV1;
} }
pub type PasswordServerSetup = ServerSetup<TensaminOpaque>; pub type PasswordServerSetup = ServerSetup<TensaminOpaque>;
@ -39,6 +106,13 @@ pub fn deserialize_server_setup(bytes: &[u8]) -> Result<PasswordServerSetup, Pas
ServerSetup::deserialize(bytes).map_err(|error| PasswordAuthError::Opaque(error.to_string())) ServerSetup::deserialize(bytes).map_err(|error| PasswordAuthError::Opaque(error.to_string()))
} }
pub fn deserialize_server_setup_owned(
bytes: Vec<u8>,
) -> Result<PasswordServerSetup, PasswordAuthError> {
let bytes = Zeroizing::new(bytes);
deserialize_server_setup(&bytes)
}
pub fn registration_start( pub fn registration_start(
setup: &PasswordServerSetup, setup: &PasswordServerSetup,
request: &[u8], request: &[u8],
@ -59,10 +133,14 @@ pub fn registration_finish(upload: &[u8]) -> Result<Vec<u8>, PasswordAuthError>
.to_vec()) .to_vec())
} }
pub fn validate_login_request(request: &[u8]) -> Result<(), PasswordAuthError> {
CredentialRequest::<TensaminOpaque>::deserialize(request)
.map(|_| ())
.map_err(|error| PasswordAuthError::Opaque(error.to_string()))
}
pub fn server_identifier(iota_id: i64) -> Vec<u8> { pub fn server_identifier(iota_id: i64) -> Vec<u8> {
let mut result = b"tensamin:iota-password:v1\0".to_vec(); OpaqueProfileV1::server_identifier(iota_id)
result.extend_from_slice(&iota_id.to_be_bytes());
result
} }
pub fn login_context( pub fn login_context(
@ -71,25 +149,16 @@ pub fn login_context(
session_id: uuid::Uuid, session_id: uuid::Uuid,
contact_key: &PublicKeyBundle, contact_key: &PublicKeyBundle,
) -> Result<Vec<u8>, PasswordAuthError> { ) -> Result<Vec<u8>, PasswordAuthError> {
let mut result = b"tensamin:password-provisioning:v1\0".to_vec(); OpaqueProfileV1::login_context(principal, iota_id, session_id, contact_key)
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 struct LoginStartResult { pub struct LoginStartResult {
pub response: Vec<u8>, pub response: Vec<u8>,
pub state: Vec<u8>, pub state: SerializedLoginState,
} }
pub type SerializedLoginState = Zeroizing<Vec<u8>>;
pub fn login_start( pub fn login_start(
setup: &PasswordServerSetup, setup: &PasswordServerSetup,
record: Option<&[u8]>, record: Option<&[u8]>,
@ -121,21 +190,29 @@ pub fn login_start(
.map_err(|error| PasswordAuthError::Opaque(error.to_string()))?; .map_err(|error| PasswordAuthError::Opaque(error.to_string()))?;
Ok(LoginStartResult { Ok(LoginStartResult {
response: result.message.serialize().to_vec(), 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( pub fn login_finish(
serialized_state: &[u8], serialized_state: &[u8],
finalization: &[u8], finalization: &[u8],
principal: &[u8], principal: &[u8],
server_id: &[u8], server_id: &[u8],
context: &[u8], context: &[u8],
) -> Result<(), PasswordAuthError> { ) -> Result<(), LoginFinishError> {
let state = ServerLogin::<TensaminOpaque>::deserialize(serialized_state) let state = ServerLogin::<TensaminOpaque>::deserialize(serialized_state)
.map_err(|error| PasswordAuthError::Opaque(error.to_string()))?; .map_err(|_| LoginFinishError::InvalidMessage)?;
let message = CredentialFinalization::<TensaminOpaque>::deserialize(finalization) let message = CredentialFinalization::<TensaminOpaque>::deserialize(finalization)
.map_err(|error| PasswordAuthError::Opaque(error.to_string()))?; .map_err(|_| LoginFinishError::InvalidMessage)?;
state state
.finish( .finish(
message, message,
@ -147,7 +224,7 @@ pub fn login_finish(
}, },
}, },
) )
.map_err(|error| PasswordAuthError::Opaque(error.to_string()))?; .map_err(|_| LoginFinishError::AuthenticationFailed)?;
Ok(()) Ok(())
} }
@ -161,6 +238,46 @@ mod tests {
}; };
use uuid::Uuid; 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::<TensaminOpaque>::start(&mut OsRng, b"password").unwrap();
assert!(validate_login_request(&request.message.serialize()).is_ok());
}
#[test] #[test]
fn registration_and_login_bind_context_identifiers_and_persisted_setup() { fn registration_and_login_bind_context_identifiers_and_persisted_setup() {
let setup = generate_server_setup(); let setup = generate_server_setup();
@ -186,7 +303,8 @@ mod tests {
.unwrap(); .unwrap();
let record = registration_finish(&upload.message.serialize()).unwrap(); let record = registration_finish(&upload.message.serialize()).unwrap();
let key = Keyring::generate().public_key_bundle(); 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::<TensaminOpaque>::start(&mut OsRng, b"correct horse").unwrap(); let client = ClientLogin::<TensaminOpaque>::start(&mut OsRng, b"correct horse").unwrap();
let result = login_start( let result = login_start(
@ -218,14 +336,39 @@ mod tests {
) )
.is_ok() .is_ok()
); );
let client = ClientLogin::<TensaminOpaque>::start(&mut OsRng, b"correct horse").unwrap(); 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(),
)
.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::<TensaminOpaque>::start(&mut OsRng, b"correct horse").unwrap();
let mismatched = login_start( let mismatched = login_start(
&restored, &restored,
Some(&record), Some(&record),
&client.message.serialize(), &client.message.serialize(),
principal, principal,
&server, &server,
b"wrong context", &changed_context,
) )
.unwrap(); .unwrap();
assert!( assert!(
@ -236,10 +379,12 @@ mod tests {
b"correct horse", b"correct horse",
CredentialResponse::<TensaminOpaque>::deserialize(&mismatched.response) CredentialResponse::<TensaminOpaque>::deserialize(&mismatched.response)
.unwrap(), .unwrap(),
ClientLoginFinishParameters::new(Some(&context), identifiers, None) ClientLoginFinishParameters::new(Some(&context), identifiers, None),
) )
.is_err() .is_err(),
"login accepted a different {binding}"
); );
}
let wrong_server = server_identifier(12); let wrong_server = server_identifier(12);
for (login_principal, login_server, password) in [ for (login_principal, login_server, password) in [

View file

@ -5,8 +5,7 @@ use crate::{storage_error::StorageError, util::db};
#[derive(Clone, Debug, PartialEq, Eq)] #[derive(Clone, Debug, PartialEq, Eq)]
pub struct PasswordCredential { pub struct PasswordCredential {
pub user_id: i64, pub user_id: i64,
pub protocol_version: i64, pub opaque_profile: i64,
pub credential_format_version: i64,
pub opaque_record: Vec<u8>, pub opaque_record: Vec<u8>,
pub encrypted_tu_credential: Vec<u8>, pub encrypted_tu_credential: Vec<u8>,
pub account_public_key_sha256: Vec<u8>, pub account_public_key_sha256: Vec<u8>,
@ -17,12 +16,11 @@ pub struct PasswordCredential {
pub fn get(user_id: i64) -> Result<Option<PasswordCredential>, StorageError> { pub fn get(user_id: i64) -> Result<Option<PasswordCredential>, StorageError> {
db::with_db(|conn| { db::with_db(|conn| {
conn.query_row( 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], params![user_id],
|row| Ok(PasswordCredential { |row| Ok(PasswordCredential {
user_id: row.get(0)?, protocol_version: row.get(1)?, credential_format_version: row.get(2)?, user_id: row.get(0)?, opaque_profile: row.get(1)?, opaque_record: row.get(2)?, encrypted_tu_credential: row.get(3)?,
opaque_record: row.get(3)?, encrypted_tu_credential: row.get(4)?, account_public_key_sha256: row.get(4)?, created_at: row.get(5)?, updated_at: row.get(6)?,
account_public_key_sha256: row.get(5)?, created_at: row.get(6)?, updated_at: row.get(7)?,
}), }),
).optional().map_err(Into::into) ).optional().map_err(Into::into)
}) })
@ -53,17 +51,15 @@ pub fn any() -> Result<bool, StorageError> {
pub fn upsert(credential: &PasswordCredential) -> Result<(), StorageError> { pub fn upsert(credential: &PasswordCredential) -> Result<(), StorageError> {
db::with_immediate_transaction(|tx| { db::with_immediate_transaction(|tx| {
tx.execute( 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) "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, ?8) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
ON CONFLICT(user_id) DO UPDATE SET ON CONFLICT(user_id) DO UPDATE SET
protocol_version = excluded.protocol_version, opaque_profile = excluded.opaque_profile,
credential_format_version = excluded.credential_format_version,
opaque_record = excluded.opaque_record, opaque_record = excluded.opaque_record,
encrypted_tu_credential = excluded.encrypted_tu_credential, encrypted_tu_credential = excluded.encrypted_tu_credential,
account_public_key_sha256 = excluded.account_public_key_sha256, account_public_key_sha256 = excluded.account_public_key_sha256,
updated_at = excluded.updated_at", updated_at = excluded.updated_at",
params![credential.user_id, credential.protocol_version, credential.credential_format_version, params![credential.user_id, credential.opaque_profile, credential.opaque_record, credential.encrypted_tu_credential,
credential.opaque_record, credential.encrypted_tu_credential,
credential.account_public_key_sha256, credential.created_at, credential.updated_at], credential.account_public_key_sha256, credential.created_at, credential.updated_at],
)?; )?;
Ok(()) Ok(())

View file

@ -2027,8 +2027,6 @@ fn run_migrations_on_connection(conn: &Connection) -> Result<(), StorageError> {
r#" r#"
CREATE TABLE user_password_credentials ( CREATE TABLE user_password_credentials (
user_id INTEGER PRIMARY KEY, user_id INTEGER PRIMARY KEY,
protocol_version INTEGER NOT NULL,
credential_format_version INTEGER NOT NULL,
opaque_record BLOB NOT NULL, opaque_record BLOB NOT NULL,
encrypted_tu_credential BLOB NOT NULL, encrypted_tu_credential BLOB NOT NULL,
account_public_key_sha256 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(()) Ok(())
} }
@ -2122,7 +2163,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))?; 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"] { for column in ["height", "reply_to", "edited_count", "deleted_by_external"] {
let mut statement = let mut statement =
conn.prepare("SELECT 1 FROM pragma_table_info('messages') WHERE name = ?1")?; 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)?;
run_migrations_on_connection(&conn)?; run_migrations_on_connection(&conn)?;
let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; 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 [ for table in [
"sync_heads", "sync_heads",
"sync_events", "sync_events",
@ -2185,6 +2226,94 @@ mod tests {
conn.prepare("SELECT 1 FROM pragma_table_info('pending_relays') WHERE name = ?1")?; conn.prepare("SELECT 1 FROM pragma_table_info('pending_relays') WHERE name = ?1")?;
assert!(statement.exists([column])?); assert!(statement.exists([column])?);
} }
let password_columns: Vec<String> = conn
.prepare("SELECT name FROM pragma_table_info('user_password_credentials')")?
.query_map([], |row| row.get(0))?
.collect::<Result<_, _>>()?;
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<u8>, Vec<u8>, Vec<u8>, 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<u8>, Vec<u8>, Vec<u8>, 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(()) Ok(())
} }
@ -2224,7 +2353,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))?; 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 [ for column in [
"id", "id",
"user_id", "user_id",
@ -2319,7 +2448,7 @@ mod tests {
)?; )?;
let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?;
assert_eq!(preserved, "remote_committed"); assert_eq!(preserved, "remote_committed");
assert_eq!(version, 47); assert_eq!(version, 49);
Ok(()) Ok(())
} }
@ -2357,7 +2486,7 @@ mod tests {
})?; })?;
assert_eq!(count, 0); assert_eq!(count, 0);
let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?;
assert_eq!(version, 47); assert_eq!(version, 49);
Ok(()) Ok(())
} }

@ -1 +1 @@
Subproject commit 4f18c7a0d9b04d38a77fbb011c4f0b21c25bf7bf Subproject commit 4f81a27820982b54da79213d8fcae358728e983a

View file

@ -34,3 +34,7 @@ base64 = "0.22.1"
rand_core = { version = "0.6", features = ["getrandom", "std"] } rand_core = { version = "0.6", features = ["getrandom", "std"] }
trust-dns-resolver = { version = "0.23", features = ["tokio-runtime"] } trust-dns-resolver = { version = "0.23", features = ["tokio-runtime"] }
url = "2" url = "2"
zeroize = "1"
[dev-dependencies]
opaque-ke = { version = "4.0.1", features = ["argon2"] }

View file

@ -2,7 +2,7 @@ pub mod client;
pub mod identity; pub mod identity;
pub mod omega_discovery; pub mod omega_discovery;
pub mod omikron_connection; pub mod omikron_connection;
mod password_provisioning; mod password;
pub mod router; pub mod router;
pub mod tauth; pub mod tauth;
pub mod user_ops; pub mod user_ops;

View file

@ -24,7 +24,7 @@ use uuid::Uuid;
use crate::client::{OmikronClient, OmikronError}; use crate::client::{OmikronClient, OmikronError};
use crate::omega_discovery; use crate::omega_discovery;
use crate::password_provisioning::PasswordAuthRuntime; use crate::password::PasswordAuthRuntime;
use iota_connection::message_common::*; use iota_connection::message_common::*;
use iota_connection::message_handlers; 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<Self>) -> Result<ConnectionAttemptResult, String> { async fn connect_once(self: Arc<Self>) -> Result<ConnectionAttemptResult, String> {
self.set_state(ConnectionState::Connecting).await; self.set_state(ConnectionState::Connecting).await;
log_t!("omikron_connecting"); log_t!("omikron_connecting");
@ -640,8 +646,8 @@ impl OmikronConnection {
.await .await
.map_err(|error| format!("Iota identity initialization failed: {error}"))?; .map_err(|error| format!("Iota identity initialization failed: {error}"))?;
let keyring = identity.keyring(); let keyring = identity.keyring();
self.password_auth.initialize(identity_path())?;
*self.keyring.write().await = Some(keyring.clone()); *self.keyring.write().await = Some(keyring.clone());
self.initialize_password_auth(identity_path());
let existing_iota_id = CONFIG.load().iota_id; let existing_iota_id = CONFIG.load().iota_id;
@ -698,6 +704,7 @@ impl OmikronConnection {
*self.sender.write().await = Some(sender_arc.clone()); *self.sender.write().await = Some(sender_arc.clone());
self.set_state(ConnectionState::Connected { identified: true }) self.set_state(ConnectionState::Connected { identified: true })
.await; .await;
self.password_auth.prune();
// Start read loop // Start read loop
let connection = Arc::new(connection); let connection = Arc::new(connection);
@ -3672,6 +3679,32 @@ mod tests {
use super::*; use super::*;
use std::fs; 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 { fn test_path(name: &str) -> PathBuf {
std::env::temp_dir().join(format!( std::env::temp_dir().join(format!(
"iota-identity-{name}-{}-{}", "iota-identity-{name}-{}-{}",

View file

@ -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
)
}
}

View file

@ -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<Self>,
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<CommunicationValue, String> {
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<Self>,
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<CommunicationValue, String> {
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<Self>,
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<Self>,
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;
}
}
}
}

View file

@ -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<String, String> {
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<PasswordAuthRuntime> {
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::<Vec<_>>();
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::<password::TensaminOpaque>::start(&mut OsRng, b"correct password")
.unwrap();
let response =
password::registration_start(&setup, &registration.message.serialize(), principal)
.unwrap();
let upload = registration
.state
.finish(
&mut OsRng,
b"correct password",
RegistrationResponse::<password::TensaminOpaque>::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::<password::TensaminOpaque>::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::<password::TensaminOpaque>::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::<password::TensaminOpaque>::start(&mut OsRng, b"correct password")
.unwrap();
let response =
password::registration_start(&setup, &registration.message.serialize(), principal)
.unwrap();
let upload = registration
.state
.finish(
&mut OsRng,
b"correct password",
RegistrationResponse::<password::TensaminOpaque>::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::<password::TensaminOpaque>::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::<password::TensaminOpaque>::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()
);
}
}

View file

@ -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<bool, ReplayError> {
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<VerifiedProtectedMessage, String> {
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<CommunicationValue, String> {
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<CommunicationValue, String> {
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;
}
}

View file

@ -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<u8>,
setup: Arc<PasswordServerSetup>,
attempt: AttemptLease,
delay: Duration,
}
impl OmikronConnection {
pub(crate) async fn handle_password_provisioning_start(
self: Arc<Self>,
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<PreparedPasswordStart, String> {
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<CommunicationValue, String> {
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<Self>,
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<CommunicationValue, String> {
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<CommunicationValue, PasswordFinishError> {
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::<password::TensaminOpaque>::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);
}
}

View file

@ -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<u8>,
pub(super) state: password::SerializedLoginState,
pub(super) real_record: bool,
pub(super) record_hash: Option<Vec<u8>>,
pub(super) created: Instant,
pub(super) attempt: Option<AttemptLease>,
}
#[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<PasswordAuthRuntime>,
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<Arc<PasswordServerSetup>>,
pub(super) omega_authority: OnceLock<String>,
pub(super) logins: DashMap<Uuid, PendingLogin>,
pub(super) enrollments: DashMap<Uuid, (i64, Instant)>,
pub(super) attempts: DashMap<i64, AccountAttemptState>,
pub(super) max_pending: usize,
exchange_slots: Arc<Semaphore>,
pub(super) pending_ttl: Duration,
pub(super) throttle: PasswordThrottleConfig,
}
fn configured_positive<T: std::str::FromStr + PartialOrd + Default>(name: &str, default: T) -> T {
std::env::var(name)
.ok()
.and_then(|value| value.parse::<T>().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<Arc<PasswordServerSetup>, 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::<Vec<_>>();
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<Self>,
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::<Vec<_>>();
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,
}

View file

@ -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<Vec<u8>, 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<i64, String> {
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, String> {
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<PublicKeyBundle, String> {
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<Vec<u8>, 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())
}

View file

@ -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<bool, ReplayError> {
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<u8>,
state: Vec<u8>,
real_record: bool,
record_hash: Option<Vec<u8>>,
created: Instant,
}
pub(super) struct PasswordAuthRuntime {
setup: OnceLock<Arc<PasswordServerSetup>>,
omega_authority: OnceLock<String>,
logins: DashMap<Uuid, PendingLogin>,
enrollments: DashMap<Uuid, (i64, Instant)>,
attempts: DashMap<i64, (Instant, u32)>,
max_pending: usize,
pending_ttl: Duration,
attempt_window: Duration,
max_attempts: u32,
}
fn configured_positive<T: std::str::FromStr + PartialOrd + Default>(name: &str, default: T) -> T {
std::env::var(name)
.ok()
.and_then(|value| value.parse::<T>().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<Vec<u8>, 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<i64, String> {
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, String> {
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<PublicKeyBundle, String> {
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<Vec<u8>, 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<String, String> {
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<VerifiedProtectedMessage, String> {
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<CommunicationValue, String> {
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<CommunicationValue, String> {
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<Self>,
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<CommunicationValue, String> {
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<Self>,
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<CommunicationValue, String> {
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<Self>,
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<Self>,
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<Self>,
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<CommunicationValue, String> {
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<Self>,
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<CommunicationValue, String> {
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()
);
}
}