use std::io::{Read, Write};
use std::net::TcpStream;
use std::time::Duration;
use kevy_resp::{ArgvView, encode_error, parse_command};
use kevy_store::Store;
use crate::scope_integration;
pub(crate) fn cmd_move_scope<A: ArgvView + ?Sized>(
store: &mut Store,
args: &A,
out: &mut Vec<u8>,
) {
let Some((prefix_owned, from_id, to_id)) = parse_move_scope_args(args, out) else {
return;
};
let Some(target_addr) = validate_move_scope_route(&from_id, &to_id, out) else {
return;
};
if let Err(e) = scope_integration::migration_start(
prefix_owned.clone(),
from_id.to_string(),
to_id.to_string(),
) {
return encode_error(out, &format!("ERR MOVE-SCOPE: {e}"));
}
match ship_prefix_to_target(store, &prefix_owned, &target_addr) {
Ok(count) => {
scope_integration::migration_commit(&prefix_owned);
let reply = format!("+OK {count}\r\n");
out.extend_from_slice(reply.as_bytes());
}
Err(e) => {
scope_integration::migration_abort(&prefix_owned);
encode_error(out, &format!("ERR MOVE-SCOPE ship failed: {e}"));
}
}
}
fn parse_move_scope_args<A: ArgvView + ?Sized>(
args: &A,
out: &mut Vec<u8>,
) -> Option<(Vec<u8>, String, String)> {
if args.len() != 6 {
encode_error(
out,
"ERR wrong number of arguments — MOVE-SCOPE <prefix> FROM <from-id> TO <to-id>",
);
return None;
}
let Some(prefix) = args.get(1) else {
wrong_syntax(out);
return None;
};
let from_kw = args.get(2).unwrap_or_default();
let from_id = args.get(3).unwrap_or_default();
let to_kw = args.get(4).unwrap_or_default();
let to_id = args.get(5).unwrap_or_default();
if !from_kw.eq_ignore_ascii_case(b"FROM") || !to_kw.eq_ignore_ascii_case(b"TO") {
wrong_syntax(out);
return None;
}
let Ok(from_id) = std::str::from_utf8(from_id) else {
wrong_syntax(out);
return None;
};
let Ok(to_id) = std::str::from_utf8(to_id) else {
wrong_syntax(out);
return None;
};
Some((prefix.to_vec(), from_id.to_string(), to_id.to_string()))
}
fn validate_move_scope_route(from_id: &str, to_id: &str, out: &mut Vec<u8>) -> Option<String> {
match scope_integration::self_node_id() {
Some(me) if me == from_id => {}
Some(me) => {
encode_error(
out,
&format!("ERR MOVE-SCOPE: from-id {from_id:?} is not this node ({me:?})"),
);
return None;
}
None => {
encode_error(
out,
"ERR MOVE-SCOPE: [cluster] node_id is not configured on this node",
);
return None;
}
}
let Some(target_addr) = scope_integration::peer_addr(to_id) else {
encode_error(
out,
&format!("ERR MOVE-SCOPE: target node {to_id:?} not in [cluster] peers"),
);
return None;
};
Some(target_addr)
}
fn wrong_syntax(out: &mut Vec<u8>) {
encode_error(
out,
"ERR MOVE-SCOPE syntax: MOVE-SCOPE <prefix> FROM <from-id> TO <to-id>",
);
}
fn ship_prefix_to_target(
store: &mut Store,
prefix: &[u8],
target_addr: &str,
) -> Result<usize, String> {
let (bulk, count) = serialize_prefix(store, prefix);
let mut s = TcpStream::connect_timeout(
&target_addr.parse().map_err(|e| format!("bad target addr {target_addr:?}: {e}"))?,
Duration::from_secs(10),
)
.map_err(|e| format!("connect {target_addr:?}: {e}"))?;
s.set_read_timeout(Some(Duration::from_secs(60)))
.map_err(|e| format!("set_read_timeout: {e}"))?;
let mut req = Vec::new();
req.extend_from_slice(b"*3\r\n");
req.extend_from_slice(b"$17\r\nMOVE-SCOPE-INGEST\r\n");
req.extend_from_slice(format!("${}\r\n", prefix.len()).as_bytes());
req.extend_from_slice(prefix);
req.extend_from_slice(b"\r\n");
req.extend_from_slice(format!("${}\r\n", bulk.len()).as_bytes());
req.extend_from_slice(&bulk);
req.extend_from_slice(b"\r\n");
s.write_all(&req).map_err(|e| format!("write: {e}"))?;
let mut buf = [0u8; 256];
let n = s.read(&mut buf).map_err(|e| format!("read: {e}"))?;
let reply = &buf[..n];
if !reply.starts_with(b"+") {
return Err(format!(
"target replied non-OK: {:?}",
String::from_utf8_lossy(reply)
));
}
Ok(count)
}
pub(crate) fn cmd_move_scope_ingest<A: ArgvView + ?Sized>(
store: &mut Store,
args: &A,
out: &mut Vec<u8>,
) {
if args.len() != 3 {
return encode_error(
out,
"ERR wrong number of arguments — MOVE-SCOPE-INGEST <prefix> <bulk>",
);
}
let Some(prefix) = args.get(1) else {
return encode_error(out, "ERR MOVE-SCOPE-INGEST: missing prefix");
};
let Some(bulk) = args.get(2) else {
return encode_error(out, "ERR MOVE-SCOPE-INGEST: missing bulk");
};
let _guard = scope_integration::IngestGuard::enter(prefix.to_vec());
let mut buf = bulk.to_vec();
let mut applied = 0usize;
let mut scratch = Vec::with_capacity(256);
loop {
match parse_command(&buf) {
Ok(Some((argv, consumed))) => {
scratch.clear();
crate::dispatch::dispatch_into(store, &argv, &mut scratch);
buf.drain(..consumed);
applied += 1;
}
Ok(None) => break,
Err(_) => {
return encode_error(out, "ERR MOVE-SCOPE-INGEST: malformed bulk");
}
}
}
let reply = format!("+OK {applied}\r\n");
out.extend_from_slice(reply.as_bytes());
}
fn serialize_prefix(store: &mut Store, prefix: &[u8]) -> (Vec<u8>, usize) {
let mut bulk = Vec::new();
let mut count = 0usize;
let keys = store.collect_keys(None, None);
for key in keys {
if !key.starts_with(prefix) {
continue;
}
let ttl_ms = store.pttl(&key);
let abs_expire = if ttl_ms > 0 {
Some(kevy_store::now_unix_ms().saturating_add(ttl_ms as u64))
} else {
None
};
match store.type_of(&key) {
"string" => emit_string(store, &key, &mut bulk, &mut count),
"hash" => emit_hash(store, &key, &mut bulk, &mut count),
"list" => emit_list(store, &key, &mut bulk, &mut count),
"set" => emit_set(store, &key, &mut bulk, &mut count),
"zset" => emit_zset(store, &key, &mut bulk, &mut count),
_ => continue, }
if let Some(ms) = abs_expire {
let ms_str = ms.to_string();
append_resp_argv(&mut bulk, &[b"PEXPIREAT", &key, ms_str.as_bytes()]);
count += 1;
}
}
(bulk, count)
}
fn emit_string(store: &mut Store, key: &[u8], bulk: &mut Vec<u8>, count: &mut usize) {
if let Ok(Some(v)) = store.get(key) {
append_resp_argv(bulk, &[b"SET", key, &v]);
*count += 1;
}
}
fn emit_hash(store: &mut Store, key: &[u8], bulk: &mut Vec<u8>, count: &mut usize) {
let Ok(pairs) = store.hgetall(key) else { return };
if pairs.is_empty() {
return;
}
let mut parts: Vec<&[u8]> = Vec::with_capacity(2 + pairs.len());
parts.push(b"HSET");
parts.push(key);
for p in &pairs {
parts.push(p);
}
append_resp_argv(bulk, &parts);
*count += 1;
}
fn emit_list(store: &mut Store, key: &[u8], bulk: &mut Vec<u8>, count: &mut usize) {
let Ok(items) = store.lrange(key, 0, -1) else { return };
if items.is_empty() {
return;
}
let mut parts: Vec<&[u8]> = Vec::with_capacity(2 + items.len());
parts.push(b"RPUSH");
parts.push(key);
for item in &items {
parts.push(item);
}
append_resp_argv(bulk, &parts);
*count += 1;
}
fn emit_set(store: &mut Store, key: &[u8], bulk: &mut Vec<u8>, count: &mut usize) {
let Ok(members) = store.smembers(key) else { return };
if members.is_empty() {
return;
}
let mut parts: Vec<&[u8]> = Vec::with_capacity(2 + members.len());
parts.push(b"SADD");
parts.push(key);
for m in &members {
parts.push(m);
}
append_resp_argv(bulk, &parts);
*count += 1;
}
fn emit_zset(store: &mut Store, key: &[u8], bulk: &mut Vec<u8>, count: &mut usize) {
let Ok(items) = store.zrange(key, 0, -1) else { return };
if items.is_empty() {
return;
}
let score_strs: Vec<String> = items.iter().map(|(_, s)| format_score(*s)).collect();
let mut parts: Vec<&[u8]> = Vec::with_capacity(2 + items.len() * 2);
parts.push(b"ZADD");
parts.push(key);
for (i, (member, _)) in items.iter().enumerate() {
parts.push(score_strs[i].as_bytes());
parts.push(member);
}
append_resp_argv(bulk, &parts);
*count += 1;
}
fn format_score(s: f64) -> String {
format!("{s}")
}
fn append_resp_argv(out: &mut Vec<u8>, parts: &[&[u8]]) {
out.extend_from_slice(format!("*{}\r\n", parts.len()).as_bytes());
for p in parts {
out.extend_from_slice(format!("${}\r\n", p.len()).as_bytes());
out.extend_from_slice(p);
out.extend_from_slice(b"\r\n");
}
}
#[cfg(test)]
mod tests {
use super::*;
use kevy_resp::Argv;
fn argv(parts: &[&[u8]]) -> Argv {
let mut a = Argv::default();
for p in parts {
a.push(p);
}
a
}
fn fresh_store() -> Store {
Store::new()
}
#[test]
fn serialize_prefix_emits_set_for_strings() {
let mut store = fresh_store();
store.set(b"app:foo", b"v1".to_vec(), None, false, false);
store.set(b"app:bar", b"v2".to_vec(), None, false, false);
store.set(b"other:k", b"v3".to_vec(), None, false, false);
let (bulk, count) = serialize_prefix(&mut store, b"app:");
assert_eq!(count, 2, "two string keys under prefix");
let s = String::from_utf8_lossy(&bulk);
assert!(s.contains("$3\r\nSET\r\n"), "wire shape has SET: {s:?}");
assert!(s.contains("app:foo"), "key 1 present");
assert!(s.contains("app:bar"), "key 2 present");
assert!(!s.contains("other:k"), "non-matching key absent");
}
#[test]
fn serialize_prefix_emits_hset_for_hash_in_order() {
let mut store = fresh_store();
store
.hset(b"app:h", &[(b"f1".to_vec(), b"v1".to_vec()), (b"f2".to_vec(), b"v2".to_vec())])
.unwrap();
let (bulk, count) = serialize_prefix(&mut store, b"app:");
assert_eq!(count, 1);
let s = String::from_utf8_lossy(&bulk);
assert!(s.contains("HSET"), "HSET emitted: {s:?}");
}
#[test]
fn serialize_prefix_skips_non_matching_keys() {
let mut store = fresh_store();
store.set(b"foo", b"v".to_vec(), None, false, false);
let (bulk, count) = serialize_prefix(&mut store, b"app:");
assert_eq!(count, 0);
assert!(bulk.is_empty());
}
#[test]
fn ingest_handler_applies_embedded_commands_and_replies_ok() {
let mut store = fresh_store();
let mut bulk = Vec::new();
append_resp_argv(&mut bulk, &[b"SET", b"app:a", b"1"]);
append_resp_argv(&mut bulk, &[b"SET", b"app:b", b"2"]);
let args = argv(&[b"MOVE-SCOPE-INGEST", b"app:", &bulk]);
let mut out = Vec::new();
cmd_move_scope_ingest(&mut store, &args, &mut out);
assert_eq!(out, b"+OK 2\r\n", "wire reply shape");
assert_eq!(
store.get(b"app:a").map(|v| v.map(|c| c.into_owned())),
Ok(Some(b"1".to_vec()))
);
assert_eq!(
store.get(b"app:b").map(|v| v.map(|c| c.into_owned())),
Ok(Some(b"2".to_vec()))
);
}
#[test]
fn ingest_handler_rejects_wrong_arity() {
let mut store = fresh_store();
let args = argv(&[b"MOVE-SCOPE-INGEST", b"only-one"]);
let mut out = Vec::new();
cmd_move_scope_ingest(&mut store, &args, &mut out);
assert!(out.starts_with(b"-ERR"), "got {:?}", String::from_utf8_lossy(&out));
}
#[test]
fn move_scope_rejects_bad_syntax() {
let mut store = fresh_store();
let args = argv(&[b"MOVE-SCOPE", b"p:", b"NOT-FROM", b"A", b"TO", b"B"]);
let mut out = Vec::new();
cmd_move_scope(&mut store, &args, &mut out);
assert!(out.starts_with(b"-ERR"));
}
#[test]
fn move_scope_rejects_when_self_node_id_not_configured() {
let mut store = fresh_store();
let args = argv(&[b"MOVE-SCOPE", b"p:", b"FROM", b"A", b"TO", b"B"]);
let mut out = Vec::new();
cmd_move_scope(&mut store, &args, &mut out);
assert!(out.starts_with(b"-ERR"));
let s = String::from_utf8_lossy(&out);
assert!(s.contains("node_id is not configured") || s.contains("from-id"), "{s}");
}
}