This commit is contained in:
0m.ax 2026-07-18 20:57:58 +02:00
parent 28d5d6db14
commit c9e10cd2c3
6 changed files with 286 additions and 74 deletions

119
Cargo.lock generated
View file

@ -149,6 +149,12 @@ dependencies = [
"cfg-if",
]
[[package]]
name = "equivalent"
version = "1.0.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "877a4ace8713b0bcf2a4e7eec82529c029f1d0619886d18145fea96c3ffe5c0f"
[[package]]
name = "fdeflate"
version = "0.3.7"
@ -177,7 +183,9 @@ dependencies = [
"libc",
"nix",
"png",
"serde",
"socket2",
"toml",
]
[[package]]
@ -190,12 +198,28 @@ dependencies = [
"weezl",
]
[[package]]
name = "hashbrown"
version = "0.17.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ed5909b6e89a2db4456e54cd5f673791d7eca6732202bbf2a9cc504fe2f9b84a"
[[package]]
name = "heck"
version = "0.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea"
[[package]]
name = "indexmap"
version = "2.14.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d466e9454f08e4a911e14806c24e16fba1b4c121d1ea474396f396069cf949d9"
dependencies = [
"equivalent",
"hashbrown",
]
[[package]]
name = "is_terminal_polyfill"
version = "1.70.2"
@ -208,6 +232,12 @@ version = "0.2.174"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1171693293099992e19cddea4e8b849964e9846f4acee11b3948bcc337be8776"
[[package]]
name = "memchr"
version = "2.8.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "cf8baf1c55e62ffcace7a9f06f4bd9cd3f0c4beb022d3b367256b91b87513d98"
[[package]]
name = "memoffset"
version = "0.9.1"
@ -277,6 +307,45 @@ dependencies = [
"proc-macro2",
]
[[package]]
name = "serde"
version = "1.0.228"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9a8e94ea7f378bd32cbbd37198a4a91436180c5bb472411e48b5ec2e2124ae9e"
dependencies = [
"serde_core",
"serde_derive",
]
[[package]]
name = "serde_core"
version = "1.0.228"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "41d385c7d4ca58e59fc732af25c3983b67ac852c1a25000afe1175de458b67ad"
dependencies = [
"serde_derive",
]
[[package]]
name = "serde_derive"
version = "1.0.228"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d540f220d3187173da220f885ab66608367b6574e925011a9353e4badda91d79"
dependencies = [
"proc-macro2",
"quote",
"syn",
]
[[package]]
name = "serde_spanned"
version = "0.6.9"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "bf41e0cfaf7226dca15e8197172c295a782857fcb97fad1808a166870dee75a3"
dependencies = [
"serde",
]
[[package]]
name = "simd-adler32"
version = "0.3.7"
@ -310,6 +379,47 @@ dependencies = [
"unicode-ident",
]
[[package]]
name = "toml"
version = "0.8.23"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "dc1beb996b9d83529a9e75c17a1686767d148d70663143c7854d8b4a09ced362"
dependencies = [
"serde",
"serde_spanned",
"toml_datetime",
"toml_edit",
]
[[package]]
name = "toml_datetime"
version = "0.6.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "22cddaf88f4fbc13c51aebbf5f8eceb5c7c5a9da2ac40a13519eb5b0a0e8f11c"
dependencies = [
"serde",
]
[[package]]
name = "toml_edit"
version = "0.22.27"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "41fe8c660ae4257887cf66394862d21dbca4a6ddd26f04a3560410406a2f819a"
dependencies = [
"indexmap",
"serde",
"serde_spanned",
"toml_datetime",
"toml_write",
"winnow",
]
[[package]]
name = "toml_write"
version = "0.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5d99f8c9a7727884afe522e9bd5edbfc91a3312b36a77b5fb8926e4c31a41801"
[[package]]
name = "unicode-ident"
version = "1.0.24"
@ -364,3 +474,12 @@ checksum = "ae137229bcbd6cdf0f7b80a31df61766145077ddf49416a728b02cb3921ff3fc"
dependencies = [
"windows-link",
]
[[package]]
name = "winnow"
version = "0.7.15"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "df79d97927682d2fd8adb29682d1140b343be4ac0f08fd68b7765d9c059d3945"
dependencies = [
"memchr",
]

View file

@ -13,3 +13,5 @@ nix = { version = "0.29", features = ["socket", "uio"] }
socket2 = "0.4.7"
clap = { version = "4", features = ["derive"] }
gif = "0.13"
serde = { version = "1", features = ["derive"] }
toml = "0.8"

14
config.toml Normal file
View file

@ -0,0 +1,14 @@
# Each [[sender]] block defines a destination with its own thread group.
# All senders share the same frame data.
[[sender]]
address = "100.65.0.2"
port = 5005
interface = "eth0"
threads = 5
# Example: send to a second display from a different interface
# [[sender]]
# address = "192.168.1.50"
# port = 5005
# threads = 3

