use yo_reactor::{BATCH_MAX, Engine, Reactor};
use crate::dispatch::table;
use crate::dispatch::{self, Flow, Server};
use crate::front::{Front, Wrote};
use crate::proto::Limits;
use yo_kv::Keyspace;
pub use crate::front::Cmd;
pub type ConnId = u32;
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: Server,
}
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 {
front: Front::new(sink),
server,
}
}
#[must_use]
pub const fn server(&self) -> &Server {
&self.server
}
pub const fn server_mut(&mut self) -> &mut Server {
&mut self.server
}
#[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.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 mut at = 0;
while at < self.server.parked() {
let p = self.server.waiters().at(at);
if !self.front.answers(p.conn, p.client) {
self.server.drop_waiter(at);
continue;
}
let served = {
let Wire { server, front } = self;
server.serve_waiter(at, now, front.out(p.conn))
};
if served {
self.server.drop_waiter(at);
self.front.unpark(p.conn);
self.front.soil(p.conn);
} else {
at += 1;
}
}
}
#[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 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) {
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.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() != 0 {
self.serve_waiters();
}
yo_reactor::Flow::Next
}
fn flush(&mut self) {
if self.server.parked() != 0 {
self.server.refresh_clock();
self.serve_waiters();
}
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().flush();
reactor.engine_mut().maintain();
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 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();
for _ in 0..1000 {
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 a thousand 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_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);
}
}