Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions sdk/core/azure_core_amqp/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ typespec = { path = "../typespec", version = "1.2.0-beta.1" }
typespec_macros = { path = "../typespec_macros", version = "1.1.0-beta.1" }

[dev-dependencies]
fe2o3-amqp = { workspace = true, features = ["acceptor"] }
serde_json.workspace = true
tracing-subscriber = { workspace = true, features = ["env-filter"] }

Expand All @@ -51,6 +52,10 @@ fe2o3_amqp = [
"serde_bytes",
"azure_core/tokio",
]
transaction = [
"fe2o3-amqp?/transaction",
"fe2o3-amqp-types?/transaction",
]

[lints]
workspace = true
Expand Down
11 changes: 11 additions & 0 deletions sdk/core/azure_core_amqp/src/error/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,10 @@ pub enum AmqpErrorKind {
/// Transport Implementation Error
TransportImplementationError(Box<dyn std::error::Error + Send + Sync>),

#[cfg(feature = "transaction")]
/// Transactions are not supported by the remote peer.
TransactionsNotSupported(Box<dyn std::error::Error + Send + Sync>),

/// A send was rejected.
SendRejected,
}
Expand Down Expand Up @@ -193,6 +197,9 @@ impl std::error::Error for AmqpError {
| AmqpErrorKind::FramingError(e)
| AmqpErrorKind::IdleTimeoutElapsed(e) => Some(e.as_ref()),

#[cfg(feature = "transaction")]
AmqpErrorKind::TransactionsNotSupported(e) => Some(e.as_ref()),

AmqpErrorKind::ManagementStatusCode(_, _)
| AmqpErrorKind::NonTerminalDeliveryState
| AmqpErrorKind::SimpleMessage(_)
Expand Down Expand Up @@ -248,6 +255,10 @@ impl std::fmt::Display for AmqpError {
AmqpErrorKind::TransportImplementationError(s) => {
write!(f, "Transport Implementation Error: {}", s)
}
#[cfg(feature = "transaction")]
AmqpErrorKind::TransactionsNotSupported(s) => {
write!(f, "Transactions not supported by the remote peer: {}", s)
}
AmqpErrorKind::ConnectionDropped(s) => {
write!(f, "Connection dropped: {}", s)
}
Expand Down
6 changes: 6 additions & 0 deletions sdk/core/azure_core_amqp/src/fe2o3/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,12 @@ impl From<&fe2o3_amqp_types::definitions::ErrorCondition> for AmqpErrorCondition
fe2o3_amqp_types::definitions::ErrorCondition::Custom(symbol) => {
AmqpErrorCondition::from(AmqpSymbol::from(symbol))
}
#[cfg(feature = "transaction")]
fe2o3_amqp_types::definitions::ErrorCondition::TransactionError(
ref transaction_error,
) => AmqpErrorCondition::from(AmqpSymbol::from(
fe2o3_amqp_types::primitives::Symbol::from(transaction_error),
)),
}
}
}
Expand Down
93 changes: 93 additions & 0 deletions sdk/core/azure_core_amqp/src/fe2o3/messaging/message_target.rs
Original file line number Diff line number Diff line change
@@ -1,2 +1,95 @@
// Copyright (c) Microsoft Corporation. All Rights reserved
// Licensed under the MIT license.

use crate::{messaging::AmqpTarget, value::AmqpValue};

pub(crate) fn to_fe2o3_target(target: AmqpTarget) -> fe2o3_amqp_types::messaging::Target {
let mut builder = fe2o3_amqp_types::messaging::Target::builder();

if let Some(address) = target.address {
builder = builder.address(address);
}
if let Some(durable) = target.durable {
builder = builder.durable(durable.into());
}
if let Some(expiry_policy) = target.expiry_policy {
builder = builder.expiry_policy(expiry_policy.into());
}
if let Some(timeout) = target.timeout {
builder = builder.timeout(timeout);
}
if let Some(dynamic) = target.dynamic {
builder = builder.dynamic(dynamic);
}
if let Some(dynamic_node_properties) = target.dynamic_node_properties {
builder = builder.dynamic_node_properties(
dynamic_node_properties
.iter()
.map(|(k, v)| {
(
fe2o3_amqp_types::primitives::Symbol::from(k.as_str()),
v.into(),
)
})
.collect::<fe2o3_amqp_types::definitions::Fields>(),
);
}
if let Some(capabilities) = target.capabilities {
builder = builder.capabilities(
capabilities
.into_iter()
.map(|v| match v {
AmqpValue::Symbol(s) => fe2o3_amqp_types::primitives::Symbol::from(s.0),
AmqpValue::String(s) => fe2o3_amqp_types::primitives::Symbol::from(s),
_ => fe2o3_amqp_types::primitives::Symbol::from(format!("{:?}", v)),
})
.collect::<fe2o3_amqp_types::primitives::Array<
fe2o3_amqp_types::primitives::Symbol,
>>(),
);
}
builder.build()
}

