| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244 |
- use async_std::sync::Arc;
- use std::convert::From;
- use std::net::SocketAddr;
- use std::path::Path;
- use super::reqrep::{PeerId, Publisher, RepProtocol, Reply, ReqProtocol, Request, Subscriber};
- use crate::{
- serial::deserialize, serial::serialize, slab::Slab, slabstore::SlabStore, Error, Result,
- };
- use async_executor::Executor;
- use log::*;
- pub type Slabs = Vec<Vec<u8>>;
- #[repr(u8)]
- enum GatewayCommand {
- PutSlab,
- GetSlab,
- GetLastIndex,
- }
- pub struct GatewayService {
- slabstore: Arc<SlabStore>,
- addr: SocketAddr,
- pub_addr: SocketAddr,
- }
- impl GatewayService {
- pub fn new(addr: SocketAddr, pub_addr: SocketAddr) -> Result<Arc<GatewayService>> {
- let slabstore = SlabStore::new(Path::new("slabstore.db"))?;
- Ok(Arc::new(GatewayService {
- slabstore,
- addr,
- pub_addr,
- }))
- }
- pub async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
- let service_name = String::from("GATEWAY DAEMON");
- let mut protocol = RepProtocol::new(self.addr.clone(), service_name.clone());
- let (send, recv) = protocol.start().await?;
- let (publish_queue, publish_recv_queue) = async_channel::unbounded::<Vec<u8>>();
- let publisher_task = executor.spawn(Self::start_publisher(
- self.pub_addr,
- service_name,
- publish_recv_queue.clone(),
- ));
- let handle_request_task = executor.spawn(self.handle_request_loop(
- send.clone(),
- recv.clone(),
- publish_queue.clone(),
- executor.clone(),
- ));
- protocol.run(executor.clone()).await?;
- let _ = publisher_task.cancel().await;
- let _ = handle_request_task.cancel().await;
- Ok(())
- }
- async fn start_publisher(
- pub_addr: SocketAddr,
- service_name: String,
- publish_recv_queue: async_channel::Receiver<Vec<u8>>,
- ) -> Result<()> {
- let mut publisher = Publisher::new(pub_addr, service_name);
- publisher.start(publish_recv_queue).await?;
- Ok(())
- }
- async fn handle_request_loop(
- self: Arc<Self>,
- send_queue: async_channel::Sender<(PeerId, Reply)>,
- recv_queue: async_channel::Receiver<(PeerId, Request)>,
- publish_queue: async_channel::Sender<Vec<u8>>,
- executor: Arc<Executor<'_>>,
- ) -> Result<()> {
- loop {
- match recv_queue.recv().await {
- Ok(msg) => {
- let slabstore = self.slabstore.clone();
- let _ = executor
- .spawn(Self::handle_request(
- msg,
- slabstore,
- send_queue.clone(),
- publish_queue.clone(),
- ))
- .detach();
- }
- Err(_) => {
- break;
- }
- }
- }
- Ok(())
- }
- async fn handle_request(
- msg: (PeerId, Request),
- slabstore: Arc<SlabStore>,
- send_queue: async_channel::Sender<(PeerId, Reply)>,
- publish_queue: async_channel::Sender<Vec<u8>>,
- ) -> Result<()> {
- let request = msg.1;
- let peer = msg.0;
- match request.get_command() {
- 0 => {
- // PUTSLAB
- let slab = request.get_payload();
- // add to slabstore
- slabstore.put(slab.clone())?;
- // send reply
- let reply = Reply::from(&request, 0, vec![]);
- send_queue.send((peer, reply)).await?;
- // publish to all subscribes
- publish_queue.send(slab).await?;
- info!("Received putslab msg");
- }
- 1 => {
- let index = request.get_payload();
- let slab = slabstore.get(index)?;
- let mut payload = vec![];
- if let Some(sb) = slab {
- payload = sb;
- }
- let reply = Reply::from(&request, 0, payload);
- send_queue.send((peer, reply)).await?;
- // GETSLAB
- info!("Received getslab msg");
- }
- 2 => {
- let index = slabstore.get_last_index_as_bytes()?;
- let reply = Reply::from(&request, 0, index);
- send_queue.send((peer, reply)).await?;
- // GETLASTINDEX
- info!("Received getlastindex msg");
- }
- _ => {
- return Err(Error::ServicesError("received wrong command"));
- }
- }
- Ok(())
- }
- }
- pub struct GatewayClient {
- protocol: ReqProtocol,
- slabstore: Arc<SlabStore>,
- }
- impl GatewayClient {
- pub fn new(addr: SocketAddr, path: &Path) -> Result<Self> {
- let protocol = ReqProtocol::new(addr, String::from("GATEWAY CLIENT"));
- let slabstore = SlabStore::new(path)?;
- Ok(GatewayClient {
- protocol,
- slabstore,
- })
- }
- pub async fn start(&mut self) -> Result<()> {
- self.protocol.start().await?;
- info!("Start Syncing");
- let local_last_index = self.slabstore.get_last_index()?;
- let last_index = self.get_last_index().await?;
- if last_index > 0 {
- for index in (local_last_index + 1)..(last_index + 1) {
- self.get_slab(index).await?;
- }
- }
- info!("End Syncing");
- Ok(())
- }
- pub async fn get_slab(&mut self, index: u64) -> Result<Vec<u8>> {
- let slab = self
- .protocol
- .request(GatewayCommand::GetSlab as u8, serialize(&index))
- .await?;
- self.slabstore.put(slab.clone())?;
- Ok(slab)
- }
- pub async fn put_slab(&mut self, mut slab: Slab) -> Result<()> {
- let last_index = self.get_last_index().await?;
- slab.set_index(last_index + 1);
- let slab = serialize(&slab);
- self.protocol
- .request(GatewayCommand::PutSlab as u8, slab.clone())
- .await?;
- Ok(())
- }
- pub async fn get_last_index(&mut self) -> Result<u64> {
- let rep = self
- .protocol
- .request(GatewayCommand::GetLastIndex as u8, vec![])
- .await?;
- Ok(deserialize(&rep)?)
- }
- pub fn get_slabstore(&self) -> Arc<SlabStore> {
- self.slabstore.clone()
- }
- pub async fn start_subscriber(sub_addr: SocketAddr) -> Result<Subscriber> {
- let mut subscriber = Subscriber::new(sub_addr, String::from("GATEWAY CLIENT"));
- subscriber.start().await?;
- Ok(subscriber)
- }
- pub async fn subscribe(mut subscriber: Subscriber, slabstore: Arc<SlabStore>) -> Result<()> {
- loop {
- let slab: Vec<u8>;
- slab = subscriber.fetch().await?;
- slabstore.put(slab)?;
- }
- }
- }
|