Skip to main content

ServerTxPublisher

Struct ServerTxPublisher 

Source
pub struct ServerTxPublisher { /* private fields */ }
Expand description

The live publisher backed by AMQP server transactions (tx.select / tx.commit / tx.rollback).

Between begin_transaction and commit messages accumulate on the broker inside the channel transaction and become visible atomically at commit; abort rolls them back server-side. Outside a transaction publish behaves like the fire-and-forget publisher.

Only the borrowed transaction kind (TransactionalPublisher) applies here, unlike ConfirmsPublisher: tx.select puts the channel itself into transactional mode, so the transaction is channel state with exactly one instance, and there is no buffer for an owned Transaction value to own.

Clones share the transactional channel and its open/closed state. Interleaving publish and begin_transaction/commit from concurrent tasks is not supported: which side of the transaction boundary a concurrent publish lands on would be a race either way. Like every live publisher it aliases the connection and may outlive it: after shutdown every operation reports AmqpError::Closed.

Trait Implementations§

Source§

impl Clone for ServerTxPublisher

Source§

fn clone(&self) -> ServerTxPublisher

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for ServerTxPublisher

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Publisher for ServerTxPublisher

Source§

async fn publish(&self, msg: OutgoingMessage<'_>) -> Result<(), Self::Error>

Publishes msg: into the open server transaction, or plainly when none is open.

§Errors

Returns AmqpError::Closed once the broker has shut down and AmqpError::Publish when the channel rejects the frame.

§Cancel safety

Not cancel safe: dropping the future may leave the message queued in the transaction or not.

Source§

type Error = AmqpError

The error type returned by publish.
Source§

impl TransactionalPublisher for ServerTxPublisher

Source§

async fn begin_transaction(&self) -> Result<(), Self::Error>

Opens a server transaction (tx.select on first use).

§Errors

Returns AmqpError::Transaction when a transaction is already open on this handle (the open transaction is left untouched), AmqpError::Closed once the broker has shut down, and AmqpError::Publish when the transactional channel cannot be set up.

Source§

async fn commit(&self) -> Result<(), Self::Error>

Commits the open server transaction.

§Errors

Returns AmqpError::Transaction when no transaction is open, and AmqpError::Publish when tx.commit fails; the transaction state on the broker is then unknown (the channel may be closed) and the publisher should be discarded.

Source§

async fn abort(&self) -> Result<(), Self::Error>

Rolls back the open server transaction.

§Errors

Returns AmqpError::Transaction when no transaction is open, and AmqpError::Publish when tx.rollback fails.

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<'a, T, E> AsTaggedExplicit<'a, E> for T
where T: 'a,

Source§

fn explicit(self, class: Class, tag: u32) -> TaggedParser<'a, Explicit, Self, E>

Source§

impl<'a, T, E> AsTaggedImplicit<'a, E> for T
where T: 'a,

Source§

fn implicit( self, class: Class, constructed: bool, tag: u32, ) -> TaggedParser<'a, Implicit, Self, E>

Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> CompatExt for T

Source§

fn compat(self) -> Compat<T>
where T: Sized,

Applies the Compat adapter by value. Read more
Source§

fn compat_ref(&self) -> Compat<&T>

Applies the Compat adapter by shared reference. Read more
Source§

fn compat_mut(&mut self) -> Compat<&mut T>

Applies the Compat adapter by mutable reference. Read more
Source§

impl<P> ErasedPublisher for P
where P: Publisher,

Source§

fn publish_bytes<'a>( &'a self, name: &'a str, payload: &'a [u8], ) -> Pin<Box<dyn Future<Output = Result<(), Box<dyn Error + Sync + Send>>> + Send + 'a>>

Publishes payload to name, with no headers. Read more
Source§

fn publish_message<'a>( &'a self, name: &'a str, payload: &'a [u8], headers: &'a Headers, ) -> Pin<Box<dyn Future<Output = Result<(), Box<dyn Error + Sync + Send>>> + Send + 'a>>

Publishes payload to name with headers. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<Reply, DeliveryCx, Pipeline, Bare> ReplySink<Reply, DeliveryCx, Pipeline> for Bare
where Reply: AsRef<[u8]> + Sync, DeliveryCx: Sync, Pipeline: Send + Sync, Bare: Publisher,

Source§

type Error = <Bare as Publisher>::Error

The error surfaced when the reply cannot be published.
Source§

async fn deliver( &self, name: &str, reply: &Reply, _pipeline: &Pipeline, _cx: &PublishContext<'_, DeliveryCx>, ) -> Result<(), <Bare as ReplySink<Reply, DeliveryCx, Pipeline>>::Error>

Publishes reply to name.
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more