ruda-runtime 0.1.11

Ruda portable host runtime.
Documentation
use super::{
    ManagedMemoryHandle, MemoryPoolOptions, PoolType,
    memory_manage::DynamicPool,
    memory_pool::{MemoryPool, SlicedPool, calculate_padding},
};
use crate::runtime::{server::IoError, storage::{ComputeStorage, StorageHandle}};
use alloc::{format, vec, vec::Vec};
use ruda_core::{backtrace::BackTrace, ir::MemoryDeviceProperties};

const MIB: u64 = 1024 * 1024;
const SMALL_SLICE: u64 = 64 * 1024;
const ARENA_START: usize = 2;
const ARENA_SLOTS: usize = 64;

pub(super) struct AdaptiveState {
    alignment: u64,
    max_page_size: u64,
    min_page_size: u64,
    small_page_size: u64,
    current: usize,
    outdated: Vec<usize>,
}

struct Relocation {
    allocation: ManagedMemoryHandle,
    target: ManagedMemoryHandle,
    source: StorageHandle,
    destination: StorageHandle,
    cursor: u64,
}

impl AdaptiveState {
    pub(super) fn new(properties: &MemoryDeviceProperties) -> Self {
        let alignment = properties.alignment;
        assert!(alignment != 0, "Memory alignment must be nonzero");
        let max_page_size = properties.max_page_size / alignment * alignment;
        assert!(max_page_size != 0, "Device page limit must fit its alignment");
        let aligned = |size: u64| size.max(alignment).next_multiple_of(alignment).min(max_page_size);
        Self {
            alignment,
            max_page_size,
            min_page_size: aligned(2 * MIB),
            small_page_size: aligned(8 * MIB),
            current: ARENA_START,
            outdated: Vec::new(),
        }
    }

    pub(super) fn pool_options(&self) -> Vec<MemoryPoolOptions> {
        vec![
            MemoryPoolOptions { pool_type: PoolType::ExclusivePages { max_alloc_size: 0 }, dealloc_period: None },
            MemoryPoolOptions { pool_type: PoolType::SlicedPages {
                page_size: self.small_page_size, max_slice_size: SMALL_SLICE.min(self.small_page_size),
            }, dealloc_period: None },
            MemoryPoolOptions { pool_type: PoolType::SlicedPages {
                page_size: self.min_page_size, max_slice_size: self.min_page_size,
            }, dealloc_period: None },
        ]
    }

