pub(crate) mod handle;
use handle::{
GlobalLockArrayHandle, GlobalLockCollectiveMutLocalDataHandle, GlobalLockLocalDataHandle,
GlobalLockMutLocalDataHandle, GlobalLockReadHandle, GlobalLockWriteHandle,
};
mod collective;
mod iteration;
pub(crate) mod operations;
mod rdma;
use crate::array::private::ArrayExecAm;
use crate::array::r#unsafe::__UnsafeByteArray;
use crate::barrier::BarrierHandle;
use crate::darc::global_rw_darc::{
GlobalRwDarc, GlobalRwDarcCollectiveWriteGuard, GlobalRwDarcReadGuard, GlobalRwDarcWriteGuard,
};
use crate::darc::DarcMode;
use crate::lamellar_request::LamellarRequest;
use crate::lamellar_team::{IntoLamellarTeam, LamellarTeamRT};
use crate::memregion::Dist;
use crate::scheduler::LamellarTask;
use crate::warnings::RuntimeWarning;
use crate::Remote;
use crate::{array::*, Darc};
use pin_project::pin_project;
use std::ops::{Deref, DerefMut};
use std::task::{Context, Poll};
#[derive(crate::Deserialize, crate::Serialize, Clone, Debug)]
#[serde(bound = "T: Dist")]
pub struct GlobalLockArray<T: Remote> {
pub(crate) lock: GlobalRwDarc<()>,
pub(crate) array: UnsafeArray<T>,
}
impl<T: Remote> crate::active_messaging::DarcSerde for GlobalLockArray<T> {
fn ser(&self, num_pes: usize, darcs: &mut Vec<RemotePtr>) {
self.lock.ser(num_pes, darcs);
self.array.ser(num_pes, darcs);
}
}
#[lamellar_impl::AmDataRT(Clone, Debug)]
pub struct __GlobalLockByteArray {
lock: GlobalRwDarc<()>,
pub(crate) array: __UnsafeByteArray,
}
impl __GlobalLockByteArray {}
#[derive(Debug)]
pub struct GlobalLockMutLocalData<T: Dist> {
pub(crate) array: GlobalLockArray<T>,
start_index: usize,
end_index: usize,
lock_guard: GlobalRwDarcWriteGuard<()>,
}
impl<T: Dist> Deref for GlobalLockMutLocalData<T> {
type Target = [T];
fn deref(&self) -> &Self::Target {
unsafe { &self.array.array.local_as_mut_slice()[self.start_index..self.end_index] }
}
}
impl<T: Dist> DerefMut for GlobalLockMutLocalData<T> {
fn deref_mut(&mut self) -> &mut Self::Target {
unsafe { &mut self.array.array.local_as_mut_slice()[self.start_index..self.end_index] }
}
}
#[derive(Debug)]
pub struct GlobalLockCollectiveMutLocalData<T: Dist> {
pub(crate) array: GlobalLockArray<T>,
start_index: usize,
end_index: usize,
_lock_guard: GlobalRwDarcCollectiveWriteGuard<()>,
}
impl<T: Dist> Deref for GlobalLockCollectiveMutLocalData<T> {
type Target = [T];
fn deref(&self) -> &Self::Target {
unsafe { &self.array.array.local_as_mut_slice()[self.start_index..self.end_index] }
}
}
impl<T: Dist> DerefMut for GlobalLockCollectiveMutLocalData<T> {
fn deref_mut(&mut self) -> &mut Self::Target {
unsafe { &mut self.array.array.local_as_mut_slice()[self.start_index..self.end_index] }
}
}
pub struct GlobalLockLocalData<T: Dist> {
pub(crate) array: GlobalLockArray<T>,
start_index: usize,
end_index: usize,
lock_guard: GlobalRwDarcReadGuard<()>,
}
impl<T: Dist + std::fmt::Debug> std::fmt::Debug for GlobalLockLocalData<T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{:?}", self.deref())
}
}
impl<T: Dist> Clone for GlobalLockLocalData<T> {
fn clone(&self) -> Self {
GlobalLockLocalData {
array: self.array.clone(),
start_index: self.start_index,
end_index: self.end_index,
lock_guard: self.lock_guard.clone(),
}
}
}
impl<T: Dist> Deref for GlobalLockLocalData<T> {
type Target = [T];
fn deref(&self) -> &Self::Target {
unsafe { &self.array.array.local_as_slice()[self.start_index..self.end_index] }
}
}
impl<T: Dist> GlobalLockLocalData<T> {
pub fn into_sub_data(self, start: usize, end: usize) -> GlobalLockLocalData<T> {
GlobalLockLocalData {
array: self.array.clone(),
start_index: start,
end_index: end,
lock_guard: self.lock_guard,
}
}
}
impl<T: Dist + serde::Serialize> serde::Serialize for GlobalLockLocalData<T> {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
unsafe { &self.array.array.local_as_mut_slice()[self.start_index..self.end_index] }
.serialize(serializer)
}
}
pub struct GlobalLockLocalDataIter<'a, T: Dist> {
data: &'a [T],
index: usize,
}
impl<'a, T: Dist> Iterator for GlobalLockLocalDataIter<'a, T> {
type Item = &'a T;
fn next(&mut self) -> Option<Self::Item> {
if self.index < self.data.len() {
self.index += 1;
Some(&self.data[self.index - 1])
} else {
None
}
}
}
impl<'a, T: Dist> IntoIterator for &'a GlobalLockLocalData<T> {
type Item = &'a T;
type IntoIter = GlobalLockLocalDataIter<'a, T>;
fn into_iter(self) -> Self::IntoIter {
GlobalLockLocalDataIter {
data: unsafe {
&self.array.array.local_as_mut_slice()[self.start_index..self.end_index]
},
index: 0,
}
}
}
#[derive(Clone)]
pub struct GlobalLockReadGuard<T: Dist> {
pub(crate) array: GlobalLockArray<T>,
lock_guard: GlobalRwDarcReadGuard<()>,
}
impl<T: Dist> GlobalLockReadGuard<T> {
pub fn local_data(&self) -> GlobalLockLocalData<T> {
GlobalLockLocalData {
array: self.array.clone(),
start_index: 0,
end_index: self.array.num_elems_local(),
lock_guard: self.lock_guard.clone(),
}
}
}
pub struct GlobalLockWriteGuard<T: Dist> {
pub(crate) array: GlobalLockArray<T>,
lock_guard: GlobalRwDarcWriteGuard<()>,
}
impl<T: Dist> From<GlobalLockMutLocalData<T>> for GlobalLockWriteGuard<T> {
fn from(data: GlobalLockMutLocalData<T>) -> Self {
GlobalLockWriteGuard {
array: data.array,
lock_guard: data.lock_guard,
}
}
}
impl<T: Dist> GlobalLockWriteGuard<T> {
pub fn local_data(self) -> GlobalLockMutLocalData<T> {
GlobalLockMutLocalData {
array: self.array.clone(),
start_index: 0,
end_index: self.array.num_elems_local(),
lock_guard: self.lock_guard,
}
}
}
impl<T: Dist + ArrayOps + std::default::Default> GlobalLockArray<T> {
#[doc(alias = "Collective")]
pub fn new<U: Clone + Into<IntoLamellarTeam>>(
team: U,
array_size: usize,
distribution: Distribution,
) -> GlobalLockArrayHandle<T> {
let team = team.into().team.clone();
GlobalLockArrayHandle {
team: team.clone(),
launched: false,
creation_future: Box::pin(async move {
let lock_task = GlobalRwDarc::new(team.clone(), ()).spawn();
GlobalLockArray {
lock: lock_task.await.expect("pe exists in team"),
array: UnsafeArray::async_new(
team.clone(),
array_size,
distribution,
DarcMode::GlobalLockArray,
)
.await,
}
}),
}
}
}
impl<T: Dist> GlobalLockArray<T> {
#[doc(alias("One-sided", "onesided"))]
pub fn use_distribution(self, distribution: Distribution) -> Self {
GlobalLockArray {
lock: self.lock.clone(),
array: self.array.use_distribution(distribution),
}
}
#[doc(alias("One-sided", "onesided"))]
pub fn read_lock(&self) -> GlobalLockReadHandle<T> {
GlobalLockReadHandle::new(self.clone())
}
#[doc(alias("One-sided", "onesided"))]
pub fn write_lock(&self) -> GlobalLockWriteHandle<T> {
GlobalLockWriteHandle::new(self.clone())
}
#[doc(alias("One-sided", "onesided"))]
pub fn read_local_data(&self) -> GlobalLockLocalDataHandle<T> {
GlobalLockLocalDataHandle {
array: self.clone(),
start_index: 0,
end_index: self.array.num_elems_local(),
lock_handle: self.lock.read(),
}
}
#[doc(alias("One-sided", "onesided"))]
pub fn write_local_data(&self) -> GlobalLockMutLocalDataHandle<T> {
GlobalLockMutLocalDataHandle {
array: self.clone(),
start_index: 0,
end_index: self.array.num_elems_local(),
lock_handle: self.lock.write(),
}
}
#[doc(alias("Collective"))]
pub fn collective_write_local_data(&self) -> GlobalLockCollectiveMutLocalDataHandle<T> {
GlobalLockCollectiveMutLocalDataHandle {
array: self.clone(),
start_index: 0,
end_index: self.array.num_elems_local(),
lock_handle: self.lock.collective_write(),
}
}
#[doc(hidden)]
pub unsafe fn __local_as_slice(&self) -> &[T] {
self.array.local_as_mut_slice()
}
#[doc(alias = "Collective")]
pub fn into_unsafe(self) -> IntoUnsafeArrayHandle<T> {
IntoUnsafeArrayHandle {
team: self.array.inner.data.team.clone(),
launched: false,
outstanding_future: Box::pin(self.async_into()),
}
}
#[doc(alias = "Collective")]
pub fn into_read_only(self) -> IntoReadOnlyArrayHandle<T> {
self.array.into_read_only()
}
#[doc(alias = "Collective")]
pub fn into_local_lock(self) -> IntoLocalLockArrayHandle<T> {
self.array.into_local_lock()
}
}
impl<T: Dist + 'static> GlobalLockArray<T> {
#[doc(alias = "Collective")]
pub fn into_atomic(self) -> IntoAtomicArrayHandle<T> {
self.array.into_atomic()
}
}
impl<T: Dist + ArrayOps + Default> AsyncTeamFrom<(Vec<T>, Distribution)> for GlobalLockArray<T> {
async fn team_from(input: (Vec<T>, Distribution), team: &Arc<LamellarTeam>) -> Self {
let array: UnsafeArray<T> = AsyncTeamInto::team_into(input, team).await;
array.async_into().await
}
}
#[async_trait]
impl<T: Dist> AsyncFrom<UnsafeArray<T>> for GlobalLockArray<T> {
async fn async_from(array: UnsafeArray<T>) -> Self {
array.await_on_outstanding(DarcMode::GlobalLockArray).await;
let lock = GlobalRwDarc::new(array.team_rt(), ())
.await
.expect("PE in team");
GlobalLockArray { lock, array }
}
}
impl<T: Dist> From<GlobalLockArray<T>> for __GlobalLockByteArray {
fn from(array: GlobalLockArray<T>) -> Self {
__GlobalLockByteArray {
lock: array.lock.clone(),
array: array.array.into(),
}
}
}
impl<T: Dist> From<GlobalLockArray<T>> for LamellarByteArray {
fn from(array: GlobalLockArray<T>) -> Self {
LamellarByteArray::GlobalLockArray(__GlobalLockByteArray {
lock: array.lock.clone(),
array: array.array.into(),
})
}
}
impl<T: Dist> From<LamellarByteArray> for GlobalLockArray<T> {
fn from(array: LamellarByteArray) -> Self {
if let LamellarByteArray::GlobalLockArray(array) = array {
array.into()
} else {
panic!("Expected LamellarByteArray::GlobalLockArray")
}
}
}
impl<T: Dist> From<__GlobalLockByteArray> for GlobalLockArray<T> {
fn from(array: __GlobalLockByteArray) -> Self {
GlobalLockArray {
lock: array.lock.clone(),
array: array.array.into(),
}
}
}
impl<T: Dist> From<&__GlobalLockByteArray> for GlobalLockArray<T> {
fn from(array: &__GlobalLockByteArray) -> Self {
array.clone().into()
}
}
impl<T: Dist> From<&mut __GlobalLockByteArray> for GlobalLockArray<T> {
fn from(array: &mut __GlobalLockByteArray) -> Self {
array.clone().into()
}
}
impl<T: Dist> private::ArrayExecAm<T> for GlobalLockArray<T> {
fn team_rt(&self) -> Darc<LamellarTeamRT> {
self.array.team_rt()
}
fn team_counters(&self) -> Arc<AMCounters> {
self.array.team_counters()
}
}
impl<T: Dist> private::LamellarArrayPrivate<T> for GlobalLockArray<T> {
fn inner_array(&self) -> &UnsafeArray<T> {
&self.array
}
fn local_as_ptr(&self) -> *const T {
self.array.local_as_mut_ptr()
}
fn local_as_mut_ptr(&self) -> *mut T {
self.array.local_as_mut_ptr()
}
fn pe_for_dist_index(&self, index: usize) -> Option<usize> {
self.array.pe_for_dist_index(index)
}
fn pe_offset_for_dist_index(&self, pe: usize, index: usize) -> Option<usize> {
self.array.pe_offset_for_dist_index(pe, index)
}
unsafe fn into_inner(self) -> UnsafeArray<T> {
self.array
}
fn as_lamellar_byte_array(&self) -> LamellarByteArray {
self.clone().into()
}
}
impl<T: Dist> ActiveMessaging for GlobalLockArray<T> {
type SinglePeAmHandle<R: AmDist> = AmHandle<R>;
type MultiAmHandle<R: AmDist> = MultiAmHandle<R>;
type LocalAmHandle<L> = LocalAmHandle<L>;
fn exec_am_all<F>(&self, am: F) -> Self::MultiAmHandle<F::Output>
where
F: RemoteActiveMessage + LamellarAM + Serde + AmDist,
{
self.array.exec_am_all_tg(am)
}
fn exec_am_pe<F>(&self, pe: usize, am: F) -> Self::SinglePeAmHandle<F::Output>
where
F: RemoteActiveMessage + LamellarAM + Serde + AmDist,
{
self.array.exec_am_pe_tg(pe, am)
}
fn exec_am_local<F>(&self, am: F) -> Self::LocalAmHandle<F::Output>
where
F: LamellarActiveMessage + LocalAM + 'static,
{
self.array.exec_am_local_tg(am)
}
fn wait_all(&self) {
self.array.wait_all()
}
fn await_all(&self) -> impl Future<Output = ()> + Send {
self.array.await_all()
}
fn barrier(&self) {
self.array.barrier()
}
fn async_barrier(&self) -> BarrierHandle {
self.array.async_barrier()
}
fn spawn<F>(&self, f: F) -> LamellarTask<F::Output>
where
F: Future + Send + 'static,
F::Output: Send,
{
self.array.spawn(f)
}
fn block_on<F: Future>(&self, f: F) -> F::Output {
self.array.block_on(f)
}
fn block_on_all<I>(&self, iter: I) -> Vec<<<I as IntoIterator>::Item as Future>::Output>
where
I: IntoIterator,
<I as IntoIterator>::Item: Future + Send + 'static,
<<I as IntoIterator>::Item as Future>::Output: Send,
{
self.array.block_on_all(iter)
}
}
impl<T: Dist> LamellarArray<T> for GlobalLockArray<T> {
fn len(&self) -> usize {
self.array.len()
}
fn num_elems_local(&self) -> usize {
self.array.num_elems_local()
}
fn pe_and_offset_for_global_index(&self, index: usize) -> Option<(usize, usize)> {
self.array.pe_and_offset_for_global_index(index)
}
fn first_global_index_for_pe(&self, pe: usize) -> Option<usize> {
self.array.first_global_index_for_pe(pe)
}
fn last_global_index_for_pe(&self, pe: usize) -> Option<usize> {
self.array.last_global_index_for_pe(pe)
}
}
impl<T: Dist> LamellarEnv for GlobalLockArray<T> {
fn my_pe(&self) -> usize {
LamellarEnv::my_pe(&self.array)
}
fn num_pes(&self) -> usize {
LamellarEnv::num_pes(&self.array)
}
fn num_threads_per_pe(&self) -> usize {
self.array.team_rt().num_threads()
}
fn world(&self) -> Arc<LamellarTeam> {
self.array.team_rt().world()
}
fn team(&self) -> Arc<LamellarTeam> {
self.array.team_rt().team()
}
}
impl<T: Dist> LamellarWrite for GlobalLockArray<T> {}
impl<T: Dist> LamellarRead for GlobalLockArray<T> {}
impl<T: Dist> SubArray<T> for GlobalLockArray<T> {
type Array = GlobalLockArray<T>;
fn sub_array<R: std::ops::RangeBounds<usize>>(&self, range: R) -> Self::Array {
GlobalLockArray {
lock: self.lock.clone(),
array: self.array.sub_array(range),
}
}
fn global_index(&self, sub_index: usize) -> usize {
self.array.global_index(sub_index)
}
}
impl<T: Dist + std::fmt::Debug> GlobalLockArray<T> {
#[doc(alias = "Collective")]
pub fn print(&self) {
self.barrier();
let _guard = self.read_local_data().block();
self.array.print();
}
}
impl<T: Dist + std::fmt::Debug> ArrayPrint<T> for GlobalLockArray<T> {
fn print(&self) {
self.barrier();
let _guard = self.read_local_data().block();
self.array.print()
}
}
#[pin_project]
pub struct GlobalLockArrayReduceHandle<T: Dist + AmDist> {
req: crate::array::ArrayReduceHandle<T>,
lock_guard: GlobalLockReadGuard<T>,
}
impl<T: Dist + AmDist> GlobalLockArrayReduceHandle<T> {
#[must_use = "this function returns a future used to poll for completion and retrieve the result. Call '.await' on the future otherwise, if it is ignored (via ' let _ = *.spawn()') or dropped the only way to ensure completion is calling 'wait_all()' on the world or array. Alternatively it may be acceptable to call '.block()' instead of 'spawn()'"]
pub fn spawn(mut self) -> LamellarTask<Option<T>> {
self.req.launch();
self.lock_guard.array.clone().spawn(self)
}
pub fn block(self) -> Option<T> {
RuntimeWarning::BlockingCall(
"GlobalLockArrayReduceHandle::block",
"<handle>.spawn() or <handle>.await",
)
.print();
self.lock_guard.array.clone().block_on(self)
}
}
impl<T: Dist + AmDist> Future for GlobalLockArrayReduceHandle<T> {
type Output = Option<T>;
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
let this = self.project();
match this.req.ready_or_set_waker(cx.waker()) {
true => Poll::Ready(this.req.val()),
false => Poll::Pending,
}
}
}
impl<T: Dist + AmDist + 'static> GlobalLockReadGuard<T> {
#[doc(alias("One-sided", "onesided"))]
#[must_use = "this function is lazy and does nothing unless awaited. Either await the returned future, or call 'spawn()' or 'block()' on it "]
pub fn registered_reduce(self, op: &str) -> GlobalLockArrayReduceHandle<T> {
GlobalLockArrayReduceHandle {
req: self
.array
.array
.reduce_data_user(op, self.array.clone().into()),
lock_guard: self,
}
}
}
impl<T: Dist + AmDist + ElementArithmeticOps + 'static> GlobalLockReadGuard<T> {
#[doc(alias("One-sided", "onesided"))]
#[must_use = "this function is lazy and does nothing unless awaited. Either await the returned future, or call 'spawn()' or 'block()' on it "]
pub fn sum(self) -> GlobalLockArrayReduceHandle<T> {
let req = match ScalarType::get_type::<T>() {
Some((scalar_type, _)) => {
self.array
.array
.reduce_data(Arc::new(ScalarBuiltinReductionAm::new(
self.array.clone().into(),
scalar_type,
BuiltinOp::Sum,
)))
}
None => self
.array
.array
.reduce_data_user("sum", self.array.clone().into()),
};
GlobalLockArrayReduceHandle {
req,
lock_guard: self,
}
}
#[doc(alias("One-sided", "onesided"))]
#[must_use = "this function is lazy and does nothing unless awaited. Either await the returned future, or call 'spawn()' or 'block()' on it "]
pub fn prod(self) -> GlobalLockArrayReduceHandle<T> {
let req = match ScalarType::get_type::<T>() {
Some((scalar_type, _)) => {
self.array
.array
.reduce_data(Arc::new(ScalarBuiltinReductionAm::new(
self.array.clone().into(),
scalar_type,
BuiltinOp::Prod,
)))
}
None => self
.array
.array
.reduce_data_user("prod", self.array.clone().into()),
};
GlobalLockArrayReduceHandle {
req,
lock_guard: self,
}
}
}
impl<T: Dist + AmDist + ElementComparePartialEqOps + 'static> GlobalLockReadGuard<T> {
#[doc(alias("One-sided", "onesided"))]
#[must_use = "this function is lazy and does nothing unless awaited. Either await the returned future, or call 'spawn()' or 'block()' on it "]
pub fn max(self) -> GlobalLockArrayReduceHandle<T> {
let req = match ScalarType::get_type::<T>() {
Some((scalar_type, _)) => {
self.array
.array
.reduce_data(Arc::new(ScalarBuiltinReductionAm::new(
self.array.clone().into(),
scalar_type,
BuiltinOp::Max,
)))
}
None => self
.array
.array
.reduce_data_user("max", self.array.clone().into()),
};
GlobalLockArrayReduceHandle {
req,
lock_guard: self,
}
}
#[doc(alias("One-sided", "onesided"))]
#[must_use = "this function is lazy and does nothing unless awaited. Either await the returned future, or call 'spawn()' or 'block()' on it "]
pub fn min(self) -> GlobalLockArrayReduceHandle<T> {
let req = match ScalarType::get_type::<T>() {
Some((scalar_type, _)) => {
self.array
.array
.reduce_data(Arc::new(ScalarBuiltinReductionAm::new(
self.array.clone().into(),
scalar_type,
BuiltinOp::Min,
)))
}
None => self
.array
.array
.reduce_data_user("min", self.array.clone().into()),
};
GlobalLockArrayReduceHandle {
req,
lock_guard: self,
}
}
}
impl<T: Dist + AmDist + ElementBitWiseOps + 'static> GlobalLockReadGuard<T> {
#[doc(alias("One-sided", "onesided"))]
#[must_use = "this function is lazy and does nothing unless awaited. Either await the returned future, or call 'spawn()' or 'block()' on it "]
pub fn and(self) -> GlobalLockArrayReduceHandle<T> {
let req = match ScalarType::get_type::<T>() {
Some((scalar_type, _)) => {
self.array
.array
.reduce_data(Arc::new(ScalarBuiltinReductionAm::new(
self.array.clone().into(),
scalar_type,
BuiltinOp::And,
)))
}
None => self
.array
.array
.reduce_data_user("and", self.array.clone().into()),
};
GlobalLockArrayReduceHandle {
req,
lock_guard: self,
}
}
#[doc(alias("One-sided", "onesided"))]
#[must_use = "this function is lazy and does nothing unless awaited. Either await the returned future, or call 'spawn()' or 'block()' on it "]
pub fn or(self) -> GlobalLockArrayReduceHandle<T> {
let req = match ScalarType::get_type::<T>() {
Some((scalar_type, _)) => {
self.array
.array
.reduce_data(Arc::new(ScalarBuiltinReductionAm::new(
self.array.clone().into(),
scalar_type,
BuiltinOp::Or,
)))
}
None => self
.array
.array
.reduce_data_user("or", self.array.clone().into()),
};
GlobalLockArrayReduceHandle {
req,
lock_guard: self,
}
}
#[doc(alias("One-sided", "onesided"))]
#[must_use = "this function is lazy and does nothing unless awaited. Either await the returned future, or call 'spawn()' or 'block()' on it "]
pub fn xor(self) -> GlobalLockArrayReduceHandle<T> {
let req = match ScalarType::get_type::<T>() {
Some((scalar_type, _)) => {
self.array
.array
.reduce_data(Arc::new(ScalarBuiltinReductionAm::new(
self.array.clone().into(),
scalar_type,
BuiltinOp::Xor,
)))
}
None => self
.array
.array
.reduce_data_user("xor", self.array.clone().into()),
};
GlobalLockArrayReduceHandle {
req,
lock_guard: self,
}
}
}