v_common_module/
remote_indv_r_storage.rs

1use 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
17// inproc storage server
18
19pub 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    // Wait for the response from the server.
89    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}