iota/iota-storage/src/util/chat_files.rs
2026-08-08 02:09:38 +02:00

793 lines
24 KiB
Rust

use crate::storage_error::StorageError;
use crate::util::db;
use crate::util::sync::{self, EntityType, Operation};
use iota_logger::log;
use rusqlite::params;
#[derive(PartialEq, Debug, Clone)]
pub enum MessageState {
Read,
Received,
Sent,
Sending,
}
impl MessageState {
pub fn as_str(&self) -> &'static str {
match self {
MessageState::Read => "read",
MessageState::Received => "received",
MessageState::Sent => "sent",
MessageState::Sending => "sending",
}
}
pub fn from_str(value: &str) -> Self {
match value.to_lowercase().as_str() {
"read" => MessageState::Read,
"received" => MessageState::Received,
"sent" => MessageState::Sent,
_ => MessageState::Sending,
}
}
pub fn upgrade(self, other: Self) -> Self {
if other == Self::Read || self == Self::Read {
Self::Read
} else if other == Self::Received || self == Self::Received {
Self::Received
} else if other == Self::Sent || self == Self::Sent {
Self::Sent
} else {
Self::Sending
}
}
}
#[derive(Debug, Clone)]
pub struct StoredMessage {
pub id: i64,
pub external_user: i64,
pub message_time: i64,
pub content: String,
pub edited: bool,
pub sent_by_self: bool,
pub message_state: String,
pub height: i64,
pub reply_to: Option<i64>,
pub reactions: Vec<StoredReaction>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StoredReaction {
pub reaction: String,
pub user_id: i64,
}
/*
* Each edit is recorded in message_edits with the before/after content and a
* timestamp. Only the original sender (sent_by_self = 1) may edit.
*/
pub fn edit_message(
storage_owner: i64,
external_user: i64,
message_time: i64,
editor_id: i64,
new_content: &str,
) -> Result<(), StorageError> {
update_message_content(
storage_owner,
external_user,
message_time,
editor_id,
new_content,
true,
)
}
/* Applies an edit received from the message sender to the recipient's copy. */
pub fn apply_remote_edit(
storage_owner: i64,
external_user: i64,
message_time: i64,
editor_id: i64,
new_content: &str,
) -> Result<(), StorageError> {
if editor_id != external_user {
return Err(StorageError::Other(
"Remote editor does not match chat partner".into(),
));
}
update_message_content(
storage_owner,
external_user,
message_time,
editor_id,
new_content,
false,
)
}
fn update_message_content(
storage_owner: i64,
external_user: i64,
message_time: i64,
editor_id: i64,
new_content: &str,
require_sent_by_self: bool,
) -> Result<(), StorageError> {
db::with_db(|conn| {
let msg = conn.query_row(
r#"
SELECT id, content, sent_by_self
FROM messages
WHERE storage_owner = ?1 AND external_user = ?2 AND message_time = ?3
ORDER BY id DESC LIMIT 1
"#,
params![storage_owner, external_user, message_time],
|row| {
Ok((
row.get::<_, i64>(0)?,
row.get::<_, String>(1)?,
row.get::<_, i64>(2)?,
))
},
)?;
let (msg_id, old_content, sent_by_self) = msg;
if require_sent_by_self && sent_by_self != 1 {
return Err(StorageError::Other(
"Only the original sender can edit this message".into(),
));
}
if !require_sent_by_self && sent_by_self != 0 {
return Err(StorageError::Other(
"Remote edits may only update received messages".into(),
));
}
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as i64;
let tx = conn.unchecked_transaction()?;
tx.execute(
r#"
INSERT INTO message_edits (message_id, content_before, content_after, edited_at, edited_by)
VALUES (?1, ?2, ?3, ?4, ?5)
"#,
params![msg_id, old_content, new_content, now, editor_id],
)?;
tx.execute(
r#"
UPDATE messages
SET content = ?1, edited_count = edited_count + 1
WHERE id = ?2
"#,
params![new_content, msg_id],
)?;
sync::record_event(
&tx,
storage_owner,
EntityType::Message,
msg_id,
Operation::Upsert,
)?;
tx.commit()?;
Ok(())
})
}
pub fn hard_delete_message(
storage_owner: i64,
external_user: i64,
message_time: i64,
) -> Result<(), StorageError> {
db::with_db(|conn| {
let msg_id: i64 = conn.query_row(
r#"
SELECT id FROM messages
WHERE storage_owner = ?1 AND external_user = ?2 AND message_time = ?3
ORDER BY id DESC LIMIT 1
"#,
params![storage_owner, external_user, message_time],
|row| row.get(0),
)?;
let tx = conn.unchecked_transaction()?;
tx.execute(
"DELETE FROM message_edits WHERE message_id = ?1",
params![msg_id],
)?;
tx.execute(
"DELETE FROM reactions WHERE message_id = ?1",
params![msg_id],
)?;
tx.execute("DELETE FROM messages WHERE id = ?1", params![msg_id])?;
sync::record_event(
&tx,
storage_owner,
EntityType::Message,
msg_id,
Operation::Delete,
)?;
tx.commit()?;
Ok(())
})
}
/* Deletes a message from the sender's local copy after checking ownership. */
pub fn delete_message(
storage_owner: i64,
external_user: i64,
message_time: i64,
) -> Result<(), StorageError> {
ensure_message_direction(storage_owner, external_user, message_time, true)?;
hard_delete_message(storage_owner, external_user, message_time)
}
/* Flags the recipient's local copy after validating its sender, preserving its history. */
pub fn apply_remote_delete(
storage_owner: i64,
external_user: i64,
message_time: i64,
sender_id: i64,
) -> Result<(), StorageError> {
if sender_id != external_user {
return Err(StorageError::Other(
"Remote sender does not match chat partner".into(),
));
}
ensure_message_direction(storage_owner, external_user, message_time, false)?;
flag_deleted_by_external(storage_owner, external_user, message_time)
}
fn ensure_message_direction(
storage_owner: i64,
external_user: i64,
message_time: i64,
expected_sent_by_self: bool,
) -> Result<(), StorageError> {
db::with_db(|conn| {
let sent_by_self: i64 = conn.query_row(
r#"
SELECT sent_by_self FROM messages
WHERE storage_owner = ?1 AND external_user = ?2 AND message_time = ?3
ORDER BY id DESC LIMIT 1
"#,
params![storage_owner, external_user, message_time],
|row| row.get(0),
)?;
if (sent_by_self != 0) != expected_sent_by_self {
return Err(StorageError::Other(
"Message sender is not authorized".into(),
));
}
Ok(())
})
}
/*
* Marks a message as deleted by the external user rather than removing the row,
* so the storage owner still sees a tombstone in the UI.
*/
pub fn flag_deleted_by_external(
storage_owner: i64,
external_user: i64,
message_time: i64,
) -> Result<(), StorageError> {
db::with_db(|conn| {
let tx = conn.unchecked_transaction()?;
let affected = tx.execute(
r#"
UPDATE messages
SET deleted_by_external = 1
WHERE storage_owner = ?1 AND external_user = ?2 AND message_time = ?3
"#,
params![storage_owner, external_user, message_time],
)?;
if affected == 0 {
return Err(StorageError::Other("Message not found".into()));
}
let msg_id: i64 = tx.query_row("SELECT id FROM messages WHERE storage_owner = ?1 AND external_user = ?2 AND message_time = ?3 ORDER BY id DESC LIMIT 1", params![storage_owner, external_user, message_time], |r| r.get(0))?;
sync::record_event(
&tx,
storage_owner,
EntityType::Message,
msg_id,
Operation::Delete,
)?;
tx.commit()?;
Ok(())
})
}
/*
* Removes the edit trail but keeps the message with edited_count > 0 so
* the UI still shows the "edited" indicator. Only the own user should
* call this.
*/
pub fn delete_edit_history(
storage_owner: i64,
external_user: i64,
message_time: i64,
) -> Result<(), StorageError> {
db::with_db(|conn| {
let msg_id: i64 = conn.query_row(
r#"
SELECT id FROM messages
WHERE storage_owner = ?1 AND external_user = ?2 AND message_time = ?3
ORDER BY id DESC LIMIT 1
"#,
params![storage_owner, external_user, message_time],
|row| row.get(0),
)?;
let tx = conn.unchecked_transaction()?;
tx.execute(
"DELETE FROM message_edits WHERE message_id = ?1",
params![msg_id],
)?;
sync::record_event(
&tx,
storage_owner,
EntityType::Message,
msg_id,
Operation::Upsert,
)?;
tx.commit()?;
Ok(())
})
}
pub fn add_reaction(
storage_owner: i64,
external_user: i64,
message_time: i64,
user_id: i64,
reaction: &str,
) -> Result<(), StorageError> {
db::with_db(|conn| {
let msg_id: i64 = conn.query_row(
r#"
SELECT id FROM messages
WHERE storage_owner = ?1 AND external_user = ?2 AND message_time = ?3
ORDER BY id DESC LIMIT 1
"#,
params![storage_owner, external_user, message_time],
|row| row.get(0),
)?;
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as i64;
let tx = conn.unchecked_transaction()?;
tx.execute(
r#"
INSERT OR IGNORE INTO reactions (message_id, user_id, reaction, created_at)
VALUES (?1, ?2, ?3, ?4)
"#,
params![msg_id, user_id, reaction, now],
)?;
sync::record_event(
&tx,
storage_owner,
EntityType::Message,
msg_id,
Operation::Upsert,
)?;
tx.commit()?;
Ok(())
})
}
pub fn remove_reaction(
storage_owner: i64,
external_user: i64,
message_time: i64,
user_id: i64,
reaction: &str,
) -> Result<(), StorageError> {
db::with_db(|conn| {
let msg_id: i64 = conn.query_row(
r#"
SELECT id FROM messages
WHERE storage_owner = ?1 AND external_user = ?2 AND message_time = ?3
ORDER BY id DESC LIMIT 1
"#,
params![storage_owner, external_user, message_time],
|row| row.get(0),
)?;
let tx = conn.unchecked_transaction()?;
tx.execute(
"DELETE FROM reactions WHERE message_id = ?1 AND user_id = ?2 AND reaction = ?3",
params![msg_id, user_id, reaction],
)?;
sync::record_event(
&tx,
storage_owner,
EntityType::Message,
msg_id,
Operation::Upsert,
)?;
tx.commit()?;
Ok(())
})
}
pub fn add_message(
send_time: u128,
storage_owner_is_sender: bool,
storage_owner: i64,
external_user: i64,
message: &str,
height: i64,
reply_to: Option<i64>,
) {
let message_time = match i64::try_from(send_time) {
Ok(v) => v,
Err(_) => {
log!("Failed to store message: send_time out of range for i64 ({send_time})");
return;
}
};
if let Err(e) = db::with_db(|conn| {
let tx = conn.unchecked_transaction()?;
tx.execute(
r#"
INSERT INTO messages (
storage_owner, external_user, message_time, content,
sent_by_self, message_state, height, reply_to
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
"#,
params![
storage_owner,
external_user,
message_time,
message,
if storage_owner_is_sender {
1_i64
} else {
0_i64
},
MessageState::Sending.as_str(),
height,
reply_to,
],
)?;
let msg_id = tx.last_insert_rowid();
sync::record_event(
&tx,
storage_owner,
EntityType::Message,
msg_id,
Operation::Upsert,
)?;
tx.commit()?;
Ok(())
}) {
log!("Failed to insert message into sqlite: {}", e);
return;
}
let mut contact = crate::users::contact::Contact::new(external_user);
contact.set_last_message_at(message_time);
crate::util::chats_util::mod_user(storage_owner, &contact);
}
pub fn change_message_state(
timestamp: i64,
storage_owner: i64,
external_user: i64,
new_state: MessageState,
) -> std::io::Result<()> {
db::with_db(|conn| {
let current: Option<String> = match conn.query_row(
r#"
SELECT message_state
FROM messages
WHERE storage_owner = ?1 AND external_user = ?2 AND message_time = ?3
ORDER BY id DESC LIMIT 1
"#,
params![storage_owner, external_user, timestamp],
|row| row.get(0),
) {
Ok(state) => Some(state),
Err(rusqlite::Error::QueryReturnedNoRows) => None,
Err(e) => return Err(e.into()),
};
let Some(current_state_raw) = current else {
return Ok(());
};
let upgraded = MessageState::from_str(&current_state_raw)
.upgrade(new_state)
.as_str()
.to_string();
let tx = conn.unchecked_transaction()?;
tx.execute(
r#"
UPDATE messages
SET message_state = ?1
WHERE id = (
SELECT id FROM messages
WHERE storage_owner = ?2 AND external_user = ?3 AND message_time = ?4
ORDER BY id DESC LIMIT 1
)
"#,
params![upgraded, storage_owner, external_user, timestamp],
)?;
let msg_id: i64 = tx.query_row(
"SELECT id FROM messages WHERE storage_owner = ?1 AND external_user = ?2 AND message_time = ?3 ORDER BY id DESC LIMIT 1",
params![storage_owner, external_user, timestamp],
|row| row.get(0),
)?;
sync::record_event(&tx, storage_owner, EntityType::Message, msg_id, Operation::Upsert)?;
tx.commit()?;
Ok(())
})
.map_err(|e: StorageError| std::io::Error::new(std::io::ErrorKind::Other, e.to_string()))
}
fn load_reactions(
conn: &rusqlite::Connection,
msg_ids: &[i64],
) -> std::collections::HashMap<i64, Vec<StoredReaction>> {
if msg_ids.is_empty() {
return std::collections::HashMap::new();
}
let placeholders: Vec<String> = msg_ids
.iter()
.enumerate()
.map(|(i, _)| format!("?{}", i + 1))
.collect();
let query = format!(
"SELECT message_id, reaction, user_id FROM reactions WHERE message_id IN ({}) ORDER BY created_at ASC, id ASC",
placeholders.join(", ")
);
let mut map: std::collections::HashMap<i64, Vec<StoredReaction>> =
std::collections::HashMap::new();
if let Ok(mut stmt) = conn.prepare(&query) {
let params: Vec<&dyn rusqlite::types::ToSql> = msg_ids
.iter()
.map(|id| id as &dyn rusqlite::types::ToSql)
.collect();
if let Ok(rows) = stmt.query_map(params.as_slice(), |row| {
Ok((
row.get::<_, i64>(0)?,
StoredReaction {
reaction: row.get(1)?,
user_id: row.get(2)?,
},
))
}) {
for row in rows.flatten() {
map.entry(row.0).or_default().push(row.1);
}
}
}
map
}
pub fn get_messages(
storage_owner: i64,
external_user: i64,
loaded_messages: i64,
amount: i64,
) -> Vec<StoredMessage> {
if amount <= 0 || loaded_messages < 0 {
return Vec::new();
}
match db::with_db(|conn| {
let mut stmt = conn.prepare(
r#"
SELECT id, message_time, content, sent_by_self, message_state, height, reply_to, edited_count
FROM messages
WHERE storage_owner = ?1 AND external_user = ?2 AND deleted_by_external = 0
ORDER BY message_time DESC, id DESC
LIMIT ?3 OFFSET ?4
"#,
)?;
let rows = stmt.query_map(
params![storage_owner, external_user, amount, loaded_messages],
|row| {
Ok(StoredMessage {
id: row.get(0)?,
external_user,
message_time: row.get(1)?,
content: row.get(2)?,
sent_by_self: row.get::<_, i64>(3)? != 0,
message_state: row.get(4)?,
height: row.get(5).unwrap_or(0),
reply_to: row.get(6).ok().flatten(),
edited: row.get::<_, i64>(7).unwrap_or(0) > 0,
reactions: Vec::new(),
})
},
)?;
let mut out = Vec::new();
for row in rows {
match row {
Ok(msg) => out.push(msg),
Err(e) => log!("Failed to read row from sqlite: {}", e),
}
}
let msg_ids: Vec<i64> = out.iter().map(|m| m.id).collect();
let reaction_map = load_reactions(conn, &msg_ids);
for msg in &mut out {
msg.reactions = reaction_map.get(&msg.id).cloned().unwrap_or_default();
}
Ok(out)
}) {
Ok(v) => v,
Err(e) => {
log!("Failed to query messages: {}", e);
Vec::new()
}
}
}
pub fn get_message(
storage_owner: i64,
message_time: i64,
external_user: Option<i64>,
) -> Result<Option<StoredMessage>, StorageError> {
db::with_db(|conn| {
let mut stmt = conn.prepare(
r#"
SELECT id, message_time, content, sent_by_self, message_state, height,
reply_to, edited_count, external_user
FROM messages
WHERE storage_owner = ?1
AND message_time = ?2
AND deleted_by_external = 0
AND (?3 IS NULL OR external_user = ?3)
ORDER BY id DESC
"#,
)?;
let rows = stmt.query_map(params![storage_owner, message_time, external_user], |row| {
Ok(StoredMessage {
id: row.get(0)?,
message_time: row.get(1)?,
content: row.get(2)?,
sent_by_self: row.get::<_, i64>(3)? != 0,
message_state: row.get(4)?,
height: row.get(5).unwrap_or(0),
reply_to: row.get(6).ok().flatten(),
edited: row.get::<_, i64>(7).unwrap_or(0) > 0,
external_user: row.get(8)?,
reactions: Vec::new(),
})
})?;
let messages: Vec<StoredMessage> = rows.collect::<Result<_, _>>()?;
if messages.is_empty() {
return Ok(None);
}
if external_user.is_none()
&& messages
.iter()
.map(|message| message.external_user)
.collect::<std::collections::HashSet<_>>()
.len()
> 1
{
return Ok(None);
}
let mut message = messages.into_iter().next().expect("checked non-empty");
let reaction_map = load_reactions(conn, &[message.id]);
message.reactions = reaction_map.get(&message.id).cloned().unwrap_or_default();
Ok(Some(message))
})
}
pub fn get_messages_by_ids(storage_owner: i64, ids: &[i64]) -> Vec<StoredMessage> {
if ids.is_empty() {
return Vec::new();
}
let wanted: std::collections::HashSet<i64> = ids.iter().copied().collect();
// A journal id uniquely identifies a row. Load all messages for this owner and retain only
// those ids; this keeps reaction hydration identical to normal message loading.
match db::with_db(|conn| {
let mut stmt = conn.prepare("SELECT id, message_time, content, sent_by_self, message_state, height, reply_to, edited_count, external_user FROM messages WHERE storage_owner = ?1 AND deleted_by_external = 0")?;
let rows = stmt.query_map([storage_owner], |row| {
let external_user: i64 = row.get(8)?;
Ok(StoredMessage {
id: row.get(0)?,
external_user,
message_time: row.get(1)?,
content: row.get(2)?,
sent_by_self: row.get::<_, i64>(3)? != 0,
message_state: row.get(4)?,
height: row.get(5).unwrap_or(0),
reply_to: row.get(6).ok().flatten(),
edited: row.get::<_, i64>(7).unwrap_or(0) > 0,
reactions: Vec::new(),
})
})?;
let mut messages = Vec::new();
for row in rows {
let message = row?;
if wanted.contains(&message.id) {
messages.push(message);
}
}
let reaction_map = load_reactions(conn, &messages.iter().map(|m| m.id).collect::<Vec<_>>());
for message in &mut messages {
message.reactions = reaction_map.get(&message.id).cloned().unwrap_or_default();
}
Ok(messages)
}) {
Ok(messages) => messages,
Err(e) => {
log!("Failed to query messages by id: {}", e);
Vec::new()
}
}
}
pub fn get_all_messages(storage_owner: i64) -> Vec<StoredMessage> {
let ids = match db::with_db(|conn| {
let mut stmt = conn.prepare(
"SELECT id FROM messages WHERE storage_owner = ?1 AND deleted_by_external = 0",
)?;
Ok(stmt
.query_map([storage_owner], |row| row.get::<_, i64>(0))?
.collect::<Result<Vec<_>, _>>()?)
}) {
Ok(ids) => ids,
Err(e) => {
log!("Failed to query all messages: {}", e);
return Vec::new();
}
};
get_messages_by_ids(storage_owner, &ids)
}
#[cfg(test)]
mod tests {
use super::MessageState;
#[test]
fn upgrade_prefers_highest_state() {
assert_eq!(
MessageState::Sending.upgrade(MessageState::Sent),
MessageState::Sent
);
assert_eq!(
MessageState::Sent.upgrade(MessageState::Received),
MessageState::Received
);
assert_eq!(
MessageState::Received.upgrade(MessageState::Read),
MessageState::Read
);
}
#[test]
fn from_str_is_case_insensitive() {
assert_eq!(MessageState::from_str("READ"), MessageState::Read);
assert_eq!(MessageState::from_str("received"), MessageState::Received);
assert_eq!(MessageState::from_str("Sent"), MessageState::Sent);
assert_eq!(MessageState::from_str("unknown"), MessageState::Sending);
}
}