1use std::sync::Arc;
75
76use yo_reactor::{BATCH_MAX, Engine, Reactor};
77
78use crate::dispatch::table;
79use crate::dispatch::{self, Flow, Parked, Reply, Server};
80use crate::front::{Front, Wrote};
81use crate::proto::Limits;
82use yo_kv::Keyspace;
83
84pub use crate::front::Cmd;
85
86pub type ConnId = u32;
88
89const SWEEP_LOOKS: usize = yo_reactor::MAINTENANCE_UNITS as usize;
96
97pub trait Sink {
103 fn write(&mut self, conn: ConnId, bytes: &[u8]) -> usize;
108
109 fn closed(&mut self, conn: ConnId) {
111 let _ = conn;
112 }
113}
114
115#[derive(Debug, Default)]
117pub struct Recorder {
118 sent: Vec<Vec<u8>>,
119 closed: Vec<ConnId>,
120}
121
122impl Recorder {
123 #[must_use]
125 pub fn new() -> Recorder {
126 Recorder::default()
127 }
128
129 #[must_use]
131 pub fn sent(&self, conn: ConnId) -> &[u8] {
132 self.sent.get(conn as usize).map_or(&[], Vec::as_slice)
133 }
134
135 #[must_use]
137 pub fn was_closed(&self, conn: ConnId) -> bool {
138 self.closed.contains(&conn)
139 }
140
141 pub fn clear(&mut self) {
143 for c in &mut self.sent {
144 c.clear();
145 }
146 self.closed.clear();
147 }
148}
149
150impl Sink for Recorder {
151 fn write(&mut self, conn: ConnId, bytes: &[u8]) -> usize {
152 yo_alloc::allow(|| {
154 if self.sent.len() <= conn as usize {
155 self.sent.resize_with(conn as usize + 1, Vec::new);
156 }
157 self.sent[conn as usize].extend_from_slice(bytes);
158 });
159 bytes.len()
160 }
161
162 fn closed(&mut self, conn: ConnId) {
163 yo_alloc::allow(|| self.closed.push(conn));
164 }
165}
166
167pub struct Wire<S> {
177 front: Front<S>,
178 server: Arc<Server>,
179 parked: Vec<Parked>,
186 post: Vec<dispatch::Envelope>,
193 held: Vec<(ConnId, u64)>,
201}
202
203impl<S: Sink> Wire<S> {
204 #[must_use]
206 pub fn new(sink: S) -> Wire<S> {
207 Wire::with_server(Server::new(), sink)
208 }
209
210 #[must_use]
213 pub fn with_server(server: Server, sink: S) -> Wire<S> {
214 Wire::over(Arc::new(server), sink)
215 }
216
217 #[must_use]
226 pub fn over(server: Arc<Server>, sink: S) -> Wire<S> {
227 server.is_behind();
231 Wire {
232 front: Front::new(sink),
233 parked: Vec::new(),
234 post: Vec::new(),
235 held: Vec::new(),
236 server,
237 }
238 }
239
240 #[must_use]
242 pub fn server(&self) -> &Server {
243 &self.server
244 }
245
246 #[must_use]
249 pub fn shared(&self) -> Arc<Server> {
250 Arc::clone(&self.server)
251 }
252
253 pub fn server_mut(&mut self) -> &mut Server {
268 self.server.forget_behind();
272 Arc::get_mut(&mut self.server)
273 .expect("the server is set up before the threads that share it are started")
274 }
275
276 #[must_use]
278 pub const fn sink(&self) -> &S {
279 self.front.sink()
280 }
281
282 pub const fn sink_mut(&mut self) -> &mut S {
284 self.front.sink_mut()
285 }
286
287 pub fn set_limits(&mut self, limits: Limits) {
289 self.front.set_limits(limits);
290 }
291
292 pub fn accept(&mut self) -> ConnId {
294 self.server.counted().opened();
295 let at = self.front.open(self.server.next_client());
296 let now = self.server.now_ms();
297 let guarded = self.server.guarded();
302 let row = if let Some(session) = self.front.session_mut(at) {
303 session.opened(now);
304 session.admit(!guarded);
305 Some(session.row().clone())
306 } else {
307 None
308 };
309 if let Some(row) = row {
314 self.server.register_client(&row);
315 }
316 self.note_buffers();
317 at
318 }
319
320 pub fn accept_from(&mut self, peer: &str, local: &str, fd: i32, unix: bool) -> ConnId {
328 let at = self.accept();
329 if let Some(session) = self.front.session_mut(at) {
330 session.set_socket(peer, local, fd, unix);
331 }
332 at
333 }
334
335 fn note_buffers(&mut self) {
340 let delta = self.front.buffer_delta();
341 if delta != 0 {
342 self.server.note_conn_bytes(delta);
343 }
344 }
345
346 pub fn hangup(&mut self, conn: ConnId) {
352 if !self.front.live(conn) {
353 return;
354 }
355 self.front.mark_gone(conn);
356 if self.front.blocked(conn) {
361 self.front.unpark(conn);
362 }
363 if self.front.pending(conn) == 0 {
364 self.release(conn);
365 }
366 self.note_buffers();
367 }
368
369 fn serve_waiters(&mut self) {
382 let now = self.server.now_ms();
383 let mine = self.server.my_slot();
384 self.server.waiters().mine(mine, &mut self.parked);
385 for at in 0..self.parked.len() {
386 let p = self.parked[at];
387 if !self.front.answers(p.conn, p.client) {
392 self.server.forget_waiters(p.client);
393 continue;
394 }
395 let served = {
399 let Wire { server, front, .. } = self;
400 server.serve_waiter(p.client, now, front.out(p.conn))
401 };
402 if served {
403 self.server.forget_waiters(p.client);
404 self.front.unpark(p.conn);
405 self.front.soil(p.conn);
406 }
407 }
408 self.parked.clear();
409 }
410
411 fn deliver(&mut self) {
420 let mut post = core::mem::take(&mut self.post);
423 self.server.take_mail(&mut post);
424 for env in post.drain(..) {
425 let conn = env.conn();
426 if !self.front.answers(conn, env.client()) {
427 continue;
428 }
429 env.write(self.front.out(conn));
430 self.front.soil(conn);
431 }
432 self.post = post;
433 }
434
435 #[must_use]
437 pub fn clients(&self) -> usize {
438 self.front.clients()
439 }
440
441 #[must_use]
443 pub fn ready(&self) -> usize {
444 self.front.ready()
445 }
446
447 #[must_use]
453 pub fn owed(&self) -> usize {
454 self.front.owed()
455 }
456
457 #[must_use]
468 pub fn waiting(&self) -> usize {
469 self.server.parked_here() + self.held.len()
470 }
471
472 #[must_use]
481 pub fn posted(&self) -> usize {
482 self.server.posted()
483 }
484
485 #[must_use]
492 pub fn stopping(&self) -> bool {
493 self.server.stopping()
494 }
495
496 #[must_use]
498 pub fn decoders(&self) -> usize {
499 self.front.decoders()
500 }
501
502 #[must_use]
504 pub fn buffer_bytes(&self) -> usize {
505 self.front.buffer_bytes()
506 }
507
508 pub fn feed(&mut self, conn: ConnId, bytes: &[u8]) {
514 self.front.feed(conn, bytes);
515 self.note_buffers();
516 }
517
518 fn release(&mut self, conn: ConnId) {
520 if let Some(session) = self.front.session_mut(conn) {
524 dispatch::forget_session(&self.server, session);
525 }
526 let Some(client) = self.front.close(conn) else {
527 return;
528 };
529 self.forget(client);
530 }
531
532 fn forget(&mut self, client: u64) {
539 self.server.forget_waiters(client);
540 self.server.forget_client(client);
541 self.server.counted().closed();
542 }
543
544 fn resume(&mut self) {
558 let mut held = core::mem::take(&mut self.held);
561 for &(conn, client) in &held {
562 if self.front.answers(conn, client) && self.front.blocked(conn) {
567 self.front.unpark(conn);
568 }
569 }
570 held.clear();
571 self.held = held;
572 }
573
574 fn reap(&mut self) {
575 for (conn, client) in self.server.my_kills() {
576 self.server.kill_done();
577 if self.front.answers(conn, client) {
578 self.hangup(conn);
579 }
580 }
581 }
582
583 pub fn take_ready(&mut self, into: &mut Vec<Cmd>, max: usize) -> usize {
588 self.front.take_ready(into, max)
589 }
590
591 pub fn tick(&mut self) {
597 self.server.refresh_clock();
598 }
599
600 pub fn maintain(&mut self) -> Option<usize> {
618 if !self.server.cron_running() {
624 return None;
625 }
626 self.server.refresh_memory_slice();
634 self.server.backup_expire();
638 self.server.expire_slice(SWEEP_LOOKS);
643 self.server.compact_step()
644 }
645}
646
647impl<S: Sink> Engine for Wire<S> {
648 type Work = Cmd;
649
650 fn key_hash(&self, cmd: &Cmd) -> Option<u64> {
651 let spec = table::at(cmd.spec)?;
656 if spec.first_key <= 0 {
657 return None;
658 }
659 let args = self.front.args(cmd);
660 let key = args.opt(spec.first_key as usize)?;
665 Some(Keyspace::hash_of(key))
666 }
667
668 fn prefetch(&self, cmd: &Cmd, hash: u64) {
669 let db = self.front.db(cmd.conn());
670 self.server.striped_ref(db).prefetch_hashed(hash);
675 }
676
677 fn run(&mut self, cmd: Cmd, _hash: Option<u64>) -> yo_reactor::Flow {
678 let conn = cmd.conn();
679 if self.front.blocked(conn) {
684 self.front.park(conn, cmd);
685 return yo_reactor::Flow::Next;
686 }
687
688 let flow = if self.front.start(&cmd) {
692 let Wire { front, server, .. } = self;
693 let (args, session, out) = front.parts(&cmd);
694 let spec = table::at(cmd.spec);
695 let mark = out.len();
696 let flow = dispatch::resolved(server, session, spec, args, out);
697 if flow != Flow::Hold {
702 session.finished();
709 session.note_proto(out.proto().version());
710 let mode = session.reply_mode();
711 session.step_reply();
712 if mode != Reply::On {
713 out.truncate(mark);
714 }
715 }
716 flow
717 } else {
718 Flow::Continue
722 };
723
724 if flow == Flow::Hold {
731 self.front.hold(conn, cmd);
732 let client = self.front.client(conn);
733 yo_alloc::allow(|| self.held.push((conn, client)));
734 return yo_reactor::Flow::Next;
735 }
736
737 self.front.done(&cmd);
738 if self.front.gone(conn) {
739 if self.front.pending(conn) == 0 {
740 self.release(conn);
741 }
742 } else {
743 match flow {
744 Flow::Close => {
745 self.front.quit(conn);
746 self.front.soil(conn);
747 }
748 Flow::Block => {
753 self.front.block(conn);
754 let client = self.front.client(conn);
755 self.server.bind_waiter(client, conn);
756 }
757 Flow::Hold => {}
760 Flow::Continue => self.front.soil(conn),
761 }
762 }
763
764 if self.server.parked_here() != 0 {
769 self.serve_waiters();
770 }
771 yo_reactor::Flow::Next
772 }
773
774 fn flush(&mut self) {
775 if self.server.parked_here() != 0 {
786 self.server.refresh_clock();
787 self.serve_waiters();
788 }
789
790 if !self.held.is_empty() {
796 self.server.refresh_clock();
797 if self.server.paused(self.server.now_ms()).is_none() {
798 self.resume();
799 }
800 }
801
802 if self.server.kills() != 0 {
807 self.reap();
808 }
809
810 if self.server.mail_here() != 0 {
815 self.deliver();
816 }
817
818 let mut dirty = self.front.take_dirty();
822 let mut at = 0;
823 while at < dirty.len() {
824 let conn = dirty[at];
825 match self.front.write_out(conn) {
826 Wrote::Owed => at += 1,
830 Wrote::Done => {
831 dirty.swap_remove(at);
832 }
833 Wrote::Ended(client) => {
834 self.forget(client);
835 dirty.swap_remove(at);
836 }
837 }
838 }
839 self.front.give_dirty(dirty);
840 self.note_buffers();
841 }
842
843 fn maintain(&mut self, budget: &mut yo_reactor::Budget) {
844 if !budget.spend(1) {
847 return;
848 }
849 self.tick();
850 let looks = budget.left() as usize;
856 let spent = self.server.expire_slice(looks);
857 budget.spend(u32::try_from(spent).unwrap_or(u32::MAX));
858 }
859}
860
861pub fn pump<S: Sink>(reactor: &mut Reactor<Wire<S>>, batch: &mut Vec<Cmd>) -> usize {
868 let mut ran = 0;
869 reactor.engine_mut().tick();
870 loop {
876 loop {
877 batch.clear();
878 if reactor.engine_mut().take_ready(batch, BATCH_MAX) == 0 {
879 break;
880 }
881 let armed = yo_alloc::guard();
891 ran += reactor.execute_all(batch.drain(..));
892 drop(armed);
893 reactor.engine_mut().flush();
894 reactor.engine_mut().maintain();
897 }
898 reactor.engine_mut().maintain();
901 reactor.engine_mut().flush();
909 if reactor.engine().ready() == 0 {
910 break;
911 }
912 }
913 ran
914}
915
916#[cfg(test)]
917mod tests {
918 use super::*;
919
920 fn wire(args: &[&[u8]]) -> Vec<u8> {
922 let mut b = format!("*{}\r\n", args.len()).into_bytes();
923 for a in args {
924 b.extend_from_slice(format!("${}\r\n", a.len()).as_bytes());
925 b.extend_from_slice(a);
926 b.extend_from_slice(b"\r\n");
927 }
928 b
929 }
930
931 fn engine() -> (Reactor<Wire<Recorder>>, ConnId, Vec<Cmd>) {
932 let mut r = Reactor::inline(Wire::new(Recorder::new()));
933 let conn = r.engine_mut().accept();
934 (r, conn, Vec::new())
935 }
936
937 const START_MS: u64 = 1_000_000;
939
940 fn timed() -> (Reactor<Wire<Recorder>>, ConnId, Vec<Cmd>) {
946 let server = crate::dispatch::Server::with_clock(yo_kv::Clock::fixed(START_MS));
947 let mut r = Reactor::inline(Wire::with_server(server, Recorder::new()));
948 let conn = r.engine_mut().accept();
949 (r, conn, Vec::new())
950 }
951
952 #[test]
953 fn a_pipelined_batch_comes_back_in_order_and_in_one_write() {
954 let (mut r, conn, mut batch) = engine();
955 let mut stream = wire(&[b"SET", b"k", b"v"]);
956 stream.extend(wire(&[b"GET", b"k"]));
957 stream.extend(wire(&[b"INCR", b"n"]));
958
959 r.engine_mut().feed(conn, &stream);
960 assert_eq!(r.engine().ready(), 3);
961 assert_eq!(pump(&mut r, &mut batch), 3);
962
963 assert_eq!(r.engine().sink().sent(conn), b"+OK\r\n$1\r\nv\r\n:1\r\n");
964 assert_eq!(r.engine().ready(), 0);
965 }
966
967 #[test]
970 fn a_command_split_across_reads_resumes_rather_than_restarts() {
971 let (mut r, conn, mut batch) = engine();
972 let bytes = wire(&[b"SET", b"key", b"value"]);
973
974 for at in 1..bytes.len() {
975 r.engine_mut().feed(conn, &bytes[at - 1..at]);
976 assert_eq!(r.engine().ready(), 0, "not a command yet at {at}");
977 }
978 r.engine_mut().feed(conn, &bytes[bytes.len() - 1..]);
979 assert_eq!(r.engine().ready(), 1);
980 assert_eq!(pump(&mut r, &mut batch), 1);
981 assert_eq!(r.engine().sink().sent(conn), b"+OK\r\n");
982
983 r.engine_mut().feed(conn, &wire(&[b"GET", b"key"]));
986 pump(&mut r, &mut batch);
987 assert_eq!(r.engine().sink().sent(conn), b"+OK\r\n$5\r\nvalue\r\n");
988 }
989
990 #[test]
991 fn two_connections_are_two_sessions_over_one_server() {
992 let (mut r, a, mut batch) = engine();
993 let b = r.engine_mut().accept();
994
995 r.engine_mut().feed(a, &wire(&[b"SELECT", b"3"]));
996 r.engine_mut().feed(a, &wire(&[b"SET", b"k", b"a"]));
997 r.engine_mut().feed(b, &wire(&[b"SET", b"k", b"b"]));
998 r.engine_mut().feed(a, &wire(&[b"GET", b"k"]));
999 r.engine_mut().feed(b, &wire(&[b"GET", b"k"]));
1000 pump(&mut r, &mut batch);
1001
1002 assert_eq!(r.engine().sink().sent(a), b"+OK\r\n+OK\r\n$1\r\na\r\n");
1003 assert_eq!(r.engine().sink().sent(b), b"+OK\r\n$1\r\nb\r\n");
1004 assert_eq!(r.engine().clients(), 2);
1005 }
1006
1007 #[test]
1019 fn two_threads_write_into_one_server() {
1020 const EACH: usize = 200;
1021
1022 let mut server = Server::new();
1023 server.set_threads(2);
1024 let first = Wire::with_server(server, Recorder::new());
1025 let second = Wire::over(first.shared(), Recorder::new());
1026 let server = first.shared();
1027
1028 std::thread::scope(|s| {
1029 for (at, engine) in [first, second].into_iter().enumerate() {
1030 s.spawn(move || {
1031 let mut r = Reactor::inline(engine);
1032 let mut batch = Vec::new();
1033 let conn = r.engine_mut().accept();
1034 for i in 0..EACH {
1035 let key = format!("t{at}:{i}");
1036 r.engine_mut()
1037 .feed(conn, &wire(&[b"SET", key.as_bytes(), b"v"]));
1038 pump(&mut r, &mut batch);
1039 }
1040 });
1041 }
1042 });
1043
1044 assert_eq!(server.striped_ref(0).len(), 2 * EACH);
1047 assert_eq!(server.totals().connections, 2);
1050 }
1051
1052 #[test]
1055 fn a_waiter_belongs_to_the_thread_that_parked_it() {
1056 let mut server = Server::new();
1057 server.set_threads(2);
1058 let first = Wire::with_server(server, Recorder::new());
1059 let second = Wire::over(first.shared(), Recorder::new());
1060 let server = first.shared();
1061
1062 let parked = std::sync::Barrier::new(2);
1063 let swept = std::sync::Barrier::new(2);
1064
1065 std::thread::scope(|s| {
1066 let (parked, swept) = (&parked, &swept);
1067 s.spawn(move || {
1068 let mut r = Reactor::inline(first);
1069 let mut batch = Vec::new();
1070 let conn = r.engine_mut().accept();
1071 r.engine_mut().feed(conn, &wire(&[b"BLPOP", b"a", b"0"]));
1072 pump(&mut r, &mut batch);
1073 parked.wait();
1074
1075 for _ in 0..50 {
1078 pump(&mut r, &mut batch);
1079 }
1080 swept.wait();
1081 assert!(r.engine().sink().sent(conn).is_empty(), "nothing to say");
1082 });
1083 s.spawn(move || {
1084 let mut r = Reactor::inline(second);
1085 let mut batch = Vec::new();
1086 let conn = r.engine_mut().accept();
1087 r.engine_mut().feed(conn, &wire(&[b"BLPOP", b"b", b"0"]));
1088 pump(&mut r, &mut batch);
1089 parked.wait();
1090 swept.wait();
1091
1092 let pusher = r.engine_mut().accept();
1095 r.engine_mut().feed(pusher, &wire(&[b"RPUSH", b"b", b"v"]));
1096 pump(&mut r, &mut batch);
1097 assert_eq!(
1098 r.engine().sink().sent(conn),
1099 b"*2\r\n$1\r\nb\r\n$1\r\nv\r\n",
1100 "served by the thread that parked it"
1101 );
1102 });
1103 });
1104
1105 assert_eq!(server.parked(), 1, "and the other one is still waiting");
1106 }
1107
1108 #[test]
1114 fn a_thread_counts_the_clients_it_blocked_and_nobody_else_s() {
1115 let mut server = Server::new();
1116 server.set_threads(2);
1117 let first = Wire::with_server(server, Recorder::new());
1118 let second = Wire::over(first.shared(), Recorder::new());
1119 let server = first.shared();
1120
1121 let parked = std::sync::Barrier::new(2);
1122 let looked = std::sync::Barrier::new(2);
1123
1124 std::thread::scope(|s| {
1125 let (parked, looked) = (&parked, &looked);
1126 s.spawn(move || {
1127 let mut r = Reactor::inline(first);
1128 let mut batch = Vec::new();
1129 let conn = r.engine_mut().accept();
1130 r.engine_mut().feed(conn, &wire(&[b"BLPOP", b"a", b"0"]));
1131 pump(&mut r, &mut batch);
1132 assert_eq!(r.engine().waiting(), 1, "the one this thread blocked");
1133 parked.wait();
1134 looked.wait();
1135
1136 let other = r.engine_mut().accept();
1140 r.engine_mut().feed(other, &wire(&[b"PING"]));
1141 pump(&mut r, &mut batch);
1142 r.engine_mut().hangup(other);
1143 pump(&mut r, &mut batch);
1144 assert_eq!(r.engine().waiting(), 1, "still just the blocked one");
1145 });
1146 s.spawn(move || {
1147 let mut r = Reactor::inline(second);
1148 let mut batch = Vec::new();
1149 parked.wait();
1150
1151 pump(&mut r, &mut batch);
1154 assert_eq!(r.engine().waiting(), 0, "none of them are this one's");
1155 assert_eq!(r.engine().server().parked(), 1, "one on the server");
1156 looked.wait();
1157 });
1158 });
1159
1160 assert_eq!(server.parked(), 1);
1161 }
1162
1163 #[test]
1166 fn client_ids_are_the_server_s_to_hand_out() {
1167 let first = Wire::new(Recorder::new());
1168 let second = Wire::over(first.shared(), Recorder::new());
1169 let mut a = Reactor::inline(first);
1170 let mut b = Reactor::inline(second);
1171
1172 let (one, two) = (a.engine_mut().accept(), b.engine_mut().accept());
1173 assert_eq!(one, two, "the same slot on each front");
1174
1175 let mut batch = Vec::new();
1180 a.engine_mut().feed(one, &wire(&[b"HELLO", b"3"]));
1181 b.engine_mut().feed(two, &wire(&[b"HELLO", b"3"]));
1182 pump(&mut a, &mut batch);
1183 pump(&mut b, &mut batch);
1184
1185 let first = String::from_utf8_lossy(a.engine().sink().sent(one)).into_owned();
1186 let second = String::from_utf8_lossy(b.engine().sink().sent(two)).into_owned();
1187 assert!(first.contains(":1\r\n"), "{first}");
1188 assert!(second.contains(":2\r\n"), "{second}");
1189 }
1190
1191 #[test]
1192 fn quit_is_answered_and_then_the_connection_goes() {
1193 let (mut r, conn, mut batch) = engine();
1194 r.engine_mut().feed(conn, &wire(&[b"PING"]));
1195 r.engine_mut().feed(conn, &wire(&[b"QUIT"]));
1196 pump(&mut r, &mut batch);
1197
1198 assert_eq!(r.engine().sink().sent(conn), b"+PONG\r\n+OK\r\n");
1199 assert!(r.engine().sink().was_closed(conn));
1200 assert_eq!(r.engine().clients(), 0);
1201
1202 let again = r.engine_mut().accept();
1204 assert_eq!(again, conn);
1205 assert_eq!(r.engine().clients(), 1);
1206 }
1207
1208 #[test]
1211 fn what_a_client_pipelined_behind_quit_is_never_run() {
1212 let (mut r, conn, mut batch) = engine();
1213 let mut stream = wire(&[b"QUIT"]);
1214 stream.extend(wire(&[b"SET", b"foo", b"bar"]));
1215 r.engine_mut().feed(conn, &stream);
1216 assert_eq!(r.engine().ready(), 2);
1218 pump(&mut r, &mut batch);
1219
1220 assert_eq!(r.engine().sink().sent(conn), b"+OK\r\n");
1222 assert!(r.engine().sink().was_closed(conn));
1223
1224 r.engine_mut().sink_mut().clear();
1229 let next = r.engine_mut().accept();
1230 r.engine_mut().feed(next, &wire(&[b"GET", b"foo"]));
1231 pump(&mut r, &mut batch);
1232 assert_eq!(r.engine().sink().sent(next), b"$-1\r\n");
1233 }
1234
1235 #[test]
1245 fn a_slot_that_last_spoke_resp3_answers_the_next_client_in_resp2() {
1246 let (mut r, conn, mut batch) = engine();
1247 r.engine_mut().feed(conn, &wire(&[b"HELLO", b"3"]));
1248 r.engine_mut().feed(conn, &wire(&[b"GET", b"nothing"]));
1249 pump(&mut r, &mut batch);
1250 assert!(r.engine().sink().sent(conn).ends_with(b"_\r\n"));
1251 r.engine_mut().feed(conn, &wire(&[b"QUIT"]));
1252 pump(&mut r, &mut batch);
1253
1254 r.engine_mut().sink_mut().clear();
1255 let next = r.engine_mut().accept();
1256 assert_eq!(next, conn, "the same slot, which is what this is about");
1257 r.engine_mut().feed(next, &wire(&[b"GET", b"nothing"]));
1258 pump(&mut r, &mut batch);
1259 assert_eq!(r.engine().sink().sent(next), b"$-1\r\n");
1260 }
1261
1262 #[test]
1264 fn commands_that_arrived_before_a_protocol_error_are_still_answered() {
1265 let (mut r, conn, mut batch) = engine();
1266 let mut stream = wire(&[b"SET", b"k", b"v"]);
1267 stream.extend(wire(&[b"GET", b"k"]));
1268 stream.extend_from_slice(b"*1\r\n+notabulk\r\n");
1269 r.engine_mut().feed(conn, &stream);
1270 pump(&mut r, &mut batch);
1271
1272 let sent = r.engine().sink().sent(conn);
1275 assert!(
1276 sent.starts_with(b"+OK\r\n$1\r\nv\r\n-ERR Protocol error: "),
1277 "{sent:?}"
1278 );
1279 assert!(r.engine().sink().was_closed(conn));
1280 }
1281
1282 #[test]
1283 fn a_protocol_error_is_written_and_closes_the_connection() {
1284 let (mut r, conn, mut batch) = engine();
1285 r.engine_mut().feed(conn, b"*1\r\n+notabulk\r\n");
1287 pump(&mut r, &mut batch);
1288
1289 let sent = r.engine().sink().sent(conn);
1290 assert!(sent.starts_with(b"-ERR Protocol error: "), "{sent:?}");
1291 assert!(r.engine().sink().was_closed(conn));
1292 assert_eq!(r.engine().clients(), 0);
1293 }
1294
1295 #[test]
1299 fn a_decoder_that_came_back_mid_command_starts_the_next_one_clean() {
1300 let (mut r, conn, mut batch) = engine();
1301 r.engine_mut()
1303 .feed(conn, b"*3\r\n$3\r\nSET\r\n$1\r\nx\r\n$blabla\r\n");
1304 pump(&mut r, &mut batch);
1305 let sent = r.engine().sink().sent(conn);
1306 assert!(
1307 sent.starts_with(b"-ERR Protocol error: invalid bulk length"),
1308 "{sent:?}"
1309 );
1310
1311 r.engine_mut().sink_mut().clear();
1315 let next = r.engine_mut().accept();
1316 r.engine_mut().feed(next, &wire(&[b"GET", b"k"]));
1317 pump(&mut r, &mut batch);
1318 assert_eq!(r.engine().sink().sent(next), b"$-1\r\n");
1319
1320 r.engine_mut().sink_mut().clear();
1321 let third = r.engine_mut().accept();
1322 r.engine_mut().feed(third, b"*1\r\n+notabulk\r\n");
1323 pump(&mut r, &mut batch);
1324 let sent = r.engine().sink().sent(third);
1325 assert!(sent.starts_with(b"-ERR Protocol error: "), "{sent:?}");
1326 }
1327
1328 #[test]
1331 fn a_hangup_with_commands_in_flight_waits_for_them() {
1332 let (mut r, conn, mut batch) = engine();
1333 r.engine_mut().feed(conn, &wire(&[b"SET", b"k", b"v"]));
1334 r.engine_mut().feed(conn, &wire(&[b"GET", b"k"]));
1335
1336 batch.clear();
1337 r.engine_mut().take_ready(&mut batch, BATCH_MAX);
1338 r.engine_mut().hangup(conn);
1339 assert_eq!(r.engine().clients(), 1, "still holding the buffer");
1340
1341 r.execute_all(batch.drain(..));
1342 r.engine_mut().flush();
1343 assert_eq!(r.engine().clients(), 0);
1344 assert!(r.engine().sink().sent(conn).is_empty(), "nobody to answer");
1345
1346 let decoders = r.engine().decoders();
1349 let again = r.engine_mut().accept();
1350 assert_eq!(again, conn);
1351 r.engine_mut().feed(again, &wire(&[b"PING"]));
1352 pump(&mut r, &mut batch);
1353 assert_eq!(r.engine().sink().sent(again), b"+PONG\r\n");
1354 assert_eq!(r.engine().decoders(), decoders);
1355 }
1356
1357 #[test]
1360 fn the_buffers_and_the_decoder_pool_stop_growing() {
1361 let (mut r, conn, mut batch) = engine();
1362 let mut stream = Vec::new();
1363 for i in 0..32 {
1364 stream.extend(wire(&[b"SET", format!("k{i}").as_bytes(), b"v"]));
1365 }
1366
1367 r.engine_mut().feed(conn, &stream);
1368 pump(&mut r, &mut batch);
1369 let decoders = r.engine().decoders();
1370 let batch_cap = batch.capacity();
1371
1372 for _ in 0..10 {
1373 r.engine_mut().feed(conn, &stream);
1374 pump(&mut r, &mut batch);
1375 }
1376 assert_eq!(r.engine().decoders(), decoders, "the pool is reused");
1377 assert_eq!(batch.capacity(), batch_cap, "the batch buffer is reused");
1378 assert!(
1379 decoders <= BATCH_MAX + 1,
1380 "{decoders} decoders for 32 commands"
1381 );
1382 }
1383
1384 #[test]
1393 fn a_pipelining_client_does_not_grow_the_read_buffer() {
1394 let (mut r, conn, mut batch) = engine();
1395 let mut round = Vec::new();
1396 for i in 0..16 {
1397 round.extend(wire(&[b"SET", format!("k{i}").as_bytes(), b"v"]));
1398 }
1399
1400 r.engine_mut().feed(conn, &round);
1401 pump(&mut r, &mut batch);
1402 r.engine_mut().sink_mut().clear();
1403 let after_one = r.engine().buffer_bytes();
1404
1405 let rounds = if cfg!(miri) { 50 } else { 1000 };
1414 for _ in 0..rounds {
1415 r.engine_mut().feed(conn, &round);
1416 pump(&mut r, &mut batch);
1417 r.engine_mut().sink_mut().clear();
1418 }
1419
1420 assert_eq!(
1421 r.engine().buffer_bytes(),
1422 after_one,
1423 "the buffers grew over {rounds} rounds of the same sixteen commands"
1424 );
1425 assert!(
1426 r.engine().server().memory_bytes() >= after_one,
1427 "the buffers are counted in what the server reports"
1428 );
1429 }
1430
1431 #[test]
1434 fn a_command_split_across_reads_survives_compaction() {
1435 let (mut r, conn, mut batch) = engine();
1436 let cmd = wire(&[b"SET", b"key", b"value"]);
1437 let (head, tail) = cmd.split_at(cmd.len() - 4);
1438
1439 r.engine_mut().feed(conn, &wire(&[b"PING"]));
1442 r.engine_mut().feed(conn, head);
1443 pump(&mut r, &mut batch);
1444 assert_eq!(r.engine().sink().sent(conn), b"+PONG\r\n");
1445
1446 r.engine_mut().feed(conn, tail);
1448 pump(&mut r, &mut batch);
1449 assert_eq!(r.engine().sink().sent(conn), b"+PONG\r\n+OK\r\n");
1450
1451 r.engine_mut().feed(conn, &wire(&[b"GET", b"key"]));
1452 pump(&mut r, &mut batch);
1453 assert!(r.engine().sink().sent(conn).ends_with(b"$5\r\nvalue\r\n"));
1454 }
1455
1456 #[test]
1459 fn the_batch_goes_through_the_reactors_two_walks() {
1460 let (mut r, conn, mut batch) = engine();
1461 for i in 0..100 {
1462 r.engine_mut()
1463 .feed(conn, &wire(&[b"INCR", format!("k{}", i % 7).as_bytes()]));
1464 }
1465 let ran = pump(&mut r, &mut batch);
1466
1467 assert_eq!(ran, 100);
1468 assert_eq!(r.commands(), 100);
1469 assert_eq!(r.turns(), 2);
1471 assert!(r.engine().sink().sent(conn).ends_with(b":15\r\n"));
1473 }
1474
1475 #[derive(Default)]
1478 struct Trickle {
1479 sent: Vec<u8>,
1480 writes: usize,
1481 }
1482
1483 impl Sink for Trickle {
1484 fn write(&mut self, _conn: ConnId, bytes: &[u8]) -> usize {
1485 self.writes += 1;
1486 let n = bytes.len().min(4);
1487 self.sent.extend_from_slice(&bytes[..n]);
1488 n
1489 }
1490 }
1491
1492 #[test]
1495 fn a_blpop_on_a_list_with_something_in_it_never_waits() {
1496 let (mut r, conn, mut batch) = engine();
1497 r.engine_mut().feed(conn, &wire(&[b"RPUSH", b"q", b"a"]));
1498 r.engine_mut().feed(conn, &wire(&[b"BLPOP", b"q", b"0"]));
1499 pump(&mut r, &mut batch);
1500
1501 assert_eq!(
1502 r.engine().sink().sent(conn),
1503 b":1\r\n*2\r\n$1\r\nq\r\n$1\r\na\r\n"
1504 );
1505 assert_eq!(r.engine().server().parked(), 0);
1506 }
1507
1508 #[test]
1511 fn a_parked_client_is_answered_by_another_connections_push() {
1512 let (mut r, a, mut batch) = engine();
1513 let b = r.engine_mut().accept();
1514
1515 r.engine_mut().feed(a, &wire(&[b"BLPOP", b"q", b"0"]));
1516 pump(&mut r, &mut batch);
1517 assert!(r.engine().sink().sent(a).is_empty(), "nothing to say yet");
1518 assert_eq!(r.engine().server().parked(), 1);
1519
1520 r.engine_mut().feed(b, &wire(&[b"RPUSH", b"q", b"one"]));
1521 pump(&mut r, &mut batch);
1522
1523 assert_eq!(r.engine().sink().sent(a), b"*2\r\n$1\r\nq\r\n$3\r\none\r\n");
1524 assert_eq!(r.engine().sink().sent(b), b":1\r\n");
1527 assert_eq!(r.engine().server().parked(), 0);
1528 }
1529
1530 #[test]
1533 fn only_a_list_arriving_under_a_named_key_wakes_a_waiter() {
1534 let (mut r, a, mut batch) = engine();
1535 let b = r.engine_mut().accept();
1536 r.engine_mut().feed(a, &wire(&[b"BLPOP", b"q", b"0"]));
1537 pump(&mut r, &mut batch);
1538
1539 r.engine_mut()
1540 .feed(b, &wire(&[b"RPUSH", b"elsewhere", b"x"]));
1541 r.engine_mut().feed(b, &wire(&[b"SADD", b"q", b"x"]));
1542 pump(&mut r, &mut batch);
1543
1544 assert!(r.engine().sink().sent(a).is_empty());
1545 assert_eq!(r.engine().server().parked(), 1, "still waiting");
1546 assert_eq!(r.engine().sink().sent(b), b":1\r\n:1\r\n");
1549 }
1550
1551 #[test]
1554 fn two_parked_clients_are_served_in_the_order_they_arrived() {
1555 let (mut r, a, mut batch) = engine();
1556 let b = r.engine_mut().accept();
1557 let c = r.engine_mut().accept();
1558
1559 r.engine_mut().feed(a, &wire(&[b"BLPOP", b"q", b"0"]));
1560 pump(&mut r, &mut batch);
1561 r.engine_mut().feed(b, &wire(&[b"BLPOP", b"q", b"0"]));
1562 pump(&mut r, &mut batch);
1563 assert_eq!(r.engine().server().parked(), 2);
1564
1565 r.engine_mut()
1566 .feed(c, &wire(&[b"RPUSH", b"q", b"first", b"second"]));
1567 pump(&mut r, &mut batch);
1568
1569 assert_eq!(
1570 r.engine().sink().sent(a),
1571 b"*2\r\n$1\r\nq\r\n$5\r\nfirst\r\n"
1572 );
1573 assert_eq!(
1574 r.engine().sink().sent(b),
1575 b"*2\r\n$1\r\nq\r\n$6\r\nsecond\r\n"
1576 );
1577 assert_eq!(r.engine().server().parked(), 0);
1578 }
1579
1580 #[test]
1583 fn what_a_client_pipelined_behind_a_block_waits_for_the_block() {
1584 let (mut r, a, mut batch) = engine();
1585 let b = r.engine_mut().accept();
1586
1587 let mut stream = wire(&[b"BLPOP", b"q", b"0"]);
1590 stream.extend(wire(&[b"PING"]));
1591 r.engine_mut().feed(a, &stream);
1592 pump(&mut r, &mut batch);
1593 assert!(
1594 r.engine().sink().sent(a).is_empty(),
1595 "the PING went out in front of the answer it was sent behind"
1596 );
1597
1598 r.engine_mut().feed(a, &wire(&[b"ECHO", b"after"]));
1600 pump(&mut r, &mut batch);
1601 assert!(r.engine().sink().sent(a).is_empty());
1602
1603 r.engine_mut().feed(b, &wire(&[b"RPUSH", b"q", b"x"]));
1604 pump(&mut r, &mut batch);
1605 assert_eq!(
1606 r.engine().sink().sent(a),
1607 b"*2\r\n$1\r\nq\r\n$1\r\nx\r\n+PONG\r\n$5\r\nafter\r\n"
1608 );
1609 }
1610
1611 #[test]
1616 fn a_waiter_is_served_between_two_pipelined_pushes() {
1617 let (mut r, a, mut batch) = engine();
1618 let b = r.engine_mut().accept();
1619 r.engine_mut()
1620 .feed(a, &wire(&[b"BLPOP", b"p1", b"p2", b"0"]));
1621 pump(&mut r, &mut batch);
1622
1623 let mut stream = wire(&[b"RPUSH", b"p2", b"second"]);
1624 stream.extend(wire(&[b"RPUSH", b"p1", b"first"]));
1625 r.engine_mut().feed(b, &stream);
1626 pump(&mut r, &mut batch);
1627
1628 assert_eq!(
1629 r.engine().sink().sent(a),
1630 b"*2\r\n$2\r\np2\r\n$6\r\nsecond\r\n"
1631 );
1632 r.engine_mut()
1634 .feed(b, &wire(&[b"LRANGE", b"p1", b"0", b"-1"]));
1635 pump(&mut r, &mut batch);
1636 assert!(
1637 r.engine()
1638 .sink()
1639 .sent(b)
1640 .ends_with(b"*1\r\n$5\r\nfirst\r\n")
1641 );
1642 }
1643
1644 #[test]
1648 fn a_waiter_woken_by_another_waiter() {
1649 let (mut r, a, mut batch) = engine();
1650 let b = r.engine_mut().accept();
1651 let c = r.engine_mut().accept();
1652
1653 r.engine_mut()
1654 .feed(a, &wire(&[b"BLMOVE", b"x", b"y", b"LEFT", b"RIGHT", b"0"]));
1655 pump(&mut r, &mut batch);
1656 r.engine_mut().feed(b, &wire(&[b"BLPOP", b"y", b"0"]));
1657 pump(&mut r, &mut batch);
1658 assert_eq!(r.engine().server().parked(), 2);
1659
1660 r.engine_mut().feed(c, &wire(&[b"RPUSH", b"x", b"chain"]));
1661 pump(&mut r, &mut batch);
1662
1663 assert_eq!(r.engine().sink().sent(a), b"$5\r\nchain\r\n");
1664 assert_eq!(
1665 r.engine().sink().sent(b),
1666 b"*2\r\n$1\r\ny\r\n$5\r\nchain\r\n"
1667 );
1668 assert_eq!(r.engine().server().parked(), 0);
1669 }
1670
1671 #[test]
1674 fn a_waiter_is_only_woken_on_the_database_it_blocked_on() {
1675 let (mut r, a, mut batch) = engine();
1676 let b = r.engine_mut().accept();
1677 r.engine_mut().feed(a, &wire(&[b"SELECT", b"3"]));
1678 r.engine_mut().feed(a, &wire(&[b"BLPOP", b"q", b"0"]));
1679 pump(&mut r, &mut batch);
1680 assert_eq!(r.engine().sink().sent(a), b"+OK\r\n");
1681
1682 r.engine_mut().feed(b, &wire(&[b"RPUSH", b"q", b"wrongdb"]));
1683 pump(&mut r, &mut batch);
1684 assert_eq!(r.engine().sink().sent(a), b"+OK\r\n", "still waiting");
1685
1686 r.engine_mut().feed(b, &wire(&[b"SELECT", b"3"]));
1687 r.engine_mut().feed(b, &wire(&[b"RPUSH", b"q", b"rightdb"]));
1688 pump(&mut r, &mut batch);
1689 assert!(r.engine().sink().sent(a).ends_with(b"$7\r\nrightdb\r\n"));
1690 }
1691
1692 #[test]
1694 fn a_client_that_waited_long_enough_gets_a_null_array() {
1695 let (mut r, conn, mut batch) = timed();
1696 r.engine_mut().feed(conn, &wire(&[b"BLPOP", b"q", b"30"]));
1697 pump(&mut r, &mut batch);
1698 assert!(r.engine().sink().sent(conn).is_empty());
1699
1700 r.engine_mut().server_mut().set_clock_ms(START_MS + 29_999);
1701 pump(&mut r, &mut batch);
1702 assert!(
1703 r.engine().sink().sent(conn).is_empty(),
1704 "a millisecond short"
1705 );
1706
1707 r.engine_mut().server_mut().set_clock_ms(START_MS + 30_000);
1708 pump(&mut r, &mut batch);
1709 assert_eq!(r.engine().sink().sent(conn), b"*-1\r\n");
1711 assert_eq!(r.engine().server().parked(), 0);
1712 }
1713
1714 #[test]
1718 fn every_blocking_command_times_out_with_the_same_null_array() {
1719 for cmd in [
1720 &[b"BLPOP".as_slice(), b"q", b"0.001"][..],
1721 &[b"BRPOP", b"q", b"0.001"],
1722 &[b"BLMOVE", b"q", b"d", b"LEFT", b"RIGHT", b"0.001"],
1723 &[b"BRPOPLPUSH", b"q", b"d", b"0.001"],
1724 &[b"BLMPOP", b"0.001", b"1", b"q", b"LEFT"],
1725 ] {
1726 let (mut r, conn, mut batch) = timed();
1727 r.engine_mut().feed(conn, &wire(cmd));
1728 pump(&mut r, &mut batch);
1729 r.engine_mut().server_mut().set_clock_ms(START_MS + 1);
1730 pump(&mut r, &mut batch);
1731 assert_eq!(r.engine().sink().sent(conn), b"*-1\r\n", "for {cmd:?}");
1732 }
1733 }
1734
1735 #[test]
1738 fn a_waiter_that_timed_out_does_not_eat_a_later_push() {
1739 let (mut r, a, mut batch) = timed();
1740 let b = r.engine_mut().accept();
1741 r.engine_mut().feed(a, &wire(&[b"BLPOP", b"q", b"1"]));
1742 pump(&mut r, &mut batch);
1743 r.engine_mut().server_mut().set_clock_ms(START_MS + 1000);
1744 pump(&mut r, &mut batch);
1745 assert_eq!(r.engine().sink().sent(a), b"*-1\r\n");
1746
1747 r.engine_mut().feed(b, &wire(&[b"RPUSH", b"q", b"late"]));
1748 r.engine_mut()
1749 .feed(b, &wire(&[b"LRANGE", b"q", b"0", b"-1"]));
1750 pump(&mut r, &mut batch);
1751 assert_eq!(r.engine().sink().sent(a), b"*-1\r\n", "nothing more");
1752 assert!(r.engine().sink().sent(b).ends_with(b"*1\r\n$4\r\nlate\r\n"));
1753 }
1754
1755 #[test]
1760 fn a_client_that_goes_away_while_it_waits_takes_its_waiter_with_it() {
1761 let (mut r, a, mut batch) = engine();
1762 let b = r.engine_mut().accept();
1763 r.engine_mut().feed(a, &wire(&[b"BLPOP", b"q", b"0"]));
1764 pump(&mut r, &mut batch);
1765 assert_eq!(r.engine().server().parked(), 1);
1766
1767 r.engine_mut().hangup(a);
1768 pump(&mut r, &mut batch);
1769 assert_eq!(r.engine().server().parked(), 0);
1770 assert_eq!(r.engine().clients(), 1);
1771
1772 let again = r.engine_mut().accept();
1775 assert_eq!(again, a);
1776 r.engine_mut().feed(b, &wire(&[b"RPUSH", b"q", b"x"]));
1777 r.engine_mut()
1778 .feed(again, &wire(&[b"LRANGE", b"q", b"0", b"-1"]));
1779 pump(&mut r, &mut batch);
1780 assert_eq!(r.engine().sink().sent(again), b"*1\r\n$1\r\nx\r\n");
1781 }
1782
1783 #[test]
1787 fn a_hangup_while_parked_gives_back_the_slot_and_the_decoders() {
1788 let (mut r, a, mut batch) = engine();
1789 let mut stream = wire(&[b"BLPOP", b"q", b"0"]);
1790 stream.extend(wire(&[b"PING"]));
1791 stream.extend(wire(&[b"PING"]));
1792 r.engine_mut().feed(a, &stream);
1793 pump(&mut r, &mut batch);
1794
1795 let decoders = r.engine().decoders();
1796 r.engine_mut().hangup(a);
1797 pump(&mut r, &mut batch);
1798
1799 assert_eq!(r.engine().clients(), 0);
1800 assert!(r.engine().sink().was_closed(a));
1801 assert_eq!(r.engine().decoders(), decoders, "the pool came back whole");
1802 let again = r.engine_mut().accept();
1803 assert_eq!(again, a);
1804 r.engine_mut().feed(again, &wire(&[b"PING"]));
1805 pump(&mut r, &mut batch);
1806 assert_eq!(r.engine().sink().sent(again), b"+PONG\r\n");
1807 }
1808
1809 #[test]
1812 fn a_published_message_lands_on_the_subscriber() {
1813 let (mut r, sub, mut batch) = engine();
1814 let pubr = r.engine_mut().accept();
1815
1816 r.engine_mut().feed(sub, &wire(&[b"SUBSCRIBE", b"news"]));
1817 pump(&mut r, &mut batch);
1818 assert_eq!(
1819 r.engine().sink().sent(sub),
1820 b"*3\r\n$9\r\nsubscribe\r\n$4\r\nnews\r\n:1\r\n"
1821 );
1822 r.engine_mut().sink_mut().clear();
1823
1824 r.engine_mut()
1825 .feed(pubr, &wire(&[b"PUBLISH", b"news", b"hi"]));
1826 pump(&mut r, &mut batch);
1827 assert_eq!(r.engine().sink().sent(pubr), b":1\r\n");
1828 assert_eq!(
1829 r.engine().sink().sent(sub),
1830 b"*3\r\n$7\r\nmessage\r\n$4\r\nnews\r\n$2\r\nhi\r\n"
1831 );
1832 }
1833
1834 #[test]
1837 fn a_pattern_subscriber_is_told_the_pattern_and_the_channel() {
1838 let (mut r, sub, mut batch) = engine();
1839 let pubr = r.engine_mut().accept();
1840
1841 r.engine_mut().feed(sub, &wire(&[b"PSUBSCRIBE", b"ne*"]));
1842 pump(&mut r, &mut batch);
1843 r.engine_mut().sink_mut().clear();
1844
1845 r.engine_mut()
1846 .feed(pubr, &wire(&[b"PUBLISH", b"news", b"hi"]));
1847 pump(&mut r, &mut batch);
1848 assert_eq!(r.engine().sink().sent(pubr), b":1\r\n");
1849 assert_eq!(
1850 r.engine().sink().sent(sub),
1851 b"*4\r\n$8\r\npmessage\r\n$3\r\nne*\r\n$4\r\nnews\r\n$2\r\nhi\r\n"
1852 );
1853 }
1854
1855 #[test]
1860 fn resp2_takes_almost_nothing_from_a_subscriber() {
1861 let (mut r, conn, mut batch) = engine();
1862
1863 r.engine_mut().feed(conn, &wire(&[b"SUBSCRIBE", b"a"]));
1864 pump(&mut r, &mut batch);
1865 r.engine_mut().sink_mut().clear();
1866
1867 r.engine_mut().feed(conn, &wire(&[b"GET", b"k"]));
1868 pump(&mut r, &mut batch);
1869 assert_eq!(
1870 r.engine().sink().sent(conn),
1871 b"-ERR Can't execute 'get': only (P|S)SUBSCRIBE / (P|S)UNSUBSCRIBE / PING / QUIT / RESET are allowed in this context\r\n"
1872 );
1873 r.engine_mut().sink_mut().clear();
1874
1875 r.engine_mut().feed(conn, &wire(&[b"PING"]));
1877 pump(&mut r, &mut batch);
1878 assert_eq!(
1879 r.engine().sink().sent(conn),
1880 b"*2\r\n$4\r\npong\r\n$0\r\n\r\n"
1881 );
1882 r.engine_mut().sink_mut().clear();
1883
1884 r.engine_mut().feed(conn, &wire(&[b"UNSUBSCRIBE", b"a"]));
1886 r.engine_mut().feed(conn, &wire(&[b"GET", b"k"]));
1887 pump(&mut r, &mut batch);
1888 assert_eq!(
1889 r.engine().sink().sent(conn),
1890 b"*3\r\n$11\r\nunsubscribe\r\n$1\r\na\r\n:0\r\n$-1\r\n"
1891 );
1892 }
1893
1894 #[test]
1898 fn a_shard_channel_and_a_pattern_do_not_hear_each_other() {
1899 let (mut r, sub, mut batch) = engine();
1900 let pubr = r.engine_mut().accept();
1901
1902 r.engine_mut().feed(sub, &wire(&[b"SSUBSCRIBE", b"sx"]));
1903 r.engine_mut().feed(sub, &wire(&[b"PSUBSCRIBE", b"s*"]));
1904 pump(&mut r, &mut batch);
1905 r.engine_mut().sink_mut().clear();
1906
1907 r.engine_mut()
1908 .feed(pubr, &wire(&[b"SPUBLISH", b"sx", b"one"]));
1909 pump(&mut r, &mut batch);
1910 assert_eq!(r.engine().sink().sent(pubr), b":1\r\n");
1911 assert_eq!(
1912 r.engine().sink().sent(sub),
1913 b"*3\r\n$8\r\nsmessage\r\n$2\r\nsx\r\n$3\r\none\r\n"
1914 );
1915 r.engine_mut().sink_mut().clear();
1916
1917 r.engine_mut()
1918 .feed(pubr, &wire(&[b"PUBLISH", b"sx", b"two"]));
1919 pump(&mut r, &mut batch);
1920 assert_eq!(r.engine().sink().sent(pubr), b":1\r\n");
1921 assert_eq!(
1922 r.engine().sink().sent(sub),
1923 b"*4\r\n$8\r\npmessage\r\n$2\r\ns*\r\n$2\r\nsx\r\n$3\r\ntwo\r\n"
1924 );
1925 }
1926
1927 #[test]
1931 fn a_subscriber_that_goes_away_leaves_the_registry() {
1932 let (mut r, sub, mut batch) = engine();
1933 let pubr = r.engine_mut().accept();
1934
1935 r.engine_mut().feed(sub, &wire(&[b"SUBSCRIBE", b"news"]));
1936 pump(&mut r, &mut batch);
1937 r.engine_mut().hangup(sub);
1938 pump(&mut r, &mut batch);
1939 r.engine_mut().sink_mut().clear();
1940
1941 r.engine_mut()
1942 .feed(pubr, &wire(&[b"PUBLISH", b"news", b"hi"]));
1943 pump(&mut r, &mut batch);
1944 assert_eq!(r.engine().sink().sent(pubr), b":0\r\n");
1945
1946 let next = r.engine_mut().accept();
1948 assert_eq!(next, sub);
1949 r.engine_mut()
1950 .feed(pubr, &wire(&[b"PUBLISH", b"news", b"hi"]));
1951 pump(&mut r, &mut batch);
1952 assert_eq!(r.engine().sink().sent(next), b"");
1953 }
1954
1955 #[test]
1964 fn resp3_delivers_a_message_as_a_push() {
1965 let (mut r, conn, mut batch) = engine();
1966
1967 r.engine_mut().feed(conn, &wire(&[b"HELLO", b"3"]));
1968 r.engine_mut().feed(conn, &wire(&[b"SUBSCRIBE", b"a"]));
1969 pump(&mut r, &mut batch);
1970 r.engine_mut().sink_mut().clear();
1971
1972 r.engine_mut().feed(conn, &wire(&[b"GET", b"k"]));
1973 r.engine_mut().feed(conn, &wire(&[b"PUBLISH", b"a", b"w"]));
1974 pump(&mut r, &mut batch);
1975 assert_eq!(
1976 r.engine().sink().sent(conn),
1977 b"_\r\n:1\r\n>3\r\n$7\r\nmessage\r\n$1\r\na\r\n$1\r\nw\r\n"
1978 );
1979 }
1980
1981 #[test]
1984 fn a_write_reaches_a_keyspace_subscriber() {
1985 let (mut r, sub, mut batch) = engine();
1986 let writer = r.engine_mut().accept();
1987
1988 r.engine_mut().feed(
1989 writer,
1990 &wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", b"KEA"]),
1991 );
1992 r.engine_mut()
1993 .feed(sub, &wire(&[b"PSUBSCRIBE", b"__key*@0__:*"]));
1994 pump(&mut r, &mut batch);
1995 r.engine_mut().sink_mut().clear();
1996
1997 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
1998 pump(&mut r, &mut batch);
1999 assert_eq!(r.engine().sink().sent(writer), b"+OK\r\n");
2000 assert_eq!(
2001 r.engine().sink().sent(sub),
2002 b"*4\r\n$8\r\npmessage\r\n$12\r\n__key*@0__:*\r\n\
2003 $16\r\n__keyspace@0__:k\r\n$3\r\nset\r\n\
2004 *4\r\n$8\r\npmessage\r\n$12\r\n__key*@0__:*\r\n\
2005 $18\r\n__keyevent@0__:set\r\n$1\r\nk\r\n"
2006 );
2007 }
2008
2009 #[test]
2012 fn a_write_says_nothing_until_the_setting_turns_it_on() {
2013 let (mut r, sub, mut batch) = engine();
2014 let writer = r.engine_mut().accept();
2015
2016 r.engine_mut()
2017 .feed(sub, &wire(&[b"PSUBSCRIBE", b"__key*@0__:*"]));
2018 pump(&mut r, &mut batch);
2019 r.engine_mut().sink_mut().clear();
2020
2021 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
2022 pump(&mut r, &mut batch);
2023 assert_eq!(r.engine().sink().sent(sub), b"");
2024 }
2025
2026 #[test]
2029 fn only_the_classes_that_were_asked_for_are_published() {
2030 let (mut r, sub, mut batch) = engine();
2031 let writer = r.engine_mut().accept();
2032
2033 r.engine_mut().feed(
2034 writer,
2035 &wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", b"Eg"]),
2036 );
2037 r.engine_mut()
2038 .feed(sub, &wire(&[b"PSUBSCRIBE", b"__key*@0__:*"]));
2039 pump(&mut r, &mut batch);
2040 r.engine_mut().sink_mut().clear();
2041
2042 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
2043 r.engine_mut().feed(writer, &wire(&[b"DEL", b"k"]));
2044 pump(&mut r, &mut batch);
2045 assert_eq!(
2046 r.engine().sink().sent(sub),
2047 b"*4\r\n$8\r\npmessage\r\n$12\r\n__key*@0__:*\r\n\
2048 $18\r\n__keyevent@0__:del\r\n$1\r\nk\r\n"
2049 );
2050 }
2051
2052 #[test]
2055 fn a_write_with_a_deadline_on_it_says_two_things() {
2056 let (mut r, sub, mut batch) = engine();
2057 let writer = r.engine_mut().accept();
2058
2059 r.engine_mut().feed(
2060 writer,
2061 &wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", b"EA"]),
2062 );
2063 r.engine_mut()
2064 .feed(sub, &wire(&[b"PSUBSCRIBE", b"__keyevent@0__:*"]));
2065 pump(&mut r, &mut batch);
2066 r.engine_mut().sink_mut().clear();
2067
2068 r.engine_mut()
2069 .feed(writer, &wire(&[b"SETEX", b"k", b"100", b"v"]));
2070 pump(&mut r, &mut batch);
2071 assert_eq!(
2072 r.engine().sink().sent(sub),
2073 b"*4\r\n$8\r\npmessage\r\n$16\r\n__keyevent@0__:*\r\n\
2074 $18\r\n__keyevent@0__:set\r\n$1\r\nk\r\n\
2075 *4\r\n$8\r\npmessage\r\n$16\r\n__keyevent@0__:*\r\n\
2076 $21\r\n__keyevent@0__:expire\r\n$1\r\nk\r\n"
2077 );
2078 }
2079
2080 fn watching() -> (Reactor<Wire<Recorder>>, ConnId, ConnId, Vec<Cmd>) {
2086 watching_flags(b"EA")
2087 }
2088
2089 fn watching_flags(flags: &[u8]) -> (Reactor<Wire<Recorder>>, ConnId, ConnId, Vec<Cmd>) {
2091 let (mut r, sub, mut batch) = engine();
2092 let writer = r.engine_mut().accept();
2093 r.engine_mut().feed(
2094 writer,
2095 &wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", flags]),
2096 );
2097 r.engine_mut()
2098 .feed(sub, &wire(&[b"PSUBSCRIBE", b"__keyevent@0__:*"]));
2099 pump(&mut r, &mut batch);
2100 r.engine_mut().sink_mut().clear();
2101 (r, sub, writer, batch)
2102 }
2103
2104 fn fired(r: &Reactor<Wire<Recorder>>, sub: ConnId) -> Vec<(String, String)> {
2110 fired_on(r, sub, 0)
2111 }
2112
2113 fn fired_on(r: &Reactor<Wire<Recorder>>, sub: ConnId, db: usize) -> Vec<(String, String)> {
2116 let head = format!("__keyevent@{db}__:");
2117 let sent = String::from_utf8_lossy(r.engine().sink().sent(sub)).into_owned();
2118 let mut out = Vec::new();
2119 let mut parts = sent.split("\r\n");
2120 while let Some(p) = parts.next() {
2121 let Some(event) = p.strip_prefix(head.as_str()) else {
2122 continue;
2123 };
2124 if event == "*" {
2127 continue;
2128 }
2129 parts.next();
2130 let key = parts.next().unwrap_or_default();
2131 out.push((event.to_owned(), key.to_owned()));
2132 }
2133 out
2134 }
2135
2136 #[test]
2139 fn taking_the_last_of_a_collection_says_the_key_went_with_it() {
2140 let (mut r, sub, writer, mut batch) = watching();
2141
2142 r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"k", b"a"]));
2143 r.engine_mut().feed(writer, &wire(&[b"LPOP", b"k"]));
2144 pump(&mut r, &mut batch);
2145 assert_eq!(
2146 fired(&r, sub),
2147 [("rpush", "k"), ("lpop", "k"), ("del", "k")]
2148 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2149 );
2150 }
2151
2152 #[test]
2155 fn a_move_onto_a_member_already_there_says_only_the_removal() {
2156 let (mut r, sub, writer, mut batch) = watching();
2157
2158 r.engine_mut().feed(writer, &wire(&[b"SADD", b"a", b"m"]));
2159 r.engine_mut()
2160 .feed(writer, &wire(&[b"SADD", b"b", b"m", b"n"]));
2161 pump(&mut r, &mut batch);
2162 r.engine_mut().sink_mut().clear();
2163
2164 r.engine_mut()
2165 .feed(writer, &wire(&[b"SMOVE", b"a", b"b", b"m"]));
2166 pump(&mut r, &mut batch);
2167 assert_eq!(
2168 fired(&r, sub),
2169 [("srem", "a"), ("del", "a")].map(|(e, k)| (e.to_owned(), k.to_owned()))
2170 );
2171 }
2172
2173 #[test]
2176 fn a_score_that_did_not_move_says_nothing() {
2177 let (mut r, sub, writer, mut batch) = watching();
2178
2179 r.engine_mut()
2180 .feed(writer, &wire(&[b"ZADD", b"z", b"4", b"m"]));
2181 pump(&mut r, &mut batch);
2182 r.engine_mut().sink_mut().clear();
2183
2184 r.engine_mut()
2185 .feed(writer, &wire(&[b"ZADD", b"z", b"4", b"m"]));
2186 r.engine_mut()
2187 .feed(writer, &wire(&[b"ZINCRBY", b"z", b"0", b"m"]));
2188 pump(&mut r, &mut batch);
2189 assert_eq!(fired(&r, sub), []);
2190
2191 r.engine_mut()
2194 .feed(writer, &wire(&[b"ZINCRBY", b"z", b"1", b"m"]));
2195 pump(&mut r, &mut batch);
2196 assert_eq!(fired(&r, sub), [("zincr".to_owned(), "z".to_owned())]);
2197 }
2198
2199 #[test]
2202 fn a_stream_write_says_what_the_trim_behind_it_took() {
2203 let (mut r, sub, writer, mut batch) = watching();
2204
2205 r.engine_mut()
2206 .feed(writer, &wire(&[b"XADD", b"s", b"1-1", b"f", b"v"]));
2207 r.engine_mut().feed(
2208 writer,
2209 &wire(&[b"XADD", b"s", b"MAXLEN", b"9", b"2-1", b"f", b"v"]),
2210 );
2211 r.engine_mut().feed(
2212 writer,
2213 &wire(&[b"XADD", b"s", b"MAXLEN", b"1", b"3-1", b"f", b"v"]),
2214 );
2215 pump(&mut r, &mut batch);
2216 assert_eq!(
2217 fired(&r, sub),
2218 [("xadd", "s"), ("xadd", "s"), ("xadd", "s"), ("xtrim", "s")]
2219 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2220 );
2221 }
2222
2223 #[test]
2227 fn acknowledging_an_entry_that_is_already_gone_says_nothing() {
2228 let (mut r, sub, writer, mut batch) = watching();
2229
2230 for cmd in [
2231 wire(&[b"XADD", b"s", b"1-1", b"f", b"v"]),
2232 wire(&[b"XGROUP", b"CREATE", b"s", b"g", b"0"]),
2233 wire(&[b"XREADGROUP", b"GROUP", b"g", b"c", b"STREAMS", b"s", b">"]),
2234 wire(&[b"XDEL", b"s", b"1-1"]),
2235 ] {
2236 r.engine_mut().feed(writer, &cmd);
2237 }
2238 pump(&mut r, &mut batch);
2239 r.engine_mut().sink_mut().clear();
2240
2241 r.engine_mut().feed(
2242 writer,
2243 &wire(&[b"XACKDEL", b"s", b"g", b"IDS", b"1", b"1-1"]),
2244 );
2245 pump(&mut r, &mut batch);
2246 assert_eq!(fired(&r, sub), []);
2247 }
2248
2249 #[test]
2252 fn emptying_a_hash_says_the_key_went_with_the_last_field() {
2253 let (mut r, sub, writer, mut batch) = watching();
2254
2255 r.engine_mut()
2256 .feed(writer, &wire(&[b"HSET", b"h", b"a", b"1", b"b", b"2"]));
2257 r.engine_mut().feed(writer, &wire(&[b"HDEL", b"h", b"a"]));
2258 r.engine_mut()
2261 .feed(writer, &wire(&[b"HDEL", b"h", b"b", b"a"]));
2262 pump(&mut r, &mut batch);
2263 assert_eq!(
2264 fired(&r, sub),
2265 [("hset", "h"), ("hdel", "h"), ("hdel", "h"), ("del", "h")]
2266 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2267 );
2268 }
2269
2270 #[test]
2273 fn a_field_deadline_already_past_reads_as_a_removal() {
2274 let (mut r, sub, writer, mut batch) = watching();
2275
2276 r.engine_mut()
2277 .feed(writer, &wire(&[b"HSET", b"h", b"a", b"1", b"b", b"2"]));
2278 pump(&mut r, &mut batch);
2279 r.engine_mut().sink_mut().clear();
2280
2281 r.engine_mut().feed(
2282 writer,
2283 &wire(&[b"HEXPIRE", b"h", b"0", b"FIELDS", b"1", b"a"]),
2284 );
2285 r.engine_mut().feed(
2286 writer,
2287 &wire(&[b"HEXPIRE", b"h", b"100", b"FIELDS", b"1", b"b"]),
2288 );
2289 pump(&mut r, &mut batch);
2290 assert_eq!(
2291 fired(&r, sub),
2292 [("hdel", "h"), ("hexpire", "h")].map(|(e, k)| (e.to_owned(), k.to_owned()))
2293 );
2294 }
2295
2296 #[test]
2299 fn a_write_under_a_deadline_already_gone_says_the_write_first() {
2300 let (mut r, sub, writer, mut batch) = watching();
2301
2302 r.engine_mut().feed(
2303 writer,
2304 &wire(&[b"HSETEX", b"h", b"EXAT", b"1", b"FIELDS", b"1", b"a", b"1"]),
2305 );
2306 pump(&mut r, &mut batch);
2307 assert_eq!(
2308 fired(&r, sub),
2309 [("hset", "h"), ("hdel", "h"), ("del", "h")].map(|(e, k)| (e.to_owned(), k.to_owned()))
2310 );
2311 }
2312
2313 #[test]
2316 fn clearing_a_deadline_that_was_never_set_says_nothing() {
2317 let (mut r, sub, writer, mut batch) = watching();
2318
2319 r.engine_mut()
2320 .feed(writer, &wire(&[b"HSET", b"h", b"a", b"1", b"b", b"2"]));
2321 r.engine_mut().feed(
2322 writer,
2323 &wire(&[b"HEXPIRE", b"h", b"100", b"FIELDS", b"1", b"a"]),
2324 );
2325 pump(&mut r, &mut batch);
2326 r.engine_mut().sink_mut().clear();
2327
2328 r.engine_mut()
2329 .feed(writer, &wire(&[b"HPERSIST", b"h", b"FIELDS", b"1", b"b"]));
2330 r.engine_mut().feed(
2331 writer,
2332 &wire(&[b"HGETEX", b"h", b"PERSIST", b"FIELDS", b"1", b"b"]),
2333 );
2334 pump(&mut r, &mut batch);
2335 assert_eq!(fired(&r, sub), []);
2336
2337 r.engine_mut().feed(
2340 writer,
2341 &wire(&[b"HGETEX", b"h", b"PERSIST", b"FIELDS", b"2", b"a", b"b"]),
2342 );
2343 pump(&mut r, &mut batch);
2344 assert_eq!(fired(&r, sub), [("hpersist".to_owned(), "h".to_owned())]);
2345 }
2346
2347 #[test]
2350 fn a_key_that_was_not_there_before_says_so() {
2351 let (mut r, sub, writer, mut batch) = watching_flags(b"En");
2352
2353 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
2354 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"w"]));
2355 r.engine_mut().feed(writer, &wire(&[b"APPEND", b"k", b"x"]));
2356 r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"l", b"a"]));
2357 r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"l", b"b"]));
2358 pump(&mut r, &mut batch);
2359 assert_eq!(
2360 fired(&r, sub),
2361 [("new", "k"), ("new", "l")].map(|(e, k)| (e.to_owned(), k.to_owned()))
2362 );
2363 }
2364
2365 #[test]
2368 fn the_news_of_a_new_key_comes_before_the_write_that_made_it() {
2369 let (mut r, sub, writer, mut batch) = watching_flags(b"EAn");
2370
2371 r.engine_mut().feed(writer, &wire(&[b"SET", b"a", b"1"]));
2372 r.engine_mut().feed(writer, &wire(&[b"SET", b"b", b"2"]));
2373 pump(&mut r, &mut batch);
2374 r.engine_mut().sink_mut().clear();
2375
2376 r.engine_mut().feed(writer, &wire(&[b"RENAME", b"a", b"b"]));
2380 pump(&mut r, &mut batch);
2381 assert_eq!(
2382 fired(&r, sub),
2383 [("new", "b"), ("rename_from", "a"), ("rename_to", "b")]
2384 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2385 );
2386 }
2387
2388 #[test]
2392 fn writing_over_a_destination_is_not_a_key_arriving() {
2393 let (mut r, sub, writer, mut batch) = watching_flags(b"EAn");
2394
2395 r.engine_mut()
2396 .feed(writer, &wire(&[b"RPUSH", b"l", b"c", b"a", b"b"]));
2397 r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"d", b"x"]));
2398 pump(&mut r, &mut batch);
2399 r.engine_mut().sink_mut().clear();
2400
2401 r.engine_mut()
2402 .feed(writer, &wire(&[b"SORT", b"l", b"ALPHA", b"STORE", b"d"]));
2403 pump(&mut r, &mut batch);
2404 assert_eq!(fired(&r, sub), [("sortstore".to_owned(), "d".to_owned())]);
2405 r.engine_mut().sink_mut().clear();
2406
2407 r.engine_mut().feed(writer, &wire(&[b"DEL", b"d"]));
2410 r.engine_mut()
2411 .feed(writer, &wire(&[b"SORT", b"l", b"ALPHA", b"STORE", b"d"]));
2412 pump(&mut r, &mut batch);
2413 assert_eq!(
2414 fired(&r, sub),
2415 [("del", "d"), ("new", "d"), ("sortstore", "d")]
2416 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2417 );
2418 }
2419
2420 #[test]
2423 fn replacing_a_value_says_what_went_and_whether_the_kind_changed() {
2424 let (mut r, sub, writer, mut batch) = watching_flags(b"Eoc");
2425
2426 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
2427 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"w"]));
2428 r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"l", b"a"]));
2429 r.engine_mut().feed(writer, &wire(&[b"SET", b"l", b"v"]));
2430 pump(&mut r, &mut batch);
2431 assert_eq!(
2432 fired(&r, sub),
2433 [
2434 ("overwritten", "k"),
2435 ("overwritten", "l"),
2436 ("type_changed", "l")
2437 ]
2438 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2439 );
2440 }
2441
2442 #[test]
2445 fn a_write_that_changes_part_of_a_value_has_not_replaced_it() {
2446 let (mut r, sub, writer, mut batch) = watching_flags(b"Eoc");
2447
2448 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"1"]));
2449 r.engine_mut().feed(writer, &wire(&[b"APPEND", b"k", b"2"]));
2450 r.engine_mut().feed(writer, &wire(&[b"INCR", b"k"]));
2451 r.engine_mut()
2452 .feed(writer, &wire(&[b"SETRANGE", b"k", b"0", b"9"]));
2453 r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"l", b"a"]));
2454 r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"l", b"b"]));
2455 pump(&mut r, &mut batch);
2456 assert_eq!(fired(&r, sub), []);
2457 }
2458
2459 #[test]
2463 fn a_rename_says_what_it_replaced_after_saying_what_it_did() {
2464 let (mut r, sub, writer, mut batch) = watching_flags(b"EAnoc");
2465
2466 r.engine_mut().feed(writer, &wire(&[b"SET", b"a", b"1"]));
2467 r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"b", b"x"]));
2468 pump(&mut r, &mut batch);
2469 r.engine_mut().sink_mut().clear();
2470
2471 r.engine_mut().feed(writer, &wire(&[b"RENAME", b"a", b"b"]));
2472 pump(&mut r, &mut batch);
2473 assert_eq!(
2474 fired(&r, sub),
2475 [
2476 ("new", "b"),
2477 ("rename_from", "a"),
2478 ("rename_to", "b"),
2479 ("overwritten", "b"),
2480 ("type_changed", "b")
2481 ]
2482 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2483 );
2484 }
2485
2486 #[test]
2489 fn a_store_form_says_what_it_replaced_before_saying_what_it_did() {
2490 let (mut r, sub, writer, mut batch) = watching_flags(b"EAnoc");
2491
2492 r.engine_mut().feed(writer, &wire(&[b"SADD", b"s", b"m"]));
2493 r.engine_mut().feed(writer, &wire(&[b"SET", b"d", b"q"]));
2494 pump(&mut r, &mut batch);
2495 r.engine_mut().sink_mut().clear();
2496
2497 r.engine_mut()
2498 .feed(writer, &wire(&[b"SINTERSTORE", b"d", b"s"]));
2499 pump(&mut r, &mut batch);
2500 assert_eq!(
2501 fired(&r, sub),
2502 [
2503 ("overwritten", "d"),
2504 ("type_changed", "d"),
2505 ("sinterstore", "d")
2506 ]
2507 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2508 );
2509 }
2510
2511 #[test]
2514 fn a_bit_write_that_left_the_value_alone_says_nothing() {
2515 let (mut r, sub, writer, mut batch) = watching();
2516
2517 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"abc"]));
2521 pump(&mut r, &mut batch);
2522 r.engine_mut().sink_mut().clear();
2523
2524 r.engine_mut()
2525 .feed(writer, &wire(&[b"SETBIT", b"k", b"1", b"1"]));
2526 r.engine_mut()
2527 .feed(writer, &wire(&[b"SETBIT", b"k", b"1", b"0"]));
2528 r.engine_mut()
2529 .feed(writer, &wire(&[b"SETBIT", b"k", b"1", b"0"]));
2530 r.engine_mut()
2534 .feed(writer, &wire(&[b"SETBIT", b"k", b"100", b"0"]));
2535 pump(&mut r, &mut batch);
2536 assert_eq!(
2537 fired(&r, sub),
2538 [("setbit", "k"), ("setbit", "k")].map(|(e, k)| (e.to_owned(), k.to_owned()))
2539 );
2540 }
2541
2542 #[test]
2545 fn a_bitfield_that_wrote_the_same_values_back_says_nothing() {
2546 let (mut r, sub, writer, mut batch) = watching();
2547
2548 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"abc"]));
2549 pump(&mut r, &mut batch);
2550 r.engine_mut().sink_mut().clear();
2551
2552 r.engine_mut().feed(
2554 writer,
2555 &wire(&[b"BITFIELD", b"k", b"SET", b"u8", b"0", b"97"]),
2556 );
2557 r.engine_mut()
2558 .feed(writer, &wire(&[b"BITFIELD", b"k", b"GET", b"u8", b"0"]));
2559 r.engine_mut().feed(
2560 writer,
2561 &wire(&[b"BITFIELD", b"k", b"INCRBY", b"u8", b"0", b"0"]),
2562 );
2563 pump(&mut r, &mut batch);
2564 assert_eq!(fired(&r, sub), []);
2565
2566 r.engine_mut().feed(
2569 writer,
2570 &wire(&[b"BITFIELD", b"k", b"SET", b"u8", b"0", b"98"]),
2571 );
2572 r.engine_mut().feed(
2573 writer,
2574 &wire(&[b"BITFIELD", b"k", b"SET", b"u8", b"800", b"0"]),
2575 );
2576 pump(&mut r, &mut batch);
2577 assert_eq!(
2578 fired(&r, sub),
2579 [("setbit", "k"), ("setbit", "k")].map(|(e, k)| (e.to_owned(), k.to_owned()))
2580 );
2581 }
2582
2583 #[test]
2586 fn the_sketch_commands_say_what_a_mass_add_says() {
2587 let (mut r, sub, writer, mut batch) = watching();
2588
2589 r.engine_mut().feed(writer, &wire(&[b"PFADD", b"h", b"a"]));
2590 r.engine_mut().feed(writer, &wire(&[b"PFADD", b"h", b"a"]));
2592 r.engine_mut().feed(writer, &wire(&[b"PFADD", b"h"]));
2594 r.engine_mut().feed(writer, &wire(&[b"PFMERGE", b"d"]));
2596 r.engine_mut()
2597 .feed(writer, &wire(&[b"PFMERGE", b"d", b"h"]));
2598 pump(&mut r, &mut batch);
2599 assert_eq!(
2600 fired(&r, sub),
2601 [("pfadd", "h"), ("pfadd", "d"), ("pfadd", "d")]
2602 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2603 );
2604 }
2605
2606 #[test]
2609 fn the_geo_commands_say_what_the_sorted_set_under_them_did() {
2610 let (mut r, sub, writer, mut batch) = watching();
2611
2612 let point: &[&[u8]] = &[b"GEOADD", b"g", b"13.361389", b"38.115556", b"P"];
2613 r.engine_mut().feed(writer, &wire(point));
2614 r.engine_mut().feed(writer, &wire(point));
2616 r.engine_mut().feed(
2618 writer,
2619 &wire(&[b"GEOADD", b"g", b"14.0", b"38.115556", b"P"]),
2620 );
2621 r.engine_mut().feed(
2622 writer,
2623 &wire(&[
2624 b"GEORADIUS",
2625 b"g",
2626 b"14.0",
2627 b"38.0",
2628 b"200",
2629 b"km",
2630 b"STORE",
2631 b"d",
2632 ]),
2633 );
2634 r.engine_mut().feed(
2635 writer,
2636 &wire(&[
2637 b"GEOSEARCHSTORE",
2638 b"e",
2639 b"g",
2640 b"FROMLONLAT",
2641 b"14.0",
2642 b"38.0",
2643 b"BYRADIUS",
2644 b"200",
2645 b"km",
2646 ]),
2647 );
2648 pump(&mut r, &mut batch);
2649 assert_eq!(
2650 fired(&r, sub),
2651 [
2652 ("zadd", "g"),
2653 ("zadd", "g"),
2654 ("georadiusstore", "d"),
2655 ("geosearchstore", "e")
2656 ]
2657 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2658 );
2659 }
2660
2661 #[test]
2664 fn a_geo_store_that_found_nothing_takes_the_destination_with_it() {
2665 let (mut r, sub, writer, mut batch) = watching();
2666
2667 r.engine_mut().feed(
2668 writer,
2669 &wire(&[b"GEOADD", b"g", b"13.361389", b"38.115556", b"P"]),
2670 );
2671 r.engine_mut().feed(writer, &wire(&[b"SET", b"d", b"x"]));
2672 pump(&mut r, &mut batch);
2673 r.engine_mut().sink_mut().clear();
2674
2675 r.engine_mut().feed(
2676 writer,
2677 &wire(&[
2678 b"GEORADIUS",
2679 b"g",
2680 b"1.0",
2681 b"1.0",
2682 b"1",
2683 b"km",
2684 b"STORE",
2685 b"d",
2686 ]),
2687 );
2688 r.engine_mut().feed(
2691 writer,
2692 &wire(&[
2693 b"GEORADIUS",
2694 b"g",
2695 b"1.0",
2696 b"1.0",
2697 b"1",
2698 b"km",
2699 b"STORE",
2700 b"d",
2701 ]),
2702 );
2703 pump(&mut r, &mut batch);
2704 assert_eq!(
2705 fired(&r, sub),
2706 [("del", "d")].map(|(e, k)| (e.to_owned(), k.to_owned()))
2707 );
2708 }
2709
2710 #[test]
2714 fn a_key_that_lands_on_another_database_is_new_over_there() {
2715 let (mut r, sub, mut batch) = engine();
2716 let writer = r.engine_mut().accept();
2717 r.engine_mut().feed(
2718 writer,
2719 &wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", b"EAnoc"]),
2720 );
2721 r.engine_mut().feed(
2722 sub,
2723 &wire(&[b"PSUBSCRIBE", b"__keyevent@0__:*", b"__keyevent@1__:*"]),
2724 );
2725 r.engine_mut().feed(writer, &wire(&[b"SET", b"a", b"v"]));
2726 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
2727 pump(&mut r, &mut batch);
2728 r.engine_mut().sink_mut().clear();
2729
2730 r.engine_mut()
2731 .feed(writer, &wire(&[b"COPY", b"a", b"b", b"DB", b"1"]));
2732 r.engine_mut().feed(writer, &wire(&[b"MOVE", b"k", b"1"]));
2733 pump(&mut r, &mut batch);
2734 assert_eq!(
2735 fired_on(&r, sub, 1),
2736 [
2737 ("new", "b"),
2738 ("copy_to", "b"),
2739 ("new", "k"),
2740 ("move_to", "k")
2741 ]
2742 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2743 );
2744 assert_eq!(
2747 fired_on(&r, sub, 0),
2748 [("move_from", "k")].map(|(e, k)| (e.to_owned(), k.to_owned()))
2749 );
2750 }
2751
2752 #[test]
2758 fn a_read_that_found_nothing_says_which_key_it_was() {
2759 let (mut r, sub, writer, mut batch) = watching_flags(b"Em");
2760
2761 r.engine_mut().feed(writer, &wire(&[b"SET", b"a", b"v"]));
2762 r.engine_mut()
2763 .feed(writer, &wire(&[b"MGET", b"a", b"nk", b"nk"]));
2764 r.engine_mut()
2765 .feed(writer, &wire(&[b"EXISTS", b"nj", b"a"]));
2766 pump(&mut r, &mut batch);
2767 assert_eq!(
2768 fired(&r, sub),
2769 [("keymiss", "nk"), ("keymiss", "nk"), ("keymiss", "nj")]
2770 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2771 );
2772 }
2773
2774 #[test]
2777 fn only_the_shape_of_a_write_that_reads_says_it() {
2778 let (mut r, sub, writer, mut batch) = watching_flags(b"Em");
2779
2780 r.engine_mut().feed(writer, &wire(&[b"SET", b"a", b"v"]));
2783 r.engine_mut().feed(
2784 writer,
2785 &wire(&[b"BITFIELD", b"nk", b"SET", b"u8", b"0", b"1"]),
2786 );
2787 r.engine_mut().feed(writer, &wire(&[b"LPOP", b"nk"]));
2788 pump(&mut r, &mut batch);
2789 assert_eq!(fired(&r, sub), []);
2790
2791 r.engine_mut()
2793 .feed(writer, &wire(&[b"SET", b"nj", b"v", b"GET"]));
2794 r.engine_mut()
2795 .feed(writer, &wire(&[b"BITFIELD", b"nl", b"GET", b"u8", b"0"]));
2796 r.engine_mut().feed(writer, &wire(&[b"GETDEL", b"nm"]));
2797 pump(&mut r, &mut batch);
2798 assert_eq!(
2799 fired(&r, sub),
2800 [("keymiss", "nj"), ("keymiss", "nl"), ("keymiss", "nm")]
2801 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2802 );
2803 }
2804
2805 #[test]
2808 fn a_store_form_says_it_only_for_its_sources() {
2809 let (mut r, sub, writer, mut batch) = watching_flags(b"Em");
2810
2811 r.engine_mut()
2812 .feed(writer, &wire(&[b"SINTERSTORE", b"dst", b"nk", b"nj"]));
2813 r.engine_mut()
2814 .feed(writer, &wire(&[b"ZUNIONSTORE", b"dst", b"1", b"nz"]));
2815 r.engine_mut()
2816 .feed(writer, &wire(&[b"ZRANGESTORE", b"dst", b"nz", b"0", b"-1"]));
2817 pump(&mut r, &mut batch);
2818 assert_eq!(
2819 fired(&r, sub),
2820 [
2821 ("keymiss", "nk"),
2822 ("keymiss", "nj"),
2823 ("keymiss", "nz"),
2824 ("keymiss", "nz")
2825 ]
2826 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2827 );
2828 }
2829
2830 #[test]
2834 fn the_stream_reads_say_it_for_the_keys_behind_the_keyword() {
2835 let (mut r, sub, writer, mut batch) = watching_flags(b"Em");
2836
2837 r.engine_mut().feed(
2838 writer,
2839 &wire(&[b"XREAD", b"STREAMS", b"nk", b"nj", b"0", b"0"]),
2840 );
2841 pump(&mut r, &mut batch);
2842 assert_eq!(
2843 fired(&r, sub),
2844 [
2845 ("keymiss", "nk"),
2846 ("keymiss", "nj"),
2847 ("keymiss", "nk"),
2848 ("keymiss", "nj")
2849 ]
2850 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2851 );
2852 }
2853
2854 #[test]
2858 fn an_argument_that_did_not_parse_takes_the_miss_back() {
2859 let (mut r, sub, writer, mut batch) = watching_flags(b"Em");
2860
2861 r.engine_mut()
2862 .feed(writer, &wire(&[b"GETRANGE", b"nk", b"x", b"-1"]));
2863 r.engine_mut()
2864 .feed(writer, &wire(&[b"LPOS", b"nk", b"a", b"RANK", b"0"]));
2865 pump(&mut r, &mut batch);
2866 assert_eq!(fired(&r, sub), []);
2867
2868 r.engine_mut().feed(writer, &wire(&[b"SET", b"s", b"v"]));
2872 pump(&mut r, &mut batch);
2873 r.engine_mut().sink_mut().clear();
2874
2875 r.engine_mut()
2876 .feed(writer, &wire(&[b"SINTER", b"nk", b"s"]));
2877 pump(&mut r, &mut batch);
2878 assert_eq!(fired(&r, sub), [("keymiss".to_owned(), "nk".to_owned())]);
2879 }
2880
2881 #[test]
2885 fn a_read_stops_missing_where_it_stops_looking() {
2886 let (mut r, sub, writer, mut batch) = watching_flags(b"Em");
2887
2888 r.engine_mut().feed(writer, &wire(&[b"SET", b"s", b"v"]));
2889 r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"l", b"x"]));
2890 pump(&mut r, &mut batch);
2891 r.engine_mut().sink_mut().clear();
2892
2893 r.engine_mut()
2894 .feed(writer, &wire(&[b"SINTER", b"s", b"nk"]));
2895 pump(&mut r, &mut batch);
2896 assert_eq!(fired(&r, sub), []);
2897
2898 r.engine_mut().feed(writer, &wire(&[b"MGET", b"l", b"nk"]));
2901 pump(&mut r, &mut batch);
2902 assert_eq!(fired(&r, sub), [("keymiss".to_owned(), "nk".to_owned())]);
2903 }
2904
2905 #[test]
2909 fn the_keys_a_sort_pattern_names_say_it_as_they_are_read() {
2910 let (mut r, sub, writer, mut batch) = watching_flags(b"Em");
2911
2912 r.engine_mut()
2913 .feed(writer, &wire(&[b"RPUSH", b"l", b"1", b"2"]));
2914 r.engine_mut().feed(writer, &wire(&[b"SET", b"w_1", b"5"]));
2915 pump(&mut r, &mut batch);
2916 r.engine_mut().sink_mut().clear();
2917
2918 r.engine_mut().feed(
2919 writer,
2920 &wire(&[b"SORT", b"l", b"BY", b"w_*", b"GET", b"p_*"]),
2921 );
2922 pump(&mut r, &mut batch);
2923 assert_eq!(
2924 fired(&r, sub),
2925 [("keymiss", "w_2"), ("keymiss", "p_2"), ("keymiss", "p_1")]
2926 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2927 );
2928 }
2929
2930 #[test]
2934 fn a_deadline_that_passed_is_news_before_the_miss_it_causes() {
2935 let (mut r, sub, mut batch) = timed();
2936 let writer = r.engine_mut().accept();
2937 r.engine_mut().feed(
2938 writer,
2939 &wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", b"EgAm"]),
2940 );
2941 r.engine_mut()
2942 .feed(sub, &wire(&[b"PSUBSCRIBE", b"__keyevent@0__:*"]));
2943 r.engine_mut()
2944 .feed(writer, &wire(&[b"SET", b"k", b"v", b"PX", b"10"]));
2945 pump(&mut r, &mut batch);
2946 r.engine_mut().sink_mut().clear();
2947
2948 r.engine().server().advance_clock_ms(50);
2949 r.engine_mut().feed(writer, &wire(&[b"GET", b"k"]));
2950 pump(&mut r, &mut batch);
2951 assert_eq!(
2952 fired(&r, sub),
2953 [("expired", "k"), ("keymiss", "k")].map(|(e, k)| (e.to_owned(), k.to_owned()))
2954 );
2955 }
2956
2957 #[test]
2960 fn a_module_read_says_it_the_way_a_core_read_does() {
2961 let (mut r, sub, writer, mut batch) = watching_flags(b"Em");
2962
2963 r.engine_mut().feed(writer, &wire(&[b"JSON.GET", b"nk"]));
2964 r.engine_mut().feed(writer, &wire(&[b"TS.GET", b"nj"]));
2965 r.engine_mut()
2966 .feed(writer, &wire(&[b"BF.EXISTS", b"nl", b"x"]));
2967 r.engine_mut().feed(writer, &wire(&[b"TDIGEST.MIN", b"nm"]));
2968 r.engine_mut().feed(writer, &wire(&[b"TOPK.LIST", b"nn"]));
2969 r.engine_mut().feed(writer, &wire(&[b"CMS.INFO", b"no"]));
2970 r.engine_mut().feed(writer, &wire(&[b"VCARD", b"np"]));
2971 pump(&mut r, &mut batch);
2972 assert_eq!(
2973 fired(&r, sub),
2974 [
2975 ("keymiss", "nk"),
2976 ("keymiss", "nj"),
2977 ("keymiss", "nl"),
2978 ("keymiss", "nm"),
2979 ("keymiss", "nn"),
2980 ("keymiss", "no"),
2981 ("keymiss", "np")
2982 ]
2983 .map(|(e, k)| (e.to_owned(), k.to_owned()))
2984 );
2985
2986 r.engine_mut().sink_mut().clear();
2989 r.engine_mut()
2990 .feed(writer, &wire(&[b"JSON.MGET", b"nk", b"nj", b"$"]));
2991 pump(&mut r, &mut batch);
2992 assert_eq!(
2993 fired(&r, sub),
2994 [("keymiss", "nk"), ("keymiss", "nj")].map(|(e, k)| (e.to_owned(), k.to_owned()))
2995 );
2996 }
2997
2998 #[test]
3002 fn the_quiet_module_commands_stay_quiet() {
3003 let (mut r, sub, writer, mut batch) = watching_flags(b"Em");
3004
3005 r.engine_mut()
3006 .feed(writer, &wire(&[b"JSON.SET", b"nk", b"$", b"1"]));
3007 r.engine_mut()
3008 .feed(writer, &wire(&[b"BF.ADD", b"nj", b"x"]));
3009 r.engine_mut()
3010 .feed(writer, &wire(&[b"TS.ADD", b"nl", b"1000", b"1"]));
3011 r.engine_mut().feed(writer, &wire(&[b"TS.INFO", b"nm"]));
3012 r.engine_mut().feed(writer, &wire(&[b"CF.COMPACT", b"nn"]));
3013 r.engine_mut()
3014 .feed(writer, &wire(&[b"JSON.DEBUG", b"HELP"]));
3015 r.engine_mut()
3016 .feed(writer, &wire(&[b"FT.GET", b"ni", b"no"]));
3017 r.engine_mut()
3018 .feed(writer, &wire(&[b"FT.SEARCH", b"ni", b"*"]));
3019 pump(&mut r, &mut batch);
3020 assert_eq!(fired(&r, sub), []);
3021
3022 r.engine_mut()
3024 .feed(writer, &wire(&[b"JSON.DEBUG", b"MEMORY", b"np"]));
3025 r.engine_mut()
3026 .feed(writer, &wire(&[b"FT.SUGGET", b"nq", b"x"]));
3027 r.engine_mut().feed(writer, &wire(&[b"FT.SUGLEN", b"nr"]));
3028 pump(&mut r, &mut batch);
3029 assert_eq!(
3030 fired(&r, sub),
3031 [("keymiss", "np"), ("keymiss", "nq"), ("keymiss", "nr")]
3032 .map(|(e, k)| (e.to_owned(), k.to_owned()))
3033 );
3034 }
3035
3036 #[test]
3043 fn a_module_read_keeps_the_miss_a_later_argument_spoiled() {
3044 let (mut r, sub, writer, mut batch) = watching_flags(b"Em");
3045
3046 r.engine_mut()
3047 .feed(writer, &wire(&[b"JSON.GET", b"nk", b"$..["]));
3048 r.engine_mut().feed(
3049 writer,
3050 &wire(&[b"VSIM", b"nj", b"ELE", b"e", b"COUNT", b"x"]),
3051 );
3052 pump(&mut r, &mut batch);
3053 assert_eq!(
3054 fired(&r, sub),
3055 [("keymiss", "nk"), ("keymiss", "nj")].map(|(e, k)| (e.to_owned(), k.to_owned()))
3056 );
3057 }
3058
3059 #[test]
3062 fn the_module_merges_stop_at_the_first_empty_source() {
3063 let (mut r, sub, writer, mut batch) = watching_flags(b"Em");
3064
3065 r.engine_mut().feed(
3066 writer,
3067 &wire(&[b"TDIGEST.MERGE", b"nk", b"2", b"nj", b"nl"]),
3068 );
3069 pump(&mut r, &mut batch);
3070 assert_eq!(
3071 fired(&r, sub),
3072 [("keymiss", "nk"), ("keymiss", "nj")].map(|(e, k)| (e.to_owned(), k.to_owned()))
3073 );
3074
3075 r.engine_mut().sink_mut().clear();
3079 r.engine_mut()
3080 .feed(writer, &wire(&[b"CMS.MERGE", b"nk", b"1", b"nj"]));
3081 pump(&mut r, &mut batch);
3082 assert_eq!(fired(&r, sub), []);
3083
3084 r.engine_mut()
3087 .feed(writer, &wire(&[b"CMS.INITBYDIM", b"cm", b"100", b"5"]));
3088 pump(&mut r, &mut batch);
3089 r.engine_mut().sink_mut().clear();
3090
3091 r.engine_mut()
3092 .feed(writer, &wire(&[b"CMS.MERGE", b"cm", b"2", b"nj", b"nl"]));
3093 pump(&mut r, &mut batch);
3094 assert_eq!(fired(&r, sub), [("keymiss".to_owned(), "nj".to_owned())]);
3095 }
3096
3097 fn watching_fields(flags: &[u8]) -> (Reactor<Wire<Recorder>>, ConnId, ConnId, Vec<Cmd>) {
3103 let (mut r, sub, mut batch) = engine();
3104 let writer = r.engine_mut().accept();
3105 r.engine_mut().feed(
3106 writer,
3107 &wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", flags]),
3108 );
3109 r.engine_mut()
3110 .feed(sub, &wire(&[b"PSUBSCRIBE", b"__subkey*@0__:*"]));
3111 pump(&mut r, &mut batch);
3112 r.engine_mut().sink_mut().clear();
3113 (r, sub, writer, batch)
3114 }
3115
3116 fn carried(r: &Reactor<Wire<Recorder>>, sub: ConnId) -> Vec<(String, String)> {
3122 let sent = String::from_utf8_lossy(r.engine().sink().sent(sub)).into_owned();
3123 let mut out = Vec::new();
3124 let mut parts = sent.split("\r\n");
3125 while let Some(p) = parts.next() {
3126 if p != "pmessage" {
3127 continue;
3128 }
3129 let mut next = || {
3132 parts.next();
3133 parts.next().unwrap_or_default().to_owned()
3134 };
3135 next();
3136 let channel = next();
3137 out.push((channel, next()));
3138 }
3139 out
3140 }
3141
3142 #[test]
3146 fn the_subkey_channels_carry_the_fields_an_event_touched() {
3147 let (mut r, sub, writer, mut batch) = watching_fields(b"ASTIV");
3148
3149 r.engine_mut()
3150 .feed(writer, &wire(&[b"HSET", b"h", b"a,b", b"1", b"c", b"2"]));
3151 pump(&mut r, &mut batch);
3152 assert_eq!(
3153 carried(&r, sub),
3154 [
3155 ("__subkeyspace@0__:h", "hset|3:a,b,1:c"),
3156 ("__subkeyevent@0__:hset", "1:h|3:a,b,1:c"),
3157 ("__subkeyspaceitem@0__:h\na,b", "hset"),
3158 ("__subkeyspaceitem@0__:h\nc", "hset"),
3159 ("__subkeyspaceevent@0__:hset|h", "3:a,b,1:c"),
3160 ]
3161 .map(|(c, p)| (c.to_owned(), p.to_owned()))
3162 );
3163 }
3164
3165 #[test]
3169 fn a_key_holding_a_newline_skips_the_per_field_channel() {
3170 let (mut r, sub, writer, mut batch) = watching_fields(b"ASTIV");
3171
3172 r.engine_mut()
3173 .feed(writer, &wire(&[b"HSET", b"h\nx", b"f", b"1"]));
3174 pump(&mut r, &mut batch);
3175 assert_eq!(
3176 carried(&r, sub),
3177 [
3178 ("__subkeyspace@0__:h\nx", "hset|1:f"),
3179 ("__subkeyevent@0__:hset", "3:h\nx|1:f"),
3180 ("__subkeyspaceevent@0__:hset|h\nx", "1:f"),
3181 ]
3182 .map(|(c, p)| (c.to_owned(), p.to_owned()))
3183 );
3184 }
3185
3186 #[test]
3190 fn an_event_with_no_fields_stays_off_the_subkey_channels() {
3191 let (mut r, sub, writer, mut batch) = watching_fields(b"AS");
3192
3193 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
3194 r.engine_mut().feed(writer, &wire(&[b"RPUSH", b"l", b"a"]));
3195 r.engine_mut()
3196 .feed(writer, &wire(&[b"HSET", b"h", b"f", b"1"]));
3197 r.engine_mut().feed(writer, &wire(&[b"HDEL", b"h", b"f"]));
3198 pump(&mut r, &mut batch);
3199 assert_eq!(
3200 carried(&r, sub),
3201 [
3202 ("__subkeyspace@0__:h", "hset|1:f"),
3203 ("__subkeyspace@0__:h", "hdel|1:f"),
3204 ]
3205 .map(|(c, p)| (c.to_owned(), p.to_owned()))
3206 );
3207 }
3208
3209 fn watching_clock() -> (Reactor<Wire<Recorder>>, ConnId, ConnId, Vec<Cmd>) {
3215 let (mut r, sub, mut batch) = timed();
3216 let writer = r.engine_mut().accept();
3217 r.engine_mut().feed(
3218 writer,
3219 &wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", b"EA"]),
3220 );
3221 r.engine_mut()
3222 .feed(sub, &wire(&[b"PSUBSCRIBE", b"__keyevent@0__:*"]));
3223 pump(&mut r, &mut batch);
3224 r.engine_mut().sink_mut().clear();
3225 (r, sub, writer, batch)
3226 }
3227
3228 #[test]
3232 fn a_deadline_that_passed_is_news_when_a_reader_finds_it() {
3233 let (mut r, sub, writer, mut batch) = watching_clock();
3234
3235 r.engine_mut()
3236 .feed(writer, &wire(&[b"SET", b"k", b"v", b"PX", b"10"]));
3237 pump(&mut r, &mut batch);
3238 r.engine_mut().sink_mut().clear();
3239
3240 r.engine().server().advance_clock_ms(50);
3241 r.engine_mut().feed(writer, &wire(&[b"GET", b"k"]));
3242 pump(&mut r, &mut batch);
3243 assert_eq!(
3244 fired(&r, sub),
3245 [("expired", "k")].map(|(e, k)| (e.to_owned(), k.to_owned())),
3246 "and not a del alongside it, which is a different piece of news"
3247 );
3248 }
3249
3250 #[test]
3258 fn a_deadline_that_passed_is_news_with_nobody_reading() {
3259 let (mut r, sub, writer, mut batch) = watching_clock();
3260
3261 r.engine_mut()
3262 .feed(writer, &wire(&[b"SET", b"k", b"v", b"PX", b"10"]));
3263 pump(&mut r, &mut batch);
3264 r.engine_mut().sink_mut().clear();
3265
3266 r.engine().server().advance_clock_ms(50);
3267 pump(&mut r, &mut batch);
3269 assert_eq!(
3270 fired(&r, sub),
3271 [("expired", "k")].map(|(e, k)| (e.to_owned(), k.to_owned()))
3272 );
3273
3274 r.engine_mut().sink_mut().clear();
3275 pump(&mut r, &mut batch);
3276 assert!(fired(&r, sub).is_empty(), "and it only goes once");
3277 }
3278
3279 #[test]
3283 fn a_key_a_limit_took_says_it_was_evicted() {
3284 let (mut r, sub, writer, mut batch) = watching();
3285
3286 let val = vec![b'v'; 256];
3287 for i in 0..2000u32 {
3288 let k = format!("key:{i:08}");
3289 r.engine_mut()
3290 .feed(writer, &wire(&[b"SET", k.as_bytes(), &val]));
3291 }
3292 pump(&mut r, &mut batch);
3293 r.engine().server().refresh_memory();
3294 let full = r.engine().server().memory_bytes();
3295 r.engine_mut().sink_mut().clear();
3296
3297 let limit = (full / 2).to_string();
3300 r.engine_mut().feed(
3301 writer,
3302 &wire(&[b"CONFIG", b"SET", b"maxmemory-policy", b"allkeys-random"]),
3303 );
3304 r.engine_mut().feed(
3305 writer,
3306 &wire(&[b"CONFIG", b"SET", b"maxmemory", limit.as_bytes()]),
3307 );
3308 r.engine_mut()
3309 .feed(writer, &wire(&[b"SET", b"newcomer", &val]));
3310 pump(&mut r, &mut batch);
3311
3312 let events = fired(&r, sub);
3313 assert!(
3314 events.iter().any(|(e, _)| e == "evicted"),
3315 "the write made room and never said so: {events:?}"
3316 );
3317 assert!(
3318 events
3319 .iter()
3320 .all(|(e, k)| e != "evicted" || k != "newcomer"),
3321 "the key the write was for is the one key it cannot have taken"
3322 );
3323 }
3324
3325 #[test]
3328 fn a_transaction_publishes_between_its_commands_and_not_after_them() {
3329 let (mut r, sub, mut batch) = engine();
3330 let writer = r.engine_mut().accept();
3331
3332 r.engine_mut().feed(
3333 writer,
3334 &wire(&[b"CONFIG", b"SET", b"notify-keyspace-events", b"EA"]),
3335 );
3336 r.engine_mut()
3337 .feed(sub, &wire(&[b"PSUBSCRIBE", b"__keyevent@0__:*"]));
3338 pump(&mut r, &mut batch);
3339 r.engine_mut().sink_mut().clear();
3340
3341 r.engine_mut().feed(writer, &wire(&[b"MULTI"]));
3342 r.engine_mut().feed(writer, &wire(&[b"SET", b"k", b"v"]));
3343 r.engine_mut().feed(writer, &wire(&[b"DEL", b"k"]));
3344 r.engine_mut().feed(writer, &wire(&[b"EXEC"]));
3345 pump(&mut r, &mut batch);
3346 assert_eq!(
3347 r.engine().sink().sent(sub),
3348 b"*4\r\n$8\r\npmessage\r\n$16\r\n__keyevent@0__:*\r\n\
3349 $18\r\n__keyevent@0__:set\r\n$1\r\nk\r\n\
3350 *4\r\n$8\r\npmessage\r\n$16\r\n__keyevent@0__:*\r\n\
3351 $18\r\n__keyevent@0__:del\r\n$1\r\nk\r\n"
3352 );
3353 }
3354
3355 #[test]
3356 fn a_reply_the_socket_would_not_take_is_offered_again() {
3357 let mut r = Reactor::inline(Wire::new(Trickle::default()));
3358 let conn = r.engine_mut().accept();
3359 let mut batch = Vec::new();
3360
3361 r.engine_mut().feed(conn, &wire(&[b"PING"]));
3362 pump(&mut r, &mut batch);
3363 assert_eq!(r.engine().sink().sent, b"+PONG\r\n");
3365 assert_eq!(r.engine().sink().writes, 2);
3366 }
3367 #[test]
3371 fn client_reply_skip_covers_the_command_after_it_and_nothing_else() {
3372 let (mut r, conn, mut batch) = engine();
3373
3374 r.engine_mut()
3375 .feed(conn, &wire(&[b"CLIENT", b"REPLY", b"SKIP"]));
3376 r.engine_mut().feed(conn, &wire(&[b"SET", b"a", b"1"]));
3377 r.engine_mut().feed(conn, &wire(&[b"SET", b"b", b"2"]));
3378 pump(&mut r, &mut batch);
3379 assert_eq!(r.engine().sink().sent(conn), b"+OK\r\n");
3382 }
3383
3384 #[test]
3385 fn client_reply_off_stays_off_until_it_is_turned_back_on() {
3386 let (mut r, conn, mut batch) = engine();
3387
3388 r.engine_mut()
3389 .feed(conn, &wire(&[b"CLIENT", b"REPLY", b"OFF"]));
3390 r.engine_mut().feed(conn, &wire(&[b"SET", b"c", b"3"]));
3391 r.engine_mut().feed(conn, &wire(&[b"PING"]));
3392 r.engine_mut()
3393 .feed(conn, &wire(&[b"CLIENT", b"REPLY", b"ON"]));
3394 r.engine_mut().feed(conn, &wire(&[b"PING"]));
3395 pump(&mut r, &mut batch);
3396 assert_eq!(r.engine().sink().sent(conn), b"+OK\r\n+PONG\r\n");
3399
3400 r.engine_mut().sink_mut().clear();
3402 r.engine_mut().feed(conn, &wire(&[b"GET", b"c"]));
3403 pump(&mut r, &mut batch);
3404 assert_eq!(r.engine().sink().sent(conn), b"$1\r\n3\r\n");
3405 }
3406
3407 #[test]
3410 fn a_connection_reports_the_socket_it_was_opened_on() {
3411 let mut r = Reactor::inline(Wire::new(Recorder::new()));
3412 let conn = r
3413 .engine_mut()
3414 .accept_from("10.0.0.7:54321", "10.0.0.1:6379", 11, false);
3415 let mut batch = Vec::new();
3416
3417 r.engine_mut().feed(conn, &wire(&[b"CLIENT", b"INFO"]));
3418 pump(&mut r, &mut batch);
3419 let sent = String::from_utf8_lossy(r.engine().sink().sent(conn)).into_owned();
3420 assert!(sent.contains("addr=10.0.0.7:54321"), "{sent}");
3421 assert!(sent.contains("laddr=10.0.0.1:6379"), "{sent}");
3422 assert!(sent.contains("fd=11"), "{sent}");
3423 }
3424
3425 #[test]
3428 fn a_connection_counts_the_bytes_and_the_commands_that_went_over_it() {
3429 let (mut r, conn, mut batch) = engine();
3430 let mut stream = wire(&[b"PING"]);
3431 stream.extend(wire(&[b"PING"]));
3432 let sent_in = stream.len();
3433
3434 r.engine_mut().feed(conn, &stream);
3435 pump(&mut r, &mut batch);
3436 r.engine_mut().sink_mut().clear();
3437
3438 r.engine_mut().feed(conn, &wire(&[b"CLIENT", b"INFO"]));
3439 pump(&mut r, &mut batch);
3440 let sent = String::from_utf8_lossy(r.engine().sink().sent(conn)).into_owned();
3441 assert!(sent.contains("tot-cmds=2"), "{sent}");
3445 assert!(sent.contains("read-events=2"), "{sent}");
3446 assert!(
3447 sent.contains(&format!("tot-net-in={}", sent_in + 26)),
3448 "{sent}"
3449 );
3450 assert!(sent.contains("tot-net-out=14"), "{sent}");
3451 }
3452
3453 #[test]
3457 fn client_list_reports_every_connection_and_not_just_the_one_asking() {
3458 let mut r = Reactor::inline(Wire::new(Recorder::new()));
3459 let one = r
3460 .engine_mut()
3461 .accept_from("10.0.0.7:1111", "10.0.0.1:6379", 11, false);
3462 let two = r
3463 .engine_mut()
3464 .accept_from("10.0.0.8:2222", "10.0.0.1:6379", 12, false);
3465 let mut batch = Vec::new();
3466
3467 r.engine_mut()
3468 .feed(two, &wire(&[b"CLIENT", b"SETNAME", b"worker"]));
3469 r.engine_mut().feed(one, &wire(&[b"CLIENT", b"LIST"]));
3470 pump(&mut r, &mut batch);
3471
3472 let sent = String::from_utf8_lossy(r.engine().sink().sent(one)).into_owned();
3473 let lines: Vec<&str> = sent.lines().filter(|l| l.starts_with("id=")).collect();
3474 assert_eq!(lines.len(), 2, "{sent}");
3475 assert!(lines[0].contains("addr=10.0.0.7:1111"), "{sent}");
3476 assert!(lines[0].contains("cmd=client|list"), "{sent}");
3477 assert!(lines[1].contains("addr=10.0.0.8:2222"), "{sent}");
3478 assert!(lines[1].contains("name=worker"), "{sent}");
3479 assert!(lines[1].contains("cmd=client|setname"), "{sent}");
3480 }
3481
3482 #[test]
3485 fn client_list_keeps_the_order_the_connections_were_opened_in() {
3486 let mut r = Reactor::inline(Wire::new(Recorder::new()));
3487 let one = r.engine_mut().accept();
3488 let two = r.engine_mut().accept();
3489 let three = r.engine_mut().accept();
3490 let mut batch = Vec::new();
3491
3492 r.engine_mut().hangup(two);
3493 r.engine_mut().feed(three, &wire(&[b"CLIENT", b"LIST"]));
3494 pump(&mut r, &mut batch);
3495
3496 let sent = String::from_utf8_lossy(r.engine().sink().sent(three)).into_owned();
3497 let ids: Vec<&str> = sent
3498 .lines()
3499 .filter(|l| l.starts_with("id="))
3500 .map(|l| l.split(' ').next().unwrap_or(""))
3501 .collect();
3502 assert_eq!(ids, vec!["id=1", "id=3"], "{sent}");
3503 let _ = one;
3504 }
3505
3506 #[test]
3510 fn client_kill_closes_the_connection_it_names_and_answers_a_count() {
3511 let mut r = Reactor::inline(Wire::new(Recorder::new()));
3512 let one = r
3513 .engine_mut()
3514 .accept_from("10.0.0.7:1111", "10.0.0.1:6379", 11, false);
3515 let two = r
3516 .engine_mut()
3517 .accept_from("10.0.0.8:2222", "10.0.0.1:6379", 12, false);
3518 let mut batch = Vec::new();
3519
3520 r.engine_mut()
3521 .feed(one, &wire(&[b"CLIENT", b"KILL", b"ADDR", b"10.0.0.8:2222"]));
3522 pump(&mut r, &mut batch);
3523
3524 assert_eq!(r.engine().sink().sent(one), b":1\r\n");
3525 assert!(r.engine().sink().was_closed(two), "the named one went away");
3526 assert!(!r.engine().sink().was_closed(one), "the caller did not");
3527 }
3528
3529 #[test]
3532 fn the_old_kill_takes_one_address_and_will_take_the_caller() {
3533 let mut r = Reactor::inline(Wire::new(Recorder::new()));
3534 let one = r
3535 .engine_mut()
3536 .accept_from("10.0.0.7:1111", "10.0.0.1:6379", 11, false);
3537 let mut batch = Vec::new();
3538
3539 r.engine_mut()
3540 .feed(one, &wire(&[b"CLIENT", b"KILL", b"10.0.0.9:9999"]));
3541 pump(&mut r, &mut batch);
3542 assert_eq!(r.engine().sink().sent(one), b"-ERR No such client\r\n");
3543
3544 r.engine_mut().sink_mut().clear();
3545 r.engine_mut()
3546 .feed(one, &wire(&[b"CLIENT", b"KILL", b"10.0.0.7:1111"]));
3547 pump(&mut r, &mut batch);
3548 assert_eq!(r.engine().sink().sent(one), b"+OK\r\n");
3549 assert!(r.engine().sink().was_closed(one), "it took itself");
3550 }
3551
3552 #[test]
3555 fn a_kill_spares_the_caller_unless_it_is_told_not_to() {
3556 let (mut r, conn, mut batch) = engine();
3557 r.engine_mut()
3558 .feed(conn, &wire(&[b"CLIENT", b"KILL", b"TYPE", b"normal"]));
3559 pump(&mut r, &mut batch);
3560 assert_eq!(r.engine().sink().sent(conn), b":0\r\n");
3561 assert!(!r.engine().sink().was_closed(conn));
3562
3563 r.engine_mut().sink_mut().clear();
3564 r.engine_mut().feed(
3565 conn,
3566 &wire(&[b"CLIENT", b"KILL", b"TYPE", b"normal", b"SKIPME", b"no"]),
3567 );
3568 pump(&mut r, &mut batch);
3569 assert_eq!(r.engine().sink().sent(conn), b":1\r\n");
3570 assert!(r.engine().sink().was_closed(conn));
3571 }
3572
3573 #[test]
3576 fn a_connection_that_went_away_is_off_the_list() {
3577 let mut r = Reactor::inline(Wire::new(Recorder::new()));
3578 let one = r.engine_mut().accept();
3579 let two = r.engine_mut().accept();
3580 let mut batch = Vec::new();
3581 assert_eq!(r.engine().server().client_count(), 2);
3582
3583 r.engine_mut().hangup(two);
3584 assert_eq!(r.engine().server().client_count(), 1);
3585
3586 r.engine_mut().feed(one, &wire(&[b"CLIENT", b"LIST"]));
3587 pump(&mut r, &mut batch);
3588 let sent = String::from_utf8_lossy(r.engine().sink().sent(one)).into_owned();
3589 assert_eq!(sent.lines().filter(|l| l.starts_with("id=")).count(), 1);
3590 }
3591
3592 #[test]
3595 fn a_write_pause_holds_the_writes_and_lets_the_reads_through() {
3596 let (mut r, one, mut batch) = timed();
3597 let two = r.engine_mut().accept();
3598
3599 r.engine_mut()
3600 .feed(one, &wire(&[b"CLIENT", b"PAUSE", b"500", b"WRITE"]));
3601 pump(&mut r, &mut batch);
3602 assert_eq!(r.engine().sink().sent(one), b"+OK\r\n");
3603
3604 r.engine_mut().sink_mut().clear();
3605 r.engine_mut().feed(two, &wire(&[b"GET", b"k"]));
3606 r.engine_mut().feed(two, &wire(&[b"SET", b"k", b"v"]));
3607 pump(&mut r, &mut batch);
3608 assert_eq!(
3609 r.engine().sink().sent(two),
3610 b"$-1\r\n",
3611 "the read answered and the write is being held"
3612 );
3613 assert_eq!(r.engine().waiting(), 1, "the held connection");
3614
3615 r.engine_mut().sink_mut().clear();
3616 r.engine().server().advance_clock_ms(500);
3617 pump(&mut r, &mut batch);
3618 assert_eq!(r.engine().sink().sent(two), b"+OK\r\n");
3619 assert_eq!(r.engine().waiting(), 0);
3620 }
3621
3622 #[test]
3624 fn an_all_pause_holds_every_command_and_cannot_be_called_off() {
3625 let (mut r, one, mut batch) = timed();
3626 let two = r.engine_mut().accept();
3627
3628 r.engine_mut()
3629 .feed(one, &wire(&[b"CLIENT", b"PAUSE", b"500", b"ALL"]));
3630 pump(&mut r, &mut batch);
3631
3632 r.engine_mut().sink_mut().clear();
3633 r.engine_mut().feed(two, &wire(&[b"PING"]));
3634 r.engine_mut().feed(two, &wire(&[b"CLIENT", b"UNPAUSE"]));
3635 pump(&mut r, &mut batch);
3636 assert_eq!(r.engine().sink().sent(two), b"", "not even the ping");
3637
3638 r.engine().server().advance_clock_ms(499);
3639 pump(&mut r, &mut batch);
3640 assert_eq!(r.engine().sink().sent(two), b"", "still inside the pause");
3641
3642 r.engine().server().advance_clock_ms(1);
3643 pump(&mut r, &mut batch);
3644 assert_eq!(
3645 r.engine().sink().sent(two),
3646 b"+PONG\r\n+OK\r\n",
3647 "both, in the order they were sent"
3648 );
3649 }
3650
3651 #[test]
3653 fn a_shorter_pause_does_not_shorten_the_one_already_running() {
3654 let (mut r, one, mut batch) = timed();
3655 let two = r.engine_mut().accept();
3656
3657 r.engine_mut()
3658 .feed(one, &wire(&[b"CLIENT", b"PAUSE", b"500", b"WRITE"]));
3659 r.engine_mut()
3660 .feed(one, &wire(&[b"CLIENT", b"PAUSE", b"10", b"WRITE"]));
3661 pump(&mut r, &mut batch);
3662 assert_eq!(r.engine().server().pause_ends(), START_MS + 500);
3663
3664 r.engine_mut().sink_mut().clear();
3665 r.engine_mut().feed(two, &wire(&[b"SET", b"k", b"v"]));
3666 r.engine().server().advance_clock_ms(100);
3667 pump(&mut r, &mut batch);
3668 assert_eq!(r.engine().sink().sent(two), b"", "the longer end holds");
3669
3670 r.engine().server().advance_clock_ms(400);
3671 pump(&mut r, &mut batch);
3672 assert_eq!(r.engine().sink().sent(two), b"+OK\r\n");
3673 }
3674
3675 #[test]
3677 fn a_write_pause_on_top_of_an_all_pause_still_holds_the_reads() {
3678 let (mut r, one, mut batch) = timed();
3679 let two = r.engine_mut().accept();
3680
3681 r.engine_mut()
3682 .feed(one, &wire(&[b"CLIENT", b"PAUSE", b"500", b"ALL"]));
3683 pump(&mut r, &mut batch);
3684 r.engine_mut()
3687 .feed(one, &wire(&[b"CLIENT", b"PAUSE", b"500", b"WRITE"]));
3688
3689 r.engine_mut().sink_mut().clear();
3690 r.engine_mut().feed(two, &wire(&[b"GET", b"k"]));
3691 r.engine().server().advance_clock_ms(200);
3692 pump(&mut r, &mut batch);
3693 assert_eq!(r.engine().sink().sent(two), b"", "still everything");
3694 }
3695
3696 #[test]
3698 fn unpause_lets_go_of_a_write_pause_at_once() {
3699 let (mut r, one, mut batch) = timed();
3700 let two = r.engine_mut().accept();
3701
3702 r.engine_mut()
3703 .feed(one, &wire(&[b"CLIENT", b"PAUSE", b"5000", b"WRITE"]));
3704 pump(&mut r, &mut batch);
3705 r.engine_mut().feed(two, &wire(&[b"SET", b"k", b"v"]));
3706 pump(&mut r, &mut batch);
3707 assert_eq!(r.engine().sink().sent(two), b"");
3708
3709 r.engine_mut().sink_mut().clear();
3710 r.engine_mut().feed(one, &wire(&[b"CLIENT", b"UNPAUSE"]));
3711 pump(&mut r, &mut batch);
3712 assert_eq!(r.engine().sink().sent(two), b"+OK\r\n");
3713 }
3714
3715 #[test]
3718 fn a_write_pause_holds_an_exec_only_when_the_transaction_writes() {
3719 let (mut r, one, mut batch) = timed();
3720 let two = r.engine_mut().accept();
3721 let three = r.engine_mut().accept();
3722
3723 for conn in [two, three] {
3724 r.engine_mut().feed(conn, &wire(&[b"MULTI"]));
3725 r.engine_mut().feed(conn, &wire(&[b"GET", b"k"]));
3726 }
3727 r.engine_mut().feed(three, &wire(&[b"SET", b"k", b"v"]));
3728 pump(&mut r, &mut batch);
3729
3730 r.engine_mut()
3731 .feed(one, &wire(&[b"CLIENT", b"PAUSE", b"500", b"WRITE"]));
3732 pump(&mut r, &mut batch);
3733
3734 r.engine_mut().sink_mut().clear();
3735 r.engine_mut().feed(two, &wire(&[b"EXEC"]));
3736 r.engine_mut().feed(three, &wire(&[b"EXEC"]));
3737 pump(&mut r, &mut batch);
3738 assert_eq!(r.engine().sink().sent(two), b"*1\r\n$-1\r\n", "reads only");
3739 assert_eq!(r.engine().sink().sent(three), b"", "one write in it");
3740
3741 r.engine_mut().sink_mut().clear();
3742 r.engine().server().advance_clock_ms(500);
3743 pump(&mut r, &mut batch);
3744 assert_eq!(r.engine().sink().sent(three), b"*2\r\n$-1\r\n+OK\r\n");
3745 }
3746
3747 #[test]
3750 fn a_held_connection_that_hangs_up_lets_go_of_its_slot() {
3751 let (mut r, one, mut batch) = timed();
3752 let two = r.engine_mut().accept();
3753
3754 r.engine_mut()
3755 .feed(one, &wire(&[b"CLIENT", b"PAUSE", b"500", b"ALL"]));
3756 pump(&mut r, &mut batch);
3757 r.engine_mut().feed(two, &wire(&[b"PING"]));
3758 pump(&mut r, &mut batch);
3759 assert_eq!(r.engine().waiting(), 1);
3760
3761 r.engine_mut().hangup(two);
3762 pump(&mut r, &mut batch);
3763 assert_eq!(r.engine().server().client_count(), 1);
3764
3765 r.engine().server().advance_clock_ms(500);
3766 pump(&mut r, &mut batch);
3767 assert_eq!(r.engine().waiting(), 0, "nothing left holding a slot");
3768 }
3769
3770 #[test]
3773 fn a_held_connection_names_the_command_it_is_waiting_to_run() {
3774 let (mut r, one, mut batch) = timed();
3775 let two = r.engine_mut().accept();
3776
3777 r.engine_mut()
3778 .feed(one, &wire(&[b"CLIENT", b"PAUSE", b"500", b"WRITE"]));
3779 pump(&mut r, &mut batch);
3780 r.engine_mut().feed(two, &wire(&[b"SET", b"k", b"v"]));
3781 pump(&mut r, &mut batch);
3782
3783 r.engine_mut().sink_mut().clear();
3784 r.engine().server().advance_clock_ms(0);
3785 r.engine_mut().feed(one, &wire(&[b"CLIENT", b"LIST"]));
3786 r.engine().server().advance_clock_ms(500);
3790 pump(&mut r, &mut batch);
3791 let sent = String::from_utf8_lossy(r.engine().sink().sent(one)).into_owned();
3792 assert!(sent.contains("cmd=set"), "{sent}");
3793 }
3794
3795 fn watched() -> (Reactor<Wire<Recorder>>, ConnId, ConnId, Vec<Cmd>) {
3800 let (mut r, eye, mut batch) = timed();
3801 let one = r.engine_mut().accept();
3802 r.engine_mut().feed(eye, &wire(&[b"MONITOR"]));
3803 pump(&mut r, &mut batch);
3804 assert_eq!(r.engine().sink().sent(eye), b"+OK\r\n");
3805 r.engine_mut().sink_mut().clear();
3806 (r, eye, one, batch)
3807 }
3808
3809 const STAMP: &str = "1000.000000";
3811
3812 fn fed(r: &Reactor<Wire<Recorder>>, eye: ConnId) -> String {
3814 String::from_utf8_lossy(r.engine().sink().sent(eye)).into_owned()
3815 }
3816
3817 #[test]
3818 fn a_monitor_is_fed_a_command_another_connection_ran() {
3819 let (mut r, eye, one, mut batch) = watched();
3820
3821 r.engine_mut().feed(one, &wire(&[b"SET", b"k", b"v"]));
3822 pump(&mut r, &mut batch);
3823
3824 assert_eq!(
3825 fed(&r, eye),
3826 format!("+{STAMP} [0 ?:0] \"SET\" \"k\" \"v\"\r\n")
3827 );
3828 }
3829
3830 #[test]
3833 fn a_select_is_reported_on_the_database_it_moved_to() {
3834 let (mut r, eye, one, mut batch) = watched();
3835
3836 r.engine_mut().feed(one, &wire(&[b"SELECT", b"3"]));
3837 r.engine_mut().feed(one, &wire(&[b"GET", b"k"]));
3838 r.engine_mut().feed(one, &wire(&[b"SELECT", b"0"]));
3839 pump(&mut r, &mut batch);
3840
3841 assert_eq!(
3842 fed(&r, eye),
3843 format!(
3844 "+{STAMP} [3 ?:0] \"SELECT\" \"3\"\r\n\
3845 +{STAMP} [3 ?:0] \"GET\" \"k\"\r\n\
3846 +{STAMP} [0 ?:0] \"SELECT\" \"0\"\r\n"
3847 )
3848 );
3849 }
3850
3851 #[test]
3854 fn an_argument_is_quoted_the_way_a_real_server_quotes_it() {
3855 let (mut r, eye, one, mut batch) = watched();
3856
3857 r.engine_mut()
3858 .feed(one, &wire(&[b"SET", b"a b", b"q\"s\\n\nt\tz\x01\xff"]));
3859 pump(&mut r, &mut batch);
3860
3861 assert_eq!(
3862 fed(&r, eye),
3863 format!("+{STAMP} [0 ?:0] \"SET\" \"a b\" \"q\\\"s\\\\n\\nt\\tz\\x01\\xff\"\r\n")
3864 );
3865 }
3866
3867 #[test]
3870 fn only_a_command_that_ran_is_reported() {
3871 let (mut r, eye, one, mut batch) = watched();
3872
3873 r.engine_mut().feed(one, &wire(&[b"NOSUCHCOMMAND", b"a"]));
3874 r.engine_mut().feed(one, &wire(&[b"GET"]));
3875 r.engine_mut().feed(one, &wire(&[b"LPUSH", b"l", b"x"]));
3876 r.engine_mut().feed(one, &wire(&[b"GET", b"l"]));
3877 pump(&mut r, &mut batch);
3878
3879 assert_eq!(
3880 fed(&r, eye),
3881 format!(
3882 "+{STAMP} [0 ?:0] \"LPUSH\" \"l\" \"x\"\r\n\
3883 +{STAMP} [0 ?:0] \"GET\" \"l\"\r\n"
3884 ),
3885 "the unknown command and the wrong arity are not commands that ran"
3886 );
3887 }
3888
3889 #[test]
3893 fn a_transaction_is_reported_at_exec_and_in_order() {
3894 let (mut r, eye, one, mut batch) = watched();
3895
3896 r.engine_mut().feed(one, &wire(&[b"MULTI"]));
3897 r.engine_mut().feed(one, &wire(&[b"SET", b"t", b"1"]));
3898 r.engine_mut().feed(one, &wire(&[b"INCR", b"t"]));
3899 r.engine_mut().feed(one, &wire(&[b"EXEC"]));
3900 pump(&mut r, &mut batch);
3901
3902 assert_eq!(
3903 fed(&r, eye),
3904 format!(
3905 "+{STAMP} [0 ?:0] \"MULTI\"\r\n\
3906 +{STAMP} [0 ?:0] \"SET\" \"t\" \"1\"\r\n\
3907 +{STAMP} [0 ?:0] \"INCR\" \"t\"\r\n\
3908 +{STAMP} [0 ?:0] \"EXEC\"\r\n"
3909 )
3910 );
3911 }
3912
3913 #[test]
3916 #[cfg_attr(miri, ignore = "a Lua state is C, and Miri interprets Rust")]
3917 fn a_script_is_reported_in_front_of_its_own_effects() {
3918 let (mut r, eye, one, mut batch) = watched();
3919
3920 r.engine_mut().feed(
3921 one,
3922 &wire(&[
3923 b"EVAL",
3924 b"redis.call('set', KEYS[1], 'z') return 1",
3925 b"1",
3926 b"sk",
3927 ]),
3928 );
3929 pump(&mut r, &mut batch);
3930
3931 assert_eq!(
3932 fed(&r, eye),
3933 format!(
3934 "+{STAMP} [0 ?:0] \"EVAL\" \"redis.call('set', KEYS[1], 'z') return 1\" \"1\" \"sk\"\r\n\
3935 +{STAMP} [0 lua] \"set\" \"sk\" \"z\"\r\n"
3936 )
3937 );
3938 }
3939
3940 #[test]
3943 fn an_administrative_command_is_not_reported() {
3944 let (mut r, eye, one, mut batch) = watched();
3945
3946 r.engine_mut().feed(one, &wire(&[b"CLIENT", b"LIST"]));
3947 r.engine_mut()
3948 .feed(one, &wire(&[b"CONFIG", b"GET", b"maxmemory"]));
3949 r.engine_mut().feed(one, &wire(&[b"CLIENT", b"ID"]));
3950 r.engine_mut().feed(one, &wire(&[b"CONFIG", b"HELP"]));
3951 pump(&mut r, &mut batch);
3952
3953 assert_eq!(
3954 fed(&r, eye),
3955 format!(
3956 "+{STAMP} [0 ?:0] \"CLIENT\" \"ID\"\r\n\
3957 +{STAMP} [0 ?:0] \"CONFIG\" \"HELP\"\r\n"
3958 )
3959 );
3960 }
3961
3962 #[test]
3965 fn a_password_is_not_echoed_to_a_monitor() {
3966 let (mut r, eye, one, mut batch) = watched();
3967
3968 r.engine_mut().feed(
3969 one,
3970 &wire(&[b"HELLO", b"3", b"AUTH", b"default", b"hunter2"]),
3971 );
3972 pump(&mut r, &mut batch);
3973
3974 let sent = fed(&r, eye);
3975 assert!(
3976 sent.contains("\"HELLO\" \"3\" \"AUTH\" \"(redacted)\" \"(redacted)\""),
3977 "{sent}"
3978 );
3979 assert!(!sent.contains("hunter2"), "{sent}");
3980 }
3981
3982 #[test]
3985 fn a_monitor_may_not_touch_the_keyspace() {
3986 let (mut r, eye, _one, mut batch) = watched();
3987
3988 r.engine_mut().feed(eye, &wire(&[b"PING"]));
3993 pump(&mut r, &mut batch);
3994 r.engine_mut().feed(eye, &wire(&[b"GET", b"k"]));
3995 pump(&mut r, &mut batch);
3996
3997 assert_eq!(
3998 fed(&r, eye),
3999 format!(
4000 "+PONG\r\n\
4001 +{STAMP} [0 ?:0] \"PING\"\r\n\
4002 -ERR Replica can't interact with the keyspace\r\n"
4003 )
4004 );
4005 }
4006
4007 #[test]
4010 fn a_monitor_is_refused_the_keyspace_at_queue_time() {
4011 let (mut r, eye, _one, mut batch) = watched();
4012
4013 r.engine_mut().feed(eye, &wire(&[b"MULTI"]));
4014 r.engine_mut().feed(eye, &wire(&[b"GET", b"k"]));
4015 r.engine_mut().feed(eye, &wire(&[b"EXEC"]));
4016 pump(&mut r, &mut batch);
4017
4018 let sent = fed(&r, eye);
4019 assert!(
4020 sent.contains("-ERR Replica can't interact with the keyspace"),
4021 "{sent}"
4022 );
4023 assert!(sent.contains("-EXECABORT"), "{sent}");
4024 }
4025
4026 #[test]
4029 fn a_monitor_runs_through_a_pause() {
4030 let (mut r, eye, one, mut batch) = watched();
4031
4032 r.engine_mut()
4033 .feed(one, &wire(&[b"CLIENT", b"PAUSE", b"5000", b"ALL"]));
4034 pump(&mut r, &mut batch);
4035 r.engine_mut().sink_mut().clear();
4036
4037 r.engine_mut().feed(eye, &wire(&[b"PING"]));
4038 r.engine_mut().feed(one, &wire(&[b"PING"]));
4039 pump(&mut r, &mut batch);
4040
4041 assert!(
4042 fed(&r, eye).starts_with("+PONG\r\n"),
4043 "the monitor is let through"
4044 );
4045 assert_eq!(r.engine().sink().sent(one), b"", "and the client is not");
4046 }
4047
4048 #[test]
4051 fn monitor_sent_twice_answers_nothing_the_second_time() {
4052 let (mut r, eye, _one, mut batch) = watched();
4053
4054 r.engine_mut().feed(eye, &wire(&[b"MONITOR"]));
4055 pump(&mut r, &mut batch);
4056
4057 assert_eq!(fed(&r, eye), "", "no reply and no line either");
4058 }
4059
4060 #[test]
4062 fn reset_takes_a_connection_out_of_monitor_mode() {
4063 let (mut r, eye, one, mut batch) = watched();
4064
4065 r.engine_mut().feed(eye, &wire(&[b"RESET"]));
4066 pump(&mut r, &mut batch);
4067 r.engine_mut().sink_mut().clear();
4068
4069 r.engine_mut().feed(one, &wire(&[b"SET", b"k", b"v"]));
4070 r.engine_mut().feed(eye, &wire(&[b"GET", b"k"]));
4071 pump(&mut r, &mut batch);
4072
4073 assert_eq!(fed(&r, eye), "$1\r\nv\r\n", "not watching and not refused");
4074 }
4075
4076 #[test]
4079 fn a_monitor_that_closes_stops_being_one() {
4080 let (mut r, eye, one, mut batch) = watched();
4081
4082 r.engine_mut().hangup(eye);
4083 pump(&mut r, &mut batch);
4084 assert!(!r.engine().server().monitored());
4085
4086 r.engine_mut().feed(one, &wire(&[b"SET", b"k", b"v"]));
4087 pump(&mut r, &mut batch);
4088 assert_eq!(r.engine().sink().sent(one), b"+OK\r\n");
4089 }
4090
4091 #[test]
4093 fn two_monitors_are_fed_the_same_bytes() {
4094 let (mut r, eye, one, mut batch) = watched();
4095 let other = r.engine_mut().accept();
4096 r.engine_mut().feed(other, &wire(&[b"MONITOR"]));
4097 pump(&mut r, &mut batch);
4098 r.engine_mut().sink_mut().clear();
4099
4100 r.engine_mut().feed(one, &wire(&[b"SET", b"k", b"v"]));
4101 pump(&mut r, &mut batch);
4102
4103 assert_eq!(fed(&r, eye), fed(&r, other));
4104 assert_eq!(
4105 fed(&r, eye),
4106 format!("+{STAMP} [0 ?:0] \"SET\" \"k\" \"v\"\r\n")
4107 );
4108 }
4109
4110 #[test]
4112 fn a_monitor_is_reported_as_one_in_client_list() {
4113 let (mut r, _eye, one, mut batch) = watched();
4114
4115 r.engine_mut().feed(one, &wire(&[b"CLIENT", b"LIST"]));
4116 pump(&mut r, &mut batch);
4117
4118 let sent = String::from_utf8_lossy(r.engine().sink().sent(one)).into_owned();
4119 assert!(sent.contains("flags=O"), "{sent}");
4120 assert!(sent.contains("flags=N"), "and the other one is not: {sent}");
4121 }
4122}