diff --git a/iota-cli/src/ipc_client.rs b/iota-cli/src/ipc_client.rs index 99e8653..bd1dd21 100644 --- a/iota-cli/src/ipc_client.rs +++ b/iota-cli/src/ipc_client.rs @@ -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() { diff --git a/iota-cli/src/screens/users/mod.rs b/iota-cli/src/screens/users/mod.rs index 62afad6..a80d226 100644 --- a/iota-cli/src/screens/users/mod.rs +++ b/iota-cli/src/screens/users/mod.rs @@ -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}")) diff --git a/iota-daemon-lib/src/command_router.rs b/iota-daemon-lib/src/command_router.rs index f28d5d7..9ad50c9 100644 --- a/iota-daemon-lib/src/command_router.rs +++ b/iota-daemon-lib/src/command_router.rs @@ -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, + 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) -> Vec { 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 { - authority, - invitation_id, - raw_token: iota_ipc::SecretString(raw_token), - expires_at, - short_url, - })) + 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, + 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); + } } diff --git a/iota-ipc/src/protocol.rs b/iota-ipc/src/protocol.rs index 8d3b40e..26e9e9a 100644 --- a/iota-ipc/src/protocol.rs +++ b/iota-ipc/src/protocol.rs @@ -371,6 +371,8 @@ pub struct InvitationCreated { pub raw_token: SecretString, pub expires_at: i64, pub short_url: Option, + #[serde(default)] + pub mirror_synced: bool, } #[derive(Clone, Debug, Deserialize, Serialize)] diff --git a/iota-storage/src/users/invitations.rs b/iota-storage/src/users/invitations.rs index b4b80d6..b5003d8 100644 --- a/iota-storage/src/users/invitations.rs +++ b/iota-storage/src/users/invitations.rs @@ -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 { 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 { 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>(2)?, + )) + })? + .collect::, _>>()?; + assert_eq!( + rows, + vec![ + ("iota".into(), "redeemed".into(), None), + ("omega".into(), "pending".into(), Some("revoke".into())), + ] + ); + Ok(()) + } } diff --git a/iota-storage/src/util/db.rs b/iota-storage/src/util/db.rs index 85e1c16..965aee7 100644 --- a/iota-storage/src/util/db.rs +++ b/iota-storage/src/util/db.rs @@ -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) = 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(()) } } diff --git a/iota/src/main.rs b/iota/src/main.rs index c0e593c..9eddc8a 100644 --- a/iota/src/main.rs +++ b/iota/src/main.rs @@ -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() { diff --git a/omikron-connector/src/omikron_connection.rs b/omikron-connector/src/omikron_connection.rs index e9e0bc3..6e05380 100644 --- a/omikron-connector/src/omikron_connection.rs +++ b/omikron-connector/src/omikron_connection.rs @@ -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) -> 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, 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, 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();