use std::collections::VecDeque;
use std::time::Instant;
use minip2p_pubsub::{PubsubAction, PubsubAgent, PubsubEvent};
use minip2p_swarm::SwarmEvent;
use crate::EndpointSwarm;
use crate::Error;
#[derive(Debug, thiserror::Error)]
pub enum PubsubError {
#[error("pubsub is not enabled on this endpoint (EndpointBuilder::pubsub)")]
NotEnabled,
#[error("cannot unsubscribe from the discovery topic while discovery is enabled")]
DiscoveryTopicReserved,
#[error(transparent)]
Publish(#[from] minip2p_pubsub::PublishError),
#[error(transparent)]
Topic(#[from] minip2p_pubsub::TopicError),
#[error(transparent)]
Driver(#[from] Error),
}
pub(crate) struct PubsubDriver {
pub(crate) agent: PubsubAgent,
pub(crate) events: VecDeque<PubsubEvent>,
epoch: Instant,
}
impl PubsubDriver {
pub(crate) fn new(agent: PubsubAgent) -> Self {
Self {
agent,
events: VecDeque::new(),
epoch: Instant::now(),
}
}
pub(crate) fn now_ms(&self) -> u64 {
self.epoch.elapsed().as_millis() as u64
}
pub(crate) fn ingest(&mut self, event: &SwarmEvent, swarm: &mut EndpointSwarm) -> bool {
let handled = self.agent.handle_event(event, self.now_ms());
self.pump(swarm);
handled
}
pub(crate) fn tick(&mut self, swarm: &mut EndpointSwarm) {
let now_ms = self.now_ms();
if self.agent.next_timeout(now_ms) != Some(0) {
return;
}
self.agent.handle_tick(now_ms);
self.pump(swarm);
}
pub(crate) fn pump(&mut self, swarm: &mut EndpointSwarm) {
loop {
let mut progressed = false;
while let Some(action) = self.agent.poll_action() {
progressed = true;
self.execute(action, swarm);
}
while let Some(event) = self.agent.poll_event() {
progressed = true;
self.events.push_back(event);
}
if !progressed {
break;
}
}
}
fn execute(&mut self, action: PubsubAction, swarm: &mut EndpointSwarm) {
match action {
PubsubAction::OpenStream {
token,
peer,
protocol_id,
} => {
let result = swarm
.open_stream(&peer, &protocol_id)
.map_err(|e| e.to_string());
let now_ms = self.now_ms();
self.agent.stream_open_result(&peer, token, result, now_ms);
}
PubsubAction::SendStream {
token,
peer,
stream_id,
data,
} => {
let result = swarm
.send_stream(&peer, stream_id, data)
.map_err(|e| e.to_string());
let now_ms = self.now_ms();
self.agent
.send_result(&peer, stream_id, token, result, now_ms);
}
PubsubAction::CloseStreamWrite { peer, stream_id } => {
let _ = swarm.close_stream_write(&peer, stream_id);
}
PubsubAction::ResetStream { peer, stream_id } => {
let _ = swarm.reset_stream(&peer, stream_id);
}
}
}
}