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
250 changes: 249 additions & 1 deletion litebox/src/broker/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,11 @@ use litebox_broker_protocol::fs::{
FileAccessMode, FileDirectoryEntry, FileError, FileMode, FileOpenFlags, FileSeekWhence,
FileStatus, FileStatusFlags, FileUser, MAX_FILE_TRANSFER_SIZE,
};
use litebox_broker_protocol::local_socket::{
CreateLocalSocketPairResponse, GetLocalSocketOptionsResponse, LOCAL_SOCKET_BUFFER_SIZE,
LocalSocketAddress, LocalSocketError, LocalSocketName, LocalSocketOption,
MAX_ENCODED_LOCAL_SOCKET_NAME_SIZE, MAX_LOCAL_SOCKET_TRANSFER_SIZE, ReceiveLocalSocketResponse,
};
use litebox_broker_protocol::pipe::{CreatePipeResponse, MAX_PIPE_TRANSFER_SIZE};
use litebox_broker_protocol::process::{
CreatedProcess, MAX_CHILD_MEMORY_WRITE_SIZE, MAX_CHILD_OBJECT_DUPLICATES,
Expand All @@ -30,7 +35,8 @@ use litebox_broker_protocol::socket::{
AcceptSocketResponse, MAX_SOCKET_TRANSFER_SIZE, MAX_UDP_DATAGRAM_SIZE,
ReceiveFlags as BrokerReceiveFlags, ReceiveFromFlags as BrokerReceiveFromFlags,
ReceiveFromSocketResponse, ReceiveSocketResponse, SendFlags as BrokerSendFlags, ShutdownMode,
SocketConnectionStatus, SocketOutcome, SocketStatusResponse, TcpOptionName, TcpOptionValue,
SocketConnectionStatus, SocketOutcome, SocketStatusResponse, SocketType, TcpOptionName,
TcpOptionValue,
};
use litebox_broker_protocol::timer::TimerSpec;
use litebox_broker_transport::channel::LocalCallChannel;
Expand Down Expand Up @@ -357,6 +363,101 @@ pub(crate) trait BrokerControl: Send + Sync {
user: FileUser,
) -> core::result::Result<core::result::Result<(), FileError>, BrokerControlError>;

fn create_local_socket(
&self,
socket_type: SocketType,
flags: FileOpenFlags,
) -> core::result::Result<ObjectHandle, BrokerControlError>;

fn create_local_socket_pair(
&self,
socket_type: SocketType,
flags: FileOpenFlags,
) -> core::result::Result<CreateLocalSocketPairResponse, BrokerControlError>;

fn bind_local_socket(
&self,
handle: ObjectHandle,
address: &LocalSocketAddress,
user: FileUser,
mode: FileMode,
) -> core::result::Result<core::result::Result<(), LocalSocketError>, BrokerControlError>;

fn listen_local_socket(
&self,
handle: ObjectHandle,
backlog: u32,
) -> core::result::Result<core::result::Result<(), LocalSocketError>, BrokerControlError>;

fn connect_local_socket(
&self,
handle: ObjectHandle,
address: &LocalSocketAddress,
user: FileUser,
) -> core::result::Result<core::result::Result<(), LocalSocketError>, BrokerControlError>;

fn accept_local_socket(
&self,
handle: ObjectHandle,
flags: FileOpenFlags,
) -> core::result::Result<
core::result::Result<ObjectHandle, LocalSocketError>,
BrokerControlError,
>;

/// Sends `data` to `address`, or to the connected peer if `address` is
/// `None`.
///
/// Data that does not fit in one transfer fails with
/// [`LocalSocketError::MessageTooLarge`], so stream callers send it in
/// chunks of at most [`LOCAL_SOCKET_BUFFER_SIZE`] bytes.
fn send_local_socket(
&self,
handle: ObjectHandle,
address: Option<&LocalSocketAddress>,
data: &[u8],
user: FileUser,
) -> core::result::Result<core::result::Result<usize, LocalSocketError>, BrokerControlError>;

/// Receives into at most [`LOCAL_SOCKET_BUFFER_SIZE`] bytes of `data`, as
/// for a non-blocking socket if `nonblocking` is set.
fn receive_local_socket(
&self,
handle: ObjectHandle,
data: &mut [u8],
peek: bool,
nonblocking: bool,
) -> core::result::Result<
core::result::Result<(ReceiveLocalSocketResponse, LocalSocketName), LocalSocketError>,
BrokerControlError,
>;

fn shutdown_local_socket(
&self,
handle: ObjectHandle,
mode: ShutdownMode,
) -> core::result::Result<core::result::Result<(), LocalSocketError>, BrokerControlError>;

fn local_socket_name(
&self,
handle: ObjectHandle,
peer: bool,
) -> core::result::Result<
core::result::Result<LocalSocketName, LocalSocketError>,
BrokerControlError,
>;

fn set_local_socket_option(
&self,
handle: ObjectHandle,
option: LocalSocketOption,
) -> core::result::Result<(), BrokerControlError>;

fn local_socket_options(
&self,
handle: ObjectHandle,
) -> core::result::Result<GetLocalSocketOptionsResponse, BrokerControlError>;

fn close_object(&self, handle: ObjectHandle) -> core::result::Result<(), BrokerControlError>;

fn fail_connection(&self);
Expand Down Expand Up @@ -1066,6 +1167,153 @@ where
self.request(|local| local.rmdir_file(lease.sequence(), path, user))
}

