use std::sync::Arc;
use yo_reactor::{BATCH_MAX, Engine, Reactor};
use crate::dispatch::table;
use crate::dispatch::{self, Flow, Parked, Server};
use crate::front::{Front, Wrote};
use crate::proto::Limits;
use yo_kv::Keyspace;
pub use crate::front::Cmd;
pub type ConnId = u32;
const SWEEP_LOOKS: usize = yo_reactor::MAINTENANCE_UNITS as usize;
pub trait Sink {
fn write(&mut self, conn: ConnId, bytes: &[u8]) -> usize;
fn closed(&mut self, conn: ConnId) {
let _ = conn;
}
}
#[derive(Debug, Default)]
pub struct Recorder {
sent: Vec<Vec<u8>>,
closed: Vec<ConnId>,
}
impl Recorder {
#[must_use]
pub fn new() -> Recorder {
Recorder::default()
}
#[must_use]
pub fn sent(&self, conn: ConnId) -> &[u8] {
self.sent.get(conn as usize).map_or(&[], Vec::as_slice)
}
#[must_use]
pub fn was_closed(&self, conn: ConnId) -> bool {
self.closed.contains(&conn)
}
pub fn clear(&mut self) {
for c in &mut self.sent {
c.clear();
}
self.closed.clear();
}
}
impl Sink for Recorder {
fn write(&mut self, conn: ConnId, bytes: &[u8]) -> usize {
yo_alloc::allow(|| {
if self.sent.len() <= conn as usize {
self.sent.resize_with(conn as usize + 1, Vec::new);
}
self.sent[conn as usize].extend_from_slice(bytes);
});
bytes.len()
}
fn closed(&mut self, conn: ConnId) {
yo_alloc::allow(|| self.closed.push(conn));
}
}
pub struct Wire<S> {
front: Front<S>,
server: Arc<Server>,
parked: Vec<Parked>,
post: Vec<dispatch::Envelope>,
}
impl<S: Sink> Wire<S> {
#[must_use]
pub fn new(sink: S) -> Wire<S> {
Wire::with_server(Server::new(), sink)
}
#[must_use]
pub fn with_server(server: Server, sink: S) -> Wire<S> {
Wire::over(Arc::new(server), sink)
}
#[must_use]
pub fn over(server: Arc<Server>, sink: S) -> Wire<S> {
Wire {
front: Front::new(sink),
parked: Vec::new(),
post: Vec::new(),
server,
}
}
#[must_use]
pub fn server(&self) -> &Server {
&self.server
}
#[must_use]
pub fn shared(&self) -> Arc<Server> {
Arc::clone(&self.server)
}
pub fn server_mut(&mut self) -> &mut Server {
Arc::get_mut(&mut self.server)
.expect("the server is set up before the threads that share it are started")
}
#[must_use]
pub const fn sink(&self) -> &S {
self.front.sink()
}
pub const fn sink_mut(&mut self) -> &mut S {
self.front.sink_mut()
}
pub fn set_limits(&mut self, limits: Limits) {
self.front.set_limits(limits);
}
pub fn accept(&mut self) -> ConnId {
self.server.counted().opened();
let at = self.front.open(self.server.next_client());
self.note_buffers();
at
}
fn note_buffers(&mut self) {
let delta = self.front.buffer_delta();
if delta != 0 {
self.server.note_conn_bytes(delta);
}
}
pub fn hangup(&mut self, conn: ConnId) {
if !self.front.live(conn) {
return;
}
self.front.mark_gone(conn);
if self.front.blocked(conn) {
self.front.unpark(conn);
}
if self.front.pending(conn) == 0 {
self.release(conn);
}
self.note_buffers();
}
fn serve_waiters(&mut self) {
let now = self.server.now_ms();
let mine = self.server.my_slot();
self.server.waiters().mine(mine, &mut self.parked);
for at in 0..self.parked.len() {
let p = self.parked[at];
if !self.front.answers(p.conn, p.client) {
self.server.forget_waiters(p.client);
continue;
}
let served = {
let Wire { server, front, .. } = self;
server.serve_waiter(p.client, now, front.out(p.conn))
};
if served {
self.server.forget_waiters(p.client);
self.front.unpark(p.conn);
self.front.soil(p.conn);
}
}
self.parked.clear();
}
fn deliver(&mut self) {
let mut post = core::mem::take(&mut self.post);
self.server.take_mail(&mut post);
for env in post.drain(..) {
let conn = env.conn();
if !self.front.answers(conn, env.client()) {
continue;
}
env.write(self.front.out(conn));
self.front.soil(conn);
}
self.post = post;
}
#[must_use]
pub fn clients(&self) -> usize {
self.front.clients()
}
#[must_use]
pub fn ready(&self) -> usize {
self.front.ready()
}
#[must_use]
pub fn owed(&self) -> usize {
self.front.owed()
}
#[must_use]
pub fn waiting(&self) -> usize {
self.server.parked_here()
}
#[must_use]
pub fn posted(&self) -> usize {
self.server.posted()
}
#[must_use]
pub fn stopping(&self) -> bool {
self.server.stopping()
}
#[must_use]
pub fn decoders(&self) -> usize {
self.front.decoders()
}
#[must_use]
pub fn buffer_bytes(&self) -> usize {
self.front.buffer_bytes()
}
pub fn feed(&mut self, conn: ConnId, bytes: &[u8]) {
self.front.feed(conn, bytes);
self.note_buffers();
}
fn release(&mut self, conn: ConnId) {
if let Some(session) = self.front.session_mut(conn) {
dispatch::forget_session(&self.server, session);
}
let Some(client) = self.front.close(conn) else {
return;
};
self.forget(client);
}
fn forget(&mut self, client: u64) {
self.server.forget_waiters(client);
self.server.counted().closed();
}
pub fn take_ready(&mut self, into: &mut Vec<Cmd>, max: usize) -> usize {
self.front.take_ready(into, max)
}
pub fn tick(&mut self) {
self.server.refresh_clock();
}
pub fn maintain(&mut self) -> Option<usize> {
self.server.refresh_memory();
self.server.backup_expire();
self.server.expire_slice(SWEEP_LOOKS);
self.server.compact_step()
}
}
impl<S: Sink> Engine for Wire<S> {
type Work = Cmd;
fn key_hash(&self, cmd: &Cmd) -> Option<u64> {
let spec = table::at(cmd.spec)?;
if spec.first_key <= 0 {
return None;
}
let args = self.front.args(cmd);
let key = args.opt(spec.first_key as usize)?;
Some(Keyspace::hash_of(key))
}
fn prefetch(&self, cmd: &Cmd, hash: u64) {
let db = self.front.db(cmd.conn());
self.server.striped_ref(db).prefetch_hashed(hash);
}
fn run(&mut self, cmd: Cmd, _hash: Option<u64>) -> yo_reactor::Flow {
let conn = cmd.conn();
if self.front.blocked(conn) {
self.front.park(conn, cmd);
return yo_reactor::Flow::Next;
}
let flow = if self.front.start(&cmd) {
let Wire { front, server, .. } = self;
let (args, session, out) = front.parts(&cmd);
let spec = table::at(cmd.spec);
dispatch::resolved(server, session, spec, args, out)
} else {
Flow::Continue
};
self.front.done(&cmd);
if self.front.gone(conn) {
if self.front.pending(conn) == 0 {
self.release(conn);
}
} else {
match flow {
Flow::Close => {
self.front.quit(conn);
self.front.soil(conn);
}
Flow::Block => {
self.front.block(conn);
let client = self.front.client(conn);
self.server.bind_waiter(client, conn);
}
Flow::Continue => self.front.soil(conn),
}
}
if self.server.parked_here() != 0 {
self.serve_waiters();
}
yo_reactor::Flow::Next
}
fn flush(&mut self) {
if self.server.parked_here() != 0 {
self.server.refresh_clock();
self.serve_waiters();
}
if self.server.mail_here() != 0 {
self.deliver();
}
let mut dirty = self.front.take_dirty();
let mut at = 0;
while at < dirty.len() {
let conn = dirty[at];
match self.front.write_out(conn) {
Wrote::Owed => at += 1,
Wrote::Done => {
dirty.swap_remove(at);
}
Wrote::Ended(client) => {
self.forget(client);
dirty.swap_remove(at);
}
}
}
self.front.give_dirty(dirty);
self.note_buffers();
}
fn maintain(&mut self, budget: &mut yo_reactor::Budget) {
if !budget.spend(1) {
return;
}
self.tick();
let looks = budget.left() as usize;
let spent = self.server.expire_slice(looks);
budget.spend(u32::try_from(spent).unwrap_or(u32::MAX));
}
}
pub fn pump<S: Sink>(reactor: &mut Reactor<Wire<S>>, batch: &mut Vec<Cmd>) -> usize {
let mut ran = 0;
reactor.engine_mut().tick();
loop {
batch.clear();
if reactor.engine_mut().take_ready(batch, BATCH_MAX) == 0 {
break;
}
let armed = yo_alloc::guard();
ran += reactor.execute_all(batch.drain(..));
drop(armed);
reactor.engine_mut().flush();
reactor.engine_mut().maintain();
}
reactor.engine_mut().maintain();
reactor.engine_mut().flush();
ran
}
#[cfg(test)]
mod tests {
use super::*;
fn wire(args: &[&[u8]]) -> Vec<u8> {
let mut b = format!("*{}\r\n", args.len()).into_bytes();
for a in args {
b.extend_from_slice(format!("${}\r\n", a.len()).as_bytes());
b.extend_from_slice(a);
b.extend_from_slice(b"\r\n");
}
b
}
fn engine() -> (Reactor<Wire<Recorder>>, ConnId, Vec<Cmd>) {
let mut r = Reactor::inline(Wire::new(Recorder::new()));
let conn = r.engine_mut().accept();
(r, conn, Vec::new())
}
const START_MS: u64 = 1_000_000;
fn timed() -> (Reactor<Wire<Recorder>>, ConnId, Vec<Cmd>) {
let server = crate::dispatch::Server::with_clock(yo_kv::Clock::fixed(START_MS));
let mut r = Reactor::inline(Wire::with_server(server, Recorder::new()));
let conn = r.engine_mut().accept();
(r, conn, Vec::new())
}
#[test]
fn a_pipelined_batch_comes_back_in_order_and_in_one_write() {
let (mut r, conn, mut batch) = engine();
let mut stream = wire(&[b"SET", b"k", b"v"]);
stream.extend(wire(&[b"GET", b"k"]));
stream.extend(wire(&[b"INCR", b"n"]));
r.engine_mut().feed(conn, &stream);
assert_eq!(r.engine().ready(), 3);
assert_eq!(pump(&mut r, &mut batch), 3);
assert_eq!(r.engine().sink().sent(conn), b"+OK\r\n$1\r\nv\r\n:1\r\n");
assert_eq!(r.engine().ready(), 0);
}
#[test]
fn a_command_split_across_reads_resumes_rather_than_restarts() {
let (mut r, conn, mut batch) = engine();
let bytes = wire(&[b"SET", b"key", b"value"]);
for at in 1..bytes.len() {
r.engine_mut().feed(conn, &bytes[at - 1..at]);
assert_eq!(r.engine().ready(), 0, "not a command yet at {at}");
}
r.engine_mut().feed(conn, &bytes[bytes.len() - 1..]);
assert_eq!(r.engine().ready(), 1);
assert_eq!(pump(&mut r, &mut batch), 1);
assert_eq!(r.engine().sink().sent(conn), b"+OK\r\n");
r.engine_mut().feed(conn, &wire(&[b"GET", b"key"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(conn), b"+OK\r\n$5\r\nvalue\r\n");
}
#[test]
fn two_connections_are_two_sessions_over_one_server() {
let (mut r, a, mut batch) = engine();
let b = r.engine_mut().accept();
r.engine_mut().feed(a, &wire(&[b"SELECT", b"3"]));
r.engine_mut().feed(a, &wire(&[b"SET", b"k", b"a"]));
r.engine_mut().feed(b, &wire(&[b"SET", b"k", b"b"]));
r.engine_mut().feed(a, &wire(&[b"GET", b"k"]));
r.engine_mut().feed(b, &wire(&[b"GET", b"k"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(a), b"+OK\r\n+OK\r\n$1\r\na\r\n");
assert_eq!(r.engine().sink().sent(b), b"+OK\r\n$1\r\nb\r\n");
assert_eq!(r.engine().clients(), 2);
}
#[test]
fn two_threads_write_into_one_server() {
const EACH: usize = 200;
let mut server = Server::new();
server.set_threads(2);
let first = Wire::with_server(server, Recorder::new());
let second = Wire::over(first.shared(), Recorder::new());
let server = first.shared();
std::thread::scope(|s| {
for (at, engine) in [first, second].into_iter().enumerate() {
s.spawn(move || {
let mut r = Reactor::inline(engine);
let mut batch = Vec::new();
let conn = r.engine_mut().accept();
for i in 0..EACH {
let key = format!("t{at}:{i}");
r.engine_mut()
.feed(conn, &wire(&[b"SET", key.as_bytes(), b"v"]));
pump(&mut r, &mut batch);
}
});
}
});
assert_eq!(server.striped_ref(0).len(), 2 * EACH);
assert_eq!(server.totals().connections, 2);
}
#[test]
fn a_waiter_belongs_to_the_thread_that_parked_it() {
let mut server = Server::new();
server.set_threads(2);
let first = Wire::with_server(server, Recorder::new());
let second = Wire::over(first.shared(), Recorder::new());
let server = first.shared();
let parked = std::sync::Barrier::new(2);
let swept = std::sync::Barrier::new(2);
std::thread::scope(|s| {
let (parked, swept) = (&parked, &swept);
s.spawn(move || {
let mut r = Reactor::inline(first);
let mut batch = Vec::new();
let conn = r.engine_mut().accept();
r.engine_mut().feed(conn, &wire(&[b"BLPOP", b"a", b"0"]));
pump(&mut r, &mut batch);
parked.wait();
for _ in 0..50 {
pump(&mut r, &mut batch);
}
swept.wait();
assert!(r.engine().sink().sent(conn).is_empty(), "nothing to say");
});
s.spawn(move || {
let mut r = Reactor::inline(second);
let mut batch = Vec::new();
let conn = r.engine_mut().accept();
r.engine_mut().feed(conn, &wire(&[b"BLPOP", b"b", b"0"]));
pump(&mut r, &mut batch);
parked.wait();
swept.wait();
let pusher = r.engine_mut().accept();
r.engine_mut().feed(pusher, &wire(&[b"RPUSH", b"b", b"v"]));
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(conn),
b"*2\r\n$1\r\nb\r\n$1\r\nv\r\n",
"served by the thread that parked it"
);
});
});
assert_eq!(server.parked(), 1, "and the other one is still waiting");
}
#[test]
fn a_thread_counts_the_clients_it_blocked_and_nobody_else_s() {
let mut server = Server::new();
server.set_threads(2);
let first = Wire::with_server(server, Recorder::new());
let second = Wire::over(first.shared(), Recorder::new());
let server = first.shared();
let parked = std::sync::Barrier::new(2);
let looked = std::sync::Barrier::new(2);
std::thread::scope(|s| {
let (parked, looked) = (&parked, &looked);
s.spawn(move || {
let mut r = Reactor::inline(first);
let mut batch = Vec::new();
let conn = r.engine_mut().accept();
r.engine_mut().feed(conn, &wire(&[b"BLPOP", b"a", b"0"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().waiting(), 1, "the one this thread blocked");
parked.wait();
looked.wait();
let other = r.engine_mut().accept();
r.engine_mut().feed(other, &wire(&[b"PING"]));
pump(&mut r, &mut batch);
r.engine_mut().hangup(other);
pump(&mut r, &mut batch);
assert_eq!(r.engine().waiting(), 1, "still just the blocked one");
});
s.spawn(move || {
let mut r = Reactor::inline(second);
let mut batch = Vec::new();
parked.wait();
pump(&mut r, &mut batch);
assert_eq!(r.engine().waiting(), 0, "none of them are this one's");
assert_eq!(r.engine().server().parked(), 1, "one on the server");
looked.wait();
});
});
assert_eq!(server.parked(), 1);
}
#[test]
fn client_ids_are_the_server_s_to_hand_out() {
let first = Wire::new(Recorder::new());
let second = Wire::over(first.shared(), Recorder::new());
let mut a = Reactor::inline(first);
let mut b = Reactor::inline(second);
let (one, two) = (a.engine_mut().accept(), b.engine_mut().accept());
assert_eq!(one, two, "the same slot on each front");
let mut batch = Vec::new();
a.engine_mut().feed(one, &wire(&[b"HELLO", b"3"]));
b.engine_mut().feed(two, &wire(&[b"HELLO", b"3"]));
pump(&mut a, &mut batch);
pump(&mut b, &mut batch);
let first = String::from_utf8_lossy(a.engine().sink().sent(one)).into_owned();
let second = String::from_utf8_lossy(b.engine().sink().sent(two)).into_owned();
assert!(first.contains(":1\r\n"), "{first}");
assert!(second.contains(":2\r\n"), "{second}");
}
#[test]
fn quit_is_answered_and_then_the_connection_goes() {
let (mut r, conn, mut batch) = engine();
r.engine_mut().feed(conn, &wire(&[b"PING"]));
r.engine_mut().feed(conn, &wire(&[b"QUIT"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(conn), b"+PONG\r\n+OK\r\n");
assert!(r.engine().sink().was_closed(conn));
assert_eq!(r.engine().clients(), 0);
let again = r.engine_mut().accept();
assert_eq!(again, conn);
assert_eq!(r.engine().clients(), 1);
}
#[test]
fn what_a_client_pipelined_behind_quit_is_never_run() {
let (mut r, conn, mut batch) = engine();
let mut stream = wire(&[b"QUIT"]);
stream.extend(wire(&[b"SET", b"foo", b"bar"]));
r.engine_mut().feed(conn, &stream);
assert_eq!(r.engine().ready(), 2);
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(conn), b"+OK\r\n");
assert!(r.engine().sink().was_closed(conn));
r.engine_mut().sink_mut().clear();
let next = r.engine_mut().accept();
r.engine_mut().feed(next, &wire(&[b"GET", b"foo"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(next), b"$-1\r\n");
}
#[test]
fn a_slot_that_last_spoke_resp3_answers_the_next_client_in_resp2() {
let (mut r, conn, mut batch) = engine();
r.engine_mut().feed(conn, &wire(&[b"HELLO", b"3"]));
r.engine_mut().feed(conn, &wire(&[b"GET", b"nothing"]));
pump(&mut r, &mut batch);
assert!(r.engine().sink().sent(conn).ends_with(b"_\r\n"));
r.engine_mut().feed(conn, &wire(&[b"QUIT"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
let next = r.engine_mut().accept();
assert_eq!(next, conn, "the same slot, which is what this is about");
r.engine_mut().feed(next, &wire(&[b"GET", b"nothing"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(next), b"$-1\r\n");
}
#[test]
fn commands_that_arrived_before_a_protocol_error_are_still_answered() {
let (mut r, conn, mut batch) = engine();
let mut stream = wire(&[b"SET", b"k", b"v"]);
stream.extend(wire(&[b"GET", b"k"]));
stream.extend_from_slice(b"*1\r\n+notabulk\r\n");
r.engine_mut().feed(conn, &stream);
pump(&mut r, &mut batch);
let sent = r.engine().sink().sent(conn);
assert!(
sent.starts_with(b"+OK\r\n$1\r\nv\r\n-ERR Protocol error: "),
"{sent:?}"
);
assert!(r.engine().sink().was_closed(conn));
}
#[test]
fn a_protocol_error_is_written_and_closes_the_connection() {
let (mut r, conn, mut batch) = engine();
r.engine_mut().feed(conn, b"*1\r\n+notabulk\r\n");
pump(&mut r, &mut batch);
let sent = r.engine().sink().sent(conn);
assert!(sent.starts_with(b"-ERR Protocol error: "), "{sent:?}");
assert!(r.engine().sink().was_closed(conn));
assert_eq!(r.engine().clients(), 0);
}
#[test]
fn a_decoder_that_came_back_mid_command_starts_the_next_one_clean() {
let (mut r, conn, mut batch) = engine();
r.engine_mut()
.feed(conn, b"*3\r\n$3\r\nSET\r\n$1\r\nx\r\n$blabla\r\n");
pump(&mut r, &mut batch);
let sent = r.engine().sink().sent(conn);
assert!(
sent.starts_with(b"-ERR Protocol error: invalid bulk length"),
"{sent:?}"
);
r.engine_mut().sink_mut().clear();
let next = r.engine_mut().accept();
r.engine_mut().feed(next, &wire(&[b"GET", b"k"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(next), b"$-1\r\n");
r.engine_mut().sink_mut().clear();
let third = r.engine_mut().accept();
r.engine_mut().feed(third, b"*1\r\n+notabulk\r\n");
pump(&mut r, &mut batch);
let sent = r.engine().sink().sent(third);
assert!(sent.starts_with(b"-ERR Protocol error: "), "{sent:?}");
}
#[test]
fn a_hangup_with_commands_in_flight_waits_for_them() {
let (mut r, conn, mut batch) = engine();
r.engine_mut().feed(conn, &wire(&[b"SET", b"k", b"v"]));
r.engine_mut().feed(conn, &wire(&[b"GET", b"k"]));
batch.clear();
r.engine_mut().take_ready(&mut batch, BATCH_MAX);
r.engine_mut().hangup(conn);
assert_eq!(r.engine().clients(), 1, "still holding the buffer");
r.execute_all(batch.drain(..));
r.engine_mut().flush();
assert_eq!(r.engine().clients(), 0);
assert!(r.engine().sink().sent(conn).is_empty(), "nobody to answer");
let decoders = r.engine().decoders();
let again = r.engine_mut().accept();
assert_eq!(again, conn);
r.engine_mut().feed(again, &wire(&[b"PING"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(again), b"+PONG\r\n");
assert_eq!(r.engine().decoders(), decoders);
}
#[test]
fn the_buffers_and_the_decoder_pool_stop_growing() {
let (mut r, conn, mut batch) = engine();
let mut stream = Vec::new();
for i in 0..32 {
stream.extend(wire(&[b"SET", format!("k{i}").as_bytes(), b"v"]));
}
r.engine_mut().feed(conn, &stream);
pump(&mut r, &mut batch);
let decoders = r.engine().decoders();
let batch_cap = batch.capacity();
for _ in 0..10 {
r.engine_mut().feed(conn, &stream);
pump(&mut r, &mut batch);
}
assert_eq!(r.engine().decoders(), decoders, "the pool is reused");
assert_eq!(batch.capacity(), batch_cap, "the batch buffer is reused");
assert!(
decoders <= BATCH_MAX + 1,
"{decoders} decoders for 32 commands"
);
}
#[test]
fn a_pipelining_client_does_not_grow_the_read_buffer() {
let (mut r, conn, mut batch) = engine();
let mut round = Vec::new();
for i in 0..16 {
round.extend(wire(&[b"SET", format!("k{i}").as_bytes(), b"v"]));
}
r.engine_mut().feed(conn, &round);
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
let after_one = r.engine().buffer_bytes();
let rounds = if cfg!(miri) { 50 } else { 1000 };
for _ in 0..rounds {
r.engine_mut().feed(conn, &round);
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
}
assert_eq!(
r.engine().buffer_bytes(),
after_one,
"the buffers grew over {rounds} rounds of the same sixteen commands"
);
assert!(
r.engine().server().memory_bytes() >= after_one,
"the buffers are counted in what the server reports"
);
}
#[test]
fn a_command_split_across_reads_survives_compaction() {
let (mut r, conn, mut batch) = engine();
let cmd = wire(&[b"SET", b"key", b"value"]);
let (head, tail) = cmd.split_at(cmd.len() - 4);
r.engine_mut().feed(conn, &wire(&[b"PING"]));
r.engine_mut().feed(conn, head);
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(conn), b"+PONG\r\n");
r.engine_mut().feed(conn, tail);
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(conn), b"+PONG\r\n+OK\r\n");
r.engine_mut().feed(conn, &wire(&[b"GET", b"key"]));
pump(&mut r, &mut batch);
assert!(r.engine().sink().sent(conn).ends_with(b"$5\r\nvalue\r\n"));
}
#[test]
fn the_batch_goes_through_the_reactors_two_walks() {
let (mut r, conn, mut batch) = engine();
for i in 0..100 {
r.engine_mut()
.feed(conn, &wire(&[b"INCR", format!("k{}", i % 7).as_bytes()]));
}
let ran = pump(&mut r, &mut batch);
assert_eq!(ran, 100);
assert_eq!(r.commands(), 100);
assert_eq!(r.turns(), 2);
assert!(r.engine().sink().sent(conn).ends_with(b":15\r\n"));
}
#[derive(Default)]
struct Trickle {
sent: Vec<u8>,
writes: usize,
}
impl Sink for Trickle {
fn write(&mut self, _conn: ConnId, bytes: &[u8]) -> usize {
self.writes += 1;
let n = bytes.len().min(4);
self.sent.extend_from_slice(&bytes[..n]);
n
}
}
#[test]
fn a_blpop_on_a_list_with_something_in_it_never_waits() {
let (mut r, conn, mut batch) = engine();
r.engine_mut().feed(conn, &wire(&[b"RPUSH", b"q", b"a"]));
r.engine_mut().feed(conn, &wire(&[b"BLPOP", b"q", b"0"]));
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(conn),
b":1\r\n*2\r\n$1\r\nq\r\n$1\r\na\r\n"
);
assert_eq!(r.engine().server().parked(), 0);
}
#[test]
fn a_parked_client_is_answered_by_another_connections_push() {
let (mut r, a, mut batch) = engine();
let b = r.engine_mut().accept();
r.engine_mut().feed(a, &wire(&[b"BLPOP", b"q", b"0"]));
pump(&mut r, &mut batch);
assert!(r.engine().sink().sent(a).is_empty(), "nothing to say yet");
assert_eq!(r.engine().server().parked(), 1);
r.engine_mut().feed(b, &wire(&[b"RPUSH", b"q", b"one"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(a), b"*2\r\n$1\r\nq\r\n$3\r\none\r\n");
assert_eq!(r.engine().sink().sent(b), b":1\r\n");
assert_eq!(r.engine().server().parked(), 0);
}
#[test]
fn only_a_list_arriving_under_a_named_key_wakes_a_waiter() {
let (mut r, a, mut batch) = engine();
let b = r.engine_mut().accept();
r.engine_mut().feed(a, &wire(&[b"BLPOP", b"q", b"0"]));
pump(&mut r, &mut batch);
r.engine_mut()
.feed(b, &wire(&[b"RPUSH", b"elsewhere", b"x"]));
r.engine_mut().feed(b, &wire(&[b"SADD", b"q", b"x"]));
pump(&mut r, &mut batch);
assert!(r.engine().sink().sent(a).is_empty());
assert_eq!(r.engine().server().parked(), 1, "still waiting");
assert_eq!(r.engine().sink().sent(b), b":1\r\n:1\r\n");
}
#[test]
fn two_parked_clients_are_served_in_the_order_they_arrived() {
let (mut r, a, mut batch) = engine();
let b = r.engine_mut().accept();
let c = r.engine_mut().accept();
r.engine_mut().feed(a, &wire(&[b"BLPOP", b"q", b"0"]));
pump(&mut r, &mut batch);
r.engine_mut().feed(b, &wire(&[b"BLPOP", b"q", b"0"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().server().parked(), 2);
r.engine_mut()
.feed(c, &wire(&[b"RPUSH", b"q", b"first", b"second"]));
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(a),
b"*2\r\n$1\r\nq\r\n$5\r\nfirst\r\n"
);
assert_eq!(
r.engine().sink().sent(b),
b"*2\r\n$1\r\nq\r\n$6\r\nsecond\r\n"
);
assert_eq!(r.engine().server().parked(), 0);
}
#[test]
fn what_a_client_pipelined_behind_a_block_waits_for_the_block() {
let (mut r, a, mut batch) = engine();
let b = r.engine_mut().accept();
let mut stream = wire(&[b"BLPOP", b"q", b"0"]);
stream.extend(wire(&[b"PING"]));
r.engine_mut().feed(a, &stream);
pump(&mut r, &mut batch);
assert!(
r.engine().sink().sent(a).is_empty(),
"the PING went out in front of the answer it was sent behind"
);
r.engine_mut().feed(a, &wire(&[b"ECHO", b"after"]));
pump(&mut r, &mut batch);
assert!(r.engine().sink().sent(a).is_empty());
r.engine_mut().feed(b, &wire(&[b"RPUSH", b"q", b"x"]));
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(a),
b"*2\r\n$1\r\nq\r\n$1\r\nx\r\n+PONG\r\n$5\r\nafter\r\n"
);
}
#[test]
fn a_waiter_is_served_between_two_pipelined_pushes() {
let (mut r, a, mut batch) = engine();
let b = r.engine_mut().accept();
r.engine_mut()
.feed(a, &wire(&[b"BLPOP", b"p1", b"p2", b"0"]));
pump(&mut r, &mut batch);
let mut stream = wire(&[b"RPUSH", b"p2", b"second"]);
stream.extend(wire(&[b"RPUSH", b"p1", b"first"]));
r.engine_mut().feed(b, &stream);
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(a),
b"*2\r\n$2\r\np2\r\n$6\r\nsecond\r\n"
);
r.engine_mut()
.feed(b, &wire(&[b"LRANGE", b"p1", b"0", b"-1"]));
pump(&mut r, &mut batch);
assert!(
r.engine()
.sink()
.sent(b)
.ends_with(b"*1\r\n$5\r\nfirst\r\n")
);
}
#[test]
fn a_waiter_woken_by_another_waiter() {
let (mut r, a, mut batch) = engine();
let b = r.engine_mut().accept();
let c = r.engine_mut().accept();
r.engine_mut()
.feed(a, &wire(&[b"BLMOVE", b"x", b"y", b"LEFT", b"RIGHT", b"0"]));
pump(&mut r, &mut batch);
r.engine_mut().feed(b, &wire(&[b"BLPOP", b"y", b"0"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().server().parked(), 2);
r.engine_mut().feed(c, &wire(&[b"RPUSH", b"x", b"chain"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(a), b"$5\r\nchain\r\n");
assert_eq!(
r.engine().sink().sent(b),
b"*2\r\n$1\r\ny\r\n$5\r\nchain\r\n"
);
assert_eq!(r.engine().server().parked(), 0);
}
#[test]
fn a_waiter_is_only_woken_on_the_database_it_blocked_on() {
let (mut r, a, mut batch) = engine();
let b = r.engine_mut().accept();
r.engine_mut().feed(a, &wire(&[b"SELECT", b"3"]));
r.engine_mut().feed(a, &wire(&[b"BLPOP", b"q", b"0"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(a), b"+OK\r\n");
r.engine_mut().feed(b, &wire(&[b"RPUSH", b"q", b"wrongdb"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(a), b"+OK\r\n", "still waiting");
r.engine_mut().feed(b, &wire(&[b"SELECT", b"3"]));
r.engine_mut().feed(b, &wire(&[b"RPUSH", b"q", b"rightdb"]));
pump(&mut r, &mut batch);
assert!(r.engine().sink().sent(a).ends_with(b"$7\r\nrightdb\r\n"));
}
#[test]
fn a_client_that_waited_long_enough_gets_a_null_array() {
let (mut r, conn, mut batch) = timed();
r.engine_mut().feed(conn, &wire(&[b"BLPOP", b"q", b"30"]));
pump(&mut r, &mut batch);
assert!(r.engine().sink().sent(conn).is_empty());
r.engine_mut().server_mut().set_clock_ms(START_MS + 29_999);
pump(&mut r, &mut batch);
assert!(
r.engine().sink().sent(conn).is_empty(),
"a millisecond short"
);
r.engine_mut().server_mut().set_clock_ms(START_MS + 30_000);
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(conn), b"*-1\r\n");
assert_eq!(r.engine().server().parked(), 0);
}
#[test]
fn every_blocking_command_times_out_with_the_same_null_array() {
for cmd in [
&[b"BLPOP".as_slice(), b"q", b"0.001"][..],
&[b"BRPOP", b"q", b"0.001"],
&[b"BLMOVE", b"q", b"d", b"LEFT", b"RIGHT", b"0.001"],
&[b"BRPOPLPUSH", b"q", b"d", b"0.001"],
&[b"BLMPOP", b"0.001", b"1", b"q", b"LEFT"],
] {
let (mut r, conn, mut batch) = timed();
r.engine_mut().feed(conn, &wire(cmd));
pump(&mut r, &mut batch);
r.engine_mut().server_mut().set_clock_ms(START_MS + 1);
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(conn), b"*-1\r\n", "for {cmd:?}");
}
}
#[test]
fn a_waiter_that_timed_out_does_not_eat_a_later_push() {
let (mut r, a, mut batch) = timed();
let b = r.engine_mut().accept();
r.engine_mut().feed(a, &wire(&[b"BLPOP", b"q", b"1"]));
pump(&mut r, &mut batch);
r.engine_mut().server_mut().set_clock_ms(START_MS + 1000);
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(a), b"*-1\r\n");
r.engine_mut().feed(b, &wire(&[b"RPUSH", b"q", b"late"]));
r.engine_mut()
.feed(b, &wire(&[b"LRANGE", b"q", b"0", b"-1"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(a), b"*-1\r\n", "nothing more");
assert!(r.engine().sink().sent(b).ends_with(b"*1\r\n$4\r\nlate\r\n"));
}
#[test]
fn a_client_that_goes_away_while_it_waits_takes_its_waiter_with_it() {
let (mut r, a, mut batch) = engine();
let b = r.engine_mut().accept();
r.engine_mut().feed(a, &wire(&[b"BLPOP", b"q", b"0"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().server().parked(), 1);
r.engine_mut().hangup(a);
pump(&mut r, &mut batch);
assert_eq!(r.engine().server().parked(), 0);
assert_eq!(r.engine().clients(), 1);
let again = r.engine_mut().accept();
assert_eq!(again, a);
r.engine_mut().feed(b, &wire(&[b"RPUSH", b"q", b"x"]));
r.engine_mut()
.feed(again, &wire(&[b"LRANGE", b"q", b"0", b"-1"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(again), b"*1\r\n$1\r\nx\r\n");
}
#[test]
fn a_hangup_while_parked_gives_back_the_slot_and_the_decoders() {
let (mut r, a, mut batch) = engine();
let mut stream = wire(&[b"BLPOP", b"q", b"0"]);
stream.extend(wire(&[b"PING"]));
stream.extend(wire(&[b"PING"]));
r.engine_mut().feed(a, &stream);
pump(&mut r, &mut batch);
let decoders = r.engine().decoders();
r.engine_mut().hangup(a);
pump(&mut r, &mut batch);
assert_eq!(r.engine().clients(), 0);
assert!(r.engine().sink().was_closed(a));
assert_eq!(r.engine().decoders(), decoders, "the pool came back whole");
let again = r.engine_mut().accept();
assert_eq!(again, a);
r.engine_mut().feed(again, &wire(&[b"PING"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(again), b"+PONG\r\n");
}
#[test]
fn a_published_message_lands_on_the_subscriber() {
let (mut r, sub, mut batch) = engine();
let pubr = r.engine_mut().accept();
r.engine_mut().feed(sub, &wire(&[b"SUBSCRIBE", b"news"]));
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(sub),
b"*3\r\n$9\r\nsubscribe\r\n$4\r\nnews\r\n:1\r\n"
);
r.engine_mut().sink_mut().clear();
r.engine_mut()
.feed(pubr, &wire(&[b"PUBLISH", b"news", b"hi"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(pubr), b":1\r\n");
assert_eq!(
r.engine().sink().sent(sub),
b"*3\r\n$7\r\nmessage\r\n$4\r\nnews\r\n$2\r\nhi\r\n"
);
}
#[test]
fn a_pattern_subscriber_is_told_the_pattern_and_the_channel() {
let (mut r, sub, mut batch) = engine();
let pubr = r.engine_mut().accept();
r.engine_mut().feed(sub, &wire(&[b"PSUBSCRIBE", b"ne*"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut()
.feed(pubr, &wire(&[b"PUBLISH", b"news", b"hi"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(pubr), b":1\r\n");
assert_eq!(
r.engine().sink().sent(sub),
b"*4\r\n$8\r\npmessage\r\n$3\r\nne*\r\n$4\r\nnews\r\n$2\r\nhi\r\n"
);
}
#[test]
fn resp2_takes_almost_nothing_from_a_subscriber() {
let (mut r, conn, mut batch) = engine();
r.engine_mut().feed(conn, &wire(&[b"SUBSCRIBE", b"a"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut().feed(conn, &wire(&[b"GET", b"k"]));
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(conn),
b"-ERR Can't execute 'get': only (P|S)SUBSCRIBE / (P|S)UNSUBSCRIBE / PING / QUIT / RESET are allowed in this context\r\n"
);
r.engine_mut().sink_mut().clear();
r.engine_mut().feed(conn, &wire(&[b"PING"]));
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(conn),
b"*2\r\n$4\r\npong\r\n$0\r\n\r\n"
);
r.engine_mut().sink_mut().clear();
r.engine_mut().feed(conn, &wire(&[b"UNSUBSCRIBE", b"a"]));
r.engine_mut().feed(conn, &wire(&[b"GET", b"k"]));
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(conn),
b"*3\r\n$11\r\nunsubscribe\r\n$1\r\na\r\n:0\r\n$-1\r\n"
);
}
#[test]
fn a_shard_channel_and_a_pattern_do_not_hear_each_other() {
let (mut r, sub, mut batch) = engine();
let pubr = r.engine_mut().accept();
r.engine_mut().feed(sub, &wire(&[b"SSUBSCRIBE", b"sx"]));
r.engine_mut().feed(sub, &wire(&[b"PSUBSCRIBE", b"s*"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut()
.feed(pubr, &wire(&[b"SPUBLISH", b"sx", b"one"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(pubr), b":1\r\n");
assert_eq!(
r.engine().sink().sent(sub),
b"*3\r\n$8\r\nsmessage\r\n$2\r\nsx\r\n$3\r\none\r\n"
);
r.engine_mut().sink_mut().clear();
r.engine_mut()
.feed(pubr, &wire(&[b"PUBLISH", b"sx", b"two"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(pubr), b":1\r\n");
assert_eq!(
r.engine().sink().sent(sub),
b"*4\r\n$8\r\npmessage\r\n$2\r\ns*\r\n$2\r\nsx\r\n$3\r\ntwo\r\n"
);
}
#[test]
fn a_subscriber_that_goes_away_leaves_the_registry() {
let (mut r, sub, mut batch) = engine();
let pubr = r.engine_mut().accept();
r.engine_mut().feed(sub, &wire(&[b"SUBSCRIBE", b"news"]));
pump(&mut r, &mut batch);
r.engine_mut().hangup(sub);
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut()
.feed(pubr, &wire(&[b"PUBLISH", b"news", b"hi"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(pubr), b":0\r\n");
let next = r.engine_mut().accept();
assert_eq!(next, sub);
r.engine_mut()
.feed(pubr, &wire(&[b"PUBLISH", b"news", b"hi"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(next), b"");
}
#[test]
fn resp3_delivers_a_message_as_a_push() {
let (mut r, conn, mut batch) = engine();
r.engine_mut().feed(conn, &wire(&[b"HELLO", b"3"]));
r.engine_mut().feed(conn, &wire(&[b"SUBSCRIBE", b"a"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut().feed(conn, &wire(&[b"GET", b"k"]));
r.engine_mut().feed(conn, &wire(&[b"PUBLISH", b"a", b"w"]));
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(conn),
b"_\r\n:1\r\n>3\r\n$7\r\nmessage\r\n$1\r\na\r\n$1\r\nw\r\n"
);
}
#[test]
fn a_write_reaches_a_keyspace_subscriber() {
let (mut r, sub, mut batch) = engine();
let writer = r.engine_mut().accept();
r.engine_mut().feed(
writer,
&wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", b"KEA"]),
);
r.engine_mut()
.feed(sub, &wire(&[b"PSUBSCRIBE", b"__key*@0__:*"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(writer), b"+OK\r\n");
assert_eq!(
r.engine().sink().sent(sub),
b"*4\r\n$8\r\npmessage\r\n$12\r\n__key*@0__:*\r\n\
$16\r\n__keyspace@0__:k\r\n$3\r\nset\r\n\
*4\r\n$8\r\npmessage\r\n$12\r\n__key*@0__:*\r\n\
$18\r\n__keyevent@0__:set\r\n$1\r\nk\r\n"
);
}
#[test]
fn a_write_says_nothing_until_the_setting_turns_it_on() {
let (mut r, sub, mut batch) = engine();
let writer = r.engine_mut().accept();
r.engine_mut()
.feed(sub, &wire(&[b"PSUBSCRIBE", b"__key*@0__:*"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent(sub), b"");
}
#[test]
fn only_the_classes_that_were_asked_for_are_published() {
let (mut r, sub, mut batch) = engine();
let writer = r.engine_mut().accept();
r.engine_mut().feed(
writer,
&wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", b"Eg"]),
);
r.engine_mut()
.feed(sub, &wire(&[b"PSUBSCRIBE", b"__key*@0__:*"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
r.engine_mut().feed(writer, &wire(&[b"DEL", b"k"]));
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(sub),
b"*4\r\n$8\r\npmessage\r\n$12\r\n__key*@0__:*\r\n\
$18\r\n__keyevent@0__:del\r\n$1\r\nk\r\n"
);
}
#[test]
fn a_write_with_a_deadline_on_it_says_two_things() {
let (mut r, sub, mut batch) = engine();
let writer = r.engine_mut().accept();
r.engine_mut().feed(
writer,
&wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", b"EA"]),
);
r.engine_mut()
.feed(sub, &wire(&[b"PSUBSCRIBE", b"__keyevent@0__:*"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut()
.feed(writer, &wire(&[b"SETEX", b"k", b"100", b"v"]));
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(sub),
b"*4\r\n$8\r\npmessage\r\n$16\r\n__keyevent@0__:*\r\n\
$18\r\n__keyevent@0__:set\r\n$1\r\nk\r\n\
*4\r\n$8\r\npmessage\r\n$16\r\n__keyevent@0__:*\r\n\
$21\r\n__keyevent@0__:expire\r\n$1\r\nk\r\n"
);
}
fn watching() -> (Reactor<Wire<Recorder>>, ConnId, ConnId, Vec<Cmd>) {
watching_flags(b"EA")
}
fn watching_flags(flags: &[u8]) -> (Reactor<Wire<Recorder>>, ConnId, ConnId, Vec<Cmd>) {
let (mut r, sub, mut batch) = engine();
let writer = r.engine_mut().accept();
r.engine_mut().feed(
writer,
&wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", flags]),
);
r.engine_mut()
.feed(sub, &wire(&[b"PSUBSCRIBE", b"__keyevent@0__:*"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
(r, sub, writer, batch)
}
fn fired(r: &Reactor<Wire<Recorder>>, sub: ConnId) -> Vec<(String, String)> {
let sent = String::from_utf8_lossy(r.engine().sink().sent(sub)).into_owned();
let mut out = Vec::new();
let mut parts = sent.split("\r\n");
while let Some(p) = parts.next() {
let Some(event) = p.strip_prefix("__keyevent@0__:") else {
continue;
};
if event == "*" {
continue;
}
parts.next();
let key = parts.next().unwrap_or_default();
out.push((event.to_owned(), key.to_owned()));
}
out
}
#[test]
fn taking_the_last_of_a_collection_says_the_key_went_with_it() {
let (mut r, sub, writer, mut batch) = watching();
r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"k", b"a"]));
r.engine_mut().feed(writer, &wire(&[b"LPOP", b"k"]));
pump(&mut r, &mut batch);
assert_eq!(
fired(&r, sub),
[("rpush", "k"), ("lpop", "k"), ("del", "k")]
.map(|(e, k)| (e.to_owned(), k.to_owned()))
);
}
#[test]
fn a_move_onto_a_member_already_there_says_only_the_removal() {
let (mut r, sub, writer, mut batch) = watching();
r.engine_mut().feed(writer, &wire(&[b"SADD", b"a", b"m"]));
r.engine_mut()
.feed(writer, &wire(&[b"SADD", b"b", b"m", b"n"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut()
.feed(writer, &wire(&[b"SMOVE", b"a", b"b", b"m"]));
pump(&mut r, &mut batch);
assert_eq!(
fired(&r, sub),
[("srem", "a"), ("del", "a")].map(|(e, k)| (e.to_owned(), k.to_owned()))
);
}
#[test]
fn a_score_that_did_not_move_says_nothing() {
let (mut r, sub, writer, mut batch) = watching();
r.engine_mut()
.feed(writer, &wire(&[b"ZADD", b"z", b"4", b"m"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut()
.feed(writer, &wire(&[b"ZADD", b"z", b"4", b"m"]));
r.engine_mut()
.feed(writer, &wire(&[b"ZINCRBY", b"z", b"0", b"m"]));
pump(&mut r, &mut batch);
assert_eq!(fired(&r, sub), []);
r.engine_mut()
.feed(writer, &wire(&[b"ZINCRBY", b"z", b"1", b"m"]));
pump(&mut r, &mut batch);
assert_eq!(fired(&r, sub), [("zincr".to_owned(), "z".to_owned())]);
}
#[test]
fn a_stream_write_says_what_the_trim_behind_it_took() {
let (mut r, sub, writer, mut batch) = watching();
r.engine_mut()
.feed(writer, &wire(&[b"XADD", b"s", b"1-1", b"f", b"v"]));
r.engine_mut().feed(
writer,
&wire(&[b"XADD", b"s", b"MAXLEN", b"9", b"2-1", b"f", b"v"]),
);
r.engine_mut().feed(
writer,
&wire(&[b"XADD", b"s", b"MAXLEN", b"1", b"3-1", b"f", b"v"]),
);
pump(&mut r, &mut batch);
assert_eq!(
fired(&r, sub),
[("xadd", "s"), ("xadd", "s"), ("xadd", "s"), ("xtrim", "s")]
.map(|(e, k)| (e.to_owned(), k.to_owned()))
);
}
#[test]
fn acknowledging_an_entry_that_is_already_gone_says_nothing() {
let (mut r, sub, writer, mut batch) = watching();
for cmd in [
wire(&[b"XADD", b"s", b"1-1", b"f", b"v"]),
wire(&[b"XGROUP", b"CREATE", b"s", b"g", b"0"]),
wire(&[b"XREADGROUP", b"GROUP", b"g", b"c", b"STREAMS", b"s", b">"]),
wire(&[b"XDEL", b"s", b"1-1"]),
] {
r.engine_mut().feed(writer, &cmd);
}
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut().feed(
writer,
&wire(&[b"XACKDEL", b"s", b"g", b"IDS", b"1", b"1-1"]),
);
pump(&mut r, &mut batch);
assert_eq!(fired(&r, sub), []);
}
#[test]
fn emptying_a_hash_says_the_key_went_with_the_last_field() {
let (mut r, sub, writer, mut batch) = watching();
r.engine_mut()
.feed(writer, &wire(&[b"HSET", b"h", b"a", b"1", b"b", b"2"]));
r.engine_mut().feed(writer, &wire(&[b"HDEL", b"h", b"a"]));
r.engine_mut()
.feed(writer, &wire(&[b"HDEL", b"h", b"b", b"a"]));
pump(&mut r, &mut batch);
assert_eq!(
fired(&r, sub),
[("hset", "h"), ("hdel", "h"), ("hdel", "h"), ("del", "h")]
.map(|(e, k)| (e.to_owned(), k.to_owned()))
);
}
#[test]
fn a_field_deadline_already_past_reads_as_a_removal() {
let (mut r, sub, writer, mut batch) = watching();
r.engine_mut()
.feed(writer, &wire(&[b"HSET", b"h", b"a", b"1", b"b", b"2"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut().feed(
writer,
&wire(&[b"HEXPIRE", b"h", b"0", b"FIELDS", b"1", b"a"]),
);
r.engine_mut().feed(
writer,
&wire(&[b"HEXPIRE", b"h", b"100", b"FIELDS", b"1", b"b"]),
);
pump(&mut r, &mut batch);
assert_eq!(
fired(&r, sub),
[("hdel", "h"), ("hexpire", "h")].map(|(e, k)| (e.to_owned(), k.to_owned()))
);
}
#[test]
fn a_write_under_a_deadline_already_gone_says_the_write_first() {
let (mut r, sub, writer, mut batch) = watching();
r.engine_mut().feed(
writer,
&wire(&[b"HSETEX", b"h", b"EXAT", b"1", b"FIELDS", b"1", b"a", b"1"]),
);
pump(&mut r, &mut batch);
assert_eq!(
fired(&r, sub),
[("hset", "h"), ("hdel", "h"), ("del", "h")].map(|(e, k)| (e.to_owned(), k.to_owned()))
);
}
#[test]
fn clearing_a_deadline_that_was_never_set_says_nothing() {
let (mut r, sub, writer, mut batch) = watching();
r.engine_mut()
.feed(writer, &wire(&[b"HSET", b"h", b"a", b"1", b"b", b"2"]));
r.engine_mut().feed(
writer,
&wire(&[b"HEXPIRE", b"h", b"100", b"FIELDS", b"1", b"a"]),
);
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut()
.feed(writer, &wire(&[b"HPERSIST", b"h", b"FIELDS", b"1", b"b"]));
r.engine_mut().feed(
writer,
&wire(&[b"HGETEX", b"h", b"PERSIST", b"FIELDS", b"1", b"b"]),
);
pump(&mut r, &mut batch);
assert_eq!(fired(&r, sub), []);
r.engine_mut().feed(
writer,
&wire(&[b"HGETEX", b"h", b"PERSIST", b"FIELDS", b"2", b"a", b"b"]),
);
pump(&mut r, &mut batch);
assert_eq!(fired(&r, sub), [("hpersist".to_owned(), "h".to_owned())]);
}
#[test]
fn a_key_that_was_not_there_before_says_so() {
let (mut r, sub, writer, mut batch) = watching_flags(b"En");
r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"w"]));
r.engine_mut().feed(writer, &wire(&[b"APPEND", b"k", b"x"]));
r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"l", b"a"]));
r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"l", b"b"]));
pump(&mut r, &mut batch);
assert_eq!(
fired(&r, sub),
[("new", "k"), ("new", "l")].map(|(e, k)| (e.to_owned(), k.to_owned()))
);
}
#[test]
fn the_news_of_a_new_key_comes_before_the_write_that_made_it() {
let (mut r, sub, writer, mut batch) = watching_flags(b"EAn");
r.engine_mut().feed(writer, &wire(&[b"SET", b"a", b"1"]));
r.engine_mut().feed(writer, &wire(&[b"SET", b"b", b"2"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut().feed(writer, &wire(&[b"RENAME", b"a", b"b"]));
pump(&mut r, &mut batch);
assert_eq!(
fired(&r, sub),
[("new", "b"), ("rename_from", "a"), ("rename_to", "b")]
.map(|(e, k)| (e.to_owned(), k.to_owned()))
);
}
#[test]
fn writing_over_a_destination_is_not_a_key_arriving() {
let (mut r, sub, writer, mut batch) = watching_flags(b"EAn");
r.engine_mut()
.feed(writer, &wire(&[b"RPUSH", b"l", b"c", b"a", b"b"]));
r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"d", b"x"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut()
.feed(writer, &wire(&[b"SORT", b"l", b"ALPHA", b"STORE", b"d"]));
pump(&mut r, &mut batch);
assert_eq!(fired(&r, sub), [("sortstore".to_owned(), "d".to_owned())]);
r.engine_mut().sink_mut().clear();
r.engine_mut().feed(writer, &wire(&[b"DEL", b"d"]));
r.engine_mut()
.feed(writer, &wire(&[b"SORT", b"l", b"ALPHA", b"STORE", b"d"]));
pump(&mut r, &mut batch);
assert_eq!(
fired(&r, sub),
[("del", "d"), ("new", "d"), ("sortstore", "d")]
.map(|(e, k)| (e.to_owned(), k.to_owned()))
);
}
fn watching_fields(flags: &[u8]) -> (Reactor<Wire<Recorder>>, ConnId, ConnId, Vec<Cmd>) {
let (mut r, sub, mut batch) = engine();
let writer = r.engine_mut().accept();
r.engine_mut().feed(
writer,
&wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", flags]),
);
r.engine_mut()
.feed(sub, &wire(&[b"PSUBSCRIBE", b"__subkey*@0__:*"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
(r, sub, writer, batch)
}
fn carried(r: &Reactor<Wire<Recorder>>, sub: ConnId) -> Vec<(String, String)> {
let sent = String::from_utf8_lossy(r.engine().sink().sent(sub)).into_owned();
let mut out = Vec::new();
let mut parts = sent.split("\r\n");
while let Some(p) = parts.next() {
if p != "pmessage" {
continue;
}
let mut next = || {
parts.next();
parts.next().unwrap_or_default().to_owned()
};
next();
let channel = next();
out.push((channel, next()));
}
out
}
#[test]
fn the_subkey_channels_carry_the_fields_an_event_touched() {
let (mut r, sub, writer, mut batch) = watching_fields(b"ASTIV");
r.engine_mut()
.feed(writer, &wire(&[b"HSET", b"h", b"a,b", b"1", b"c", b"2"]));
pump(&mut r, &mut batch);
assert_eq!(
carried(&r, sub),
[
("__subkeyspace@0__:h", "hset|3:a,b,1:c"),
("__subkeyevent@0__:hset", "1:h|3:a,b,1:c"),
("__subkeyspaceitem@0__:h\na,b", "hset"),
("__subkeyspaceitem@0__:h\nc", "hset"),
("__subkeyspaceevent@0__:hset|h", "3:a,b,1:c"),
]
.map(|(c, p)| (c.to_owned(), p.to_owned()))
);
}
#[test]
fn a_key_holding_a_newline_skips_the_per_field_channel() {
let (mut r, sub, writer, mut batch) = watching_fields(b"ASTIV");
r.engine_mut()
.feed(writer, &wire(&[b"HSET", b"h\nx", b"f", b"1"]));
pump(&mut r, &mut batch);
assert_eq!(
carried(&r, sub),
[
("__subkeyspace@0__:h\nx", "hset|1:f"),
("__subkeyevent@0__:hset", "3:h\nx|1:f"),
("__subkeyspaceevent@0__:hset|h\nx", "1:f"),
]
.map(|(c, p)| (c.to_owned(), p.to_owned()))
);
}
#[test]
fn an_event_with_no_fields_stays_off_the_subkey_channels() {
let (mut r, sub, writer, mut batch) = watching_fields(b"AS");
r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"l", b"a"]));
r.engine_mut()
.feed(writer, &wire(&[b"HSET", b"h", b"f", b"1"]));
r.engine_mut().feed(writer, &wire(&[b"HDEL", b"h", b"f"]));
pump(&mut r, &mut batch);
assert_eq!(
carried(&r, sub),
[
("__subkeyspace@0__:h", "hset|1:f"),
("__subkeyspace@0__:h", "hdel|1:f"),
]
.map(|(c, p)| (c.to_owned(), p.to_owned()))
);
}
fn watching_clock() -> (Reactor<Wire<Recorder>>, ConnId, ConnId, Vec<Cmd>) {
let (mut r, sub, mut batch) = timed();
let writer = r.engine_mut().accept();
r.engine_mut().feed(
writer,
&wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", b"EA"]),
);
r.engine_mut()
.feed(sub, &wire(&[b"PSUBSCRIBE", b"__keyevent@0__:*"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
(r, sub, writer, batch)
}
#[test]
fn a_deadline_that_passed_is_news_when_a_reader_finds_it() {
let (mut r, sub, writer, mut batch) = watching_clock();
r.engine_mut()
.feed(writer, &wire(&[b"SET", b"k", b"v", b"PX", b"10"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine().server().advance_clock_ms(50);
r.engine_mut().feed(writer, &wire(&[b"GET", b"k"]));
pump(&mut r, &mut batch);
assert_eq!(
fired(&r, sub),
[("expired", "k")].map(|(e, k)| (e.to_owned(), k.to_owned())),
"and not a del alongside it, which is a different piece of news"
);
}
#[test]
fn a_deadline_that_passed_is_news_with_nobody_reading() {
let (mut r, sub, writer, mut batch) = watching_clock();
r.engine_mut()
.feed(writer, &wire(&[b"SET", b"k", b"v", b"PX", b"10"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine().server().advance_clock_ms(50);
pump(&mut r, &mut batch);
assert_eq!(
fired(&r, sub),
[("expired", "k")].map(|(e, k)| (e.to_owned(), k.to_owned()))
);
r.engine_mut().sink_mut().clear();
pump(&mut r, &mut batch);
assert!(fired(&r, sub).is_empty(), "and it only goes once");
}
#[test]
fn a_key_a_limit_took_says_it_was_evicted() {
let (mut r, sub, writer, mut batch) = watching();
let val = vec![b'v'; 256];
for i in 0..2000u32 {
let k = format!("key:{i:08}");
r.engine_mut()
.feed(writer, &wire(&[b"SET", k.as_bytes(), &val]));
}
pump(&mut r, &mut batch);
r.engine().server().refresh_memory();
let full = r.engine().server().memory_bytes();
r.engine_mut().sink_mut().clear();
let limit = (full / 2).to_string();
r.engine_mut().feed(
writer,
&wire(&[b"CONFIG", b"SET", b"maxmemory-policy", b"allkeys-random"]),
);
r.engine_mut().feed(
writer,
&wire(&[b"CONFIG", b"SET", b"maxmemory", limit.as_bytes()]),
);
r.engine_mut()
.feed(writer, &wire(&[b"SET", b"newcomer", &val]));
pump(&mut r, &mut batch);
let events = fired(&r, sub);
assert!(
events.iter().any(|(e, _)| e == "evicted"),
"the write made room and never said so: {events:?}"
);
assert!(
events
.iter()
.all(|(e, k)| e != "evicted" || k != "newcomer"),
"the key the write was for is the one key it cannot have taken"
);
}
#[test]
fn a_transaction_publishes_between_its_commands_and_not_after_them() {
let (mut r, sub, mut batch) = engine();
let writer = r.engine_mut().accept();
r.engine_mut().feed(
writer,
&wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", b"EA"]),
);
r.engine_mut()
.feed(sub, &wire(&[b"PSUBSCRIBE", b"__keyevent@0__:*"]));
pump(&mut r, &mut batch);
r.engine_mut().sink_mut().clear();
r.engine_mut().feed(writer, &wire(&[b"MULTI"]));
r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
r.engine_mut().feed(writer, &wire(&[b"DEL", b"k"]));
r.engine_mut().feed(writer, &wire(&[b"EXEC"]));
pump(&mut r, &mut batch);
assert_eq!(
r.engine().sink().sent(sub),
b"*4\r\n$8\r\npmessage\r\n$16\r\n__keyevent@0__:*\r\n\
$18\r\n__keyevent@0__:set\r\n$1\r\nk\r\n\
*4\r\n$8\r\npmessage\r\n$16\r\n__keyevent@0__:*\r\n\
$18\r\n__keyevent@0__:del\r\n$1\r\nk\r\n"
);
}
#[test]
fn a_reply_the_socket_would_not_take_is_offered_again() {
let mut r = Reactor::inline(Wire::new(Trickle::default()));
let conn = r.engine_mut().accept();
let mut batch = Vec::new();
r.engine_mut().feed(conn, &wire(&[b"PING"]));
pump(&mut r, &mut batch);
assert_eq!(r.engine().sink().sent, b"+PONG\r\n");
assert_eq!(r.engine().sink().writes, 2);
}
}