use super::{error::BaseNodeServiceError, service::BaseNodeState};
use futures::{stream::Fuse, StreamExt};
use std::sync::Arc;
use tari_comms::peer_manager::Peer;
use tari_common_types::chain_metadata::ChainMetadata;
use tari_service_framework::reply_channel::SenderService;
use tokio::sync::broadcast;
use tower::Service;
pub type BaseNodeEventSender = broadcast::Sender<Arc<BaseNodeEvent>>;
pub type BaseNodeEventReceiver = broadcast::Receiver<Arc<BaseNodeEvent>>;
#[derive(Debug)]
pub enum BaseNodeServiceRequest {
GetChainMetadata,
SetBaseNodePeer(Box<Peer>),
}
#[derive(Debug)]
pub enum BaseNodeServiceResponse {
ChainMetadata(Option<ChainMetadata>),
BaseNodePeerSet,
}
#[derive(Clone, Debug, Hash, PartialEq, Eq)]
pub enum BaseNodeEvent {
BaseNodeState(BaseNodeState),
BaseNodePeerSet(Box<Peer>),
}
#[derive(Clone)]
pub struct BaseNodeServiceHandle {
handle: SenderService<BaseNodeServiceRequest, Result<BaseNodeServiceResponse, BaseNodeServiceError>>,
event_stream_sender: BaseNodeEventSender,
}
impl BaseNodeServiceHandle {
pub fn new(
handle: SenderService<BaseNodeServiceRequest, Result<BaseNodeServiceResponse, BaseNodeServiceError>>,
event_stream_sender: BaseNodeEventSender,
) -> Self
{
Self {
handle,
event_stream_sender,
}
}
pub fn get_event_stream_fused(&self) -> Fuse<BaseNodeEventReceiver> {
self.event_stream_sender.subscribe().fuse()
}
pub async fn set_base_node_peer(&mut self, peer: Peer) -> Result<(), BaseNodeServiceError> {
match self
.handle
.call(BaseNodeServiceRequest::SetBaseNodePeer(Box::new(peer)))
.await??
{
BaseNodeServiceResponse::BaseNodePeerSet => Ok(()),
_ => Err(BaseNodeServiceError::UnexpectedApiResponse),
}
}
}