Initial commit: Project structure, Coordinator Server v1.0.0, and Documentation
This commit is contained in:
commit
0ead98d0c4
18 files changed
+2073
No files matched your search
+95
@@ -0,0 +1,95 @@
|
||||
use bytes::{Bytes, BytesMut};
|
||||
use crossbeam_queue::ArrayQueue;
|
||||
use std::sync::Arc;
|
||||
use tokio::net::UdpSocket;
|
||||
use std::io;
|
||||
|
||||
mod noise;
|
||||
|
||||
/// Configuration for the data plane memory limits.
|
||||
const MAX_PACKET_SIZE: usize = 16384; // 16KB
|
||||
const PACKET_POOL_SIZE: usize = 10000; // ~160MB (10k * 16KB)
|
||||
const MAX_RAM_USAGE_MB: usize = 256;
|
||||
|
||||
/// A fixed-size packet pool to prevent excessive allocations and fragmentation.
|
||||
struct PacketPool {
|
||||
pool: Arc<ArrayQueue<BytesMut>>,
|
||||
}
|
||||
|
||||
impl PacketPool {
|
||||
fn new(capacity: usize, packet_size: usize) -> Self {
|
||||
let queue = ArrayQueue::new(capacity);
|
||||
for _ in 0..capacity {
|
||||
let _ = queue.push(BytesMut::with_capacity(packet_size));
|
||||
}
|
||||
Self {
|
||||
pool: Arc::new(queue),
|
||||
}
|
||||
}
|
||||
|
||||
/// Acquires a buffer from the pool or creates a new one if the pool is empty
|
||||
/// (though in a strict memory-constrained environment, we might prefer to drop packets).
|
||||
fn acquire(&self) -> BytesMut {
|
||||
self.pool.pop().unwrap_or_else(|| BytesMut::with_capacity(MAX_PACKET_SIZE))
|
||||
}
|
||||
|
||||
/// Returns a buffer to the pool for reuse.
|
||||
fn release(&self, mut buf: BytesMut) {
|
||||
buf.clear();
|
||||
let _ = self.pool.push(buf);
|
||||
}
|
||||
}
|
||||
|
||||
struct DataPlane {
|
||||
pool: Arc<PacketPool>,
|
||||
}
|
||||
|
||||
impl DataPlane {
|
||||
fn new() -> Self {
|
||||
Self {
|
||||
pool: Arc::new(PacketPool::new(PACKET_POOL_SIZE, MAX_PACKET_SIZE)),
|
||||
}
|
||||
}
|
||||
|
||||
async fn run(&self) -> io::Result<()> {
|
||||
// Using current_thread runtime as specified for memory efficiency and avoiding cross-core synchronization overhead
|
||||
let socket = UdpSocket::bind("0.0.0.0:0").await?;
|
||||
println!("Data plane listening on {}", socket.local_addr()?);
|
||||
|
||||
loop {
|
||||
let mut buf = self.pool.acquire();
|
||||
|
||||
// Zero-copy receiving: reading directly into the pooled buffer
|
||||
match socket.recv_from(&mut buf).await {
|
||||
Ok((len, addr)) => {
|
||||
// Use BytesMut::split_to to create a zero-copy 'Bytes' view for processing
|
||||
let packet = buf.split_to(len).freeze();
|
||||
|
||||
// Simulate processing the packet without copying data
|
||||
self.process_packet(packet, addr).await;
|
||||
|
||||
// Return the remaining buffer to the pool
|
||||
self.pool.release(buf);
|
||||
}
|
||||
Err(e) => {
|
||||
eprintln!("Error receiving packet: {}", e);
|
||||
self.pool.release(buf);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn process_packet(&self, packet: Bytes, addr: std::net::SocketAddr) {
|
||||
// Logic for handling packets would go here.
|
||||
// 'packet' is a reference-counted view into the original buffer.
|
||||
let _ = packet.len();
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::main(flavor = "current_thread")]
|
||||
async fn main() -> io::Result<()> {
|
||||
let dp = DataPlane::new();
|
||||
println!("Starting memory-efficient data plane...");
|
||||
println!("Estimated pool memory: {} MB", (PACKET_POOL_SIZE * MAX_PACKET_SIZE) / 1024 / 1024);
|
||||
dp.run().await
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
use snow::params::NoiseParams;
|
||||
use snow::{HandshakeState, Noise};
|
||||
|
||||
pub struct NoiseSession {
|
||||
pub state: HandshakeState,
|
||||
}
|
||||
|
||||
impl NoiseSession {
|
||||
pub fn new_initiator(static_key: &[u8], remote_static_key: &[u8]) -> Self {
|
||||
let params = NoiseParams::Noise_XX_25519_ChaChaPoly_BLAKE2s;
|
||||
let builder = Noise::new(¶ms).unwrap();
|
||||
let state = builder.initiate_handshake(
|
||||
snow::keys::StaticKey::from_slice(static_key).unwrap(),
|
||||
snow::keys::PublicKey::from_slice(remote_static_key).unwrap(),
|
||||
).unwrap();
|
||||
|
||||
NoiseSession { state }
|
||||
}
|
||||
|
||||
pub fn new_responder(static_key: &[u8]) -> Self {
|
||||
let params = NoiseParams::Noise_XX_25519_ChaChaPoly_BLAKE2s;
|
||||
let builder = Noise::new(¶ms).unwrap();
|
||||
let state = builder.respond_handshake(
|
||||
snow::keys::StaticKey::from_slice(static_key).unwrap(),
|
||||
None,
|
||||
).unwrap();
|
||||
|
||||
NoiseSession { state }
|
||||
}
|
||||
|
||||
pub fn write_message(&mut self, payload: &[u8]) -> Vec<u8> {
|
||||
let mut buf = vec![0u8; 65535];
|
||||
let len = self.state.write_message(payload, &mut buf).unwrap();
|
||||
buf.truncate(len);
|
||||
buf
|
||||
}
|
||||
|
||||
pub fn read_message(&mut self, ciphertext: &[u8]) -> Vec<u8> {
|
||||
let mut buf = vec![0u8; 65535];
|
||||
let len = self.state.read_message(ciphertext, &mut buf).unwrap();
|
||||
buf.truncate(len);
|
||||
buf
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn test_handshake() {
|
||||
let alice_static = [0u8; 32];
|
||||
let bob_static = [1u8; 32];
|
||||
|
||||
let mut alice = NoiseSession::new_initiator(&alice_static, &bob_static);
|
||||
let mut bob = NoiseSession::new_responder(&bob_static);
|
||||
|
||||
let msg1 = alice.write_message(b"Hello Bob");
|
||||
let res1 = bob.read_message(&msg1);
|
||||
assert_eq!(res1, b"Hello Bob");
|
||||
|
||||
let msg2 = bob.write_message(b"Hello Alice");
|
||||
let res2 = alice.read_message(&msg2);
|
||||
assert_eq!(res2, b"Hello Alice");
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user