Skip to content
Closed
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
1 change: 1 addition & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,7 @@ jobs:
AARCH64_CRATES: >-
-p litebox
-p litebox_broker_core
-p litebox_broker_transport
-p litebox_broker_userland
-p litebox_common_linux
-p litebox_egress_proxy
Expand Down
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion dev_tests/src/ratchet.rs
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,7 @@ fn ratchet_globals() -> Result<()> {
("litebox/", 7),
("litebox_broker_core/", 1),
("litebox_broker_transport_linux_userland/", 1),
("litebox_broker_userland/", 1),
("litebox_broker_userland/", 2),
("litebox_platform/", 2),
("litebox_platform_linux_kernel/", 5),
("litebox_platform_linux_userland/", 5),
Expand Down
68 changes: 42 additions & 26 deletions litebox/src/broker/shared_buffer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,13 +15,15 @@ use litebox_platform::sync::RawMutex as _;

use crate::sync::{Mutex, RawSyncPrimitivesProvider};

/// Leases shared-buffer slots, always the lowest free ones, so that a light
/// workload keeps reusing a few slots whose pages stay mapped and cache-warm
/// in both processes.
pub(super) struct SlotAllocator<Platform: RawSyncPrimitivesProvider> {
state: Mutex<Platform, AllocatorState<Platform>>,
}

struct AllocatorState<Platform: RawSyncPrimitivesProvider> {
allocated_slots: Vec<bool>,
next_slot: usize,
failed: bool,
waiters: VecDeque<Arc<SlotWaiter<Platform>>>,
}
Expand Down Expand Up @@ -49,7 +51,6 @@ impl<Platform: RawSyncPrimitivesProvider> SlotAllocator<Platform> {
Self {
state: Mutex::new(AllocatorState {
allocated_slots: vec![false; SHARED_BUFFER_SLOT_COUNT as usize],
next_slot: 0,
failed: false,
waiters: VecDeque::new(),
}),
Expand Down Expand Up @@ -223,40 +224,30 @@ impl<Platform: RawSyncPrimitivesProvider> Drop for SlotLease<'_, Platform> {

impl<Platform: RawSyncPrimitivesProvider> AllocatorState<Platform> {
fn allocate(&mut self, length: u32, slot_count: usize) -> Option<SharedBufferSequence> {
if self
let mut slot_indices = [SharedBufferSlotIndex::default(); MAX_SHARED_BUFFER_SEQUENCE_SLOTS];
let free_slots = self
.allocated_slots
.iter()
.filter(|allocated| !**allocated)
.count()
< slot_count
{
return None;
}

let mut slot_indices = [SharedBufferSlotIndex::default(); MAX_SHARED_BUFFER_SEQUENCE_SLOTS];
let mut next_slot = self.next_slot;
for stored_slot in &mut slot_indices[..slot_count] {
let slot_index = self
.next_free_slot(next_slot)
.expect("validated shared-buffer capacity must contain a free slot");
self.allocated_slots[slot_index] = true;
.enumerate()
.filter_map(|(slot_index, allocated)| (!*allocated).then_some(slot_index));
let mut found = 0;
for (stored_slot, slot_index) in slot_indices[..slot_count].iter_mut().zip(free_slots) {
*stored_slot = SharedBufferSlotIndex(
u32::try_from(slot_index).expect("shared-buffer slot index must fit in u32"),
);
next_slot = (slot_index + 1) % self.allocated_slots.len();
found += 1;
}
if found < slot_count {
return None;
}
for slot_index in &slot_indices[..slot_count] {
self.allocated_slots[slot_index.0 as usize] = true;
}
self.next_slot = next_slot;
Some(
SharedBufferSequence::new(&slot_indices[..slot_count], length)
.expect("allocated shared-buffer sequence must be valid"),
)
}

fn next_free_slot(&self, next_slot: usize) -> Option<usize> {
(0..self.allocated_slots.len())
.map(|offset| (next_slot + offset) % self.allocated_slots.len())
.find(|slot_index| !self.allocated_slots[*slot_index])
}
}

#[cfg(test)]
Expand Down Expand Up @@ -290,6 +281,32 @@ mod tests {
);
}

#[test]
fn leases_reuse_the_lowest_free_slots() {
let allocator = SlotAllocator::<MockPlatform>::new();
let first = allocator.acquire(2 * SHARED_BUFFER_SLOT_SIZE).unwrap();
let second = allocator.acquire(1).unwrap();
assert_eq!(
first.sequence().slot_indices(),
&[SharedBufferSlotIndex(0), SharedBufferSlotIndex(1)]
);
assert_eq!(
second.sequence().slot_indices(),
&[SharedBufferSlotIndex(2)]
);

drop(first);
let reused = allocator.acquire(3 * SHARED_BUFFER_SLOT_SIZE).unwrap();
assert_eq!(
reused.sequence().slot_indices(),
&[
SharedBufferSlotIndex(0),
SharedBufferSlotIndex(1),
SharedBufferSlotIndex(3)
]
);
}

#[test]
fn oversized_acquisitions_do_not_fail_the_allocator() {
let allocator = SlotAllocator::<MockPlatform>::new();
Expand All @@ -305,7 +322,6 @@ mod tests {
fn allocator_state_supports_slots_beyond_bitmap_widths() {
let mut state = AllocatorState::<MockPlatform> {
allocated_slots: alloc::vec![true; 65],
next_slot: 64,
failed: false,
waiters: VecDeque::new(),
};
Expand Down
10 changes: 10 additions & 0 deletions litebox/src/platform/page_mgmt.rs
Original file line number Diff line number Diff line change
Expand Up @@ -220,6 +220,16 @@ pub trait PageManagementProvider<const ALIGN: usize>: RawPointerProvider {
/// Note that the returned ranges should be `ALIGN`-aligned.
fn reserved_pages(&self) -> impl Iterator<Item = &Range<usize>>;

/// Hints that the caller is about to fill the allocated pages in `range` with data from their
/// start, possibly stopping short of the end (e.g., at the end of a file copied into them).
///
/// A platform may, for example, back these pages with larger host pages, which take fewer
/// faults to fill but can leave up to one larger page past the filled part resident.
///
/// The default implementation does nothing.
#[expect(unused_variables, reason = "default body")]
fn advise_fill(&self, range: Range<usize>) {}

/// Attempt to allocate pages with copy-on-write semantics backed by static data.
///
/// This method allows platforms that support it to create CoW mappings instead of performing
Expand Down
11 changes: 11 additions & 0 deletions litebox_broker_core/src/fs/backend.rs
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,17 @@ pub trait Backend: Send + Sync + Any {
/// Read directory entries at `dir`.
fn list_dir_at(&self, handle: DirHandle) -> Result<Vec<DirEntry>, ReadDirError>;

/// Look up the single entry `name` at `dir`, as [`Self::list_dir_at`] would report it.
///
/// The default implementation scans the full listing; backends that can find one entry more
/// cheaply should override it.
fn lookup_at(&self, dir: &DirHandle, name: &str) -> Result<Option<DirEntry>, ReadDirError> {
Ok(self
.list_dir_at(dir.clone())?
.into_iter()
.find(|entry| entry.name == name))
}

/// Read at `offset` into `buf`, returning the number of bytes read.
///
/// Backends do not have an internal notion of offsets; instead the resolver maintains offsets
Expand Down
41 changes: 30 additions & 11 deletions litebox_broker_core/src/fs/in_mem.rs
Original file line number Diff line number Diff line change
Expand Up @@ -381,20 +381,24 @@ impl<Platform: sync::RawSyncPrimitivesProvider> super::backend::Backend for InMe
.read()
.children
.iter()
.map(|(name, child)| {
let (file_type, node_info) = match child {
Node::File(file) => (FileType::RegularFile, file.read().node_info),
Node::Dir(dir) => (FileType::Directory, dir.read().node_info),
};
DirEntry {
name: name.clone(),
file_type,
ino_info: Some(node_info),
}
})
.map(|(name, child)| child.dir_entry(name))
.collect())
}

fn lookup_at(
&self,
dir: &super::backend::DirHandle,
name: &str,
) -> Result<Option<DirEntry>, ReadDirError> {
Ok(dir
.get_typed::<Self>()
.dir
.read()
.children
.get_key_value(name)
.map(|(name, child)| child.dir_entry(name)))
}

fn read(
&self,
h: &super::backend::FileHandle,
Expand Down Expand Up @@ -657,6 +661,21 @@ impl<Platform: sync::RawSyncPrimitivesProvider> Clone for Node<Platform> {
}
}

impl<Platform: sync::RawSyncPrimitivesProvider> Node<Platform> {
/// The directory entry for this node, named `name` in its parent.
fn dir_entry(&self, name: &str) -> DirEntry {
let (file_type, node_info) = match self {
Node::File(file) => (FileType::RegularFile, file.read().node_info),
Node::Dir(dir) => (FileType::Directory, dir.read().node_info),
};
DirEntry {
name: name.into(),
file_type,
ino_info: Some(node_info),
}
}
}

type DirNode<Platform> = Arc<sync::RwLock<Platform, DirData<Platform>>>;
struct DirData<Platform: sync::RawSyncPrimitivesProvider> {
perms: Permissions,
Expand Down
Loading
Loading