    fn pool<'a>(&self, pools: &'a [DynamicPool]) -> &'a SlicedPool {
        match &pools[self.current] { DynamicPool::Sliced(pool) => pool, _ => unreachable!() }
    }

    fn metadata_pool(&self, size: u64) -> Option<usize> {
        if size == 0 { Some(0) }
        else if size <= SMALL_SLICE.min(self.small_page_size) { Some(1) }
        else { None }
    }

    pub(super) fn pending(&self, pools: &[DynamicPool], size: u64) -> bool {
        if self.metadata_pool(size).is_some() { return false; }
        let pool = self.pool(pools);
        let Some(needed) = size.checked_add(calculate_padding(size, self.alignment)) else { return false; };
        needed <= self.max_page_size &&
            (needed > pool.page_size() || (!self.outdated.is_empty() && !pool.can_reserve(size)))
    }

    pub(super) fn reserve<Storage: ComputeStorage>(
        &mut self, pools: &mut Vec<DynamicPool>, storage: &mut Storage, size: u64,
    ) -> Result<ManagedMemoryHandle, IoError> {
        if !storage.supports_relocation() {
            return Err(IoError::UnsupportedIoOperation { backtrace: BackTrace::capture() });
        }
        let needed = size.checked_add(calculate_padding(size, self.alignment))
            .filter(|&needed| needed <= self.max_page_size)
            .ok_or_else(|| IoError::BufferTooBig { size, backtrace: BackTrace::capture() })?;
        if let Some(index) = self.metadata_pool(size) {
            if let Some(handle) = pools[index].try_reserve(size) { return Ok(handle); }
            return pools[index].alloc(storage, size);
        }

        if needed > self.pool(pools).page_size() {
            self.cleanup_outdated(pools, storage);
            let page_size = size.saturating_add(MIB).checked_next_multiple_of(MIB)
                .and_then(|size| size.checked_next_multiple_of(self.alignment))
                .unwrap_or(self.max_page_size).clamp(self.min_page_size, self.max_page_size);
            // Empty slots can be reused: no live descriptor can still name them.
            let vacant = (ARENA_START..pools.len()).find(|&index| {
                matches!(&pools[index], DynamicPool::Sliced(pool) if pool.is_empty())
            });
            let index = match vacant {
                Some(index) => index,
                None if pools.len() < ARENA_START + ARENA_SLOTS => pools.len(),
                None => return Err(IoError::Unknown {
                    description: "Adaptive page-size slots are all occupied by live allocations".into(),
                    backtrace: BackTrace::capture(),
                }),
            };
            let pool = DynamicPool::Sliced(SlicedPool::new(page_size, page_size, self.alignment, index as u8));
            if index == pools.len() { pools.push(pool); } else { pools[index] = pool; }
            if index != self.current { self.outdated.push(self.current); }
            self.current = index;
        }

        if let Some(handle) = pools[self.current].try_reserve(size) { return Ok(handle); }
        // Relocations only use room already held. An OOM never allocates a second
        // page just to attempt to recover the first reservation.
        let allocated = pools[self.current].alloc(storage, size);
        let reclaimed = self.relocate(pools, storage)?;
        match allocated {
            Ok(handle) => Ok(handle),
            Err(error) => match pools[self.current].try_reserve(size) {
                Some(handle) => Ok(handle),
                None if reclaimed => pools[self.current].alloc(storage, size),
                None => Err(error),
            },
        }
    }

    fn cleanup_outdated<Storage: ComputeStorage>(&mut self, pools: &mut [DynamicPool], storage: &mut Storage) {
        for &index in &self.outdated { pools[index].cleanup(storage, 0, true); }
        self.outdated.retain(|&index| !matches!(&pools[index], DynamicPool::Sliced(pool) if pool.is_empty()));
    }

    fn relocate<Storage: ComputeStorage>(&mut self, pools: &mut [DynamicPool], storage: &mut Storage) -> Result<bool, IoError> {
        if self.outdated.is_empty() { return Ok(false); }
        if !storage.supports_relocation() {
            return Err(IoError::UnsupportedIoOperation { backtrace: BackTrace::capture() });
        }
        let reserved_before: u64 = self.outdated.iter().map(|&index| pools[index].get_memory_usage().bytes_reserved).sum();
        let mut moves = Vec::new();
        for &index in &self.outdated {
            let allocations = match &pools[index] {
                DynamicPool::Sliced(pool) => pool.movable_allocations(), _ => unreachable!(),
            };
            for (allocation, source, cursor) in allocations {
                let Some(target) = pools[self.current].try_reserve(source.size()) else { continue; };
                let destination = pools[self.current].find_at(target.descriptor().location())?.storage.clone();
                moves.push(Relocation { allocation, target, source, destination, cursor });
            }
        }
        if moves.is_empty() {
            self.cleanup_outdated(pools, storage);
            return Ok(self.outdated.iter().map(|&index| pools[index].get_memory_usage().bytes_reserved).sum::<u64>() < reserved_before);
        }
        storage.relocation_barrier()?;
        let copied = moves.iter().try_for_each(|relocation| {
            storage.relocation_copy(&relocation.source, &relocation.destination)
        });
        // A failed enqueue can leave earlier copies active. No target may be
        // reused unless the completion wait proves those writes have finished.
        if let Err(error) = storage.relocation_complete() {
            core::mem::forget(moves);
            return Err(error);
        }
        copied?;
        for relocation in moves {
            let source_pool = relocation.allocation.descriptor().location().pool as usize;
            match &mut pools[source_pool] {
                DynamicPool::Sliced(pool) => pool.release_relocated(&relocation.allocation)?,
                _ => unreachable!(),
            }
            pools[self.current].bind(relocation.target, relocation.allocation, relocation.cursor)?;
        }
        self.cleanup_outdated(pools, storage);
        storage.flush();
        Ok(self.outdated.iter().map(|&index| pools[index].get_memory_usage().bytes_reserved).sum::<u64>() < reserved_before)
    }
}