54
src/config.rs Normal file
View file

@ -0,0 +1,54 @@
use serde::Deserialize;
use std::fs;
use std::path::Path;
#[derive(Deserialize)]
pub struct Config {
pub sender: Vec<SenderConfig>,
}
#[derive(Deserialize, Clone)]
pub struct SenderConfig {
/// Destination IP address or hostname.
pub address: String,
/// Destination port (default: 5005).
#[serde(default = "default_port")]
pub port: u16,
/// Source network interface (e.g. "eth0"). Auto-discovered if omitted.
pub interface: Option<String>,
/// Number of packer/sender thread pairs (default: 5).
#[serde(default = "default_threads")]
pub threads: usize,
}
fn default_port() -> u16 {
5005
}
fn default_threads() -> usize {
5
}
impl Config {
pub fn load(path: &Path) -> Result<Self, String> {
let content = fs::read_to_string(path)
.map_err(|e| format!("Failed to read config file '{}': {}", path.display(), e))?;
let config: Config = toml::from_str(&content)
.map_err(|e| format!("Failed to parse config file '{}': {}", path.display(), e))?;
if config.sender.is_empty() {
return Err("Config must define at least one [[sender]]".into());
}
for (i, s) in config.sender.iter().enumerate() {
if s.threads == 0 {
return Err(format!("sender[{}]: threads must be >= 1", i));
}
}
Ok(config)
}
}

View file

