use std::collections::VecDeque;
use yo_reactor::{BATCH_MAX, Engine, Reactor};
use crate::dispatch::table::{self, lookup_index};
use crate::dispatch::{self, Args, Flow, Server, Session};
use crate::error::ProtocolError;
use crate::proto::{Limits, Proto};
use crate::reply::Out;
use crate::request::{Argv, Step};
use yo_kv::Keyspace;
pub type ConnId = u32;
const READ_BUF: usize = 16 * 1024;
const OUT_BUF: usize = 16 * 1024;
const ARGV_HINT: usize = 8;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Cmd {
conn: ConnId,
slot: u32,
base: usize,
spec: u16,
}
impl Cmd {
#[must_use]
pub const fn conn(&self) -> ConnId {
self.conn
}
}
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));
}
}
struct Conn {
live: bool,
session: Session,
out: Out,
buf: Vec<u8>,
head: usize,
partial: Option<u32>,
pending: u32,
closing: bool,
deferred: Option<ProtocolError>,
skip: bool,
gone: bool,
dirty: bool,
blocked: bool,
parked: Vec<Cmd>,
held: usize,
}
impl Conn {
fn new(id: u64) -> Conn {
yo_alloc::allow(|| Conn {
live: true,
session: Session::new(id),
out: Out::with_capacity(Proto::Resp2, OUT_BUF),
buf: Vec::with_capacity(READ_BUF),
head: 0,
partial: None,
pending: 0,
closing: false,
deferred: None,
skip: false,
gone: false,
dirty: false,
blocked: false,
parked: Vec::new(),
held: 0,
})
}
fn size(&self) -> usize {
self.buf.capacity() + self.out.capacity()
}
fn reset(&mut self, id: u64) {
self.live = true;
self.session = Session::new(id);
self.out.clear();
self.buf.clear();
self.head = 0;
self.partial = None;
self.pending = 0;
self.closing = false;
self.deferred = None;
self.skip = false;
self.gone = false;
self.dirty = false;
self.blocked = false;
self.parked.clear();
}
fn compact(&mut self) {
if self.pending > 0 || self.head == 0 {
return;
}
if self.head == self.buf.len() {
self.buf.clear();
} else {
self.buf.drain(..self.head);
}
self.head = 0;
}
}
pub struct Wire<S> {
server: Server,
sink: S,
conns: Vec<Conn>,
free: Vec<ConnId>,
argvs: Vec<Argv>,
spare: Vec<u32>,
ready: VecDeque<Cmd>,
dirty: Vec<ConnId>,
scratch: Vec<u8>,
limits: Limits,
next_id: u64,
}
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 {
server,
sink,
conns: Vec::new(),
free: Vec::new(),
argvs: Vec::new(),
spare: Vec::new(),
ready: VecDeque::with_capacity(BATCH_MAX),
dirty: Vec::with_capacity(16),
scratch: Vec::with_capacity(128),
limits: Limits::default(),
next_id: 1,
}
}
#[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.sink
}
pub const fn sink_mut(&mut self) -> &mut S {
&mut self.sink
}
pub fn set_limits(&mut self, limits: Limits) {
self.limits = limits;
}
pub fn accept(&mut self) -> ConnId {
let id = self.next_id;
self.next_id += 1;
self.server.stats.clients += 1;
self.server.stats.connections += 1;
let at = match self.free.pop() {
Some(at) => {
self.conns[at as usize].reset(id);
at
}
None => {
let conn = Conn::new(id);
yo_alloc::allow(|| self.conns.push(conn));
(self.conns.len() - 1) as ConnId
}
};
self.note_size(at);
at
}
pub fn hangup(&mut self, conn: ConnId) {
let c = &mut self.conns[conn as usize];
if !c.live {
return;
}
c.gone = true;
c.closing = true;
if c.blocked {
self.unpark(conn);
}
if self.conns[conn as usize].pending == 0 {
self.release(conn);
}
}
fn unpark(&mut self, conn: ConnId) {
let mut parked = {
let c = &mut self.conns[conn as usize];
c.blocked = false;
core::mem::take(&mut c.parked)
};
while let Some(cmd) = parked.pop() {
if self.ready.len() == self.ready.capacity() {
yo_alloc::allow(|| self.ready.reserve(BATCH_MAX));
}
self.ready.push_front(cmd);
}
self.conns[conn as usize].parked = parked;
if !self.conns[conn as usize].closing {
self.frame(conn);
}
}
fn serve_waiters(&mut self) {
let now = self.server.now_ms();
let mut at = 0;
while at < self.server.waiters().len() {
let p = self.server.waiters().at(at);
{
let c = &self.conns[p.conn as usize];
if !c.live || c.session.id() != p.client {
self.server.waiters_mut().drop_at(at);
continue;
}
}
let served = {
let Wire { server, conns, .. } = self;
server.serve_waiter(at, now, &mut conns[p.conn as usize].out)
};
if served {
self.server.waiters_mut().drop_at(at);
self.unpark(p.conn);
self.soil(p.conn);
} else {
at += 1;
}
}
}
#[must_use]
pub fn clients(&self) -> usize {
self.conns.iter().filter(|c| c.live).count()
}
#[must_use]
pub fn ready(&self) -> usize {
self.ready.len()
}
#[must_use]
pub fn owed(&self) -> usize {
self.dirty.len()
}
#[must_use]
pub fn decoders(&self) -> usize {
self.argvs.len()
}
#[must_use]
pub fn buffer_bytes(&self) -> usize {
self.conns.iter().map(Conn::size).sum()
}
pub fn feed(&mut self, conn: ConnId, bytes: &[u8]) {
{
let c = &mut self.conns[conn as usize];
if !c.live || c.closing {
return;
}
yo_alloc::allow(|| c.buf.extend_from_slice(bytes));
}
self.frame(conn);
self.note_size(conn);
}
fn note_size(&mut self, conn: ConnId) {
let c = &mut self.conns[conn as usize];
let now = c.size();
if now == c.held {
return;
}
let delta = now as isize - c.held as isize;
c.held = now;
self.server.note_conn_bytes(delta);
}
fn frame(&mut self, conn: ConnId) {
if self.conns[conn as usize].blocked {
return;
}
loop {
let base = self.conns[conn as usize].head;
let slot = match self.conns[conn as usize].partial.take() {
Some(slot) => slot,
None => self.take_decoder(),
};
let step = {
let c = &self.conns[conn as usize];
self.argvs[slot as usize].decode(&c.buf[base..], &self.limits)
};
match step {
Ok(Step::Command { consumed }) => {
self.conns[conn as usize].head += consumed;
if self.argvs[slot as usize].is_empty() {
self.spare.push(slot);
} else {
if self.ready.len() == self.ready.capacity() {
yo_alloc::allow(|| self.ready.reserve(BATCH_MAX));
}
let spec = {
let c = &self.conns[conn as usize];
let args = Args::new(&self.argvs[slot as usize], &c.buf[base..]);
lookup_index(args.name())
};
self.ready.push_back(Cmd {
conn,
slot,
base,
spec,
});
self.conns[conn as usize].pending += 1;
}
}
Ok(Step::Incomplete) => {
self.conns[conn as usize].partial = Some(slot);
break;
}
Err(e) => {
self.spare.push(slot);
let c = &mut self.conns[conn as usize];
c.deferred = Some(e);
c.closing = true;
self.soil(conn);
break;
}
}
}
self.conns[conn as usize].compact();
}
fn take_decoder(&mut self) -> u32 {
match self.spare.pop() {
Some(slot) => {
self.argvs[slot as usize].reset();
slot
}
None => yo_alloc::allow(|| {
self.argvs.push(Argv::with_capacity(ARGV_HINT));
self.spare.reserve(self.argvs.len());
(self.argvs.len() - 1) as u32
}),
}
}
fn soil(&mut self, conn: ConnId) {
let c = &mut self.conns[conn as usize];
if !c.dirty {
c.dirty = true;
if self.dirty.len() == self.dirty.capacity() {
yo_alloc::allow(|| self.dirty.reserve(16));
}
self.dirty.push(conn);
}
}
fn release(&mut self, conn: ConnId) {
{
let c = &mut self.conns[conn as usize];
if !c.live {
return;
}
if let Some(slot) = c.partial.take() {
self.spare.push(slot);
}
c.live = false;
c.dirty = false;
c.blocked = false;
c.out.clear();
c.buf.clear();
c.head = 0;
}
let client = self.conns[conn as usize].session.id();
self.server.waiters_mut().forget(client);
self.server.stats.clients = self.server.stats.clients.saturating_sub(1);
self.sink.closed(conn);
yo_alloc::allow(|| self.free.push(conn));
}
pub fn take_ready(&mut self, into: &mut Vec<Cmd>, max: usize) -> usize {
let n = max.min(self.ready.len());
into.extend(self.ready.drain(..n));
n
}
fn write_out(&mut self, conn: ConnId) -> bool {
{
let c = &self.conns[conn as usize];
if !c.live {
return false;
}
}
if self.conns[conn as usize].pending == 0
&& let Some(e) = self.conns[conn as usize].deferred.take()
{
self.scratch.clear();
e.write_reply(&mut self.scratch);
self.conns[conn as usize].out.raw(&self.scratch);
}
let taken = {
let c = &self.conns[conn as usize];
if c.out.is_empty() {
0
} else {
self.sink.write(conn, c.out.as_slice())
}
};
let c = &mut self.conns[conn as usize];
if taken >= c.out.len() {
c.out.clear();
} else {
c.out.consume(taken);
}
if !c.out.is_empty() {
return true;
}
c.dirty = false;
if c.closing && c.pending == 0 {
self.release(conn);
} else {
c.compact();
}
self.note_size(conn);
false
}
pub fn tick(&mut self) {
self.server.refresh_clock();
}
pub fn maintain(&mut self) -> Option<usize> {
self.server.refresh_memory();
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 c = &self.conns[cmd.conn as usize];
let args = Args::new(&self.argvs[cmd.slot as usize], &c.buf[cmd.base..]);
let key = args.opt(spec.first_key as usize)?;
Some(Keyspace::hash_of(key))
}
fn prefetch(&self, cmd: &Cmd, hash: u64) {
let db = self.conns[cmd.conn as usize].session.db();
self.server.db_ref(db).prefetch(hash);
}
fn run(&mut self, cmd: Cmd, _hash: Option<u64>) -> yo_reactor::Flow {
if self.conns[cmd.conn as usize].blocked {
yo_alloc::allow(|| self.conns[cmd.conn as usize].parked.push(cmd));
return yo_reactor::Flow::Next;
}
let flow = {
let c = &mut self.conns[cmd.conn as usize];
c.pending -= 1;
if c.gone || c.skip {
Flow::Continue
} else {
let args = Args::new(&self.argvs[cmd.slot as usize], &c.buf[cmd.base..]);
let spec = table::at(cmd.spec);
dispatch::resolved(&mut self.server, &mut c.session, spec, args, &mut c.out)
}
};
self.spare.push(cmd.slot);
let c = &self.conns[cmd.conn as usize];
if c.gone {
if c.pending == 0 {
self.release(cmd.conn);
}
} else {
match flow {
Flow::Close => {
let c = &mut self.conns[cmd.conn as usize];
c.closing = true;
c.skip = true;
self.soil(cmd.conn);
}
Flow::Block => {
self.conns[cmd.conn as usize].blocked = true;
let client = self.conns[cmd.conn as usize].session.id();
self.server.waiters_mut().bind(client, cmd.conn);
}
Flow::Continue => self.soil(cmd.conn),
}
}
if !self.server.waiters().is_empty() {
self.serve_waiters();
}
yo_reactor::Flow::Next
}
fn flush(&mut self) {
if !self.server.waiters().is_empty() {
self.server.refresh_clock();
self.serve_waiters();
}
let mut dirty = core::mem::take(&mut self.dirty);
let mut at = 0;
while at < dirty.len() {
let conn = dirty[at];
let owed = self.write_out(conn);
if owed {
at += 1;
} else {
dirty.swap_remove(at);
}
}
self.dirty = dirty;
}
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 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().waiters().len(), 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().waiters().len(), 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().waiters().len(), 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().waiters().len(), 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().waiters().len(), 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().waiters().len(), 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().waiters().len(), 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().waiters().len(), 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().waiters().len(), 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().waiters().len(), 1);
r.engine_mut().hangup(a);
pump(&mut r, &mut batch);
assert_eq!(r.engine().server().waiters().len(), 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);
}
}