use std::sync::Arc;
use crate::CloseReason;
use crate::messages::{ControllerMessage, ServiceMessage};
#[derive(Debug, thiserror::Error)]
#[error("transport send failed")]
pub struct TransportError;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TransportClosePolicy {
Reconnect { reason: Option<CloseReason> },
Shutdown,
}
#[async_trait::async_trait]
pub trait ServiceTransport: Send {
async fn transport_send(&mut self, msg: ServiceMessage) -> Result<(), TransportError>;
async fn transport_send_best_effort(&mut self, msg: ServiceMessage);
async fn transport_send_auto_paginate(
&mut self,
msg: ServiceMessage,
) -> Result<(), TransportError>;
async fn transport_recv(&mut self) -> Option<ControllerMessage>;
fn close_policy(&self) -> TransportClosePolicy {
TransportClosePolicy::Reconnect { reason: None }
}
fn is_yielded(&self) -> bool {
false
}
fn yield_change_notifier(&self) -> Option<Arc<tokio::sync::Notify>> {
None
}
}
#[cfg(test)]
mod tests {
use super::{ServiceTransport, TransportClosePolicy, TransportError};
use crate::CloseReason;
use crate::messages::{ControllerMessage, ServiceMessage};
struct DefaultTransport;
#[async_trait::async_trait]
impl ServiceTransport for DefaultTransport {
async fn transport_send(&mut self, _msg: ServiceMessage) -> Result<(), TransportError> {
Ok(())
}
async fn transport_send_best_effort(&mut self, _msg: ServiceMessage) {}
async fn transport_send_auto_paginate(
&mut self,
_msg: ServiceMessage,
) -> Result<(), TransportError> {
Ok(())
}
async fn transport_recv(&mut self) -> Option<ControllerMessage> {
None
}
}
#[test]
fn default_close_policy_is_reconnect_with_no_reason() {
let transport = DefaultTransport;
assert_eq!(
transport.close_policy(),
TransportClosePolicy::Reconnect { reason: None }
);
}
#[test]
fn default_is_yielded_is_false() {
let transport = DefaultTransport;
assert!(!transport.is_yielded());
}
#[test]
fn close_policy_reconnect_none_eq_reconnect_none() {
let a = TransportClosePolicy::Reconnect { reason: None };
let b = TransportClosePolicy::Reconnect { reason: None };
assert_eq!(a, b);
}
#[test]
fn close_policy_shutdown_eq_shutdown() {
assert_eq!(
TransportClosePolicy::Shutdown,
TransportClosePolicy::Shutdown
);
}
#[test]
fn close_policy_reconnect_ne_shutdown() {
assert_ne!(
TransportClosePolicy::Reconnect { reason: None },
TransportClosePolicy::Shutdown
);
}
#[test]
fn close_policy_clone_preserves_equality() {
let original = TransportClosePolicy::Reconnect {
reason: Some(CloseReason::CertificateRotated),
};
let cloned = original.clone();
assert_eq!(original, cloned);
}
#[test]
fn close_policy_debug_contains_variant_name() {
let policy = TransportClosePolicy::Shutdown;
let debug = format!("{policy:?}");
assert!(debug.contains("Shutdown"), "Debug output was: {debug}");
}
#[test]
fn default_yield_change_notifier_returns_none() {
let t = DefaultTransport;
assert!(t.yield_change_notifier().is_none());
}
}