#[cfg(test)]
mod tests {
use super::*;
use crate::{
messaging::{TerminusDurability, TerminusExpiryPolicy},
value::{AmqpOrderedMap, AmqpSymbol},
};

#[test]
fn message_target_conversion_to_fe2o3() {
let mut dynamic_node_properties = AmqpOrderedMap::new();
dynamic_node_properties.insert("prop".to_string(), AmqpValue::from("val"));

let amqp_target = AmqpTarget::builder()
.with_address("test".to_string())
.with_durable(TerminusDurability::UnsettledState)
.with_expiry_policy(TerminusExpiryPolicy::SessionEnd)
.with_timeout(95)
.with_dynamic(false)
.with_dynamic_node_properties(dynamic_node_properties)
.with_capabilities(vec![AmqpValue::Symbol(AmqpSymbol::from("capability"))])
.build();

let fe2o3_target = to_fe2o3_target(amqp_target);

assert_eq!(fe2o3_target.address.unwrap().as_str(), "test");
assert_eq!(
fe2o3_target.durable,
fe2o3_amqp_types::messaging::TerminusDurability::UnsettledState
);
assert_eq!(
fe2o3_target.expiry_policy,
fe2o3_amqp_types::messaging::TerminusExpiryPolicy::SessionEnd
);
assert_eq!(fe2o3_target.timeout, 95);
assert_eq!(fe2o3_target.dynamic, false);
assert_eq!(
fe2o3_target.capabilities.unwrap().as_slice(),
&[fe2o3_amqp_types::primitives::Symbol::from("capability")]
);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,10 @@ impl From<fe2o3_amqp_types::messaging::Outcome> for AmqpOutcome {
fe2o3_amqp_types::messaging::Outcome::Released(_) => AmqpOutcome::Released,
fe2o3_amqp_types::messaging::Outcome::Rejected(_) => AmqpOutcome::Rejected,
fe2o3_amqp_types::messaging::Outcome::Modified(_) => AmqpOutcome::Modified,
#[cfg(feature = "transaction")]
fe2o3_amqp_types::messaging::Outcome::Declared(declared) => {
AmqpOutcome::Declared(declared.txn_id.into_vec())
}
}
}
}
Expand All @@ -98,18 +102,27 @@ impl From<AmqpOutcome> for fe2o3_amqp_types::messaging::Outcome {
message_annotations: None,
},
),
#[cfg(feature = "transaction")]
AmqpOutcome::Declared(txn_id) => fe2o3_amqp_types::messaging::Outcome::Declared(
fe2o3_amqp_types::transaction::Declared {
txn_id: serde_bytes::ByteBuf::from(txn_id),
},
),
}
}
}

