use std::sync::Arc;
use std::sync::atomic::Ordering::{Acquire, Relaxed, Release};
use std::sync::atomic::{AtomicI32, AtomicI64, AtomicU32, AtomicU64, AtomicUsize};
use yo_common::lock::Lock;
pub(super) const SUBSCRIBED: u32 = 1;
pub(super) const IN_MULTI: u32 = 2;
pub(super) const UNIX: u32 = 4;
pub(super) const NO_EVICT: u32 = 8;
pub(super) const NO_TOUCH: u32 = 16;
pub(super) const KILLED: u32 = 32;
pub(super) const REAPED: u32 = 64;
#[derive(Default)]
pub(super) struct Text {
pub(super) peer: Vec<u8>,
pub(super) local: Vec<u8>,
pub(super) name: Vec<u8>,
pub(super) lib_name: Vec<u8>,
pub(super) lib_ver: Vec<u8>,
pub(super) sub: Vec<u8>,
}
pub struct Client {
pub(super) id: u64,
pub(super) conn: AtomicU32,
pub(super) thread: AtomicUsize,
pub(super) since_ms: AtomicU64,
pub(super) fd: AtomicI32,
pub(super) text: Lock<Text>,
pub(super) last_ms: AtomicU64,
pub(super) net_in: AtomicU64,
pub(super) net_out: AtomicU64,
pub(super) cmds: AtomicU64,
pub(super) reads: AtomicU64,
pub(super) qbuf: AtomicU64,
pub(super) qbuf_free: AtomicU64,
pub(super) rbs: AtomicU64,
pub(super) rbp: AtomicU64,
pub(super) obl: AtomicU64,
pub(super) argv_mem: AtomicU64,
pub(super) spec: AtomicU32,
pub(super) has_sub: AtomicU32,
pub(super) db: AtomicU32,
pub(super) resp: AtomicU32,
pub(super) sub: AtomicU32,
pub(super) psub: AtomicU32,
pub(super) ssub: AtomicU32,
pub(super) watch: AtomicU32,
pub(super) multi: AtomicI64,
pub(super) multi_mem: AtomicU64,
pub(super) flags: AtomicU32,
}
impl Client {
pub(super) fn new(id: u64) -> Client {
Client {
id,
conn: AtomicU32::new(u32::MAX),
thread: AtomicUsize::new(0),
since_ms: AtomicU64::new(0),
fd: AtomicI32::new(-1),
text: Lock::default(),
last_ms: AtomicU64::new(0),
net_in: AtomicU64::new(0),
net_out: AtomicU64::new(0),
cmds: AtomicU64::new(0),
reads: AtomicU64::new(0),
qbuf: AtomicU64::new(0),
qbuf_free: AtomicU64::new(0),
rbs: AtomicU64::new(0),
rbp: AtomicU64::new(0),
obl: AtomicU64::new(0),
argv_mem: AtomicU64::new(0),
spec: AtomicU32::new(u32::MAX),
has_sub: AtomicU32::new(0),
db: AtomicU32::new(0),
resp: AtomicU32::new(2),
sub: AtomicU32::new(0),
psub: AtomicU32::new(0),
ssub: AtomicU32::new(0),
watch: AtomicU32::new(0),
multi: AtomicI64::new(-1),
multi_mem: AtomicU64::new(0),
flags: AtomicU32::new(0),
}
}
pub(super) fn set_flag(&self, bit: u32, on: bool) {
let was = self.flags.load(Relaxed);
let now = if on { was | bit } else { was & !bit };
if now != was {
self.flags.store(now, Relaxed);
}
}
pub(super) fn flag(&self, bit: u32) -> bool {
self.flags.load(Relaxed) & bit != 0
}
pub(super) fn kill(&self) -> bool {
let was = self.flags.load(Relaxed);
if was & KILLED != 0 {
return false;
}
self.flags.store(was | KILLED, Release);
true
}
pub(super) fn killed(&self) -> bool {
self.flags.load(Acquire) & KILLED != 0
}
pub(super) fn note_command(&self, at: usize, sub: Option<&[u8]>) {
match sub {
Some(sub) => {
yo_alloc::allow(|| {
let mut text = self.text.lock();
text.sub.clear();
text.sub.extend_from_slice(sub);
});
self.spec.store(at as u32, Relaxed);
self.has_sub.store(1, Release);
}
None => {
self.has_sub.store(0, Relaxed);
self.spec.store(at as u32, Relaxed);
}
}
}
pub(super) fn set_text(&self, pick: fn(&mut Text) -> &mut Vec<u8>, value: &[u8]) {
yo_alloc::allow(|| {
let mut text = self.text.lock();
let into = pick(&mut text);
into.clear();
into.extend_from_slice(value);
});
}
}
#[derive(Default)]
pub(super) struct Clients {
rows: Vec<Arc<Client>>,
}
impl Clients {
fn add(&mut self, row: &Arc<Client>) {
yo_alloc::allow(|| self.rows.push(Arc::clone(row)));
}
fn remove(&mut self, id: u64) {
if let Some(at) = self.rows.iter().position(|row| row.id == id) {
self.rows.remove(at);
}
}
fn len(&self) -> usize {
self.rows.len()
}
}
impl super::Server {
pub(crate) fn register_client(&self, row: &Arc<Client>) {
row.thread.store(self.my_slot(), Relaxed);
self.clients.lock().add(row);
}
pub(crate) fn forget_client(&self, id: u64) {
self.clients.lock().remove(id);
}
pub(super) fn client_rows(&self) -> Vec<Arc<Client>> {
let rows = self.clients.lock();
yo_alloc::allow(|| rows.rows.clone())
}
#[must_use]
pub fn client_count(&self) -> usize {
self.clients.lock().len()
}
pub(super) fn note_kills(&self, n: usize) {
if n != 0 {
self.kills.fetch_add(n, Release);
}
}
#[must_use]
pub fn kills(&self) -> usize {
self.kills.load(Acquire)
}
pub fn kill_done(&self) {
self.kills.fetch_sub(1, Release);
}
pub fn my_kills(&self) -> Vec<(u32, u64)> {
let mine = self.my_slot();
let rows = self.clients.lock();
yo_alloc::allow(|| {
rows.rows
.iter()
.filter(|row| row.thread.load(Relaxed) == mine && row.killed() && !row.flag(REAPED))
.inspect(|row| row.set_flag(REAPED, true))
.map(|row| (row.conn.load(Relaxed), row.id))
.collect()
})
}
pub fn pause(&self, until_ms: u64, all: bool) {
let until_ms = until_ms.min(u64::MAX >> 1);
let want = (until_ms << 1) | u64::from(all);
let mut have = self.pause.load(Relaxed);
loop {
let live = have != 0 && (have >> 1) > self.now_ms();
let next = if live {
((have >> 1).max(until_ms) << 1) | (have & 1) | u64::from(all)
} else {
want
};
match self
.pause
.compare_exchange_weak(have, next, Release, Relaxed)
{
Ok(_) => return,
Err(seen) => have = seen,
}
}
}
pub fn unpause(&self) {
self.pause.store(0, Release);
}
#[must_use]
pub fn paused(&self, now_ms: u64) -> Option<bool> {
let word = self.pause.load(Relaxed);
if word == 0 || (word >> 1) <= now_ms {
return None;
}
Some(word & 1 == 1)
}
#[must_use]
pub fn pause_ends(&self) -> u64 {
self.pause.load(Relaxed) >> 1
}
}