diff --git a/Cargo.lock b/Cargo.lock index a5b9c46..ee52b1c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1614,7 +1614,7 @@ dependencies = [ [[package]] name = "mtp" version = "0.2.0" -source = "git+https://git.methanium.net/methanium/mtp#fa271e62bec7918a8c994eea59ecdd62be76ead4" +source = "git+https://git.methanium.net/methanium/mtp#bcf8aee3716f1690f1284748ac1ff3cd22799fb0" dependencies = [ "mtp-client", "mtp-codec", @@ -1630,7 +1630,7 @@ dependencies = [ [[package]] name = "mtp-client" version = "0.2.0" -source = "git+https://git.methanium.net/methanium/mtp#fa271e62bec7918a8c994eea59ecdd62be76ead4" +source = "git+https://git.methanium.net/methanium/mtp#bcf8aee3716f1690f1284748ac1ff3cd22799fb0" dependencies = [ "mtp-codec", "mtp-common", @@ -1643,7 +1643,7 @@ dependencies = [ [[package]] name = "mtp-codec" version = "0.2.0" -source = "git+https://git.methanium.net/methanium/mtp#fa271e62bec7918a8c994eea59ecdd62be76ead4" +source = "git+https://git.methanium.net/methanium/mtp#bcf8aee3716f1690f1284748ac1ff3cd22799fb0" dependencies = [ "base64", "byteorder", @@ -1656,7 +1656,7 @@ dependencies = [ [[package]] name = "mtp-common" version = "0.2.0" -source = "git+https://git.methanium.net/methanium/mtp#fa271e62bec7918a8c994eea59ecdd62be76ead4" +source = "git+https://git.methanium.net/methanium/mtp#bcf8aee3716f1690f1284748ac1ff3cd22799fb0" dependencies = [ "quinn", "rustls", @@ -1667,7 +1667,7 @@ dependencies = [ [[package]] name = "mtp-crypto" version = "0.2.0" -source = "git+https://git.methanium.net/methanium/mtp#fa271e62bec7918a8c994eea59ecdd62be76ead4" +source = "git+https://git.methanium.net/methanium/mtp#bcf8aee3716f1690f1284748ac1ff3cd22799fb0" dependencies = [ "base64", "chacha20poly1305", @@ -1689,7 +1689,7 @@ dependencies = [ [[package]] name = "mtp-files" version = "0.2.0" -source = "git+https://git.methanium.net/methanium/mtp#fa271e62bec7918a8c994eea59ecdd62be76ead4" +source = "git+https://git.methanium.net/methanium/mtp#bcf8aee3716f1690f1284748ac1ff3cd22799fb0" dependencies = [ "mtp-crypto", "rand 0.10.2", @@ -1700,7 +1700,7 @@ dependencies = [ [[package]] name = "mtp-host" version = "0.2.0" -source = "git+https://git.methanium.net/methanium/mtp#fa271e62bec7918a8c994eea59ecdd62be76ead4" +source = "git+https://git.methanium.net/methanium/mtp#bcf8aee3716f1690f1284748ac1ff3cd22799fb0" dependencies = [ "mtp-codec", "mtp-common", @@ -1715,7 +1715,7 @@ dependencies = [ [[package]] name = "mtp-transport" version = "0.2.0" -source = "git+https://git.methanium.net/methanium/mtp#fa271e62bec7918a8c994eea59ecdd62be76ead4" +source = "git+https://git.methanium.net/methanium/mtp#bcf8aee3716f1690f1284748ac1ff3cd22799fb0" dependencies = [ "async-trait", "mtp-codec", @@ -1733,7 +1733,7 @@ dependencies = [ [[package]] name = "mtp-type-map" version = "0.2.0" -source = "git+https://git.methanium.net/methanium/mtp#fa271e62bec7918a8c994eea59ecdd62be76ead4" +source = "git+https://git.methanium.net/methanium/mtp#bcf8aee3716f1690f1284748ac1ff3cd22799fb0" dependencies = [ "serde", "serde_yaml", @@ -1742,7 +1742,7 @@ dependencies = [ [[package]] name = "mtp-webserver" version = "0.2.0" -source = "git+https://git.methanium.net/methanium/mtp#fa271e62bec7918a8c994eea59ecdd62be76ead4" +source = "git+https://git.methanium.net/methanium/mtp#bcf8aee3716f1690f1284748ac1ff3cd22799fb0" dependencies = [ "async-trait", "bytes", @@ -3515,7 +3515,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.48.0", ] [[package]] diff --git a/mtp-type-maps b/mtp-type-maps index ece6e2c..594646a 160000 --- a/mtp-type-maps +++ b/mtp-type-maps @@ -1 +1 @@ -Subproject commit ece6e2c3b4e925f3cefe46f4a048fbfc8f823093 +Subproject commit 594646ac39d986f0787aa614a99d580035a67318 diff --git a/src/transport/omikron_connection.rs b/src/transport/omikron_connection.rs index e86f4a2..54a3bf9 100644 --- a/src/transport/omikron_connection.rs +++ b/src/transport/omikron_connection.rs @@ -6,7 +6,7 @@ use crate::{ }; use dashmap::DashMap; use mtp::{ - codec::{CommunicationType, CommunicationValue}, + codec::{CommunicationType, CommunicationValue, DataType, DataValue}, crypto::PublicKeyBundle, host::{AuthenticationPolicy, HostConfig, Policy, SendMode}, webserver::{MTPWebServer, WebMtpReceiver, WebMtpSender}, @@ -19,7 +19,10 @@ use std::{ }, time::{Duration, Instant}, }; -use tokio::{sync::Mutex, time::interval}; +use tokio::{ + sync::{Mutex, RwLock}, + time::interval, +}; const CLEANUP_INTERVAL: Duration = Duration::from_secs(30); const MAX_WAITING_AGE: Duration = Duration::from_secs(60); @@ -49,6 +52,7 @@ pub struct WaitingTask { pub struct OmikronConnection { id: u64, sender: Mutex>, + pub ping: RwLock, waiting_tasks: DashMap, cleanup_handle: std::sync::Mutex>>, } @@ -65,6 +69,7 @@ impl OmikronConnection { Arc::new(Self { id, sender: Mutex::new(Some(sender)), + ping: RwLock::new(-1), waiting_tasks: DashMap::new(), cleanup_handle: std::sync::Mutex::new(None), }) @@ -122,11 +127,16 @@ impl OmikronConnection { } async fn process_message(self: Arc, value: CommunicationValue) -> OmikronResult<()> { - log_cv_in!(PrintType::Omikron, &value); + if !value.is_type(CommunicationType::Pong) && !value.is_type(CommunicationType::Ping) { + log_cv_in!(PrintType::Omikron, &value); + } if let Some((_, task)) = self.waiting_tasks.remove(&value.get_id()) { let _ = (task.task)(self.clone(), value); return Ok(()); } + if value.is_type(CommunicationType::Ping) { + return self.ping(value).await; + } self.dispatch(value).await } @@ -205,8 +215,17 @@ impl OmikronConnection { } } + async fn ping(self: Arc, value: CommunicationValue) -> OmikronResult<()> { + if let DataValue::SignedNumber(last_ping) = value.get_data(DataType::LastPing) { + *self.ping.write().await = *last_ping as i64; + } + let response = CommunicationValue::new(CommunicationType::Pong).with_id(value.get_id()); + self.send(&response).await + } pub(crate) async fn send(self: Arc, value: &CommunicationValue) -> OmikronResult<()> { - log_cv_out!(PrintType::Omikron, value); + if !value.is_type(CommunicationType::Pong) && !value.is_type(CommunicationType::Ping) { + log_cv_out!(PrintType::Omikron, value); + } let guard = self.sender.lock().await; let sender = guard .as_ref()