Fixes Patches Creation & More
Working with files / Paths fixed Websocket PingPongs (There is no data behind this just PingPong) Storing User Profiles correctly Added necessary dependencies and initial code support for: - X.509 certificate handling via x509/pkcs8 crates - Key generation and crypto operations via x448/sha2 - Structured hex encoding and base64 encoding
This commit is contained in:
parent
bba3f548ab
commit
0d07ddf851
23 changed files with 887 additions and 266 deletions
|
|
@ -1,10 +1,12 @@
|
|||
use crate::data::communication::{CommunicationType, CommunicationValue, DataTypes};
|
||||
use crate::omikron::ping_pong_task::PingPongTask;
|
||||
use crate::users::contact::Contact;
|
||||
use crate::users::user_community_util::UserCommunityUtil;
|
||||
use crate::util::chat_files::ChatFiles;
|
||||
use crate::util::chats_util::{get_user, get_users, mod_user};
|
||||
use futures_util::{SinkExt, StreamExt};
|
||||
use json::JsonValue;
|
||||
use json::number::Number;
|
||||
use std::collections::HashMap;
|
||||
use std::str::FromStr;
|
||||
use std::sync::Arc;
|
||||
|
|
@ -50,6 +52,7 @@ impl OmikronConnection {
|
|||
let (write_half, read_half) = ws_stream.split();
|
||||
*self.writer.lock().await = Some(write_half);
|
||||
self.spawn_listener(read_half);
|
||||
self.start_ping_pong_task().await;
|
||||
break;
|
||||
}
|
||||
Err(e) => {
|
||||
|
|
@ -59,8 +62,22 @@ impl OmikronConnection {
|
|||
}
|
||||
}
|
||||
|
||||
/// Gracefully close connection
|
||||
pub async fn close(&self) {
|
||||
// Start the PingPongTask
|
||||
pub async fn start_ping_pong_task(&self) {
|
||||
let ping_pong_task = PingPongTask::new(Arc::new(self.clone()));
|
||||
|
||||
// Store the task handle in `pingpong` so we can manage it
|
||||
let mut pingpong_handle = self.pingpong.lock().await;
|
||||
*pingpong_handle = Some(tokio::spawn(async move {
|
||||
ping_pong_task.run_ping_loop();
|
||||
}));
|
||||
}
|
||||
|
||||
pub async fn send_message(&self, msg: String) {
|
||||
Self::send_message_static(&self.writer, msg).await;
|
||||
}
|
||||
|
||||
pub async fn disconnect(&self) {
|
||||
if let Some(handle) = self.pingpong.lock().await.take() {
|
||||
handle.abort();
|
||||
}
|
||||
|
|
@ -88,14 +105,14 @@ impl OmikronConnection {
|
|||
}
|
||||
Ok(Message::Text(text)) => {
|
||||
let mut cv = CommunicationValue::from_json(&text); // needs CommunicationValue parser
|
||||
println!("[Omikron] Received message: {:?}", cv);
|
||||
if cv.is_type(CommunicationType::Pong) {
|
||||
println!("[Omikron] Pong received");
|
||||
// handle pingpong reset here
|
||||
continue;
|
||||
}
|
||||
// ************************************************ //
|
||||
// Direct messages //
|
||||
// ************************************************ //
|
||||
println!("[Omikron] Received message: {:?}", cv);
|
||||
if cv.is_type(CommunicationType::MessageOtherIota) {
|
||||
let sender_id = &cv.get_sender();
|
||||
let receiver_id = &cv.get_receiver();
|
||||
|
|
@ -105,8 +122,8 @@ impl OmikronConnection {
|
|||
.as_i64()
|
||||
.unwrap_or(-1),
|
||||
false,
|
||||
receiver_id.unwrap(),
|
||||
sender_id.unwrap(),
|
||||
*receiver_id,
|
||||
*sender_id,
|
||||
cv.get_data(DataTypes::MessageContent)
|
||||
.unwrap()
|
||||
.as_str()
|
||||
|
|
@ -114,7 +131,7 @@ impl OmikronConnection {
|
|||
);
|
||||
let response = CommunicationValue::new(CommunicationType::MessageLive)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(cv.get_receiver().unwrap())
|
||||
.with_receiver(cv.get_receiver())
|
||||
.add_data(
|
||||
DataTypes::SendTime,
|
||||
cv.get_data(DataTypes::SendTime).unwrap().clone(),
|
||||
|
|
@ -125,7 +142,7 @@ impl OmikronConnection {
|
|||
)
|
||||
.add_data(
|
||||
DataTypes::SenderId,
|
||||
JsonValue::String(cv.get_sender().unwrap().clone().to_string()),
|
||||
JsonValue::String(cv.get_sender().clone().to_string()),
|
||||
);
|
||||
Self::send_message_static(
|
||||
&writer.clone(),
|
||||
|
|
@ -143,7 +160,7 @@ impl OmikronConnection {
|
|||
.unwrap()
|
||||
.as_millis() as i64,
|
||||
true,
|
||||
my_id.unwrap(),
|
||||
my_id,
|
||||
Uuid::from_str(
|
||||
&*cv.get_data(DataTypes::ReceiverId).unwrap().to_string(),
|
||||
)
|
||||
|
|
@ -165,23 +182,23 @@ impl OmikronConnection {
|
|||
}
|
||||
|
||||
if cv.is_type(CommunicationType::MessageGet) {
|
||||
let my_id = cv.get_sender().unwrap();
|
||||
let my_id = cv.get_sender();
|
||||
let partner_id = Uuid::from_str(
|
||||
&*cv.get_data(DataTypes::ChatPartnerId).unwrap().to_string(),
|
||||
)
|
||||
.unwrap();
|
||||
let offset = cv
|
||||
.get_data(DataTypes::LoadedMessages)
|
||||
.unwrap()
|
||||
.unwrap_or(&JsonValue::Null)
|
||||
.to_string()
|
||||
.parse::<i64>()
|
||||
.unwrap();
|
||||
.unwrap_or(0);
|
||||
let amount = cv
|
||||
.get_data(DataTypes::MessageAmount)
|
||||
.unwrap()
|
||||
.unwrap_or(&JsonValue::Null)
|
||||
.to_string()
|
||||
.parse::<i64>()
|
||||
.unwrap();
|
||||
.unwrap_or(0);
|
||||
|
||||
let messages =
|
||||
ChatFiles::get_messages(my_id, partner_id, offset, amount); // needs ChatFiles
|
||||
|
|
@ -198,11 +215,12 @@ impl OmikronConnection {
|
|||
|
||||
if cv.is_type(CommunicationType::GetChats) {
|
||||
let user_id = cv.get_sender();
|
||||
let users = get_users(user_id.unwrap()); // needs ChatsUtil
|
||||
let users = get_users(user_id); // needs ChatsUtil
|
||||
let resp = CommunicationValue::new(CommunicationType::GetChats)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(user_id.unwrap())
|
||||
.with_receiver(user_id)
|
||||
.add_data(DataTypes::UserIds, users);
|
||||
println!("ALARM: {}", resp.to_json().to_string());
|
||||
Self::send_message_static(&writer.clone(), resp.to_json().to_string())
|
||||
.await;
|
||||
continue;
|
||||
|
|
@ -214,18 +232,18 @@ impl OmikronConnection {
|
|||
&*cv.get_data(DataTypes::UserId).unwrap().to_string(),
|
||||
)
|
||||
.unwrap();
|
||||
let mut contact = get_user(user_id.unwrap(), other_id)
|
||||
.unwrap_or(Contact::new(other_id)); // needs ChatsUtil + Contact
|
||||
let mut contact =
|
||||
get_user(user_id, other_id).unwrap_or(Contact::new(other_id)); // needs ChatsUtil + Contact
|
||||
contact.set_last_message_at(
|
||||
SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap()
|
||||
.as_millis() as i64,
|
||||
);
|
||||
mod_user(user_id.unwrap(), &contact);
|
||||
mod_user(user_id, &contact);
|
||||
let resp = CommunicationValue::new(CommunicationType::AddChat)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(user_id.unwrap());
|
||||
.with_receiver(user_id);
|
||||
Self::send_message_static(&writer.clone(), resp.to_json().to_string())
|
||||
.await;
|
||||
continue;
|
||||
|
|
@ -233,7 +251,7 @@ impl OmikronConnection {
|
|||
|
||||
if cv.is_type(CommunicationType::AddCommunity) {
|
||||
UserCommunityUtil::add_community(
|
||||
cv.get_sender().unwrap(),
|
||||
cv.get_sender(),
|
||||
cv.get_data(DataTypes::CommunityAddress)
|
||||
.unwrap()
|
||||
.to_string(),
|
||||
|
|
@ -242,7 +260,7 @@ impl OmikronConnection {
|
|||
);
|
||||
let resp = CommunicationValue::new(CommunicationType::AddCommunity)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(cv.get_sender().unwrap());
|
||||
.with_receiver(cv.get_sender());
|
||||
Self::send_message_static(&writer.clone(), resp.to_json().to_string())
|
||||
.await;
|
||||
continue;
|
||||
|
|
@ -251,10 +269,10 @@ impl OmikronConnection {
|
|||
if cv.is_type(CommunicationType::GetCommunities) {
|
||||
let resp = CommunicationValue::new(CommunicationType::GetCommunities)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(cv.get_sender().unwrap())
|
||||
.with_receiver(cv.get_sender())
|
||||
.add_data(
|
||||
DataTypes::Communities,
|
||||
UserCommunityUtil::get_communities(cv.get_sender().unwrap()),
|
||||
UserCommunityUtil::get_communities(cv.get_sender()),
|
||||
); // needs UserCommunityUtil
|
||||
Self::send_message_static(&writer.clone(), resp.to_json().to_string())
|
||||
.await;
|
||||
|
|
@ -263,14 +281,14 @@ impl OmikronConnection {
|
|||
|
||||
if cv.is_type(CommunicationType::RemoveCommunity) {
|
||||
UserCommunityUtil::remove_community(
|
||||
cv.get_sender().unwrap(),
|
||||
cv.get_sender(),
|
||||
cv.get_data(DataTypes::CommunityAddress)
|
||||
.unwrap()
|
||||
.to_string(),
|
||||
); // needs UserCommunityUtil
|
||||
let resp = CommunicationValue::new(CommunicationType::RemoveCommunity)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(cv.get_sender().unwrap());
|
||||
.with_receiver(cv.get_sender());
|
||||
Self::send_message_static(&writer.clone(), resp.to_json().to_string())
|
||||
.await;
|
||||
continue;
|
||||
|
|
@ -285,9 +303,6 @@ impl OmikronConnection {
|
|||
}
|
||||
});
|
||||
}
|
||||
pub async fn send_message(&self, msg: String) {
|
||||
Self::send_message_static(&self.writer, msg).await;
|
||||
}
|
||||
pub async fn send_message_static(
|
||||
writer: &Arc<
|
||||
Mutex<
|
||||
|
|
@ -322,10 +337,18 @@ impl OmikronConnection {
|
|||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() {
|
||||
let conn = OmikronConnection::new();
|
||||
conn.connect().await;
|
||||
pub async fn send_ping_message(&self, uuid: Uuid) {
|
||||
// Send the ping message over the connection
|
||||
let ping_message = CommunicationValue::new(CommunicationType::Ping)
|
||||
.with_id(uuid)
|
||||
.add_data_num(DataTypes::LastPing, Number::from(2))
|
||||
.to_json()
|
||||
.to_string();
|
||||
self.send_message(ping_message).await;
|
||||
}
|
||||
pub async fn reconnect(&self) {
|
||||
self.disconnect().await;
|
||||
self.connect().await;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue