use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use fe2o3_amqp::connection::{Connection, ConnectionHandle};
use fe2o3_amqp::session::{Session, SessionHandle};
use fe2o3_amqp::{Receiver, Sender};
use fe2o3_amqp_types::messaging::Source;
use fe2o3_amqp_types::primitives::{Array, Symbol};
use ruststream::{Broker, ConnectedBroker, DefaultPublish, DescribeServer, ServerSpec, Subscribe};
use tokio::sync::{Mutex, OnceCell};
use crate::address::{AmqpAddress, Settle};
use crate::config::Sasl;
use crate::error::{AmqpError, box_err};
use crate::publisher::{AmqpPublish, AmqpPublisher};
use crate::subscriber::AmqpSubscriber;
pub(crate) struct AmqpCore {
pub(crate) conn: Mutex<ConnectionHandle<()>>,
pub(crate) session: Mutex<SessionHandle<()>>,
pub(crate) closed: AtomicBool,
pub(crate) container_id: String,
link_seq: AtomicU64,
}
impl AmqpCore {
pub(crate) fn ensure_open(&self) -> Result<(), AmqpError> {
if self.closed.load(Ordering::Acquire) {
return Err(AmqpError::NotConnected);
}
Ok(())
}
pub(crate) fn link_name(&self, role: &str) -> String {
let seq = self.link_seq.fetch_add(1, Ordering::Relaxed);
format!("{}-{role}-{seq}", self.container_id)
}
pub(crate) fn correlation_id(&self) -> String {
let seq = self.link_seq.fetch_add(1, Ordering::Relaxed);
format!("{}-corr-{seq}", self.container_id)
}
}
impl std::fmt::Debug for AmqpCore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AmqpCore")
.field("container_id", &self.container_id)
.field("closed", &self.closed.load(Ordering::Relaxed))
.finish_non_exhaustive()
}
}
pub(crate) type CoreCell = Arc<OnceCell<Arc<AmqpCore>>>;
#[derive(Debug, Clone)]
#[must_use]
pub struct AmqpBroker {
url: String,
container_id: Option<String>,
sasl: Option<Sasl>,
cell: CoreCell,
}
impl AmqpBroker {
pub fn new(url: impl Into<String>) -> Self {
Self {
url: url.into(),
container_id: None,
sasl: None,
cell: Arc::new(OnceCell::new()),
}
}
pub fn sasl(mut self, sasl: Sasl) -> Self {
self.sasl = Some(sasl);
self
}
pub fn container_id(mut self, id: impl Into<String>) -> Self {
self.container_id = Some(id.into());
self
}
#[must_use]
pub fn publisher(&self) -> AmqpPublisher {
AmqpPublisher::new(Arc::clone(&self.cell))
}
}
impl Broker for AmqpBroker {
type Error = AmqpError;
type Connected = ConnectedAmqpBroker;
async fn connect(self) -> Result<Self::Connected, Self::Error> {
let container_id = self
.container_id
.clone()
.unwrap_or_else(|| "ruststream".to_owned());
let core = self
.cell
.get_or_try_init(async || {
let mut builder = Connection::builder().container_id(container_id.clone());
if let Some(sasl) = &self.sasl {
builder = builder.sasl_profile(sasl.profile.clone());
}
let mut conn = builder
.open(self.url.as_str())
.await
.map_err(|e| AmqpError::Connect(box_err(e)))?;
let session = Session::begin(&mut conn)
.await
.map_err(|e| AmqpError::Session(box_err(e)))?;
Ok::<_, AmqpError>(Arc::new(AmqpCore {
conn: Mutex::new(conn),
session: Mutex::new(session),
closed: AtomicBool::new(false),
container_id,
link_seq: AtomicU64::new(0),
}))
})
.await?
.clone();
Ok(ConnectedAmqpBroker {
core,
cell: self.cell,
})
}
}
impl DescribeServer for AmqpBroker {
fn describe_server(&self) -> ServerSpec {
ServerSpec::new(
self.url
.trim_start_matches("amqps://")
.trim_start_matches("amqp://"),
"amqp",
)
}
}
#[derive(Debug)]
pub struct ConnectedAmqpBroker {
pub(crate) core: Arc<AmqpCore>,
cell: CoreCell,
}
impl ConnectedAmqpBroker {
#[must_use]
pub fn publisher(&self) -> AmqpPublisher {
AmqpPublisher::new(Arc::clone(&self.cell))
}
pub async fn subscribe_address(
&self,
address: AmqpAddress,
) -> Result<AmqpSubscriber, AmqpError> {
address.validate()?;
self.core.ensure_open()?;
let session = {
let mut conn = self.core.conn.lock().await;
Session::begin(&mut conn)
.await
.map_err(|e| AmqpError::Session(box_err(e)))?
};
let subscriber = AmqpSubscriber::attach(&self.core, session, address).await?;
Ok(subscriber)
}
pub(crate) async fn attach_sender(core: &AmqpCore, address: &str) -> Result<Sender, AmqpError> {
let mut session = core.session.lock().await;
Sender::attach(&mut session, core.link_name("sender"), address)
.await
.map_err(|e| AmqpError::Attach {
address: address.to_owned(),
source: box_err(e),
})
}
pub(crate) async fn attach_dynamic_receiver(core: &AmqpCore) -> Result<Receiver, AmqpError> {
let mut session = core.session.lock().await;
Receiver::builder()
.name(core.link_name("reply"))
.source(Source::builder().dynamic(true).build())
.auto_accept(true)
.attach(&mut session)
.await
.map_err(|e| AmqpError::Attach {
address: "(dynamic)".to_owned(),
source: box_err(e),
})
}
}
pub(crate) fn source_for(address: &AmqpAddress) -> Source {
let mut source = Source::builder().address(address.address()).build();
if let Some(capability) = address.capability() {
source.capabilities = Some(Array::from(vec![Symbol::from(capability)]));
}
source
}
impl ConnectedBroker for ConnectedAmqpBroker {
type Error = AmqpError;
type Closed = ();
async fn shutdown(self) -> Result<(), Self::Error> {
self.core.closed.store(true, Ordering::Release);
let session_result = {
let mut session = self.core.session.lock().await;
session.end().await
};
let conn_result = {
let mut conn = self.core.conn.lock().await;
conn.close().await
};
conn_result.map_err(|e| AmqpError::Connect(box_err(e)))?;
session_result.map_err(|e| AmqpError::Session(box_err(e)))?;
Ok(())
}
}
impl Subscribe for ConnectedAmqpBroker {
type Subscriber = AmqpSubscriber;
async fn subscribe(&self, name: &str) -> Result<Self::Subscriber, Self::Error> {
self.subscribe_address(AmqpAddress::raw(name)).await
}
}
impl DefaultPublish for ConnectedAmqpBroker {
type Policy = AmqpPublish;
}
pub(crate) fn is_at_most_once(settle: Settle) -> bool {
matches!(settle, Settle::AtMostOnce)
}