[Upd] mtp update

This commit is contained in:
Alex Emmet 2026-07-20 15:28:15 +02:00
commit f82500ea7d
24 changed files with 2535 additions and 1800 deletions

599
Cargo.lock generated

File diff suppressed because it is too large Load diff

View file

@ -1,3 +1,19 @@
[workspace] [workspace]
members = ["iota-storage", "client", "iota-auth", "other-iota", "iota-updater", "iota-terms", "iota-state", "iota-cli", "iota-core", "omikron-connector", "web-server", "web-ui", "iota-logger", "iota-util"] members = [
"iota-storage",
"iota-connection",
"client",
"iota-auth",
"other-iota",
"iota-updater",
"iota-terms",
"iota-state",
"iota-cli",
"iota-core",
"omikron-connector",
"web-server",
"web-ui",
"iota-logger",
"iota-util",
]
resolver = "3" resolver = "3"

View file

@ -5,6 +5,7 @@ edition = "2024"
[dependencies] [dependencies]
mtp = { git = "https://git.methanium.net/Methanium/mtp.git", features = ["client"] } mtp = { git = "https://git.methanium.net/Methanium/mtp.git", features = ["client"] }
iota-connection = { path = "../iota-connection" }
iota-logger = { path = "../iota-logger" } iota-logger = { path = "../iota-logger" }
iota-util = { path = "../iota-util" } iota-util = { path = "../iota-util" }
iota-storage = { path = "../iota-storage" } iota-storage = { path = "../iota-storage" }

View file

