use std::time::Duration;
use weida_core::{Error, Limits};
use crate::{
ClientTls, CursorLevel, CursorSet, Drained, Identity, IncomingMeta, ReportMode, Reported,
RuntimeConfig, ServerTls, TransferMeta, Trust,
};
pub struct Runtime {
inner: crate::Runtime,
}
impl Runtime {
pub fn new(config: RuntimeConfig) -> Result<Runtime, Error> {
outside_a_reactor()?;
Ok(Runtime {
inner: crate::Runtime::owned(config)?,
})
}
pub fn runtime(&self) -> &crate::Runtime {
&self.inner
}
pub fn requester(&self, trust: impl Into<ClientTls>) -> Requester {
Requester {
endpoint: self.inner.requester(trust),
}
}
pub fn pusher(&self, trust: impl Into<ClientTls>) -> Pusher {
Pusher {
endpoint: self.inner.pusher(trust),
}
}
pub fn subscriber(&self, trust: impl Into<ClientTls>) -> Subscriber {
Subscriber {
endpoint: self.inner.subscriber(trust),
}
}
pub fn pair(&self, trust: impl Into<ClientTls>) -> Paired {
Paired {
endpoint: self.inner.pair(trust),
}
}
pub fn surveyor(&self, trust: impl Into<ClientTls>) -> Surveyor {
Surveyor {
endpoint: self.inner.surveyor(trust),
}
}
pub fn bind_quic(
&self,
addr: std::net::SocketAddr,
tls: impl Into<ServerTls>,
) -> Result<Binding, Error> {
outside_a_reactor()?;
let listener = self.inner.listener();
let binding = drive(listener.bind_quic(addr, tls))?;
Ok(Binding { listener, binding })
}
pub fn limits(&self) -> &Limits {
&self.inner.config().limits
}
pub fn drain(self, deadline: Duration) -> Result<Drained, Error> {
outside_a_reactor()?;
Ok(drive(self.inner.drain(deadline)))
}
pub fn shutdown(self) -> Result<(), Error> {
outside_a_reactor()?;
drive(self.inner.shutdown());
Ok(())
}
}
pub struct Binding {
listener: crate::Listener,
binding: crate::Binding,
}
impl Binding {
pub fn local_addr(&self) -> std::net::SocketAddr {
self.binding.local_addr()
}
pub fn replier(&self, path: &str) -> Result<Replier, Error> {
Ok(Replier {
endpoint: self.listener.replier(path)?,
})
}
pub fn puller(&self, path: &str) -> Result<Puller, Error> {
Ok(Puller {
endpoint: self.listener.puller(path)?,
})
}
pub fn publisher(&self, path: &str) -> Result<Publisher, Error> {
Ok(Publisher {
endpoint: self.listener.publisher(path)?,
})
}
pub fn pair(&self, path: &str) -> Result<Paired, Error> {
Ok(Paired {
endpoint: self.listener.pair(path)?,
})
}
pub fn respondent(&self, path: &str) -> Result<Respondent, Error> {
Ok(Respondent {
endpoint: self.listener.respondent(path)?,
})
}
pub fn bus(&self, path: &str, trust: impl Into<ClientTls>) -> Result<BusMember, Error> {
Ok(BusMember {
endpoint: self.listener.bus(path, trust)?,
})
}
pub fn listener(&self) -> &crate::Listener {
&self.listener
}
}
#[derive(Clone, Debug)]
pub struct Message {
pub payload: Vec<u8>,
pub meta: IncomingMeta,
}
pub struct Cursors {
inner: crate::Cursors,
}
impl Cursors {
pub fn snapshot(&self) -> CursorSet {
self.inner.snapshot()
}
pub fn offset(&self, level: CursorLevel) -> Option<u64> {
self.inner.offset(level)
}
pub fn changed(&mut self, deadline: Duration) -> Result<Reported, Error> {
outside_a_reactor()?;
Ok(drive(self.inner.changed_within(deadline)))
}
pub fn cursors(&mut self) -> &mut crate::Cursors {
&mut self.inner
}
}
pub struct Reporter {
inner: crate::Reporter,
}
impl Reporter {
pub fn levels(&self) -> &[CursorLevel] {
self.inner.levels()
}
pub fn mode(&self) -> ReportMode {
self.inner.mode()
}
pub fn report(&mut self, level: CursorLevel, offset: u64) -> Result<(), Error> {
outside_a_reactor()?;
drive(self.inner.report(level, offset))
}
pub fn finish(self) -> Result<(), Error> {
outside_a_reactor()?;
drive(self.inner.finish())
}
}
fn drive<F: Future>(future: F) -> F::Output {
futures::executor::block_on(future)
}
fn outside_a_reactor() -> Result<(), Error> {
match tokio::runtime::Handle::try_current() {
Err(_) => Ok(()),
Ok(_) => Err(Error::Runtime(
"weida::blocking blocks the calling thread and was called from inside a Tokio \
runtime, which would deadlock: use the asynchronous API there, or call this from \
a thread the runtime does not own"
.to_owned(),
)),
}
}
macro_rules! dialling {
($name:ident, $inner:ty) => {
impl $name {
pub fn connect(&self, url: &str) -> Result<(), Error> {
outside_a_reactor()?;
drive(self.endpoint.connect(url))
}
pub fn endpoint(&self) -> &$inner {
&self.endpoint
}
}
};
}
pub struct Requester {
endpoint: crate::Requester,
}
dialling!(Requester, crate::Requester);
impl Requester {
pub fn request(&self, body: &[u8], max_reply_bytes: usize) -> Result<Vec<u8>, Error> {
self.request_with(TransferMeta::default(), body, max_reply_bytes)
}
pub fn request_with(
&self,
meta: TransferMeta,
body: &[u8],
max_reply_bytes: usize,
) -> Result<Vec<u8>, Error> {
outside_a_reactor()?;
drive(async {
let reply = self.endpoint.request_with(meta, body).await?;
reply.collect(max_reply_bytes).await
})
}
}
pub struct Pusher {
endpoint: crate::Pusher,
}
dialling!(Pusher, crate::Pusher);
impl Pusher {
pub fn send(&self, body: &[u8]) -> Result<(), Error> {
self.send_with(
TransferMeta::default().with_content_len(body.len() as u64),
body,
)
}
pub fn send_with(&self, meta: TransferMeta, body: &[u8]) -> Result<(), Error> {
outside_a_reactor()?;
drive(async {
let mut transfer = self.endpoint.open(meta).await?;
transfer.write_all(body).await?;
transfer.finish()?.delivered().await
})
}
pub fn send_reporting(
&self,
meta: TransferMeta,
body: &[u8],
) -> Result<Option<Cursors>, Error> {
outside_a_reactor()?;
drive(async {
let mut transfer = self.endpoint.open(meta).await?;
let cursors = transfer.cursors();
transfer.write_all(body).await?;
transfer.finish()?.delivered().await?;
Ok(cursors.map(|inner| Cursors { inner }))
})
}
}
pub struct Subscriber {
endpoint: crate::Subscriber,
}
dialling!(Subscriber, crate::Subscriber);
impl Subscriber {
pub fn subscribe(&self, filter: &str) -> Result<(), Error> {
outside_a_reactor()?;
drive(self.endpoint.subscribe(filter))
}
pub fn unsubscribe(&self, filter: &str) -> Result<(), Error> {
outside_a_reactor()?;
drive(self.endpoint.unsubscribe(filter))
}
pub fn recv(&self, max_bytes: usize) -> Result<Message, Error> {
outside_a_reactor()?;
drive(async {
let transfer = self.endpoint.recv().await?;
let meta = transfer.meta().clone();
let payload = transfer.collect(max_bytes).await?;
Ok(Message { payload, meta })
})
}
}
pub struct Replier {
endpoint: crate::Replier,
}
impl Replier {
pub fn path(&self) -> &str {
self.endpoint.path()
}
pub fn accept(&self, max_bytes: usize) -> Result<Request, Error> {
outside_a_reactor()?;
drive(async {
let mut request = self.endpoint.accept().await?;
let meta = request.meta().clone();
let payload = request.take_body().collect(max_bytes).await?;
Ok(Request {
request,
message: Message { payload, meta },
})
})
}
pub fn endpoint(&self) -> &crate::Replier {
&self.endpoint
}
}
pub struct Request {
request: crate::IncomingRequest,
message: Message,
}
impl Request {
pub fn message(&self) -> &Message {
&self.message
}
pub fn into_payload(self) -> Vec<u8> {
self.message.payload
}
pub fn reply(self, body: &[u8]) -> Result<(), Error> {
self.reply_with(
TransferMeta::default().with_content_len(body.len() as u64),
body,
)
}
pub fn reply_with(self, meta: TransferMeta, body: &[u8]) -> Result<(), Error> {
outside_a_reactor()?;
drive(async {
let mut out = self.request.reply(meta).await?;
out.write_all(body).await?;
out.finish()?;
Ok(())
})
}
pub fn refuse(self, code: weida_core::ErrorCode) -> Result<(), Error> {
outside_a_reactor()?;
drive(self.request.refuse(code));
Ok(())
}
pub fn into_inner(self) -> crate::IncomingRequest {
self.request
}
}
pub struct Puller {
endpoint: crate::Puller,
}
impl Puller {
pub fn path(&self) -> &str {
self.endpoint.path()
}
pub fn recv(&self, max_bytes: usize) -> Result<Message, Error> {
outside_a_reactor()?;
drive(async {
let transfer = self.endpoint.recv().await?;
let meta = transfer.meta().clone();
let payload = transfer.collect(max_bytes).await?;
Ok(Message { payload, meta })
})
}
pub fn recv_reporting(&self, max_bytes: usize) -> Result<(Message, Option<Reporter>), Error> {
outside_a_reactor()?;
drive(async {
let transfer = self.endpoint.recv().await?;
let meta = transfer.meta().clone();
let reporter = transfer.reporter().map(|inner| Reporter { inner });
let payload = transfer.collect(max_bytes).await?;
Ok((Message { payload, meta }, reporter))
})
}
pub fn endpoint(&self) -> &crate::Puller {
&self.endpoint
}
}
pub struct Publisher {
endpoint: crate::Publisher,
}
impl Publisher {
pub fn path(&self) -> &str {
self.endpoint.path()
}
pub fn publish(&self, topic: &str, payload: impl Into<bytes::Bytes>) -> Result<usize, Error> {
self.endpoint.publish(topic, payload)
}
pub fn subscriber_count(&self) -> usize {
self.endpoint.subscriber_count()
}
pub fn dropped(&self) -> u64 {
self.endpoint.dropped()
}
pub fn endpoint(&self) -> &crate::Publisher {
&self.endpoint
}
}
pub struct Paired {
endpoint: crate::Paired,
}
dialling!(Paired, crate::Paired);
impl Paired {
pub fn path(&self) -> &str {
self.endpoint.path()
}
pub fn peer_count(&self) -> usize {
self.endpoint.peer_count()
}
pub fn send(&self, body: &[u8]) -> Result<(), Error> {
self.send_with(
TransferMeta::default().with_content_len(body.len() as u64),
body,
)
}
pub fn send_with(&self, meta: TransferMeta, body: &[u8]) -> Result<(), Error> {
outside_a_reactor()?;
drive(async {
let mut transfer = self.endpoint.open(meta).await?;
transfer.write_all(body).await?;
transfer.finish()?.delivered().await
})
}
pub fn send_reporting(
&self,
meta: TransferMeta,
body: &[u8],
) -> Result<Option<Cursors>, Error> {
outside_a_reactor()?;
drive(async {
let mut transfer = self.endpoint.open(meta).await?;
let cursors = transfer.cursors();
transfer.write_all(body).await?;
transfer.finish()?.delivered().await?;
Ok(cursors.map(|inner| Cursors { inner }))
})
}
pub fn recv(&self, max_bytes: usize) -> Result<Message, Error> {
outside_a_reactor()?;
drive(async {
let transfer = self.endpoint.recv().await?;
let meta = transfer.meta().clone();
let payload = transfer.collect(max_bytes).await?;
Ok(Message { payload, meta })
})
}
pub fn recv_reporting(&self, max_bytes: usize) -> Result<(Message, Option<Reporter>), Error> {
outside_a_reactor()?;
drive(async {
let transfer = self.endpoint.recv().await?;
let meta = transfer.meta().clone();
let reporter = transfer.reporter().map(|inner| Reporter { inner });
let payload = transfer.collect(max_bytes).await?;
Ok((Message { payload, meta }, reporter))
})
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct Survey {
pub replies: Vec<Vec<u8>>,
pub asked: usize,
pub failed: usize,
pub late: u64,
}
impl Survey {
pub fn silent(&self) -> usize {
self.asked
.saturating_sub(self.replies.len())
.saturating_sub(self.failed)
}
}
pub struct Surveyor {
endpoint: crate::Surveyor,
}
dialling!(Surveyor, crate::Surveyor);
impl Surveyor {
pub fn peer_count(&self) -> usize {
self.endpoint.peer_count()
}
pub fn survey(
&self,
body: &[u8],
deadline: Duration,
max_reply_bytes: usize,
) -> Result<Survey, Error> {
self.survey_with(TransferMeta::default(), body, deadline, max_reply_bytes)
}
pub fn survey_with(
&self,
meta: TransferMeta,
body: &[u8],
deadline: Duration,
max_reply_bytes: usize,
) -> Result<Survey, Error> {
outside_a_reactor()?;
drive(async {
let mut run = self.endpoint.survey_with(meta, body, deadline).await?;
let asked = run.respondents();
let mut replies = Vec::new();
let mut failed = 0usize;
while let Some(answer) = run.next(max_reply_bytes).await {
match answer {
Ok(reply) => replies.push(reply),
Err(_) => failed += 1,
}
}
Ok(Survey {
replies,
asked,
failed,
late: run.late(),
})
})
}
}
pub struct Respondent {
endpoint: crate::Respondent,
}
impl Respondent {
pub fn path(&self) -> &str {
self.endpoint.path()
}
pub fn accept(&self, max_bytes: usize) -> Result<Request, Error> {
outside_a_reactor()?;
drive(async {
let mut request = self.endpoint.accept().await?;
let meta = request.meta().clone();
let payload = request.take_body().collect(max_bytes).await?;
Ok(Request {
request,
message: Message { payload, meta },
})
})
}
pub fn endpoint(&self) -> &crate::Respondent {
&self.endpoint
}
}
pub struct BusMember {
endpoint: crate::BusMember,
}
dialling!(BusMember, crate::BusMember);
impl BusMember {
pub fn path(&self) -> &str {
self.endpoint.path()
}
pub fn send(&self, body: &[u8]) -> Result<usize, Error> {
outside_a_reactor()?;
drive(self.endpoint.send(body))
}
pub fn send_with(&self, meta: TransferMeta, body: &[u8]) -> Result<usize, Error> {
outside_a_reactor()?;
drive(self.endpoint.send_with(meta, body))
}
pub fn recv(&self, max_bytes: usize) -> Result<Message, Error> {
outside_a_reactor()?;
drive(async {
let transfer = self.endpoint.recv().await?;
let meta = transfer.meta().clone();
let payload = transfer.collect(max_bytes).await?;
Ok(Message { payload, meta })
})
}
pub fn peer_count(&self) -> usize {
self.endpoint.peer_count()
}
pub fn dropped(&self) -> u64 {
self.endpoint.dropped()
}
}
#[cfg(feature = "generate")]
pub fn identity() -> Result<Identity, Error> {
Identity::generate()
}
pub fn trust_by_address() -> Trust {
Trust::by_address()
}