diff --git a/client/src/client_connection.rs b/client/src/client_connection.rs index 040cb99..d9fd8c4 100644 --- a/client/src/client_connection.rs +++ b/client/src/client_connection.rs @@ -222,6 +222,15 @@ impl ClientConnection { return; } + if cv.is_type(CommunicationType::ReadNotification) { + let mutation = message_handlers::handle_read_notification(&cv); + self.send_message(&mutation.response).await; + if let Some(changed) = mutation.changed { + self.send_message(&changed).await; + } + return; + } + // ************************************************ // // Direct messages // // ************************************************ // diff --git a/iota-connection/src/message_handlers.rs b/iota-connection/src/message_handlers.rs index 86102d4..6e24b23 100644 --- a/iota-connection/src/message_handlers.rs +++ b/iota-connection/src/message_handlers.rs @@ -46,6 +46,12 @@ pub struct PolicyMutation { pub changed: Option, } +#[derive(Debug)] +pub struct NotificationMutation { + pub response: CommunicationValue, + pub changed: Option, +} + struct SettingLocator { scope: SettingScope, scope_key: String, @@ -860,6 +866,10 @@ fn contact_value( DataValue::SignedNumber(last_message_at as i128), )); } + fields.push(( + DataType::Notifications, + DataValue::SignedNumber(contact.notifications.into()), + )); fields.push(( DataType::Messages, DataValue::Array( @@ -1328,6 +1338,10 @@ pub fn handle_get_chats(cv: &CommunicationValue) -> CommunicationValue { if let Some(ts) = user.last_message_at { container.push((DataType::LastMessageAt, DataValue::SignedNumber(ts as i128))); } + container.push(( + DataType::Notifications, + DataValue::SignedNumber(user.notifications.into()), + )); user_array.push(typed_container(container)); } CommunicationValue::new(CommunicationType::GetChats) @@ -1336,6 +1350,76 @@ pub fn handle_get_chats(cv: &CommunicationValue) -> CommunicationValue { .add_typed_default(DataType::UserIds, DataValue::Array(user_array)) } +fn notification_response( + ty: CommunicationType, + request: &CommunicationValue, + owner: i64, + contact: &iota_storage::users::contact::Contact, +) -> CommunicationValue { + let mut response = CommunicationValue::new(ty) + .with_request_id(request) + .with_receiver(sender_wire_id(owner)) + .add_typed_default( + DataType::ChatPartnerId, + DataValue::SignedNumber(contact.user_id.into()), + ) + .add_typed_default( + DataType::Notifications, + DataValue::SignedNumber(contact.notifications.into()), + ); + if let Some(last_message_at) = contact.last_message_at { + response = response.add_typed_default( + DataType::LastMessageAt, + DataValue::SignedNumber(last_message_at.into()), + ); + } + response +} + +pub fn handle_read_notification(cv: &CommunicationValue) -> NotificationMutation { + let Ok(owner) = required_sender_id(cv) else { + return NotificationMutation { + response: error_response(cv, CommunicationType::ErrorInvalidData), + changed: None, + }; + }; + let Some(partner_id) = data_i64(cv, DataType::ChatPartnerId).filter(|id| *id > 0) else { + return NotificationMutation { + response: error_response(cv, CommunicationType::ErrorInvalidData), + changed: None, + }; + }; + let Some(through) = data_i64(cv, DataType::LastMessageAt).filter(|time| *time >= 0) else { + return NotificationMutation { + response: error_response(cv, CommunicationType::ErrorInvalidData), + changed: None, + }; + }; + + match chats_util::read_notifications(owner, partner_id, through) { + Ok(Some(contact)) => NotificationMutation { + response: notification_response( + CommunicationType::ReadNotification, + cv, + owner, + &contact, + ), + changed: Some( + notification_response(CommunicationType::PushNotification, cv, owner, &contact) + .with_id(next_notification_id()), + ), + }, + Ok(None) => NotificationMutation { + response: error_response(cv, CommunicationType::ErrorNotFound), + changed: None, + }, + Err(_) => NotificationMutation { + response: error_response(cv, CommunicationType::ErrorInternal), + changed: None, + }, + } +} + pub fn handle_add_community(cv: &CommunicationValue) -> CommunicationValue { let sender_id = match required_sender_id(cv) { Ok(sender_id) => sender_id, diff --git a/iota-storage/src/users/contact.rs b/iota-storage/src/users/contact.rs index 4d2ca43..c2d471e 100644 --- a/iota-storage/src/users/contact.rs +++ b/iota-storage/src/users/contact.rs @@ -4,6 +4,8 @@ pub struct Contact { pub user_name: Option, pub created_at: i64, pub last_message_at: Option, + pub notifications: i64, + pub notifications_read_at: i64, } impl Default for Contact { @@ -13,6 +15,8 @@ impl Default for Contact { user_name: None, created_at: 0, last_message_at: None, + notifications: 0, + notifications_read_at: 0, } } } @@ -28,6 +32,8 @@ impl Contact { user_name: None, created_at, last_message_at: None, + notifications: 0, + notifications_read_at: 0, } } pub fn set_last_message_at(&mut self, p0: i64) { diff --git a/iota-storage/src/util/chat_files.rs b/iota-storage/src/util/chat_files.rs index aee6032..3da9699 100644 --- a/iota-storage/src/util/chat_files.rs +++ b/iota-storage/src/util/chat_files.rs @@ -658,6 +658,12 @@ pub fn add_message(message: NewMessage<'_>) -> Result { .unwrap_or(authored_at), ); crate::util::chats_util::upsert_contact(&tx, storage_owner, &contact)?; + if !sent_by_self { + tx.execute( + "UPDATE contacts SET notifications = CASE WHEN notifications < 9223372036854775807 THEN notifications + 1 ELSE notifications END WHERE storage_owner = ?1 AND user_id = ?2", + params![storage_owner, external_user], + )?; + } tx.commit()?; Ok(msg_id) }) diff --git a/iota-storage/src/util/chats_util.rs b/iota-storage/src/util/chats_util.rs index 3eb4043..9257f76 100644 --- a/iota-storage/src/util/chats_util.rs +++ b/iota-storage/src/util/chats_util.rs @@ -2,7 +2,7 @@ use crate::storage_error::StorageError; use crate::users::contact::Contact; use crate::util::db; use crate::util::sync::{self, EntityType, Operation}; -use rusqlite::params; +use rusqlite::{OptionalExtension, params}; pub(crate) fn upsert_contact( tx: &rusqlite::Transaction<'_>, @@ -11,8 +11,11 @@ pub(crate) fn upsert_contact( ) -> Result<(), StorageError> { tx.execute( r#" - INSERT INTO contacts (storage_owner, user_id, user_name, created_at, last_message_at) - VALUES (?1, ?2, ?3, ?4, ?5) + INSERT INTO contacts ( + storage_owner, user_id, user_name, created_at, last_message_at, + notifications, notifications_read_at + ) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) ON CONFLICT(storage_owner, user_id) DO UPDATE SET user_name = COALESCE(excluded.user_name, contacts.user_name), created_at = MIN(contacts.created_at, excluded.created_at), @@ -28,6 +31,8 @@ pub(crate) fn upsert_contact( contact.user_name, contact.created_at, contact.last_message_at, + contact.notifications, + contact.notifications_read_at, ], )?; sync::record_event( @@ -69,7 +74,8 @@ pub fn get_user(storage_owner: i64, user_id: i64) -> Result, Sto db::with_db(|conn| { match conn.query_row( r#" - SELECT user_id, user_name, created_at, last_message_at + SELECT user_id, user_name, created_at, last_message_at, + notifications, notifications_read_at FROM contacts WHERE storage_owner = ?1 AND user_id = ?2 LIMIT 1 @@ -81,6 +87,8 @@ pub fn get_user(storage_owner: i64, user_id: i64) -> Result, Sto user_name: r.get(1)?, created_at: r.get(2)?, last_message_at: r.get(3)?, + notifications: r.get(4)?, + notifications_read_at: r.get(5)?, }) }, ) { @@ -95,7 +103,8 @@ pub fn get_users(storage_owner: i64) -> Result, StorageError> { db::with_db(|conn| { let mut stmt = conn.prepare( r#" - SELECT user_id, user_name, created_at, last_message_at + SELECT user_id, user_name, created_at, last_message_at, + notifications, notifications_read_at FROM contacts WHERE storage_owner = ?1 ORDER BY @@ -111,6 +120,8 @@ pub fn get_users(storage_owner: i64) -> Result, StorageError> { user_name: r.get(1)?, created_at: r.get(2)?, last_message_at: r.get(3)?, + notifications: r.get(4)?, + notifications_read_at: r.get(5)?, }) })?; @@ -121,3 +132,68 @@ pub fn get_users(storage_owner: i64) -> Result, StorageError> { Ok(out) }) } + +pub fn read_notifications( + storage_owner: i64, + user_id: i64, + through: i64, +) -> Result, StorageError> { + db::with_immediate_transaction(|tx| { + let Some(current_read_at) = tx + .query_row( + "SELECT notifications_read_at FROM contacts WHERE storage_owner = ?1 AND user_id = ?2", + params![storage_owner, user_id], + |row| row.get::<_, i64>(0), + ) + .optional()? + else { + return Ok(None); + }; + let read_at = current_read_at.max(through); + let notifications = tx.query_row( + r#" + SELECT COUNT(*) + FROM messages + WHERE storage_owner = ?1 + AND external_user = ?2 + AND sent_by_self = 0 + AND deleted_by_external = 0 + AND history_deleted = 0 + AND COALESCE(destination_iota_received_at, stored_at, authored_at, message_time) > ?3 + "#, + params![storage_owner, user_id, read_at], + |row| row.get::<_, i64>(0), + )?; + tx.execute( + "UPDATE contacts SET notifications = ?3, notifications_read_at = ?4 WHERE storage_owner = ?1 AND user_id = ?2", + params![storage_owner, user_id, notifications, read_at], + )?; + sync::record_event( + tx, + storage_owner, + EntityType::Contact, + user_id, + Operation::Upsert, + )?; + + Ok(Some(tx.query_row( + r#" + SELECT user_id, user_name, created_at, last_message_at, + notifications, notifications_read_at + FROM contacts + WHERE storage_owner = ?1 AND user_id = ?2 + "#, + params![storage_owner, user_id], + |row| { + Ok(Contact { + user_id: row.get(0)?, + user_name: row.get(1)?, + created_at: row.get(2)?, + last_message_at: row.get(3)?, + notifications: row.get(4)?, + notifications_read_at: row.get(5)?, + }) + }, + )?)) + }) +} diff --git a/iota-storage/src/util/db.rs b/iota-storage/src/util/db.rs index c12e14c..025cf84 100644 --- a/iota-storage/src/util/db.rs +++ b/iota-storage/src/util/db.rs @@ -182,6 +182,8 @@ fn run_migrations_on_connection(conn: &Connection) -> Result<(), StorageError> { user_name TEXT, created_at INTEGER NOT NULL, last_message_at INTEGER, + notifications INTEGER NOT NULL DEFAULT 0, + notifications_read_at INTEGER NOT NULL DEFAULT 0, UNIQUE(storage_owner, user_id) ); CREATE INDEX IF NOT EXISTS idx_contacts_owner @@ -858,6 +860,33 @@ fn run_migrations_on_connection(conn: &Connection) -> Result<(), StorageError> { conn.pragma_update(None, "user_version", 24)?; } + if current_version < 25 { + let contacts_exist: bool = conn.query_row( + "SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'contacts')", + [], + |row| row.get(0), + )?; + if contacts_exist { + add_table_column_if_missing( + conn, + "contacts", + "notifications", + "notifications INTEGER NOT NULL DEFAULT 0", + )?; + add_table_column_if_missing( + conn, + "contacts", + "notifications_read_at", + "notifications_read_at INTEGER NOT NULL DEFAULT 0", + )?; + conn.execute( + "UPDATE contacts SET notifications_read_at = COALESCE(last_message_at, 0) WHERE notifications_read_at = 0 AND notifications = 0", + [], + )?; + } + conn.pragma_update(None, "user_version", 25)?; + } + Ok(()) } @@ -927,7 +956,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, 24); + assert_eq!(version, 25); 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")?; @@ -946,7 +975,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, 24); + assert_eq!(version, 25); for table in [ "sync_heads", "sync_events", @@ -992,7 +1021,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, 24); + assert_eq!(version, 25); for column in [ "id", "user_id", @@ -1087,7 +1116,7 @@ mod tests { )?; let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?; assert_eq!(preserved, "remote_committed"); - assert_eq!(version, 24); + assert_eq!(version, 25); Ok(()) } } diff --git a/iota-storage/src/util/sync.rs b/iota-storage/src/util/sync.rs index 78cb6d6..e4aa539 100644 --- a/iota-storage/src/util/sync.rs +++ b/iota-storage/src/util/sync.rs @@ -4,7 +4,7 @@ use crate::util::db; use rusqlite::{OptionalExtension, Transaction, params}; use std::collections::BTreeMap; -pub const CACHE_SCHEMA_VERSION: i64 = 4; +pub const CACHE_SCHEMA_VERSION: i64 = 5; pub const STALE_CLIENT_SYNC_STATE_MS: i64 = 90 * 24 * 60 * 60 * 1_000; #[derive(Debug, Clone, Copy, PartialEq, Eq)] diff --git a/omikron-connector/src/omikron_connection.rs b/omikron-connector/src/omikron_connection.rs index b529305..3258883 100644 --- a/omikron-connector/src/omikron_connection.rs +++ b/omikron-connector/src/omikron_connection.rs @@ -3,7 +3,9 @@ use iota_logger::{log, log_cv_in, log_cv_out, log_t}; use iota_state::AppState; use iota_storage::util::config_util::{CONFIG, modify_config}; use iota_storage::util::relay_replay; -use iota_storage::util::{chat_files, client_relay_delivery, outgoing_relay, relay_queue}; +use iota_storage::util::{ + chat_files, chats_util, client_relay_delivery, outgoing_relay, relay_queue, +}; use iota_util::crypto_helper::{self, keyring_from_base64}; use iota_util::crypto_util::{self}; use mtp::client::{Client, ClientConfig, MTPConnection, Policy, SendMode, Sender}; @@ -952,6 +954,7 @@ impl OmikronConnection { signer_id: u64, recipient_id: u64, relay_message_id: &str, + accepted_at: i64, ) -> Option { let sender = DataValue::UnsignedNumber(u128::from(signer_id)); let receiver = recipient_id; @@ -977,14 +980,32 @@ impl OmikronConnection { if let Some(reply_id) = frame.get_data(DataType::ReplyId) { message.push((DataType::ReplyId, reply_id.clone())); } - Some( - CommunicationValue::new(CommunicationType::MessageLive) - .with_id(event_id) - .with_sender(signer_id) - .with_receiver(receiver) - .add_typed_default(DataType::SenderId, sender) - .add_typed_default(DataType::Message, typed_container(message)), - ) + let mut event = CommunicationValue::new(CommunicationType::MessageLive) + .with_id(event_id) + .with_sender(signer_id) + .with_receiver(receiver) + .add_typed_default(DataType::SenderId, sender) + .add_typed_default(DataType::Message, typed_container(message)) + .add_typed_default( + DataType::LastMessageAt, + DataValue::SignedNumber(accepted_at.into()), + ); + if let (Ok(owner), Ok(partner)) = + (i64::try_from(recipient_id), i64::try_from(signer_id)) + && let Ok(Some(contact)) = chats_util::get_user(owner, partner) + { + event = event.add_typed_default( + DataType::Notifications, + DataValue::SignedNumber(contact.notifications.into()), + ); + if let Some(last_message_at) = contact.last_message_at { + event = event.add_typed_default( + DataType::LastMessageAt, + DataValue::SignedNumber(last_message_at.into()), + ); + } + } + Some(event) } CommunicationType::SetChatSecret => { let frame = CommunicationValue::new(CommunicationType::SetChatSecret) @@ -1801,6 +1822,7 @@ impl OmikronConnection { verified.context.signer_id, destination, &verified.context.message_id, + accepted_at, ); if let Err(error) = relay_replay::mark_applied(verified.context.signer_id, &verified.context.message_id) @@ -2301,6 +2323,7 @@ impl OmikronConnection { dispatch!(DeleteApp, handle_delete_app); dispatch!(ClientConnected, handle_client_connected); dispatch!(ClientStateAck, handle_client_state_ack); + dispatch!(ReadNotification, handle_read_notification); dispatch!(MessageEdit, handle_message_edit); dispatch!(MessageEditLive, handle_message_edit_live); dispatch!(MessageReactionAdd, handle_message_reaction_add); @@ -2662,6 +2685,14 @@ impl OmikronConnection { } } + async fn handle_read_notification(self: Arc, cv: &CommunicationValue) { + let mutation = message_handlers::handle_read_notification(cv); + let _ = self.send_message(&mutation.response).await; + if let Some(changed) = mutation.changed { + let _ = self.send_message(&changed).await; + } + } + fn mutation_live_message( ty: CommunicationType, request: &CommunicationValue,