use std::fmt;
use std::sync::Arc;
use super::{Model, Transport};
use crate::error::ProviderError;
use crate::observe::AdapterContext;
use crate::streaming::Streamed;
use crate::wasm_compat::{WasmCompatSend, WasmCompatSync};
use crate::wire::{Capabilities, Descriptor, Mode, Operation, Wire};
pub(crate) trait ErasedModel<Op: Operation>: WasmCompatSend + WasmCompatSync {
fn describe(&self) -> Descriptor<'_>;
fn open(
&self,
request: Op::Request,
mode: Mode,
observation: Option<AdapterContext>,
) -> Result<Streamed<Op>, ProviderError>;
}
impl<W, T> ErasedModel<W::Op> for Model<W, T>
where
W: Wire,
T: Transport<W>,
{
fn describe(&self) -> Descriptor<'_> {
self.wire.describe()
}
fn open(
&self,
request: <W::Op as Operation>::Request,
mode: Mode,
observation: Option<AdapterContext>,
) -> Result<Streamed<W::Op>, ProviderError> {
Model::open(self, request, mode, observation)
}
}
pub struct DynModel<Op: Operation> {
inner: Arc<dyn ErasedModel<Op>>,
}
impl<W, T> Model<W, T>
where
W: Wire,
T: Transport<W>,
{
pub fn erase(self) -> DynModel<W::Op> {
DynModel {
inner: Arc::new(self),
}
}
}
impl<W, T> From<Model<W, T>> for DynModel<W::Op>
where
W: Wire,
T: Transport<W>,
{
fn from(model: Model<W, T>) -> Self {
model.erase()
}
}
impl<Op: Operation> Clone for DynModel<Op> {
fn clone(&self) -> Self {
Self {
inner: Arc::clone(&self.inner),
}
}
}
impl<Op: Operation> fmt::Debug for DynModel<Op> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("DynModel")
.field("name", &self.name())
.field("id", &self.id())
.finish()
}
}
impl<Op: Operation> DynModel<Op> {
pub fn name(&self) -> &str {
self.inner.describe().name
}
pub fn id(&self) -> Option<&str> {
self.inner.describe().model
}
pub fn capabilities(&self) -> Capabilities {
self.inner.describe().capabilities
}
pub fn call(
&self,
request: impl Into<Op::Request>,
) -> impl Future<Output = Result<Op::Response, ProviderError>> + WasmCompatSend + 'static {
self.finished(request.into(), None)
}
pub fn call_observed(
&self,
request: impl Into<Op::Request>,
observation: AdapterContext,
) -> impl Future<Output = Result<Op::Response, ProviderError>> + WasmCompatSend + 'static {
self.finished(request.into(), Some(observation))
}
fn finished(
&self,
request: Op::Request,
observation: Option<AdapterContext>,
) -> impl Future<Output = Result<Op::Response, ProviderError>> + WasmCompatSend + 'static {
let inner = self.inner.clone();
async move {
inner
.open(request, Mode::Unary, observation)?
.finish()
.await
}
}
pub fn stream(&self, request: impl Into<Op::Request>) -> Result<Streamed<Op>, ProviderError> {
self.inner.open(request.into(), Mode::Streaming, None)
}
pub fn stream_observed(
&self,
request: impl Into<Op::Request>,
observation: AdapterContext,
) -> Result<Streamed<Op>, ProviderError> {
self.inner
.open(request.into(), Mode::Streaming, Some(observation))
}
}
#[cfg(test)]
mod tests;