1use std::collections::{HashMap, HashSet, VecDeque};
2
3use raftpb::Message;
4
5#[derive(Debug, Clone, PartialEq)]
11pub struct ReadState {
12 pub index: u64,
13 pub request_ctx: Vec<u8>,
14}
15
16#[derive(Debug, PartialEq, Clone, Copy)]
17pub enum ReadOnlyOption {
18 Safe,
21 LeaseBased,
27}
28
29impl Default for ReadOnlyOption {
30 fn default() -> ReadOnlyOption {
31 ReadOnlyOption::Safe
32 }
33}
34
35#[derive(Default, Debug, Clone)]
36pub struct ReadIndexStatus {
37 pub req: Message,
38 pub index: u64,
39 acks: HashSet<u64>,
40}
41
42#[derive(Default, Debug, Clone)]
43pub struct ReadOnly {
44 pub option: ReadOnlyOption,
45 pub pending_read_index: HashMap<Vec<u8>, ReadIndexStatus>,
46 pub read_index_queue: VecDeque<Vec<u8>>,
47}
48
49impl ReadOnly {
50 pub(crate) fn new(option: ReadOnlyOption) -> ReadOnly {
51 ReadOnly {
52 option,
53 pending_read_index: HashMap::new(),
54 read_index_queue: VecDeque::new(),
55 }
56 }
57
58 pub(crate) fn last_pending_request_ctx(&mut self) -> Option<Vec<u8>> {
59 self.read_index_queue.back().cloned()
60 }
61
62 pub(crate) fn add_request(&mut self, index: u64, msg: Message) {
67 let ctx = msg.get_entries()[0].get_data().to_vec();
68 if self.pending_read_index.contains_key(&ctx) {
69 return;
70 }
71 let ris = ReadIndexStatus {
72 index,
73 req: msg,
74 acks: HashSet::new(),
75 };
76 self.pending_read_index.insert(ctx.clone(), ris);
77 self.read_index_queue.push_back(ctx);
78 }
79
80 pub(crate) fn recv_ack(&mut self, msg: &Message) -> usize {
84 if let Some(rs) = self.pending_read_index.get_mut(msg.get_context()) {
85 rs.acks.insert(msg.get_from());
86 rs.acks.len() + 1
87 } else {
88 0
89 }
90 }
91
92 pub(crate) fn advance(&mut self, msg: &Message) -> Vec<ReadIndexStatus> {
96 let mut rss = vec![];
97 if let Some(i) = self.read_index_queue.iter().position(|x| {
98 if !self.pending_read_index.contains_key(x) {
99 panic!("cannot find correspond read state from pending map");
100 }
101 *x == msg.get_context()
102 }) {
103 for _ in 0..i + 1 {
104 let rs = self.read_index_queue.pop_front().unwrap();
105 let status = self.pending_read_index.remove(&rs).unwrap();
106 rss.push(status);
107 }
108 }
109 rss
110 }
111}