@ -1,160 +1,22 @@
use dashmap::DashMap; use dashmap::DashMap;
use iota_connection::message_common::*;
use iota_connection::message_handlers;
use iota_logger::{log_cv_in, log_cv_out, log_t}; use iota_logger::{log_cv_in, log_cv_out, log_t};
use iota_state::SHUTDOWN; use iota_state::SHUTDOWN;
use iota_storage::users::contact::Contact; use iota_storage::util::chat_files::{self, MessageState, change_message_state};
use iota_storage::util::chat_files::{MessageState, change_message_state};
use iota_storage::util::chats_util::{get_user, mod_user};
use iota_storage::util::communities_util::CommunitiesUtil;
use iota_storage::util::config_util::CONFIG; use iota_storage::util::config_util::CONFIG;
use iota_storage::util::e2ee_storage::{self, ChatSecretQuery, StoredChatSecret}; use iota_storage::util::e2ee_storage::{self, StoredChatSecret};
use iota_storage::util::{chat_files, chats_util};
use iota_util::crypto_helper::keyring_from_base64; use iota_util::crypto_helper::keyring_from_base64;
use iota_util::crypto_util::{self}; use iota_util::crypto_util::{self};
use iota_util::file_util::{get_children, load_file, save_file};
use mtp::client::{Receiver, Sender}; use mtp::client::{Receiver, Sender};
use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue}; use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue};
use mtp::type_map::TypeMap;
use std::sync::Arc; use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH}; use std::time::Duration;
use tokio::sync::{Mutex, RwLock, mpsc, watch}; use tokio::sync::{Mutex, RwLock, mpsc, watch};
use tokio::task::JoinHandle; use tokio::task::JoinHandle;
use uuid::Uuid; use uuid::Uuid;
fn typed_container(items: Vec<(DataType, DataValue)>) -> DataValue {
use mtp::type_map::{DataTypeId, TypeMap};
let tm = TypeMap::latest();
DataValue::Container(
items
.into_iter()
.filter_map(|(dt, dv)| tm.data_id_enum(dt).map(|id| (DataTypeId(id), dv)))
.collect(),
)
}
fn data_string(cv: &CommunicationValue, dt: DataType) -> Option<String> {
cv.get_data(dt)
.as_str()
.map(|s| s.to_string())
.or_else(|| cv.get_data(dt).as_number().map(|n| n.to_string()))
.or_else(|| cv.get_data(dt).as_signed_number().map(|n| n.to_string()))
}
fn data_i64(cv: &CommunicationValue, dt: DataType) -> Option<i64> {
cv.get_data(dt)
.as_number()
.and_then(|n| i64::try_from(n).ok())
.or_else(|| {
cv.get_data(dt)
.as_signed_number()
.and_then(|n| i64::try_from(n).ok())
})
.or_else(|| cv.get_data(dt).as_str().and_then(|s| s.parse::<i64>().ok()))
}
#[derive(Debug, Clone)]
struct ChatSecretRecipient {
user_id: String,
encrypted_secret: Vec<u8>,
kem_ciphertext: Vec<u8>,
}
fn recipient_from_value(value: &DataValue) -> Option<ChatSecretRecipient> {
let tm = TypeMap::latest();
let user_id = value
.get_field(DataType::UserId.to_id(&tm))?
.as_str()
.map(|s| s.to_string())
.or_else(|| {
value
.get_field(DataType::UserId.to_id(&tm))?
.as_number()
.map(|n| n.to_string())
})?;
let encrypted_secret = value
.get_field(DataType::EncryptedSecret.to_id(&tm))?
.as_bytes()?;
let kem_ciphertext = value
.get_field(DataType::KemCiphertext.to_id(&tm))?
.as_bytes()?;
Some(ChatSecretRecipient {
user_id,
encrypted_secret,
kem_ciphertext,
})
}
fn chat_secret_recipients(cv: &CommunicationValue) -> Option<Vec<ChatSecretRecipient>> {
let recipients = cv.get_data(DataType::Recipients).as_array()?;
let parsed = recipients
.iter()
.map(recipient_from_value)
.collect::<Option<Vec<_>>>()?;
if parsed.is_empty() {
None
} else {
Some(parsed)
}
}
fn set_chat_secret_cv_for_recipient(
source: &CommunicationValue,
recipient: &ChatSecretRecipient,
) -> CommunicationValue {
let recipient_value = typed_container(vec![
(DataType::UserId, DataValue::Str(recipient.user_id.clone())),
(
DataType::EncryptedSecret,
DataValue::Bytes(recipient.encrypted_secret.clone()),
),
(
DataType::KemCiphertext,
DataValue::Bytes(recipient.kem_ciphertext.clone()),
),
]);
CommunicationValue::new(CommunicationType::SetChatSecret)
.with_id(source.get_id())
.with_sender(source.get_sender())
.with_receiver(recipient.user_id.parse::<u64>().unwrap_or(0))
.add_typed_default(DataType::ChatId, source.get_data(DataType::ChatId).clone())
.add_typed_default(
DataType::SecretId,
source.get_data(DataType::SecretId).clone(),
)
.add_typed_default(
DataType::VersionNumber,
source.get_data(DataType::VersionNumber).clone(),
)
.add_typed_default(
DataType::WrappingScheme,
source.get_data(DataType::WrappingScheme).clone(),
)
.add_typed_default(
DataType::CreatedAt,
source.get_data(DataType::CreatedAt).clone(),
)
.add_typed_default(
DataType::Recipients,
DataValue::Array(vec![recipient_value]),
)
}
fn now_millis_i64() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as i64
}
fn error_response(request: &CommunicationValue, ty: CommunicationType) -> CommunicationValue {
CommunicationValue::new(ty)
.with_id(request.get_id())
.with_receiver(request.get_sender())
}
// ============================================================================ // ============================================================================
// Waiting Task System // Waiting Task System
// ============================================================================ // ============================================================================
@ -223,7 +85,7 @@ impl ClientConnection {
} }
if let Some(sender) = self.sender.read().await.as_ref() { if let Some(sender) = self.sender.read().await.as_ref() {
sender.close(); sender.close().await;
} }
*self.sender.write().await = None; *self.sender.write().await = None;
@ -235,12 +97,9 @@ impl ClientConnection {
async fn handle_ping(self: Arc<Self>, cv: CommunicationValue) { async fn handle_ping(self: Arc<Self>, cv: CommunicationValue) {
// Update our ping if provided // Update our ping if provided
if let DataValue::SignedNumber(last_ping) = cv.get_data(DataType::LastPing) { if let DataValue::SignedNumber(last_ping) = cv.get_data(DataType::LastPing) {
let current = SystemTime::now() let current = now_millis_i64();
.duration_since(UNIX_EPOCH)
.unwrap()
.as_millis();
let mut ping_guard = self.ping.write().await; let mut ping_guard = self.ping.write().await;
*ping_guard = current as i64 - *last_ping as i64; *ping_guard = current - *last_ping as i64;
} }
// Send pong response // Send pong response
@ -324,71 +183,9 @@ impl ClientConnection {
} }
if cv.is_type(CommunicationType::GetChatSecret) { if cv.is_type(CommunicationType::GetChatSecret) {
let Some(user_id) = data_string(&cv, DataType::UserId) else { self.send_message(&message_handlers::handle_get_chat_secret(&cv))
self.send_message(&error_response(&cv, CommunicationType::ErrorInvalidData))
.await; .await;
return; return;
};
let sender_id = cv.get_sender().to_string();
if user_id != sender_id {
self.send_message(&error_response(&cv, CommunicationType::ErrorNotFound))
.await;
return;
}
let Some(chat_id) = data_string(&cv, DataType::ChatId) else {
self.send_message(&error_response(&cv, CommunicationType::ErrorInvalidData))
.await;
return;
};
match e2ee_storage::get_chat_secret(ChatSecretQuery {
user_id,
chat_id,
secret_id: data_string(&cv, DataType::SecretId),
}) {
Ok(Some(record)) => {
let response = CommunicationValue::new(CommunicationType::ChatSecretResponse)
.with_id(cv.get_id())
.with_receiver(cv.get_sender())
.add_typed_default(DataType::UserId, DataValue::Str(record.user_id))
.add_typed_default(DataType::ChatId, DataValue::Str(record.chat_id))
.add_typed_default(DataType::SecretId, DataValue::Str(record.secret_id))
.add_typed_default(
DataType::VersionNumber,
DataValue::SignedNumber(record.version as i128),
)
.add_typed_default(
DataType::EncryptedSecret,
DataValue::Bytes(record.encrypted_secret),
)
.add_typed_default(
DataType::KemCiphertext,
DataValue::Bytes(record.kem_ciphertext),
)
.add_typed_default(
DataType::WrappingScheme,
DataValue::Str(record.wrapping_scheme),
)
.add_typed_default(
DataType::CreatedAt,
DataValue::SignedNumber(record.created_at as i128),
)
.add_typed_default(
DataType::UpdatedAt,
DataValue::SignedNumber(record.updated_at as i128),
);
self.send_message(&response).await;
}
Ok(None) => {
self.send_message(&error_response(&cv, CommunicationType::ErrorNotSet))
.await;
}
Err(_) => {
self.send_message(&error_response(&cv, CommunicationType::ErrorInvalidData))
.await;
}
}
return;
} }
if cv.is_type(CommunicationType::ChatSecretForward) { if cv.is_type(CommunicationType::ChatSecretForward) {
@ -434,126 +231,20 @@ impl ClientConnection {
} }
if cv.is_type(CommunicationType::CreateApp) { if cv.is_type(CommunicationType::CreateApp) {
let sender_id = cv.get_sender() as i64; self.send_message(&message_handlers::handle_create_app(&cv))
let app_identifier = cv .await;
.get_data(DataType::AppIdentifier)
.as_str()
.unwrap_or("")
.to_string();
let app_public_key = cv
.get_data(DataType::AppPublicKey)
.as_str()
.unwrap_or("")
.to_string();
if !app_identifier.is_empty() && !app_public_key.is_empty() {
if let Some(mut user) = iota_storage::users::user_manager::get_user(sender_id) {
if !user.trusted_apps.contains_key(&app_identifier) {
user.trusted_apps.insert(app_identifier, app_public_key);
iota_storage::users::user_manager::update_user(user);
}
}
}
let res = CommunicationValue::new(CommunicationType::CreateApp)
.with_id(cv.get_id())
.with_receiver(sender_id as u64);
self.send_message(&res).await;
return; return;
} }
if cv.is_type(CommunicationType::DeleteApp) { if cv.is_type(CommunicationType::DeleteApp) {
let sender_id = cv.get_sender() as i64; self.send_message(&message_handlers::handle_delete_app(&cv))
let app_identifier = cv .await;
.get_data(DataType::AppIdentifier)
.as_str()
.unwrap_or("")
.to_string();
if !app_identifier.is_empty() {
if let Some(mut user) = iota_storage::users::user_manager::get_user(sender_id) {
if user.trusted_apps.contains_key(&app_identifier) {
user.trusted_apps.remove(&app_identifier);
iota_storage::users::user_manager::update_user(user);
}
}
}
let res = CommunicationValue::new(CommunicationType::DeleteApp)
.with_id(cv.get_id())
.with_receiver(sender_id as u64);
self.send_message(&res).await;
return; return;
} }
if cv.is_type(CommunicationType::ClientConnected) { if cv.is_type(CommunicationType::ClientConnected) {
let user_id = cv.get_data(DataType::UserId).as_number().unwrap_or(0) as i64; self.send_message(&message_handlers::handle_client_connected(&cv))
let _session_id = cv.get_data(DataType::SessionId).as_number().unwrap_or(0) as i64; .await;
let contacts = chats_util::get_users(user_id);
let mut contacts_array = Vec::new();
for (i, contact) in contacts.iter().enumerate() {
let mut contact_container = Vec::new();
contact_container.push((
DataType::UserId,
DataValue::SignedNumber(contact.user_id as i128),
));
contact_container.push((
DataType::LastMessageAt,
DataValue::SignedNumber(contact.last_message_at.unwrap_or(0) as i128),
));
if let Some(ref name) = contact.user_name {
contact_container.push((DataType::Username, DataValue::Str(name.clone())));
}
let amount = if i < 10 { 20 } else { 1 };
let messages = chat_files::get_messages(user_id, contact.user_id, 0, amount);
let mut msg_array = Vec::new();
for m in &messages {
let mut msg_container = Vec::new();
msg_container.push((
DataType::SendTime,
DataValue::SignedNumber(m.message_time as i128),
));
msg_container.push((DataType::Content, DataValue::Str(m.content.clone())));
msg_container.push((DataType::MessageState, DataValue::Str(m.message_state.clone())));
msg_container.push((DataType::Height, DataValue::SignedNumber(m.height as i128)));
msg_container.push((
DataType::SenderId,
DataValue::UnsignedNumber(if m.sent_by_self {
user_id as u128
} else {
contact.user_id as u128
}),
));
msg_array.push(typed_container(msg_container));
if msg_array.len() == 1 {
let sender_id = if m.sent_by_self {
user_id
} else {
contact.user_id
};
let mut last_msg = Vec::new();
last_msg.push((DataType::Content, DataValue::Str(m.content.clone())));
last_msg.push((
DataType::SenderId,
DataValue::SignedNumber(sender_id as i128),
));
contact_container.push((DataType::LastMessage, typed_container(last_msg)));
}
}
contact_container.push((DataType::Messages, DataValue::Array(msg_array)));
contacts_array.push(typed_container(contact_container));
}
let resp = CommunicationValue::new(CommunicationType::ClientConnected)
.with_id(cv.get_id())
.add_typed_default(DataType::Contacts, DataValue::Array(contacts_array));
self.send_message(&resp).await;
return; return;
} }
@ -562,32 +253,32 @@ impl ClientConnection {
// ************************************************ // // ************************************************ //
if cv.is_type(CommunicationType::MessageState) { if cv.is_type(CommunicationType::MessageState) {
let sender_id = &cv.get_sender(); message_handlers::handle_message_state(&cv);
let receiver_id = match cv.get_data(DataType::ChatPartnerId).as_number() { return;
Some(id) => id, }
_ => return,
};
// Parse send_time robustly: accept numeric or string, fallback to current time if cv.is_type(CommunicationType::MessageEdit) {
let send_time_val = cv.get_data(DataType::SendTime); self.send_message(&message_handlers::handle_message_edit(&cv))
let now_i64 = SystemTime::now() .await;
.duration_since(UNIX_EPOCH) return;
.unwrap_or_default() }
.as_millis() as i64;
let timestamp_i64 = if let Some(n) = send_time_val.as_number() {
n as i64
} else if let Some(s) = send_time_val.as_str() {
s.parse::<i64>().unwrap_or(now_i64)
} else {
now_i64
};
let _ = chat_files::change_message_state( if cv.is_type(CommunicationType::MessageReactionAdd) {
timestamp_i64, self.send_message(&message_handlers::handle_message_reaction(&cv, true))
receiver_id as i64, .await;
*sender_id as i64, return;
MessageState::from_str(cv.get_data(DataType::MessageState).as_str().unwrap_or("")), }
);
if cv.is_type(CommunicationType::MessageReactionRemove) {
self.send_message(&message_handlers::handle_message_reaction(&cv, false))
.await;
return;
}
if cv.is_type(CommunicationType::MessageDeleteLive) {
self.send_message(&message_handlers::handle_message_delete(&cv))
.await;
return;
} }
// Incoming storsed message: store for the recipient, attempt local delivery, notify sender. // Incoming storsed message: store for the recipient, attempt local delivery, notify sender.
@ -603,10 +294,7 @@ impl ClientConnection {
// parse send_time safely (number or string), fallback to now // parse send_time safely (number or string), fallback to now
let send_time_val = cv.get_data(DataType::SendTime); let send_time_val = cv.get_data(DataType::SendTime);
let now_i64 = SystemTime::now() let now_i64 = now_millis_i64();
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as i64;
let timestamp = if let Some(n) = send_time_val.as_number() { let timestamp = if let Some(n) = send_time_val.as_number() {
n as i64 n as i64
} else if let Some(s) = send_time_val.as_str() { } else if let Some(s) = send_time_val.as_str() {
@ -732,175 +420,38 @@ impl ClientConnection {
} }
if cv.is_type(CommunicationType::MessagesGet) { if cv.is_type(CommunicationType::MessagesGet) {
let my_id = cv.get_sender(); self.send_message(&message_handlers::handle_messages_get(&cv))
let partner_id = cv.get_data(DataType::UserId).as_number().unwrap_or(0); .await;
let offset = cv.get_data(DataType::Offset).as_number().unwrap_or(0);
let amount = cv.get_data(DataType::Amount).as_number().unwrap_or(0);
let messages = chat_files::get_messages(
my_id as i64,
partner_id as i64,
offset as i64,
amount as i64,
);
let mut msg_array: Vec<DataValue> = Vec::new();
for m in &messages {
let sender_id: i64 = if m.sent_by_self {
my_id as i64
} else {
if let Some(n) = cv.get_data(DataType::ChatPartnerId).as_number() {
n as i64
} else if let Some(s) = cv.get_data(DataType::ChatPartnerId).as_str() {
s.parse::<i64>().unwrap_or(partner_id as i64)
} else {
partner_id as i64
}
};
let mut container = Vec::new();
container.push((
DataType::SendTime,
DataValue::SignedNumber(m.message_time as i128),
));
container.push((DataType::Content, DataValue::Str(m.content.clone())));
container.push((
DataType::SenderId,
DataValue::SignedNumber(sender_id as i128),
));
container.push((DataType::MessageState, DataValue::Str(m.message_state.clone())));
container.push((DataType::Height, DataValue::SignedNumber(m.height as i128)));
container.push((
DataType::SenderId,
DataValue::UnsignedNumber(if m.sent_by_self {
my_id as u128
} else {
partner_id as u128
}),
));
msg_array.push(typed_container(container));
}
let resp = CommunicationValue::new(CommunicationType::MessagesGet)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(DataType::Messages, DataValue::Array(msg_array));
self.send_message(&resp).await;
return; return;
} }
if cv.is_type(CommunicationType::GetChats) { if cv.is_type(CommunicationType::GetChats) {
let user_id = cv.get_sender(); self.send_message(&message_handlers::handle_get_chats(&cv))
let users = chats_util::get_users(user_id as i64); .await;
let mut user_array = Vec::new();
for user in users {
let mut container = Vec::new();
container.push((
DataType::UserId,
DataValue::SignedNumber(user.user_id as i128),
));
if let Some(name) = user.user_name {
container.push((DataType::Username, DataValue::Str(name)));
}
if let Some(ts) = user.last_message_at {
container.push((DataType::LastMessageAt, DataValue::SignedNumber(ts as i128)));
}
user_array.push(typed_container(container));
}
let resp = CommunicationValue::new(CommunicationType::GetChats)
.with_id(cv.get_id())
.with_receiver(user_id)
.add_typed_default(DataType::UserIds, DataValue::Array(user_array));
self.send_message(&resp).await;
return; return;
} }
if cv.is_type(CommunicationType::AddConversation) { if cv.is_type(CommunicationType::AddConversation) {
let user_id = cv.get_sender(); self.send_message(&message_handlers::handle_add_conversation(&cv))
let other_id = match cv.get_data(DataType::ChatPartnerId).as_number() { .await;
Some(n) => n as i64,
None => cv
.get_data(DataType::ChatPartnerId)
.as_str()
.unwrap_or("0")
.parse()
.unwrap_or(0),
};
let mut contact = get_user(user_id as i64, other_id).unwrap_or(Contact::new(other_id));
if let Some(name) = cv.get_data(DataType::ChatPartnerName).as_str() {
contact.user_name = Some(name.to_string());
}
contact.set_last_message_at(
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_millis() as i64,
);
mod_user(user_id as i64, &contact);
let resp = CommunicationValue::new(CommunicationType::AddConversation)
.with_id(cv.get_id())
.with_receiver(user_id);
self.send_message(&resp).await;
return; return;
} }
if cv.is_type(CommunicationType::AddCommunity) { if cv.is_type(CommunicationType::AddCommunity) {
CommunitiesUtil::add_community( self.send_message(&message_handlers::handle_add_community(&cv))
cv.get_sender() as i64, .await;
cv.get_data(DataType::CommunityAddress)
.as_str()
.unwrap()
.to_string(),
cv.get_data(DataType::CommunityTitle)
.as_str()
.unwrap()
.to_string(),
cv.get_data(DataType::Position)
.as_str()
.unwrap()
.to_string(),
);
let resp = CommunicationValue::new(CommunicationType::AddCommunity)
.with_id(cv.get_id())
.with_receiver(cv.get_sender());
self.send_message(&resp).await;
return; return;
} }
if cv.is_type(CommunicationType::GetCommunities) { if cv.is_type(CommunicationType::GetCommunities) {
let mut comm_array = Vec::new(); self.send_message(&message_handlers::handle_get_communities(&cv))
for c in CommunitiesUtil::get_communities(cv.get_sender() as i64) { .await;
let mut container: Vec<(DataType, DataValue)> = Vec::new();
container.push((
DataType::CommunityAddress,
DataValue::Str(c.address.clone()),
));
container.push((DataType::CommunityTitle, DataValue::Str(c.title.clone())));
container.push((DataType::Position, DataValue::Str(c.position.clone())));
comm_array.push(typed_container(container));
}
let resp = CommunicationValue::new(CommunicationType::GetCommunities)
.with_id(cv.get_id())
.with_receiver(cv.get_sender())
.add_typed_default(DataType::Communities, DataValue::Array(comm_array));
self.send_message(&resp).await;
return; return;
} }
if cv.is_type(CommunicationType::RemoveCommunity) { if cv.is_type(CommunicationType::RemoveCommunity) {
CommunitiesUtil::remove_community( self.send_message(&message_handlers::handle_remove_community(&cv))
cv.get_sender() as i64, .await;
cv.get_data(DataType::CommunityAddress)
.as_str()
.unwrap()
.to_string(),
);
let resp = CommunicationValue::new(CommunicationType::RemoveCommunity)
.with_id(cv.get_id())
.with_receiver(cv.get_sender());
self.send_message(&resp).await;
return; return;
} }
@ -909,10 +460,11 @@ impl ClientConnection {
let settings_name = cv.get_data(DataType::SettingsName).as_str().unwrap(); let settings_name = cv.get_data(DataType::SettingsName).as_str().unwrap();
let settings_value = cv.get_data(DataType::Payload).as_str().unwrap(); let settings_value = cv.get_data(DataType::Payload).as_str().unwrap();
save_file( let _ = iota_storage::util::settings::save(
&format!("users/{}/settings/", my_id), my_id as i64,
&format!("{}.settings", settings_name), iota_storage::util::settings::GLOBAL_SESSION_ID,
&settings_value, settings_name,
settings_value,
); );
let response = CommunicationValue::new(CommunicationType::SettingsSave) let response = CommunicationValue::new(CommunicationType::SettingsSave)
@ -926,10 +478,14 @@ impl ClientConnection {
if cv.is_type(CommunicationType::SettingsLoad) { if cv.is_type(CommunicationType::SettingsLoad) {
let my_id = cv.get_sender(); let my_id = cv.get_sender();
let settings_name = cv.get_data(DataType::SettingsName).as_string().unwrap(); let settings_name = cv.get_data(DataType::SettingsName).as_string().unwrap();
let settings_value_str = load_file( let settings_value_str = iota_storage::util::settings::load(
&format!("users/{}/settings/", my_id), my_id as i64,
&format!("{}.settings", settings_name), iota_storage::util::settings::GLOBAL_SESSION_ID,
); &settings_name,
)
.ok()
.flatten()
.unwrap_or_default();
let response = CommunicationValue::new(CommunicationType::SettingsLoad) let response = CommunicationValue::new(CommunicationType::SettingsLoad)
.with_id(cv.get_id()) .with_id(cv.get_id())
.with_receiver(my_id) .with_receiver(my_id)
@ -942,15 +498,12 @@ impl ClientConnection {
if cv.is_type(CommunicationType::SettingsList) { if cv.is_type(CommunicationType::SettingsList) {
let my_id = cv.get_sender(); let my_id = cv.get_sender();
let settings = get_children(&format!("users/{}/settings/", my_id)); let settings = iota_storage::util::settings::list(
let mut settings_json = Vec::new(); my_id as i64,
for s in settings { iota_storage::util::settings::GLOBAL_SESSION_ID,
let s = s.replace(".settings", ""); )
if s.is_empty() { .unwrap_or_default();
continue; let settings_json = settings.into_iter().map(DataValue::Str).collect();
}
let _ = settings_json.push(DataValue::Str(s));
}
let response = CommunicationValue::new(CommunicationType::SettingsList) let response = CommunicationValue::new(CommunicationType::SettingsList)
.with_id(cv.get_id()) .with_id(cv.get_id())
.with_receiver(my_id) .with_receiver(my_id)
@ -997,7 +550,7 @@ impl ClientConnection {
if !sender.is_open() { if !sender.is_open() {
drop(sender_guard); drop(sender_guard);
if let Some(sender) = self.sender.write().await.take() { if let Some(sender) = self.sender.write().await.take() {
sender.close(); sender.close().await;
} }
return Err("connection closed".to_string()); return Err("connection closed".to_string());
} }

View file

@ -0,0 +1,9 @@
[package]
name = "iota-connection"
version = "0.1.0"
edition = "2024"
[dependencies]
iota-storage = { path = "../iota-storage" }
iota-util = { path = "../iota-util" }
mtp = { git = "https://git.methanium.net/Methanium/mtp.git" }

View file

@ -0,0 +1,36 @@
use mtp::codec::CommunicationValue;
use std::future::Future;
use std::time::Duration;
/// Unified interface for all connection types (Omikron, Direct, future modes).
///
/// Provides the common messaging API that the rest of the codebase uses,
/// regardless of whether the connection goes through Omikron or is direct.
pub trait ConnectionHandler: Send + Sync {
/// Send a message to the remote end.
fn send_message(
&self,
cv: &CommunicationValue,
) -> impl Future<Output = Result<(), String>> + Send;
/// Send a message and wait for a correlated response.
///
/// The implementation correlates requests/responses by message ID and
/// enforces the given `timeout`. Returns an error on timeout or if the
/// connection drops while waiting.
fn await_response(
&self,
cv: &CommunicationValue,
timeout: Option<Duration>,
) -> impl Future<Output = Result<CommunicationValue, String>> + Send;
/// Returns `true` when the connection is alive and ready for traffic.
fn is_connected(&self) -> impl Future<Output = bool> + Send;
/// Returns `true` when the connection has completed identification /
/// registration and is fully operational.
fn is_identified(&self) -> impl Future<Output = bool> + Send;
/// Gracefully tear down the connection.
fn stop(&self) -> impl Future<Output = ()> + Send;
}

View file

@ -0,0 +1,3 @@
pub mod connection_handler;
pub mod message_common;
pub mod message_handlers;

View file

@ -0,0 +1,137 @@
use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue};
use mtp::type_map::TypeMap;
use std::time::{SystemTime, UNIX_EPOCH};
pub fn typed_container(items: Vec<(DataType, DataValue)>) -> DataValue {
use mtp::type_map::{DataTypeId, TypeMap};
let tm = TypeMap::latest();
DataValue::Container(
items
.into_iter()
.filter_map(|(dt, dv)| tm.data_id_enum(dt).map(|id| (DataTypeId(id), dv)))
.collect(),
)
}
pub fn data_string(cv: &CommunicationValue, dt: DataType) -> Option<String> {
cv.get_data(dt)
.as_str()
.map(|s| s.to_string())
.or_else(|| cv.get_data(dt).as_number().map(|n| n.to_string()))
.or_else(|| cv.get_data(dt).as_signed_number().map(|n| n.to_string()))
}
pub fn data_i64(cv: &CommunicationValue, dt: DataType) -> Option<i64> {
cv.get_data(dt)
.as_number()
.and_then(|n| i64::try_from(n).ok())
.or_else(|| {
cv.get_data(dt)
.as_signed_number()
.and_then(|n| i64::try_from(n).ok())
})
.or_else(|| cv.get_data(dt).as_str().and_then(|s| s.parse::<i64>().ok()))
}
#[derive(Debug, Clone)]
pub struct ChatSecretRecipient {
pub user_id: String,
pub encrypted_secret: Vec<u8>,
pub kem_ciphertext: Vec<u8>,
}
pub fn recipient_from_value(value: &DataValue) -> Option<ChatSecretRecipient> {
let tm = TypeMap::latest();
let user_id = value
.get_field(DataType::UserId.try_to_id(&tm)?)?
.as_str()
.map(|s| s.to_string())
.or_else(|| {
value
.get_field(DataType::UserId.try_to_id(&tm)?)?
.as_number()
.map(|n| n.to_string())
})?;
let encrypted_secret = value
.get_field(DataType::EncryptedSecret.try_to_id(&tm)?)?
.as_bytes()?;
let kem_ciphertext = value
.get_field(DataType::KemCiphertext.try_to_id(&tm)?)?
.as_bytes()?;
Some(ChatSecretRecipient {
user_id,
encrypted_secret,
kem_ciphertext,
})
}
pub fn chat_secret_recipients(cv: &CommunicationValue) -> Option<Vec<ChatSecretRecipient>> {
let recipients = cv.get_data(DataType::Recipients).as_array()?;
let parsed = recipients
.iter()
.map(recipient_from_value)
.collect::<Option<Vec<_>>>()?;
if parsed.is_empty() {
None
} else {
Some(parsed)
}
}
pub fn set_chat_secret_cv_for_recipient(
source: &CommunicationValue,
recipient: &ChatSecretRecipient,
) -> CommunicationValue {
let recipient_value = typed_container(vec![
(DataType::UserId, DataValue::Str(recipient.user_id.clone())),
(
DataType::EncryptedSecret,
DataValue::Bytes(recipient.encrypted_secret.clone()),
),
(
DataType::KemCiphertext,
DataValue::Bytes(recipient.kem_ciphertext.clone()),
),
]);
CommunicationValue::new(CommunicationType::SetChatSecret)
.with_id(source.get_id())
.with_sender(source.get_sender())
.with_receiver(recipient.user_id.parse::<u64>().unwrap_or(0))
.add_typed_default(DataType::ChatId, source.get_data(DataType::ChatId).clone())
.add_typed_default(
DataType::SecretId,
source.get_data(DataType::SecretId).clone(),
)
.add_typed_default(
DataType::VersionNumber,
source.get_data(DataType::VersionNumber).clone(),
)
.add_typed_default(
DataType::WrappingScheme,
source.get_data(DataType::WrappingScheme).clone(),
)
.add_typed_default(
DataType::CreatedAt,
source.get_data(DataType::CreatedAt).clone(),
)
.add_typed_default(
DataType::Recipients,
DataValue::Array(vec![recipient_value]),
)
}
pub fn now_millis_i64() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as i64
}
pub fn error_response(request: &CommunicationValue, ty: CommunicationType) -> CommunicationValue {
CommunicationValue::new(ty)
.with_id(request.get_id())
.with_receiver(request.get_sender())
}

View file

@ -0,0 +1,766 @@
use crate::message_common::*;
use iota_storage::util::chat_files::{self, MessageState};
use iota_storage::util::chats_util::{self, get_user, mod_user};
use iota_storage::util::communities_util::CommunitiesUtil;
use iota_storage::util::e2ee_storage::{self, ChatSecretQuery};
use iota_storage::util::settings;
use mtp::codec::{CommunicationType, CommunicationValue, DataType, DataValue};
pub struct MessageMutation {
pub sender_id: i64,
pub partner_id: i64,
pub send_time: i64,
}
pub fn message_mutation(cv: &CommunicationValue) -> Result<MessageMutation, CommunicationValue> {
let sender_id = i64::try_from(cv.get_sender())
.map_err(|_| error_response(cv, CommunicationType::ErrorInvalidData))?;
let partner_id = data_i64(cv, DataType::ChatPartnerId)
.filter(|id| *id > 0)
.ok_or_else(|| error_response(cv, CommunicationType::ErrorInvalidData))?;
let send_time = data_i64(cv, DataType::SendTime)
.filter(|time| *time > 0)
.ok_or_else(|| error_response(cv, CommunicationType::ErrorInvalidData))?;
Ok(MessageMutation {
sender_id,
partner_id,
send_time,
})
}
pub fn success_response(cv: &CommunicationValue) -> CommunicationValue {
error_response(cv, CommunicationType::Success)
}
pub fn handle_message_edit(cv: &CommunicationValue) -> CommunicationValue {
let mutation = match message_mutation(cv) {
Ok(mutation) => mutation,
Err(response) => return response,
};
let Some(content) = cv.get_data(DataType::Content).as_str() else {
return error_response(cv, CommunicationType::ErrorInvalidData);
};
match chat_files::edit_message(
mutation.sender_id,
mutation.partner_id,
mutation.send_time,
mutation.sender_id,
content,
) {
Ok(()) => success_response(cv),
Err(_) => error_response(cv, CommunicationType::ErrorNotFound),
}
}
pub fn handle_message_reaction(cv: &CommunicationValue, add: bool) -> CommunicationValue {
let mutation = match message_mutation(cv) {
Ok(mutation) => mutation,
Err(response) => return response,
};
let Some(reaction) = cv.get_data(DataType::Reaction).as_str() else {
return error_response(cv, CommunicationType::ErrorInvalidData);
};
if reaction.is_empty() || reaction.len() > 64 {
return error_response(cv, CommunicationType::ErrorInvalidData);
}
let result = if add {
chat_files::add_reaction(
mutation.sender_id,
mutation.partner_id,
mutation.send_time,
mutation.sender_id,
reaction,
)
} else {
chat_files::remove_reaction(
mutation.sender_id,
mutation.partner_id,
mutation.send_time,
mutation.sender_id,
reaction,
)
};
match result {
Ok(()) => success_response(cv),
Err(_) => error_response(cv, CommunicationType::ErrorNotFound),
}
}
pub fn handle_message_delete(cv: &CommunicationValue) -> CommunicationValue {
let mutation = match message_mutation(cv) {
Ok(mutation) => mutation,
Err(response) => return response,
};
match chat_files::delete_message(mutation.sender_id, mutation.partner_id, mutation.send_time) {
Ok(()) => success_response(cv),
Err(_) => error_response(cv, CommunicationType::ErrorNotFound),
}
}
fn stored_message_value(
message: &chat_files::StoredMessage,
storage_owner: i64,
partner_id: i64,
) -> DataValue {
let mut fields = vec![
(
DataType::SendTime,
DataValue::SignedNumber(message.message_time as i128),
),
(DataType::Content, DataValue::Str(message.content.clone())),
(
DataType::MessageState,
DataValue::Str(message.message_state.clone()),
),
(
DataType::Height,
DataValue::SignedNumber(message.height as i128),
),
(
DataType::SenderId,
DataValue::UnsignedNumber(if message.sent_by_self {
storage_owner as u128
} else {
partner_id as u128
}),
),
];
if message.edited {
fields.push((DataType::Edited, DataValue::Bool(true)));
}
if let Some(reply_to) = message.reply_to {
fields.push((
DataType::ReplyId,
DataValue::UnsignedNumber(reply_to as u64 as u128),
));
}
if !message.reactions.is_empty() {
let reactions = message
.reactions
.iter()
.map(|reaction| {
typed_container(vec![
(
DataType::Reaction,
DataValue::Str(reaction.reaction.clone()),
),
(
DataType::SenderId,
DataValue::SignedNumber(reaction.user_id as i128),
),
])
})
.collect();
fields.push((DataType::Reactions, DataValue::Array(reactions)));
}
typed_container(fields)
}
pub fn handle_get_chat_secret(cv: &CommunicationValue) -> CommunicationValue {
let Some(user_id) = data_string(cv, DataType::UserId) else {
return error_response(cv, CommunicationType::ErrorInvalidData);
};
if user_id != cv.get_sender().to_string() {
return error_response(cv, CommunicationType::ErrorNotFound);
}
let Some(chat_id) = data_string(cv, DataType::ChatId) else {
return error_response(cv, CommunicationType::ErrorInvalidData);
};
match e2ee_storage::get_chat_secret(ChatSecretQuery {
user_id,
chat_id,
secret_id: data_string(cv, DataType::SecretId),
}) {
Ok(Some(record)) => CommunicationValue::new(CommunicationType::ChatSecretResponse)
.with_id(cv.get_id())
.with_receiver(cv.get_sender())
.add_typed_default(DataType::UserId, DataValue::Str(record.user_id))
.add_typed_default(DataType::ChatId, DataValue::Str(record.chat_id))
.add_typed_default(DataType::SecretId, DataValue::Str(record.secret_id))
.add_typed_default(
DataType::VersionNumber,
DataValue::SignedNumber(record.version as i128),
)
.add_typed_default(
DataType::EncryptedSecret,
DataValue::Bytes(record.encrypted_secret),
)
.add_typed_default(
DataType::KemCiphertext,
DataValue::Bytes(record.kem_ciphertext),
)
.add_typed_default(
DataType::WrappingScheme,
DataValue::Str(record.wrapping_scheme),
)
.add_typed_default(
DataType::CreatedAt,
DataValue::SignedNumber(record.created_at as i128),
)
.add_typed_default(
DataType::UpdatedAt,
DataValue::SignedNumber(record.updated_at as i128),
),
Ok(None) => error_response(cv, CommunicationType::ErrorNotSet),
Err(_) => error_response(cv, CommunicationType::ErrorInvalidData),
}
}
pub fn handle_create_app(cv: &CommunicationValue) -> CommunicationValue {
let sender_id = cv.get_sender() as i64;
let app_identifier = cv
.get_data(DataType::AppIdentifier)
.as_str()
.unwrap_or("")
.to_string();
let app_public_key = cv
.get_data(DataType::AppPublicKey)
.as_str()
.unwrap_or("")
.to_string();
if !app_identifier.is_empty() && !app_public_key.is_empty() {
if let Some(mut user) = iota_storage::users::user_manager::get_user(sender_id) {
if !user.trusted_apps.contains_key(&app_identifier) {
user.trusted_apps.insert(app_identifier, app_public_key);
iota_storage::users::user_manager::update_user(user);
}
}
}
CommunicationValue::new(CommunicationType::CreateApp)
.with_id(cv.get_id())
.with_receiver(sender_id as u64)
}
pub fn handle_delete_app(cv: &CommunicationValue) -> CommunicationValue {
let sender_id = cv.get_sender() as i64;
let app_identifier = cv
.get_data(DataType::AppIdentifier)
.as_str()
.unwrap_or("")
.to_string();
if !app_identifier.is_empty() {
if let Some(mut user) = iota_storage::users::user_manager::get_user(sender_id) {
if user.trusted_apps.contains_key(&app_identifier) {
user.trusted_apps.remove(&app_identifier);
iota_storage::users::user_manager::update_user(user);
}
}
}
CommunicationValue::new(CommunicationType::DeleteApp)
.with_id(cv.get_id())
.with_receiver(sender_id as u64)
}
pub fn handle_client_connected(cv: &CommunicationValue) -> CommunicationValue {
let user_id = cv.get_data(DataType::UserId).as_number().unwrap_or(0) as i64;
let contacts = chats_util::get_users(user_id);
let mut contacts_array = Vec::new();
for (i, contact) in contacts.iter().enumerate() {
let mut contact_container = Vec::new();
contact_container.push((
DataType::UserId,
DataValue::SignedNumber(contact.user_id as i128),
));
contact_container.push((
DataType::LastMessageAt,
DataValue::SignedNumber(contact.last_message_at.unwrap_or(0) as i128),
));
if let Some(ref name) = contact.user_name {
contact_container.push((DataType::Username, DataValue::Str(name.clone())));
}
let amount = if i < 10 { 20 } else { 1 };
let messages = chat_files::get_messages(user_id, contact.user_id, 0, amount);
let mut msg_array = Vec::new();
for m in &messages {
msg_array.push(stored_message_value(m, user_id, contact.user_id));
if msg_array.len() == 1 {
let sender_id = if m.sent_by_self {
user_id
} else {
contact.user_id
};
let mut last_msg = Vec::new();
last_msg.push((DataType::Content, DataValue::Str(m.content.clone())));
last_msg.push((
DataType::SenderId,
DataValue::SignedNumber(sender_id as i128),
));
contact_container.push((DataType::LastMessage, typed_container(last_msg)));
}
}
contact_container.push((DataType::Messages, DataValue::Array(msg_array)));
contacts_array.push(typed_container(contact_container));
}
CommunicationValue::new(CommunicationType::ClientConnected)
.with_id(cv.get_id())
.add_typed_default(DataType::Contacts, DataValue::Array(contacts_array))
}
pub fn handle_message_state(cv: &CommunicationValue) {
let sender_id = &cv.get_sender();
let receiver_id = match cv.get_data(DataType::ChatPartnerId).as_number() {
Some(id) => id,
_ => return,
};
let timestamp_i64 = if let Some(n) = cv.get_data(DataType::SendTime).as_number() {
n as i64
} else if let Some(s) = cv.get_data(DataType::SendTime).as_str() {
s.parse::<i64>().unwrap_or_else(|_| now_millis_i64())
} else {
now_millis_i64()
};
let _ = chat_files::change_message_state(
timestamp_i64,
receiver_id as i64,
*sender_id as i64,
MessageState::from_str(cv.get_data(DataType::MessageState).as_str().unwrap_or("")),
);
}
pub fn handle_messages_get(cv: &CommunicationValue) -> CommunicationValue {
let my_id = cv.get_sender();
let partner_id = cv.get_data(DataType::UserId).as_number().unwrap_or(0);
let offset = cv.get_data(DataType::Offset).as_number().unwrap_or(0);
let amount = cv.get_data(DataType::Amount).as_number().unwrap_or(0);
let messages = chat_files::get_messages(
my_id as i64,
partner_id as i64,
offset as i64,
amount as i64,
);
let mut msg_array: Vec<DataValue> = Vec::new();
for m in &messages {
msg_array.push(stored_message_value(m, my_id as i64, partner_id as i64));
}
CommunicationValue::new(CommunicationType::MessagesGet)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(DataType::Messages, DataValue::Array(msg_array))
}
pub fn handle_get_chats(cv: &CommunicationValue) -> CommunicationValue {
let user_id = cv.get_sender();
let users = chats_util::get_users(user_id as i64);
let mut user_array = Vec::new();
for user in users {
let mut container = Vec::new();
container.push((
DataType::UserId,
DataValue::SignedNumber(user.user_id as i128),
));
if let Some(name) = user.user_name {
container.push((DataType::Username, DataValue::Str(name)));
}
if let Some(ts) = user.last_message_at {
container.push((DataType::LastMessageAt, DataValue::SignedNumber(ts as i128)));
}
user_array.push(typed_container(container));
}
CommunicationValue::new(CommunicationType::GetChats)
.with_id(cv.get_id())
.with_receiver(user_id)
.add_typed_default(DataType::UserIds, DataValue::Array(user_array))
}
pub fn handle_add_conversation(cv: &CommunicationValue) -> CommunicationValue {
let user_id = cv.get_sender();
let other_id = match cv.get_data(DataType::ChatPartnerId).as_number() {
Some(n) => n as i64,
None => cv
.get_data(DataType::ChatPartnerId)
.as_str()
.unwrap_or("0")
.parse()
.unwrap_or(0),
};
let mut contact = get_user(user_id as i64, other_id)
.unwrap_or(iota_storage::users::contact::Contact::new(other_id));
if let Some(name) = cv.get_data(DataType::ChatPartnerName).as_str() {
contact.user_name = Some(name.to_string());
}
contact.set_last_message_at(now_millis_i64());
mod_user(user_id as i64, &contact);
CommunicationValue::new(CommunicationType::AddConversation)
.with_id(cv.get_id())
.with_receiver(user_id)
}
pub fn handle_add_community(cv: &CommunicationValue) -> CommunicationValue {
CommunitiesUtil::add_community(
cv.get_sender() as i64,
cv.get_data(DataType::CommunityAddress)
.as_str()
.unwrap()
.to_string(),
cv.get_data(DataType::CommunityTitle)
.as_str()
.unwrap()
.to_string(),
cv.get_data(DataType::Position)
.as_str()
.unwrap()
.to_string(),
);
CommunicationValue::new(CommunicationType::AddCommunity)
.with_id(cv.get_id())
.with_receiver(cv.get_sender())
}
pub fn handle_get_communities(cv: &CommunicationValue) -> CommunicationValue {
let mut comm_array = Vec::new();
for c in CommunitiesUtil::get_communities(cv.get_sender() as i64) {
let mut container: Vec<(DataType, DataValue)> = Vec::new();
container.push((
DataType::CommunityAddress,
DataValue::Str(c.address.clone()),
));
container.push((DataType::CommunityTitle, DataValue::Str(c.title.clone())));
container.push((DataType::Position, DataValue::Str(c.position.clone())));
comm_array.push(typed_container(container));
}
CommunicationValue::new(CommunicationType::GetCommunities)
.with_id(cv.get_id())
.with_receiver(cv.get_sender())
.add_typed_default(DataType::Communities, DataValue::Array(comm_array))
}
pub fn handle_remove_community(cv: &CommunicationValue) -> CommunicationValue {
CommunitiesUtil::remove_community(
cv.get_sender() as i64,
cv.get_data(DataType::CommunityAddress)
.as_str()
.unwrap()
.to_string(),
);
CommunicationValue::new(CommunicationType::RemoveCommunity)
.with_id(cv.get_id())
.with_receiver(cv.get_sender())
}
pub fn handle_global_settings_save(cv: &CommunicationValue) -> CommunicationValue {
let my_id = cv.get_sender();
let Some(settings_value) = cv.get_data(DataType::Payload).as_str() else {
return CommunicationValue::new(CommunicationType::ErrorInvalidData)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Message,
DataValue::Str("Missing settings payload".to_string()),
);
};
if settings::save_global(my_id as i64, settings_value).is_err() {
return error_response(cv, CommunicationType::ErrorInvalidData);
}
let mut response = CommunicationValue::new(CommunicationType::GlobalSettingsSave)
.with_receiver(my_id)
.with_id(cv.get_id());
if let Some(session_id) = cv.get_data(DataType::SessionId).as_number() {
response = response.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
);
}
response
}
pub fn handle_global_settings_load(cv: &CommunicationValue) -> CommunicationValue {
let my_id = cv.get_sender();
let Ok(settings_value) = settings::load_global(my_id as i64) else {
return error_response(cv, CommunicationType::ErrorInvalidData);
};
let Some(settings_value_str) = settings_value else {
let mut response = CommunicationValue::new(CommunicationType::ErrorNotFound)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Path,
DataValue::Str("global.settings".to_string()),
);
if let Some(session_id) = cv.get_data(DataType::SessionId).as_number() {
response = response.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
);
}
return response;
};
let mut response = CommunicationValue::new(CommunicationType::GlobalSettingsLoad)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(DataType::Payload, DataValue::Str(settings_value_str));
if let Some(session_id) = cv.get_data(DataType::SessionId).as_number() {
response = response.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
);
}
response
}
pub fn handle_settings_save(
cv: &CommunicationValue,
_expected_session_id: i128,
) -> CommunicationValue {
let my_id = cv.get_sender();
let Some(session_id) = cv.get_data(DataType::SessionId).as_number() else {
return CommunicationValue::new(CommunicationType::ErrorInvalidData)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Message,
DataValue::Str("Missing session_id".to_string()),
);
};
if session_id == 0 || session_id > 1_000_000 {
return CommunicationValue::new(CommunicationType::ErrorInvalidData)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Message,
DataValue::Str("Invalid session_id".to_string()),
)
.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
);
};
let Some(settings_name) = cv.get_data(DataType::SettingsName).as_str() else {
return CommunicationValue::new(CommunicationType::ErrorInvalidData)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Message,
DataValue::Str("Missing settings_name".to_string()),
)
.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
);
};
let Some(settings_value) = cv.get_data(DataType::Payload).as_str() else {
return CommunicationValue::new(CommunicationType::ErrorInvalidData)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Message,
DataValue::Str("Missing settings payload".to_string()),
)
.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
);
};
if !settings_name
.chars()
.all(|c| c.is_alphanumeric() || c == '_' || c == '-' || c == '.')
|| settings_name.contains("..")
{
return CommunicationValue::new(CommunicationType::ErrorInvalidData)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Message,
DataValue::Str("Invalid settings_name".to_string()),
)
.add_typed_default(
DataType::SettingsName,
DataValue::Str(settings_name.to_string()),
)
.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
);
}
if settings::save(
my_id as i64,
session_id as i64,
settings_name,
settings_value,
)
.is_err()
{
return error_response(cv, CommunicationType::ErrorInvalidData);
}
CommunicationValue::new(CommunicationType::SettingsSave)
.with_receiver(my_id)
.with_id(cv.get_id())
.add_typed_default(
DataType::SettingsName,
DataValue::Str(settings_name.to_string()),
)
.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
)
}
pub fn handle_settings_load(
cv: &CommunicationValue,
_expected_session_id: i128,
) -> CommunicationValue {
let my_id = cv.get_sender();
let Some(session_id) = cv.get_data(DataType::SessionId).as_number() else {
return CommunicationValue::new(CommunicationType::ErrorInvalidData)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Message,
DataValue::Str("Missing session_id".to_string()),
);
};
if session_id == 0 || session_id > 1_000_000 {
return CommunicationValue::new(CommunicationType::ErrorInvalidData)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Message,
DataValue::Str("Invalid session_id".to_string()),
)
.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
);
}
let Some(settings_name) = cv.get_data(DataType::SettingsName).as_str() else {
return CommunicationValue::new(CommunicationType::ErrorInvalidData)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Message,
DataValue::Str("Missing settings_name".to_string()),
)
.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
);
};
if !settings_name
.chars()
.all(|c| c.is_alphanumeric() || c == '_' || c == '-' || c == '.')
|| settings_name.contains("..")
{
return CommunicationValue::new(CommunicationType::ErrorInvalidData)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Message,
DataValue::Str("Invalid settings_name".to_string()),
)
.add_typed_default(
DataType::SettingsName,
DataValue::Str(settings_name.to_string()),
)
.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
);
}
let Ok(settings_value) = settings::load(my_id as i64, session_id as i64, settings_name) else {
return error_response(cv, CommunicationType::ErrorInvalidData);
};
let Some(settings_value_str) = settings_value else {
return CommunicationValue::new(CommunicationType::ErrorNotFound)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::SettingsName,
DataValue::Str(settings_name.to_string()),
)
.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
);
};
CommunicationValue::new(CommunicationType::SettingsLoad)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(DataType::Payload, DataValue::Str(settings_value_str))
.add_typed_default(
DataType::SettingsName,
DataValue::Str(settings_name.to_string()),
)
.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
)
}
pub fn handle_settings_list(
cv: &CommunicationValue,
_expected_session_id: i128,
) -> CommunicationValue {
let my_id = cv.get_sender();
let Some(session_id) = cv.get_data(DataType::SessionId).as_number() else {
return CommunicationValue::new(CommunicationType::ErrorInvalidData)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Message,
DataValue::Str("Missing session_id".to_string()),
);
};
if session_id == 0 || session_id > 1_000_000 {
return CommunicationValue::new(CommunicationType::ErrorInvalidData)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(
DataType::Message,
DataValue::Str("Invalid session_id".to_string()),
)
.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
);
}
let Ok(settings) = settings::list(my_id as i64, session_id as i64) else {
return error_response(cv, CommunicationType::ErrorInvalidData);
};
let settings_json = settings.into_iter().map(DataValue::Str).collect();
CommunicationValue::new(CommunicationType::SettingsList)
.with_id(cv.get_id())
.with_receiver(my_id)
.add_typed_default(DataType::Settings, DataValue::Array(settings_json))
.add_typed_default(
DataType::SessionId,
DataValue::SignedNumber(session_id as i128),
)
}

