diff --git a/client/src/client_connection.rs b/client/src/client_connection.rs index 2e30bb7..6a793e8 100644 --- a/client/src/client_connection.rs +++ b/client/src/client_connection.rs @@ -210,9 +210,8 @@ impl ClientConnection { return; } - if cv.is_type(CommunicationType::ClientStateGet) { - // Durable client synchronization is routed through Omikron, which - // is the authority that assigns the numeric transport SessionId. + if cv.is_type(CommunicationType::AccountStateRequest) { + // Account initialization is routed through Omikron. self.send_message( &CommunicationValue::new(CommunicationType::ErrorInvalidData).with_request_id(&cv), ) @@ -220,7 +219,7 @@ impl ClientConnection { return; } - if cv.is_type(CommunicationType::ClientStateAck) { + if cv.is_type(CommunicationType::AccountStateApplied) { self.send_message( &CommunicationValue::new(CommunicationType::ErrorInvalidData).with_request_id(&cv), ) diff --git a/iota-connection/src/message_handlers.rs b/iota-connection/src/message_handlers.rs index e331f78..7882c79 100644 --- a/iota-connection/src/message_handlers.rs +++ b/iota-connection/src/message_handlers.rs @@ -17,17 +17,6 @@ use std::sync::atomic::{AtomicU32, Ordering}; static NEXT_NOTIFICATION_ID: AtomicU32 = AtomicU32::new(1); -fn valid_device_id(id: &str) -> bool { - id.len() == 36 - && id - .chars() - .enumerate() - .all(|(index, character)| match index { - 8 | 13 | 18 | 23 => character == '-', - _ => character.is_ascii_hexdigit(), - }) -} - fn next_notification_id() -> u32 { NEXT_NOTIFICATION_ID.fetch_add(1, Ordering::Relaxed).max(1) } @@ -953,138 +942,40 @@ mod presence_tests { } } -fn sync_error(cv: &CommunicationValue) -> CommunicationValue { - error_response(cv, CommunicationType::ErrorInvalidData).add_typed_default( - DataType::SessionId, - cv.get_data(DataType::SessionId) - .cloned() - .unwrap_or(DataValue::Null), - ) +fn account_state_error(cv: &CommunicationValue) -> CommunicationValue { + error_response(cv, CommunicationType::ErrorInvalidData) } /// The sender is authenticated by MTP; a UserId embedded by a client is never trusted here. -pub fn handle_client_state_get(cv: &CommunicationValue) -> CommunicationValue { - use iota_storage::util::sync::{self, CACHE_SCHEMA_VERSION}; +pub fn handle_account_state_request(cv: &CommunicationValue) -> CommunicationValue { let user_id = match required_sender_id(cv) { Ok(id) if id > 0 => id, - _ => return sync_error(cv), + _ => return account_state_error(cv), }; - let session_id = match data_i64(cv, DataType::SessionId) { - Some(id) if id > 0 => id, - _ => return sync_error(cv), + let contacts = match chats_util::get_users(user_id) { + Ok(contacts) => contacts, + Err(_) => return account_state_error(cv), }; - let device_id = match cv.get_data(DataType::DeviceId).and_then(DataValue::as_str) { - Some(id) if valid_device_id(id) => id.to_owned(), - _ => return sync_error(cv), + let messages = chat_files::get_all_messages(user_id); + let settings = match synced_settings::list(user_id) { + Ok(settings) => settings, + Err(_) => return account_state_error(cv), }; - let reported_version = match data_i64(cv, DataType::VersionNumber) { - Some(version) if version >= 0 => version, - _ => return sync_error(cv), - }; - let cache_valid = match cv.get_data(DataType::CacheValid).as_bool() { - Some(value) => value, - None => return sync_error(cv), - }; - let schema = match data_i64(cv, DataType::CacheSchemaVersion) { - Some(value) if value >= 0 => value, - _ => return sync_error(cv), - }; - let head = match sync::head(user_id) { - Ok(version) => version, - Err(_) => return sync_error(cv), - }; - let known_device = sync::has_device(user_id, &device_id).unwrap_or(false); - let acknowledged_version = sync::acknowledged_version(user_id, &device_id).unwrap_or(None); - let full = !cache_valid - || reported_version == 0 - || !known_device - || acknowledged_version.is_some_and(|version| reported_version < version) - || reported_version > head - || schema != CACHE_SCHEMA_VERSION; - let (contacts, messages, settings, deleted_messages, deleted_contacts, deleted_settings, mode) = - if full { - let settings = match synced_settings::list(user_id) { - Ok(settings) => settings, - Err(_) => return sync_error(cv), - }; - ( - match chats_util::get_users(user_id) { - Ok(contacts) => contacts, - Err(_) => return sync_error(cv), - }, - chat_files::get_all_messages(user_id), - settings, - Vec::new(), - Vec::new(), - Vec::new(), - "full", - ) - } else { - match sync::delta(user_id, reported_version, head) { - Ok(delta) => { - let settings = - match synced_settings::list_by_ids(user_id, &delta.setting_upserts) { - Ok(settings) => settings, - Err(_) => return sync_error(cv), - }; - ( - match chats_util::get_users_by_ids(user_id, &delta.contact_upserts) { - Ok(contacts) => contacts, - Err(_) => return sync_error(cv), - }, - chat_files::get_messages_by_ids(user_id, &delta.message_upserts), - settings, - delta.deleted_message_ids, - delta.deleted_contact_ids, - delta.deleted_setting_ids, - "delta", - ) - } - Err(_) => { - let settings = match synced_settings::list(user_id) { - Ok(settings) => settings, - Err(_) => return sync_error(cv), - }; - ( - match chats_util::get_users(user_id) { - Ok(contacts) => contacts, - Err(_) => return sync_error(cv), - }, - chat_files::get_all_messages(user_id), - settings, - Vec::new(), - Vec::new(), - Vec::new(), - "full", - ) - } - } - }; let message_values = messages .iter() .map(|message| stored_message_value(message, user_id, message.external_user)) .collect(); - if iota_storage::util::client_message_delivery::record_sync_delivery( - user_id, - session_id, - head, - messages.iter().map(|message| message.id), - ) - .is_err() - { - return sync_error(cv); - } let blocked_users = match blocked_users::list(user_id) { Ok(users) => users, - Err(_) => return sync_error(cv), + Err(_) => return account_state_error(cv), }; let receipt_policy = match receipt_policy::get(user_id) { Ok(policy) => policy, - Err(_) => return sync_error(cv), + Err(_) => return account_state_error(cv), }; let message_storage_policy = match message_storage_policy::get(user_id) { Ok(policy) => policy, - Err(_) => return sync_error(cv), + Err(_) => return account_state_error(cv), }; let contact_ids = match current_contact_ids(user_id) { Ok(contact_ids) => contact_ids, @@ -1094,22 +985,9 @@ pub fn handle_client_state_get(cv: &CommunicationValue) -> CommunicationValue { message_storage_policy::MessageRetention::Forever => None, message_storage_policy::MessageRetention::Duration { duration_ms } => Some(duration_ms), }; - let mut response = CommunicationValue::new(CommunicationType::ClientStateSync) + let mut response = CommunicationValue::new(CommunicationType::AccountStateSnapshot) .with_request_id(cv) .with_receiver(sender_wire_id(user_id)) - .add_typed_default( - DataType::SessionId, - DataValue::SignedNumber(session_id as i128), - ) - .add_typed_default( - DataType::VersionNumber, - DataValue::SignedNumber(head as i128), - ) - .add_typed_default( - DataType::CacheSchemaVersion, - DataValue::SignedNumber(CACHE_SCHEMA_VERSION as i128), - ) - .add_typed_default(DataType::SyncMode, DataValue::Str(mode.into())) .add_typed_default( DataType::Contacts, DataValue::Array( @@ -1154,33 +1032,6 @@ pub fn handle_client_state_get(cv: &CommunicationValue) -> CommunicationValue { DataType::Communities, DataValue::Array(community_values(user_id)), ) - .add_typed_default( - DataType::DeletedMessageIds, - DataValue::Array( - deleted_messages - .into_iter() - .map(|id| DataValue::SignedNumber(id as i128)) - .collect(), - ), - ) - .add_typed_default( - DataType::DeletedContactIds, - DataValue::Array( - deleted_contacts - .into_iter() - .map(|id| DataValue::SignedNumber(id as i128)) - .collect(), - ), - ) - .add_typed_default( - DataType::DeletedSettingIds, - DataValue::Array( - deleted_settings - .into_iter() - .map(|id| DataValue::SignedNumber(id as i128)) - .collect(), - ), - ) .add_typed_default(DataType::UserIds, contact_ids) .add_typed_default(DataType::Calls, DataValue::Array(Vec::new())); if let Some(duration_ms) = retention_duration { @@ -1192,44 +1043,13 @@ pub fn handle_client_state_get(cv: &CommunicationValue) -> CommunicationValue { response } -pub fn handle_client_state_ack(cv: &CommunicationValue) -> CommunicationValue { - use iota_storage::util::sync::CACHE_SCHEMA_VERSION; +pub fn handle_account_state_applied(cv: &CommunicationValue) -> CommunicationValue { let user_id = match required_sender_id(cv) { Ok(id) if id > 0 => id, - _ => return sync_error(cv), + _ => return account_state_error(cv), }; - let session_id = match data_i64(cv, DataType::SessionId) { - Some(id) if id > 0 => id, - _ => return sync_error(cv), - }; - let device_id = match cv.get_data(DataType::DeviceId).and_then(DataValue::as_str) { - Some(id) if valid_device_id(id) => id.to_owned(), - _ => return sync_error(cv), - }; - let version = match data_i64(cv, DataType::VersionNumber) { - Some(version) if version >= 0 => version, - _ => return sync_error(cv), - }; - if iota_storage::util::client_message_delivery::acknowledge_client_state( - user_id, - &device_id, - session_id, - version, - CACHE_SCHEMA_VERSION, - ) - .is_err() - { - return sync_error(cv); - } + let _ = user_id; success_response(cv) - .add_typed_default( - DataType::SessionId, - DataValue::SignedNumber(session_id as i128), - ) - .add_typed_default( - DataType::VersionNumber, - DataValue::SignedNumber(version as i128), - ) } pub fn handle_messages_get(cv: &CommunicationValue) -> CommunicationValue { diff --git a/mtp-type-maps b/mtp-type-maps index 2de3be6..2388c22 160000 --- a/mtp-type-maps +++ b/mtp-type-maps @@ -1 +1 @@ -Subproject commit 2de3be6016410304d6505d12ef37c35e63a9c731 +Subproject commit 2388c225b50db4d566653fb220786931eb3989c7 diff --git a/omikron-connector/src/omikron_connection.rs b/omikron-connector/src/omikron_connection.rs index d0531a7..3ce2426 100644 --- a/omikron-connector/src/omikron_connection.rs +++ b/omikron-connector/src/omikron_connection.rs @@ -2321,8 +2321,8 @@ impl OmikronConnection { dispatch!(LoadAppData, handle_load_app_data); dispatch!(CreateApp, handle_create_app); dispatch!(DeleteApp, handle_delete_app); - dispatch!(ClientStateGet, handle_client_state_get); - dispatch!(ClientStateAck, handle_client_state_ack); + dispatch!(AccountStateRequest, handle_account_state_request); + dispatch!(AccountStateApplied, handle_account_state_applied); dispatch!(ReadNotification, handle_read_notification); dispatch!(MessageEdit, handle_message_edit); dispatch!(MessageEditLive, handle_message_edit_live); @@ -2622,22 +2622,16 @@ impl OmikronConnection { .await; } - async fn handle_client_state_get(self: Arc, cv: &CommunicationValue) { - let response = message_handlers::handle_client_state_get(cv); + async fn handle_account_state_request(self: Arc, cv: &CommunicationValue) { + let response = message_handlers::handle_account_state_request(cv); let user_id = cv .require_sender() .ok() .and_then(|id| i64::try_from(id).ok()); - let session_id = cv - .get_data(DataType::SessionId) - .as_number() - .and_then(|id| i64::try_from(id).ok()); - - if response.is_type(CommunicationType::ClientStateSync) { + if response.is_type(CommunicationType::AccountStateSnapshot) { log!( - "ClientStateSync user={:?} session={:?} stage=generated", - user_id, - session_id + "AccountStateSnapshot user={:?} stage=generated", + user_id ); if let Some(user_id) = user_id && let Err(error) = relay_queue::pause_client_deliveries(user_id) @@ -2647,28 +2641,26 @@ impl OmikronConnection { if let Err(error) = self.send_message(&response).await { log!( - "Initial ClientStateSync delivery failed for user {:?}, session {:?}: {}", + "Initial AccountStateSnapshot delivery failed for user {:?}: {}", user_id, - session_id, error ); } else { log!( - "ClientStateSync user={:?} session={:?} stage=sent_to_omikron", - user_id, - session_id + "AccountStateSnapshot user={:?} stage=sent_to_omikron", + user_id ); } return; } if let Err(error) = self.send_message(&response).await { - log!("ClientStateGet response delivery failed: {}", error); + log!("AccountStateRequest response delivery failed: {}", error); } } - async fn handle_client_state_ack(self: Arc, cv: &CommunicationValue) { - let response = message_handlers::handle_client_state_ack(cv); + async fn handle_account_state_applied(self: Arc, cv: &CommunicationValue) { + let response = message_handlers::handle_account_state_applied(cv); let user_id = cv .require_sender() .ok()