use core::marker::PhantomData;
use futures_util::future::join_all;
use serde::{Deserialize, Deserializer};
use std::cmp::PartialEq;
use std::fmt;
use std::hash::Hash;
use std::ops::Deref;
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::Instant;
use tracing::trace;
use crate::{
active_messaging::{AMCounters, RemotePtr},
barrier::Barrier,
env_var::config,
lamellae::{
AllocationType, Backend, CommAlloc, CommAllocAddr, CommAllocRdma, CommInfo, CommMem,
CommProgress, CommSlice,
},
lamellar_team::{IntoLamellarTeam, LamellarTeamRT},
lamellar_world::LAMELLAES,
scheduler::LamellarTask,
warnings::RuntimeWarning,
IdError, LamellarEnv, LamellarTeam,
};
pub mod prelude;
pub(crate) mod local_rw_darc;
pub use local_rw_darc::LocalRwDarc;
pub(crate) mod global_rw_darc;
pub use global_rw_darc::GlobalRwDarc;
use self::handle::{DarcHandle, IntoGlobalRwDarcHandle, IntoLocalRwDarcHandle};
pub(crate) mod handle;
static DARC_ID: AtomicUsize = AtomicUsize::new(0);
#[repr(u64)]
#[derive(PartialEq, Copy, Clone, Debug)]
pub(crate) enum DarcMode {
Darc,
LocalRw,
GlobalRw,
UnsafeArray,
ReadOnlyArray,
GenericAtomicArray,
NativeAtomicArray,
NetworkAtomicArray,
LocalLockArray,
GlobalLockArray,
WorldTeam,
Dropping,
Dropped,
RestartDrop,
}
#[lamellar_prof::prof]
impl DarcMode {
fn drop_am_launched(&self) -> bool {
matches!(
self,
DarcMode::Dropping | DarcMode::Dropped | DarcMode::RestartDrop
)
}
}
#[lamellar_prof::prof]
impl Default for DarcMode {
fn default() -> Self {
DarcMode::Darc
}
}
#[lamellar_prof::prof]
impl From<u64> for DarcMode {
fn from(val: u64) -> Self {
match val {
x if x == DarcMode::Darc as u64 => DarcMode::Darc,
x if x == DarcMode::LocalRw as u64 => DarcMode::LocalRw,
x if x == DarcMode::GlobalRw as u64 => DarcMode::GlobalRw,
x if x == DarcMode::UnsafeArray as u64 => DarcMode::UnsafeArray,
x if x == DarcMode::ReadOnlyArray as u64 => DarcMode::ReadOnlyArray,
x if x == DarcMode::GenericAtomicArray as u64 => DarcMode::GenericAtomicArray,
x if x == DarcMode::NativeAtomicArray as u64 => DarcMode::NativeAtomicArray,
x if x == DarcMode::NetworkAtomicArray as u64 => DarcMode::NetworkAtomicArray,
x if x == DarcMode::LocalLockArray as u64 => DarcMode::LocalLockArray,
x if x == DarcMode::GlobalLockArray as u64 => DarcMode::GlobalLockArray,
x if x == DarcMode::Dropping as u64 => DarcMode::Dropping,
x if x == DarcMode::Dropped as u64 => DarcMode::Dropped,
x if x == DarcMode::RestartDrop as u64 => DarcMode::RestartDrop,
x if x == DarcMode::WorldTeam as u64 => DarcMode::WorldTeam,
_ => panic!("invalid darc mode value {}", val),
}
}
}
#[lamellar_impl::AmDataRT(Debug)]
struct FinishedAm {
cnt: usize,
src_pe: usize,
inner_addr: CommAllocAddr, }
#[lamellar_impl::rt_am]
impl LamellarAM for FinishedAm {
async fn exec() {
trace!(target: "drop", "in finished! {:?}", self);
let inner: &DarcInner<()> = unsafe { &*(self.inner_addr.as_ptr()) }; inner.dist_cnt.fetch_sub(self.cnt, Ordering::SeqCst);
}
}
#[doc(hidden)]
#[repr(C)]
pub struct DarcInner<T> {
team: *const DarcInner<LamellarTeamRT>,
item: *const T,
id: usize,
my_pe: usize, num_pes: usize, local_cnt: AtomicUsize, total_local_cnt: AtomicUsize,
weak_local_cnt: AtomicUsize, dist_cnt: AtomicUsize, total_dist_cnt: AtomicUsize,
ref_cnt_slice: CommSlice<usize>, total_ref_cnt_slice: CommSlice<usize>,
mode_slice: CommSlice<DarcMode>,
mode_ref_cnt_slice: CommSlice<usize>,
mode_barrier_slice: CommSlice<usize>,
barrier: *mut Barrier,
am_counters: *const AMCounters,
drop: Option<fn(&mut T) -> bool>,
valid: AtomicBool,
}
unsafe impl<T> Send for DarcInner<T> {} unsafe impl<T> Sync for DarcInner<T> {}
pub struct Darc<T: 'static> {
inner: DarcCommPtr<T>,
src_pe: usize,
id: usize,
}
unsafe impl<T: Sync + Send> Send for Darc<T> {}
unsafe impl<T: Sync + Send> Sync for Darc<T> {}
#[lamellar_prof::prof]
impl<T> LamellarEnv for Darc<T> {
fn my_pe(&self) -> usize {
self.inner().my_pe
}
fn num_pes(&self) -> usize {
self.inner().num_pes
}
fn num_threads_per_pe(&self) -> usize {
self.inner().darc_rt_team().rt_num_threads_per_pe()
}
fn world(&self) -> Arc<LamellarTeam> {
self.inner().darc_rt_team().user_world()
}
fn team(&self) -> Arc<LamellarTeam> {
self.inner().darc_rt_team().user_team()
}
}
#[lamellar_prof::prof]
impl<T: 'static> serde::Serialize for Darc<T> {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
__NetworkDarc::from(self).serialize(serializer)
}
}
#[lamellar_prof::prof]
impl<'de, T: 'static> Deserialize<'de> for Darc<T> {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let ndarc: __NetworkDarc = Deserialize::deserialize(deserializer)?;
Ok(ndarc.into())
}
}
#[derive(Debug)]
pub struct WeakDarc<T: 'static> {
inner: DarcCommPtr<T>,
src_pe: usize,
}
unsafe impl<T: Send> Send for WeakDarc<T> {}
unsafe impl<T: Sync> Sync for WeakDarc<T> {}
#[lamellar_prof::prof]
impl<T> WeakDarc<T> {
pub fn upgrade(&self) -> Option<Darc<T>> {
let inner = &*self.inner;
inner.local_cnt.fetch_add(1, Ordering::SeqCst);
let id = inner.total_local_cnt.fetch_add(1, Ordering::SeqCst);
if inner.valid.load(Ordering::SeqCst) {
Some(Darc {
inner: self.inner.clone(),
src_pe: self.src_pe,
id,
})
} else {
let cnt = inner.local_cnt.fetch_sub(1, Ordering::SeqCst);
if cnt == 0 {
panic!("darc dropped too many times");
}
None
}
}
}
#[lamellar_prof::prof]
impl<T> Drop for WeakDarc<T> {
fn drop(&mut self) {
trace!(target: "drop", "begin drop WeakDarc");
let inner = &*self.inner;
inner.weak_local_cnt.fetch_sub(1, Ordering::SeqCst);
trace!(target: "drop", "end drop WeakDarc");
}
}
#[lamellar_prof::prof]
impl<T> Clone for WeakDarc<T> {
fn clone(&self) -> Self {
let inner = &*self.inner;
inner.weak_local_cnt.fetch_add(1, Ordering::SeqCst);
WeakDarc {
inner: self.inner.clone(),
src_pe: self.src_pe,
}
}
}
#[lamellar_prof::prof]
impl<T> crate::active_messaging::DarcSerde for Darc<T> {
fn ser(&self, num_pes: usize, darcs: &mut Vec<RemotePtr>) {
trace!(target:"darc_clone", "darc ser {:?} ", self.inner());
self.serialize_update_cnts(num_pes);
trace!(target:"darc_clone", "darc ser {:?} ", self.inner());
darcs.push(RemotePtr::NetworkDarc(self.clone().into()));
}
}
#[lamellar_prof::prof]
impl<T: 'static> DarcInner<T> {
pub(crate) fn darc_rt_team(&self) -> Darc<LamellarTeamRT> {
unsafe { Darc::cloned_team_from_raw(self.team) }
}
pub(crate) fn rt_team(&self) -> &LamellarTeamRT {
unsafe { &*self.team }.item()
}
fn am_counters(&self) -> Arc<AMCounters> {
unsafe {
Arc::increment_strong_count(self.am_counters);
Arc::from_raw(self.am_counters)
}
}
fn inc_pe_ref_count(&self, pe: usize, amt: usize) -> usize {
trace!("inc_pe_ref_count pe: {} amt: {} {:?}", pe, amt, self);
let team_pe = pe;
let tot_ref_cnt = unsafe {
(&self.total_ref_cnt_slice[team_pe] as *const _ as *const AtomicUsize)
.as_ref()
.expect("invalid darc addr")
};
tot_ref_cnt.fetch_add(amt, Ordering::SeqCst);
let ref_cnt = unsafe {
(&self.ref_cnt_slice[team_pe] as *const _ as *const AtomicUsize)
.as_ref()
.expect("invalid darc addr")
};
ref_cnt.fetch_add(amt, Ordering::SeqCst)
}
fn update_item(&mut self, item: *const T) {
self.item = item;
}
fn set_dropping(&self) -> DarcMode {
let mode = unsafe {
(&self.mode_slice[self.my_pe] as *const _ as *const AtomicU64)
.as_ref()
.expect("invalid darc addr")
};
mode.swap(DarcMode::Dropping as u64, Ordering::SeqCst)
.into()
}
fn send_finished(&self) -> Vec<LamellarTask<()>> {
trace!(
"[{:?}] in send_finished {:?}",
std::thread::current().id(),
self
);
let ref_cnts = unsafe {
self.ref_cnt_slice
.as_casted_slice::<AtomicUsize>()
.expect("invalid ref cnt slice")
};
let team = self.darc_rt_team();
let mut reqs = vec![];
for pe in 0..ref_cnts.len() {
let cnt = ref_cnts[pe].swap(0, Ordering::SeqCst);
if cnt > 0 {
let my_addr = &*self as *const DarcInner<T> as usize;
let pe_addr = team.lamellae.comm().remote_addr(
team.arch.world_pe(pe).expect("invalid team member"),
my_addr,
);
trace!(
"[{:?}] sending finished to {:?} {:?} team {:?} {:x}",
std::thread::current().id(),
pe,
cnt,
team.team_hash,
my_addr
);
reqs.push(
team.spawn_am_pe_tg(
pe,
FinishedAm {
cnt,
src_pe: pe,
inner_addr: pe_addr,
},
Some(self.am_counters()),
)
.spawn(),
);
} else {
trace!(
"[{:?}] no finished to send to {:?} {:?} team {:?} {:x}",
std::thread::current().id(),
pe,
cnt,
team.team_hash,
&*self as *const DarcInner<T> as usize
);
}
}
reqs
}
async fn wait_on_state(
inner: DarcCommPtr<T>,
state: DarcMode,
extra_cnt: usize,
reset: bool,
) -> bool {
let team = inner.rt_team();
let rdma = team.lamellae.comm();
rdma.thread_flush();
for pe in inner.mode_slice.iter() {
let mut timer = std::time::Instant::now();
while *pe != state {
if inner.local_cnt.load(Ordering::SeqCst) == 1 + extra_cnt {
join_all(inner.send_finished()).await;
}
if !reset && timer.elapsed().as_secs_f64() > config().deadlock_warning_timeout {
println!("[{:?}][{:?}][WARNING] -- Potential deadlock detected.\n\
The runtime is currently waiting for all remaining references to this distributed object to be dropped.\n\
The object is likely a {:?} with {:?} remaining local references and {:?} remaining remote references, ref cnts by pe {:?}\n\
An example where this can occur can be found at https://docs.rs/lamellar/latest/lamellar/array/struct.ReadOnlyArray.html#method.into_local_lock\n\
The deadlock timeout can be set via the LAMELLAR_DEADLOCK_WARNING_TIMEOUT environment variable, the current timeout is {} seconds\n\
To view backtrace set RUST_LIB_BACKTRACE=1\n\
{}",
inner.my_pe,
std::thread::current().id(),
inner.mode_slice.as_slice(),
inner.local_cnt.load(Ordering::SeqCst),
inner.dist_cnt.load(Ordering::SeqCst),
inner.ref_cnt_slice,
config().deadlock_warning_timeout,
std::backtrace::Backtrace::capture()
);
timer = std::time::Instant::now();
}
if reset && timer.elapsed().as_secs_f64() > config().deadlock_warning_timeout / 2.0
{
println!("[{:?}][{:?}][WARNING] -- Sending RestartDrop.\n\
The runtime is currently waiting for all remaining references to this distributed object to be dropped. The object is likely a {:?}\n\
with {:?} remaining local references and {:?} remaining remote references, ref cnts by pe {:?}\n\
To view backtrace set RUST_LIB_BACKTRACE=1\n\
{}",
inner.my_pe,
std::thread::current().id(),
inner.mode_slice.as_slice(),
inner.local_cnt.load(Ordering::SeqCst),
inner.dist_cnt.load(Ordering::SeqCst),
inner.ref_cnt_slice,
std::backtrace::Backtrace::capture()
);
return false;
}
if reset && inner.mode_slice.iter().any(|x| *x == DarcMode::RestartDrop) {
return false;
}
rdma.thread_flush();
async_std::task::yield_now().await;
}
}
true
}
async fn broadcast_state(
inner: DarcCommPtr<T>,
team: Darc<LamellarTeamRT>,
state: DarcMode,
) {
let rdma = team.lamellae.comm();
let my_pe = inner.my_pe;
trace!(
"broadcast state {:?} {:?} {:?}",
state,
inner.mode_slice.as_ptr(),
inner.mode_slice.index_addr(my_pe)
);
for pe in team.arch.team_iter() {
trace!("putting state {:?} to pe {}", state, pe);
inner.mode_slice.put_unmanaged(state, pe, my_pe);
}
rdma.thread_wait(); trace!("broadcasted state {:?}", state);
}
async fn block_on_outstanding(mut inner: DarcCommPtr<T>, state: DarcMode, extra_cnt: usize) {
trace!(
"[{:?}] entering block_on_outstanding {:?} {:?}",
std::thread::current().id(),
inner.as_ptr(),
inner.mode_slice.as_ptr(),
);
let team = inner.darc_rt_team();
let orig_state = inner.mode_slice[inner.my_pe];
inner.await_all().await;
if team.num_pes() == 1 {
trace!(
"[{:?}] single pe block_on_outstanding {:?} {:?}",
std::thread::current().id(),
inner.as_ptr(),
inner.mode_slice.as_ptr(),
);
while inner.local_cnt.load(Ordering::SeqCst) > 1 + extra_cnt {
async_std::task::yield_now().await;
}
unsafe {
(*(((&inner.mode_slice[inner.my_pe]) as *const DarcMode) as *const AtomicU64)) .store(state as u64, Ordering::SeqCst)
};
} else {
trace!(
"[{:?}] multi pe block_on_outstanding {:?} {:?}",
std::thread::current().id(),
inner.as_ptr(),
inner.mode_slice.as_ptr(),
);
let mut outstanding_refs = true;
let mut prev_ref_cnts = vec![0usize; inner.num_pes];
let mut barrier_id = 1usize;
trace!(
"[{:?}] starting block_on_outstanding loop initial barrier_id: {:?} {:?} {:?} {:?}",
std::thread::current().id(),
barrier_id,
inner.as_ref(),
inner.local_cnt.load(Ordering::SeqCst),
extra_cnt
);
while inner.local_cnt.load(Ordering::SeqCst) > 1 + extra_cnt {
async_std::task::yield_now().await;
}
trace!(
"[{:?}] finished waiting for local cnt to meet threshold: {:?} {:?} {:?} {:?}",
std::thread::current().id(),
barrier_id,
inner.as_ref(),
inner.local_cnt.load(Ordering::SeqCst),
extra_cnt + 1
);
join_all(inner.send_finished()).await;
trace!(
"[{:?}] entering initial block_on barrier() {:?}",
std::thread::current().id(),
inner.as_ref()
);
if !Self::wait_on_state(inner.clone(), orig_state, extra_cnt, false).await {
panic!("deadlock waiting for original state");
}
let barrier_fut = unsafe { inner.barrier.as_ref().unwrap().async_barrier() };
barrier_fut.await;
trace!(
"[{:?}] leaving initial block_on barrier() {:?}",
std::thread::current().id(),
inner.as_ref()
);
while outstanding_refs {
if inner.mode_slice.iter().any(|x| *x == DarcMode::RestartDrop) {
Self::broadcast_state(inner.clone(), team.clone(), DarcMode::RestartDrop).await;
if !(Self::wait_on_state(
inner.clone(),
DarcMode::RestartDrop,
extra_cnt,
false,
)
.await)
{
panic!("deadlock");
}
Self::broadcast_state(inner.clone(), team.clone(), orig_state).await;
Box::pin(DarcInner::block_on_outstanding(
inner.clone(),
state,
extra_cnt,
))
.await;
return;
}
trace!(
"[{:?}] starting block_on loop iteration barrier_id: {:?} {:?}",
std::thread::current().id(),
barrier_id,
inner.as_ref()
);
outstanding_refs = false;
for id in inner.mode_barrier_slice.iter_mut() {
*id = 0;
}
let old_barrier_id = barrier_id; while inner.local_cnt.load(Ordering::SeqCst) > 1 + extra_cnt {
async_std::task::yield_now().await;
}
join_all(inner.send_finished()).await;
let barrier_fut = unsafe { inner.barrier.as_ref().unwrap().async_barrier() };
barrier_fut.await;
trace!(
"[{:?}] finished initial barrier {:?}",
std::thread::current().id(),
inner.as_ref()
);
let old_ref_cnts = inner.total_ref_cnt_slice.to_vec();
let old_local_cnt = inner.total_local_cnt.load(Ordering::SeqCst);
let old_dist_cnt = inner.total_dist_cnt.load(Ordering::SeqCst);
let rdma = team.lamellae.comm();
for pe in 0..inner.num_pes {
if prev_ref_cnts[pe] != old_ref_cnts[pe] {
let send_pe = team.arch.single_iter(pe).next().unwrap();
trace!(
"[{:?}] putting ref cnt {:?} to pe {:?} at offset {:?} {:?}",
std::thread::current().id(),
old_ref_cnts[pe],
send_pe,
inner.my_pe,
inner.as_ref()
);
inner.mode_ref_cnt_slice.put_unmanaged(
old_ref_cnts[pe],
send_pe,
inner.my_pe,
);
outstanding_refs = true;
barrier_id = 0;
}
}
rdma.thread_wait(); trace!(
"[{:?}] finished putting ref cnts {:?}",
std::thread::current().id(),
inner.as_ref()
);
rdma.thread_flush();
let barrier_fut = unsafe { inner.barrier.as_ref().unwrap().async_barrier() };
barrier_fut.await;
trace!(
"[{:?}] finished barrier after putting ref cnts {:?}",
std::thread::current().id(),
inner.as_ref()
);
outstanding_refs |= old_local_cnt != inner.total_local_cnt.load(Ordering::SeqCst);
outstanding_refs |= old_dist_cnt != inner.total_dist_cnt.load(Ordering::SeqCst);
let mut barrier_sum = 0;
for pe in 0..inner.num_pes {
outstanding_refs |= old_ref_cnts[pe] != inner.total_ref_cnt_slice[pe];
barrier_sum += inner.mode_ref_cnt_slice[pe];
}
outstanding_refs |= barrier_sum != old_dist_cnt;
if outstanding_refs {
barrier_id = 0;
}
rdma.thread_flush();
let barrier_fut = unsafe { inner.barrier.as_ref().unwrap().async_barrier() };
barrier_fut.await;
trace!(
"[{:?}] finished barrier after checking ref cnts {:?}",
std::thread::current().id(),
inner.as_ref()
);
for pe in 0..inner.num_pes {
let send_pe = team.arch.single_iter(pe).next().unwrap();
trace!(
"[{:?}] putting barrier_id {:?} to pe {:?} at offset {:?} {:?}",
std::thread::current().id(),
barrier_id,
send_pe,
inner.my_pe,
inner.as_ref()
);
inner
.mode_barrier_slice
.put_unmanaged(barrier_id, send_pe, inner.my_pe);
}
rdma.thread_wait(); trace!(
"[{:?}] after putting mode_barrier_slice{:?}",
std::thread::current().id(),
inner.as_ref()
);
rdma.thread_flush();
let barrier_fut = unsafe { inner.barrier.as_ref().unwrap().async_barrier() };
barrier_fut.await;
trace!(
"[{:?}] after barrier after putting mode_barrier_slice{:?}",
std::thread::current().id(),
inner.as_ref()
);
for id in inner.mode_barrier_slice.iter_mut() {
outstanding_refs |= *id == 0;
}
barrier_id = old_barrier_id + 1;
prev_ref_cnts = old_ref_cnts;
let barrier_fut = unsafe { inner.barrier.as_ref().unwrap().async_barrier() };
barrier_fut.await;
trace!(
"[{:?}] end of block_on loop iteration barrier_id: {:?} outstanding_refs: {:?} {:?}",
std::thread::current().id(),
barrier_id,
outstanding_refs,
inner.as_ref()
);
}
trace!(
"[{:?}] all outstanding refs are resolved {:?}",
std::thread::current().id(),
inner.as_ref(),
);
Self::broadcast_state(inner.clone(), team.clone(), state).await;
let barrier_fut = unsafe { inner.barrier.as_ref().unwrap().async_barrier() };
barrier_fut.await;
trace!(
"[{:?}] after barrier after putting dropped{:?}",
std::thread::current().id(),
inner.mode_slice.as_slice()
);
if !Self::wait_on_state(inner.clone(), state, extra_cnt, true).await {
Self::broadcast_state(inner.clone(), team.clone(), DarcMode::RestartDrop).await;
if !(Self::wait_on_state(inner.clone(), DarcMode::RestartDrop, extra_cnt, false)
.await)
{
panic!("deadlock");
}
Self::broadcast_state(inner.clone(), team.clone(), orig_state).await;
Box::pin(DarcInner::block_on_outstanding(
inner.clone(),
state,
extra_cnt,
))
.await;
return;
}
let barrier_fut = unsafe { inner.barrier.as_ref().unwrap().async_barrier() };
barrier_fut.await;
}
}
pub(crate) async fn await_all(&self) {
self.rt_team().lamellae.comm().wait_all(); let mut temp_now = Instant::now();
let am_counters = self.am_counters();
let mut orig_reqs = am_counters.send_req_cnt.load(Ordering::SeqCst);
let mut orig_launched = am_counters.launched_req_cnt.load(Ordering::SeqCst);
let mut done = false;
while !done {
while self.rt_team().panic.load(Ordering::SeqCst) == 0
&& ((am_counters.outstanding_reqs.load(Ordering::SeqCst) > 0)
|| orig_reqs != am_counters.send_req_cnt.load(Ordering::SeqCst)
|| orig_launched != am_counters.launched_req_cnt.load(Ordering::SeqCst))
{
orig_reqs = am_counters.send_req_cnt.load(Ordering::SeqCst);
orig_launched = am_counters.launched_req_cnt.load(Ordering::SeqCst);
async_std::task::yield_now().await;
if temp_now.elapsed().as_secs_f64() > config().deadlock_warning_timeout {
println!(
"in darc await_all mype: {:?} cnt: {:?} {:?}",
self.rt_team().world_pe,
am_counters.send_req_cnt.load(Ordering::SeqCst),
am_counters.outstanding_reqs.load(Ordering::SeqCst),
);
temp_now = Instant::now();
}
}
if am_counters.send_req_cnt.load(Ordering::SeqCst)
!= am_counters.launched_req_cnt.load(Ordering::SeqCst)
{
if am_counters.outstanding_reqs.load(Ordering::SeqCst) > 0
|| orig_reqs != am_counters.send_req_cnt.load(Ordering::SeqCst)
|| orig_launched != am_counters.launched_req_cnt.load(Ordering::SeqCst)
{
continue;
}
println!(
"in darc await_all mype: {:?} cnt: {:?} {:?} {:?}",
self.rt_team().world_pe,
am_counters.send_req_cnt.load(Ordering::SeqCst),
am_counters.outstanding_reqs.load(Ordering::SeqCst),
am_counters.launched_req_cnt.load(Ordering::SeqCst)
);
RuntimeWarning::UnspawnedTask(
"`await_all` before all tasks/active messages have been spawned",
)
.print();
}
done = true;
}
}
}
impl<T: 'static> DarcInner<T> {
#[allow(dead_code)]
fn item(&self) -> &T {
unsafe { &(*self.item) }
}
}
impl<T: 'static> fmt::Debug for DarcInner<T> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "[{:}/{:?}] ", self.my_pe, self.num_pes)?;
write!(f, "id: {} ", self.id)?;
write!(f, "lc: {:?} ", self.local_cnt.load(Ordering::SeqCst))?;
write!(f, "dc: {:?} ", self.dist_cnt.load(Ordering::SeqCst))?;
write!(f, "wc: {:?} ", self.weak_local_cnt.load(Ordering::SeqCst))?;
write!(f, "ref_cnt: {:?} ", self.ref_cnt_slice.as_slice())?;
write!(
f,
"am_cnt ({:?},{:?}) ",
self.am_counters().outstanding_reqs.load(Ordering::Relaxed),
self.am_counters().send_req_cnt.load(Ordering::Relaxed)
)?;
write!(f, "mode {:?} ", self.mode_slice.as_slice())?;
write!(f, "item addr: {:x} ", self.item as usize)?;
write!(f, "my addr: {:x} ", self as *const _ as usize)?;
Ok(())
}
}
enum TeamAndItem<T> {
Team(LamellarTeamRT),
NonTeam(Darc<LamellarTeamRT>, T),
}
#[lamellar_prof::prof]
impl<T> TeamAndItem<T> {
fn team(&self) -> &LamellarTeamRT {
match self {
TeamAndItem::Team(team_rt) => team_rt,
TeamAndItem::NonTeam(team_darc, _) => &team_darc,
}
}
fn into_raw(self) -> (*const DarcInner<LamellarTeamRT>, *const T) {
match self {
TeamAndItem::Team(team_rt) => {
let ptr = Box::into_raw(Box::new(team_rt));
(ptr as *const DarcInner<LamellarTeamRT>, ptr as *const T)
}
TeamAndItem::NonTeam(team_darc, item) => (
Darc::into_raw_team(team_darc),
Box::into_raw(Box::new(item)) as *const T,
),
}
}
}
#[lamellar_prof::prof]
impl Darc<LamellarTeamRT> {
pub(crate) async fn async_try_new_team_darc(
team_rt: LamellarTeamRT,
) -> Result<Darc<LamellarTeamRT>, IdError> {
trace!("creating team darc");
Darc::async_try_new_with_drop_inner(TeamAndItem::Team(team_rt), DarcMode::WorldTeam, None)
.await
}
pub(crate) unsafe fn team_from_raw(
ptr: *const DarcInner<LamellarTeamRT>,
) -> Darc<LamellarTeamRT> {
let alloc = (*ptr)
.rt_team()
.lamellae
.comm()
.get_alloc_cloned(CommAllocAddr(ptr as usize))
.expect("invalid darc ptr");
let inner: DarcCommPtr<LamellarTeamRT> = DarcCommPtr {
alloc: alloc.clone(),
_phantom: PhantomData,
};
let src_pe = (*ptr).my_pe;
let id = (*ptr).total_local_cnt.fetch_add(1, Ordering::SeqCst);
let d = Darc { inner, src_pe, id };
trace!("reconstructed team darc {:?}", d.inner());
d
}
pub(crate) unsafe fn cloned_team_from_raw(
ptr: *const DarcInner<LamellarTeamRT>,
) -> Darc<LamellarTeamRT> {
let alloc = (*ptr)
.rt_team()
.lamellae
.comm()
.get_alloc_cloned(CommAllocAddr(ptr as usize))
.expect("invalid darc ptr");
let inner: DarcCommPtr<LamellarTeamRT> = DarcCommPtr {
alloc: alloc.clone(),
_phantom: PhantomData,
};
let src_pe = (*ptr).my_pe;
(*ptr).local_cnt.fetch_add(1, Ordering::SeqCst);
let id = (*ptr).total_local_cnt.fetch_add(1, Ordering::SeqCst);
let d = Darc { inner, src_pe, id };
trace!(
"[{}][{}] reconstructed team darc cloned {:?}",
d.inner().id,
d.id,
d.inner()
);
d
}
pub(crate) fn into_raw_team(self) -> *const DarcInner<LamellarTeamRT> {
trace!(
"[{}][{}] {:?} into_raw_team {:?}",
self.inner().id,
self.id,
self.inner(),
self.inner.as_ptr()
);
self.inc_local_cnt(1); self.inner.as_ptr()
}
}
impl<T> Darc<T> {
#[lamellar_prof::prof]
pub fn downgrade(the_darc: &Darc<T>) -> WeakDarc<T> {
trace!("downgrading darc {:?}", the_darc.id);
the_darc
.inner()
.weak_local_cnt
.fetch_add(1, Ordering::SeqCst);
let weak = WeakDarc {
inner: the_darc.inner.clone(),
src_pe: the_darc.src_pe,
};
weak
}
pub(crate) fn darc_addr(&self) -> usize {
self.inner.addr().into()
}
#[lamellar_prof::prof]
pub(crate) async fn into_inner(self) -> T {
DarcInner::block_on_outstanding(
self.inner.clone(),
DarcMode::Dropped,
1, )
.await;
trace!("into_inner after block on outstanding {:?}", self.id);
let this = std::mem::ManuallyDrop::new(self);
let item = unsafe { Box::from_raw(this.inner().item as *mut T) };
trace!("into_inner after box from raw");
let _ = unsafe { Arc::from_raw(this.inner().am_counters) };
trace!("into_inner after am_counters");
let _ = unsafe { Box::from_raw(this.inner().barrier) };
trace!("into_inner after barrier");
let mut darc_comm_ptr = std::mem::MaybeUninit::uninit();
unsafe {
std::ptr::copy_nonoverlapping(&this.inner, darc_comm_ptr.as_mut_ptr(), 1);
let darc_comm_ptr = darc_comm_ptr.assume_init();
let mut darc_inner = std::mem::MaybeUninit::<DarcInner<T>>::uninit();
std::ptr::copy_nonoverlapping(
darc_comm_ptr.alloc.as_ptr::<DarcInner<T>>(),
darc_inner.as_mut_ptr(),
1,
);
let _darc_inner = darc_inner.assume_init();
}
trace!("into_inner after inner_temp");
*item
}
pub(crate) fn inner(&self) -> &DarcInner<T> {
self.inner.as_ref().expect("invalid darc inner ptr")
}
fn inner_mut(&self) -> &mut DarcInner<T> {
self.inner.as_mut().expect("invalid darc inner ptr")
}
#[doc(hidden)]
#[lamellar_prof::prof]
pub fn serialize_update_cnts(&self, cnt: usize) {
trace!(target: "darc_clone", "darc[{:?}] serialize darc cnts {:?}", self.id, self.inner());
self.inner()
.dist_cnt
.fetch_add(cnt, std::sync::atomic::Ordering::SeqCst);
self.inner()
.total_dist_cnt
.fetch_add(cnt, std::sync::atomic::Ordering::SeqCst);
}
#[doc(hidden)]
#[lamellar_prof::prof]
pub fn deserialize_update_cnts(&self) {
trace!(target: "darc_clone", "darc[{:?}] deserialize darc cnts {:?}", self.id, self.inner());
self.inner().inc_pe_ref_count(self.src_pe, 1);
self.inner().local_cnt.fetch_add(1, Ordering::SeqCst);
self.inner().total_local_cnt.fetch_add(1, Ordering::SeqCst);
}
#[doc(hidden)]
#[lamellar_prof::prof]
pub fn inc_local_cnt(&self, cnt: usize) -> usize {
trace!("darc[{:?}] inc_local_cnt {:?}", self.id, self.inner());
self.inner().local_cnt.fetch_add(cnt, Ordering::SeqCst);
self.inner()
.total_local_cnt
.fetch_add(cnt, Ordering::SeqCst)
}
#[doc(hidden)]
pub fn print(&self) {
println!(
"[{:?}]--------\nid: {:?}:{:?} orig: {:?} 0x{:x} item_addr {:?} {:?}\n--------[{:?}]",
std::thread::current().id(),
self.inner().id,
self.id,
self.src_pe,
self.inner.addr(),
self.inner().item,
self.inner(),
std::thread::current().id(),
);
}
}
fn calc_padding(addr: usize, align: usize) -> usize {
let rem = addr % align;
if rem == 0 {
0
} else {
align - rem
}
}
#[lamellar_prof::prof]
impl<T: Send + Sync> Darc<T> {
#[doc(alias = "Collective")]
pub fn new<U: Into<IntoLamellarTeam>>(team: U, item: T) -> DarcHandle<T> {
let team = team.into().team.clone();
DarcHandle {
team: team.clone(),
launched: false,
creation_future: Box::pin(Darc::async_try_new_with_drop(
team,
item,
DarcMode::Darc,
None,
)),
}
}
pub(crate) async fn async_try_new_with_drop<U: Into<IntoLamellarTeam>>(
team: U,
item: T,
state: DarcMode,
drop: Option<fn(&mut T) -> bool>,
) -> Result<Darc<T>, IdError> {
trace!("creating darc with drop");
Darc::async_try_new_with_drop_inner(
TeamAndItem::NonTeam(team.into().team.clone(), item),
state,
drop,
)
.await
}
async fn async_try_new_with_drop_inner(
team_and_item: TeamAndItem<T>,
state: DarcMode,
custom_drop: Option<fn(&mut T) -> bool>,
) -> Result<Darc<T>, IdError> {
trace!(target: "new", "creating darc");
let timer = Instant::now();
let team_rt = team_and_item.team();
let my_pe = team_rt.team_pe?;
let num_pes = team_rt.num_pes;
let alloc = if team_rt.num_pes == team_rt.num_world_pes {
trace!("using global alloc for darc");
AllocationType::Global
} else {
trace!("using sub team alloc for darc");
AllocationType::Sub(team_rt.get_pes())
};
let mut size = std::mem::size_of::<DarcInner<T>>();
let padding = calc_padding(size, std::mem::align_of::<usize>());
let ref_cnt_offset = size + padding;
size += padding + team_rt.num_pes * std::mem::size_of::<usize>();
let padding = calc_padding(size, std::mem::align_of::<usize>());
let total_ref_cnt_offset = size + padding;
size += padding + team_rt.num_pes * std::mem::size_of::<usize>();
let padding = calc_padding(size, std::mem::align_of::<DarcMode>());
let mode_offset = size + padding;
size += padding + team_rt.num_pes * std::mem::size_of::<DarcMode>();
let padding = calc_padding(size, std::mem::align_of::<usize>());
let mode_ref_cnt_offset = size + padding;
size += padding + team_rt.num_pes * std::mem::size_of::<usize>();
let padding = calc_padding(size, std::mem::align_of::<usize>());
let mode_barrier_offset = size + padding;
size += padding + team_rt.num_pes * std::mem::size_of::<usize>();
trace!("Darc::new before barrier time: {:?}", timer.elapsed());
team_rt.async_barrier().await;
trace!("Darc::new after barrier time: {:?}", timer.elapsed());
let darc_alloc = team_rt
.lamellae
.comm()
.alloc(size, alloc, std::mem::align_of::<DarcInner<T>>())
.expect("out of memory");
trace!("Darc::new after alloc time: {:?}", timer.elapsed());
trace!(
"[{:?}] creating new darc[{:?}] {:?} alloc: {:?}",
std::thread::current().id(),
DARC_ID.load(Ordering::Relaxed),
team_rt.team_hash,
darc_alloc
);
trace!("Darc::new after team_ptr time: {:?}", timer.elapsed());
let am_counters = Arc::new(AMCounters::new());
let am_counters_ptr = Arc::into_raw(am_counters);
let barrier = Box::new(Barrier::new(
team_rt.world_pe,
team_rt.num_world_pes,
team_rt.lamellae.clone(),
team_rt.arch.clone(),
team_rt.scheduler.clone(),
team_rt.panic.clone(),
));
trace!(target: "lamellae_debug", "Darc::new after barrier creation lamellae cnt: {:?}", Arc::strong_count(&team_rt.lamellae));
let barrier_ptr = Box::into_raw(barrier);
trace!("Darc::new after barrier_creation {:?}", timer.elapsed());
unsafe {
let darc_temp_ptr = darc_alloc.as_mut_ptr::<DarcInner<T>>();
let (mut team, item) = team_and_item.into_raw();
trace!(
"team ptr: {:p}, item ptr: {:p}, setting darc inner ptrs, start addr: {:?}",
team,
item,
darc_temp_ptr
);
if team.addr() == item.addr() {
team = darc_temp_ptr as *const DarcInner<LamellarTeamRT>;
}
(*darc_temp_ptr).team = team; trace!("darc team ptr addr: {:p}", &(*darc_temp_ptr).team);
(*darc_temp_ptr).item = item;
trace!("darc item ptr addr: {:p}", &(*darc_temp_ptr).item);
trace!("id ptr: {:?}", &(*darc_temp_ptr).id as *const usize);
(*darc_temp_ptr).id = DARC_ID.fetch_add(1, Ordering::Relaxed);
trace!("my_pe ptr: {:?}", &(*darc_temp_ptr).my_pe as *const usize);
(*darc_temp_ptr).my_pe = my_pe;
trace!(
"num_pes ptr: {:?}",
&(*darc_temp_ptr).num_pes as *const usize
);
(*darc_temp_ptr).num_pes = num_pes;
trace!(
"local_cnt ptr: {:?}",
&(*darc_temp_ptr).local_cnt as *const AtomicUsize
);
(*darc_temp_ptr).local_cnt = AtomicUsize::new(1);
trace!(
"total_local_cnt ptr: {:?}",
&(*darc_temp_ptr).total_local_cnt as *const AtomicUsize
);
(*darc_temp_ptr).total_local_cnt = AtomicUsize::new(1);
trace!(
"weak_local_cnt ptr: {:?}",
&(*darc_temp_ptr).weak_local_cnt as *const AtomicUsize
);
(*darc_temp_ptr).weak_local_cnt = AtomicUsize::new(0);
trace!(
"dist_cnt ptr: {:?}",
&(*darc_temp_ptr).dist_cnt as *const AtomicUsize
);
(*darc_temp_ptr).dist_cnt = AtomicUsize::new(0);
trace!(
"total_dist_cnt ptr: {:?}",
&(*darc_temp_ptr).total_dist_cnt as *const AtomicUsize
);
(*darc_temp_ptr).total_dist_cnt = AtomicUsize::new(0);
trace!(
"am_counters ptr: {:?}",
&(*darc_temp_ptr).am_counters as *const *const AMCounters
);
(*darc_temp_ptr).am_counters = std::ptr::null();
trace!(
"creating slices, ref_cnt_slice ptr: {:?}",
&(*darc_temp_ptr).ref_cnt_slice as *const CommSlice<usize>
);
std::ptr::write(
&mut (*darc_temp_ptr).ref_cnt_slice,
darc_alloc.comm_slice_at_byte_offset(ref_cnt_offset, num_pes),
);
trace!(
"ref_cnt_slice {:?} padding: {:?}",
(*darc_temp_ptr).ref_cnt_slice,
calc_padding(
(*darc_temp_ptr).ref_cnt_slice.usize_addr(),
std::mem::align_of::<usize>()
)
);
std::ptr::write(
&mut (*darc_temp_ptr).total_ref_cnt_slice,
darc_alloc.comm_slice_at_byte_offset(total_ref_cnt_offset, num_pes),
);
trace!(
"total ref cnt slice {:?} padding: {:?}",
(*darc_temp_ptr).total_ref_cnt_slice,
calc_padding(
(*darc_temp_ptr).total_ref_cnt_slice.usize_addr(),
std::mem::align_of::<usize>()
)
);
std::ptr::write(
&mut (*darc_temp_ptr).mode_slice,
darc_alloc.comm_slice_at_byte_offset(mode_offset, num_pes),
);
trace!(
"mode slice {:?} padding: {:?}",
(*darc_temp_ptr).mode_slice,
calc_padding(
(*darc_temp_ptr).mode_slice.usize_addr(),
std::mem::align_of::<DarcMode>()
)
);
std::ptr::write(
&mut (*darc_temp_ptr).mode_ref_cnt_slice,
darc_alloc.comm_slice_at_byte_offset(mode_ref_cnt_offset, num_pes),
);
trace!(
"mode ref cnt slice {:?} padding: {:?}",
(*darc_temp_ptr).mode_ref_cnt_slice,
calc_padding(
(*darc_temp_ptr).mode_ref_cnt_slice.usize_addr(),
std::mem::align_of::<usize>()
)
);
std::ptr::write(
&mut (*darc_temp_ptr).mode_barrier_slice,
darc_alloc.comm_slice_at_byte_offset(mode_barrier_offset, num_pes),
);
trace!(
"mode barrier slice {:?} padding: {:?}",
(*darc_temp_ptr).mode_barrier_slice,
calc_padding(
(*darc_temp_ptr).mode_barrier_slice.usize_addr(),
std::mem::align_of::<usize>()
)
);
(*darc_temp_ptr).barrier = barrier_ptr;
(*darc_temp_ptr).am_counters = am_counters_ptr;
(*darc_temp_ptr).drop = custom_drop;
(*darc_temp_ptr).valid = AtomicBool::new(true);
}
let d = Darc {
inner: DarcCommPtr {
alloc: darc_alloc,
_phantom: std::marker::PhantomData::<DarcInner<T>>,
}, src_pe: my_pe,
id: 0,
};
for elem in d.inner().ref_cnt_slice.clone().iter_mut() {
*elem = 0;
}
for elem in d.inner().total_ref_cnt_slice.clone().iter_mut() {
*elem = 0;
}
for elem in d.inner().mode_slice.clone().iter_mut() {
*elem = state;
}
for elem in d.inner().mode_ref_cnt_slice.clone().iter_mut() {
*elem = 0;
}
for elem in d.inner().mode_barrier_slice.clone().iter_mut() {
*elem = 0;
}
trace!(
target: "new",
" [{:?}] created new darc[{:?}] , next_inner_id: {:?} {:?} ",
std::thread::current().id(),
d.id,
DARC_ID.load(Ordering::Relaxed),
d.inner(),
);
d.inner().rt_team().async_barrier().await;
Ok(d)
}
pub(crate) async fn block_on_outstanding(self, state: DarcMode, extra_cnt: usize) {
DarcInner::block_on_outstanding(self.inner.clone(), state, extra_cnt).await;
}
#[doc(alias = "Collective")]
pub fn into_localrw(self) -> IntoLocalRwDarcHandle<T> {
let inner = self.inner.clone();
let team = self.inner().darc_rt_team();
IntoLocalRwDarcHandle {
darc: self.into(),
team,
launched: false,
outstanding_future: Box::pin(async move {
DarcInner::block_on_outstanding(inner, DarcMode::LocalRw, 0).await;
}),
}
}
#[doc(alias = "Collective")]
pub fn into_globalrw(self) -> IntoGlobalRwDarcHandle<T> {
let inner = self.inner.clone();
let team = self.inner().darc_rt_team();
IntoGlobalRwDarcHandle {
darc: self.into(),
team,
launched: false,
outstanding_future: Box::pin(async move {
DarcInner::block_on_outstanding(inner, DarcMode::GlobalRw, 0).await;
}),
}
}
}
#[lamellar_prof::prof]
impl<T> Clone for Darc<T> {
fn clone(&self) -> Self {
self.inner().local_cnt.fetch_add(1, Ordering::SeqCst);
let id = self.inner().total_local_cnt.fetch_add(1, Ordering::SeqCst);
trace! {target:"darc_clone", "[{:?}] darc[{:?}][{id}] cloned from [{:?}] {:?} {:?} {:?}", std::thread::current().id(),self.inner().id,self.id,self.inner,self.inner().local_cnt.load(Ordering::SeqCst),self.inner().total_local_cnt.load(Ordering::SeqCst)};
Darc {
inner: self.inner.clone(),
src_pe: self.src_pe,
id,
}
}
}
impl<T> Deref for Darc<T> {
type Target = T;
#[inline]
fn deref(&self) -> &T {
self.inner().item()
}
}
#[lamellar_prof::prof]
impl<T: Hash> Hash for Darc<T> {
fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
(**self).hash(state);
}
}
#[lamellar_prof::prof]
impl<T: PartialEq> PartialEq for Darc<T> {
fn eq(&self, other: &Self) -> bool {
(**self).eq(&**other)
}
}
impl<T: Eq> Eq for Darc<T> {}
impl<T: fmt::Display> fmt::Display for Darc<T> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
fmt::Display::fmt(&**self, f)
}
}
impl<T: fmt::Debug> fmt::Debug for Darc<T> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
fmt::Debug::fmt(&**self, f)
}
}
macro_rules! launch_drop {
($mode:ty, $inner:ident, $inner_ptr:expr) => {{
trace!("launching drop task as {}", stringify!($mode));
let team = $inner.darc_rt_team();
for pe in team.arch.team_iter() {
let dropping = DarcMode::Dropping;
trace!(
"[{:?}] putting dropping to pe {:?} at offset {:?} mode_slice ptr: {:?}",
std::thread::current().id(),
pe,
$inner.my_pe,
$inner.mode_slice.as_ptr()
);
$inner
.mode_slice
.put_unmanaged::<DarcMode>(dropping, pe, $inner.my_pe);
}
team.team_counters.inc_outstanding(1);
team.world_counters.inc_outstanding(1); let am = team.exec_am_local(DroppedWaitAM {
inner: $inner_ptr,
my_pe: $inner.my_pe,
num_pes: $inner.num_pes,
team: team.clone(),
phantom: PhantomData::<T>,
});
let task = am.spawn();
team.team_counters.dec_outstanding(1);
team.world_counters.dec_outstanding(1);
task
}};
}
#[lamellar_prof::prof]
impl<T: 'static> Drop for Darc<T> {
fn drop(&mut self) {
let inner = self.inner();
let cnt = inner.local_cnt.fetch_sub(1, Ordering::SeqCst);
trace! {target: "drop", "[{:?}] darc[{:?}][{:?}] dropped {:?} {:?} {:?}: team: {:?}",std::thread::current().id(),self.inner().id,self.id,self.inner,self.inner().local_cnt.load(Ordering::SeqCst),inner.total_local_cnt.load( Ordering::SeqCst), unsafe{&*inner.team}};
if cnt == 0 {
panic!("darc dropped too many times");
}
if cnt == 1 {
trace!("last local darc ref dropped on pe {:?}", inner.my_pe);
if self.inner().ref_cnt_slice.iter().any(|&x| x > 0) {
trace!(
"last local darc ref dropped on pe {:?} need to send finished.",
inner.my_pe
);
inner.send_finished();
}
}
if inner.local_cnt.load(Ordering::SeqCst) == 0 {
trace!(target: "drop",
"no more local references on pe {:?} launching dropped darc am",
inner.my_pe
);
let cur_mode = inner.set_dropping();
if !cur_mode.drop_am_launched() {
let task = launch_drop!(cur_mode, inner, self.inner.clone());
if cur_mode == DarcMode::WorldTeam {
task.block();
}
} else {
trace!(
target: "drop",
"darc drop already in progress on pe {:?} skipping launch",
inner.my_pe,
);
}
}
trace!(target: "drop", "end drop Darc");
}
}
#[lamellar_impl::AmLocalDataRT]
struct DroppedWaitAM<T> {
inner: DarcCommPtr<T>,
my_pe: usize,
num_pes: usize,
team: Darc<LamellarTeamRT>, phantom: PhantomData<T>,
}
impl<T: 'static> std::fmt::Debug for DroppedWaitAM<T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"DroppedWaitAM {{ inner: {:?}, my_pe: {:?}, num_pes: {:?}, team: {:?} }}",
self.inner, self.my_pe, self.num_pes, self.team
)
}
}
unsafe impl<T> Send for DroppedWaitAM<T> {}
unsafe impl<T> Sync for DroppedWaitAM<T> {}
pub(crate) struct DarcCommPtr<T> {
alloc: CommAlloc,
_phantom: std::marker::PhantomData<DarcInner<T>>,
}
impl<T> std::fmt::Debug for DarcCommPtr<T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "DarcCommPtr {{ alloc: {:?} }}", self.alloc)
}
}
unsafe impl<T> Send for DarcCommPtr<T> {}
unsafe impl<T> Sync for DarcCommPtr<T> {}
impl<T> DarcCommPtr<T> {
pub(crate) fn as_ptr(&self) -> *const DarcInner<T> {
unsafe { self.alloc.as_ptr() }
}
pub(crate) fn as_mut_ptr(&self) -> *mut DarcInner<T> {
unsafe { self.alloc.as_mut_ptr() }
}
pub(crate) fn as_ref(&self) -> Option<&DarcInner<T>> {
unsafe { self.as_ptr().as_ref() }
}
pub(crate) fn as_mut(&self) -> Option<&mut DarcInner<T>> {
unsafe { self.as_mut_ptr().as_mut() }
}
pub(crate) fn addr(&self) -> CommAllocAddr {
self.alloc.comm_addr()
}
pub(crate) fn transmute<U>(&self) -> DarcCommPtr<U> {
DarcCommPtr {
alloc: CommAlloc {
inner_alloc: self.alloc.inner_alloc.clone(),
},
_phantom: std::marker::PhantomData,
}
}
}
#[lamellar_prof::prof]
impl<T> Clone for DarcCommPtr<T> {
fn clone(&self) -> Self {
DarcCommPtr {
alloc: CommAlloc {
inner_alloc: self.alloc.inner_alloc.clone(),
},
_phantom: self._phantom,
}
}
}
impl<T> std::ops::Deref for DarcCommPtr<T> {
type Target = DarcInner<T>;
fn deref(&self) -> &Self::Target {
unsafe { self.as_ptr().as_ref().unwrap() }
}
}
impl<T> std::ops::DerefMut for DarcCommPtr<T> {
fn deref_mut(&mut self) -> &mut Self::Target {
unsafe { self.as_mut_ptr().as_mut().unwrap() }
}
}
#[lamellar_impl::rt_am_local]
impl<T: 'static> LamellarAM for DroppedWaitAM<T> {
async fn exec(self) {
let mut timeout = std::time::Instant::now();
let is_world_team = unsafe { (&*self.inner.team).item.addr() } == self.inner.item.addr();
trace!(
target: "drop", "[{:?}] in DroppedWaitAM {:?} {:?} {:?} {:?} {:?} {:?} {:x} {:x} ",
std::thread::current().id(),
self.inner,
self.inner.mode_slice.as_slice(), self.inner.mode_slice,
self.inner.id,
self.inner.local_cnt.load(Ordering::SeqCst),
self.inner.total_local_cnt.load(Ordering::SeqCst),
unsafe { (&*self.inner.team).item.addr() },
self.inner.item.addr(),
);
let extra_cnt = if is_world_team { 6 } else { 0 };
let block_on_fut =
{ DarcInner::block_on_outstanding(self.inner.clone(), DarcMode::Dropped, extra_cnt) };
block_on_fut.await;
trace!(
"[{:?}] past block_on_outstanding {:?}",
std::thread::current().id(),
self.inner.as_ref()
);
for pe in self.inner.mode_slice.iter() {
while *pe != DarcMode::Dropped {
async_std::task::yield_now().await;
if self.inner.local_cnt.load(Ordering::SeqCst) == 0 {
join_all(self.inner.send_finished()).await;
}
if timeout.elapsed().as_secs_f64() > config().deadlock_warning_timeout {
println!("[{:?}][WARNING] -- Potential deadlock detected when trying to free distributed object.\n\
The runtime is currently waiting for all remaining references to this distributed object to be dropped.\n\
The current status of the object on each pe is {:?} with {:?} remaining local references and {:?} remaining remote references, ref cnts by pe {:?}\n\
the deadlock timeout can be set via the LAMELLAR_DEADLOCK_WARNING_TIMEOUT environment variable, the current timeout is {} seconds\n\
To view backtrace set RUST_LIB_BACKTRACE=1\n\
{}",
std::thread::current().id(),
self.inner.mode_slice.as_slice(),
self.inner.local_cnt.load(Ordering::SeqCst),
self.inner.ref_cnt_slice.as_slice(),
self.inner.dist_cnt.load(Ordering::SeqCst),
config().deadlock_warning_timeout,
std::backtrace::Backtrace::capture()
);
timeout = std::time::Instant::now();
}
}
}
trace!("after DarcMode::Dropped");
unsafe {
self.inner.valid.store(false, Ordering::SeqCst);
while self.inner.dist_cnt.load(Ordering::SeqCst) != 0
|| self.inner.local_cnt.load(Ordering::SeqCst) != extra_cnt
{
if self.inner.local_cnt.load(Ordering::SeqCst) == extra_cnt {
join_all(self.inner.send_finished()).await;
}
if timeout.elapsed().as_secs_f64() > config().deadlock_warning_timeout {
println!("[{:?}][WARNING] --- Potential deadlock detected when trying to free distributed object.\n\
The runtime is currently waiting for all remaining references to this distributed object to be dropped.\n\
The current status of the object on each pe is {:?} with {:?} remaining local references and {:?} remaining remote references, ref cnts by pe {:?}\n\
the deadlock timeout can be set via the LAMELLAR_DEADLOCK_WARNING_TIMEOUT environment variable, the current timeout is {} seconds\n\
To view backtrace set RUST_LIB_BACKTRACE=1\n\
{}",
std::thread::current().id(),
self.inner.mode_slice.as_slice(),
self.inner.local_cnt.load(Ordering::SeqCst),
self.inner.ref_cnt_slice.as_slice(),
self.inner.dist_cnt.load(Ordering::SeqCst),
config().deadlock_warning_timeout,
std::backtrace::Backtrace::capture()
);
timeout = std::time::Instant::now();
}
async_std::task::yield_now().await;
}
trace!("going to drop object");
if let Some(my_drop) = self.inner.drop {
let mut dropped_done = false;
while !dropped_done {
dropped_done = my_drop(&mut *(self.inner.item as *mut T));
async_std::task::yield_now().await;
}
}
let _ = Box::from_raw(self.inner.item as *mut T);
trace!("afterdrop object");
while self.inner.weak_local_cnt.load(Ordering::SeqCst) != 0 {
async_std::task::yield_now().await;
}
let _am_counters = Arc::from_raw(self.inner.am_counters);
let _barrier = Box::from_raw(self.inner.barrier);
let mut darc_temp = std::mem::MaybeUninit::<DarcInner<T>>::uninit();
std::ptr::copy_nonoverlapping(
self.inner.alloc.as_ptr::<DarcInner<T>>(),
darc_temp.as_mut_ptr(),
1,
);
darc_temp.assume_init();
trace!("after darc_temp {:?}", &*self.inner.team);
Darc::team_from_raw(self.inner.team);
trace!("after team from raw {:?}", &*self.inner.team);
trace!(
target: "drop",
"[{:?}]leaving DroppedWaitAM {:?}",
std::thread::current().id(),
self
);
}
}
}
#[doc(hidden)]
#[derive(serde::Deserialize, serde::Serialize, Clone)]
pub struct __NetworkDarc {
inner_addr: CommAllocAddr,
backend: Backend,
orig_world_pe: usize,
orig_team_pe: usize,
}
impl std::fmt::Debug for __NetworkDarc {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "NetworkDarc {{ inner_addr: {:x}, backend: {:?}, orig_world_pe: {:?}, orig_team_pe: {:?} }}", self.inner_addr, self.backend, self.orig_world_pe, self.orig_team_pe)
}
}
#[lamellar_prof::prof]
impl<T> From<Darc<T>> for __NetworkDarc {
fn from(darc: Darc<T>) -> Self {
trace!(target:"darc_clone", "net darc from darc id: {:?} {:?}", darc.id, darc.inner());
let team = darc.inner().rt_team();
let ndarc = __NetworkDarc {
inner_addr: darc.inner.addr(),
backend: team.lamellae.comm().backend(),
orig_world_pe: team.world_pe,
orig_team_pe: team.team_pe.expect("darcs only valid on team members"),
};
ndarc
}
}
#[lamellar_prof::prof]
impl<T> From<&Darc<T>> for __NetworkDarc {
fn from(darc: &Darc<T>) -> Self {
trace!(target:"darc_clone", "net darc from &darc id: {:?} {:?}\n{:?}", darc.id, darc.inner(), std::backtrace::Backtrace::capture());
let team = darc.inner().rt_team();
let ndarc = __NetworkDarc {
inner_addr: darc.inner.addr(),
backend: team.lamellae.comm().backend(),
orig_world_pe: team.world_pe,
orig_team_pe: team.team_pe.expect("darcs only valid on team members"),
};
ndarc
}
}
#[lamellar_prof::prof]
impl<T> From<__NetworkDarc> for Darc<T> {
fn from(ndarc: __NetworkDarc) -> Self {
if let Some(lamellae) = LAMELLAES.read().get(&ndarc.backend) {
trace!(target: "lamellae_debug", "creating darc from network darc lamellae cnt: {:?} ", Arc::strong_count(lamellae));
let local_addr = lamellae
.comm()
.local_addr(ndarc.orig_world_pe, ndarc.inner_addr.0);
let alloc = lamellae
.comm()
.get_alloc_cloned(local_addr)
.expect("alloc should be valid on remote PE");
trace!("found alloc: {:?}", alloc);
let inner = DarcCommPtr {
alloc,
_phantom: PhantomData,
};
trace!("inner: {:?}", inner.as_ref());
inner.inc_pe_ref_count(ndarc.orig_team_pe, 1);
inner.local_cnt.fetch_add(1, Ordering::SeqCst);
let id = inner.total_local_cnt.fetch_add(1, Ordering::SeqCst);
trace!(target:"darc_clone","darc from network darc id: {} {:?}", id, inner.as_ref());
let darc = Darc {
inner,
src_pe: ndarc.orig_team_pe,
id,
};
darc
} else {
println!(
"ndarc: 0x{:x} {:?} {:?} {:?} ",
ndarc.inner_addr, ndarc.backend, ndarc.orig_world_pe, ndarc.orig_team_pe
);
panic!("unexepected lamellae backend {:?}", &ndarc.backend);
}
}
}