Fixed warnings and migrated to TTP.
This commit is contained in:
parent
b4e8e42267
commit
7a331d4450
20 changed files with 107 additions and 55 deletions
|
|
@ -6,12 +6,12 @@ use crate::rho::{rho_connection::RhoConnection, rho_manager};
|
|||
use crate::util::logger::PrintType;
|
||||
use crate::{data::user::UserStatus, omega::omega_connection::OmegaConnection};
|
||||
use crate::{log_cv_in, log_cv_out, log_out};
|
||||
use epsilon_core::{CommunicationType, CommunicationValue, DataTypes, DataValue};
|
||||
use epsilon_native::{Receiver, Sender};
|
||||
use std::str::FromStr;
|
||||
use std::sync::Arc;
|
||||
use std::time::{Duration, SystemTime, UNIX_EPOCH};
|
||||
use tokio::sync::RwLock;
|
||||
use ttp_core::{CommunicationType, CommunicationValue, DataTypes, DataValue};
|
||||
use ttp_native::{Receiver, Sender};
|
||||
use uuid::Uuid;
|
||||
|
||||
pub struct ClientConnection {
|
||||
|
|
@ -455,6 +455,7 @@ impl ClientConnection {
|
|||
}
|
||||
|
||||
/// Close the connection
|
||||
#[allow(dead_code)]
|
||||
pub async fn close(&self) {
|
||||
let mut is_open_guard = self.is_open.write().await;
|
||||
if *is_open_guard {
|
||||
|
|
@ -470,12 +471,14 @@ impl ClientConnection {
|
|||
let mut interested_guard = self.interested_users.write().await;
|
||||
*interested_guard = interested_ids;
|
||||
}
|
||||
#[allow(dead_code)]
|
||||
pub async fn get_interested_users(self: Arc<Self>) -> Vec<i64> {
|
||||
let interested_guard = self.interested_users.read().await;
|
||||
interested_guard.clone()
|
||||
}
|
||||
|
||||
/// Check if interested in a user and send notification
|
||||
#[allow(dead_code)]
|
||||
pub async fn are_you_interested(self: Arc<Self>, user_id: i64) {
|
||||
let interested_guard = self.clone().get_interested_users().await;
|
||||
if interested_guard.contains(&user_id) {
|
||||
|
|
@ -488,6 +491,7 @@ impl ClientConnection {
|
|||
}
|
||||
|
||||
/// Handle connection close
|
||||
#[allow(dead_code)]
|
||||
pub async fn handle_close(&self) {
|
||||
let user_id = self.get_user_id().await;
|
||||
if let Some(rho_conn) = rho_manager::get_rho_con_for_user(user_id as i64).await {
|
||||
|
|
|
|||
|
|
@ -1,8 +1,8 @@
|
|||
use epsilon_core::{CommunicationType, CommunicationValue, DataTypes, DataValue};
|
||||
use epsilon_native::{Receiver, Sender};
|
||||
use rand::{Rng, distributions::Alphanumeric};
|
||||
use std::{sync::Arc, time::Duration};
|
||||
use tokio::sync::RwLock;
|
||||
use ttp_core::{CommunicationType, CommunicationValue, DataTypes, DataValue};
|
||||
use ttp_native::{Receiver, Sender};
|
||||
|
||||
use crate::{
|
||||
anonymous_clients::anonymous_client_connection::AnonymousClientConnection,
|
||||
|
|
@ -20,6 +20,7 @@ use crate::{
|
|||
};
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
#[allow(dead_code)]
|
||||
pub enum ConnectionKind {
|
||||
Client,
|
||||
Iota,
|
||||
|
|
|
|||
|
|
@ -7,12 +7,6 @@ use crate::omega::omega_connection::get_omega_connection;
|
|||
use crate::rho::connection::GeneralConnection;
|
||||
use crate::util::logger::PrintType;
|
||||
use dashmap::DashMap;
|
||||
use epsilon_core::CommunicationType;
|
||||
use epsilon_core::CommunicationValue;
|
||||
use epsilon_core::DataTypes;
|
||||
use epsilon_core::DataValue;
|
||||
use epsilon_native::Receiver;
|
||||
use epsilon_native::Sender;
|
||||
use std::collections::BTreeMap;
|
||||
use std::{
|
||||
collections::HashMap,
|
||||
|
|
@ -21,11 +15,18 @@ use std::{
|
|||
};
|
||||
use tokio::sync::RwLock;
|
||||
use tokio::sync::mpsc;
|
||||
use ttp_core::CommunicationType;
|
||||
use ttp_core::CommunicationValue;
|
||||
use ttp_core::DataTypes;
|
||||
use ttp_core::DataValue;
|
||||
use ttp_native::Receiver;
|
||||
use ttp_native::Sender;
|
||||
use x448::PublicKey;
|
||||
|
||||
use super::{rho_connection::RhoConnection, rho_manager};
|
||||
use crate::omega::omega_connection::OmegaConnection;
|
||||
|
||||
#[allow(dead_code)]
|
||||
pub struct IotaConnection {
|
||||
pub iota_id: u64,
|
||||
pub sender: Arc<Sender>,
|
||||
|
|
@ -72,6 +73,7 @@ impl IotaConnection {
|
|||
self.iota_id
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
pub async fn get_public_key(&self) -> Option<PublicKey> {
|
||||
if let Some(public_key) = self.pub_key.read().await.clone() {
|
||||
PublicKey::from_bytes(&public_key)
|
||||
|
|
@ -163,11 +165,13 @@ impl IotaConnection {
|
|||
self.forward_to_client(cv).await;
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
async fn send_error_response(&self, message_id: u32, error_type: CommunicationType) {
|
||||
let error = CommunicationValue::new(error_type).with_id(message_id);
|
||||
self.send_message(&error).await;
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
async fn close(&self) {
|
||||
let _ = self.sender.close();
|
||||
}
|
||||
|
|
@ -346,12 +350,14 @@ impl IotaConnection {
|
|||
}
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
pub async fn handle_close(&self) {
|
||||
if let Some(rho_conn) = self.get_rho_connection().await {
|
||||
rho_conn.close_iota_connection().await;
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
pub async fn await_response(
|
||||
self: Arc<IotaConnection>,
|
||||
cv: &CommunicationValue,
|
||||
|
|
|
|||
|
|
@ -2,10 +2,10 @@ use super::{client_connection::ClientConnection, iota_connection::IotaConnection
|
|||
|
||||
use crate::data::user::UserStatus;
|
||||
use crate::omega::omega_connection::OmegaConnection;
|
||||
use epsilon_core::{CommunicationType, CommunicationValue, DataTypes, DataValue};
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
use tokio::sync::RwLock;
|
||||
use ttp_core::{CommunicationType, CommunicationValue, DataTypes, DataValue};
|
||||
|
||||
pub struct RhoConnection {
|
||||
iota_connection: Arc<IotaConnection>,
|
||||
|
|
@ -58,6 +58,7 @@ impl RhoConnection {
|
|||
}
|
||||
|
||||
/// Add a client connection
|
||||
#[allow(dead_code)]
|
||||
pub async fn add_client_connection(&self, connection: Arc<ClientConnection>) {
|
||||
let notification = CommunicationValue::new(CommunicationType::client_connected).add_data(
|
||||
DataTypes::user_id,
|
||||
|
|
@ -145,6 +146,7 @@ impl RhoConnection {
|
|||
}
|
||||
|
||||
/// Check if clients are interested in a user
|
||||
#[allow(dead_code)]
|
||||
pub async fn are_they_interested(&self, user_id: i64) {
|
||||
let connections = self.client_connections.read().await;
|
||||
for connection in connections.iter() {
|
||||
|
|
@ -166,11 +168,13 @@ impl RhoConnection {
|
|||
}
|
||||
|
||||
/// Check if this RhoConnection contains a specific user ID
|
||||
#[allow(dead_code)]
|
||||
pub fn contains_user(&self, user_id: &i64) -> bool {
|
||||
self.user_ids.contains(user_id)
|
||||
}
|
||||
|
||||
/// Get count of active client connections
|
||||
#[allow(dead_code)]
|
||||
pub async fn client_count(&self) -> usize {
|
||||
let connections = self.client_connections.read().await;
|
||||
connections.len()
|
||||
|
|
|
|||
|
|
@ -26,6 +26,7 @@ pub async fn get_rho_con_for_user(user_id: i64) -> Option<Arc<RhoConnection>> {
|
|||
None
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
pub async fn contains_iota(iota_id: i64) -> bool {
|
||||
let connections = RHO_CONNECTIONS.read().await;
|
||||
connections.contains_key(&iota_id)
|
||||
|
|
@ -45,6 +46,7 @@ pub async fn add_rho(rho_connection: Arc<RhoConnection>) {
|
|||
}
|
||||
|
||||
/// Get a RhoConnection by Iota ID directly
|
||||
#[allow(dead_code)]
|
||||
pub async fn get_rho_by_iota(iota_id: i64) -> Option<Arc<RhoConnection>> {
|
||||
let connections = RHO_CONNECTIONS.read().await;
|
||||
connections.get(&iota_id).map(Arc::clone)
|
||||
|
|
|
|||
|
|
@ -3,14 +3,14 @@ use crate::{
|
|||
rho::connection::GeneralConnection,
|
||||
util::{file_util::load_file_vec, logger::PrintType},
|
||||
};
|
||||
use epsilon_native::Host;
|
||||
use ttp_native::Host;
|
||||
|
||||
pub async fn start(port: u16) -> Result<(), Box<dyn std::error::Error>> {
|
||||
let cert_pem = load_file_vec("certs", "cert.pem").expect("Error loading Pemfile");
|
||||
|
||||
let key_pem = load_file_vec("certs", "key.pem").expect("Error loading Keyfile");
|
||||
|
||||
let mut host: Host = epsilon_native::host(port, cert_pem, key_pem).await?;
|
||||
let mut host: Host = ttp_native::host(port, cert_pem, key_pem).await?;
|
||||
log!(0, PrintType::General, "Server listening on port {}", port);
|
||||
|
||||
while let Some((sender, receiver)) = host.next().await {
|
||||
|
|
|
|||
Loading…
Reference in a new issue