Skip to content
Merged
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
43 changes: 16 additions & 27 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

5 changes: 4 additions & 1 deletion wallguard-cli/src/update.rs
Original file line number Diff line number Diff line change
Expand Up @@ -205,7 +205,10 @@ async fn apply_update(version: &str) -> AnyResult<()> {
// for the lock at all.
let new_version = poll_agent_version(20, Duration::from_millis(500)).await;

if new_version.as_deref().is_some_and(|v| versions_match(v, version)) {
if new_version
.as_deref()
.is_some_and(|v| versions_match(v, version))
{
let _ = std::fs::remove_file(&backup_path);
println!("WallGuard successfully updated to v{version}.");
return Ok(());
Expand Down
2 changes: 1 addition & 1 deletion wallguard-server/src/token.rs
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@ mod tests {
fn test_token() {
let jwt = "eyJ0eXAiOiJKV1QiLCJhbGciOiJIUzI1NiJ9.eyJhY2NvdW50Ijp7ImFjY291bnRfaWQiOiJxd2JxNDZqcWNsZXYiLCJhY2NvdW50X29yZ2FuaXphdGlvbl9pZCI6bnVsbCwiYWNjb3VudF9zdGF0dXMiOiJBY3RpdmUiLCJjb250YWN0Ijp7fSwiZGV2aWNlIjp7fSwiaWQiOiIwMUtSUE41RUpLUEhUSldUM1Y2WUM3NjhXQSIsIm9yZ2FuaXphdGlvbiI6eyJjYXRlZ29yaWVzIjpbIlBlcnNvbmFsIl0sImNvZGUiOiJPMDAwMDIxIiwiaWQiOiIwMUtSUE41RVA3QTdOWjJWQTU2SjBBOVNZSyIsIm5hbWUiOiJQZXJzb25hbCBPcmdhbml6YXRpb24iLCJvcmdhbml6YXRpb25faWQiOiIwMUtSUE41RVA3QTdOWjJWQTU2SjBBOVNZSyIsInBhcmVudF9vcmdhbml6YXRpb25faWQiOm51bGwsInN0YXR1cyI6IkFjdGl2ZSJ9LCJvcmdhbml6YXRpb25faWQiOiIwMUtSUE41RVA3QTdOWjJWQTU2SjBBOVNZSyIsInByb2ZpbGUiOnsiYWNjb3VudF9pZCI6IjAxS1JQTjVFSktQSFRKV1QzVjZZQzc2OFdBIiwiY2F0ZWdvcmllcyI6W10sImNvZGUiOm51bGwsImVtYWlsIjoicXdicTQ2anFjbGV2IiwiZmlyc3RfbmFtZSI6IiIsImlkIjoiMDFLUlBONUZCQ1cwRU5GWDNDNTBGOE01TVYiLCJsYXN0X25hbWUiOiIiLCJvcmdhbml6YXRpb25faWQiOiIwMUtSUE41RVA3QTdOWjJWQTU2SjBBOVNZSyIsInN0YXR1cyI6IkFjdGl2ZSJ9LCJyb2xlX2lkIjpudWxsLCJzZXNzaW9uSUQiOiIifSwiZXhwIjoxNzc4OTYzOTM4LCJpYXQiOjE3Nzg4Nzc1MzgsInJvbGVfbmFtZSI6IiIsInNlbnNpdGl2aXR5X2xldmVsIjoxMDAwLCJzZXNzaW9uSUQiOiIiLCJzaWduZWRfaW5fYWNjb3VudCI6eyJhY2NvdW50X2lkIjoicXdicTQ2anFjbGV2IiwiYWNjb3VudF9vcmdhbml6YXRpb25faWQiOm51bGwsImFjY291bnRfc3RhdHVzIjoiQWN0aXZlIiwiY29udGFjdCI6e30sImRldmljZSI6e30sImlkIjoiMDFLUlBONUVKS1BIVEpXVDNWNllDNzY4V0EiLCJvcmdhbml6YXRpb24iOnsiY2F0ZWdvcmllcyI6WyJQZXJzb25hbCJdLCJjb2RlIjoiTzAwMDAyMSIsImlkIjoiMDFLUlBONUVQN0E3TloyVkE1NkowQTlTWUsiLCJuYW1lIjoiUGVyc29uYWwgT3JnYW5pemF0aW9uIiwib3JnYW5pemF0aW9uX2lkIjoiMDFLUlBONUVQN0E3TloyVkE1NkowQTlTWUsiLCJwYXJlbnRfb3JnYW5pemF0aW9uX2lkIjpudWxsLCJzdGF0dXMiOiJBY3RpdmUifSwib3JnYW5pemF0aW9uX2lkIjoiMDFLUlBONUVQN0E3TloyVkE1NkowQTlTWUsiLCJwcm9maWxlIjp7ImFjY291bnRfaWQiOiIwMUtSUE41RUpLUEhUSldUM1Y2WUM3NjhXQSIsImNhdGVnb3JpZXMiOltdLCJjb2RlIjpudWxsLCJlbWFpbCI6InF3YnE0NmpxY2xldiIsImZpcnN0X25hbWUiOiIiLCJpZCI6IjAxS1JQTjVGQkNXMEVORlgzQzUwRjhNNU1WIiwibGFzdF9uYW1lIjoiIiwib3JnYW5pemF0aW9uX2lkIjoiMDFLUlBONUVQN0E3TloyVkE1NkowQTlTWUsiLCJzdGF0dXMiOiJBY3RpdmUifSwicm9sZV9pZCI6bnVsbCwic2Vzc2lvbklEIjoiIn19.uCwpxbfDo6-3v-2hkbgPisbEo0GzMaQUv9SxKXIhWfo";

let result = Token::from_jwt(&jwt);
let result = Token::from_jwt(jwt);

assert!(result.is_ok());
}
Expand Down
2 changes: 1 addition & 1 deletion wallguard-server/src/tunneling/timeout_controller.rs
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ impl TimeoutController {

let lock = self.tunnels.lock().await;

for (_, tunnel) in lock.iter() {
for tunnel in lock.values() {
match tunnel {
WallguardTunnel::Http(http_tunnel) => {
let tun = http_tunnel.lock().await;
Expand Down
6 changes: 3 additions & 3 deletions wallguard/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,8 @@ chrono = "0.4.41"
once_cell = "1.21.4"
nullnet-traffic-monitor = "0.1.6"
etherparse = "0.19.0"
sysinfo = { version = "0.37.2", default-features = false, features = ["disk"] }
sysinfo = { version = "0.37.2", default-features = false, features = ["disk", "system", "component"] }
async-channel = "2.3.1"
nullnet-libresmon = "0.1.2"
wallguard-common = { path = "../wallguard-common" }
xmltree = "0.12.0"
md5 = "0.8.0"
Expand All @@ -37,6 +36,7 @@ nullnet-liberror.workspace = true
rustls.workspace = true
tokio-rustls.workspace = true
whoami = "2.0.2"
listeners = "0.6"

[target.'cfg(any(target_os = "linux", target_os = "freebsd"))'.dependencies]
x11rb = "0.13"
Expand All @@ -52,7 +52,7 @@ evdev = "0.12"

[target.'cfg(windows)'.dependencies]
is_elevated = "0.1.2"
winapi = {version = "0.3.9", features = ["winerror", "iphlpapi", "handleapi", "tlhelp32", "wingdi", "winuser", "windef", "minwindef"]}
winapi = {version = "0.3.9", features = ["wingdi", "winuser", "windef", "minwindef"]}

[target.'cfg(unix)'.dependencies]
nix = { version = "0.31.3", features = ["user", "mman", "fs"] }
27 changes: 17 additions & 10 deletions wallguard/src/data_transmission/dump_dir.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,22 +55,29 @@ impl DumpDir {
pub(crate) async fn dump_item_to_file(&self, dump_item: DumpItem) {
let now = chrono::Utc::now().to_rfc3339();
let file_path = self.get_file_path(&now, &dump_item);
tokio::fs::write(
file_path,
serde_json::to_string(&dump_item).expect("Failed to serialize item"),
)
// Serializing a dump item (up to the full queue, e.g. 1M records) is
// CPU-bound work; run it on the blocking pool so it can't stall the
// tokio runtime that also drives gRPC/heartbeat traffic.
let json = tokio::task::spawn_blocking(move || {
serde_json::to_string(&dump_item).expect("Failed to serialize item")
})
.await
.expect("Failed to write dump file");
.expect("Serialization task panicked");
tokio::fs::write(file_path, json)
.await
.expect("Failed to write dump file");
}

pub(crate) async fn update_items_dump_file(&self, file_path: PathBuf, mut dump: DumpItem) {
dump.set_token(String::new());
tokio::fs::write(
file_path,
serde_json::to_string(&dump).expect("Failed to serialize items"),
)
let json = tokio::task::spawn_blocking(move || {
serde_json::to_string(&dump).expect("Failed to serialize items")
})
.await
.expect("Failed to write dump file");
.expect("Serialization task panicked");
tokio::fs::write(file_path, json)
.await
.expect("Failed to write dump file");
}
}

Expand Down
47 changes: 32 additions & 15 deletions wallguard/src/data_transmission/grpc_handler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,30 +37,42 @@ pub(crate) async fn handle_connection_and_retransmission(
let Ok(string) = fs::read_to_string(file.path()).await else {
continue;
};
let Ok(mut dump) = serde_json::from_str::<DumpItem>(&string) else {
// Deserializing a dump file (up to the full queue, e.g. 1M
// records) is CPU-bound; keep it off the tokio runtime so it
// can't stall gRPC/heartbeat traffic sharing the same executor.
let Ok(Ok(mut dump)) =
tokio::task::spawn_blocking(move || serde_json::from_str::<DumpItem>(&string))
.await
else {
continue;
};
// update auth token of items retrieved from disk
dump.set_token(token.clone());

while dump.size() != 0 {
let range = ..min(dump.size(), BATCH_SIZE);
// `dump.set_token` above already updated the token field in
// place, so only the (cheap) token string needs cloning here
// — cloning the whole item via `..c.clone()` used to clone
// the entire, not-yet-drained items vector on every batch.
// Batches are sliced by a `sent` offset rather than drained from
// the front on every iteration: draining a Vec's front repeatedly
// shifts the remaining tail down each time (O(remaining) per
// batch), which turns replaying a large backlog file into an
// O(n^2) sequence of memmoves. Slicing leaves the vector
// untouched until a single drain(..sent) at the end.
let total = dump.size();
let mut sent = 0;
let mut failed = false;

while sent < total {
let range = ..min(total - sent, BATCH_SIZE);
let send_res = match &dump {
DumpItem::Connections(c) => {
let msg = ConnectionsData {
token: c.token.clone(),
connections: c.connections.get(range).unwrap_or_default().to_vec(),
connections: c.connections[sent..][range].to_vec(),
};
interface.handle_connections_data(msg).await
}
DumpItem::Resources(r) => {
let msg = SystemResourcesData {
token: r.token.clone(),
resources: r.resources.get(range).unwrap_or_default().to_vec(),
resources: r.resources[sent..][range].to_vec(),
};
interface.handle_system_resources_data(msg).await
}
Expand All @@ -75,13 +87,18 @@ pub(crate) async fn handle_connection_and_retransmission(
// back off before retrying instead of immediately
// re-reading and re-sending the same file in a tight loop.
log::warn!("Failed to send dump. Reconnecting...",);
// update dump file with unsent items
dump_dir.update_items_dump_file(file.path(), dump).await;
tokio::time::sleep(Duration::from_secs(10)).await;
break 'file_loop;
failed = true;
break;
}
// remove sent items from dump
dump.drain(range);
sent += range.end;
}

if failed {
// remove the items that did get sent, in one shot, and persist the rest
dump.drain(..sent);
dump_dir.update_items_dump_file(file.path(), dump).await;
tokio::time::sleep(Duration::from_secs(10)).await;
break 'file_loop;
}

log::info!("Dump file '{:?}' sent successfully", file.file_name());
Expand Down
16 changes: 11 additions & 5 deletions wallguard/src/data_transmission/item_buffer.rs
Original file line number Diff line number Diff line change
@@ -1,28 +1,34 @@
use std::collections::VecDeque;
use std::ops::RangeTo;

// Backed by a VecDeque (not a Vec) so that repeatedly draining a batch off
// the front — the access pattern every caller uses — costs O(batch), not
// O(remaining length): a Vec::drain(..batch) has to shift the whole
// remaining tail down on every call, which turns catching up a large backlog
// into an O(n^2) sequence of memmoves.
pub(crate) struct ItemBuffer<T> {
buffer: Vec<T>,
buffer: VecDeque<T>,
size: usize,
}

impl<T: Clone> ItemBuffer<T> {
pub(crate) fn new(size: usize) -> Self {
Self {
buffer: Vec::with_capacity(size),
buffer: VecDeque::with_capacity(size),
size,
}
}

pub(crate) fn push(&mut self, item: T) {
self.buffer.push(item);
self.buffer.push_back(item);
}

pub(crate) fn take(&mut self) -> Vec<T> {
std::mem::take(&mut self.buffer)
Vec::from(std::mem::take(&mut self.buffer))
}

pub(crate) fn get(&mut self, range: RangeTo<usize>) -> Vec<T> {
self.buffer.get(range).unwrap_or_default().to_vec()
self.buffer.iter().take(range.end).cloned().collect()
}

pub(crate) fn extend(&mut self, items: Vec<T>) {
Expand Down
7 changes: 6 additions & 1 deletion wallguard/src/data_transmission/packets/transmitter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,12 @@ pub(crate) async fn transmit_packets(
if raw_batch.len() >= batch_size || timer.is_expired() {
timer.reset();

let connections = parse_packets(std::mem::take(&mut raw_batch));
// Swap in a fresh, pre-sized buffer rather than `mem::take`ing
// (which would leave a 0-capacity Vec behind) so the next
// accumulation cycle doesn't have to regrow from scratch back up
// to `batch_size` on every flush.
let batch = std::mem::replace(&mut raw_batch, Vec::with_capacity(batch_size));
let connections = parse_packets(batch);
connection_queue.extend(connections);

send_connections(&client, &mut connection_queue, &token_provider, batch_size).await;
Expand Down
Loading
Loading