1
  2
  3
  4
  5
  6
  7
  8
  9
 10
 11
 12
 13
 14
 15
 16
 17
 18
 19
 20
 21
 22
 23
 24
 25
 26
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
use std::fmt::Display;
use std::hash::Hash;
use std::marker::PhantomData;

use tokio::sync::mpsc;
use tokio::time::Duration;

use super::{MutableServiceHandle, ServiceHandle, Shard, ShardConfig};

use crate::{DefaultCommitPolicy, MutableUpstreamFactory, ServiceData, UpstreamFactory};

mod defaults {
    use tokio::time::Duration;

    pub const PERSIST_QUEUE_CAPACITY: usize = 0;
    pub const SHARD_SENDER_CHANNEL_BUFFER_SIZE: usize = 4096;
    pub const LRU_CANDIDATES_NUM_PROBES: u16 = 5;
    pub const HANDLE_POOL_CAPACITY: usize = 128;
    pub const CACHE_EXPIRATION_PROBE_INTERVAL: Duration = Duration::from_millis(100);
    pub const CACHE_EXPIRATION_PROBE_KEYS_PER_TICK: usize = 25;
}

pub struct ServiceBuilder<Key, Data> {
    max_data_capacity: usize,
    num_shards: Option<u8>,
    shard_persist_queue_capacity: Option<usize>,
    shard_sender_channel_buffer_size: Option<usize>,
    lru_candidates_num_probes: Option<u16>,
    handle_pool_capacity: Option<usize>,
    cache_expiration_probe_interval: Option<Duration>,
    cache_expiration_probe_keys_per_tick: Option<usize>,
    _phantom: PhantomData<(Key, Data)>,
}

pub fn service_builder<Key: Send + Clone + Hash + Eq + Display + 'static, Data: ServiceData>(
    max_data_capacity: usize,
) -> ServiceBuilder<Key, Data> {
    ServiceBuilder::new(max_data_capacity)
}