fn create_local_socket(
&self,
socket_type: SocketType,
flags: FileOpenFlags,
) -> core::result::Result<ObjectHandle, BrokerControlError> {
self.request(|local| local.create_local_socket(socket_type, flags))
}

fn create_local_socket_pair(
&self,
socket_type: SocketType,
flags: FileOpenFlags,
) -> core::result::Result<CreateLocalSocketPairResponse, BrokerControlError> {
self.request(|local| local.create_local_socket_pair(socket_type, flags))
}

fn bind_local_socket(
&self,
handle: ObjectHandle,
address: &LocalSocketAddress,
user: FileUser,
mode: FileMode,
) -> core::result::Result<core::result::Result<(), LocalSocketError>, BrokerControlError> {
let Some(address) = address.encode() else {
return Ok(Err(LocalSocketError::InvalidArgument));
};
let lease = self.acquire_shared_buffer(address.len())?;
self.request(|local| {
local.bind_local_socket(handle, lease.sequence(), &address, user, mode)
})
}

fn listen_local_socket(
&self,
handle: ObjectHandle,
backlog: u32,
) -> core::result::Result<core::result::Result<(), LocalSocketError>, BrokerControlError> {
self.request(|local| local.listen_local_socket(handle, backlog))
}

fn connect_local_socket(
&self,
handle: ObjectHandle,
address: &LocalSocketAddress,
user: FileUser,
) -> core::result::Result<core::result::Result<(), LocalSocketError>, BrokerControlError> {
let Some(address) = address.encode() else {
return Ok(Err(LocalSocketError::InvalidArgument));
};
let lease = self.acquire_shared_buffer(address.len())?;
self.request(|local| local.connect_local_socket(handle, lease.sequence(), &address, user))
}

fn accept_local_socket(
&self,
handle: ObjectHandle,
flags: FileOpenFlags,
) -> core::result::Result<
core::result::Result<ObjectHandle, LocalSocketError>,
BrokerControlError,
> {
self.request(|local| local.accept_local_socket(handle, flags))
}

fn send_local_socket(
&self,
handle: ObjectHandle,
address: Option<&LocalSocketAddress>,
data: &[u8],
user: FileUser,
) -> core::result::Result<core::result::Result<usize, LocalSocketError>, BrokerControlError>
{
let mut staged = match address {
Some(address) => match address.encode() {
Some(address) => address,
None => return Ok(Err(LocalSocketError::InvalidArgument)),
},
None => Vec::new(),
};
let address_length = staged.len();
if data.len() > MAX_LOCAL_SOCKET_TRANSFER_SIZE as usize - address_length {
return Ok(Err(LocalSocketError::MessageTooLarge));
}
staged
.try_reserve_exact(data.len())
.map_err(|_| BrokerControlError::Broker(ErrorCode::OutOfMemory))?;
staged.extend_from_slice(data);
let lease = self.acquire_shared_buffer(staged.len())?;
self.request(|local| {
local.send_local_socket(handle, lease.sequence(), &staged, address_length, user)
})
}

fn receive_local_socket(
&self,
handle: ObjectHandle,
data: &mut [u8],
peek: bool,
nonblocking: bool,
) -> core::result::Result<
core::result::Result<(ReceiveLocalSocketResponse, LocalSocketName), LocalSocketError>,
BrokerControlError,
> {
let capacity = data.len().min(LOCAL_SOCKET_BUFFER_SIZE as usize);
let data = &mut data[..capacity];
let lease =
self.acquire_shared_buffer(capacity + MAX_ENCODED_LOCAL_SOCKET_NAME_SIZE as usize)?;
self.request(|local| {
local.receive_local_socket(handle, lease.sequence(), data, peek, nonblocking)
})
}

fn shutdown_local_socket(
&self,
handle: ObjectHandle,
mode: ShutdownMode,
) -> core::result::Result<core::result::Result<(), LocalSocketError>, BrokerControlError> {
self.request(|local| local.shutdown_local_socket(handle, mode))
}

fn local_socket_name(
&self,
handle: ObjectHandle,
peer: bool,
) -> core::result::Result<
core::result::Result<LocalSocketName, LocalSocketError>,
BrokerControlError,
> {
let lease = self.acquire_shared_buffer(MAX_ENCODED_LOCAL_SOCKET_NAME_SIZE as usize)?;
self.request(|local| local.local_socket_name(handle, peer, lease.sequence()))
}

fn set_local_socket_option(
&self,
handle: ObjectHandle,
option: LocalSocketOption,
) -> core::result::Result<(), BrokerControlError> {
self.request(|local| local.set_local_socket_option(handle, option))
}

fn local_socket_options(
&self,
handle: ObjectHandle,
) -> core::result::Result<GetLocalSocketOptionsResponse, BrokerControlError> {
self.request(|local| local.local_socket_options(handle))
}

fn close_object(&self, handle: ObjectHandle) -> core::result::Result<(), BrokerControlError> {
self.request(|local| local.close_object(handle))
}
Expand Down
1 change: 1 addition & 0 deletions litebox/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ extern crate alloc;
pub mod event;
pub mod fd;
pub mod fs;
pub mod local_sockets;
pub mod mm;
pub mod net;
pub mod path;
Expand Down
Loading
Loading