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
6 changes: 0 additions & 6 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -40,9 +40,6 @@ jobs:
uses: taiki-e/install-action@v2
with:
tool: nextest@${{ env.NEXTEST_VERSION }}
- name: Set up tun
run: |
sudo ./litebox_platform_linux_userland/scripts/tun-setup.sh
- name: Install iperf3
run: |
sudo apt install -y iperf3
Expand Down Expand Up @@ -136,9 +133,6 @@ jobs:
uses: taiki-e/install-action@v2
with:
tool: nextest@${{ env.NEXTEST_VERSION }}
- name: Set up tun
run: |
sudo ./litebox_platform_linux_userland/scripts/tun-setup.sh
- uses: Swatinem/rust-cache@v2
- name: Cache custom out directories
uses: actions/cache@v5
Expand Down
3 changes: 0 additions & 3 deletions .github/workflows/copilot-setup-steps.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,3 @@ jobs:
- name: Set up Nextest
run: |
curl -LsSf https://get.nexte.st/latest/linux | tar zxf - -C ${CARGO_HOME:-~/.cargo}/bin
- name: Set up tun device for Linux userland testing
run: |
sudo ./litebox_platform_linux_userland/scripts/tun-setup.sh
147 changes: 111 additions & 36 deletions dev_bench/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ use clap::Parser;
use std::sync::atomic::Ordering::Relaxed;
use std::{
collections::{BTreeMap, BTreeSet},
net::{Ipv4Addr, TcpListener},
path::{Path, PathBuf},
sync::atomic::AtomicBool,
time::Duration,
Expand Down Expand Up @@ -695,55 +696,129 @@ fn run_rewritten_iperf3(ctx: BenchCtx<'_>) -> Result<()> {
"cargo build -p litebox_runner_linux_userland {release...} {features...}"
)
.run()?;
cmd!(sh, "cargo build -p litebox_broker_userland {release...}").run()?;
} else {
let mode = if release_mode { "release" } else { "debug" };
let iperf3_host = locate_command(sh, "iperf3")?;
let runner = format!(
"{}/target/{mode}/litebox_runner_linux_userland",
project_root.display()
);
let broker = format!(
"{}/target/{mode}/litebox-broker-userland",
project_root.display()
);
let iperf3_rewritten = sh.current_dir().join("iperf3_rewritten");

// Spawn the sandboxed iperf3 server in a background thread so we can run the client
// from the host side. The server uses `-1` to exit after handling one client.
let server_handle = std::thread::spawn(move || -> Result<()> {
let sh = xshell::Shell::new()?;
cmd!(
sh,
"{runner} --unstable --env LD_LIBRARY_PATH=/lib64:/lib32:/lib --env HOME=/ --tun-device-name tun99 --initial-files {tar_file} {iperf3_rewritten} -s -1 -B 10.0.0.2"
).run()?;
Ok(())
});

// Retry the client connection until the server is ready, using a short
// connect-timeout so we don't waste time sleeping for a fixed duration.
debug!("Connecting iperf3 client to sandboxed server");
let client_sh = xshell::Shell::new()?;
let max_attempts = 50;
for attempt in 1..=max_attempts {
let result = cmd!(
client_sh,
"{iperf3_host} -c 10.0.0.2 --bytes 1G --connect-timeout 50"
)
.quiet()
.ignore_stdout()
.ignore_stderr()
.run();
if result.is_ok() {
break;
let max_server_attempts = 5;
let max_client_attempts = 50;
let mut last_failure = "sandboxed server did not become ready".to_owned();

for server_attempt in 1..=max_server_attempts {
let port = TcpListener::bind((Ipv4Addr::LOCALHOST, 0))?
.local_addr()?
.port()
.to_string();
let mut server_command = std::process::Command::new(&broker);
server_command
.arg("--runner")
.arg(&runner)
.arg("--")
.args([
"--env",
"LD_LIBRARY_PATH=/lib64:/lib32:/lib",
"--env",
"HOME=/",
"--initial-files",
])
.arg(&tar_file)
.arg(&iperf3_rewritten)
.args(["-s", "-1", "-B", "127.0.0.1", "-p", &port]);
if COMMAND_EXECUTION_IS_QUIET.load(Relaxed) {
server_command
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null());
}
if attempt == max_attempts {
return Err(anyhow!(
"iperf3 client failed to connect after {max_attempts} attempts"
));
let mut server = server_command.spawn()?;

debug!(
server_attempt,
"Connecting iperf3 client to sandboxed server"
);
let mut client_succeeded = false;
for client_attempt in 1..=max_client_attempts {
let result = cmd!(
client_sh,
"{iperf3_host} -c 127.0.0.1 -p {port} --bytes 1G --connect-timeout 50"
)
.quiet()
.ignore_stdout()
.ignore_stderr()
.run();
if result.is_ok() {
client_succeeded = true;
break;
}
match server.try_wait() {
Ok(Some(status)) => {
last_failure = format!("sandboxed server exited with {status}");
break;
}
Ok(None) => {}
Err(error) => {
let _ = server.kill();
let _ = server.wait();
return Err(error.into());
}
}
if client_attempt == max_client_attempts {
last_failure = format!(
"iperf3 client failed to connect after {max_client_attempts} attempts"
);
break;
}
debug!(client_attempt, "iperf3 client connection failed, retrying");
std::thread::sleep(Duration::from_millis(100));
}

if !client_succeeded {
match server.try_wait() {
Ok(Some(_)) => {}
Ok(None) => {
let _ = server.kill();
let _ = server.wait();
}
Err(error) => {
let _ = server.kill();
let _ = server.wait();
return Err(error.into());
}
}
debug!(
server_attempt,
%last_failure,
"Relaunching sandboxed iperf3 server"
);
continue;
}

let status = match server.wait() {
Ok(status) => status,
Err(error) => {
let _ = server.kill();
let _ = server.wait();
return Err(error.into());
}
};
if !status.success() {
return Err(anyhow!("sandboxed iperf3 server exited with {status}"));
}
debug!(attempt, "iperf3 client connection failed, retrying");
std::thread::sleep(Duration::from_millis(100));
return Ok(());
}

server_handle
.join()
.map_err(|e| anyhow!("iperf3 server thread panicked: {e:?}"))??;
return Err(anyhow!(
"sandboxed iperf3 server failed after {max_server_attempts} attempts: {last_failure}"
));
}
Ok(())
}
2 changes: 2 additions & 0 deletions litebox/src/net/errors.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@ use thiserror::Error;
pub enum SocketError {
#[error("Unsupported protocol {0}")]
UnsupportedProtocol(u8),
#[error("Brokered networking is unavailable")]
BrokerUnavailable,
#[error("Socket resources are exhausted")]
ResourceExhausted,
#[error("Socket creation was denied")]
Expand Down
77 changes: 37 additions & 40 deletions litebox/src/net/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ use crate::sync::RawSyncPrimitivesProvider;
use crate::{LiteBox, platform, sync};

use bitflags::bitflags;
use smoltcp::socket::{icmp, raw, tcp, udp};
use smoltcp::socket::{tcp, udp};

mod broker_socket;
pub mod errors;
Expand Down Expand Up @@ -56,6 +56,7 @@ const INTERFACE_IP_ADDR: Ipv4Addr = Ipv4Addr::new(10, 0, 0, 2);
/// IP address for the gateway
// TODO: Make this configurable
const GATEWAY_IP_ADDR: Ipv4Addr = Ipv4Addr::new(10, 0, 0, 1);
const ICMP_PROTOCOL_NUMBER: u8 = 1;

/// Maximum size of rx/tx buffers for sockets
pub const SOCKET_BUFFER_SIZE: usize = 65536 * 4;
Expand All @@ -77,6 +78,7 @@ pub enum ShutdownDirection {
}

/// Limits maximum number of packets in a buffer
#[cfg(test)]
const MAX_PACKET_COUNT: usize = 32;

/// TCP connection timeout.
Expand Down Expand Up @@ -793,24 +795,37 @@ where
)
}

