[Fix] Calling & Invites

This commit is contained in:
Alex Emmet 2026-09-11 18:31:27 +02:00
commit 9e9e3597da
No known key found for this signature in database
8 changed files with 220 additions and 39 deletions

View file

@ -485,7 +485,7 @@ impl IpcClient {
}
}
ResponsePayload::InvitationCreated(invitation) => format!(
"Created {:?} invitation {}\nToken: {}\nExpires: {}{}",
"Created {:?} invitation {}\nToken: {}\nExpires: {}{}{}",
invitation.authority,
invitation.invitation_id,
invitation.raw_token.0,
@ -495,6 +495,11 @@ impl IpcClient {
.as_ref()
.map(|url| format!("\nLink: {url}"))
.unwrap_or_default(),
if invitation.mirror_synced {
String::new()
} else {
"\nWarning: local invitation metadata is still synchronizing.".into()
},
),
ResponsePayload::Invitations(invitations) => {
if invitations.is_empty() {

View file

@ -418,10 +418,15 @@ impl UsersScreen {
Ok(iota_ipc::ResponseResult::Ok(
iota_ipc::ResponsePayload::InvitationCreated(invitation),
)) => Ok(format!(
"Invitation {} created. Token: {} Link: {}",
"Invitation {} created. Token: {} Link: {}{}",
invitation.invitation_id,
invitation.raw_token.0,
invitation.short_url.as_deref().unwrap_or("unavailable")
invitation.short_url.as_deref().unwrap_or("unavailable"),
if invitation.mirror_synced {
""
} else {
" Local metadata is still synchronizing."
}
)),
Ok(iota_ipc::ResponseResult::Error(error)) => {
Err(format!("Create invitation failed: {error}"))

View file

@ -63,6 +63,35 @@ fn now_millis() -> i64 {
.as_millis() as i64
}
fn invitation_created_after_mirror(
authority: InvitationAuthority,
invitation_id: i64,
raw_token: String,
expires_at: i64,
short_url: Option<String>,
mirror_result: Result<(), iota_storage::storage_error::StorageError>,
) -> InvitationCreated {
let mirror_synced = match mirror_result {
Ok(()) => true,
Err(error) => {
log!(
"Invitation {} was created by Omega, but its local mirror could not be stored: {}",
invitation_id,
error
);
false
}
};
InvitationCreated {
authority,
invitation_id,
raw_token: iota_ipc::SecretString(raw_token),
expires_at,
short_url,
mirror_synced,
}
}
fn bounded_log_entries(mut entries: Vec<LogEntry>) -> Vec<LogEntry> {
entries.truncate(MAX_LOG_ENTRIES_PER_RESPONSE);
while !entries.is_empty() {
@ -343,16 +372,18 @@ impl CommandRouter {
local_provisioned_user_id: None,
local_provisioned_at: None,
};
if iota_storage::users::invitations::insert(&summary, None, now_millis()).is_err() {
return ResponseResult::Error(IpcErrorCode::StorageFailure);
}
ResponseResult::Ok(ResponsePayload::InvitationCreated(InvitationCreated {
let mirror_result =
iota_storage::users::invitations::insert(&summary, None, now_millis());
ResponseResult::Ok(ResponsePayload::InvitationCreated(
invitation_created_after_mirror(
authority,
invitation_id,
raw_token: iota_ipc::SecretString(raw_token),
raw_token,
expires_at,
short_url,
}))
mirror_result,
),
))
}
LocalRequest::ListInvitations { authority } => {
if authority != Some(InvitationAuthority::Iota)
@ -1036,7 +1067,10 @@ impl CommandRouter {
#[cfg(test)]
mod tests {
use super::{IpcRole, LocalRequest, bounded_log_entries};
use super::{
InvitationAuthority, IpcRole, LocalRequest, bounded_log_entries,
invitation_created_after_mirror,
};
use iota_ipc::{ExitIntent, LogEntry, SecretString};
#[test]
@ -1119,4 +1153,25 @@ mod tests {
assert!(bounded_log_entries(entries).is_empty());
}
#[test]
fn omega_creation_credentials_survive_a_local_mirror_failure() {
let created = invitation_created_after_mirror(
InvitationAuthority::Omega,
7,
"raw-token".into(),
99,
Some("https://omega/direct/short".into()),
Err(iota_storage::storage_error::StorageError::Other(
"injected mirror failure".into(),
)),
);
assert_eq!(created.raw_token.0, "raw-token");
assert_eq!(
created.short_url.as_deref(),
Some("https://omega/direct/short")
);
assert!(!created.mirror_synced);
}
}

View file

@ -371,6 +371,8 @@ pub struct InvitationCreated {
pub raw_token: SecretString,
pub expires_at: i64,
pub short_url: Option<String>,
#[serde(default)]
pub mirror_synced: bool,
}
#[derive(Clone, Debug, Deserialize, Serialize)]

View file

@ -128,7 +128,7 @@ fn insert_in_tx(
synced_at: i64,
) -> Result<(), StorageError> {
tx.execute(
"INSERT INTO user_invitations (invitation_id, authority, token_hash, label, password_protected, created_at, expires_at, state, redeemed_user_id, redeemed_at, revoked_at, pending_action, pending_action_at, remote_revision, last_synced_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, NULL, NULL, ?12, ?13) ON CONFLICT(invitation_id) DO UPDATE SET authority = excluded.authority, token_hash = excluded.token_hash, label = excluded.label, password_protected = excluded.password_protected, expires_at = excluded.expires_at, state = excluded.state, redeemed_user_id = excluded.redeemed_user_id, redeemed_at = excluded.redeemed_at, revoked_at = excluded.revoked_at, remote_revision = excluded.remote_revision, last_synced_at = excluded.last_synced_at WHERE excluded.remote_revision >= user_invitations.remote_revision",
"INSERT INTO user_invitations (invitation_id, authority, token_hash, label, password_protected, created_at, expires_at, state, redeemed_user_id, redeemed_at, revoked_at, pending_action, pending_action_at, remote_revision, last_synced_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, NULL, NULL, ?12, ?13) ON CONFLICT(authority, invitation_id) DO UPDATE SET token_hash = excluded.token_hash, label = excluded.label, password_protected = excluded.password_protected, expires_at = excluded.expires_at, state = excluded.state, redeemed_user_id = excluded.redeemed_user_id, redeemed_at = excluded.redeemed_at, revoked_at = excluded.revoked_at, remote_revision = excluded.remote_revision, last_synced_at = excluded.last_synced_at WHERE excluded.remote_revision >= user_invitations.remote_revision",
params![summary.invitation_id, summary.authority.as_str(), token_hash, summary.label, summary.password_protected, summary.created_at, summary.expires_at, summary.state.as_str(), summary.redeemed_user_id, summary.redeemed_at, summary.revoked_at, summary.remote_revision, synced_at],
)?;
Ok(())
@ -207,7 +207,7 @@ pub fn apply_revoke_result(
) -> Result<bool, StorageError> {
db::with_immediate_transaction(|tx| {
let changed = tx.execute(
"UPDATE user_invitations SET state = ?2, remote_revision = ?3, revoked_at = ?4, pending_action = NULL, pending_action_at = NULL, last_synced_at = ?5 WHERE invitation_id = ?1 AND ?3 >= remote_revision",
"UPDATE user_invitations SET state = ?2, remote_revision = ?3, revoked_at = ?4, pending_action = NULL, pending_action_at = NULL, last_synced_at = ?5 WHERE invitation_id = ?1 AND authority = 'omega' AND ?3 >= remote_revision",
params![invitation_id, state.as_str(), remote_revision, revoked_at, synced_at],
)?;
Ok(changed == 1)
@ -272,7 +272,7 @@ fn apply_external_invitation_provisioning_in_tx(
return Ok(ProvisioningResult::Conflict);
}
tx.execute(
"UPDATE user_invitations SET local_provisioned_user_id = ?2, local_provisioned_at = COALESCE(local_provisioned_at, ?3) WHERE invitation_id = ?1",
"UPDATE user_invitations SET local_provisioned_user_id = ?2, local_provisioned_at = COALESCE(local_provisioned_at, ?3) WHERE invitation_id = ?1 AND authority = 'omega'",
params![invitation_id, user.user_id, changed_at],
)?;
return Ok(ProvisioningResult::AlreadyApplied);
@ -286,7 +286,7 @@ fn apply_external_invitation_provisioning_in_tx(
params![user.user_id, user.username, changed_at],
)?;
tx.execute(
"UPDATE user_invitations SET local_provisioned_user_id = ?2, local_provisioned_at = ?3 WHERE invitation_id = ?1",
"UPDATE user_invitations SET local_provisioned_user_id = ?2, local_provisioned_at = ?3 WHERE invitation_id = ?1 AND authority = 'omega'",
params![invitation_id, user.user_id, changed_at],
)?;
Ok(ProvisioningResult::Created)
@ -300,7 +300,7 @@ mod tests {
fn database() -> Result<Connection, StorageError> {
let connection = Connection::open_in_memory()?;
connection.execute_batch(
"CREATE TABLE user_invitations (invitation_id INTEGER PRIMARY KEY, authority TEXT NOT NULL, token_hash BLOB, label TEXT, password_protected INTEGER NOT NULL, created_at INTEGER NOT NULL, expires_at INTEGER, state TEXT NOT NULL, redeemed_user_id INTEGER, redeemed_at INTEGER, revoked_at INTEGER, pending_action TEXT, pending_action_at INTEGER, remote_revision INTEGER NOT NULL DEFAULT 0, last_synced_at INTEGER, local_provisioned_user_id INTEGER, local_provisioned_at INTEGER); CREATE TABLE users (user_id INTEGER PRIMARY KEY, username TEXT NOT NULL, public_key TEXT NOT NULL, private_key_hash TEXT, reset_token TEXT, created_at INTEGER NOT NULL, display_name TEXT); CREATE TABLE user_residency (user_id INTEGER PRIMARY KEY, username TEXT NOT NULL, lifecycle_state TEXT NOT NULL, data_state TEXT NOT NULL, credential_origin TEXT NOT NULL, updated_at INTEGER NOT NULL);",
"CREATE TABLE user_invitations (invitation_id INTEGER NOT NULL, authority TEXT NOT NULL, token_hash BLOB, label TEXT, password_protected INTEGER NOT NULL, created_at INTEGER NOT NULL, expires_at INTEGER, state TEXT NOT NULL, redeemed_user_id INTEGER, redeemed_at INTEGER, revoked_at INTEGER, pending_action TEXT, pending_action_at INTEGER, remote_revision INTEGER NOT NULL DEFAULT 0, last_synced_at INTEGER, local_provisioned_user_id INTEGER, local_provisioned_at INTEGER, PRIMARY KEY (authority, invitation_id)); CREATE TABLE users (user_id INTEGER PRIMARY KEY, username TEXT NOT NULL, public_key TEXT NOT NULL, private_key_hash TEXT, reset_token TEXT, created_at INTEGER NOT NULL, display_name TEXT); CREATE TABLE user_residency (user_id INTEGER PRIMARY KEY, username TEXT NOT NULL, lifecycle_state TEXT NOT NULL, data_state TEXT NOT NULL, credential_origin TEXT NOT NULL, updated_at INTEGER NOT NULL);",
)?;
Ok(connection)
}
@ -440,4 +440,46 @@ mod tests {
);
Ok(())
}
#[test]
fn equal_numeric_ids_are_isolated_by_authority() -> Result<(), StorageError> {
let mut connection = database()?;
let tx = connection.transaction()?;
let omega = summary(3, InvitationState::Pending);
let mut local = omega.clone();
local.authority = InvitationAuthority::Iota;
local.state = InvitationState::Redeemed;
insert_in_tx(&tx, &omega, None, 30)?;
insert_in_tx(&tx, &local, Some(&[1, 2]), 31)?;
assert_eq!(
tx.query_row(
"SELECT COUNT(*) FROM user_invitations WHERE invitation_id = 7",
[],
|row| row.get::<_, i64>(0),
)?,
2
);
assert!(mark_revoke_pending_in_tx(&tx, 7, 32)?);
let rows = tx
.prepare(
"SELECT authority, state, pending_action FROM user_invitations WHERE invitation_id = 7 ORDER BY authority",
)?
.query_map([], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, Option<String>>(2)?,
))
})?
.collect::<Result<Vec<_>, _>>()?;
assert_eq!(
rows,
vec![
("iota".into(), "redeemed".into(), None),
("omega".into(), "pending".into(), Some("revoke".into())),
]
);
Ok(())
}
}

View file

@ -958,14 +958,9 @@ fn run_migrations_on_connection(conn: &Connection) -> Result<(), StorageError> {
revoked_at INTEGER,
local_revocation_pending INTEGER NOT NULL DEFAULT 0
);
INSERT INTO user_invitations (
invitation_id, authority, token_hash, created_at, expires_at,
state, redeemed_user_id, redeemed_at, revoked_at
)
SELECT CAST(invitation_id AS INTEGER), 'omega', token_hash, created_at, expires_at,
state, redeemed_user_id, redeemed_at, revoked_at
FROM user_invitations_before_authority
WHERE CAST(invitation_id AS INTEGER) > 0;
/* Omega owns these mirror records. Historical string identifiers
cannot be mapped safely into Omega's numeric namespace, so a
fresh authoritative snapshot rebuilds them after reconnect. */
DROP TABLE user_invitations_before_authority;
CREATE INDEX idx_user_invitations_state_created
ON user_invitations (state, created_at);
@ -1047,6 +1042,40 @@ fn run_migrations_on_connection(conn: &Connection) -> Result<(), StorageError> {
conn.pragma_update(None, "user_version", 29)?;
}
if current_version < 30 {
conn.execute_batch(
r#"
ALTER TABLE user_invitations RENAME TO user_invitations_before_composite_key;
CREATE TABLE user_invitations (
invitation_id INTEGER NOT NULL,
authority TEXT NOT NULL CHECK (authority IN ('omega', 'iota')),
token_hash BLOB,
label TEXT,
password_protected INTEGER NOT NULL DEFAULT 0,
created_at INTEGER NOT NULL,
expires_at INTEGER,
state TEXT NOT NULL CHECK (state IN ('pending', 'provisioning', 'redeemed', 'revoked', 'expired')),
redeemed_user_id INTEGER,
redeemed_at INTEGER,
revoked_at INTEGER,
pending_action TEXT CHECK (pending_action IS NULL OR pending_action = 'revoke'),
pending_action_at INTEGER,
remote_revision INTEGER NOT NULL DEFAULT 0,
last_synced_at INTEGER,
local_provisioned_user_id INTEGER,
local_provisioned_at INTEGER,
PRIMARY KEY (authority, invitation_id)
);
INSERT INTO user_invitations
SELECT * FROM user_invitations_before_composite_key;
DROP TABLE user_invitations_before_composite_key;
CREATE INDEX idx_user_invitations_state_created
ON user_invitations (authority, state, created_at);
PRAGMA user_version = 30;
"#,
)?;
}
Ok(())
}
@ -1121,7 +1150,7 @@ mod tests {
run_migrations_on_connection(&conn)?;
let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?;
assert_eq!(version, 29);
assert_eq!(version, 30);
for column in ["height", "reply_to", "edited_count", "deleted_by_external"] {
let mut statement =
conn.prepare("SELECT 1 FROM pragma_table_info('messages') WHERE name = ?1")?;
@ -1140,7 +1169,7 @@ mod tests {
run_migrations_on_connection(&conn)?;
run_migrations_on_connection(&conn)?;
let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?;
assert_eq!(version, 29);
assert_eq!(version, 30);
for table in [
"sync_heads",
"sync_events",
@ -1187,7 +1216,7 @@ mod tests {
run_migrations_on_connection(&conn)?;
let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?;
assert_eq!(version, 29);
assert_eq!(version, 30);
for column in [
"id",
"user_id",
@ -1282,12 +1311,13 @@ mod tests {
)?;
let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?;
assert_eq!(preserved, "remote_committed");
assert_eq!(version, 29);
assert_eq!(version, 30);
Ok(())
}
#[test]
fn invitation_authority_migration_preserves_numeric_invitations() -> Result<(), StorageError> {
fn invitation_authority_migration_discards_unmappable_omega_mirrors() -> Result<(), StorageError>
{
let conn = Connection::open_in_memory()?;
conn.execute_batch(
r#"
@ -1303,19 +1333,23 @@ mod tests {
);
INSERT INTO user_invitations (
invitation_id, token_hash, created_at, expires_at, state
) VALUES ('42', X'0102', 10, 20, 'pending');
) VALUES
('42', X'0102', 10, 20, 'pending'),
('4f40e654-1234-1234-1234-123456789abc', X'0102', 10, 20, 'pending'),
('af40e654-1234-1234-1234-123456789abc', X'0102', 10, 20, 'pending'),
('malformed', X'0102', 10, 20, 'pending');
PRAGMA user_version = 26;
"#,
)?;
run_migrations_on_connection(&conn)?;
let row: (i64, String, String, Option<String>) = conn.query_row(
"SELECT invitation_id, authority, state, pending_action FROM user_invitations",
[],
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
)?;
assert_eq!(row, (42, "omega".into(), "pending".into(), None));
let count: i64 = conn.query_row("SELECT COUNT(*) FROM user_invitations", [], |row| {
row.get(0)
})?;
assert_eq!(count, 0);
let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?;
assert_eq!(version, 30);
Ok(())
}
}

View file

@ -850,6 +850,9 @@ async fn run_command(
println!("Link: {url}");
}
println!("Expires: {}", invitation.expires_at);
if !invitation.mirror_synced {
println!("Warning: local invitation metadata is still synchronizing.");
}
}
ResponsePayload::Invitations(invitations) => {
if invitations.is_empty() {

View file

@ -150,6 +150,15 @@ fn container_signed_i64(
.and_then(|value| i64::try_from(value).ok())
}
fn targets_iota(value: &CommunicationValue, expected_iota_id: Option<u64>) -> bool {
let target = value
.get_data(DataType::IotaId)
.as_signed_number()
.and_then(|id| u64::try_from(id).ok())
.filter(|id| *id > 0);
target.is_some() && target == expected_iota_id
}
fn parse_omega_invitation(
fields: &[(DataTypeId, DataValue)],
type_map: &TypeMap,
@ -2665,6 +2674,12 @@ impl OmikronConnection {
/// Omega-authorized account cleanup. The storage operation is idempotent;
/// acknowledgement is therefore safe to retry after a reconnect.
async fn handle_erase_hosted_user_data(self: Arc<Self>, cv: &CommunicationValue) {
if !targets_iota(cv, CONFIG.load().iota_id) {
let _ = self
.send_message(&error_response(cv, CommunicationType::ErrorInvalidData))
.await;
return;
}
let Some(user_id) = cv
.get_data(DataType::UserId)
.as_signed_number()
@ -2693,6 +2708,12 @@ impl OmikronConnection {
/* Omega sends only public account metadata here. The locally created
* profile records an external credential origin and never receives a TU. */
async fn handle_iota_user_provisioning(self: Arc<Self>, cv: &CommunicationValue) {
if !targets_iota(cv, CONFIG.load().iota_id) {
let _ = self
.send_message(&error_response(cv, CommunicationType::ErrorInvalidData))
.await;
return;
}
let user_id = cv
.get_data(DataType::UserId)
.as_signed_number()
@ -4098,6 +4119,20 @@ mod tests {
assert!(parse_omega_invitation(&fields, &type_map, 42, 20).is_err());
}
#[test]
fn omega_control_target_requires_signed_local_iota_id() {
let matching = CommunicationValue::new(CommunicationType::ProvisionIotaUser)
.add_typed_default(DataType::IotaId, DataValue::SignedNumber(42));
let wrong = CommunicationValue::new(CommunicationType::EraseHostedUserData)
.add_typed_default(DataType::IotaId, DataValue::SignedNumber(43));
let unsigned = CommunicationValue::new(CommunicationType::EraseHostedUserData)
.add_typed_default(DataType::IotaId, DataValue::UnsignedNumber(42));
assert!(targets_iota(&matching, Some(42)));
assert!(!targets_iota(&wrong, Some(42)));
assert!(!targets_iota(&unsigned, Some(42)));
assert!(!targets_iota(&matching, None));
}
#[test]
fn omega_invitation_snapshot_parser_rejects_wrong_authority_or_iota() {
let type_map = TypeMap::latest();