use kevy_resp::{Argv, ArgvView};
use kevy_rt::{BlockKind, Store};
pub(crate) fn block_serve_argv<A: ArgvView + ?Sized>(
args: &A,
kind: BlockKind,
key: &[u8],
) -> Argv {
match kind {
BlockKind::Blpop => pop_serve(b"BLPOP", key),
BlockKind::Brpop => pop_serve(b"BRPOP", key),
BlockKind::Bzpopmin => pop_serve(b"BZPOPMIN", key),
BlockKind::Brpoplpush => brpoplpush_serve(args, key),
BlockKind::XReadBlock => xread_serve(args, key).unwrap_or_else(|| args.to_argv()),
BlockKind::XReadGroupBlock => {
xreadgroup_serve(args, key).unwrap_or_else(|| args.to_argv())
}
}
}
fn brpoplpush_serve<A: ArgvView + ?Sized>(args: &A, key: &[u8]) -> Argv {
let mut a = Argv::default();
a.push(b"BRPOPLPUSH");
a.push(key);
if let Some(dst) = args.get(2) {
a.push(dst);
} else {
a.push(b"");
}
a.push(b"0");
a
}
fn pop_serve(verb: &[u8], key: &[u8]) -> Argv {
let mut a = Argv::default();
a.push(verb);
a.push(key);
a.push(b"0");
a
}
#[derive(Default)]
struct StreamOpts {
count: Option<Vec<u8>>,
block_ms: Option<Vec<u8>>,
noack: bool,
streams_at: usize,
}
fn scan_stream_opts<A: ArgvView + ?Sized>(args: &A, from: usize) -> Option<StreamOpts> {
let mut o = StreamOpts::default();
let mut i = from;
loop {
match args.get(i)?.to_ascii_uppercase().as_slice() {
b"COUNT" => {
o.count = Some(args.get(i + 1)?.to_vec());
i += 2;
}
b"BLOCK" => {
o.block_ms = Some(args.get(i + 1)?.to_vec());
i += 2;
}
b"NOACK" => {
o.noack = true;
i += 1;
}
b"STREAMS" => {
o.streams_at = i;
return Some(o);
}
_ => return None,
}
}
}
fn push_stream_tail(a: &mut Argv, o: &StreamOpts, key: &[u8], id: &[u8]) {
if let Some(c) = &o.count {
a.push(b"COUNT");
a.push(c);
}
if o.noack {
a.push(b"NOACK");
}
if let Some(b) = &o.block_ms {
a.push(b"BLOCK");
a.push(b);
}
a.push(b"STREAMS");
a.push(key);
a.push(id);
}
fn xread_serve<A: ArgvView + ?Sized>(args: &A, key: &[u8]) -> Option<Argv> {
let o = scan_stream_opts(args, 1)?;
let id = id_for_key(args, o.streams_at + 1, key)?;
let mut a = Argv::default();
a.push(b"XREAD");
push_stream_tail(&mut a, &o, key, &id);
Some(a)
}
fn xreadgroup_serve<A: ArgvView + ?Sized>(args: &A, key: &[u8]) -> Option<Argv> {
if args.len() < 4 || !args[1].eq_ignore_ascii_case(b"GROUP") {
return None;
}
let group = args[2].to_vec();
let consumer = args[3].to_vec();
let o = scan_stream_opts(args, 4)?;
let id = id_for_key(args, o.streams_at + 1, key)?;
let mut a = Argv::default();
a.push(b"XREADGROUP");
a.push(b"GROUP");
a.push(&group);
a.push(&consumer);
push_stream_tail(&mut a, &o, key, &id);
Some(a)
}
fn id_for_key<A: ArgvView + ?Sized>(args: &A, keys_start: usize, key: &[u8]) -> Option<Vec<u8>> {
let rest = args.len().checked_sub(keys_start)?;
if rest == 0 || !rest.is_multiple_of(2) {
return None;
}
let n = rest / 2;
let pos = (keys_start..keys_start + n).position(|i| &args[i] == key)?;
args.get(keys_start + n + pos).map(<[u8]>::to_vec)
}
pub(crate) fn block_ready<A: ArgvView + ?Sized>(
store: &mut Store,
serve_argv: &A,
kind: BlockKind,
) -> bool {
match kind {
BlockKind::Blpop | BlockKind::Brpop | BlockKind::Brpoplpush => serve_argv
.get(1)
.is_some_and(|k| store.llen(k).is_ok_and(|n| n > 0)),
BlockKind::Bzpopmin => serve_argv
.get(1)
.is_some_and(|k| store.zcard(k).is_ok_and(|n| n > 0)),
BlockKind::XReadBlock => {
let mut tmp = Vec::new();
crate::dispatch::dispatch_into(store, serve_argv, &mut tmp);
!tmp.is_empty()
}
BlockKind::XReadGroupBlock => xreadgroup_ready(store, serve_argv),
}
}
fn xreadgroup_ready<A: ArgvView + ?Sized>(store: &mut Store, serve_argv: &A) -> bool {
if serve_argv.len() < 3 || !serve_argv[1].eq_ignore_ascii_case(b"GROUP") {
return false;
}
let group = serve_argv[2].to_vec();
let mut i = 4usize;
while i < serve_argv.len() {
if serve_argv[i].eq_ignore_ascii_case(b"STREAMS") {
let Some(key) = serve_argv.get(i + 1) else {
return false;
};
return store.xreadgroup_has_new(key, &group).unwrap_or(false);
}
i += 1;
}
false
}