use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use super::error::{report_fault, RuntimeError};
use super::mailbox::Mailbox;
use super::process::FlowId;
use super::sync_lock;
const SHARDS: usize = 64;
pub struct Directory {
shards: Vec<Mutex<HashMap<FlowId, Arc<Mailbox>>>>,
}
impl Directory {
pub fn new() -> Self {
let mut shards = Vec::with_capacity(SHARDS);
for _ in 0..SHARDS {
shards.push(Mutex::new(HashMap::new()));
}
Directory { shards }
}
#[inline]
fn shard_for(&self, id: FlowId) -> &Mutex<HashMap<FlowId, Arc<Mailbox>>> {
&self.shards[(id.as_u64() as usize) % SHARDS]
}
pub fn register(&self, id: FlowId, mailbox: Arc<Mailbox>) -> Result<(), RuntimeError> {
sync_lock::lock(self.shard_for(id), "Directory::register")?.insert(id, mailbox);
Ok(())
}
pub fn unregister(&self, id: FlowId) -> Result<(), RuntimeError> {
sync_lock::lock(self.shard_for(id), "Directory::unregister")?.remove(&id);
Ok(())
}
pub fn lookup(&self, id: FlowId) -> Result<Option<Arc<Mailbox>>, RuntimeError> {
Ok(sync_lock::lock(self.shard_for(id), "Directory::lookup")?
.get(&id)
.cloned())
}
pub fn len(&self) -> usize {
let mut n = 0;
for s in &self.shards {
match sync_lock::lock(s, "Directory::len") {
Ok(g) => n += g.len(),
Err(e) => {
report_fault(e);
return n;
}
}
}
n
}
}
impl Default for Directory {
fn default() -> Self {
Self::new()
}
}