@ -4,6 +4,7 @@ use std::sync::{Arc, Condvar, Mutex};
use std::thread;
use std::time::Instant;
use crate::config::SenderConfig;
use crate::raw_socket::RawSender;
/// Per-thread packet queue capacity.
@ -18,9 +19,6 @@ const MSGSIZE: usize = 2 + 7 * PIXELS_PER_PACKET;
/// Maximum number of packets to batch into a single sendmmsg call.
const SEND_BATCH_SIZE: usize = 26;
/// Number of packer/sender thread pairs.
const NUM_THREADS: usize = 5;
// ---------------------------------------------------------------------------
// Packet + per-thread queue
// ---------------------------------------------------------------------------
@ -100,11 +98,14 @@ fn fisher_yates_shuffle(arr: &mut [usize], rng: &mut Xorshift64) {
/// |
/// publishes Arc<Vec<[u8;7]>> frame snapshot
/// |
/// 5 packer threads (each with unique noise permutation)
/// N packer threads (each with unique noise permutation)
/// continuously iterate the frame, packing pixels into packets
/// |
/// 5 sender threads (one per packer, own packet queue)
/// N sender threads (one per packer, own packet queue)
/// send via RawSender (AF_PACKET) or fallback UDP
///
/// Multiple sender configs are supported -- each config spawns its own group
/// of packer/sender thread pairs targeting a different destination.
pub struct Display {
/// Frame being assembled by the main thread.
building_frame: Vec<[u8; 7]>,
@ -114,68 +115,91 @@ pub struct Display {
}
impl Display {
/// Creates a new Display and spawns packer + sender thread pairs.
pub fn new(host: &str, port: u16, interface: Option<&str>) -> Self {
/// Creates a new Display and spawns packer + sender thread pairs for
/// each sender in the config.
pub fn new(senders: &[SenderConfig]) -> Self {
let current_frame = Arc::new(Mutex::new(Arc::new(Vec::new())));
for thread_idx in 0..NUM_THREADS {
// Each packer-sender pair gets its own queue.
let queue = Arc::new(SenderQueue {
pending: Mutex::new(VecDeque::with_capacity(QUEUE_LEN)),
condvar: Condvar::new(),
pool: Mutex::new({
let mut v = Vec::with_capacity(QUEUE_LEN);
for _ in 0..QUEUE_LEN {
v.push(alloc_buf());
let mut global_thread_idx: usize = 0;
for (sender_idx, sender) in senders.iter().enumerate() {
let host = &sender.address;
let port = sender.port;
let interface = sender.interface.as_deref();
let num_threads = sender.threads;
eprintln!(
"[display] Sender {}: {}:{}, interface={}, threads={}",
sender_idx,
host,
port,
interface.unwrap_or("auto"),
num_threads,
);
for t in 0..num_threads {
// Each packer-sender pair gets its own queue.
let queue = Arc::new(SenderQueue {
pending: Mutex::new(VecDeque::with_capacity(QUEUE_LEN)),
condvar: Condvar::new(),
pool: Mutex::new({
let mut v = Vec::with_capacity(QUEUE_LEN);
for _ in 0..QUEUE_LEN {
v.push(alloc_buf());
}
v
}),
});
// Spawn packer thread.
let frame_source = Arc::clone(&current_frame);
let packer_queue = Arc::clone(&queue);
// Each thread gets a globally unique seed.
let seed = (global_thread_idx as u64 + 1).wrapping_mul(0x517cc1b727220a95);
let packer_name = format!("packer-{}-{}", sender_idx, t);
thread::Builder::new()
.name(packer_name)
.spawn(move || {
packer_loop(frame_source, packer_queue, seed);
})
.unwrap();
// Spawn sender thread.
let sender_queue = Arc::clone(&queue);
let sender_name = format!("sender-{}-{}", sender_idx, t);
match RawSender::new(host, port, 0, interface) {
Ok(raw_sender) => {
eprintln!("[display] {}: raw AF_PACKET sender", sender_name);
thread::Builder::new()
.name(sender_name)
.spawn(move || {
sender_loop_raw(raw_sender, sender_queue);
})
.unwrap();
}
Err(e) => {
eprintln!(
"[display] {}: raw unavailable ({}), falling back to UDP",
sender_name, e
);
let remote_addr = (host.as_str(), port)
.to_socket_addrs()
.expect("Invalid remote address")
.next()
.expect("Could not resolve host");
let socket = UdpSocket::bind("0.0.0.0:0").expect("Could not bind");
socket.connect(remote_addr).expect("Could not connect");
thread::Builder::new()
.name(sender_name)
.spawn(move || {
sender_loop_udp(socket, sender_queue);
})
.unwrap();
}
v
}),
});
// Spawn packer thread.
let frame_source = Arc::clone(&current_frame);
let packer_queue = Arc::clone(&queue);
// Each thread gets a unique seed for a different permutation.
let seed = (thread_idx as u64 + 1).wrapping_mul(0x517cc1b727220a95);
thread::Builder::new()
.name(format!("packer-{}", thread_idx))
.spawn(move || {
packer_loop(frame_source, packer_queue, seed);
})
.unwrap();
// Spawn sender thread.
let sender_queue = Arc::clone(&queue);
match RawSender::new(host, port, 0, interface) {
Ok(raw_sender) => {
eprintln!("[display] Thread {}: raw AF_PACKET sender", thread_idx);
thread::Builder::new()
.name(format!("sender-{}", thread_idx))
.spawn(move || {
sender_loop_raw(raw_sender, sender_queue);
})
.unwrap();
}
Err(e) => {
eprintln!(
"[display] Thread {}: raw unavailable ({}), falling back to UDP",
thread_idx, e
);
let remote_addr = (host, port)
.to_socket_addrs()
.expect("Invalid remote address")
.next()
.expect("Could not resolve host");
let socket = UdpSocket::bind("0.0.0.0:0").expect("Could not bind");
socket.connect(remote_addr).expect("Could not connect");
thread::Builder::new()
.name(format!("sender-{}", thread_idx))
.spawn(move || {
sender_loop_udp(socket, sender_queue);
})
.unwrap();
}
global_thread_idx += 1;
}
}

View file

@ -1,4 +1,5 @@
use std::net::{SocketAddr, UdpSocket};
use std::path::PathBuf;
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::{Arc, Barrier, Mutex, RwLock};
use std::thread;
@ -10,6 +11,7 @@ use socket2::{Domain, Socket, Type};
mod bouncing_image;
mod circle;
mod color;
mod config;
mod display;
mod drawable;
mod pixel_buf;
@ -20,6 +22,7 @@ mod gif_image;
use bouncing_image::BouncingImage;
use circle::Circle;
use config::Config;
use display::Display;
use drawable::Drawable;
use gif_image::GifImage;
@ -31,16 +34,8 @@ pub const DISPLAY_HEIGHT: i32 = 1080;
#[derive(Parser)]
#[command(name = "flood-rs", about = "Pixel flooding tool")]
struct Args {
/// Target IP address or hostname
target: String,
/// Target port
#[arg(short, long, default_value_t = 5005)]
port: u16,
/// Source network interface (e.g. eth0). Auto-discovered if not set.
#[arg(short = 'i', long)]
interface: Option<String>,
/// Path to the config file
config: PathBuf,
}
/// Unpacks a 4-byte slice into two u16 values (little-endian).
@ -57,6 +52,10 @@ fn unpack_coordinates(buffer: &[u8]) -> Option<(u16, u16)> {
fn main() {
let args = Args::parse();
let config = Config::load(&args.config).unwrap_or_else(|e| {
eprintln!("Error: {}", e);
std::process::exit(1);
});
let x: Arc<Mutex<u32>> = Arc::new(Mutex::new(0));
let x_thread = x.clone();
@ -79,7 +78,7 @@ fn main() {
Box::new(GifImage::new("images/netto.gif", 200, 200)),
];
let mut display = Display::new(&args.target, args.port, args.interface.as_deref());
let mut display = Display::new(&config.sender);
// Spawn a UDP listener thread for receiving coordinates
thread::spawn(move || {