use std::collections::HashSet;
use std::pin::Pin;
use std::sync::Arc;
use std::time::{Duration, Instant};
use bytes::Bytes;
use tokio::sync::{Mutex, Semaphore, mpsc};
use weida_core::{Error, TraceContext};
use weida_protocol::{CreditHeader, FrameKind, SubscriptionHeader};
use crate::config::ClientTls;
use crate::conn::{ConnHandle, Ctl, write_control};
use crate::listener::{PairOwner, Route};
use crate::pubsub::{FanOut, SubRegistry};
use crate::reconnect::PeerEvents;
use crate::runtime::RuntimeInner;
use crate::stream::{Attach, Peer};
use crate::transfer::{
IncomingRequest, IncomingTransfer, OutgoingTransfer, ReplyStream, TransferMeta,
};
mod sealed {
pub trait Sealed {}
}
pub trait Pattern: sealed::Sealed {
type State;
}
pub enum Req {}
pub enum Rep {}
pub enum Push {}
pub enum Pull {}
pub enum Pub {}
pub enum Sub {}
pub enum Pair {}
pub enum Survey {}
pub enum Respond {}
pub enum Bus {}
impl sealed::Sealed for Req {}
impl Pattern for Req {
type State = ReqState;
}
impl sealed::Sealed for Rep {}
impl Pattern for Rep {
type State = RepState;
}
impl sealed::Sealed for Push {}
impl Pattern for Push {
type State = PushState;
}
impl sealed::Sealed for Pull {}
impl Pattern for Pull {
type State = PullState;
}
impl sealed::Sealed for Pub {}
impl Pattern for Pub {
type State = PubState;
}
impl sealed::Sealed for Sub {}
impl Pattern for Sub {
type State = SubState;
}
impl sealed::Sealed for Pair {}
impl Pattern for Pair {
type State = PairState;
}
impl sealed::Sealed for Survey {}
impl Pattern for Survey {
type State = SurveyState;
}
impl sealed::Sealed for Respond {}
impl Pattern for Respond {
type State = RespondState;
}
impl sealed::Sealed for Bus {}
impl Pattern for Bus {
type State = BusState;
}
pub struct Endpoint<P: Pattern> {
state: P::State,
}
impl<P: Pattern> Endpoint<P> {
pub(crate) fn from_state(state: P::State) -> Endpoint<P> {
Endpoint { state }
}
}
pub type Requester = Endpoint<Req>;
pub type Replier = Endpoint<Rep>;
pub type Pusher = Endpoint<Push>;
pub type Puller = Endpoint<Pull>;
pub type Publisher = Endpoint<Pub>;
pub type Subscriber = Endpoint<Sub>;
pub type Paired = Endpoint<Pair>;
pub type BusMember = Endpoint<Bus>;
pub struct ReqState {
peer: Peer,
}
pub type Surveyor = Endpoint<Survey>;
pub type Respondent = Endpoint<Respond>;
impl ReqState {
pub(crate) fn new(runtime: Arc<RuntimeInner>, tls: Arc<ClientTls>) -> ReqState {
ReqState {
peer: Peer::new(runtime, tls),
}
}
}
impl Requester {
pub async fn connect(&self, url: &str) -> Result<(), Error> {
self.state.peer.connect(url).await
}
pub fn disconnect(&self, url: &str) -> bool {
self.state.peer.disconnect(url)
}
pub fn peer_count(&self) -> usize {
self.state.peer.peer_count()
}
pub fn events(&self) -> PeerEvents {
self.state.peer.events()
}
pub async fn open(&self, meta: TransferMeta) -> Result<(OutgoingTransfer, ReplyStream), Error> {
self.state.peer.open_bi(meta).await
}
pub async fn request(&self, body: &[u8]) -> Result<IncomingTransfer, Error> {
self.request_with(TransferMeta::default(), body).await
}
pub async fn request_with(
&self,
meta: TransferMeta,
body: &[u8],
) -> Result<IncomingTransfer, Error> {
let (mut transfer, reply) = self.open(meta).await?;
transfer.write_all(body).await?;
transfer.finish()?;
reply.recv().await
}
}
pub struct RepState {
path: Arc<str>,
queue: Mutex<mpsc::Receiver<IncomingRequest>>,
}
impl RepState {
pub(crate) fn new(path: &str, queue: mpsc::Receiver<IncomingRequest>) -> RepState {
RepState {
path: Arc::from(path),
queue: Mutex::new(queue),
}
}
}
impl Replier {
pub fn path(&self) -> &str {
&self.state.path
}
pub async fn accept(&self) -> Result<IncomingRequest, Error> {
let mut queue = self.state.queue.lock().await;
queue.recv().await.ok_or(Error::NotConnected)
}
}
pub struct PushState {
peer: Peer,
}
impl PushState {
pub(crate) fn new(runtime: Arc<RuntimeInner>, tls: Arc<ClientTls>) -> PushState {
PushState {
peer: Peer::new(runtime, tls),
}
}
}
impl Pusher {
pub async fn connect(&self, url: &str) -> Result<(), Error> {
self.state.peer.connect(url).await
}
pub fn disconnect(&self, url: &str) -> bool {
self.state.peer.disconnect(url)
}
pub fn peer_count(&self) -> usize {
self.state.peer.peer_count()
}
pub fn events(&self) -> PeerEvents {
self.state.peer.events()
}
pub fn dropped(&self) -> u64 {
self.state.peer.dropped()
}
pub async fn open(&self, meta: TransferMeta) -> Result<OutgoingTransfer, Error> {
self.state.peer.open(meta).await
}
pub async fn send(&self, body: &[u8]) -> Result<(), Error> {
self.send_with(TransferMeta::default(), body).await
}
pub async fn send_with(&self, meta: TransferMeta, body: &[u8]) -> Result<(), Error> {
self.state.peer.send(meta, body).await
}
}
pub struct PullState {
path: Arc<str>,
queue: Mutex<mpsc::Receiver<IncomingTransfer>>,
}
impl PullState {
pub(crate) fn new(path: &str, queue: mpsc::Receiver<IncomingTransfer>) -> PullState {
PullState {
path: Arc::from(path),
queue: Mutex::new(queue),
}
}
}
impl Puller {
pub fn path(&self) -> &str {
&self.state.path
}
pub async fn recv(&self) -> Result<IncomingTransfer, Error> {
let mut queue = self.state.queue.lock().await;
queue.recv().await.ok_or(Error::NotConnected)
}
}
pub struct PubState {
path: Arc<str>,
registry: Arc<SubRegistry>,
max_payload: usize,
}
impl PubState {
pub(crate) fn new(path: &str, registry: Arc<SubRegistry>, max_payload: usize) -> PubState {
PubState {
path: Arc::from(path),
registry,
max_payload,
}
}
}
impl Publisher {
pub fn path(&self) -> &str {
&self.state.path
}
pub fn publish(&self, topic: &str, payload: impl Into<Bytes>) -> Result<usize, Error> {
self.publish_inner(topic, payload.into(), None)
}
pub fn publish_with_trace(
&self,
topic: &str,
payload: impl Into<Bytes>,
trace: TraceContext,
) -> Result<usize, Error> {
self.publish_inner(topic, payload.into(), Some(trace))
}
fn publish_inner(
&self,
topic: &str,
payload: Bytes,
trace: Option<TraceContext>,
) -> Result<usize, Error> {
if payload.len() > self.state.max_payload {
return Err(Error::LimitExceeded);
}
let want = u32::try_from(payload.len()).map_err(|_| Error::LimitExceeded)?;
Ok(self
.state
.registry
.publish(&self.state.path, topic, payload, trace, want))
}
pub fn open(&self, topic: &str) -> FanOut {
self.state.registry.open(&self.state.path, topic, None)
}
pub fn open_with_trace(&self, topic: &str, trace: TraceContext) -> FanOut {
self.state
.registry
.open(&self.state.path, topic, Some(trace))
}
pub fn subscriber_count(&self) -> usize {
self.state.registry.subscriber_count(&self.state.path)
}
pub fn filter_count(&self) -> usize {
self.state.registry.filter_count(&self.state.path)
}
pub fn dropped(&self) -> u64 {
self.state.registry.dropped(&self.state.path)
}
pub fn dropped_on(&self, topic: &str) -> Option<crate::TopicDrops> {
self.state.registry.dropped_on(&self.state.path, topic)
}
pub fn drops(&self) -> Vec<crate::TopicDrops> {
self.state.registry.drops(&self.state.path)
}
}
struct SubAttach {
filters: Arc<std::sync::Mutex<HashSet<String>>>,
queue_tx: mpsc::Sender<IncomingTransfer>,
}
impl Attach for SubAttach {
fn attach<'a>(
&'a self,
conn: &'a ConnHandle,
path: &'a str,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>> {
Box::pin(async move {
if conn.conn.needs_reverse_pool() && conn.conn.park_reverse().await? == 0 {
return Err(Error::Unsupported);
}
conn.namespace
.register(path, Route::Transfer(self.queue_tx.clone()))?;
if conn.conn.needs_reverse_pool() {
let maintaining = Arc::clone(conn);
conn.exec
.spawn(async move { maintaining.conn.maintain_reverse().await });
}
let filters: Vec<String> = self
.filters
.lock()
.expect("filter set poisoned")
.iter()
.cloned()
.collect();
for filter in filters {
send_subscription(conn, FrameKind::Subscribe, path, &filter).await?;
}
Ok(())
})
}
}
pub struct SubState {
peer: Peer,
filters: Arc<std::sync::Mutex<HashSet<String>>>,
queue: Mutex<mpsc::Receiver<IncomingTransfer>>,
}
impl SubState {
pub(crate) fn new(runtime: Arc<RuntimeInner>, tls: Arc<ClientTls>, depth: usize) -> SubState {
let (queue_tx, queue) = mpsc::channel(depth);
let filters = Arc::new(std::sync::Mutex::new(HashSet::new()));
let attach = SubAttach {
filters: Arc::clone(&filters),
queue_tx,
};
SubState {
peer: Peer::with_attach(runtime, tls, Some(Arc::new(attach))),
filters,
queue: Mutex::new(queue),
}
}
}
impl Subscriber {
pub async fn connect(&self, url: &str) -> Result<(), Error> {
self.state.peer.connect(url).await
}
pub fn disconnect(&self, url: &str) -> bool {
self.state.peer.disconnect(url)
}
pub fn peer_count(&self) -> usize {
self.state.peer.peer_count()
}
pub fn events(&self) -> PeerEvents {
self.state.peer.events()
}
pub async fn subscribe(&self, filter: &str) -> Result<(), Error> {
weida_protocol::filter::validate(filter)?;
let fresh = self
.state
.filters
.lock()
.expect("filter set poisoned")
.insert(filter.to_owned());
if !fresh {
return Ok(());
}
self.broadcast(FrameKind::Subscribe, filter).await
}
pub async fn unsubscribe(&self, filter: &str) -> Result<(), Error> {
let known = self
.state
.filters
.lock()
.expect("filter set poisoned")
.remove(filter);
if !known {
return Ok(());
}
self.broadcast(FrameKind::Unsubscribe, filter).await
}
pub fn filter_count(&self) -> usize {
self.state
.filters
.lock()
.expect("filter set poisoned")
.len()
}
pub async fn recv(&self) -> Result<IncomingTransfer, Error> {
let mut queue = self.state.queue.lock().await;
queue.recv().await.ok_or(Error::NotConnected)
}
pub async fn grant(&self, filter: &str, limit: u64) -> Result<(), Error> {
weida_protocol::filter::validate(filter)?;
for (conn, path) in self.state.peer.live_peers() {
let header = CreditHeader::new(path.as_ref(), filter, limit).encode();
write_control(&conn.conn, FrameKind::Credit, &header).await?;
}
Ok(())
}
async fn broadcast(&self, kind: FrameKind, filter: &str) -> Result<(), Error> {
for (conn, path) in self.state.peer.live_peers() {
send_subscription(&conn, kind, &path, filter).await?;
}
Ok(())
}
}
async fn send_subscription(
conn: &ConnHandle,
kind: FrameKind,
path: &str,
filter: &str,
) -> Result<(), Error> {
let header = SubscriptionHeader::new(path, filter).encode();
write_control(&conn.conn, kind, &header).await
}
impl Drop for SubState {
fn drop(&mut self) {
let filters: Vec<String> = self
.filters
.lock()
.expect("filter set poisoned")
.iter()
.cloned()
.collect();
for (conn, path) in self.peer.live_peers() {
conn.namespace.unregister(&path);
for filter in &filters {
conn.notify(Ctl::SendUnsubscribe {
path: Arc::clone(&path),
filter: filter.clone(),
});
}
}
}
}
pub struct PairState {
peer: Option<Peer>,
owner: Option<Arc<PairOwner>>,
path: Arc<str>,
queue: Mutex<mpsc::Receiver<IncomingTransfer>>,
}
struct PairAttach {
queue_tx: mpsc::Sender<IncomingTransfer>,
}
impl Attach for PairAttach {
fn attach<'a>(
&'a self,
conn: &'a ConnHandle,
path: &'a str,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>> {
Box::pin(async move {
if conn.conn.needs_reverse_pool() && conn.conn.park_reverse().await? == 0 {
return Err(Error::Unsupported);
}
conn.namespace.register(
path,
Route::Pair {
queue: self.queue_tx.clone(),
owner: Arc::new(PairOwner::new()),
},
)?;
if conn.conn.needs_reverse_pool() {
let maintaining = ConnHandle::clone(conn);
conn.exec
.spawn(async move { maintaining.conn.maintain_reverse().await });
}
Ok(())
})
}
}
impl PairState {
pub(crate) fn bound(
path: &str,
owner: Arc<PairOwner>,
queue: mpsc::Receiver<IncomingTransfer>,
) -> PairState {
PairState {
peer: None,
owner: Some(owner),
path: Arc::from(path),
queue: Mutex::new(queue),
}
}
pub(crate) fn dialling(
runtime: Arc<RuntimeInner>,
tls: Arc<ClientTls>,
depth: usize,
) -> PairState {
let (queue_tx, queue) = mpsc::channel(depth);
let attach = PairAttach { queue_tx };
PairState {
peer: Some(Peer::with_attach(runtime, tls, Some(Arc::new(attach)))),
owner: None,
path: Arc::from(""),
queue: Mutex::new(queue),
}
}
}
impl Paired {
pub fn path(&self) -> &str {
&self.state.path
}
pub async fn connect(&self, url: &str) -> Result<(), Error> {
let Some(peer) = self.state.peer.as_ref() else {
return Err(Error::Unsupported);
};
if peer.slot_count() > 0 {
return Err(Error::LimitExceeded);
}
peer.connect(url).await
}
pub fn peer_count(&self) -> usize {
self.state.peer.as_ref().map_or(0, Peer::peer_count)
}
pub fn events(&self) -> Option<PeerEvents> {
self.state.peer.as_ref().map(Peer::events)
}
pub fn dropped(&self) -> u64 {
self.state.peer.as_ref().map_or(0, Peer::dropped)
}
pub async fn open(&self, meta: TransferMeta) -> Result<OutgoingTransfer, Error> {
if let Some(peer) = self.state.peer.as_ref() {
return peer.open(meta).await;
}
let owner = self
.state
.owner
.as_ref()
.expect("a pair is either dialling or bound");
let conn = owner.peer().await?;
crate::stream::open_transfer_on(&conn, &self.state.path, &meta).await
}
pub async fn send(&self, body: &[u8]) -> Result<(), Error> {
self.send_with(TransferMeta::default(), body).await
}
pub async fn send_with(&self, meta: TransferMeta, body: &[u8]) -> Result<(), Error> {
if let Some(peer) = self.state.peer.as_ref() {
return peer.send(meta, body).await;
}
let mut transfer = self.open(meta).await?;
transfer.write_all(body).await?;
transfer.finish()?;
Ok(())
}
pub async fn recv(&self) -> Result<IncomingTransfer, Error> {
let mut queue = self.state.queue.lock().await;
queue.recv().await.ok_or(Error::NotConnected)
}
}
impl Drop for PairState {
fn drop(&mut self) {
let Some(peer) = self.peer.as_ref() else {
return;
};
for (conn, path) in peer.live_peers() {
conn.namespace.unregister(&path);
}
}
}
pub struct SurveyState {
peer: Peer,
}
impl SurveyState {
pub(crate) fn new(runtime: Arc<RuntimeInner>, tls: Arc<ClientTls>) -> SurveyState {
SurveyState {
peer: Peer::new(runtime, tls),
}
}
}
pub struct RespondState {
path: Arc<str>,
queue: Mutex<mpsc::Receiver<IncomingRequest>>,
}
impl RespondState {
pub(crate) fn new(path: &str, queue: mpsc::Receiver<IncomingRequest>) -> RespondState {
RespondState {
path: Arc::from(path),
queue: Mutex::new(queue),
}
}
}
impl Surveyor {
pub async fn connect(&self, url: &str) -> Result<(), Error> {
self.state.peer.connect(url).await
}
pub fn peer_count(&self) -> usize {
self.state.peer.peer_count()
}
pub async fn survey(&self, body: &[u8], deadline: Duration) -> Result<SurveyRun, Error> {
self.survey_with(TransferMeta::default(), body, deadline)
.await
}
pub async fn survey_with(
&self,
meta: TransferMeta,
body: &[u8],
deadline: Duration,
) -> Result<SurveyRun, Error> {
let targets = self.state.peer.live_peers();
let exec = self.state.peer.exec().clone();
let late = Arc::new(std::sync::atomic::AtomicU64::new(0));
let expires = Instant::now() + deadline;
let (tx, rx) = mpsc::channel(targets.len().max(1));
let mut respondents = 0usize;
let mut collecting = Vec::with_capacity(targets.len());
for (conn, path) in targets {
let remaining = expires.saturating_duration_since(Instant::now());
let asked = exec.within(remaining, ask(&conn, &path, &meta, body)).await;
let Some(asked) = asked else {
tracing::debug!(%path, "the survey deadline passed before this respondent was asked");
break;
};
let reply = match asked {
Ok(reply) => reply,
Err(e) => {
tracing::debug!(error = %e, %path, "a respondent could not be asked");
continue;
}
};
respondents += 1;
let tx = tx.clone();
let late = Arc::clone(&late);
collecting.push(exec.spawn(async move {
let answer = reply.recv().await;
if Instant::now() >= expires || tx.send(answer).await.is_err() {
late.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
}));
}
drop(tx);
Ok(SurveyRun {
rx,
exec,
expires,
late,
respondents,
collecting,
done: false,
})
}
}
async fn ask(
conn: &ConnHandle,
path: &str,
meta: &TransferMeta,
body: &[u8],
) -> Result<ReplyStream, Error> {
let (mut request, reply) = crate::stream::open_exchange_on(conn, path, meta.clone()).await?;
request.write_all(body).await?;
request.finish()?;
Ok(reply)
}
pub struct SurveyRun {
rx: mpsc::Receiver<Result<IncomingTransfer, Error>>,
exec: crate::runtime::Exec,
expires: Instant,
late: Arc<std::sync::atomic::AtomicU64>,
respondents: usize,
collecting: Vec<tokio::task::JoinHandle<()>>,
done: bool,
}
impl SurveyRun {
pub async fn next(&mut self, max_bytes: usize) -> Option<Result<Vec<u8>, Error>> {
if self.done {
return None;
}
match self.rx.try_recv() {
Ok(Ok(transfer)) => return Some(transfer.collect(max_bytes).await),
Ok(Err(e)) => return Some(Err(e)),
Err(mpsc::error::TryRecvError::Disconnected) => {
self.done = true;
return None;
}
Err(mpsc::error::TryRecvError::Empty) => {}
}
let remaining = self.expires.saturating_duration_since(Instant::now());
match self.exec.within(remaining, self.rx.recv()).await {
None => {
self.done = true;
while self.rx.try_recv().is_ok() {
self.late.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
None
}
Some(None) => {
self.done = true;
None
}
Some(Some(Ok(transfer))) => Some(transfer.collect(max_bytes).await),
Some(Some(Err(e))) => Some(Err(e)),
}
}
pub fn late(&self) -> u64 {
self.late.load(std::sync::atomic::Ordering::Relaxed)
}
pub fn respondents(&self) -> usize {
self.respondents
}
}
impl Drop for SurveyRun {
fn drop(&mut self) {
for handle in &self.collecting {
handle.abort();
}
}
}
impl std::fmt::Debug for SurveyRun {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SurveyRun")
.field("respondents", &self.respondents)
.field("late", &self.late())
.finish_non_exhaustive()
}
}
impl Respondent {
pub fn path(&self) -> &str {
&self.state.path
}
pub async fn accept(&self) -> Result<IncomingRequest, Error> {
let mut queue = self.state.queue.lock().await;
queue.recv().await.ok_or(Error::NotConnected)
}
}
pub struct BusState {
peer: Peer,
path: Arc<str>,
queue: Mutex<mpsc::Receiver<IncomingTransfer>>,
writers: std::sync::Mutex<Vec<BusWriter>>,
dropped: Arc<std::sync::atomic::AtomicU64>,
depth: usize,
budget: usize,
}
struct BusWriter {
tx: mpsc::Sender<BusMsg>,
budget: Arc<Semaphore>,
}
struct BusMsg {
meta: TransferMeta,
body: Bytes,
}
impl BusState {
pub(crate) fn new(
path: &str,
runtime: Arc<RuntimeInner>,
tls: Arc<ClientTls>,
queue: mpsc::Receiver<IncomingTransfer>,
depth: usize,
budget: usize,
) -> BusState {
BusState {
peer: Peer::new(runtime, tls),
path: Arc::from(path),
queue: Mutex::new(queue),
writers: std::sync::Mutex::new(Vec::new()),
dropped: Arc::new(std::sync::atomic::AtomicU64::new(0)),
depth,
budget,
}
}
}
impl BusMember {
pub fn path(&self) -> &str {
&self.state.path
}
pub async fn connect(&self, url: &str) -> Result<(), Error> {
let (conn, path) = self.state.peer.dial(url).await?;
let (tx, rx) = mpsc::channel(self.state.depth);
let budget = Arc::new(Semaphore::new(self.state.budget));
conn.exec.spawn(bus_writer(
conn.clone(),
path,
rx,
Arc::clone(&budget),
Arc::clone(&self.state.dropped),
));
self.state
.writers
.lock()
.expect("bus writers poisoned")
.push(BusWriter { tx, budget });
Ok(())
}
pub fn peer_count(&self) -> usize {
self.state.peer.peer_count()
}
pub async fn send(&self, body: &[u8]) -> Result<usize, Error> {
self.send_with(TransferMeta::default(), body).await
}
pub async fn send_with(&self, meta: TransferMeta, body: &[u8]) -> Result<usize, Error> {
if body.len() > self.state.budget {
return Err(Error::LimitExceeded);
}
let body = Bytes::copy_from_slice(body);
let want = body.len() as u32;
let mut reached = 0usize;
let mut dropped = 0u64;
{
let mut writers = self.state.writers.lock().expect("bus writers poisoned");
writers.retain(|writer| !writer.tx.is_closed());
for writer in writers.iter() {
let Ok(permit) = writer.budget.try_acquire_many(want) else {
dropped += 1;
continue;
};
let msg = BusMsg {
meta: meta.clone(),
body: body.clone(),
};
match writer.tx.try_send(msg) {
Ok(()) => {
permit.forget();
reached += 1;
}
Err(_) => dropped += 1,
}
}
}
if dropped > 0 {
tracing::debug!(dropped, "a bus member could not take a copy");
self.state
.dropped
.fetch_add(dropped, std::sync::atomic::Ordering::Relaxed);
}
Ok(reached)
}
pub async fn recv(&self) -> Result<IncomingTransfer, Error> {
let mut queue = self.state.queue.lock().await;
queue.recv().await.ok_or(Error::NotConnected)
}
pub fn dropped(&self) -> u64 {
self.state
.dropped
.load(std::sync::atomic::Ordering::Relaxed)
}
}
async fn bus_writer(
conn: ConnHandle,
path: Arc<str>,
mut rx: mpsc::Receiver<BusMsg>,
budget: Arc<Semaphore>,
dropped: Arc<std::sync::atomic::AtomicU64>,
) {
loop {
let msg = tokio::select! {
msg = rx.recv() => match msg {
Some(msg) => msg,
None => return,
},
_ = conn.conn.closed() => {
dropped.fetch_add(rx.len() as u64, std::sync::atomic::Ordering::Relaxed);
return;
}
};
let len = msg.body.len();
let outcome = send_one(&conn, &path, msg.meta, &msg.body).await;
budget.add_permits(len);
if let Err(e) = outcome {
tracing::debug!(error = %e, %path, "a bus write failed; copy dropped");
dropped.fetch_add(1 + rx.len() as u64, std::sync::atomic::Ordering::Relaxed);
return;
}
}
}
async fn send_one(
conn: &ConnHandle,
path: &str,
meta: TransferMeta,
body: &[u8],
) -> Result<(), Error> {
let mut transfer = crate::stream::open_transfer_on(conn, path, &meta).await?;
transfer.write_all(body).await?;
transfer.finish()?;
Ok(())
}