use crate::BlockKind;
use kevy_resp::{Argv, RespVersion};
pub use crate::message_kinds::{MultiOp, ZCombine};
pub(crate) use crate::message_part::Part;
pub(crate) use crate::message_kinds::{DispatchMeta, GatherKind, Gathered};
use std::collections::HashMap;
use std::sync::{Arc, RwLock};
pub(crate) type KvPairs = Vec<(Vec<u8>, Vec<u8>)>;
pub(crate) type PubSubReg = Arc<RwLock<HashMap<Vec<u8>, (u32, u64)>>>;
pub(crate) type PubSubPatternReg = Arc<RwLock<Vec<(Vec<u8>, u32, u64)>>>;
pub(crate) type PubMsg = Arc<(Vec<u8>, Vec<u8>)>;
pub(crate) enum Op {
Del(Vec<Vec<u8>>),
Exists(Vec<Vec<u8>>),
Dbsize,
Flush,
Save,
BgSave,
RewriteAof,
MSet(KvPairs),
Gather(GatherKind, Vec<Vec<u8>>),
ZStoreResult { dst: Vec<u8>, pairs: Vec<(Vec<u8>, f64)> },
SetStoreResult { dst: Vec<u8>, members: Vec<Vec<u8>> },
FeedRead {
cursor_gen: u64,
offset: u64,
count: usize,
prefixes: Vec<Vec<u8>>,
},
FeedTail,
Extension { argv: std::sync::Arc<[Vec<u8>]> },
GeoSearch { argv: Vec<Vec<u8>> },
ReplToken,
PrefixStats(Vec<u8>),
ClientList,
ClientKill(crate::client_ops::ClientKillFilter),
CollectKeys(Option<Vec<u8>>, Option<usize>),
RandomKey,
ScanStep {
cursor: u64,
count: usize,
pattern: Option<Vec<u8>>,
type_filter: Option<Vec<u8>>,
},
CollectWatchVersions(Vec<Vec<u8>>),
CheckWatch(Vec<(Vec<u8>, u64)>),
Rename {
src: Vec<u8>,
dst: Vec<u8>,
nx: bool,
},
RenameTake(Vec<u8>),
RenamePut {
dst: Vec<u8>,
value: kevy_store::Value,
ttl_ms: Option<u64>,
nx: bool,
},
BitOpResult { key: Vec<u8>, value: Vec<u8> },
Copy { src: Vec<u8>, dst: Vec<u8>, replace: bool },
CopyRead(Vec<u8>),
CopyPut {
dst: Vec<u8>,
value: kevy_store::Value,
ttl_ms: Option<u64>,
replace: bool,
},
ListMove {
src: Vec<u8>,
dst: Vec<u8>,
from_left: bool,
to_left: bool,
},
ListMoveTake { key: Vec<u8>, from_left: bool },
ListMovePush {
key: Vec<u8>,
value: Vec<u8>,
to_left: bool,
},
ListMoveRestore {
key: Vec<u8>,
value: Vec<u8>,
from_left: bool,
},
SlowlogGet,
SlowlogLen,
SlowlogReset,
XReadOne { index: u32, argv: Argv, write: bool },
}
pub(crate) enum SmallReply {
Inline { len: u8, buf: [u8; 30] },
Heap(Vec<u8>),
}
impl SmallReply {
#[inline]
pub(crate) fn from_slice(b: &[u8]) -> Self {
if b.len() <= 30 {
let mut buf = [0u8; 30];
buf[..b.len()].copy_from_slice(b);
SmallReply::Inline { len: b.len() as u8, buf }
} else {
SmallReply::Heap(b.to_vec())
}
}
#[inline]
pub(crate) fn from_vec(v: Vec<u8>) -> Self {
SmallReply::Heap(v)
}
#[inline]
pub(crate) fn as_slice(&self) -> &[u8] {
match self {
SmallReply::Inline { len, buf } => &buf[..*len as usize],
SmallReply::Heap(v) => v,
}
}
}
pub(crate) type ReqBatch = Vec<(u64, u64, Argv, RespVersion, DispatchMeta)>;
pub(crate) type RespBatch = Vec<(u64, u64, Part, Argv)>;
pub(crate) enum Inbound {
Request {
origin: usize,
conn: u64,
seq: u64,
op: Op,
},
Response {
conn: u64,
seq: u64,
part: Part,
},
RequestBatch {
origin: usize,
reqs: ReqBatch,
},
ResponseBatch(RespBatch),
DeliverPublish(Vec<PubMsg>),
BlockArm {
origin: usize,
conn: u64,
key: Vec<u8>,
kind: BlockKind,
serve_argv: Argv,
proto: RespVersion,
},
BlockReady { conn: u64, key: Vec<u8> },
BlockServeReq {
origin: usize,
conn: u64,
key: Vec<u8>,
},
BlockServeResp {
conn: u64,
key: Vec<u8>,
reply: Vec<u8>,
},
BlockServeAck { origin: usize, conn: u64 },
RenameCommitted { src: Vec<u8> },
BlockServeAbort { origin: usize, conn: u64 },
BlockCancel { origin: usize, conn: u64 },
ReplWaitArm {
origin: usize,
conn: u64,
seq: u64,
need: u32,
deadline_ms: u64,
},
ReplApplyArm {
origin: usize,
conn: u64,
seq: u64,
min_offset: u64,
deadline_ms: u64,
},
ReplDone { conn: u64, seq: u64, n: i64 },
}
pub(crate) use crate::message_agg::{Agg, PendingSlot, RenameStep};