use std::sync::{Arc, Mutex};
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::{Duration, Timestamp};
use ridl_rt::trace::TraceContext;
mod handle;
mod store;
pub use handle::{
CallerHandle, HandlerHandle, ReaderHandle, SinkHandle, SourceHandle, WriterHandle,
};
use handle::{Shared, lock};
use store::Store;
pub struct Handles {
pub reader: ReaderHandle,
pub writer: WriterHandle,
pub source: SourceHandle,
pub sink: SinkHandle,
pub caller: CallerHandle,
pub handler: HandlerHandle,
}
pub struct Loopback {
shared: Shared,
catalog: CatalogRef,
handles: Handles,
}
impl Loopback {
pub const SLOTS: usize = 16;
#[must_use]
pub fn new(catalog: CatalogRef) -> Self {
Loopback::over(Arc::new(Mutex::new(Store::new())), catalog)
}
fn over(shared: Shared, catalog: CatalogRef) -> Self {
let handles = Handles {
reader: ReaderHandle::new(Arc::clone(&shared), catalog),
writer: WriterHandle::new(Arc::clone(&shared), catalog),
source: SourceHandle::new(Arc::clone(&shared), catalog),
sink: SinkHandle::new(Arc::clone(&shared), catalog),
caller: CallerHandle::new(Arc::clone(&shared), catalog),
handler: HandlerHandle::new(Arc::clone(&shared), catalog),
};
Loopback {
shared,
catalog,
handles,
}
}
#[must_use]
pub fn split(self) -> Handles {
self.handles
}
#[must_use]
pub fn attach(&self) -> Loopback {
Loopback::over(Arc::clone(&self.shared), self.catalog)
}
#[must_use]
pub fn reader(&self) -> ReaderHandle {
ReaderHandle::new(Arc::clone(&self.shared), self.catalog)
}
#[must_use]
pub fn writer(&self) -> WriterHandle {
WriterHandle::new(Arc::clone(&self.shared), self.catalog)
}
#[must_use]
pub fn source(&self) -> SourceHandle {
SourceHandle::new(Arc::clone(&self.shared), self.catalog)
}
#[must_use]
pub fn sink(&self) -> SinkHandle {
SinkHandle::new(Arc::clone(&self.shared), self.catalog)
}
#[must_use]
pub fn caller(&self) -> CallerHandle {
CallerHandle::new(Arc::clone(&self.shared), self.catalog)
}
#[must_use]
pub fn handler(&self) -> HandlerHandle {
HandlerHandle::new(Arc::clone(&self.shared), self.catalog)
}
pub fn advance(&mut self, by: Duration) {
lock(&self.shared).advance(by);
}
pub fn provision_fixed(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) {
lock(&self.shared).provision_fixed(iface, ord, bytes);
}
pub fn fail_next_settle(&mut self) {
lock(&self.shared).fail_next_settle();
}
}
impl Attached for Loopback {
fn catalog(&self) -> &CatalogRef {
self.handles.reader.catalog()
}
}
impl Clock for Loopback {
fn now(&self) -> Timestamp {
self.handles.reader.now()
}
}
impl SignalReader for Loopback {
fn read(
&self,
iface: InterfaceNo,
ord: Ordinal,
out: &mut [u8],
) -> Result<RawSample, ReadError> {
self.handles.reader.read(iface, ord, out)
}
}
impl FixedReader for Loopback {
fn read_fixed(
&self,
iface: InterfaceNo,
ord: Ordinal,
out: &mut [u8],
) -> Result<usize, ReadError> {
self.handles.reader.read_fixed(iface, ord, out)
}
}
impl ScannableSignals for Loopback {
fn generation(&self, iface: InterfaceNo) -> u64 {
self.handles.reader.generation(iface)
}
fn scan(&self, marks: &mut [Watermark], out: &mut [Changed]) -> usize {
self.handles.reader.scan(marks, out)
}
}
impl CoherentSignals for Loopback {
fn read_coherent(
&self,
iface: InterfaceNo,
ords: &[Ordinal],
out: &mut [u8],
samples: &mut [RawSample],
) -> Result<usize, ReadError> {
self.handles.reader.read_coherent(iface, ords, out, samples)
}
}
impl SignalWriter for Loopback {
fn set(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) -> Result<(), WriteError> {
self.handles.writer.set(iface, ord, bytes)
}
fn invalidate(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
self.handles.writer.invalidate(iface, ord)
}
fn touch(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
self.handles.writer.touch(iface, ord)
}
fn commit(&mut self) {
self.handles.writer.commit();
}
}
impl EventSource for Loopback {
fn subscribe(&mut self, iface: InterfaceNo, ords: &[Ordinal]) -> Result<(), SubscribeError> {
self.handles.source.subscribe(iface, ords)
}
fn unsubscribe(&mut self, iface: InterfaceNo, ords: &[Ordinal]) {
self.handles.source.unsubscribe(iface, ords);
}
fn next(&mut self, out: &mut [u8]) -> Result<Option<RawOccurrence>, ReadError> {
self.handles.source.next(out)
}
}
impl EventSink for Loopback {
fn raise(
&mut self,
iface: InterfaceNo,
ord: Ordinal,
bytes: &[u8],
trace: Option<TraceContext>,
) -> Result<(), RaiseError> {
self.handles.sink.raise(iface, ord, bytes, trace)
}
}
impl Caller for Loopback {
fn command(
&mut self,
iface: InterfaceNo,
ord: Ordinal,
args: &[u8],
trace: Option<TraceContext>,
) -> Result<Correlation, SendError> {
self.handles.caller.command(iface, ord, args, trace)
}
fn query(
&mut self,
iface: InterfaceNo,
ord: Ordinal,
args: &[u8],
trace: Option<TraceContext>,
) -> Result<Correlation, SendError> {
self.handles.caller.query(iface, ord, args, trace)
}
fn ack(&mut self, c: Correlation) -> Option<Result<(), CallError>> {
self.handles.caller.ack(c)
}
fn reply(
&mut self,
c: Correlation,
out: &mut [u8],
) -> Result<Option<Result<usize, CallError>>, ReadError> {
self.handles.caller.reply(c, out)
}
fn forget(&mut self, c: Correlation) {
self.handles.caller.forget(c);
}
}
impl Handler for Loopback {
fn serve(&mut self, iface: InterfaceNo, ords: &[Ordinal]) -> Result<(), ServeError> {
self.handles.handler.serve(iface, ords)
}
fn next_claim(&mut self, out: &mut [u8]) -> Result<Option<Claim>, ReadError> {
self.handles.handler.next_claim(out)
}
fn settle(
&mut self,
claim: ClaimId,
outcome: Result<&[u8], CallError>,
) -> Result<(), SettleError> {
self.handles.handler.settle(claim, outcome)
}
}
impl Wakeable for Loopback {
fn wake_on(&self, what: Interest, waker: &Waker) {
match what {
Interest::Outcome(_) | Interest::Slot => self.handles.caller.wake_on(what, waker),
Interest::Event(_) => self.handles.source.wake_on(what, waker),
Interest::Claim(_) => self.handles.handler.wake_on(what, waker),
}
}
}
const _: () = {
const fn assert_sync<T: Sync>() {}
const fn assert_send<T: Send>() {}
assert_sync::<ReaderHandle>();
assert_send::<ReaderHandle>();
assert_send::<WriterHandle>();
assert_send::<SourceHandle>();
assert_send::<SinkHandle>();
assert_send::<CallerHandle>();
assert_send::<HandlerHandle>();
assert_send::<Loopback>();
};