oa_gateway/adapters/
stomp.rsuse std::net::SocketAddr;
use std::sync::Arc;
use oa_gateway_adapter::tls::ClientTls;
use oa_gateway_adapter::Adapter;
use oa_gateway_core::Engine;
use oa_gateway_stomp::{StompAdapter, StompConfig};
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use tracing::{error, info};
use crate::config::StompSection;
pub(crate) fn start(
section: &StompSection,
broker: SocketAddr,
tls: Option<ClientTls>,
engine: Arc<Engine>,
shutdown: CancellationToken,
) -> Result<JoinHandle<()>, String> {
let login = if section.login.is_empty() {
None
} else {
Some(section.login.clone())
};
let passcode = if section.passcode.is_empty() {
None
} else {
Some(section.passcode.clone())
};
let tls_on = tls.is_some();
let adapter = Arc::new(StompAdapter::new(
section.id.clone(),
StompConfig {
broker,
host: section.host.clone(),
login,
passcode,
destination_prefix: section.destination_prefix.clone(),
topics: section.topics.clone(),
unwrap_ma_payloads: section.unwrap_ma_payloads,
reconnect: section.reconnect,
reconnect_delay: std::time::Duration::from_secs(section.reconnect_delay_secs),
connect_timeout: std::time::Duration::from_secs(section.connect_timeout_secs),
suppress_echo: section.suppress_echo,
on_panic: section.on_panic_mode()?,
max_frame_size: section.max_frame_size,
tls,
},
));
info!(id = %adapter.id(), %broker, tls = tls_on, "starting stomp adapter");
Ok(tokio::spawn(async move {
if let Err(err) = adapter.run(engine, shutdown).await {
error!(error = %err, "stomp adapter failed");
}
}))
}