use std::collections::BTreeMap;
use std::sync::{Arc, Mutex, MutexGuard};
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, RaiseError, RawOccurrence, RawSample, ReadError,
ScannableSignals, SendError, ServeError, SettleError, SignalReader, SignalWriter,
SubscribeError, Watermark, WriteError,
};
use ridl_rt::sample::Timestamp;
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 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 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);
}
}
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) {
lock(&self.shared).close_source(self.id);
}
}
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)
}
}
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]) -> Result<(), RaiseError> {
let seq = self.take_seq((iface, ord));
lock(&self.shared).raise(iface, ord, bytes, seq);
Ok(())
}
}
pub struct CallerHandle {
shared: Shared,
catalog: CatalogRef,
next_seq: u64,
}
impl CallerHandle {
pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
CallerHandle {
shared,
catalog,
next_seq: 0,
}
}
fn take_seq(&mut self) -> u64 {
self.next_seq += 1;
self.next_seq
}
}
impl Attached for CallerHandle {
fn catalog(&self) -> &CatalogRef {
&self.catalog
}
}
impl Caller for CallerHandle {
fn command(
&mut self,
iface: InterfaceNo,
ord: Ordinal,
args: &[u8],
) -> Result<Correlation, SendError> {
let seq = self.take_seq();
Ok(lock(&self.shared).send(CallKind::Command, iface, ord, args, seq))
}
fn query(
&mut self,
iface: InterfaceNo,
ord: Ordinal,
args: &[u8],
) -> Result<Correlation, SendError> {
let seq = self.take_seq();
Ok(lock(&self.shared).send(CallKind::Query, iface, ord, args, seq))
}
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) {
lock(&self.shared).forget(c);
}
}
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 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));
}
}
Ok(())
}
fn next_claim(&mut self, out: &mut [u8]) -> Result<Option<Claim>, ReadError> {
let served = if self.served.is_empty() {
None
} else {
Some(self.served.as_slice())
};
lock(&self.shared).next_claim(self.id, served, out)
}
fn settle(
&mut self,
claim: ClaimId,
outcome: Result<&[u8], CallError>,
) -> Result<(), SettleError> {
lock(&self.shared).settle(self.id, claim, outcome)
}
}