v_common_module/
remote_indv_r_storage.rs1use crate::common::get_queue_status;
2use crate::veda_backend::get_storage_use_prop;
3use nng::{Message, Protocol, Socket};
4use std::cell::RefCell;
5use std::str;
6use std::sync::Mutex;
7use uuid::Uuid;
8use v_api::v_onto::individual::{Individual, RawObj};
9use v_api::v_onto::individual2msgpack::to_msgpack;
10use v_storage::remote_storage_client::*;
11use v_storage::storage::*;
12
13lazy_static! {
14 pub static ref STORAGE: Mutex<RefCell<StorageROClient>> = Mutex::new(RefCell::new(StorageROClient::default()));
15}
16
17pub fn inproc_storage_manager() -> std::io::Result<()> {
20 let ro_storage_url = "inproc://nng/".to_owned() + &Uuid::new_v4().to_hyphenated().to_string();
21 STORAGE.lock().unwrap().get_mut().addr = ro_storage_url.to_owned();
22
23 let mut storage = get_storage_use_prop(StorageMode::ReadOnly);
24
25 let server = Socket::new(Protocol::Rep0)?;
26 if let Err(e) = server.listen(&ro_storage_url) {
27 error!("fail listen, {:?}", e);
28 return Ok(());
29 }
30
31 loop {
32 if let Ok(recv_msg) = server.recv() {
33 let res = req_prepare(&recv_msg, &mut storage);
34 if let Err(e) = server.send(res) {
35 error!("fail send {:?}", e);
36 }
37 }
38 }
39}
40
41fn req_prepare(request: &Message, storage: &mut VStorage) -> Message {
42 if let Ok(id) = str::from_utf8(request.as_slice()) {
43 if id.starts_with("srv:queue-state-") {
44 let indv = get_queue_status(id);
45
46 let mut binobj: Vec<u8> = Vec::new();
47 if let Err(e) = to_msgpack(&indv, &mut binobj) {
48 error!("failed to serialize, err = {:?}", e);
49 return Message::from("[]".as_bytes());
50 }
51
52 return Message::from(binobj.as_slice());
53 }
54
55 let binobj = storage.get_raw_value(StorageId::Individuals, id);
56 if binobj.is_empty() {
57 return Message::from("[]".as_bytes());
58 }
59 return Message::from(binobj.as_slice());
60 }
61
62 Message::default()
63}
64
65pub fn get_individual(id: &str) -> Option<Individual> {
66 if id.starts_with("srv:queue-state-") {
67 return Some(get_queue_status(id));
68 }
69
70 let req = Message::from(id.to_string().as_bytes());
71
72 let mut sh_client = STORAGE.lock().unwrap();
73 let client = sh_client.get_mut();
74
75 if !client.is_ready {
76 client.connect();
77 }
78
79 if !client.is_ready {
80 return None;
81 }
82
83 if let Err(e) = client.soc.send(req) {
84 error!("fail send to storage_manager, err={:?}", e);
85 return None;
86 }
87
88 let wmsg = client.soc.recv();
90 if let Err(e) = wmsg {
91 error!("fail recv from main module, err={:?}", e);
92 return None;
93 }
94
95 drop(sh_client);
96
97 if let Ok(msg) = wmsg {
98 let data = msg.as_slice();
99 if data == b"[]" {
100 return None;
101 }
102 return Some(Individual::new_raw(RawObj::new(data.to_vec())));
103 }
104
105 None
106}