#[test]
fn test_outcome_round_trip() {
let outcomes = vec![
#[cfg_attr(not(feature = "transaction"), allow(unused_mut))]
let mut outcomes = vec![
AmqpOutcome::Accepted,
AmqpOutcome::Released,
AmqpOutcome::Rejected,
AmqpOutcome::Modified,
];
#[cfg(feature = "transaction")]
outcomes.push(AmqpOutcome::Declared(vec![1, 2, 3, 4]));

for outcome in outcomes {
let fe2o3_outcome: fe2o3_amqp_types::messaging::Outcome = outcome.clone().into();
Expand Down
2 changes: 2 additions & 0 deletions sdk/core/azure_core_amqp/src/fe2o3/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,4 +8,6 @@ pub(crate) mod messaging;
pub(crate) mod receiver;
pub(crate) mod sender;
pub(crate) mod session;
#[cfg(feature = "transaction")]
pub(crate) mod transaction;
pub(crate) mod value;
120 changes: 115 additions & 5 deletions sdk/core/azure_core_amqp/src/fe2o3/receiver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,11 +13,20 @@ use std::sync::OnceLock;
use tokio::sync::Mutex;
use tracing::{info, trace, warn};

#[cfg(feature = "transaction")]
use crate::TransactionId;
#[cfg(feature = "transaction")]
use fe2o3_amqp::transaction::OwnedTransaction;
#[cfg(feature = "transaction")]
use std::{collections::HashMap, sync::Arc};

use super::error::Fe2o3ReceiverAttachError;

#[derive(Default)]
pub(crate) struct Fe2o3AmqpReceiver {
receiver: OnceLock<Mutex<fe2o3_amqp::Receiver>>,
#[cfg(feature = "transaction")]
transactions: OnceLock<Arc<Mutex<HashMap<TransactionId, OwnedTransaction>>>>,
}

/// The fe2o3 link builder for a receiver, after the name, source and target are set.
Expand Down Expand Up @@ -95,6 +104,8 @@ impl AmqpReceiverApis for Fe2o3AmqpReceiver {
self.receiver
.set(Mutex::new(receiver))
.map_err(|_| Self::could_not_set_message_receiver())?;
#[cfg(feature = "transaction")]
let _ = self.transactions.set(session.implementation.transactions());
Ok(())
}

Expand Down Expand Up @@ -130,7 +141,8 @@ impl AmqpReceiverApis for Fe2o3AmqpReceiver {

async fn credit_mode(&self) -> Result<ReceiverCreditMode> {
let receiver = self.receiver.get().ok_or_else(Self::receiver_not_set)?;
Ok(receiver.lock().await.credit_mode().into())
let guard = receiver.lock().await;
Ok(guard.credit_mode().into())
}

async fn receive_delivery(&self) -> Result<AmqpDelivery> {
Expand All @@ -141,10 +153,7 @@ impl AmqpReceiverApis for Fe2o3AmqpReceiver {
.lock()
.await;

let delivery: fe2o3_amqp::link::delivery::Delivery<
fe2o3_amqp_types::messaging::Body<fe2o3_amqp_types::primitives::Value>,
> = receiver.recv().await.map_err(AmqpError::from)?;
trace!("Received delivery: {:?}", delivery);
let delivery = receiver.recv().await.map_err(AmqpError::from)?;
Ok(delivery.into())
}

Expand Down Expand Up @@ -201,12 +210,47 @@ impl AmqpReceiverApis for Fe2o3AmqpReceiver {

Ok(())
}

#[cfg(feature = "transaction")]
async fn settle_with_transaction(
&self,
delivery: &AmqpDelivery,
outcome: crate::messaging::AmqpOutcome,
txn_id: crate::TransactionId,
) -> Result<()> {
use fe2o3_amqp::transaction::TransactionalRetirement;

let mut receiver = self
.receiver
.get()
.ok_or_else(Self::receiver_not_set)?
.lock()
.await;

let transactions = self.transactions.get().ok_or_else(Self::receiver_not_set)?;

trace!("Settling delivery with transaction.");
let fe2o3_outcome: fe2o3_amqp_types::messaging::Outcome = outcome.into();
let active_txns = transactions.lock().await;
let txn = active_txns.get(&txn_id).ok_or_else(|| {
AmqpError::with_message("Transaction not found or already discharged")
})?;

txn.retire(&mut receiver, &delivery.0.delivery, fe2o3_outcome)
.await
.map_err(AmqpError::from)?;
trace!("Settled delivery with transaction.");

Ok(())
}
}

impl Fe2o3AmqpReceiver {
pub fn new() -> Self {
Self {
receiver: OnceLock::new(),
#[cfg(feature = "transaction")]
transactions: OnceLock::new(),
}
}

Expand Down Expand Up @@ -276,6 +320,10 @@ mod tests {
use super::*;
use crate::messaging::AmqpSource;
use crate::value::{AmqpOrderedMap, AmqpSymbol, AmqpValue};
#[cfg(feature = "transaction")]
use fe2o3_amqp::{acceptor::ConnectionAcceptor, connection::Connection};
#[cfg(feature = "transaction")]
use tokio::{io::duplex, sync::oneshot};

// Makes sure the link properties survive the fe2o3 builder chain. The
// `name` method clears the properties, so it must run before
Expand Down Expand Up @@ -413,4 +461,66 @@ mod tests {
"the reported error must carry the LinkStateError, not the enclosing RecvError"
);
}

#[tokio::test]
#[cfg(feature = "transaction")]
async fn settle_with_transaction_unattached_receiver_fails() {
let (client_io, server_io) = duplex(65536);
let (done_tx, done_rx) = oneshot::channel::<()>();

let server_task = tokio::spawn(async move {
let acceptor = ConnectionAcceptor::new("test-container");
let mut connection = acceptor.accept(server_io).await.unwrap();
let session_acceptor = fe2o3_amqp::acceptor::SessionAcceptor::new();
let mut session = session_acceptor.accept(&mut connection).await.unwrap();
let link_acceptor = fe2o3_amqp::acceptor::LinkAcceptor::new();
let fe2o3_amqp::acceptor::LinkEndpoint::Sender(mut sender_link) =
link_acceptor.accept(&mut session).await.unwrap()
else {
panic!("expected Sender link");
};
let message = fe2o3_amqp_types::messaging::Message::builder()
.value(42i32)
.build();
tokio::spawn(async move {
let _ = sender_link.send(message).await;
});
let _ = done_rx.await;
});

let mut connection = Connection::builder()
.container_id("client-container")
.open_with_stream(client_io)
.await
.unwrap();
let mut session = fe2o3_amqp::session::Session::begin(&mut connection)
.await
.unwrap();
let mut receiver_link = fe2o3_amqp::Receiver::builder()
.name("test-receiver")
.source("test-source")
.credit_mode(fe2o3_amqp::link::receiver::CreditMode::Auto(10))
.attach(&mut session)
.await
.unwrap();

let fe2o3_delivery = receiver_link
.recv::<fe2o3_amqp_types::messaging::Body<fe2o3_amqp_types::primitives::Value>>()
.await
.unwrap();
let delivery: AmqpDelivery = fe2o3_delivery.into();

let receiver = Fe2o3AmqpReceiver::new();
let result = receiver
.settle_with_transaction(
&delivery,
crate::messaging::AmqpOutcome::Accepted,
vec![1, 2, 3],
)
.await;
assert!(result.is_err());

let _ = done_tx.send(());
let _ = server_task.await;
}
}
Loading
Loading