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
280 changes: 271 additions & 9 deletions Cargo.lock

Large diffs are not rendered by default.

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

[features]
windows-service = []

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

# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html

[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 +28,21 @@ 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]
bincode = "1.3.3"
windows = { version = "0.62", features = [
"Win32_System_RemoteDesktop",
"Win32_Foundation",
"Win32_Security_Authorization",
"Win32_System_Threading",
"Win32_Security",
"Win32_System_Environment",
] }
windows-core = "0.62"
windows-sys = { version = "0.61", features = [ "Win32" ]}
windows-service = "0.8"

[package.metadata.rpm]
package = "rkvm-client"
Expand Down
247 changes: 140 additions & 107 deletions rkvm-client/src/client.rs
Original file line number Diff line number Diff line change
@@ -1,37 +1,70 @@
use rkvm_input::writer::Writer;
use crate::config::Config;
use crate::stream::{RkvmStream, RkvmWriter};

use rkvm_input::writer::{Writer,WriterPlatform,WriterBuilderPlatform};
use rkvm_net::auth::{AuthChallenge, AuthStatus};
use rkvm_net::message::Message;
use rkvm_net::version::Version;
use rkvm_net::{Pong, Update};
use rkvm_net::Update;
use std::collections::hash_map::Entry;
use std::collections::HashMap;
use std::io;
use std::fs::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::fs;
use tokio::io::{AsyncRead, AsyncWriteExt, BufStream};
use tokio::net::TcpStream;
use tokio::time;
use tokio_rustls::rustls::ServerName;
use tokio_rustls::rustls::{self, ServerName};
use tokio_rustls::TlsConnector;
use tracing_subscriber::{fmt, Registry,EnvFilter};
use tracing_subscriber::fmt::time::LocalTime;
use tracing_subscriber::prelude::*;
#[cfg(target_os="windows")]
use windows::core;

#[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] core::Error),
#[error("Incompatible server version (got {server}, expected {client})")]
Version { server: Version, client: Version },
#[error("Invalid password")]
Auth,
}

pub async fn run(
hostname: &ServerName,
port: u16,
connector: TlsConnector,
password: &str,
) -> Result<(), Error> {
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 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();
}
}

pub async fn init_config<P: AsRef<Path> + ?Sized> (path: &P) -> Result<Config,Error> {
let config = fs::read_to_string(path).await?;
let config = toml::from_str::<Config>(&config)?;
return Ok(config);
}

pub async fn init_stream(hostname: &ServerName, port: u16, connector: &TlsConnector, password: &str) -> Result<RkvmStream,Error> {
// Intentionally don't impose any timeout for TCP connect.
let stream = match hostname {
ServerName::DnsName(name) => TcpStream::connect(&(name.as_ref(), port)).await,
Expand Down Expand Up @@ -98,111 +131,111 @@ pub async fn run(
}

tracing::info!("Authenticated successfully");
Ok(RkvmStream::Tcp(stream))
}

let mut start = Instant::now();
pub async fn run<R,W,H>(mut reader: R, mut writer: W, mut handler: H) -> Result<(), Error>
where
R: AsyncRead + Send + Unpin,
W: RkvmWriter + Send,
H: AsyncFnMut(Update) -> Result<(), Error>, {

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

// 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)?;

entry.or_insert(writer);

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",
)));
}

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)?;
let update = match time::timeout(timeout_duration, Update::decode(&mut reader)).await {
Err(_) => Err(Error::Network(io::Error::new(io::ErrorKind::TimedOut, "Ping timeout"))),
Ok(res) => res.map_err(Error::Network)
}?;

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

if let Update::Ping = &update {
writer.send(Update::Pong).await?;
let duration = start.elapsed();
tracing::debug!(duration = ?duration, "Sent pong");
}
handler(update).await?
}
}

tracing::trace!(id = %id, "Wrote an event to device");
pub async fn handler(writers: &mut HashMap<usize,Writer>, update: Update) -> Result<(), Error> {
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",
)));
}
Update::Ping => {
let duration = start.elapsed();
tracing::debug!(duration = ?duration, "Received ping");

start = Instant::now();
interval.reset();
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?;

entry.or_insert(writer);

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",
)));
}

rkvm_net::timeout(rkvm_net::WRITE_TIMEOUT, async {
Pong.encode(&mut stream).await?;
stream.flush().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",
))
})?;

Ok(())
})
.await
.map_err(Error::Network)?;
writer.write(&event).await?;

let duration = start.elapsed();
tracing::debug!(duration = ?duration, "Sent pong");
}
tracing::trace!(id = %id, "Wrote an event to device");
}
Update::Ping => {}
Update::Pong => {}
}
Ok(())
}
Loading