Omikron Connects Data Sends Identified Users loaded
This commit is contained in:
parent
6cd859953e
commit
a1d8bcbeff
10 changed files with 228 additions and 209 deletions
|
|
@ -1,9 +1,9 @@
|
|||
use futures_util::{SinkExt, StreamExt};
|
||||
use std::collections::HashMap;
|
||||
use std::ptr::write;
|
||||
use std::str::FromStr;
|
||||
use std::sync::Arc;
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
use json::JsonValue;
|
||||
use tokio::sync::Mutex;
|
||||
use tokio::net::TcpStream;
|
||||
use tokio::time::{sleep, Duration};
|
||||
|
|
@ -37,7 +37,7 @@ impl OmikronConnection {
|
|||
}
|
||||
|
||||
/// Connect loop with retry
|
||||
pub async fn connect(&self) {
|
||||
pub async fn connect<'a>(&'a self) {
|
||||
loop {
|
||||
match connect_async("wss://tensamin.methanium.net/ws/iota/").await {
|
||||
Ok((ws_stream, _)) => {
|
||||
|
|
@ -85,14 +85,6 @@ impl OmikronConnection {
|
|||
println!("[Omikron] Pong received");
|
||||
// handle pingpong reset here
|
||||
}
|
||||
|
||||
if waiting.lock().await.contains_key(&cv.get_id()) {
|
||||
if let Some(callback) = waiting.lock().await.remove(&cv.get_id()) {
|
||||
callback(cv);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
// ************************************************ //
|
||||
// Direct messages //
|
||||
// ************************************************ //
|
||||
|
|
@ -100,19 +92,20 @@ impl OmikronConnection {
|
|||
let sender_id = &cv.get_sender();
|
||||
let receiver_id = &cv.get_receiver();
|
||||
ChatFiles::add_message(
|
||||
cv.get_data(DataTypes::SendTime).unwrap().parse::<i64>().unwrap_or(-1),
|
||||
|
||||
cv.get_data(DataTypes::SendTime).unwrap().as_i64().unwrap_or(-1),
|
||||
false,
|
||||
receiver_id.unwrap(),
|
||||
sender_id.unwrap(),
|
||||
cv.get_data(DataTypes::MessageContent).unwrap().as_str(),
|
||||
cv.get_data(DataTypes::MessageContent).unwrap().as_str().unwrap(),
|
||||
);
|
||||
let response = CommunicationValue::new(CommunicationType::MessageLive)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(cv.get_receiver().unwrap())
|
||||
.add_data(DataTypes::SendTime, cv.get_data(DataTypes::SendTime).unwrap().to_string())
|
||||
.add_data(DataTypes::Message, cv.get_data(DataTypes::MessageContent).unwrap().to_string())
|
||||
.add_data(DataTypes::SenderId, cv.get_sender().unwrap().to_string());
|
||||
Self::send_message_static(&writer, response.to_json().to_string()).await;
|
||||
.add_data(DataTypes::SendTime, cv.get_data(DataTypes::SendTime).unwrap().clone())
|
||||
.add_data(DataTypes::Message, cv.get_data(DataTypes::MessageContent).unwrap().clone())
|
||||
.add_data(DataTypes::SenderId, JsonValue::String(cv.get_sender().unwrap().clone().to_string()));
|
||||
Self::send_message_static(&writer.clone(), response.to_json().to_string()).await;
|
||||
continue;
|
||||
}
|
||||
|
||||
|
|
@ -129,10 +122,10 @@ impl OmikronConnection {
|
|||
);
|
||||
// ack
|
||||
let ack = CommunicationValue::ack_message(cv.get_id(), my_id);
|
||||
Self::send_message_static(&writer, ack.to_json().to_string()).await;
|
||||
Self::send_message_static(&writer.clone(), ack.to_json().to_string()).await;
|
||||
// forward
|
||||
let forward = CommunicationValue::forward_to_other_iota(&mut cv);
|
||||
Self::send_message_static(&writer, forward.to_json().to_string()).await;
|
||||
Self::send_message_static(&writer.clone(), forward.to_json().to_string()).await;
|
||||
continue;
|
||||
}
|
||||
|
||||
|
|
@ -147,9 +140,9 @@ impl OmikronConnection {
|
|||
.with_id(cv.get_id())
|
||||
.with_receiver(my_id);
|
||||
if !messages.is_empty() {
|
||||
resp = resp.add_data(DataTypes::MessageChunk, messages.to_string());
|
||||
resp = resp.add_data(DataTypes::MessageChunk, messages);
|
||||
}
|
||||
Self::send_message_static(&writer, resp.to_json().to_string()).await;
|
||||
Self::send_message_static(&writer.clone(), resp.to_json().to_string()).await;
|
||||
continue;
|
||||
}
|
||||
|
||||
|
|
@ -159,14 +152,14 @@ impl OmikronConnection {
|
|||
let resp = CommunicationValue::new(CommunicationType::GetChats)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(user_id.unwrap())
|
||||
.add_data(DataTypes::UserIds, users.to_string());
|
||||
Self::send_message_static(&writer, resp.to_json().to_string()).await;
|
||||
.add_data(DataTypes::UserIds, users);
|
||||
Self::send_message_static(&writer.clone(), resp.to_json().to_string()).await;
|
||||
continue;
|
||||
}
|
||||
|
||||
if cv.is_type(CommunicationType::AddChat) {
|
||||
let user_id = cv.get_sender();
|
||||
let other_id = Uuid::from_str(&*cv.get_data(DataTypes::ReceiverId).unwrap().to_string()).unwrap();
|
||||
let other_id = Uuid::from_str(&*cv.get_data(DataTypes::UserId).unwrap().to_string()).unwrap();
|
||||
let mut contact = ChatsUtil::get_user(user_id.unwrap(), other_id).unwrap_or(Contact::new(other_id)); // needs ChatsUtil + Contact
|
||||
contact.set_last_message_at(SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
|
|
@ -176,7 +169,7 @@ impl OmikronConnection {
|
|||
let resp = CommunicationValue::new(CommunicationType::AddChat)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(user_id.unwrap());
|
||||
Self::send_message_static(&writer, resp.to_json().to_string()).await;
|
||||
Self::send_message_static(&writer.clone(), resp.to_json().to_string()).await;
|
||||
continue;
|
||||
}
|
||||
|
||||
|
|
@ -189,7 +182,7 @@ impl OmikronConnection {
|
|||
let resp = CommunicationValue::new(CommunicationType::AddCommunity)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(cv.get_sender().unwrap());
|
||||
Self::send_message_static(&writer, resp.to_json().to_string()).await;
|
||||
Self::send_message_static(&writer.clone(), resp.to_json().to_string()).await;
|
||||
continue;
|
||||
}
|
||||
|
||||
|
|
@ -197,8 +190,8 @@ impl OmikronConnection {
|
|||
let resp = CommunicationValue::new(CommunicationType::GetCommunities)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(cv.get_sender().unwrap())
|
||||
.add_data(DataTypes::Communities, UserCommunityUtil::get_communities(cv.get_sender().unwrap()).to_string()); // needs UserCommunityUtil
|
||||
Self::send_message_static(&writer, resp.to_json().to_string()).await;
|
||||
.add_data(DataTypes::Communities, UserCommunityUtil::get_communities(cv.get_sender().unwrap())); // needs UserCommunityUtil
|
||||
Self::send_message_static(&writer.clone(), resp.to_json().to_string()).await;
|
||||
continue;
|
||||
}
|
||||
|
||||
|
|
@ -207,7 +200,7 @@ impl OmikronConnection {
|
|||
let resp = CommunicationValue::new(CommunicationType::RemoveCommunity)
|
||||
.with_id(cv.get_id())
|
||||
.with_receiver(cv.get_sender().unwrap());
|
||||
Self::send_message_static(&writer, resp.to_json().to_string()).await;
|
||||
Self::send_message_static(&writer.clone(), resp.to_json().to_string()).await;
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
|
@ -220,16 +213,20 @@ impl OmikronConnection {
|
|||
}
|
||||
});
|
||||
}
|
||||
pub fn send_message(&self, msg: String){
|
||||
OmikronConnection::send_message_static(&self.writer, msg);
|
||||
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<Option<futures_util::stream::SplitSink<WebSocketStream<MaybeTlsStream<TcpStream>>, Message>>>>,
|
||||
msg: String,
|
||||
) {
|
||||
) -> Result<(), tokio_tungstenite::tungstenite::Error> {
|
||||
let mut guard = writer.lock().await;
|
||||
if let Some(writer) = guard.as_mut() {
|
||||
let _ = writer.send(Message::Text(msg)).await;
|
||||
writer.send(Message::Text(msg)).await?;
|
||||
writer.flush().await?;
|
||||
Ok(())
|
||||
} else {
|
||||
Err(tokio_tungstenite::tungstenite::Error::ConnectionClosed)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Reference in a new issue