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
47 changes: 43 additions & 4 deletions litebox/src/broker/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,10 +22,11 @@ use litebox_broker_protocol::process::{
CreatedProcess, MAX_CHILD_MEMORY_WRITE_SIZE, MAX_CHILD_OBJECT_DUPLICATES,
MAX_PROCESS_BOOTSTRAP_SIZE, ProcessExitStatus, ProcessTermination,
};
use litebox_broker_protocol::process_group::ProcessGroupMembership;
use litebox_broker_protocol::random::MAX_RANDOM_TRANSFER_SIZE;
use litebox_broker_protocol::readiness::ReadinessFlags;
use litebox_broker_protocol::shared_buffer::SHARED_BUFFER_SLOT_SIZE;
use litebox_broker_protocol::signal::PendingSignal;
use litebox_broker_protocol::signal::{PendingSignal, SignalTarget};
use litebox_broker_protocol::socket::{
AcceptSocketResponse, MAX_SOCKET_TRANSFER_SIZE, MAX_UDP_DATAGRAM_SIZE,
ReceiveFlags as BrokerReceiveFlags, ReceiveFromFlags as BrokerReceiveFromFlags,
Expand Down Expand Up @@ -97,7 +98,7 @@ pub(crate) trait BrokerControl: Send + Sync {

fn send_signal(
&self,
process_id: litebox_broker_protocol::ProcessId,
target: SignalTarget,
signal: u32,
) -> core::result::Result<(), BrokerControlError>;

Expand All @@ -106,6 +107,22 @@ pub(crate) trait BrokerControl: Send + Sync {
handle: ObjectHandle,
) -> core::result::Result<PendingSignal, BrokerControlError>;

fn process_group(
&self,
process_id: litebox_broker_protocol::ProcessId,
) -> core::result::Result<ProcessGroupMembership, BrokerControlError>;

fn set_process_group(
&self,
process_id: litebox_broker_protocol::ProcessId,
process_group: litebox_broker_protocol::ProcessId,
) -> core::result::Result<(), BrokerControlError>;

fn create_session(
&self,
process_id: litebox_broker_protocol::ProcessId,
) -> core::result::Result<(), BrokerControlError>;

fn create_thread(&self) -> core::result::Result<ThreadId, BrokerControlError>;

fn exit_thread(&self, thread_id: ThreadId) -> core::result::Result<(), BrokerControlError>;
Expand Down Expand Up @@ -594,10 +611,10 @@ where

fn send_signal(
&self,
process_id: litebox_broker_protocol::ProcessId,
target: SignalTarget,
signal: u32,
) -> core::result::Result<(), BrokerControlError> {
self.request(|local| local.send_signal(process_id, signal))
self.request(|local| local.send_signal(target, signal))
}

fn take_signal(
Expand All @@ -607,6 +624,28 @@ where
self.request(|local| local.take_signal(handle))
}

fn process_group(
&self,
process_id: litebox_broker_protocol::ProcessId,
) -> core::result::Result<ProcessGroupMembership, BrokerControlError> {
self.request(|local| local.process_group(process_id))
}

fn set_process_group(
&self,
process_id: litebox_broker_protocol::ProcessId,
process_group: litebox_broker_protocol::ProcessId,
) -> core::result::Result<(), BrokerControlError> {
self.request(|local| local.set_process_group(process_id, process_group))
}

fn create_session(
&self,
process_id: litebox_broker_protocol::ProcessId,
) -> core::result::Result<(), BrokerControlError> {
self.request(|local| local.create_session(process_id))
}

fn create_thread(&self) -> core::result::Result<ThreadId, BrokerControlError> {
self.request(BrokerLocal::create_thread)
}
Expand Down
74 changes: 58 additions & 16 deletions litebox/src/process.rs
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
// Copyright (c) Microsoft Corporation.
// Licensed under the MIT license.

//! Broker-backed guest process creation, child termination, and signals
//! between processes.
//! Broker-backed guest process creation, child termination, process groups
//! and sessions, and signals between processes.

use alloc::sync::Arc;
use alloc::vec::Vec;
Expand All @@ -11,7 +11,8 @@ use litebox_broker_protocol::error::ErrorCode;
use litebox_broker_protocol::process::{
MAX_CHILD_MEMORY_WRITE_SIZE, ProcessExitStatus, ProcessIdentity, ProcessTermination,
};
use litebox_broker_protocol::signal::PendingSignal;
use litebox_broker_protocol::process_group::ProcessGroupMembership;
use litebox_broker_protocol::signal::{PendingSignal, SignalTarget};
use litebox_broker_protocol::{ObjectHandle, ProcessId};
use litebox_platform::time::TimeProvider;

Expand All @@ -33,8 +34,8 @@ pub enum ProcessError {
/// The broker association failed.
#[error("process service failed")]
ServiceFailed,
/// Process duplication is disabled by policy.
#[error("process duplication is denied")]
/// The operation is denied by policy or by process group rules.
#[error("process operation is denied")]
PolicyDenied,
/// Another pending child process already exists.
#[error("a child process is already pending")]
Expand Down Expand Up @@ -80,19 +81,52 @@ impl<Platform: RawSyncPrimitivesProvider + TimeProvider> LiteBox<Platform> {
Ok(broker.set_child_reaping(enabled)?)
}

/// Sends `signal` to process `process_id`, or only checks that the process
/// exists if `signal` is zero.
/// Sends `signal` to the processes `target` selects, or only checks that
/// one exists if `signal` is zero.
///
/// The target takes the signal through [`Signals`]. A signal already
/// pending for the target is not sent again.
pub fn send_signal(&self, process_id: ProcessId, signal: u32) -> Result<(), ProcessError> {
/// Each target takes the signal through [`Signals`]. A signal already
/// pending for a target is not sent again. Returns
/// [`ProcessError::NoSuchProcess`] if no process is targeted.
pub fn send_signal(&self, target: SignalTarget, signal: u32) -> Result<(), ProcessError> {
let broker = self.broker_control().ok_or(ProcessError::Unavailable)?;
match broker.send_signal(process_id, signal) {
Err(BrokerControlError::Broker(ErrorCode::UnknownObject)) => {
Err(ProcessError::NoSuchProcess)
}
result => Ok(result?),
}
broker.send_signal(target, signal).map_err(no_such_process)
}

/// Returns the process group and session of process `process_id`.
pub fn process_group(
&self,
process_id: ProcessId,
) -> Result<ProcessGroupMembership, ProcessError> {
let broker = self.broker_control().ok_or(ProcessError::Unavailable)?;
broker.process_group(process_id).map_err(no_such_process)
}

/// Moves process `process_id`, which is this process or one of its
/// children, into `process_group`, creating the group if it is
/// `process_id`.
///
/// Returns [`ProcessError::PolicyDenied`] if the process is in another
/// session than this process or leads a session, or if `process_group` is
/// neither `process_id` nor an existing group in this process's session.
pub fn set_process_group(
&self,
process_id: ProcessId,
process_group: ProcessId,
) -> Result<(), ProcessError> {
let broker = self.broker_control().ok_or(ProcessError::Unavailable)?;
broker
.set_process_group(process_id, process_group)
.map_err(no_such_process)
}

/// Makes process `process_id`, which is this process or its pending
/// child, the leader of a new session and of a new process group in it.
///
/// Returns [`ProcessError::PolicyDenied`] if a process group already has
/// the ID `process_id`.
pub fn create_session(&self, process_id: ProcessId) -> Result<(), ProcessError> {
let broker = self.broker_control().ok_or(ProcessError::Unavailable)?;
broker.create_session(process_id).map_err(no_such_process)
}

/// Opens the signals other processes send to this process, including
Expand Down Expand Up @@ -353,6 +387,14 @@ impl<Platform: RawSyncPrimitivesProvider + TimeProvider> IOPollable for Signals<
}
}

/// Converts an error from a request that targets processes by ID.
fn no_such_process(error: BrokerControlError) -> ProcessError {
match error {
BrokerControlError::Broker(ErrorCode::UnknownObject) => ProcessError::NoSuchProcess,
error => error.into(),
}
}

impl From<BrokerControlError> for ProcessError {
fn from(error: BrokerControlError) -> Self {
match error {
Expand Down
40 changes: 37 additions & 3 deletions litebox_broker_core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ mod object;
pub mod pipe;
mod policy;
mod process;
pub mod process_group;
pub mod random;
pub mod readiness;
pub mod signal;
Expand All @@ -41,7 +42,9 @@ pub mod test_support;
use alloc::sync::{Arc, Weak};
use core::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};

use alloc::vec::Vec;
use hashbrown::HashMap;
use litebox_broker_protocol::process_group::ProcessGroupMembership;
use litebox_broker_protocol::{ObjectHandle, ProcessId};
use spin::{Mutex, rwlock::RwLock};

Expand Down Expand Up @@ -234,6 +237,9 @@ pub struct BrokerCore {
pub(crate) limits: BrokerCoreLimits,
pub(crate) ids: Arc<Mutex<IdAllocator>>,
pub(crate) processes: Arc<RwLock<HashMap<ProcessId, Weak<BrokerProcess>>>>,
/// Serializes changes to process group and session membership, and
/// selecting the members of a group.
pub(crate) process_groups: Arc<Mutex<()>>,
/// Number of broker threads created and not normally retired.
pub(crate) active_thread_count: Arc<AtomicUsize>,
pub(crate) next_reference_handle: Arc<RwLock<u64>>,
Expand Down Expand Up @@ -297,6 +303,7 @@ impl BrokerCore {
limits,
ids: Arc::new(Mutex::new(ids)),
processes: Arc::new(RwLock::new(HashMap::new())),
process_groups: Arc::new(Mutex::new(())),
active_thread_count: Arc::new(AtomicUsize::new(0)),
next_reference_handle: Arc::new(RwLock::new(1)),
references: Arc::new(RwLock::new(HashMap::new())),
Expand All @@ -322,6 +329,25 @@ impl BrokerCore {
}
}

/// Returns the registered process `id`, or `UnknownObject` if none.
pub(crate) fn registered_process(&self, id: ProcessId) -> Result<Arc<BrokerProcess>> {
// The registry lock is released before the process can drop, since a
// final process drop removes itself from the registry.
self.processes
.read()
.get(&id)
.and_then(Weak::upgrade)
.ok_or(BrokerError::UnknownObject)
}

/// Returns every registered process.
pub(crate) fn registered_processes(&self) -> Vec<Arc<BrokerProcess>> {
// The registry lock is released before any process can drop, since a
// final process drop removes itself from the registry.
let processes = self.processes.read();
processes.values().filter_map(Weak::upgrade).collect()
}

/// Returns whether any broker process remains registered.
#[must_use]
pub fn has_processes(&self) -> bool {
Expand Down Expand Up @@ -371,7 +397,7 @@ impl BrokerCore {
caller_credential: CallerCredential,
parent_id: Option<ProcessId>,
) -> Result<Arc<BrokerProcess>> {
let allocate_process = |parent: Option<Weak<BrokerProcess>>| {
let allocate_process = |parent: Option<&Arc<BrokerProcess>>| {
let mut processes = self.processes.write();
if processes.len() >= self.limits.max_processes {
return Err(BrokerError::ResourceExhausted);
Expand All @@ -381,10 +407,18 @@ impl BrokerCore {
.map_err(|_| BrokerError::OutOfMemory)?;
let raw_id = self.ids.lock().allocate()?;
let id = ProcessId(raw_id);
let membership = parent.map_or(
ProcessGroupMembership {
process_group: id,
session: id,
},
|parent| parent.membership(),
);
let process = Arc::new(BrokerProcess::new(
self.clone(),
id,
parent,
parent.map(Arc::downgrade),
membership,
caller_credential,
));
assert!(
Expand All @@ -401,7 +435,7 @@ impl BrokerCore {
.get(&parent_id)
.and_then(Weak::upgrade)
.ok_or(BrokerError::UnknownObject)?;
return parent.with_live_owner(|| allocate_process(Some(Arc::downgrade(&parent))))?;
return parent.with_live_owner(|| allocate_process(Some(&parent)))?;
}
allocate_process(None)
}
Expand Down
32 changes: 31 additions & 1 deletion litebox_broker_core/src/process.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ use litebox_broker_protocol::fs::{FileOpenFlags, FileStatusFlags, SetStatusFlags
use litebox_broker_protocol::process::{
CreatedProcess, ProcessExitStatus, ProcessIdentity, ProcessTermination,
};
use litebox_broker_protocol::process_group::ProcessGroupMembership;
use litebox_broker_protocol::readiness::ReadinessFlags;
use litebox_broker_protocol::{ObjectHandle, ProcessId, ThreadId};
use spin::{Mutex, Once, rwlock::RwLock};
Expand Down Expand Up @@ -215,6 +216,11 @@ pub struct BrokerProcess {
initial_thread_id: Once<ThreadId>,
/// Creating parent process.
parent: Option<Weak<BrokerProcess>>,
/// Process group and session.
///
/// No other lock is taken while this one is held. Changes are serialized
/// by the core's process group lock.
pub(crate) membership: Mutex<ProcessGroupMembership>,
/// Whether this process's children are reaped when they terminate.
reap_children: AtomicBool,
state: Mutex<BrokerProcessState>,
Expand Down Expand Up @@ -344,13 +350,15 @@ impl BrokerProcess {
core: BrokerCore,
id: ProcessId,
parent: Option<Weak<BrokerProcess>>,
membership: ProcessGroupMembership,
caller_credential: CallerCredential,
) -> Self {
Self {
core,
id,
initial_thread_id: Once::new(),
parent,
membership: Mutex::new(membership),
reap_children: AtomicBool::new(false),
state: Mutex::new(BrokerProcessState {
status: ProcessStatus::Starting,
Expand Down Expand Up @@ -709,6 +717,28 @@ impl BrokerProcess {
self.duplicate_object_references_to(handles, child)
}

/// Returns this process's group and session.
pub(crate) fn membership(&self) -> ProcessGroupMembership {
*self.membership.lock()
}

/// Returns whether this process has a parent, even one that no longer
/// exists.
pub(crate) fn has_parent(&self) -> bool {
self.parent.is_some()
}

/// Returns the pending child selected by `child_process_id`.
pub(crate) fn pending_child_process(
&self,
child_process_id: ProcessId,
) -> Result<Arc<BrokerProcess>> {
let mut state = self.state.lock();
Ok(Arc::clone(
&self.pending_child(&mut state, child_process_id)?.process,
))
}

/// Returns this process's parent if it still exists.
///
/// Callers upgrade before taking this process's state lock and drop the
Expand Down Expand Up @@ -969,7 +999,7 @@ impl BrokerProcess {
self.is_active(state) && matches!(state.status, ProcessStatus::Starting)
}

fn is_child_of(&self, parent: &BrokerProcess) -> bool {
pub(crate) fn is_child_of(&self, parent: &BrokerProcess) -> bool {
// The weak reference keeps the parent's allocation, so its address
// cannot be reused by another process.
self.parent
Expand Down
Loading
Loading