Compare commits

..
3 changed files with 35 additions and 16 deletions

22
Cargo.lock generated
View file

@ -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]]

@ -1 +1 @@
Subproject commit ece6e2c3b4e925f3cefe46f4a048fbfc8f823093
Subproject commit 594646ac39d986f0787aa614a99d580035a67318

View file

@ -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<Option<WebMtpSender>>,
pub ping: RwLock<i64>,
waiting_tasks: DashMap<u32, WaitingTask>,
cleanup_handle: std::sync::Mutex<Option<tokio::task::JoinHandle<()>>>,
}
@ -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<Self>, value: CommunicationValue) -> OmikronResult<()> {
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<Self>, 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<Self>, value: &CommunicationValue) -> OmikronResult<()> {
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()