#![allow(unused_imports, clippy::ptr_arg, clippy::needless_lifetimes)]
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use std::{borrow::Cow, io::Write, string::ToString};
use wasmbus_rpc::{
deserialize, serialize, Context, Message, MessageDispatch, RpcError, RpcResult, SendOpts,
Timestamp, Transport,
};
pub const SMITHY_VERSION: &str = "1.0";
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
pub struct PubMessage {
#[serde(default)]
pub subject: String,
#[serde(rename = "replyTo")]
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reply_to: Option<String>,
#[serde(with = "serde_bytes")]
#[serde(default)]
pub body: Vec<u8>,
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
pub struct ReplyMessage {
#[serde(default)]
pub subject: String,
#[serde(rename = "replyTo")]
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reply_to: Option<String>,
#[serde(with = "serde_bytes")]
#[serde(default)]
pub body: Vec<u8>,
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
pub struct RequestMessage {
#[serde(default)]
pub subject: String,
#[serde(with = "serde_bytes")]
#[serde(default)]
pub body: Vec<u8>,
#[serde(rename = "timeoutMs")]
pub timeout_ms: u32,
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
pub struct SubMessage {
#[serde(default)]
pub subject: String,
#[serde(rename = "replyTo")]
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reply_to: Option<String>,
#[serde(with = "serde_bytes")]
#[serde(default)]
pub body: Vec<u8>,
}
#[async_trait]
pub trait Messaging {
fn contract_id() -> &'static str {
"wasmcloud:messaging"
}
async fn publish(&self, ctx: &Context, arg: &PubMessage) -> RpcResult<()>;
async fn request(&self, ctx: &Context, arg: &RequestMessage) -> RpcResult<ReplyMessage>;
}
#[doc(hidden)]
#[async_trait]
pub trait MessagingReceiver: MessageDispatch + Messaging {
async fn dispatch(&self, ctx: &Context, message: &Message<'_>) -> RpcResult<Message<'_>> {
match message.method {
"Publish" => {
let value: PubMessage = deserialize(message.arg.as_ref())
.map_err(|e| RpcError::Deser(format!("message '{}': {}", message.method, e)))?;
let _resp = Messaging::publish(self, ctx, &value).await?;
let buf = Vec::new();
Ok(Message {
method: "Messaging.Publish",
arg: Cow::Owned(buf),
})
}
"Request" => {
let value: RequestMessage = deserialize(message.arg.as_ref())
.map_err(|e| RpcError::Deser(format!("message '{}': {}", message.method, e)))?;
let resp = Messaging::request(self, ctx, &value).await?;
let buf = serialize(&resp)?;
Ok(Message {
method: "Messaging.Request",
arg: Cow::Owned(buf),
})
}
_ => Err(RpcError::MethodNotHandled(format!(
"Messaging::{}",
message.method
))),
}
}
}
#[derive(Debug)]
pub struct MessagingSender<T: Transport> {
transport: T,
}
impl<T: Transport> MessagingSender<T> {
pub fn via(transport: T) -> Self {
Self { transport }
}
pub fn set_timeout(&self, interval: std::time::Duration) {
self.transport.set_timeout(interval);
}
}
#[cfg(target_arch = "wasm32")]
impl MessagingSender<wasmbus_rpc::actor::prelude::WasmHost> {
pub fn new() -> Self {
let transport =
wasmbus_rpc::actor::prelude::WasmHost::to_provider("wasmcloud:messaging", "default")
.unwrap();
Self { transport }
}
pub fn new_with_link(link_name: &str) -> wasmbus_rpc::RpcResult<Self> {
let transport =
wasmbus_rpc::actor::prelude::WasmHost::to_provider("wasmcloud:messaging", link_name)?;
Ok(Self { transport })
}
}
#[async_trait]
impl<T: Transport + std::marker::Sync + std::marker::Send> Messaging for MessagingSender<T> {
#[allow(unused)]
async fn publish(&self, ctx: &Context, arg: &PubMessage) -> RpcResult<()> {
let buf = serialize(arg)?;
let resp = self
.transport
.send(
ctx,
Message {
method: "Messaging.Publish",
arg: Cow::Borrowed(&buf),
},
None,
)
.await?;
Ok(())
}
#[allow(unused)]
async fn request(&self, ctx: &Context, arg: &RequestMessage) -> RpcResult<ReplyMessage> {
let buf = serialize(arg)?;
let resp = self
.transport
.send(
ctx,
Message {
method: "Messaging.Request",
arg: Cow::Borrowed(&buf),
},
None,
)
.await?;
let value = deserialize(&resp)
.map_err(|e| RpcError::Deser(format!("response to {}: {}", "Request", e)))?;
Ok(value)
}
}
#[async_trait]
pub trait MessageSubscriber {
fn contract_id() -> &'static str {
"wasmcloud:messaging"
}
async fn handle_message(&self, ctx: &Context, arg: &SubMessage) -> RpcResult<()>;
}
#[doc(hidden)]
#[async_trait]
pub trait MessageSubscriberReceiver: MessageDispatch + MessageSubscriber {
async fn dispatch(&self, ctx: &Context, message: &Message<'_>) -> RpcResult<Message<'_>> {
match message.method {
"HandleMessage" => {
let value: SubMessage = deserialize(message.arg.as_ref())
.map_err(|e| RpcError::Deser(format!("message '{}': {}", message.method, e)))?;
let _resp = MessageSubscriber::handle_message(self, ctx, &value).await?;
let buf = Vec::new();
Ok(Message {
method: "MessageSubscriber.HandleMessage",
arg: Cow::Owned(buf),
})
}
_ => Err(RpcError::MethodNotHandled(format!(
"MessageSubscriber::{}",
message.method
))),
}
}
}
#[derive(Debug)]
pub struct MessageSubscriberSender<T: Transport> {
transport: T,
}
impl<T: Transport> MessageSubscriberSender<T> {
pub fn via(transport: T) -> Self {
Self { transport }
}
pub fn set_timeout(&self, interval: std::time::Duration) {
self.transport.set_timeout(interval);
}
}
#[cfg(not(target_arch = "wasm32"))]
impl<'send> MessageSubscriberSender<wasmbus_rpc::provider::ProviderTransport<'send>> {
pub fn for_actor(ld: &'send wasmbus_rpc::core::LinkDefinition) -> Self {
Self {
transport: wasmbus_rpc::provider::ProviderTransport::new(ld, None),
}
}
}
#[cfg(target_arch = "wasm32")]
impl MessageSubscriberSender<wasmbus_rpc::actor::prelude::WasmHost> {
pub fn to_actor(actor_id: &str) -> Self {
let transport =
wasmbus_rpc::actor::prelude::WasmHost::to_actor(actor_id.to_string()).unwrap();
Self { transport }
}
}
#[async_trait]
impl<T: Transport + std::marker::Sync + std::marker::Send> MessageSubscriber
for MessageSubscriberSender<T>
{
#[allow(unused)]
async fn handle_message(&self, ctx: &Context, arg: &SubMessage) -> RpcResult<()> {
let buf = serialize(arg)?;
let resp = self
.transport
.send(
ctx,
Message {
method: "MessageSubscriber.HandleMessage",
arg: Cow::Borrowed(&buf),
},
None,
)
.await?;
Ok(())
}
}