/// Creates a socket.
/// Creates a broker-owned TCP or UDP socket.
///
/// By default, the created socket has no associated proxy; to set a proxy, use
/// [`attach_socket_proxy`](Self::attach_socket_proxy).
pub fn socket(&mut self, protocol: Protocol) -> Result<SocketFd<Platform>, SocketError> {
let broker_socket = match (&protocol, self.litebox.broker_control()) {
(Protocol::Tcp, Some(broker)) => Some(BrokerSocket::Tcp(
(Protocol::Tcp, Some(broker)) => BrokerSocket::Tcp(
BrokerTcpSocket::new(broker, self.litebox.broker_pollable_registry())
.map_err(SocketError::from)?,
)),
(Protocol::Udp, Some(broker)) => Some(BrokerSocket::Udp(
),
(Protocol::Udp, Some(broker)) => BrokerSocket::Udp(
BrokerUdpSocket::new(broker, self.litebox.broker_pollable_registry())
.map_err(SocketError::from)?,
)),
_ => None,
),
(Protocol::Tcp | Protocol::Udp, None) => {
return Err(SocketError::BrokerUnavailable);
}
(Protocol::Icmp, _) => {
return Err(SocketError::UnsupportedProtocol(ICMP_PROTOCOL_NUMBER));
}
(Protocol::Raw { protocol }, _) => {
return Err(SocketError::UnsupportedProtocol(*protocol));
}
};

Ok(self.new_socket_fd(protocol, None, Some(broker_socket)))
}

