use async_lock::{RwLock, RwLockReadGuardArc, RwLockWriteGuardArc};
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use std::fmt;
use std::sync::atomic::Ordering;
use std::sync::Arc;
use crate::{
active_messaging::RemotePtr,
darc::{Darc, DarcInner, DarcMode, __NetworkDarc},
lamellar_team::IntoLamellarTeam,
LamellarEnv, LamellarTeam,
};
use super::handle::LocalRwDarcHandle;
pub(crate) use super::handle::{
IntoDarcHandle, IntoGlobalRwDarcHandle, LocalRwDarcReadHandle, LocalRwDarcWriteHandle,
};
#[derive(Debug)]
pub struct LocalRwDarcReadGuard<T: 'static> {
pub(crate) _darc: LocalRwDarc<T>,
pub(crate) lock: RwLockReadGuardArc<T>,
}
impl<T: fmt::Display> fmt::Display for LocalRwDarcReadGuard<T> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
fmt::Display::fmt(&self.lock, f)
}
}
impl<T> std::ops::Deref for LocalRwDarcReadGuard<T> {
type Target = T;
fn deref(&self) -> &T {
&self.lock
}
}
#[derive(Debug)]
pub struct LocalRwDarcWriteGuard<T: 'static> {
pub(crate) _darc: LocalRwDarc<T>,
pub(crate) lock: RwLockWriteGuardArc<T>,
}
impl<T: fmt::Display> fmt::Display for LocalRwDarcWriteGuard<T> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
fmt::Display::fmt(&self.lock, f)
}
}
impl<T> std::ops::Deref for LocalRwDarcWriteGuard<T> {
type Target = T;
fn deref(&self) -> &T {
&self.lock
}
}
impl<T> std::ops::DerefMut for LocalRwDarcWriteGuard<T> {
fn deref_mut(&mut self) -> &mut T {
&mut self.lock
}
}
#[derive(serde::Serialize, serde::Deserialize, Debug)]
pub struct LocalRwDarc<T: 'static> {
#[serde(
serialize_with = "localrw_serialize2",
deserialize_with = "localrw_from_ndarc2"
)]
pub(crate) darc: Darc<Arc<RwLock<T>>>, }
unsafe impl<T: Send> Send for LocalRwDarc<T> {} unsafe impl<T: Send> Sync for LocalRwDarc<T> {}
impl<T> LamellarEnv for LocalRwDarc<T> {
fn my_pe(&self) -> usize {
self.darc.my_pe()
}
fn num_pes(&self) -> usize {
self.darc.num_pes()
}
fn num_threads_per_pe(&self) -> usize {
self.darc.num_threads_per_pe()
}
fn world(&self) -> Arc<LamellarTeam> {
self.darc.world()
}
fn team(&self) -> Arc<LamellarTeam> {
self.darc.team().team()
}
}
impl<T> crate::active_messaging::DarcSerde for LocalRwDarc<T> {
fn ser(&self, num_pes: usize, darcs: &mut Vec<RemotePtr>) {
self.darc.serialize_update_cnts(num_pes);
darcs.push(RemotePtr::NetworkDarc(self.darc.clone().into()));
}
}
impl<T> LocalRwDarc<T> {
pub(crate) fn inner(&self) -> &DarcInner<Arc<RwLock<T>>> {
self.darc.inner()
}
#[doc(hidden)]
pub fn serialize_update_cnts(&self, cnt: usize) {
self.inner()
.dist_cnt
.fetch_add(cnt, std::sync::atomic::Ordering::SeqCst);
}
#[doc(hidden)]
pub fn deserialize_update_cnts(&self) {
tracing::trace!(
"localrwdarc[{:?}] deserialize_update_cnts {:?}",
self.darc.id,
self.inner()
);
self.inner().inc_pe_ref_count(self.darc.src_pe, 1); self.inner().local_cnt.fetch_add(1, Ordering::SeqCst);
}
#[doc(hidden)]
pub fn print(&self) {
println!(
"--------\norig: {:?} 0x{:x} {:?}\n--------",
self.darc.src_pe,
self.darc.inner.addr(),
self.inner()
);
}
}
impl<T> Drop for LocalRwDarc<T> {
fn drop(&mut self) {
tracing::trace!(target: "drop", "drop LocalRwDarc");
}
}
impl<T: Sync + Send> LocalRwDarc<T> {
#[doc(alias("One-sided", "onesided"))]
pub fn read(&self) -> LocalRwDarcReadHandle<T> {
LocalRwDarcReadHandle::new(self.clone())
}
#[doc(alias("One-sided", "onesided"))]
pub fn write(&self) -> LocalRwDarcWriteHandle<T> {
LocalRwDarcWriteHandle::new(self.clone())
}
#[doc(alias = "Collective")]
pub fn new<U: Into<IntoLamellarTeam>>(team: U, item: T) -> LocalRwDarcHandle<T> {
let team = team.into().team.clone();
LocalRwDarcHandle {
team: team.clone(),
launched: false,
creation_future: Box::pin(Darc::async_try_new_with_drop(
team,
Arc::new(RwLock::new(item)),
DarcMode::LocalRw,
None,
)),
}
}
#[doc(alias = "Collective")]
pub fn into_globalrw(self) -> IntoGlobalRwDarcHandle<T> {
let inner = self.darc.inner.clone();
let team = self.darc.inner().darc_rt_team();
IntoGlobalRwDarcHandle {
darc: self.into(),
team,
launched: false,
outstanding_future: Box::pin(DarcInner::block_on_outstanding(
inner,
DarcMode::GlobalRw,
0,
)),
}
}
}
impl<T: Send + Sync> LocalRwDarc<T> {
#[doc(alias = "Collective")]
pub fn into_darc(self) -> IntoDarcHandle<T> {
let inner = self.darc.inner.clone();
let team = self.darc.inner().darc_rt_team();
IntoDarcHandle {
darc: self.into(),
team,
launched: false,
outstanding_future: Box::pin(async move {
DarcInner::block_on_outstanding(inner, DarcMode::Darc, 0).await;
}),
}
}
}
impl<T> Clone for LocalRwDarc<T> {
fn clone(&self) -> Self {
tracing::trace!(
"LocalRwDarc[{:?}] Clone {:?}",
self.darc.id,
self.darc.inner()
);
LocalRwDarc {
darc: self.darc.clone(),
}
}
}
impl<T: fmt::Display + Sync + Send> fmt::Display for LocalRwDarc<T> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let lock: LocalRwDarc<T> = self.clone();
fmt::Display::fmt(&lock.read().block(), f)
}
}
pub(crate) fn localrw_serialize2<S, T>(
localrw: &Darc<Arc<RwLock<T>>>,
s: S,
) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
let ndarc = __NetworkDarc::from(localrw);
ndarc.serialize(s)
}
pub(crate) fn localrw_from_ndarc2<'de, D, T>(
deserializer: D,
) -> Result<Darc<Arc<RwLock<T>>>, D::Error>
where
D: Deserializer<'de>,
{
tracing::trace!("lrwdarc2 from net darc");
let ndarc: __NetworkDarc = Deserialize::deserialize(deserializer)?;
Ok(Darc::from(ndarc))
}