use std::marker::PhantomData;
use std::sync::{Arc, OnceLock};
static PROCESS_PARALLELISM: OnceLock<usize> = OnceLock::new();
pub(crate) fn sequential_fallback_permitted(
fan_out: &Result<(), moirai_core::error::ExecutorError>,
) -> bool {
match fan_out {
Ok(()) => false,
Err(moirai_core::error::ExecutorError::ShuttingDown) => true,
Err(error) => panic!(
"invariant: indexed fan-out failed after partial execution ({error}); \
retrying would duplicate caller side effects"
),
}
}
#[inline]
pub(crate) fn process_parallelism() -> usize {
*PROCESS_PARALLELISM.get_or_init(|| std::thread::available_parallelism().map_or(1, usize::from))
}
#[derive(Debug)]
pub(crate) struct SendPtr<T>(pub(crate) *mut T);
impl<T> Clone for SendPtr<T> {
#[inline]
fn clone(&self) -> Self {
*self
}
}
impl<T> Copy for SendPtr<T> {}
unsafe impl<T> Send for SendPtr<T> {}
unsafe impl<T> Sync for SendPtr<T> {}
impl<T> SendPtr<T> {
pub(crate) unsafe fn as_ptr(&self) -> *mut T {
self.0
}
}
pub fn tree_reduce<T, F>(mut items: Vec<T>, func: F) -> Option<T>
where
T: Send + Clone,
F: Fn(T, T) -> T + Send + Sync + Clone,
{
if items.is_empty() {
return None;
}
while items.len() > 1 {
let mut next = Vec::with_capacity(items.len().div_ceil(2));
for chunk in items.chunks(2) {
if chunk.len() == 2 {
next.push(func(chunk[0].clone(), chunk[1].clone()));
} else {
next.push(chunk[0].clone());
}
}
items = next;
}
items.into_iter().next()
}
pub fn process_in_batches<T, R, F>(items: Vec<T>, batch_size: usize, func: F) -> Vec<R>
where
T: Send + Clone,
R: Send,
F: Fn(&[T]) -> Vec<R> + Send + Sync,
{
items.chunks(batch_size).flat_map(func).collect()
}
pub struct BaseIterator<I, C> {
pub(crate) inner: I,
pub(crate) context: Arc<C>,
}
impl<I, C> BaseIterator<I, C> {
pub fn new(inner: I, context: C) -> Self {
Self {
inner,
context: Arc::new(context),
}
}
pub fn with_context(inner: I, context: Arc<C>) -> Self {
Self { inner, context }
}
#[must_use]
pub const fn inner(&self) -> &I {
&self.inner
}
#[must_use]
pub fn context(&self) -> &Arc<C> {
&self.context
}
#[must_use]
pub fn into_parts(self) -> (I, Arc<C>) {
(self.inner, self.context)
}
}
pub trait FromMoiraiIterator<T>: Send {
fn from_iter<I: IntoIterator<Item = T>>(iter: I) -> Self;
fn from_iter_with_hint<I: IntoIterator<Item = T>>(iter: I, size_hint: usize) -> Self
where
Self: Sized,
{
let _ = size_hint; Self::from_iter(iter)
}
}
impl<T: Send> FromMoiraiIterator<T> for Vec<T> {
fn from_iter<I: IntoIterator<Item = T>>(iter: I) -> Self {
iter.into_iter().collect()
}
fn from_iter_with_hint<I: IntoIterator<Item = T>>(iter: I, size_hint: usize) -> Self {
let mut vec = Vec::with_capacity(size_hint);
vec.extend(iter);
vec
}
}
pub struct MapAdapter<I, F, T, R> {
pub(crate) inner: I,
pub(crate) func: F,
pub(crate) _phantom: PhantomData<(T, R)>,
}
impl<I, F, T, R> MapAdapter<I, F, T, R> {
pub fn new(inner: I, func: F) -> Self {
Self {
inner,
func,
_phantom: PhantomData,
}
}
#[must_use]
pub const fn inner(&self) -> &I {
&self.inner
}
#[must_use]
pub const fn function(&self) -> &F {
&self.func
}
#[must_use]
pub fn into_parts(self) -> (I, F) {
(self.inner, self.func)
}
}
pub struct FilterAdapter<I, F, T> {
pub(crate) inner: I,
pub(crate) predicate: F,
pub(crate) _phantom: PhantomData<T>,
}
impl<I, F, T> FilterAdapter<I, F, T> {
pub fn new(inner: I, predicate: F) -> Self {
Self {
inner,
predicate,
_phantom: PhantomData,
}
}
#[must_use]
pub const fn inner(&self) -> &I {
&self.inner
}
#[must_use]
pub const fn predicate(&self) -> &F {
&self.predicate
}
#[must_use]
pub fn into_parts(self) -> (I, F) {
(self.inner, self.predicate)
}
}
pub struct BatchAdapter<I> {
pub(crate) inner: I,
pub(crate) size: usize,
}
impl<I> BatchAdapter<I> {
pub fn new(inner: I, size: usize) -> Self {
Self {
inner,
size: size.max(1),
}
}
#[must_use]
pub const fn inner(&self) -> &I {
&self.inner
}
#[must_use]
pub const fn size(&self) -> usize {
self.size
}
#[must_use]
pub fn into_parts(self) -> (I, usize) {
(self.inner, self.size)
}
}
#[derive(Debug, Clone)]
pub struct PerformanceMetrics {
pub total_items: usize,
pub execution_time_ns: u64,
pub memory_used_bytes: usize,
pub strategy_used: String,
}
impl PerformanceMetrics {
pub fn throughput_per_sec(&self) -> f64 {
if self.execution_time_ns == 0 {
0.0
} else {
(self.total_items as f64 * 1_000_000_000.0) / self.execution_time_ns as f64
}
}
}
#[cfg(test)]
#[path = "base/tests.rs"]
mod tests;