[WIP] User Invites

This commit is contained in:
Alex Emmet 2026-09-10 23:14:25 +02:00
commit 7164f4671d
No known key found for this signature in database
24 changed files with 2829 additions and 170 deletions

View file

@ -8,6 +8,7 @@ pub enum OmikronError {
Timeout(String),
Authentication(String),
Rejected(CommunicationType, String),
Storage(String),
Internal(String),
}
@ -18,6 +19,7 @@ impl std::fmt::Display for OmikronError {
| Self::Timeout(v)
| Self::Authentication(v)
| Self::Rejected(_, v)
| Self::Storage(v)
| Self::Internal(v) => f.write_str(v),
}
}
@ -42,6 +44,12 @@ pub trait OmikronClient: Send + Sync {
value: &CommunicationValue,
timeout: Duration,
) -> Result<CommunicationValue, OmikronError>;
async fn sync_omega_invitations(&self) -> Result<(), OmikronError> {
Err(OmikronError::Internal(
"invitation synchronization is unavailable".into(),
))
}
async fn flush_pending_invitation_actions(&self) {}
async fn reconnect(&self) -> Result<(), OmikronError>;
/// Replace the local Iota identity and wait for the new identity to
/// register/authenticate. This is deliberately available while offline:

View file

@ -9,13 +9,13 @@ use iota_storage::util::{
use iota_util::crypto_helper::{self, keyring_from_base64};
use iota_util::crypto_util::{self};
use mtp::client::{Client, ClientConfig, MTPConnection, Policy, SendMode, Sender};
use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue};
use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataTypeId, DataValue, TypeMap};
use mtp::crypto::{Keyring, PublicKeyBundle};
use rand_core::RngCore;
use std::env;
use std::fs;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
use std::sync::{Arc, LazyLock};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use tokio::sync::{Mutex, RwLock, Semaphore, oneshot, watch};
@ -128,6 +128,119 @@ const RECONNECT_DELAY: Duration = Duration::from_secs(5);
const MAX_RECONNECT_DELAY: Duration = Duration::from_secs(300);
const CONNECTION_TIMEOUT: Duration = Duration::from_secs(45);
const MAINTENANCE_INTERVAL: Duration = Duration::from_secs(5);
fn container_value<'a>(
fields: &'a [(DataTypeId, DataValue)],
kind: DataType,
type_map: &TypeMap,
) -> Option<&'a DataValue> {
let id = kind.try_to_id(type_map)?;
fields
.iter()
.find_map(|(field_id, value)| (*field_id == id).then_some(value))
}
fn container_signed_i64(
fields: &[(DataTypeId, DataValue)],
kind: DataType,
type_map: &TypeMap,
) -> Option<i64> {
container_value(fields, kind, type_map)
.and_then(DataValue::as_signed_number)
.and_then(|value| i64::try_from(value).ok())
}
fn parse_omega_invitation(
fields: &[(DataTypeId, DataValue)],
type_map: &TypeMap,
expected_iota_id: i64,
synced_at: i64,
) -> Result<iota_storage::users::invitations::InvitationSummary, OmikronError> {
use iota_storage::users::invitations::{InvitationState, InvitationSummary};
let invitation_id = container_signed_i64(fields, DataType::InvitationId, type_map)
.filter(|value| *value > 0)
.ok_or_else(|| OmikronError::Internal("invalid Omega invitation ID".into()))?;
if container_value(fields, DataType::InvitationAuthority, type_map).and_then(DataValue::as_str)
!= Some("omega")
{
return Err(OmikronError::Internal(
"invalid Omega invitation authority".into(),
));
}
if container_signed_i64(fields, DataType::IotaId, type_map) != Some(expected_iota_id) {
return Err(OmikronError::Internal(
"Omega invitation snapshot contained the wrong Iota ID".into(),
));
}
let remote_revision = container_signed_i64(fields, DataType::InvitationRevision, type_map)
.filter(|value| *value > 0)
.ok_or_else(|| OmikronError::Internal("invalid Omega invitation revision".into()))?;
let state = match container_value(fields, DataType::InvitationState, type_map)
.and_then(DataValue::as_str)
{
Some("pending") => InvitationState::Pending,
Some("provisioning") => InvitationState::Provisioning,
Some("redeemed") => InvitationState::Redeemed,
Some("revoked") => InvitationState::Revoked,
Some("expired") => InvitationState::Expired,
_ => {
return Err(OmikronError::Internal(
"invalid Omega invitation state".into(),
));
}
};
let label = match container_value(fields, DataType::InvitationLabel, type_map) {
Some(value) => Some(
value
.as_str()
.ok_or_else(|| OmikronError::Internal("invalid Omega invitation label".into()))?,
),
None => None,
};
let password_protected =
container_value(fields, DataType::InvitationPasswordProtected, type_map)
.and_then(DataValue::as_bool)
.ok_or_else(|| {
OmikronError::Internal("invalid Omega invitation password protection flag".into())
})?;
let created_at = container_signed_i64(fields, DataType::InvitationCreatedAt, type_map)
.filter(|value| *value > 0)
.ok_or_else(|| OmikronError::Internal("invalid Omega invitation creation time".into()))?;
let expires_at = container_signed_i64(fields, DataType::InvitationExpiresAt, type_map)
.filter(|value| *value > 0)
.ok_or_else(|| OmikronError::Internal("invalid Omega invitation expiry time".into()))?;
let redeemed_user_id = match container_value(fields, DataType::UserId, type_map) {
Some(value) => Some(
value
.as_signed_number()
.and_then(|value| i64::try_from(value).ok())
.filter(|value| *value > 0)
.ok_or_else(|| {
OmikronError::Internal("invalid Omega invitation redeemed user ID".into())
})?,
),
None => None,
};
Ok(InvitationSummary {
invitation_id,
authority: iota_storage::users::invitations::InvitationAuthority::Omega,
label: label.map(str::to_owned),
password_protected,
created_at,
expires_at: Some(expires_at),
state,
remote_revision,
redeemed_user_id,
redeemed_at: None,
revoked_at: None,
pending_action: None,
pending_action_at: None,
last_synced_at: Some(synced_at),
local_provisioned_user_id: None,
local_provisioned_at: None,
})
}
const TASK_CLEANUP_INTERVAL: Duration = Duration::from_secs(60);
const TASK_MAX_AGE: Duration = Duration::from_secs(60);
const MAX_CONCURRENT_HANDLERS: usize = 20;
@ -140,6 +253,26 @@ struct ResolvedOmikronEndpoint {
public_key: PublicKeyBundle,
}
struct InvitationSyncRuntime {
flush_in_progress: AtomicBool,
next_retry_at: Mutex<Option<Instant>>,
retry_delay: Mutex<Duration>,
}
impl InvitationSyncRuntime {
fn new() -> Self {
Self {
flush_in_progress: AtomicBool::new(false),
next_retry_at: Mutex::new(None),
retry_delay: Mutex::new(Duration::from_secs(5)),
}
}
}
fn next_invitation_retry_delay(current: Duration) -> Duration {
(current * 2).min(Duration::from_secs(300))
}
#[derive(Debug)]
pub enum IdentityError {
Storage(mtp::files::FileError),
@ -320,6 +453,7 @@ pub struct OmikronConnection {
cancellation: CancellationToken,
pub(crate) active_tasks: Arc<DashSet<String>>,
pub(crate) app: Arc<std::sync::Mutex<AppState>>,
invitation_sync: Arc<InvitationSyncRuntime>,
}
impl OmikronConnection {
@ -353,6 +487,7 @@ impl OmikronConnection {
cancellation,
active_tasks,
app,
invitation_sync: Arc::new(InvitationSyncRuntime::new()),
}
}
@ -546,6 +681,13 @@ impl OmikronConnection {
log_t!("omikron_authenticated");
self.classify_legacy_pending_relays().await;
if let Err(error) = self.sync_omega_invitations().await {
log!(
"Omega invitation snapshot synchronization failed: {}",
error
);
}
self.flush_pending_invitation_actions().await;
let maintenance_self = self.clone();
let maintenance_handle = tokio::spawn(async move {
@ -585,6 +727,155 @@ impl OmikronConnection {
load_or_migrate_keyring_at(identity_path(), CONFIG.load().keyring.clone())
}
async fn sync_pending_invitation_revocations(&self) -> Result<(), String> {
let invitations = iota_storage::users::invitations::list()
.map_err(|error| format!("failed to load pending invitation actions: {error}"))?;
for invitation in invitations.into_iter().filter(|invitation| {
invitation.authority == iota_storage::users::invitations::InvitationAuthority::Omega
&& invitation.pending_action
== Some(iota_storage::users::invitations::PendingAction::Revoke)
}) {
let request = CommunicationValue::new(CommunicationType::RevokeUserInvitation)
.add_typed_default(
DataType::InvitationAuthority,
DataValue::Str("omega".into()),
)
.add_typed_default(
DataType::InvitationId,
DataValue::SignedNumber(invitation.invitation_id.into()),
);
let response = self
.await_response(&request, Some(Duration::from_secs(20)))
.await?;
if response.is_type(CommunicationType::Success) {
let Some(revision) = response
.get_data(DataType::InvitationRevision)
.as_signed_number()
.and_then(|value| i64::try_from(value).ok())
else {
return Err("Omega revoke response omitted invitation revision".into());
};
let state = match response.get_data(DataType::InvitationState).as_str() {
Some("pending") => iota_storage::users::invitations::InvitationState::Pending,
Some("provisioning") => {
iota_storage::users::invitations::InvitationState::Provisioning
}
Some("redeemed") => iota_storage::users::invitations::InvitationState::Redeemed,
Some("revoked") => iota_storage::users::invitations::InvitationState::Revoked,
Some("expired") => iota_storage::users::invitations::InvitationState::Expired,
_ => return Err("Omega revoke response contained invalid state".into()),
};
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as i64;
iota_storage::users::invitations::apply_revoke_result(
invitation.invitation_id,
state,
revision,
(state == iota_storage::users::invitations::InvitationState::Revoked)
.then_some(now),
now,
)
.map_err(|error| format!("failed to apply invitation revoke result: {error}"))?;
} else {
return Err("Omega returned the wrong invitation revoke response type".into());
}
}
Ok(())
}
pub async fn flush_pending_invitation_actions(&self) {
if self
.invitation_sync
.flush_in_progress
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return;
}
let now = Instant::now();
if self
.invitation_sync
.next_retry_at
.lock()
.await
.is_some_and(|retry_at| retry_at > now)
{
self.invitation_sync
.flush_in_progress
.store(false, Ordering::Release);
return;
}
match self.sync_pending_invitation_revocations().await {
Ok(()) => {
*self.invitation_sync.next_retry_at.lock().await = None;
*self.invitation_sync.retry_delay.lock().await = Duration::from_secs(5);
}
Err(error) => {
log!("Pending invitation action flush failed: {}", error);
let mut delay = self.invitation_sync.retry_delay.lock().await;
*self.invitation_sync.next_retry_at.lock().await = Some(now + *delay);
*delay = next_invitation_retry_delay(*delay);
}
}
self.invitation_sync
.flush_in_progress
.store(false, Ordering::Release);
}
pub async fn sync_omega_invitations(&self) -> Result<(), OmikronError> {
let request = CommunicationValue::new(CommunicationType::ListUserInvitations)
.add_typed_default(
DataType::InvitationAuthority,
DataValue::Str("omega".into()),
);
let response = self
.await_response(&request, Some(Duration::from_secs(20)))
.await
.map_err(map_await_response_error)?;
if !response.is_type(CommunicationType::ListUserInvitations) {
return Err(OmikronError::Internal(
"Omega returned the wrong invitation snapshot response type".into(),
));
}
let Some(DataValue::Array(remote)) = response.get_data(DataType::Invitations) else {
return Err(OmikronError::Internal(
"Omega invitation snapshot omitted Invitations".into(),
));
};
let type_map = TypeMap::latest();
let synced_at = now_millis_i64();
let expected_iota_id = CONFIG
.load()
.iota_id
.and_then(|value| i64::try_from(value).ok())
.filter(|value| *value > 0)
.ok_or_else(|| OmikronError::Internal("Iota identity is not configured".into()))?;
let mut invitations = Vec::with_capacity(remote.len());
for invitation in remote {
let DataValue::Container(fields) = invitation else {
return Err(OmikronError::Internal(
"Omega invitation snapshot contained a non-container record".into(),
));
};
invitations.push(parse_omega_invitation(
fields,
&type_map,
expected_iota_id,
synced_at,
)?);
}
iota_storage::users::invitations::merge_omega_snapshot(&invitations, synced_at).map_err(
|error| {
OmikronError::Storage(format!(
"failed to persist Omega invitation snapshot: {error}"
))
},
)?;
Ok(())
}
// -------------------------------------------------------------------------
// Omikron discovery (via Omega's HTTP API, replacing the static
// host/port/public-key-file model)
@ -827,6 +1118,10 @@ impl OmikronConnection {
self.classify_legacy_pending_relays().await;
self.flush_pending_relays().await;
let invitation_self = self.clone();
tokio::spawn(async move {
invitation_self.flush_pending_invitation_actions().await;
});
if let Err(error) = relay_replay::prune_completed(
now_millis_i64().saturating_sub(RELAY_RETENTION_MILLIS),
) {
@ -2407,14 +2702,22 @@ impl OmikronConnection {
let public_key = cv.get_data(DataType::PublicKey).as_str().map(str::to_owned);
let invitation_id = cv
.get_data(DataType::InvitationId)
.as_str()
.filter(|value| uuid::Uuid::parse_str(value).is_ok())
.map(str::to_owned);
let Some((user_id, username, public_key, invitation_id)) = user_id
.as_signed_number()
.and_then(|value| i64::try_from(value).ok())
.filter(|value| *value > 0);
let invitation_revision = cv
.get_data(DataType::InvitationRevision)
.as_signed_number()
.and_then(|value| i64::try_from(value).ok())
.filter(|value| *value > 0);
let Some((user_id, username, public_key, invitation_id, invitation_revision)) = user_id
.zip(username)
.zip(public_key)
.zip(invitation_id)
.map(|(((id, username), key), invitation_id)| (id, username, key, invitation_id))
.zip(invitation_revision)
.map(|((((id, username), key), invitation_id), revision)| {
(id, username, key, invitation_id, revision)
})
else {
let _ = self
.send_message(&error_response(cv, CommunicationType::ErrorInvalidData))
@ -2424,22 +2727,65 @@ impl OmikronConnection {
let profile = iota_storage::users::user_profile::UserProfile::new(
user_id, username, None, public_key, None, None,
);
if iota_storage::users::user_manager::try_add_user_with_credential_origin(
profile,
iota_storage::users::user_manager::CredentialOrigin::External,
)
.is_err()
{
let _ = self
.send_message(&error_response(cv, CommunicationType::ErrorInternal))
.await;
return;
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as i64;
let mut result = iota_storage::users::invitations::apply_external_invitation_provisioning(
invitation_id,
invitation_revision,
&profile,
now,
);
if matches!(
result,
Ok(iota_storage::users::invitations::ProvisioningResult::MissingInvitation)
) {
if self.sync_omega_invitations().await.is_ok() {
result = iota_storage::users::invitations::apply_external_invitation_provisioning(
invitation_id,
invitation_revision,
&profile,
now,
);
}
}
match result {
Ok(
iota_storage::users::invitations::ProvisioningResult::Created
| iota_storage::users::invitations::ProvisioningResult::AlreadyApplied,
) => {}
Ok(iota_storage::users::invitations::ProvisioningResult::RevocationPending) => {
self.flush_pending_invitation_actions().await;
return;
}
Ok(iota_storage::users::invitations::ProvisioningResult::Conflict) => {
let _ = self
.send_message(&error_response(cv, CommunicationType::ErrorInvalidData))
.await;
return;
}
Ok(iota_storage::users::invitations::ProvisioningResult::MissingInvitation) => {
let _ = self
.send_message(&error_response(cv, CommunicationType::ErrorInvalidData))
.await;
return;
}
Err(_) => {
let _ = self
.send_message(&error_response(cv, CommunicationType::ErrorInternal))
.await;
return;
}
}
let acknowledgement =
CommunicationValue::new(CommunicationType::AcknowledgeIotaUserProvision)
.with_request_id(cv)
.add_typed_default(DataType::UserId, DataValue::SignedNumber(user_id.into()))
.add_typed_default(DataType::InvitationId, DataValue::Str(invitation_id));
.add_typed_default(
DataType::InvitationId,
DataValue::SignedNumber(invitation_id.into()),
);
let _ = self.send_message(&acknowledgement).await;
}
@ -3555,6 +3901,14 @@ impl OmikronClient for OmikronConnection {
.map_err(map_await_response_error)
}
async fn sync_omega_invitations(&self) -> Result<(), OmikronError> {
Self::sync_omega_invitations(self).await
}
async fn flush_pending_invitation_actions(&self) {
Self::flush_pending_invitation_actions(self).await;
}
async fn reconnect(&self) -> Result<(), OmikronError> {
let this = Arc::new(Self {
state: self.state.clone(),
@ -3574,6 +3928,7 @@ impl OmikronClient for OmikronConnection {
cancellation: self.cancellation.clone(),
active_tasks: self.active_tasks.clone(),
app: self.app.clone(),
invitation_sync: self.invitation_sync.clone(),
});
Self::reconnect(&this).await;
Ok(())
@ -3598,6 +3953,7 @@ impl OmikronClient for OmikronConnection {
cancellation: self.cancellation.clone(),
active_tasks: self.active_tasks.clone(),
app: self.app.clone(),
invitation_sync: self.invitation_sync.clone(),
});
Self::rotate_identity(&this).await
}
@ -3704,4 +4060,87 @@ mod tests {
assert!(recipient_block_policy_applies(true, 42, 43));
assert!(!recipient_block_policy_applies(false, 42, 43));
}
#[test]
fn omega_invitation_snapshot_parser_requires_authoritative_fields() {
let type_map = TypeMap::latest();
let mut fields = Vec::new();
for (kind, value) in [
(
DataType::InvitationAuthority,
DataValue::Str("omega".into()),
),
(DataType::InvitationId, DataValue::SignedNumber(7)),
(DataType::IotaId, DataValue::SignedNumber(42)),
(DataType::InvitationRevision, DataValue::SignedNumber(3)),
(
DataType::InvitationState,
DataValue::Str("provisioning".into()),
),
(DataType::InvitationCreatedAt, DataValue::SignedNumber(10)),
(DataType::InvitationExpiresAt, DataValue::SignedNumber(20)),
(
DataType::InvitationPasswordProtected,
DataValue::Bool(false),
),
] {
fields.push((kind.try_to_id(&type_map).expect("type is mapped"), value));
}
let parsed = parse_omega_invitation(&fields, &type_map, 42, 20).expect("snapshot parses");
assert_eq!(parsed.invitation_id, 7);
assert_eq!(parsed.remote_revision, 3);
assert_eq!(
parsed.authority,
iota_storage::users::invitations::InvitationAuthority::Omega
);
fields.retain(|(id, _)| Some(*id) != DataType::InvitationRevision.try_to_id(&type_map));
assert!(parse_omega_invitation(&fields, &type_map, 42, 20).is_err());
}
#[test]
fn omega_invitation_snapshot_parser_rejects_wrong_authority_or_iota() {
let type_map = TypeMap::latest();
let fields = [
(DataType::InvitationAuthority, DataValue::Str("iota".into())),
(DataType::InvitationId, DataValue::SignedNumber(7)),
(DataType::IotaId, DataValue::SignedNumber(42)),
(DataType::InvitationRevision, DataValue::SignedNumber(3)),
(DataType::InvitationState, DataValue::Str("pending".into())),
(DataType::InvitationCreatedAt, DataValue::SignedNumber(10)),
(DataType::InvitationExpiresAt, DataValue::SignedNumber(20)),
(
DataType::InvitationPasswordProtected,
DataValue::Bool(false),
),
]
.into_iter()
.map(|(kind, value)| (kind.try_to_id(&type_map).expect("type is mapped"), value))
.collect::<Vec<_>>();
assert!(parse_omega_invitation(&fields, &type_map, 42, 20).is_err());
let mut fields = fields;
let authority_id = DataType::InvitationAuthority
.try_to_id(&type_map)
.expect("type is mapped");
let authority = fields
.iter_mut()
.find(|(id, _)| *id == authority_id)
.expect("authority field exists");
authority.1 = DataValue::Str("omega".into());
assert!(parse_omega_invitation(&fields, &type_map, 43, 20).is_err());
}
#[test]
fn invitation_action_retry_delay_backs_off_and_caps() {
assert_eq!(
next_invitation_retry_delay(Duration::from_secs(5)),
Duration::from_secs(10)
);
assert_eq!(
next_invitation_retry_delay(Duration::from_secs(300)),
Duration::from_secs(300)
);
}
}