[Add] Persistent notification counts
This commit is contained in:
parent
8a54a360ca
commit
d460cb0677
8 changed files with 260 additions and 19 deletions
|
|
@ -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 //
|
||||
// ************************************************ //
|
||||
|
|
|
|||
|
|
@ -46,6 +46,12 @@ pub struct PolicyMutation {
|
|||
pub changed: Option<CommunicationValue>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct NotificationMutation {
|
||||
pub response: CommunicationValue,
|
||||
pub changed: Option<CommunicationValue>,
|
||||
}
|
||||
|
||||
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,
|
||||
|
|
|
|||
|
|
@ -4,6 +4,8 @@ pub struct Contact {
|
|||
pub user_name: Option<String>,
|
||||
pub created_at: i64,
|
||||
pub last_message_at: Option<i64>,
|
||||
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) {
|
||||
|
|
|
|||
|
|
@ -658,6 +658,12 @@ pub fn add_message(message: NewMessage<'_>) -> Result<i64, StorageError> {
|
|||
.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)
|
||||
})
|
||||
|
|
|
|||
|
|
@ -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<Option<Contact>, 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<Option<Contact>, 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<Vec<Contact>, 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<Vec<Contact>, 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<Vec<Contact>, StorageError> {
|
|||
Ok(out)
|
||||
})
|
||||
}
|
||||
|
||||
pub fn read_notifications(
|
||||
storage_owner: i64,
|
||||
user_id: i64,
|
||||
through: i64,
|
||||
) -> Result<Option<Contact>, 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)?,
|
||||
})
|
||||
},
|
||||
)?))
|
||||
})
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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(())
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)]
|
||||
|
|
|
|||
|
|
@ -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<CommunicationValue> {
|
||||
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)
|
||||
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::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<Self>, 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,
|
||||
|
|
|
|||
Loading…
Reference in a new issue