|
@@ -8,7 +8,7 @@ use crate::{serial::deserialize, serial::serialize, Error, Result};
|
|
|
use async_executor::Executor;
|
|
use async_executor::Executor;
|
|
|
use log::*;
|
|
use log::*;
|
|
|
|
|
|
|
|
-pub type Slabs = Vec<Vec<u8>>;
|
|
|
|
|
|
|
+pub type GatewaySlabsSubscriber = async_channel::Receiver<Slab>;
|
|
|
|
|
|
|
|
#[repr(u8)]
|
|
#[repr(u8)]
|
|
|
enum GatewayError {
|
|
enum GatewayError {
|
|
@@ -179,6 +179,8 @@ impl GatewayService {
|
|
|
pub struct GatewayClient {
|
|
pub struct GatewayClient {
|
|
|
protocol: ReqProtocol,
|
|
protocol: ReqProtocol,
|
|
|
slabstore: Arc<SlabStore>,
|
|
slabstore: Arc<SlabStore>,
|
|
|
|
|
+ gateway_slabs_sub_s: async_channel::Sender<Slab>,
|
|
|
|
|
+ gateway_slabs_sub_rv: GatewaySlabsSubscriber,
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
impl GatewayClient {
|
|
impl GatewayClient {
|
|
@@ -187,9 +189,13 @@ impl GatewayClient {
|
|
|
|
|
|
|
|
let slabstore = SlabStore::new(rocks)?;
|
|
let slabstore = SlabStore::new(rocks)?;
|
|
|
|
|
|
|
|
|
|
+ let (gateway_slabs_sub_s, gateway_slabs_sub_rv) = async_channel::unbounded::<Slab>();
|
|
|
|
|
+
|
|
|
Ok(GatewayClient {
|
|
Ok(GatewayClient {
|
|
|
protocol,
|
|
protocol,
|
|
|
slabstore,
|
|
slabstore,
|
|
|
|
|
+ gateway_slabs_sub_s,
|
|
|
|
|
+ gateway_slabs_sub_rv,
|
|
|
})
|
|
})
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -221,14 +227,16 @@ impl GatewayClient {
|
|
|
Ok(last_index)
|
|
Ok(last_index)
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- pub async fn get_slab(&mut self, index: u64) -> Result<Option<Vec<u8>>> {
|
|
|
|
|
|
|
+ pub async fn get_slab(&mut self, index: u64) -> Result<Option<Slab>> {
|
|
|
let rep = self
|
|
let rep = self
|
|
|
.protocol
|
|
.protocol
|
|
|
.request(GatewayCommand::GetSlab as u8, serialize(&index))
|
|
.request(GatewayCommand::GetSlab as u8, serialize(&index))
|
|
|
.await?;
|
|
.await?;
|
|
|
|
|
|
|
|
if let Some(slab) = rep {
|
|
if let Some(slab) = rep {
|
|
|
- self.slabstore.put(deserialize(&slab)?)?;
|
|
|
|
|
|
|
+ let slab: Slab = deserialize(&slab)?;
|
|
|
|
|
+ self.gateway_slabs_sub_s.send(slab.clone()).await?;
|
|
|
|
|
+ self.slabstore.put(slab.clone())?;
|
|
|
return Ok(Some(slab));
|
|
return Ok(Some(slab));
|
|
|
}
|
|
}
|
|
|
Ok(None)
|
|
Ok(None)
|
|
@@ -267,9 +275,32 @@ impl GatewayClient {
|
|
|
self.slabstore.clone()
|
|
self.slabstore.clone()
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- pub async fn start_subscriber(sub_addr: SocketAddr) -> Result<Subscriber> {
|
|
|
|
|
|
|
+ pub async fn start_subscriber(
|
|
|
|
|
+ &self,
|
|
|
|
|
+ sub_addr: SocketAddr,
|
|
|
|
|
+ executor: Arc<Executor<'_>>,
|
|
|
|
|
+ ) -> Result<GatewaySlabsSubscriber> {
|
|
|
let mut subscriber = Subscriber::new(sub_addr, String::from("GATEWAY CLIENT"));
|
|
let mut subscriber = Subscriber::new(sub_addr, String::from("GATEWAY CLIENT"));
|
|
|
subscriber.start().await?;
|
|
subscriber.start().await?;
|
|
|
- Ok(subscriber)
|
|
|
|
|
|
|
+ executor
|
|
|
|
|
+ .spawn(Self::subscribe_loop(
|
|
|
|
|
+ subscriber,
|
|
|
|
|
+ self.slabstore.clone(),
|
|
|
|
|
+ self.gateway_slabs_sub_s.clone(),
|
|
|
|
|
+ ))
|
|
|
|
|
+ .detach();
|
|
|
|
|
+ Ok(self.gateway_slabs_sub_rv.clone())
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ async fn subscribe_loop(
|
|
|
|
|
+ mut subscriber: Subscriber,
|
|
|
|
|
+ slabstore: Arc<SlabStore>,
|
|
|
|
|
+ gateway_slabs_sub_s: async_channel::Sender<Slab>,
|
|
|
|
|
+ ) -> Result<()> {
|
|
|
|
|
+ loop {
|
|
|
|
|
+ let slab = subscriber.fetch::<Slab>().await?;
|
|
|
|
|
+ gateway_slabs_sub_s.send(slab.clone()).await?;
|
|
|
|
|
+ slabstore.put(slab)?;
|
|
|
|
|
+ }
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|