mod config;
#[cfg(feature = "std")]
mod framing;
mod handles;
#[cfg(feature = "embedded")]
mod loopback;
mod participant;
mod protocol;
#[cfg(feature = "std")]
mod tcp;
pub mod websocket;
#[cfg(feature = "std")]
pub use tcp::{
DeliveredMessage, FlushMode, FlushOutcome, OBSERVABILITY_CHANNEL, PendingPushConnect,
PublishRejection, PushClient, PushWriter, PushedFrame, SubscriptionStream, TcpRemoteTransport,
};
#[cfg(feature = "std")]
pub use websocket::{
WebSocketDeliveredMessage, WebSocketRemoteTransport, WebSocketSubscriptionStream,
};
pub use config::{SdkConfig, build_channel_handle, build_conversation_handle};
pub use handles::{
RemoteChannelHandle, RemoteConversationHandle, RemoteParticipantHandle, SdkChannelHandle,
SdkConversationHandle,
};
pub use participant::{
PARTICIPANT_PUMP_WINDOW, ParticipantResponseProvenance, ParticipantResumeStore,
RemoteDetachReplayOutcome, RemoteExpectedOperationRecovery, RemoteLostOperationResolution,
RemoteLostReconnectResolution, RemoteOperationRecordOutcome, RemoteOperationTransportFate,
RemoteParticipantError, RemoteParticipantInbound, RemoteParticipantOperation,
RemoteParticipantSendOutcome, RemoteReconnectAttemptOutcome, RemoteReconnectPermit,
RemoteReconnectPermitOutcome, RemoteReconnectPermitRecovery, RemoteReplayApplyOutcome,
RemoteTransportLossOutcome,
};
#[cfg(test)]
mod tests;
use alloc::string::{String, ToString};
use alloc::sync::Arc;
use crate::connection::ConnectionPoolConfig;
use crate::{ConversationId, SdkError};
use self::protocol::{ProtocolRemoteTransport, RemoteTransport};
#[cfg(feature = "std")]
pub(crate) const SETUP_TIMEOUT: core::time::Duration = core::time::Duration::from_secs(5);
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ServerAddress(String);
impl ServerAddress {
pub fn new(value: impl Into<String>) -> Result<Self, SdkError> {
let value = value.into();
if value.trim().is_empty() {
return Err(connection_error("remote mode requires a server address"));
}
Ok(Self(value))
}
#[must_use]
pub fn as_str(&self) -> &str {
self.0.as_str()
}
}
#[derive(Clone, Debug)]
pub struct RemoteConfig {
pub server_address: ServerAddress,
pub channel_name: String,
pub conversation_id: ConversationId,
pub pool_config: ConnectionPoolConfig,
transport: Arc<dyn RemoteTransport>,
#[cfg(feature = "std")]
websocket: Option<Arc<websocket::WebSocketRemoteTransport>>,
}
impl RemoteConfig {
pub fn new(
server_address: impl Into<String>,
channel_name: impl Into<String>,
conversation_id: impl Into<ConversationId>,
pool_config: ConnectionPoolConfig,
) -> Result<Self, SdkError> {
Ok(Self {
server_address: ServerAddress::new(server_address)?,
channel_name: channel_name.into(),
conversation_id: conversation_id.into(),
pool_config: pool_config.validate()?,
transport: Arc::new(ProtocolRemoteTransport),
#[cfg(feature = "std")]
websocket: None,
})
}
#[cfg(feature = "std")]
pub fn connect_tcp(mut self) -> Result<Self, SdkError> {
let transport = self::tcp::TcpRemoteTransport::connect(&self.server_address)?;
self.transport = Arc::new(transport);
self.websocket = None;
Ok(self)
}
#[cfg(feature = "std")]
pub fn connect_tcp_with_auth(mut self, auth_token: &[u8]) -> Result<Self, SdkError> {
let transport =
self::tcp::TcpRemoteTransport::connect_with_auth(&self.server_address, auth_token)?;
self.transport = Arc::new(transport);
self.websocket = None;
Ok(self)
}
#[cfg(feature = "embedded")]
pub fn connect_loopback(
self,
server: Arc<liminal_server::server::embedded::EmbeddedServer>,
) -> Result<Self, SdkError> {
self.connect_loopback_with_auth(server, &[])
}
#[cfg(feature = "embedded")]
pub fn connect_loopback_with_auth(
mut self,
server: Arc<liminal_server::server::embedded::EmbeddedServer>,
auth_token: &[u8],
) -> Result<Self, SdkError> {
let transport =
self::loopback::LoopbackRemoteTransport::connect_with_auth(server, auth_token)?;
self.transport = Arc::new(transport);
self.websocket = None;
Ok(self)
}
#[cfg(feature = "std")]
pub fn connect_websocket(self) -> Result<Self, SdkError> {
self.connect_websocket_with_auth(&[])
}
#[cfg(feature = "std")]
pub fn connect_websocket_with_auth(mut self, auth_token: &[u8]) -> Result<Self, SdkError> {
let transport = Arc::new(websocket::WebSocketRemoteTransport::connect_with_auth(
&self.server_address,
auth_token,
)?);
self.transport = Arc::clone(&transport) as Arc<dyn RemoteTransport>;
self.websocket = Some(transport);
Ok(self)
}
#[cfg(feature = "std")]
#[must_use]
pub fn websocket_transport(&self) -> Option<Arc<websocket::WebSocketRemoteTransport>> {
self.websocket.clone()
}
}
fn connection_error(description: &str) -> SdkError {
SdkError::Connection {
description: description.to_string(),
}
}