use pamoja_core::{Error, Result, Transport};
pub struct Faulty<T> {
inner: T,
upcoming_failures: usize,
}
impl<T> Faulty<T> {
pub fn new(inner: T, failures: usize) -> Self {
Self {
inner,
upcoming_failures: failures,
}
}
pub fn fail_next(&mut self, count: usize) {
self.upcoming_failures = count;
}
}
impl<T: Transport> Transport for Faulty<T> {
async fn connect(&mut self) -> Result<()> {
self.inner.connect().await
}
async fn send(&mut self, topic: &str, payload: &[u8]) -> Result<()> {
if self.upcoming_failures > 0 {
self.upcoming_failures -= 1;
return Err(Error::Transport("simulated link failure".to_owned()));
}
self.inner.send(topic, payload).await
}
async fn subscribe(&mut self, topic: &str) -> Result<()> {
self.inner.subscribe(topic).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{LoopbackBroker, LoopbackTransport};
#[tokio::test]
async fn fails_the_configured_sends_then_passes_through() {
let broker = LoopbackBroker::new();
let mut gateway = LoopbackTransport::new(broker.clone());
gateway.connect().await.expect("connect");
gateway.subscribe("#").await.expect("subscribe");
let mut node = Faulty::new(LoopbackTransport::new(broker), 2);
node.connect().await.expect("connect");
assert!(node.send("t", b"1").await.is_err());
assert!(node.send("t", b"2").await.is_err());
node.send("t", b"3")
.await
.expect("third send passes through");
let message = gateway.recv().await.expect("recv").expect("a message");
assert_eq!(message.payload, b"3");
}
}