View file

@ -88,6 +88,9 @@ async fn main() {
if let Err(_) = user_manager::load_users().await { if let Err(_) = user_manager::load_users().await {
log_t!("user_load_failed"); log_t!("user_load_failed");
} }
if let Err(e) = iota_storage::util::settings::migrate_legacy_files() {
log!("Failed to migrate legacy settings: {}", e);
}
let mut sb = "".to_string(); let mut sb = "".to_string();
@ -102,7 +105,11 @@ async fn main() {
} }
log!( log!(
"IOTA ID: {}", "IOTA ID: {}",
CONFIG.load().iota_id.map(|id| id.to_string()).unwrap_or_else(|| "N/A".to_string()) CONFIG
.load()
.iota_id
.map(|id| id.to_string())
.unwrap_or_else(|| "N/A".to_string())
); );
log!("User IDS: {}", sb); log!("User IDS: {}", sb);
@ -121,7 +128,7 @@ async fn main() {
sb1 = sb1 + ","; sb1 = sb1 + ",";
} }
log!("Community IDS: {}", sb1); */ log!("Community IDS: {}", sb1); */
let _port = CONFIG.load().port; let port = CONFIG.load().port;
let mut _ip = "0.0.0.0".to_string(); let mut _ip = "0.0.0.0".to_string();
for iface in pnet::datalink::interfaces() { for iface in pnet::datalink::interfaces() {
let iface: NetworkInterface = iface; let iface: NetworkInterface = iface;
@ -133,17 +140,6 @@ async fn main() {
} }
} }
} }
/* Community port activation is used for activating the port for communities.
* Code is currently commented because communities have not been implemented yet.
if start(port).await {
log_t!("community_active", ip, port.to_string());
} else {
if port < 1024 {
log_t!("community_start_error_admin", port.to_string());
} else {
log_t!("community_start_error", port.to_string());
}
} */
if !has_dir("web") { if !has_dir("web") {
download_and_extract_zip( download_and_extract_zip(
"https://omega.tensamin.net/api/download/iota_frontend", "https://omega.tensamin.net/api/download/iota_frontend",
@ -151,6 +147,9 @@ async fn main() {
) )
.await; .await;
} }
if !web_server::start(port).await {
log!("Failed to start the MTP web server on port {}", port);
}
let _ = omikron::omikron_connection::get_omikron_connection().await; let _ = omikron::omikron_connection::get_omikron_connection().await;
log_t!("setup_completed"); log_t!("setup_completed");
@ -162,7 +161,9 @@ async fn main() {
if OMIKRON_CONNECTION.has_auth_failure().await { if OMIKRON_CONNECTION.has_auth_failure().await {
if let Some(reason) = OMIKRON_CONNECTION.get_auth_failure().await { if let Some(reason) = OMIKRON_CONNECTION.get_auth_failure().await {
log!("Authentication failed: {}", reason); log!("Authentication failed: {}", reason);
log!("Use /reconnect to try again or /regenerate private-key to create a new key pair"); log!(
"Use /reconnect to try again or /regenerate private-key to create a new key pair"
);
OMIKRON_CONNECTION.clear_auth_failure().await; OMIKRON_CONNECTION.clear_auth_failure().await;
} }
} }

View file

@ -3,8 +3,8 @@ use crate::util::db;
use base64::{Engine as _, engine::general_purpose::STANDARD}; use base64::{Engine as _, engine::general_purpose::STANDARD};
use iota_util::crypto_helper::{self, hex_hash, keyring_from_base64, public_key_bundle_to_base64}; use iota_util::crypto_helper::{self, hex_hash, keyring_from_base64, public_key_bundle_to_base64};
use iota_util::file_util::{load_file, save_file}; use iota_util::file_util::{load_file, save_file};
use rusqlite::params;
use rand_core::{OsRng, RngCore}; use rand_core::{OsRng, RngCore};
use rusqlite::params;
pub fn add_user(user: UserProfile) { pub fn add_user(user: UserProfile) {
if let Err(e) = db::with_db(|conn| { if let Err(e) = db::with_db(|conn| {
@ -166,9 +166,8 @@ pub fn get_users() -> Vec<UserProfile> {
fn load_trusted_apps(user_id: i64) -> std::collections::HashMap<String, String> { fn load_trusted_apps(user_id: i64) -> std::collections::HashMap<String, String> {
match db::with_db(|conn| { match db::with_db(|conn| {
let mut stmt = conn.prepare( let mut stmt =
"SELECT app_id, app_secret FROM trusted_apps WHERE user_id = ?1", conn.prepare("SELECT app_id, app_secret FROM trusted_apps WHERE user_id = ?1")?;
)?;
let rows = stmt.query_map(params![user_id], |r| { let rows = stmt.query_map(params![user_id], |r| {
Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?)) Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))
})?; })?;
@ -191,7 +190,10 @@ fn load_trusted_apps(user_id: i64) -> std::collections::HashMap<String, String>
pub fn remove_user(user_id: i64) { pub fn remove_user(user_id: i64) {
if let Err(e) = db::with_db(|conn| { if let Err(e) = db::with_db(|conn| {
conn.execute("DELETE FROM trusted_apps WHERE user_id = ?1", params![user_id])?; conn.execute(
"DELETE FROM trusted_apps WHERE user_id = ?1",
params![user_id],
)?;
conn.execute("DELETE FROM users WHERE user_id = ?1", params![user_id])?; conn.execute("DELETE FROM users WHERE user_id = ?1", params![user_id])?;
Ok(()) Ok(())
}) { }) {

View file

@ -1,7 +1,7 @@
use crate::storage_error::StorageError;
use crate::util::db; use crate::util::db;
use iota_logger::log; use iota_logger::log;
use rusqlite::params; use rusqlite::params;
use crate::storage_error::StorageError;
#[derive(PartialEq, Debug, Clone)] #[derive(PartialEq, Debug, Clone)]
pub enum MessageState { pub enum MessageState {
@ -53,7 +53,13 @@ pub struct StoredMessage {
pub message_state: String, pub message_state: String,
pub height: i64, pub height: i64,
pub reply_to: Option<i64>, pub reply_to: Option<i64>,
pub reactions: Vec<String>, pub reactions: Vec<StoredReaction>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StoredReaction {
pub reaction: String,
pub user_id: i64,
} }
/* /*
@ -66,6 +72,48 @@ pub fn edit_message(
message_time: i64, message_time: i64,
editor_id: i64, editor_id: i64,
new_content: &str, 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> { ) -> Result<(), StorageError> {
db::with_db(|conn| { db::with_db(|conn| {
let msg = conn.query_row( let msg = conn.query_row(
@ -86,11 +134,16 @@ pub fn edit_message(
)?; )?;
let (msg_id, old_content, sent_by_self) = msg; let (msg_id, old_content, sent_by_self) = msg;
if sent_by_self != 1 { if require_sent_by_self && sent_by_self != 1 {
return Err(StorageError::Other( return Err(StorageError::Other(
"Only the original sender can edit this message".into(), "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() let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH) .duration_since(std::time::UNIX_EPOCH)
@ -118,7 +171,11 @@ pub fn edit_message(
}) })
} }
pub fn hard_delete_message(storage_owner: i64, external_user: i64, message_time: i64) -> Result<(), StorageError> { pub fn hard_delete_message(
storage_owner: i64,
external_user: i64,
message_time: i64,
) -> Result<(), StorageError> {
db::with_db(|conn| { db::with_db(|conn| {
let msg_id: i64 = conn.query_row( let msg_id: i64 = conn.query_row(
r#" r#"
@ -130,18 +187,79 @@ pub fn hard_delete_message(storage_owner: i64, external_user: i64, message_time:
|row| row.get(0), |row| row.get(0),
)?; )?;
conn.execute("DELETE FROM message_edits WHERE message_id = ?1", params![msg_id])?; conn.execute(
conn.execute("DELETE FROM reactions WHERE message_id = ?1", params![msg_id])?; "DELETE FROM message_edits WHERE message_id = ?1",
params![msg_id],
)?;
conn.execute(
"DELETE FROM reactions WHERE message_id = ?1",
params![msg_id],
)?;
conn.execute("DELETE FROM messages WHERE id = ?1", params![msg_id])?; conn.execute("DELETE FROM messages WHERE id = ?1", params![msg_id])?;
Ok(()) 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, * 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. * 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> { pub fn flag_deleted_by_external(
storage_owner: i64,
external_user: i64,
message_time: i64,
) -> Result<(), StorageError> {
db::with_db(|conn| { db::with_db(|conn| {
let affected = conn.execute( let affected = conn.execute(
r#" r#"
@ -163,7 +281,11 @@ pub fn flag_deleted_by_external(storage_owner: i64, external_user: i64, message_
* the UI still shows the "edited" indicator. Only the own user should * the UI still shows the "edited" indicator. Only the own user should
* call this. * call this.
*/ */
pub fn delete_edit_history(storage_owner: i64, external_user: i64, message_time: i64) -> Result<(), StorageError> { pub fn delete_edit_history(
storage_owner: i64,
external_user: i64,
message_time: i64,
) -> Result<(), StorageError> {
db::with_db(|conn| { db::with_db(|conn| {
let msg_id: i64 = conn.query_row( let msg_id: i64 = conn.query_row(
r#" r#"
@ -175,7 +297,10 @@ pub fn delete_edit_history(storage_owner: i64, external_user: i64, message_time:
|row| row.get(0), |row| row.get(0),
)?; )?;
conn.execute("DELETE FROM message_edits WHERE message_id = ?1", params![msg_id])?; conn.execute(
"DELETE FROM message_edits WHERE message_id = ?1",
params![msg_id],
)?;
Ok(()) Ok(())
}) })
} }
@ -339,26 +464,39 @@ pub fn change_message_state(
.map_err(|e: StorageError| std::io::Error::new(std::io::ErrorKind::Other, e.to_string())) .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<String>> { fn load_reactions(
conn: &rusqlite::Connection,
msg_ids: &[i64],
) -> std::collections::HashMap<i64, Vec<StoredReaction>> {
if msg_ids.is_empty() { if msg_ids.is_empty() {
return std::collections::HashMap::new(); return std::collections::HashMap::new();
} }
let placeholders: Vec<String> = msg_ids.iter().enumerate() let placeholders: Vec<String> = msg_ids
.iter()
.enumerate()
.map(|(i, _)| format!("?{}", i + 1)) .map(|(i, _)| format!("?{}", i + 1))
.collect(); .collect();
let query = format!( let query = format!(
"SELECT message_id, reaction || ':' || COUNT(*) FROM reactions WHERE message_id IN ({}) GROUP BY message_id, reaction", "SELECT message_id, reaction, user_id FROM reactions WHERE message_id IN ({}) ORDER BY created_at ASC, id ASC",
placeholders.join(", ") placeholders.join(", ")
); );
let mut map: std::collections::HashMap<i64, Vec<String>> = std::collections::HashMap::new(); let mut map: std::collections::HashMap<i64, Vec<StoredReaction>> =
std::collections::HashMap::new();
if let Ok(mut stmt) = conn.prepare(&query) { if let Ok(mut stmt) = conn.prepare(&query) {
let params: Vec<&dyn rusqlite::types::ToSql> = msg_ids.iter() let params: Vec<&dyn rusqlite::types::ToSql> = msg_ids
.iter()
.map(|id| id as &dyn rusqlite::types::ToSql) .map(|id| id as &dyn rusqlite::types::ToSql)
.collect(); .collect();
if let Ok(rows) = stmt.query_map(params.as_slice(), |row| { if let Ok(rows) = stmt.query_map(params.as_slice(), |row| {
Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)) Ok((
row.get::<_, i64>(0)?,
StoredReaction {
reaction: row.get(1)?,
user_id: row.get(2)?,
},
))
}) { }) {
for row in rows.flatten() { for row in rows.flatten() {
map.entry(row.0).or_default().push(row.1); map.entry(row.0).or_default().push(row.1);

View file

@ -64,6 +64,30 @@ fn db_file_path(db_name: &str) -> String {
fn run_migrations(pool: &r2d2::Pool<SqliteManager>) -> Result<(), StorageError> { fn run_migrations(pool: &r2d2::Pool<SqliteManager>) -> Result<(), StorageError> {
let conn = pool.get().map_err(|e| StorageError::Pool(e.to_string()))?; let conn = pool.get().map_err(|e| StorageError::Pool(e.to_string()))?;
run_migrations_on_connection(&conn)
}
/*
* Older builds could apply a schema change without advancing user_version.
* Check each added column so those databases can resume upgrading.
*/
fn add_column_if_missing(
conn: &Connection,
column: &str,
definition: &str,
) -> Result<(), StorageError> {
let mut statement =
conn.prepare("SELECT 1 FROM pragma_table_info('messages') WHERE name = ?1")?;
let exists = statement.exists([column])?;
if !exists {
conn.execute_batch(&format!("ALTER TABLE messages ADD COLUMN {definition};"))?;
}
Ok(())
}
fn run_migrations_on_connection(conn: &Connection) -> Result<(), StorageError> {
let current_version: i64 = conn let current_version: i64 = conn
.pragma_query_value(None, "user_version", |r| r.get(0)) .pragma_query_value(None, "user_version", |r| r.get(0))
.unwrap_or(0); .unwrap_or(0);
@ -78,8 +102,7 @@ fn run_migrations(pool: &r2d2::Pool<SqliteManager>) -> Result<(), StorageError>
message_time INTEGER NOT NULL, message_time INTEGER NOT NULL,
content TEXT NOT NULL, content TEXT NOT NULL,
sent_by_self INTEGER NOT NULL, sent_by_self INTEGER NOT NULL,
message_state TEXT NOT NULL, message_state TEXT NOT NULL
height INTEGER NOT NULL DEFAULT 0
); );
CREATE INDEX IF NOT EXISTS idx_messages_lookup CREATE INDEX IF NOT EXISTS idx_messages_lookup
ON messages (storage_owner, external_user, message_time DESC); ON messages (storage_owner, external_user, message_time DESC);
@ -129,29 +152,28 @@ fn run_migrations(pool: &r2d2::Pool<SqliteManager>) -> Result<(), StorageError>
} }
if current_version < 2 { if current_version < 2 {
conn.execute_batch( add_column_if_missing(conn, "height", "height INTEGER NOT NULL DEFAULT 0")?;
r#" conn.execute_batch("PRAGMA user_version = 2;")?;
ALTER TABLE messages ADD COLUMN height INTEGER NOT NULL DEFAULT 0;
PRAGMA user_version = 2;
"#,
)?;
} }
if current_version < 3 { if current_version < 3 {
conn.execute_batch( add_column_if_missing(conn, "reply_to", "reply_to INTEGER")?;
r#" conn.execute_batch("PRAGMA user_version = 3;")?;
ALTER TABLE messages ADD COLUMN reply_to INTEGER;
PRAGMA user_version = 3;
"#,
)?;
} }
if current_version < 4 { if current_version < 4 {
add_column_if_missing(
conn,
"edited_count",
"edited_count INTEGER NOT NULL DEFAULT 0",
)?;
add_column_if_missing(
conn,
"deleted_by_external",
"deleted_by_external INTEGER NOT NULL DEFAULT 0",
)?;
conn.execute_batch( conn.execute_batch(
r#" r#"
ALTER TABLE messages ADD COLUMN edited_count INTEGER NOT NULL DEFAULT 0;
ALTER TABLE messages ADD COLUMN deleted_by_external INTEGER NOT NULL DEFAULT 0;
CREATE TABLE IF NOT EXISTS message_edits ( CREATE TABLE IF NOT EXISTS message_edits (
id INTEGER PRIMARY KEY AUTOINCREMENT, id INTEGER PRIMARY KEY AUTOINCREMENT,
message_id INTEGER NOT NULL REFERENCES messages(id), message_id INTEGER NOT NULL REFERENCES messages(id),
@ -179,6 +201,24 @@ fn run_migrations(pool: &r2d2::Pool<SqliteManager>) -> Result<(), StorageError>
)?; )?;
} }
if current_version < 5 {
conn.execute_batch(
r#"
CREATE TABLE IF NOT EXISTS settings (
user_id INTEGER NOT NULL,
session_id INTEGER NOT NULL,
name TEXT NOT NULL,
payload TEXT NOT NULL,
PRIMARY KEY (user_id, session_id, name)
);
CREATE INDEX IF NOT EXISTS idx_settings_lookup
ON settings (user_id, session_id, name);
PRAGMA user_version = 5;
"#,
)?;
}
Ok(()) Ok(())
} }
@ -221,3 +261,40 @@ where
pub fn create_general_messages_db() -> Result<Arc<std::sync::Mutex<Connection>>, String> { pub fn create_general_messages_db() -> Result<Arc<std::sync::Mutex<Connection>>, String> {
create_shared_connection(DB_NAME, "") create_shared_connection(DB_NAME, "")
} }
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn resumes_migration_when_height_exists_before_its_version() -> Result<(), StorageError> {
let conn = Connection::open_in_memory()?;
conn.execute_batch(
r#"
CREATE TABLE messages (
id INTEGER PRIMARY KEY AUTOINCREMENT,
storage_owner INTEGER NOT NULL,
external_user INTEGER NOT NULL,
message_time INTEGER NOT NULL,
content TEXT NOT NULL,
sent_by_self INTEGER NOT NULL,
message_state TEXT NOT NULL
);
ALTER TABLE messages ADD COLUMN height INTEGER NOT NULL DEFAULT 0;
PRAGMA user_version = 1;
"#,
)?;
run_migrations_on_connection(&conn)?;
let version: i64 = conn.pragma_query_value(None, "user_version", |row| row.get(0))?;
assert_eq!(version, 5);
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")?;
assert!(statement.exists([column])?);
}
Ok(())
}
}

View file

@ -1,5 +1,5 @@
use crate::util::db; use crate::util::db;
use rusqlite::{params, OptionalExtension}; use rusqlite::{OptionalExtension, params};
use std::sync::{Arc, LazyLock, Mutex}; use std::sync::{Arc, LazyLock, Mutex};
pub type StorageError = String; pub type StorageError = String;
@ -226,9 +226,7 @@ fn chat_secret_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<StoredChatS
}) })
} }
fn pending_forward_from_row( fn pending_forward_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<PendingChatSecretForward> {
row: &rusqlite::Row<'_>,
) -> rusqlite::Result<PendingChatSecretForward> {
Ok(PendingChatSecretForward { Ok(PendingChatSecretForward {
recipient_user_id: row.get(0)?, recipient_user_id: row.get(0)?,
chat_id: row.get(1)?, chat_id: row.get(1)?,

View file

@ -4,3 +4,4 @@ pub mod communities_util;
pub mod config_util; pub mod config_util;
pub mod db; pub mod db;
pub mod e2ee_storage; pub mod e2ee_storage;
pub mod settings;

View file

@ -0,0 +1,163 @@
use crate::storage_error::StorageError;
use crate::util::db;
use iota_util::file_util::get_directory;
use rusqlite::{OptionalExtension, params};
use std::fs;
use std::path::Path;
pub const GLOBAL_SESSION_ID: i64 = 0;
const GLOBAL_SETTINGS_NAME: &str = "__global__";
pub fn save(user_id: i64, session_id: i64, name: &str, payload: &str) -> Result<(), StorageError> {
db::with_db(|conn| {
conn.execute(
"INSERT INTO settings (user_id, session_id, name, payload) VALUES (?1, ?2, ?3, ?4)\n ON CONFLICT(user_id, session_id, name) DO UPDATE SET payload = excluded.payload",
params![user_id, session_id, name, payload],
)?;
Ok(())
})
}
pub fn load(user_id: i64, session_id: i64, name: &str) -> Result<Option<String>, StorageError> {
db::with_db(|conn| {
conn.query_row(
"SELECT payload FROM settings WHERE user_id = ?1 AND session_id = ?2 AND name = ?3",
params![user_id, session_id, name],
|row| row.get(0),
)
.optional()
.map_err(StorageError::from)
})
}
pub fn list(user_id: i64, session_id: i64) -> Result<Vec<String>, StorageError> {
db::with_db(|conn| {
let mut statement = conn.prepare(
"SELECT name FROM settings WHERE user_id = ?1 AND session_id = ?2 ORDER BY name",
)?;
let rows = statement.query_map(params![user_id, session_id], |row| row.get(0))?;
rows.collect::<Result<Vec<String>, _>>()
.map_err(StorageError::from)
})
}
pub fn save_global(user_id: i64, payload: &str) -> Result<(), StorageError> {
save(user_id, GLOBAL_SESSION_ID, GLOBAL_SETTINGS_NAME, payload)
}
pub fn load_global(user_id: i64) -> Result<Option<String>, StorageError> {
load(user_id, GLOBAL_SESSION_ID, GLOBAL_SETTINGS_NAME)
}
pub fn migrate_legacy_files() -> Result<(), StorageError> {
let users_dir = Path::new(&get_directory()).join("users");
let Ok(users) = fs::read_dir(users_dir) else {
return Ok(());
};
for user_entry in users {
let user_entry = user_entry?;
let Ok(user_id) = user_entry.file_name().to_string_lossy().parse::<i64>() else {
continue;
};
let user_dir = user_entry.path();
migrate_file_if_missing(
user_id,
GLOBAL_SESSION_ID,
GLOBAL_SETTINGS_NAME,
&user_dir.join("global.settings"),
)?;
let settings_dir = user_dir.join("settings");
let Ok(settings_entries) = fs::read_dir(settings_dir) else {
continue;
};
for settings_entry in settings_entries {
let settings_entry = settings_entry?;
let path = settings_entry.path();
if path.is_file() {
if let Some(name) = setting_name(&path) {
migrate_file_if_missing(user_id, GLOBAL_SESSION_ID, &name, &path)?;
}
continue;
}
let Ok(session_id) = settings_entry.file_name().to_string_lossy().parse::<i64>() else {
continue;
};
let Ok(device_settings) = fs::read_dir(path) else {
continue;
};
for setting_entry in device_settings {
let setting_entry = setting_entry?;
let path = setting_entry.path();
if let Some(name) = setting_name(&path) {
migrate_file_if_missing(user_id, session_id, &name, &path)?;
}
}
}
}
Ok(())
}
fn setting_name(path: &Path) -> Option<String> {
(path.extension()?.to_str()? == "settings").then(|| {
path.file_stem()
.and_then(|name| name.to_str())
.unwrap_or_default()
.to_string()
})
}
fn migrate_file_if_missing(
user_id: i64,
session_id: i64,
name: &str,
path: &Path,
) -> Result<(), StorageError> {
if !path.is_file() || load(user_id, session_id, name)?.is_some() {
return Ok(());
}
let payload = fs::read_to_string(path)?;
save(user_id, session_id, name, &payload)
}
#[cfg(test)]
mod tests {
use super::*;
use rusqlite::Connection;
#[test]
fn settings_schema_supports_user_and_session_keys() -> Result<(), StorageError> {
let conn = Connection::open_in_memory()?;
conn.execute_batch(
"CREATE TABLE settings (
user_id INTEGER NOT NULL,
session_id INTEGER NOT NULL,
name TEXT NOT NULL,
payload TEXT NOT NULL,
PRIMARY KEY (user_id, session_id, name)
);",
)?;
conn.execute(
"INSERT INTO settings VALUES (?1, ?2, ?3, ?4)",
params![7, 11, "theme", "dark"],
)?;
conn.execute(
"INSERT INTO settings VALUES (?1, ?2, ?3, ?4)",
params![7, 12, "theme", "light"],
)?;
let payload: String = conn.query_row(
"SELECT payload FROM settings WHERE user_id = 7 AND session_id = 11 AND name = 'theme'",
[],
|row| row.get(0),
)?;
assert_eq!(payload, "dark");
Ok(())
}
}

View file

@ -4,8 +4,9 @@ version = "0.1.0"
edition = "2024" edition = "2024"
[dependencies] [dependencies]
mtp = { git = "https://git.methanium.net/Methanium/mtp.git" } mtp = { git = "https://git.methanium.net/Methanium/mtp.git", features = [
mtp-crypto = { git = "https://git.methanium.net/Methanium/mtp.git", features = ["pqc"] } "crypto"
] }
reqwest = "0.13.2" reqwest = "0.13.2"
tokio = { version = "1.50.0", features = ["full"] } tokio = { version = "1.50.0", features = ["full"] }

View file

@ -1,5 +1,5 @@
use base64::{Engine as _, engine::general_purpose::STANDARD}; use base64::{Engine as _, engine::general_purpose::STANDARD};
use mtp_crypto::{Keyring, PublicKeyBundle}; use mtp::crypto::{Keyring, PublicKeyBundle};
pub fn generate_keyring() -> Keyring { pub fn generate_keyring() -> Keyring {
Keyring::generate() Keyring::generate()
@ -24,5 +24,5 @@ pub fn public_key_bundle_from_base64(s: &str) -> Option<PublicKeyBundle> {
} }
pub fn hex_hash(input: &str) -> String { pub fn hex_hash(input: &str) -> String {
hex::encode(mtp_crypto::sha256(input.as_bytes())) hex::encode(mtp::crypto::sha256(input.as_bytes()))
} }

View file

@ -1,5 +1,5 @@
use base64::{Engine as _, engine::general_purpose::STANDARD}; use base64::{Engine as _, engine::general_purpose::STANDARD};
use mtp_crypto::{EncryptionType, Keyring, PublicKeyBundle, encrypt_for, decrypt_with}; use mtp::crypto::{EncryptionType, Keyring, PublicKeyBundle, decrypt_with, encrypt_for};
#[derive(Clone, Copy, Debug)] #[derive(Clone, Copy, Debug)]
pub enum DataFormat { pub enum DataFormat {
@ -22,13 +22,8 @@ pub fn encrypt(
.map_err(|e| format!("encryption error: {:?}", e)) .map_err(|e| format!("encryption error: {:?}", e))
} }
pub fn decrypt( pub fn decrypt(ciphertext: &[u8], aad: &[u8], keyring: &Keyring) -> Result<Vec<u8>, String> {
ciphertext: &[u8], decrypt_with(ciphertext, keyring, aad).map_err(|e| format!("decryption error: {:?}", e))
aad: &[u8],
keyring: &Keyring,
) -> Result<Vec<u8>, String> {
decrypt_with(ciphertext, keyring, aad)
.map_err(|e| format!("decryption error: {:?}", e))
} }
pub fn encrypt_challenge( pub fn encrypt_challenge(
@ -39,12 +34,10 @@ pub fn encrypt_challenge(
Ok(STANDARD.encode(&blob)) Ok(STANDARD.encode(&blob))
} }
pub fn decrypt_challenge( pub fn decrypt_challenge(encrypted: &str, keyring: &Keyring) -> Result<String, String> {
encrypted: &str, let blob = STANDARD
keyring: &Keyring, .decode(encrypted)
) -> Result<String, String> { .map_err(|e| format!("base64 decode error: {}", e))?;
let blob =
STANDARD.decode(encrypted).map_err(|e| format!("base64 decode error: {}", e))?;
let pt = decrypt(&blob, b"challenge", keyring)?; let pt = decrypt(&blob, b"challenge", keyring)?;
String::from_utf8(pt).map_err(|e| format!("utf8 decode error: {}", e)) String::from_utf8(pt).map_err(|e| format!("utf8 decode error: {}", e))
} }

View file

@ -4,6 +4,7 @@ version = "0.1.0"
edition = "2024" edition = "2024"
[dependencies] [dependencies]
iota-connection = { path = "../iota-connection" }
iota-logger = { path = "../iota-logger" } iota-logger = { path = "../iota-logger" }
iota-state = { path = "../iota-state" } iota-state = { path = "../iota-state" }
iota-storage = { path = "../iota-storage" } iota-storage = { path = "../iota-storage" }

File diff suppressed because it is too large Load diff

View file

@ -148,6 +148,7 @@ type_maps:
MessageReactionAdd: 146 MessageReactionAdd: 146
MessageReactionRemove: 147 MessageReactionRemove: 147
MessageReactionLive: 148 MessageReactionLive: 148
MessageDeleteLive: 150
DataTypes: DataTypes:
ErrorType: 32 ErrorType: 32
ErrorProtocol: 33 ErrorProtocol: 33
@ -169,7 +170,7 @@ type_maps:
CallState: 49 CallState: 49
ScreenShare: 50 ScreenShare: 50
PrivateKeyHash: 51 PrivateKeyHash: 51
Accepted: 52 # Accepted: 52 now part of default MTP
AcceptedProfiles: 53 AcceptedProfiles: 53
DeniedProfiles: 54 DeniedProfiles: 54
Content: 55 Content: 55

View file

@ -5,4 +5,10 @@ version = "0.1.0"
edition = "2024" edition = "2024"
[dependencies] [dependencies]
mtp = { git = "https://git.methanium.net/Methanium/mtp.git" } mtp = { git = "https://git.methanium.net/Methanium/mtp.git", features = ["web-server"] }
bytes = "1"
http = "1"
iota-state = { path = "../iota-state" }
iota-util = { path = "../iota-util" }
iota-logger = { path = "../iota-logger" }
tokio = { version = "1.50.0", features = ["full"] }

View file

@ -1,6 +1,133 @@
// The web server is a TTP host & identification system, use bytes::Bytes;
// it "upgrades" connections after identification to use iota_logger::log;
// use iota_state::{ACTIVE_TASKS, SHUTDOWN};
// either Own User (Cut down version of the Omikron Connection), use iota_util::file_util::load_file_vec;
// or Community (Custom Connection), use mtp::host::HostConfig;
// or Iota (Custom Connection). use mtp::webserver::{Http3Request, Http3Response, MTPWebServer, WebServerConfig};
use std::net::{IpAddr, Ipv4Addr};
use tokio::time::{Duration, sleep};
const CERT_PATH: &str = "certs/cert.pem";
const KEY_PATH: &str = "certs/cert.key";
async fn root(_request: Http3Request, response: Http3Response) -> Http3Response {
static_file("index.html", response).await
}
async fn static_file(path: &str, response: Http3Response) -> Http3Response {
let file = path.trim_start_matches('/');
let file = if file.is_empty() { "index.html" } else { file };
if file.split('/').any(|component| component == "..") {
return response
.status(http::StatusCode::BAD_REQUEST)
.body("invalid path");
}
let path = std::path::Path::new("web").join(file);
let Some(parent) = path.parent().and_then(|path| path.to_str()) else {
return response
.status(http::StatusCode::NOT_FOUND)
.body("not found");
};
let Some(name) = path.file_name().and_then(|name| name.to_str()) else {
return response
.status(http::StatusCode::NOT_FOUND)
.body("not found");
};
match load_file_vec(parent, name) {
Ok(body) => response
.status(http::StatusCode::OK)
.header("content-type", content_type(name))
.body(Bytes::from(body)),
Err(_) => response
.status(http::StatusCode::NOT_FOUND)
.body("not found"),
}
}
fn content_type(name: &str) -> &'static str {
match std::path::Path::new(name)
.extension()
.and_then(|ext| ext.to_str())
{
Some("html") => "text/html; charset=utf-8",
Some("css") => "text/css; charset=utf-8",
Some("js") => "application/javascript; charset=utf-8",
Some("json") => "application/json",
Some("png") => "image/png",
Some("ico") => "image/x-icon",
Some("woff2") => "font/woff2",
_ => "application/octet-stream",
}
}
pub async fn start(port: u16) -> bool {
let certificate = match tokio::fs::read(CERT_PATH).await {
Ok(certificate) => certificate,
Err(error) => {
log!("MTP web server certificate load failed: {}", error);
return false;
}
};
let key = match tokio::fs::read(KEY_PATH).await {
Ok(key) => key,
Err(error) => {
log!("MTP web server key load failed: {}", error);
return false;
}
};
let host_config = HostConfig::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), port, certificate, key);
let web_config = match WebServerConfig::new().route("/", root).and_then(|config| {
config.fallback(|request, response| async move {
static_file(request.uri.path(), response).await
})
}) {
Ok(config) => config,
Err(error) => {
log!("MTP web server route setup failed: {}", error);
return false;
}
};
let mut server = match MTPWebServer::new(host_config, web_config).await {
Ok(server) => server,
Err(error) => {
log!("MTP web server startup failed: {}", error);
return false;
}
};
log!("MTP web server running on port {}", port);
tokio::spawn(async move {
ACTIVE_TASKS.insert("WebServer".into());
loop {
tokio::select! {
result = server.accept() => {
match result {
Ok(Some(_connection)) => {}
Ok(None) => break,
Err(error) => log!("MTP webserver connection failed: {}", error),
}
}
_ = wait_for_shutdown() => {
server.shutdown().await;
break;
}
}
}
ACTIVE_TASKS.remove("WebServer");
});
true
}
async fn wait_for_shutdown() {
loop {
if *SHUTDOWN.read().await {
break;
}
sleep(Duration::from_millis(100)).await;
}
}