impl<Key: Send + Clone + Hash + Eq + Display + 'static, Data: ServiceData>
    ServiceBuilder<Key, Data>
{
    fn new(max_data_capacity: usize) -> Self {
        assert!(max_data_capacity > 0, "max_data_capacity must be > 0");

        Self {
            max_data_capacity,
            num_shards: None,
            shard_persist_queue_capacity: None,
            shard_sender_channel_buffer_size: None,
            lru_candidates_num_probes: None,
            handle_pool_capacity: None,
            cache_expiration_probe_interval: None,
            cache_expiration_probe_keys_per_tick: None,
            _phantom: PhantomData,
        }
    }

    pub fn num_shards(mut self, num_shards: u8) -> Self {
        assert!(num_shards > 0, "num_shards must be > 0");
        self.num_shards = Some(num_shards);
        self
    }

    pub fn shard_persist_queue_capacity(mut self, shard_persist_queue_capacity: usize) -> Self {
        self.shard_persist_queue_capacity = Some(shard_persist_queue_capacity);
        self
    }

    pub fn lru_candidates_num_probes(mut self, lru_candidates_num_probes: u16) -> Self {
        self.lru_candidates_num_probes = Some(lru_candidates_num_probes);
        self
    }

    pub fn shard_sender_channel_buffer_size(
        mut self,
        shard_sender_channel_buffer_size: usize,
    ) -> Self {
        self.shard_sender_channel_buffer_size = Some(shard_sender_channel_buffer_size);
        self
    }

    pub fn cache_expiration_probe_interval(
        mut self,
        cache_expiration_probe_interval: Duration,
    ) -> Self {
        self.cache_expiration_probe_interval = Some(cache_expiration_probe_interval);
        self
    }
    pub fn cache_expiration_probe_keys_per_tick(
        mut self,
        cache_expiration_probe_keys_per_tick: usize,
    ) -> Self {
        self.cache_expiration_probe_keys_per_tick = Some(cache_expiration_probe_keys_per_tick);
        self
    }

    pub fn build<UpstreamFactoryT: UpstreamFactory<Key, Data>>(
        self,
        mut upstream_factory: UpstreamFactoryT,
    ) -> ServiceHandle<Key, Data> {
        let num_shards = self.get_num_shards();
        let shard_sender_channel_buffer_size = self.get_shard_sender_buffer_size();

        let mut shard_senders = Vec::with_capacity(num_shards as usize);

        for shard_id in 0..num_shards {
            let (sender, receiver) = mpsc::channel(shard_sender_channel_buffer_size);
            let upstream = upstream_factory.create();
            let shard_config = self.get_shard_config(shard_id, None);
            let shard = Shard::new(
                receiver,
                super::marker::Immutable::new(upstream),
                shard_config,
            );
            tokio::spawn(shard.run());
            shard_senders.push(sender);
        }

        ServiceHandle::with_senders_and_pool_capacity(
            shard_senders,
            self.get_handle_pool_capacity(),
        )
    }

    pub fn build_mutable<UpstreamFactoryT: MutableUpstreamFactory<Key, Data>>(
        self,
        mut upstream_factory: UpstreamFactoryT,
        default_commit_policy: DefaultCommitPolicy,
    ) -> MutableServiceHandle<Key, Data> {
        let num_shards = self.get_num_shards();
        let shard_sender_channel_buffer_size = self.get_shard_sender_buffer_size();

        let mut shard_senders = Vec::with_capacity(num_shards as usize);

        for shard_id in 0..num_shards {
            let (sender, receiver) = mpsc::channel(shard_sender_channel_buffer_size);
            let upstream = upstream_factory.create();
            let shard_config = self.get_shard_config(shard_id, Some(default_commit_policy));
            let shard = Shard::new(
                receiver,
                super::marker::Mutable::new(upstream),
                shard_config,
            );
            tokio::spawn(shard.run());
            shard_senders.push(sender);
        }

        MutableServiceHandle::from_service_handle(ServiceHandle::with_senders_and_pool_capacity(
            shard_senders,
            self.get_handle_pool_capacity(),
        ))
    }

    fn get_num_shards(&self) -> u8 {
        self.num_shards.unwrap_or_else(|| num_cpus::get() as u8)
    }

    fn get_handle_pool_capacity(&self) -> usize {
        self.handle_pool_capacity
            .unwrap_or(defaults::HANDLE_POOL_CAPACITY)
    }

    fn get_shard_sender_buffer_size(&self) -> usize {
        self.shard_sender_channel_buffer_size
            .unwrap_or(defaults::SHARD_SENDER_CHANNEL_BUFFER_SIZE)
    }

    fn get_shard_config(
        &self,
        shard_id: u8,
        default_commit_policy: Option<DefaultCommitPolicy>,
    ) -> ShardConfig {
        let shard_persist_queue_capacity = self
            .shard_persist_queue_capacity
            .unwrap_or(defaults::PERSIST_QUEUE_CAPACITY);
        let max_data_capacity = self.max_data_capacity / (self.get_num_shards() as usize);
        let lru_candidates_num_probes = self
            .lru_candidates_num_probes
            .unwrap_or(defaults::LRU_CANDIDATES_NUM_PROBES);
        let cache_expiration_probe_interval = self
            .cache_expiration_probe_interval
            .unwrap_or(defaults::CACHE_EXPIRATION_PROBE_INTERVAL);
        let cache_expiration_probe_keys_per_tick = self
            .cache_expiration_probe_keys_per_tick
            .unwrap_or(defaults::CACHE_EXPIRATION_PROBE_KEYS_PER_TICK);

        ShardConfig {
            shard_id,
            max_data_capacity,
            persist_queue_capacity: shard_persist_queue_capacity,
            lru_candidates_num_probes,
            default_commit_policy,
            cache_expiration_probe_interval,
            cache_expiration_probe_keys_per_tick,
        }
    }
}