use std::collections::BTreeMap;
use std::sync::{Arc, Mutex, MutexGuard};
use std::task::Waker;
use ridl_rt::contract::{CatalogRef, InterfaceNo, Ordinal};
use ridl_rt::error::CallError;
use ridl_rt::port::{
Attached, Caller, Changed, Claim, ClaimId, Clock, CoherentSignals, Correlation, EventSink,
EventSource, FixedReader, Handler, Interest, RaiseError, RawOccurrence, RawSample, ReadError,
ScannableSignals, SendError, ServeError, SettleError, SignalReader, SignalWriter,
SubscribeError, Wakeable, Watermark, WriteError,
};
use ridl_rt::sample::Timestamp;
use ridl_rt::trace::TraceContext;
use crate::store::{CallKind, Key, Staged, Store};
pub(crate) type Shared = Arc<Mutex<Store>>;
pub(crate) fn lock(shared: &Shared) -> MutexGuard<'_, Store> {
shared
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
pub(crate) fn locked<R>(shared: &Shared, f: impl FnOnce(&mut Store, &mut Vec<Waker>) -> R) -> R {
let mut wake = Vec::new();
let result = {
let mut store = lock(shared);
f(&mut store, &mut wake)
};
for waker in wake {
waker.wake();
}
result
}
fn wake_at_once(waker: &Waker) {
waker.wake_by_ref();
}
pub struct ReaderHandle {
shared: Shared,
catalog: CatalogRef,
}
impl ReaderHandle {
pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
ReaderHandle { shared, catalog }
}
}
impl Attached for ReaderHandle {
fn catalog(&self) -> &CatalogRef {
&self.catalog
}
}
impl Clock for ReaderHandle {
fn now(&self) -> Timestamp {
lock(&self.shared).now()
}
}
impl Wakeable for ReaderHandle {
fn wake_on(&self, _: Interest, waker: &Waker) {
wake_at_once(waker);
}
}
impl SignalReader for ReaderHandle {
fn read(
&self,
iface: InterfaceNo,
ord: Ordinal,
out: &mut [u8],
) -> Result<RawSample, ReadError> {
lock(&self.shared).read(iface, ord, out)
}
}
impl FixedReader for ReaderHandle {
fn read_fixed(
&self,
iface: InterfaceNo,
ord: Ordinal,
out: &mut [u8],
) -> Result<usize, ReadError> {
lock(&self.shared).read_fixed(iface, ord, out)
}
}
impl ScannableSignals for ReaderHandle {
fn generation(&self, iface: InterfaceNo) -> u64 {
lock(&self.shared).generation(iface)
}
fn scan(&self, marks: &mut [Watermark], out: &mut [Changed]) -> usize {
lock(&self.shared).scan(marks, out)
}
}
impl CoherentSignals for ReaderHandle {
fn read_coherent(
&self,
iface: InterfaceNo,
ords: &[Ordinal],
out: &mut [u8],
samples: &mut [RawSample],
) -> Result<usize, ReadError> {
lock(&self.shared).read_coherent(iface, ords, out, samples)
}
}
pub struct WriterHandle {
shared: Shared,
catalog: CatalogRef,
staged: BTreeMap<Key, Staged>,
seqs: BTreeMap<Key, u64>,
}
impl WriterHandle {
pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
WriterHandle {
shared,
catalog,
staged: BTreeMap::new(),
seqs: BTreeMap::new(),
}
}
}
impl Attached for WriterHandle {
fn catalog(&self) -> &CatalogRef {
&self.catalog
}
}
impl SignalWriter for WriterHandle {
fn set(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) -> Result<(), WriteError> {
self.staged
.insert((iface, ord), Staged::Set(bytes.to_vec()));
Ok(())
}
fn invalidate(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
self.staged.insert((iface, ord), Staged::Invalidate);
Ok(())
}
fn touch(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
self.staged.entry((iface, ord)).or_insert(Staged::Touch);
Ok(())
}
fn commit(&mut self) {
lock(&self.shared).commit(&mut self.staged, &mut self.seqs);
}
}
impl Wakeable for WriterHandle {
fn wake_on(&self, _: Interest, waker: &Waker) {
wake_at_once(waker);
}
}
pub struct SourceHandle {
shared: Shared,
catalog: CatalogRef,
id: usize,
}
impl SourceHandle {
pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
let id = lock(&shared).open_source();
SourceHandle {
shared,
catalog,
id,
}
}
}
impl Drop for SourceHandle {
fn drop(&mut self) {
let waiters = lock(&self.shared).close_source(self.id);
drop(waiters);
}
}
impl Attached for SourceHandle {
fn catalog(&self) -> &CatalogRef {
&self.catalog
}
}
impl EventSource for SourceHandle {
fn subscribe(&mut self, iface: InterfaceNo, ords: &[Ordinal]) -> Result<(), SubscribeError> {
lock(&self.shared).subscribe(self.id, iface, ords);
Ok(())
}
fn unsubscribe(&mut self, iface: InterfaceNo, ords: &[Ordinal]) {
lock(&self.shared).unsubscribe(self.id, iface, ords);
}
fn next(&mut self, out: &mut [u8]) -> Result<Option<RawOccurrence>, ReadError> {
lock(&self.shared).next_event(self.id, out)
}
}
impl Wakeable for SourceHandle {
fn wake_on(&self, what: Interest, waker: &Waker) {
match what {
Interest::Event(_) => locked(&self.shared, |store, wake| {
store.wait_event(self.id, what, waker, wake);
}),
Interest::Outcome(_) | Interest::Slot | Interest::Claim(_) => wake_at_once(waker),
}
}
}
pub struct SinkHandle {
shared: Shared,
catalog: CatalogRef,
seqs: BTreeMap<Key, u64>,
}
impl SinkHandle {
pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
SinkHandle {
shared,
catalog,
seqs: BTreeMap::new(),
}
}
fn take_seq(&mut self, key: Key) -> u64 {
let counter = self.seqs.entry(key).or_insert(0);
*counter += 1;
*counter
}
}
impl Attached for SinkHandle {
fn catalog(&self) -> &CatalogRef {
&self.catalog
}
}
impl EventSink for SinkHandle {
fn raise(
&mut self,
iface: InterfaceNo,
ord: Ordinal,
bytes: &[u8],
trace: Option<TraceContext>,
) -> Result<(), RaiseError> {
let seq = self.take_seq((iface, ord));
locked(&self.shared, |store, wake| {
store.raise(iface, ord, bytes, seq, trace, wake);
});
Ok(())
}
}
impl Wakeable for SinkHandle {
fn wake_on(&self, _: Interest, waker: &Waker) {
wake_at_once(waker);
}
}
pub struct CallerHandle {
shared: Shared,
catalog: CatalogRef,
id: usize,
next_seq: u64,
}
impl CallerHandle {
pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
let id = lock(&shared).open_caller();
CallerHandle {
shared,
catalog,
id,
next_seq: 0,
}
}
fn send(
&mut self,
kind: CallKind,
iface: InterfaceNo,
ord: Ordinal,
args: &[u8],
trace: Option<TraceContext>,
) -> Result<Correlation, SendError> {
let seq = self.next_seq + 1;
let c = locked(&self.shared, |store, wake| {
store.send(self.id, kind, (iface, ord), args, seq, trace, wake)
})?;
self.next_seq = seq;
Ok(c)
}
}
impl Drop for CallerHandle {
fn drop(&mut self) {
let waiters = locked(&self.shared, |store, wake| {
store.close_caller(self.id, wake)
});
drop(waiters);
}
}
impl Attached for CallerHandle {
fn catalog(&self) -> &CatalogRef {
&self.catalog
}
}
impl Clock for CallerHandle {
fn now(&self) -> Timestamp {
lock(&self.shared).now()
}
}
impl Wakeable for CallerHandle {
fn wake_on(&self, what: Interest, waker: &Waker) {
match what {
Interest::Outcome(c) => locked(&self.shared, |store, wake| {
store.wait_outcome(c, waker, wake);
}),
Interest::Slot => locked(&self.shared, |store, wake| {
store.wait_slot(self.id, waker, wake);
}),
Interest::Event(_) | Interest::Claim(_) => wake_at_once(waker),
}
}
}
impl Caller for CallerHandle {
fn command(
&mut self,
iface: InterfaceNo,
ord: Ordinal,
args: &[u8],
trace: Option<TraceContext>,
) -> Result<Correlation, SendError> {
self.send(CallKind::Command, iface, ord, args, trace)
}
fn query(
&mut self,
iface: InterfaceNo,
ord: Ordinal,
args: &[u8],
trace: Option<TraceContext>,
) -> Result<Correlation, SendError> {
self.send(CallKind::Query, iface, ord, args, trace)
}
fn ack(&mut self, c: Correlation) -> Option<Result<(), CallError>> {
lock(&self.shared).ack(c)
}
fn reply(
&mut self,
c: Correlation,
out: &mut [u8],
) -> Result<Option<Result<usize, CallError>>, ReadError> {
lock(&self.shared).reply(c, out)
}
fn forget(&mut self, c: Correlation) {
locked(&self.shared, |store, wake| store.forget(c, wake));
}
}
pub struct HandlerHandle {
shared: Shared,
catalog: CatalogRef,
id: usize,
served: Vec<Key>,
}
impl HandlerHandle {
pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
let id = lock(&shared).open_handler();
HandlerHandle {
shared,
catalog,
id,
served: Vec::new(),
}
}
#[must_use]
pub fn served(&self) -> &[(InterfaceNo, Ordinal)] {
&self.served
}
}
impl Drop for HandlerHandle {
fn drop(&mut self) {
let waiters = locked(&self.shared, |store, wake| {
store.close_handler(self.id, wake)
});
drop(waiters);
}
}
impl Attached for HandlerHandle {
fn catalog(&self) -> &CatalogRef {
&self.catalog
}
}
impl Handler for HandlerHandle {
fn serve(&mut self, iface: InterfaceNo, ords: &[Ordinal]) -> Result<(), ServeError> {
for ord in ords {
if !self.served.contains(&(iface, *ord)) {
self.served.push((iface, *ord));
}
}
locked(&self.shared, |store, wake| {
store.serve(self.id, iface, ords, wake);
});
Ok(())
}
fn next_claim(&mut self, out: &mut [u8]) -> Result<Option<Claim>, ReadError> {
lock(&self.shared).next_claim(self.id, out)
}
fn settle(
&mut self,
claim: ClaimId,
outcome: Result<&[u8], CallError>,
) -> Result<(), SettleError> {
locked(&self.shared, |store, wake| {
store.settle(self.id, claim, outcome, wake)
})
}
}
impl Wakeable for HandlerHandle {
fn wake_on(&self, what: Interest, waker: &Waker) {
match what {
Interest::Claim(_) => locked(&self.shared, |store, wake| {
store.wait_claim(self.id, what, waker, wake);
}),
Interest::Outcome(_) | Interest::Slot | Interest::Event(_) => wake_at_once(waker),
}
}
}