use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use chrono::Utc;
use tokio::sync::{mpsc, oneshot};
use tokio::time::Instant;
use crate::EventData;
#[async_trait]
pub(crate) trait Transport: Send + Sync + 'static {
async fn connect(&mut self) -> Result<(), TransportError>;
async fn send(&mut self, event: EventData) -> Result<(), TransportError>;
async fn close(&mut self) -> Result<(), TransportError>;
}
#[derive(Debug, thiserror::Error)]
pub enum TransportError {
#[error("Connection error: {0}")]
Connection(String),
#[error("Send error: {0}")]
Send(String),
#[error("Configuration error: {0}")]
Configuration(String),
}
pub(crate) struct TransportLoopConfig {
pub retry_backoff: Vec<Duration>,
pub shutdown_retry_budget: Duration,
pub flush_interval: Duration,
}
impl Default for TransportLoopConfig {
fn default() -> Self {
Self {
retry_backoff: vec![
Duration::from_millis(100),
Duration::from_secs(1),
Duration::from_secs(5),
],
shutdown_retry_budget: Duration::from_secs(3),
flush_interval: Duration::from_secs(5),
}
}
}
async fn send_with_retry(
transport: &mut dyn Transport,
event: &EventData,
backoff: &[Duration],
deadline: Option<Instant>,
) -> Result<(), TransportError> {
let mut last_err = match transport.send(event.clone()).await {
Ok(()) => return Ok(()),
Err(e) => e,
};
for delay in backoff {
let delay = match deadline {
Some(deadline) => {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return Err(last_err);
}
(*delay).min(remaining)
}
None => *delay,
};
tokio::time::sleep(delay).await;
match transport.send(event.clone()).await {
Ok(()) => return Ok(()),
Err(e) => last_err = e,
}
}
Err(last_err)
}
fn dropped_events_warning(dropped_count: u64) -> EventData {
EventData {
event_type: "event".to_string(),
event_data: serde_json::json!({
"level": "WARN",
"target": "eyes_subscriber",
"fields": {
"message": "events dropped by eyes-subscriber backpressure",
"dropped_count": dropped_count,
},
}),
event_timestamp: Utc::now(),
process_instance_id: None,
}
}
async fn report_dropped_events(
transport: &mut dyn Transport,
dropped: &AtomicU64,
backoff: &[Duration],
deadline: Option<Instant>,
) {
let count = dropped.swap(0, Ordering::Relaxed);
if count == 0 {
return;
}
if let Err(e) =
send_with_retry(transport, &dropped_events_warning(count), backoff, deadline).await
{
eprintln!("Failed to report {} dropped events: {}", count, e);
dropped.fetch_add(count, Ordering::Relaxed);
}
}
async fn shutdown_flush(
transport: &mut dyn Transport,
receiver: &mut mpsc::Receiver<EventData>,
dropped: &AtomicU64,
config: &TransportLoopConfig,
deadline: Instant,
pending: Option<EventData>,
) {
let mut next = pending;
loop {
let event = match next.take() {
Some(event) => event,
None => match receiver.try_recv() {
Ok(event) => event,
Err(_) => break,
},
};
if let Err(e) =
send_with_retry(transport, &event, &config.retry_backoff, Some(deadline)).await
{
eprintln!("Failed to send event during shutdown, dropping: {}", e);
dropped.fetch_add(1, Ordering::Relaxed);
}
}
}
pub(crate) async fn run_transport_loop(
mut transport: Box<dyn Transport>,
mut receiver: mpsc::Receiver<EventData>,
mut shutdown_rx: oneshot::Receiver<()>,
completion_tx: oneshot::Sender<()>,
dropped: Arc<AtomicU64>,
config: TransportLoopConfig,
) {
if let Err(e) = transport.connect().await {
eprintln!("Failed to connect transport: {}", e);
let _ = completion_tx.send(());
return;
}
let mut flush_interval = tokio::time::interval(config.flush_interval);
flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
flush_interval.tick().await;
let mut shutdown_deadline: Option<Instant> = None;
let mut shutdown_signal_done = false;
loop {
tokio::select! {
event = receiver.recv() => {
let Some(event) = event else {
break;
};
let result = tokio::select! {
result = send_with_retry(
transport.as_mut(),
&event,
&config.retry_backoff,
None,
) => Some(result),
res = &mut shutdown_rx, if !shutdown_signal_done => {
shutdown_signal_done = true;
match res {
Ok(()) => None,
Err(_) => Some(send_with_retry(
transport.as_mut(),
&event,
&config.retry_backoff,
None,
).await),
}
}
};
match result {
Some(Ok(())) => {}
Some(Err(e)) => {
eprintln!("Failed to send event after retries, dropping: {}", e);
dropped.fetch_add(1, Ordering::Relaxed);
}
None => {
let deadline = Instant::now() + config.shutdown_retry_budget;
shutdown_deadline = Some(deadline);
shutdown_flush(
transport.as_mut(),
&mut receiver,
&dropped,
&config,
deadline,
Some(event),
)
.await;
break;
}
}
}
_ = flush_interval.tick() => {
report_dropped_events(transport.as_mut(), &dropped, &[], None).await;
}
res = &mut shutdown_rx, if !shutdown_signal_done => {
shutdown_signal_done = true;
if res.is_ok() {
let deadline = Instant::now() + config.shutdown_retry_budget;
shutdown_deadline = Some(deadline);
shutdown_flush(
transport.as_mut(),
&mut receiver,
&dropped,
&config,
deadline,
None,
)
.await;
break;
}
}
}
}
let deadline =
shutdown_deadline.unwrap_or_else(|| Instant::now() + config.shutdown_retry_budget);
report_dropped_events(
transport.as_mut(),
&dropped,
&config.retry_backoff,
Some(deadline),
)
.await;
if let Err(e) = transport.close().await {
eprintln!("Failed to close transport: {}", e);
}
let _ = completion_tx.send(());
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Mutex;
fn test_event(event_type: &str) -> EventData {
EventData {
event_type: event_type.to_string(),
event_data: serde_json::json!({}),
event_timestamp: Utc::now(),
process_instance_id: None,
}
}
fn fast_config() -> TransportLoopConfig {
TransportLoopConfig {
retry_backoff: vec![
Duration::from_millis(1),
Duration::from_millis(1),
Duration::from_millis(1),
],
shutdown_retry_budget: Duration::from_millis(50),
flush_interval: Duration::from_millis(10),
}
}
struct MockTransport {
sent: Arc<Mutex<Vec<EventData>>>,
attempts: Arc<AtomicU64>,
failures_remaining: Arc<AtomicU64>,
}
impl MockTransport {
fn new(failures: u64) -> Self {
Self {
sent: Arc::new(Mutex::new(Vec::new())),
attempts: Arc::new(AtomicU64::new(0)),
failures_remaining: Arc::new(AtomicU64::new(failures)),
}
}
fn handles(&self) -> (Arc<Mutex<Vec<EventData>>>, Arc<AtomicU64>) {
(self.sent.clone(), self.attempts.clone())
}
}
#[async_trait]
impl Transport for MockTransport {
async fn connect(&mut self) -> Result<(), TransportError> {
Ok(())
}
async fn send(&mut self, event: EventData) -> Result<(), TransportError> {
self.attempts.fetch_add(1, Ordering::Relaxed);
let remaining = self.failures_remaining.load(Ordering::Relaxed);
if remaining > 0 {
self.failures_remaining
.store(remaining - 1, Ordering::Relaxed);
return Err(TransportError::Send("mock failure".to_string()));
}
self.sent.lock().unwrap().push(event);
Ok(())
}
async fn close(&mut self) -> Result<(), TransportError> {
Ok(())
}
}
#[tokio::test]
async fn test_send_with_retry_recovers_from_transient_failure() {
let mut transport = MockTransport::new(2);
let (sent, attempts) = transport.handles();
let backoff = [Duration::from_millis(1); 3];
let result = send_with_retry(&mut transport, &test_event("event"), &backoff, None).await;
assert!(result.is_ok());
assert_eq!(attempts.load(Ordering::Relaxed), 3);
assert_eq!(sent.lock().unwrap().len(), 1);
}
#[tokio::test]
async fn test_retry_then_drop_on_persistent_failure() {
let mut transport = MockTransport::new(u64::MAX);
let (sent, attempts) = transport.handles();
let backoff = [Duration::from_millis(1); 3];
let result = send_with_retry(&mut transport, &test_event("event"), &backoff, None).await;
assert!(result.is_err());
assert_eq!(attempts.load(Ordering::Relaxed), 4);
assert!(sent.lock().unwrap().is_empty());
}
#[tokio::test]
async fn test_transport_loop_counts_drops_on_persistent_failure() {
let transport = MockTransport::new(u64::MAX);
let (sent, _) = transport.handles();
let (sender, receiver) = mpsc::channel(16);
let (_shutdown_tx, shutdown_rx) = oneshot::channel();
let (completion_tx, completion_rx) = oneshot::channel();
let dropped = Arc::new(AtomicU64::new(0));
let loop_handle = tokio::spawn(run_transport_loop(
Box::new(transport),
receiver,
shutdown_rx,
completion_tx,
dropped.clone(),
fast_config(),
));
sender.try_send(test_event("event")).unwrap();
drop(sender);
loop_handle.await.unwrap();
completion_rx.await.unwrap();
assert_eq!(dropped.load(Ordering::Relaxed), 1);
assert!(sent.lock().unwrap().is_empty());
}
#[tokio::test]
async fn test_drop_counter_emits_synthetic_event_on_recovery() {
let transport = MockTransport::new(0);
let (sent, _) = transport.handles();
let (sender, receiver) = mpsc::channel(16);
let (shutdown_tx, shutdown_rx) = oneshot::channel();
let (completion_tx, completion_rx) = oneshot::channel();
let dropped = Arc::new(AtomicU64::new(7));
let loop_handle = tokio::spawn(run_transport_loop(
Box::new(transport),
receiver,
shutdown_rx,
completion_tx,
dropped.clone(),
fast_config(),
));
tokio::time::sleep(Duration::from_millis(50)).await;
assert_eq!(dropped.load(Ordering::Relaxed), 0);
{
let sent = sent.lock().unwrap();
assert_eq!(sent.len(), 1);
let event = &sent[0];
assert_eq!(event.event_type, "event");
assert_eq!(event.event_data["level"], "WARN");
assert_eq!(event.event_data["target"], "eyes_subscriber");
assert_eq!(
event.event_data["fields"]["message"],
"events dropped by eyes-subscriber backpressure"
);
assert_eq!(event.event_data["fields"]["dropped_count"], 7);
}
let _ = shutdown_tx.send(());
loop_handle.await.unwrap();
completion_rx.await.unwrap();
assert_eq!(dropped.load(Ordering::Relaxed), 0);
assert_eq!(sent.lock().unwrap().len(), 1);
drop(sender);
}
#[tokio::test]
async fn test_failed_drop_report_restores_counter() {
let transport = MockTransport::new(u64::MAX);
let (sent, _) = transport.handles();
let (sender, receiver) = mpsc::channel(16);
let (shutdown_tx, shutdown_rx) = oneshot::channel();
let (completion_tx, completion_rx) = oneshot::channel();
let dropped = Arc::new(AtomicU64::new(5));
let loop_handle = tokio::spawn(run_transport_loop(
Box::new(transport),
receiver,
shutdown_rx,
completion_tx,
dropped.clone(),
fast_config(),
));
tokio::time::sleep(Duration::from_millis(50)).await;
let _ = shutdown_tx.send(());
loop_handle.await.unwrap();
completion_rx.await.unwrap();
assert_eq!(dropped.load(Ordering::Relaxed), 5);
assert!(sent.lock().unwrap().is_empty());
drop(sender);
}
#[tokio::test]
async fn test_shutdown_flush_respects_retry_budget() {
let transport = MockTransport::new(u64::MAX);
let (sender, receiver) = mpsc::channel(16);
let (shutdown_tx, shutdown_rx) = oneshot::channel();
let (completion_tx, completion_rx) = oneshot::channel();
let dropped = Arc::new(AtomicU64::new(0));
let mut config = fast_config();
config.retry_backoff = vec![Duration::from_secs(60)];
config.shutdown_retry_budget = Duration::from_millis(20);
let loop_handle = tokio::spawn(run_transport_loop(
Box::new(transport),
receiver,
shutdown_rx,
completion_tx,
dropped.clone(),
config,
));
for _ in 0..3 {
sender.try_send(test_event("event")).unwrap();
}
let _ = shutdown_tx.send(());
let start = std::time::Instant::now();
loop_handle.await.unwrap();
completion_rx.await.unwrap();
assert!(
start.elapsed() < Duration::from_secs(5),
"shutdown flush should not hang on a dead transport"
);
assert_eq!(dropped.load(Ordering::Relaxed), 3);
}
#[tokio::test]
async fn test_shutdown_cancels_in_flight_retry() {
let transport = MockTransport::new(u64::MAX);
let (sender, receiver) = mpsc::channel(16);
let (shutdown_tx, shutdown_rx) = oneshot::channel();
let (completion_tx, completion_rx) = oneshot::channel();
let dropped = Arc::new(AtomicU64::new(0));
let mut config = fast_config();
config.retry_backoff = vec![Duration::from_secs(60)];
config.shutdown_retry_budget = Duration::from_millis(20);
let loop_handle = tokio::spawn(run_transport_loop(
Box::new(transport),
receiver,
shutdown_rx,
completion_tx,
dropped.clone(),
config,
));
sender.try_send(test_event("event")).unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
let _ = shutdown_tx.send(());
let start = std::time::Instant::now();
loop_handle.await.unwrap();
completion_rx.await.unwrap();
assert!(
start.elapsed() < Duration::from_secs(5),
"shutdown should cancel an in-flight retry instead of waiting out the backoff"
);
assert_eq!(dropped.load(Ordering::Relaxed), 1);
}
#[tokio::test]
async fn test_dropped_handle_does_not_stop_the_loop() {
let transport = MockTransport::new(0);
let (sent, _) = transport.handles();
let (sender, receiver) = mpsc::channel(16);
let (shutdown_tx, shutdown_rx) = oneshot::channel();
let (completion_tx, completion_rx) = oneshot::channel();
let loop_handle = tokio::spawn(run_transport_loop(
Box::new(transport),
receiver,
shutdown_rx,
completion_tx,
Arc::new(AtomicU64::new(0)),
fast_config(),
));
drop(shutdown_tx);
tokio::time::sleep(Duration::from_millis(20)).await;
sender.send(test_event("after_drop_1")).await.unwrap();
sender.send(test_event("after_drop_2")).await.unwrap();
tokio::time::sleep(Duration::from_millis(50)).await;
assert_eq!(
sent.lock().unwrap().len(),
2,
"events after handle drop must still be delivered"
);
drop(sender);
completion_rx.await.expect("loop should complete");
loop_handle.await.unwrap();
}
}