Skip to main content

ConfirmsTransaction

Struct ConfirmsTransaction 

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

An owned confirm-transaction, opened by transaction on a ConfirmsPublisher.

A private publish buffer: publish appends to this value rather than to the publisher, commit flushes the whole buffer on the confirm channel and awaits every acknowledgement, and abort discards it without touching the broker. Both settle by consuming self, so a double commit or a publish after settling is a compile error.

Unlike the handle-level buffer of the borrowed kind, any number of these can be open on one publisher at a time, and the publisher keeps publishing directly while they are.

§Examples

use ruststream::{Broker, OutgoingMessage, OwnedTransactions, Transaction};
use ruststream_lapin::{LapinBroker, LapinPublish};

let connected = LapinBroker::new("amqp://localhost:5672").connect().await?;
let publisher = connected.publisher(LapinPublish::default().confirms());

let mut orders = publisher.transaction().await?;
let mut audit = publisher.transaction().await?; // concurrent with `orders`
orders.publish(OutgoingMessage::new("orders", b"{}".as_slice())).await?;
audit.publish(OutgoingMessage::new("audit", b"{}".as_slice())).await?;
orders.commit().await?;
audit.commit().await?;

Trait Implementations§

Source§

impl Debug for ConfirmsTransaction

Source§

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

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

impl Drop for ConfirmsTransaction

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more
Source§

impl Transaction for ConfirmsTransaction

Source§

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

Buffers msg in this transaction; nothing reaches the broker before commit.

§Errors

Infallible in practice: buffering is local to this value, and a closed connection or a rejected frame surfaces at the commit, which is the visibility point.

Source§

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

Publishes the buffered messages in order and awaits every confirm.

§Errors

Returns AmqpError::Closed once the broker has shut down, and AmqpError::Publish when a message fails to publish or the broker returns a negative confirm. A failed commit has still consumed the transaction and its buffer is lost: redelivery of the inputs, not resubmission of the buffer, is the recovery path. Messages already flushed stay published - publisher confirms give durability per message, not atomicity across them (use ServerTxPublish for that).

§Cancel safety

Not cancel safe: dropping the future mid-flush leaves an unknown prefix of the buffer published.

Source§

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

Discards the buffered messages without touching the broker.

§Errors

Never fails: nothing was staged on the broker.

Source§

type Error = AmqpError

The error type returned by transaction operations.

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> 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<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<T> Same for T

Source§

type Output = T

Should always be Self
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