Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
269 changes: 260 additions & 9 deletions Cargo.lock

Large diffs are not rendered by default.

21 changes: 19 additions & 2 deletions rkvm-client/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,22 @@ version = "0.6.1"
authors = ["Jan Trefil <8711792+htrefil@users.noreply.github.com>"]
edition = "2021"

# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
[features]
windows-service = [
"dep:windows-service",
"windows/Win32_System_Diagnostics_ToolHelp",
"windows/Win32_System_Threading",
"windows/Win32_Security_Authorization",
"windows/Win32_System_RemoteDesktop",
]

[[bin]]
name = "rkvm-service"
path = "src/windows-service.rs"
required-features = ["windows-service"]

[dependencies]
async-trait = "0.1.89"
tokio = { version = "1.0.1", features = ["macros", "time", "fs", "net", "signal", "rt-multi-thread", "sync"] }
rkvm-input = { path = "../rkvm-input" }
rkvm-net = { path = "../rkvm-net" }
Expand All @@ -19,7 +32,11 @@ thiserror = "1.0.40"
tokio-rustls = "0.24.0"
rustls-pemfile = "1.0.2"
tracing = "0.1.37"
tracing-subscriber = { version = "0.3.17", features = ["env-filter"] }
tracing-subscriber = { version = "0.3.17", features = ["env-filter", "local-time"] }

[target.'cfg(windows)'.dependencies]
windows = "0.62"
windows-service = { version = "0.8", optional = true }

[package.metadata.rpm]
package = "rkvm-client"
Expand Down
262 changes: 94 additions & 168 deletions rkvm-client/src/client.rs
Original file line number Diff line number Diff line change
@@ -1,208 +1,134 @@
use rkvm_input::writer::Writer;
use rkvm_net::auth::{AuthChallenge, AuthStatus};
use rkvm_input::writer::{DeviceWriter};
use rkvm_net::message::Message;
use rkvm_net::version::Version;
use rkvm_net::{Pong, Update};
use std::collections::hash_map::Entry;
use std::collections::HashMap;
use std::io;
use rkvm_net::Update;

use async_trait::async_trait;
use std::fs::{rename, OpenOptions};
use std::io::{self, stdout, BufWriter};
use std::path::Path;
use std::time::Instant;
use thiserror::Error;
use tokio::io::{AsyncWriteExt, BufStream};
use tokio::net::TcpStream;
use tokio::io::{AsyncRead, AsyncWrite, AsyncWriteExt};
use tokio::time;
use tokio_rustls::rustls::ServerName;
use tokio_rustls::TlsConnector;
use tokio_rustls::rustls;
use tracing_subscriber::{fmt, Registry,EnvFilter};
use tracing_subscriber::fmt::time::LocalTime;
use tracing_subscriber::prelude::*;

