use std::sync::Arc;
use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering::{Relaxed, Release};
use super::args::Args;
use super::clients::{self, Client};
use super::pubsub::Envelope;
use super::{Server, Session};
use std::io::Write;
use yo_common::lock::Lock;
pub(super) const SCRIPTS: [&str; 6] = [
"eval",
"eval_ro",
"evalsha",
"evalsha_ro",
"fcall",
"fcall_ro",
];
pub(super) fn replica() -> yo_common::Error {
yo_common::Error::new(
yo_common::Code::Invalid,
"Replica can't interact with the keyspace",
)
}
#[derive(Default)]
pub(crate) struct Monitors {
rows: Lock<Vec<Arc<Client>>>,
live: AtomicUsize,
}
impl Server {
#[must_use]
pub(crate) fn monitored(&self) -> bool {
self.monitors.live.load(Relaxed) != 0
}
pub(crate) fn watch_all(&self, row: &Arc<Client>) -> bool {
let mut rows = self.monitors.rows.lock();
if rows.iter().any(|r| r.id == row.id) {
return false;
}
yo_alloc::allow(|| rows.push(Arc::clone(row)));
self.monitors.live.store(rows.len(), Release);
row.set_flag(clients::MONITOR, true);
self.note_here(row.thread.load(Relaxed), 1);
true
}
pub(crate) fn watch_no_more(&self, row: &Arc<Client>) {
let mut rows = self.monitors.rows.lock();
let Some(at) = rows.iter().position(|r| r.id == row.id) else {
return;
};
rows.remove(at);
self.monitors.live.store(rows.len(), Release);
row.set_flag(clients::MONITOR, false);
self.note_here(row.thread.load(Relaxed), -1);
}
fn monitor_rows(&self) -> Vec<Arc<Client>> {
let rows = self.monitors.rows.lock();
yo_alloc::allow(|| rows.clone())
}
fn now_us(&self) -> u64 {
self.clock.now_us()
}
}
pub(super) fn feed(server: &Server, session: &Session, args: Args<'_>) {
let rows = server.monitor_rows();
if rows.is_empty() {
return;
}
let line = yo_alloc::allow(|| Arc::new(render(server, session, args)));
for row in rows {
server.post(
row.thread.load(Relaxed),
Envelope::line(row.conn.load(Relaxed), row.id, Arc::clone(&line)),
);
}
}
fn render(server: &Server, session: &Session, args: Args<'_>) -> Vec<u8> {
yo_alloc::allow(|| {
let mut line = Vec::with_capacity(64 + args.len() * 16);
let us = server.now_us();
let _ = write!(
line,
"{}.{:06} [{} ",
us / 1_000_000,
us % 1_000_000,
session.db
);
who(session, &mut line);
line.extend_from_slice(b"] ");
let hidden = redacted(args);
for i in 0..args.len() {
if i != 0 {
line.push(b' ');
}
if hidden & (1u32 << i.min(31)) != 0 {
line.extend_from_slice(b"\"(redacted)\"");
} else {
quote(args.get(i), &mut line);
}
}
line
})
}
fn who(session: &Session, line: &mut Vec<u8>) {
if session.scripted() {
line.extend_from_slice(b"lua");
return;
}
let row = session.row();
let text = row.text.lock();
if text.peer.is_empty() {
line.extend_from_slice(b"?:0");
} else if row.flag(clients::UNIX) {
line.extend_from_slice(b"unix:");
let path = text.peer.strip_suffix(b":0").unwrap_or(&text.peer);
line.extend_from_slice(path);
} else {
line.extend_from_slice(&text.peer);
}
}
fn redacted(args: Args<'_>) -> u32 {
if super::args::is(args.name(), b"auth") {
return ((1u32 << args.len().min(31)) - 1) & !1;
}
if !super::args::is(args.name(), b"hello") {
return 0;
}
let mut mask = 0;
for i in 2..args.len() {
if super::args::is(args.get(i), b"AUTH") && i + 2 < args.len() {
mask |= (1 << (i + 1)) | (1 << (i + 2));
}
}
mask
}
fn quote(arg: &[u8], line: &mut Vec<u8>) {
line.push(b'"');
for &b in arg {
match b {
b'\\' | b'"' => {
line.push(b'\\');
line.push(b);
}
b'\n' => line.extend_from_slice(b"\\n"),
b'\r' => line.extend_from_slice(b"\\r"),
b'\t' => line.extend_from_slice(b"\\t"),
0x07 => line.extend_from_slice(b"\\a"),
0x08 => line.extend_from_slice(b"\\b"),
0x20..=0x7e => line.push(b),
other => {
let _ = write!(line, "\\x{other:02x}");
}
}
}
line.push(b'"');
}
pub(super) fn hidden(spec: &super::table::Spec, args: Args<'_>) -> bool {
if spec.flags.contains(&"admin") {
return true;
}
let subs: &[&[u8]] = match spec.name {
"client" => &[b"kill", b"list", b"no-evict", b"pause", b"unpause"],
"config" => &[b"get", b"set", b"resetstat", b"rewrite"],
"backup" => &[b"start", b"seal", b"abort", b"cleanup", b"status", b"list"],
_ => return false,
};
args.len() > 1 && subs.iter().any(|sub| super::args::is(args.get(1), sub))
}