#![allow(dead_code)]
use ocpp_client::{TransportError, TransportEvent, TransportSink, TransportStream};
use std::future::Future;
use std::pin::Pin;
use tokio::sync::mpsc;
pub struct FakeSink(mpsc::UnboundedSender<TransportEvent>);
pub struct FakeSource(mpsc::UnboundedReceiver<TransportEvent>);
impl FakeSink {
pub fn send_event(&self, event: TransportEvent) {
self.0.send(event).expect("other end still open");
}
pub fn duplicate(&self) -> FakeSink {
FakeSink(self.0.clone())
}
}
impl FakeSource {
pub async fn recv_event(&mut self) -> Option<TransportEvent> {
self.0.recv().await
}
}
impl TransportSink for FakeSink {
fn send<'a>(
&'a mut self,
frame: String,
) -> Pin<Box<dyn Future<Output = Result<(), TransportError>> + Send + 'a>> {
Box::pin(async move {
self.0
.send(TransportEvent::Frame(frame))
.map_err(|e| Box::new(e) as TransportError)
})
}
fn ping<'a>(
&'a mut self,
payload: Vec<u8>,
) -> Pin<Box<dyn Future<Output = Result<(), TransportError>> + Send + 'a>> {
Box::pin(async move {
self.0
.send(TransportEvent::Ping(payload))
.map_err(|e| Box::new(e) as TransportError)
})
}
fn pong<'a>(
&'a mut self,
payload: Vec<u8>,
) -> Pin<Box<dyn Future<Output = Result<(), TransportError>> + Send + 'a>> {
Box::pin(async move {
self.0
.send(TransportEvent::Pong(payload))
.map_err(|e| Box::new(e) as TransportError)
})
}
fn close<'a>(
&'a mut self,
) -> Pin<Box<dyn Future<Output = Result<(), TransportError>> + Send + 'a>> {
Box::pin(async move { Ok(()) })
}
}
impl TransportStream for FakeSource {
fn recv<'a>(
&'a mut self,
) -> Pin<Box<dyn Future<Output = Result<Option<TransportEvent>, TransportError>> + Send + 'a>>
{
Box::pin(async move { Ok(self.0.recv().await) })
}
}
pub fn fake_transport_pair() -> ((FakeSink, FakeSource), (FakeSink, FakeSource)) {
let (a_tx, a_rx) = mpsc::unbounded_channel();
let (b_tx, b_rx) = mpsc::unbounded_channel();
(
(FakeSink(a_tx), FakeSource(b_rx)),
(FakeSink(b_tx), FakeSource(a_rx)),
)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PongBehavior {
Echo,
Silent,
WrongPayload,
}
pub fn spawn_peer(
sink: FakeSink,
source: FakeSource,
behavior: PongBehavior,
) -> mpsc::UnboundedReceiver<Vec<u8>> {
let (observed_tx, observed_rx) = mpsc::unbounded_channel();
let mut source = source;
tokio::spawn(async move {
while let Some(event) = source.0.recv().await {
if let TransportEvent::Ping(payload) = event {
if observed_tx.send(payload.clone()).is_err() {
break;
}
match behavior {
PongBehavior::Echo => {
let _ = sink.0.send(TransportEvent::Pong(payload));
}
PongBehavior::WrongPayload => {
let _ = sink.0.send(TransportEvent::Pong(Vec::new()));
}
PongBehavior::Silent => {}
}
}
}
});
observed_rx
}