#[cfg(test)]
fn local_socket(&mut self, protocol: Protocol) -> Result<SocketFd<Platform>, SocketError> {
let handle = match protocol {
Protocol::Tcp | Protocol::Udp if broker_socket.is_some() => None,
Protocol::Tcp => Some(self.socket_set.add(tcp::Socket::new(
smoltcp::storage::RingBuffer::new(vec![0u8; SOCKET_BUFFER_SIZE]),
smoltcp::storage::RingBuffer::new(vec![0u8; SOCKET_BUFFER_SIZE]),
Expand All @@ -825,42 +840,24 @@ where
vec![0u8; SOCKET_BUFFER_SIZE],
),
))),
Protocol::Icmp => Some(self.socket_set.add(icmp::Socket::new(
smoltcp::storage::PacketBuffer::new(
vec![smoltcp::storage::PacketMetadata::EMPTY; MAX_PACKET_COUNT],
vec![0u8; SOCKET_BUFFER_SIZE],
),
smoltcp::storage::PacketBuffer::new(
vec![smoltcp::storage::PacketMetadata::EMPTY; MAX_PACKET_COUNT],
vec![0u8; SOCKET_BUFFER_SIZE],
),
))),
Protocol::Icmp => {
return Err(SocketError::UnsupportedProtocol(ICMP_PROTOCOL_NUMBER));
}
Protocol::Raw { protocol } => {
// TODO: Should we maintain a specific allow-list of protocols for raw sockets?
// Should we allow everything except TCP/UDP/ICMP? Should we allow everything? These
// questions should be resolved; for now I am disallowing everything else.
return Err(SocketError::UnsupportedProtocol(protocol));

#[expect(
unreachable_code,
reason = "currently raw is just directly disallowed; we might bring this code back in the future"
)]
Some(self.socket_set.add(raw::Socket::new(
smoltcp::wire::IpVersion::Ipv4,
smoltcp::wire::IpProtocol::from(protocol),
smoltcp::storage::PacketBuffer::new(
vec![smoltcp::storage::PacketMetadata::EMPTY; MAX_PACKET_COUNT],
vec![0u8; SOCKET_BUFFER_SIZE],
),
smoltcp::storage::PacketBuffer::new(
vec![smoltcp::storage::PacketMetadata::EMPTY; MAX_PACKET_COUNT],
vec![0u8; SOCKET_BUFFER_SIZE],
),
)))
}
};

Ok(self.new_socket_fd_for(SocketHandle {
Ok(self.new_socket_fd(protocol, handle, None))
}

fn new_socket_fd(
&mut self,
protocol: Protocol,
handle: Option<smoltcp::iface::SocketHandle>,
broker_socket: Option<BrokerSocket<Platform>>,
) -> SocketFd<Platform> {
self.new_socket_fd_for(SocketHandle {
consider_closed: false,
handle,
broker_socket,
Expand All @@ -878,7 +875,7 @@ where
Protocol::Raw { protocol: _ } => unimplemented!(),
},
proxy: None,
}))
})
}

/// Creates a new [`SocketFd`] for a newly-created [`SocketHandle`].
Expand Down
Loading