#[derive(Error, Debug)]
pub enum Error {
#[error("Network error: {0}")]
Network(io::Error),
#[error("Input error: {0}")]
Input(io::Error),
#[error("Io error: {0}")]
Io(#[from] io::Error),
#[error(transparent)]
Rustls(#[from] rustls::Error),
#[error("Toml error: {0}")]
Toml(#[from] toml::de::Error),
#[cfg(target_os="windows")]
#[error("Windows API error: {0}")]
Windows(#[from] windows::core::Error),
#[cfg(all(target_os="windows",feature="windows-service"))]
#[error("Windows Service error: {0}")]
WindowsService(#[from] windows_service::Error),
#[allow(dead_code)]
#[error("Incompatible server version (got {server}, expected {client})")]
Version { server: Version, client: Version },
#[allow(dead_code)]
#[error("Invalid password")]
Auth,
}

pub async fn run(
hostname: &ServerName,
port: u16,
connector: TlsConnector,
password: &str,
) -> Result<(), Error> {
// Intentionally don't impose any timeout for TCP connect.
let stream = match hostname {
ServerName::DnsName(name) => TcpStream::connect(&(name.as_ref(), port)).await,
ServerName::IpAddress(address) => TcpStream::connect(&(*address, port)).await,
_ => unimplemented!("Unhandled rustls ServerName variant: {:?}", hostname),
}
.map_err(Error::Network)?;

tracing::info!("Connected to server");

let stream = rkvm_net::timeout(
rkvm_net::TLS_TIMEOUT,
connector.connect(hostname.clone(), stream),
)
.await
.map_err(Error::Network)?;

tracing::info!("TLS connected");

let mut stream = BufStream::with_capacity(1024, 1024, stream);

rkvm_net::timeout(rkvm_net::WRITE_TIMEOUT, async {
Version::CURRENT.encode(&mut stream).await?;
stream.flush().await?;

Ok(())
})
.await
.map_err(Error::Network)?;

let version = rkvm_net::timeout(rkvm_net::READ_TIMEOUT, Version::decode(&mut stream))
.await
.map_err(Error::Network)?;

if version != Version::CURRENT {
return Err(Error::Version {
server: Version::CURRENT,
client: version,
});
}

let challenge = rkvm_net::timeout(rkvm_net::READ_TIMEOUT, AuthChallenge::decode(&mut stream))
.await
.map_err(Error::Network)?;

let response = challenge.respond(password);

rkvm_net::timeout(rkvm_net::WRITE_TIMEOUT, async {
response.encode(&mut stream).await?;
stream.flush().await?;

Ok(())
})
.await
.map_err(Error::Network)?;

let status = rkvm_net::timeout(rkvm_net::READ_TIMEOUT, AuthStatus::decode(&mut stream))
.await
.map_err(Error::Network)?;
pub fn init_tracing<P: AsRef<Path>>(log_level: &String, log_file: &Option<P>) {
let filter = EnvFilter::new(log_level);
if let Some(path) = log_file {
let path = path.as_ref();
// rotate log
for i in (1..10).rev() {
let old = format!("{}.{}", path.display(), i);
let old = Path::new(&old);
if old.exists() {
let new = format!("{}.{}", path.display(), i + 1);
let _ = rename(&old, &new);
}
}

match status {
AuthStatus::Passed => {}
AuthStatus::Failed => return Err(Error::Auth),
if path.exists() {
let new = format!("{}.1", path.display());
let _ = rename(path, &new);
}
let file = OpenOptions::new().create(true).append(true).open(path).unwrap();
let fmt_layer = fmt::layer().with_ansi(false).with_timer(LocalTime::rfc_3339()).with_writer(move || BufWriter::new(file.try_clone().unwrap()));
let registry = Registry::default().with(filter).with(fmt_layer);
tracing::subscriber::set_global_default(registry).unwrap();
} else {
let fmt_layer = fmt::layer().with_writer(stdout).without_time();
let registry = Registry::default().with(filter).with(fmt_layer);
tracing::subscriber::set_global_default(registry).unwrap();
}
}

tracing::info!("Authenticated successfully");
pub async fn run<R,W,H>(reader: &mut R, writer: &mut W, mut handler: H) -> Result<(), Error>
where
R: AsyncRead + Send + Unpin,
W: RkvmWriter + Send,
H: DeviceWriter {

let mut start = Instant::now();

let mut interval = time::interval(rkvm_net::PING_INTERVAL + rkvm_net::READ_TIMEOUT);
let mut writers = HashMap::new();

// Interval ticks immediately after creation.
interval.tick().await;
let timeout_duration = rkvm_net::PING_INTERVAL + rkvm_net::READ_TIMEOUT;

loop {
let update = tokio::select! {
update = Update::decode(&mut stream) => update.map_err(Error::Network)?,
_ = interval.tick() => return Err(Error::Network(io::Error::new(io::ErrorKind::TimedOut, "Ping timed out"))),
};

match update {
Update::CreateDevice {
id,
name,
vendor,
product,
version,
rel,
abs,
keys,
delay,
period,
} => {
let entry = writers.entry(id);
if let Entry::Occupied(_) = entry {
return Err(Error::Network(io::Error::new(
io::ErrorKind::InvalidData,
"Server created the same device twice",
)));
}

let writer = async {
Writer::builder()?
.name(&name)
.vendor(vendor)
.product(product)
.version(version)
.rel(rel)?
.abs(abs)?
.key(keys)?
.delay(delay)?
.period(period)?
.build()
.await
}
.await
.map_err(Error::Input)?;
let update = match time::timeout(timeout_duration, Update::decode(reader)).await {
Err(_) => Err(Error::Network(io::Error::new(io::ErrorKind::TimedOut, "Ping timeout"))),
Ok(res) => res.map_err(Error::Network)
}?;

entry.or_insert(writer);
let duration = start.elapsed();
tracing::debug!(duration = ?duration, "received {:?}", update);
start = Instant::now();

tracing::info!(
id = %id,
name = ?name,
vendor = %vendor,
product = %product,
version = %version,
"Created new device"
);
match update {
Update::CreateDevice { id,name,vendor,product,version,rel,abs,keys,delay,period,} => {
handler.create_device(id, &name, vendor, product, version, rel, abs, keys, delay, period).await?;
tracing::info!(id = %id, name = ?name, vendor = %vendor, product = %product, version = %version, "Created new device");
}
Update::DestroyDevice { id } => {
if writers.remove(&id).is_none() {
return Err(Error::Network(io::Error::new(
io::ErrorKind::InvalidData,
"Server destroyed a nonexistent device",
)));
}

handler.destroy_device(id).await?;
tracing::info!(id = %id, "Destroyed device");
}
Update::Event { id, event } => {
let writer = writers.get_mut(&id).ok_or_else(|| {
Error::Network(io::Error::new(
io::ErrorKind::InvalidData,
"Server sent an event to a nonexistent device",
))
})?;

writer.write(&event).await.map_err(Error::Input)?;

handler.event(id, event).await?;
tracing::trace!(id = %id, "Wrote an event to device");
}
Update::Ping => {
let duration = start.elapsed();
tracing::debug!(duration = ?duration, "Received ping");

start = Instant::now();
interval.reset();

rkvm_net::timeout(rkvm_net::WRITE_TIMEOUT, async {
Pong.encode(&mut stream).await?;
stream.flush().await?;

Ok(())
})
.await
.map_err(Error::Network)?;

writer.send(Update::Pong).await?;
let duration = start.elapsed();
tracing::debug!(duration = ?duration, "Sent pong");
}
Update::Pong => {}
Update::Stop => {
tracing::info!("Stoping..");
return Ok(());
}
}
}
}

#[async_trait]
pub trait RkvmWriter {
async fn send(&mut self, update: Update) -> Result<(), Error>;
}

#[async_trait]
impl<W> RkvmWriter for W
where
W: AsyncWrite + Unpin + Send,
{
async fn send(&mut self, update: Update) -> Result<(), Error> {
update.encode(self).await?;
self.flush().await?;
Ok(())
}
}
Loading