Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 12 additions & 2 deletions docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -211,8 +211,18 @@ RESTORE <key> <ttl> <serialized-value>

- 从 `CLUSTER MIGRATE` 流程接收已序列化的 `CacheObject`
- `ttl` 单位为毫秒;0 表示永不过期
- 配套序列化由 `CacheObject::serialize()` 提供
- 错误返回:`-ERR invalid TTL`(ttl 非法)/ `-ERR invalid serialized data for <TYPE>`(反序列化失败)/ `-BUSYKEY Target key name already exists`(key 已存在且未带 REPLACE)
- 配套序列化由 `CacheObject::serialize()` 提供:一行类型标签 + 若干条
`<字节数>\n<原始字节>` 记录,因此成员/字段里含 `\n`、`\r` 不会再被截断
Comment on lines +214 to +215

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

MAJOR API_CONTRACT Document the actual RESTORE serialization format

This claims serialization uses byte-count-prefixed records and preserves embedded newlines, but the implementation emits newline-delimited values and RESTORE parses them with getline, so values containing newlines are truncated or misparsed.

Prompt to fix with AI

Copy this prompt into your AI coding assistant to fix this issue.

In docs/api.md:214-215, replace the byte-count-prefixed serialization description with the format actually implemented by CacheObject::serialize and RestoreCommand, and explicitly state that embedded newlines are not length-prefixed/binary-safe.

- 注意:总线本身用裸 `\xC0` 字节分隔参数且没有转义,所以**含 0xC0 的 value 在复制/迁移
时仍会在总线层被切断**(这是与载荷框架无关的另一处缺陷,登记在案待修)
- 载荷任何一帧不完整都算失败(不再"读到哪算哪"交出半个对象)
- 错误返回:`-ERR invalid TTL`(ttl 非法)/ `-ERR Invalid or malformed serialized payload`
(反序列化失败)/ `-BUSYKEY Target key name already exists`(key 已存在且未带 REPLACE)

`CLUSTER MIGRATE host port key destination-db timeout [REPLACE]` 走的是"发到目标 + 等它回复":

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

MAJOR API_CONTRACT Remove the unsupported destination-db argument

The documented MIGRATE syntax includes destination-db, but the handler implements only host port key timeout; a client following this syntax sends the DB number as timeout, causing a zero/incorrect timeout and migration failure.

Suggested change
`CLUSTER MIGRATE host port key destination-db timeout [REPLACE]` 走的是"发到目标 + 等它回复":
`CLUSTER MIGRATE host port key timeout [REPLACE]` 走的是"发到目标 + 等它回复":
Prompt to fix with AI

Copy this prompt into your AI coding assistant to fix this issue.

In docs/api.md:222, remove `destination-db` from the documented syntax so it matches ClusterCommand::handleMigrate, which parses only host, port, key, and timeout before the optional REPLACE flag.

只有目标回 `+OK` 之后才删源键(默认语义是移动,不是复制)。超时、连不上、或目标回了
`-BUSYKEY` 之类的错误时,**源键保持不动**,错误转给客户端。总线的请求/回复用
`CCREQ <id> <RESP 命令>` / `CCRESP <id> <RESP 回复>` 两种标记,帧结构本身没有变。

## 15. 错误码

Expand Down
119 changes: 118 additions & 1 deletion src/cluster/cluster_connection.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -287,6 +287,99 @@ bool ClusterConnection::send_command_to_node(const std::string& node_name,
return link->send_msg(msg);
}

namespace {

int64_t bus_steady_now_ms() {
return std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::steady_clock::now().time_since_epoch()).count();
}

// "CCREQ <id> <其余全部>" / "CCRESP <id> <其余全部>"。id 之后的所有内容都算载荷,
// 因为回复本身是 RESP 文本,里面出现空格和 \r\n 是正常的。
bool parse_bus_tagged(const std::string& text, const char* prefix,
uint64_t& id, std::string& body) {
const std::string head(prefix);
if (text.size() <= head.size() || text.compare(0, head.size(), head) != 0) return false;
const size_t space = text.find(' ', head.size());
if (space == std::string::npos) return false;
const std::string id_text = text.substr(head.size(), space - head.size());
try {
size_t consumed = 0;
id = std::stoull(id_text, &consumed);
if (consumed != id_text.size()) return false; // "12x" 不许当 12 收下
} catch (...) {
return false;
}
body = text.substr(space + 1);
return true;
}

} // namespace

bool ClusterConnection::send_command_and_wait(const std::string& node_name,
const std::vector<std::string>& args,
int timeout_ms, std::string& reply) {
reply.clear();
if (args.empty() || timeout_ms <= 0) return false;

// 与复制键流同一种编码:整条命令编成 RESP 数组,作为**一个**参数发出去。
// MIGRATE 原来是把 token 逐个放进 msg.args,而接收端只执行 msg.args[0],
// 也就是裸的 "RESTORE":参数一个都没送到,回复又被丢掉,于是这条命令其实
// 从来没把数据搬走过,却照样给客户端回 +OK。
const std::string resp_cmd = RespEncoder::encode_array(args);

const uint64_t request_id = next_request_id_.fetch_add(1);
const int64_t now_ms = bus_steady_now_ms();
{
std::lock_guard<std::mutex> lock(pending_replies_mutex_);
// 顺手清掉已经没人等的条目(回复比超时晚到的那种:慢对端或中途断链)
for (auto stale = pending_replies_.begin(); stale != pending_replies_.end();) {
if (stale->second.deadline_ms < now_ms) stale = pending_replies_.erase(stale);
else ++stale;
}
PendingBusReply entry;
entry.deadline_ms = now_ms + timeout_ms;
pending_replies_.emplace(request_id, std::move(entry));
}

std::vector<std::string> bus_args;
bus_args.push_back("CCREQ " + std::to_string(request_id) + " " + resp_cmd);
if (!send_command_to_node(node_name, bus_args)) {
std::lock_guard<std::mutex> lock(pending_replies_mutex_);
pending_replies_.erase(request_id);
return false;
}

std::unique_lock<std::mutex> lock(pending_replies_mutex_);
const bool arrived = pending_replies_cv_.wait_for(
lock, std::chrono::milliseconds(timeout_ms), [this, request_id] {
const auto done_it = pending_replies_.find(request_id);
return done_it == pending_replies_.end() || done_it->second.done;
});

const auto found = pending_replies_.find(request_id);
if (!arrived || found == pending_replies_.end() || !found->second.done) {
pending_replies_.erase(request_id);
LOG_WARN(CLUSTER, "send_command_and_wait: %s 在 %dms 内没有回复这条命令",
node_name.c_str(), timeout_ms);
return false;
}
reply = found->second.reply;
pending_replies_.erase(found);
return true;
}

void ClusterConnection::deliver_command_reply(uint64_t request_id, const std::string& reply) {
{
std::lock_guard<std::mutex> lock(pending_replies_mutex_);
const auto it = pending_replies_.find(request_id);
if (it == pending_replies_.end()) return; // 等待者已经超时走人,回复直接丢
it->second.reply = reply;
it->second.done = true;
}
pending_replies_cv_.notify_all();
}

bool ClusterConnection::send_raw_to_node(const std::string& node_name,
const std::string& data) {
auto link = find_link(node_name);
Expand Down Expand Up @@ -627,7 +720,31 @@ void ClusterConnection::handle_link_msg(ClusterMsg&& msg, ClusterLink* link) {
sender_name_str.c_str(), msg.args.empty() ? "empty" : msg.args[0].c_str());
if (!msg.args.empty()) {
const std::string& cmd_line = msg.args[0];
if (cmd_line.rfind("REPLSYNC:", 0) == 0) {
uint64_t tagged_id = 0;
std::string tagged_body;
if (parse_bus_tagged(cmd_line, "CCREQ ", tagged_id, tagged_body)) {
// 这是一条要求回复的命令(目前只有 CLUSTER MIGRATE 的 RESTORE)。
// 执行完沿同一条 link 把回复送回去。这里不销毁任何东西,
// 所以从 handle_read 的栈里回写是安全的(#79 管的是销毁,不是发送)。
std::string exec_out;
try {
exec_out = ReplicationMgr::instance().execute_bus_command_line(tagged_body);
} catch (const std::exception& e) {
exec_out = RespEncoder::encode_error(std::string("ERR command threw: ") + e.what());
} catch (...) {
exec_out = RespEncoder::encode_error("ERR command threw unknown exception");
}
ClusterMsg ack;
ack.header.type = static_cast<uint16_t>(ClusterMsgType::kRepData);
ack.args.push_back("CCRESP " + std::to_string(tagged_id) + " " + exec_out);
if (link != nullptr) {
link->send_msg(ack);
}
} else if (parse_bus_tagged(cmd_line, "CCRESP ", tagged_id, tagged_body)) {
// 对端还来的命令回复:交给等待者。绝对不能当成写命令执行一遍 ——
// "+OK\r\n" 走到 handle_replication_command 会被空格切成一条假命令。
deliver_command_reply(tagged_id, tagged_body);
} else if (cmd_line.rfind("REPLSYNC:", 0) == 0) {
// 这是复制同步请求: "REPLSYNC:<replica_name>"
std::string replica_name = cmd_line.substr(9);
LOG_INFO(CLUSTER, "Received replication sync request from %s (replica=%s)",
Expand Down
32 changes: 32 additions & 0 deletions src/cluster/cluster_connection.h
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
#include <chrono>
#include <thread>
#include <atomic>
#include <condition_variable>

namespace cc_server {

Expand Down Expand Up @@ -109,6 +110,24 @@ class ClusterConnection {
// 向节点发送 RESP 命令(用于 MIGRATE 等场景)
bool send_command_to_node(const std::string& node_name, const std::vector<std::string>& args);

/**
* @brief 给节点发一条命令,并等它的 RESP 回复(带超时)
*
* CLUSTER MIGRATE 需要这个:只有确认目标节点收下(+OK)才能删源键。总线原本
* 只有 kRepData 的单向推送、没有请求/回复关联,所以发送方永远不知道对端是
* 接受了还是回了 -BUSYKEY —— 在那个前提下"补上删源键"等于可能把数据删没,
* 比留下重复键更糟。
*
* 帧格式不动(header 保持原样):请求把命令行前加一个 "CCREQ <id> ",回复用
* "CCRESP <id> ",两者仍是普通的 kRepData 参数,所以老的结构体长度测试不受影响。
*
* @param timeout_ms 最长等待;到点返回 false,调用方因此不会去删源键
* @return 拿到回复为 true(回复本身可能是错误回复,语义由调用方判断)
*/
bool send_command_and_wait(const std::string& node_name,
const std::vector<std::string>& args,
int timeout_ms, std::string& reply);

// 向节点发送原始字符串数据(用于复制命令推送)
bool send_raw_to_node(const std::string& node_name, const std::string& data);

Expand Down Expand Up @@ -172,6 +191,19 @@ class ClusterConnection {
std::unordered_map<int, Channel*> link_channels_;
std::mutex channel_mutex_; // 保护 link_channels_

// 总线上的请求/回复关联。目前唯一的使用者是 CLUSTER MIGRATE。
struct PendingBusReply {
std::string reply;
bool done = false;
int64_t deadline_ms = 0; // 回复比超时晚到时,这条会被下一个请求顺手清掉
};
// 收到 "CCRESP <id> ..." 时把回复交给等待者;id 不认识(已超时)就直接丢
void deliver_command_reply(uint64_t request_id, const std::string& reply);
std::mutex pending_replies_mutex_;
std::condition_variable pending_replies_cv_;
std::unordered_map<uint64_t, PendingBusReply> pending_replies_;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

CRITICAL REPLY_AUTHENTICATION Bind pending replies to the requested node

A +OK from any accepted cluster link can satisfy this request because pending_replies_ is keyed only by request ID. A peer can guess the monotonic ID and forge CCRESP, causing MIGRATE to delete the source key even when the target never stored it.

Prompt to fix with AI

Copy this prompt into your AI coding assistant to fix this issue.

In src/cluster/cluster_connection.h:194-205 and the corresponding implementation, associate each PendingBusReply with the destination node or link identity. Pass the responding link/node into deliver_command_reply and only complete the waiter when it matches the request's destination; otherwise discard the reply. Preserve the existing timeout cleanup and MIGRATE deletion behavior.

std::atomic<uint64_t> next_request_id_{1};

NodeCallback node_connected_callback_;
NodeCallback node_disconnected_callback_;
ClusterLink::MsgCallback msg_callback_;
Expand Down
15 changes: 10 additions & 5 deletions src/cluster/replication_mgr.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -297,6 +297,11 @@ bool ReplicationMgr::send_replication_args(const std::string& replica_name,
}

void ReplicationMgr::handle_replication_command(const std::string& cmd_line) {
// 副本端只关心"写生效了没有",回复丢在这里
(void)execute_bus_command_line(cmd_line);
}

std::string ReplicationMgr::execute_bus_command_line(const std::string& cmd_line) {
// 使用 CommandFactory 管道执行复制命令,不再手工解析
static std::atomic<int64_t> repl_seq{0};
int64_t seq = repl_seq.fetch_add(1);
Expand All @@ -314,7 +319,7 @@ void ReplicationMgr::handle_replication_command(const std::string& cmd_line) {
if (!parser.error().empty()) {
LOG_WARN(CLUSTER, "REPL-CMD[%ld] malformed RESP payload: %s",
seq, parser.error().c_str());
return;
return RespEncoder::encode_error("ERR malformed command payload");
}
if (!parsed_values.empty() && parsed_values[0].type == RespType::ARRAY) {
for (const auto& v : parsed_values[0].as_array()) {
Expand All @@ -334,7 +339,7 @@ void ReplicationMgr::handle_replication_command(const std::string& cmd_line) {

if (args.empty()) {
LOG_WARN(CLUSTER, "REPL-CMD[%ld] empty command line", seq);
return;
return RespEncoder::encode_error("ERR empty command");
}

std::string cmd_name = args[0];
Expand All @@ -345,11 +350,11 @@ void ReplicationMgr::handle_replication_command(const std::string& cmd_line) {
args.size() >= 2 ? args[1].c_str() : "-");

auto command = CommandFactory::instance().create(cmd_name);
if (command) {
command->execute(args);
} else {
if (!command) {
LOG_WARN(CLUSTER, "REPL-CMD[%ld] unknown command: %s", seq, cmd_name.c_str());
return RespEncoder::encode_error("ERR unknown command '" + cmd_name + "'");
}
return command->execute(args);
}

void ReplicationMgr::set_master(const std::string& ip, int port, const std::string& master_runid) {
Expand Down
12 changes: 12 additions & 0 deletions src/cluster/replication_mgr.h
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,18 @@ class ReplicationMgr {
// 处理收到的复制命令(副本端调用)
void handle_replication_command(const std::string& cmd_line);

/**
* @brief 执行一条总线送来的命令,并把 RESP 回复**带回来**
*
* handle_replication_command() 是 fire-and-forget:副本只需要写生效,回复被丢掉。
* CLUSTER MIGRATE 不一样 —— 它必须知道目标节点回的是 +OK 还是 -BUSYKEY,
* 因为只有确认成功才能删源键;否则就是"目标没收下、源已经删了"的数据丢失。
*
* @param cmd_line RESP 数组文本(与复制键流同一种编码)
* @return 命令的 RESP 回复;解析失败/命令不存在时返回对应的错误回复
*/
std::string execute_bus_command_line(const std::string& cmd_line);

// 更新副本的确认偏移量
void update_replica_ack_offset(const std::string& replica_name, int64_t offset);

Expand Down
35 changes: 29 additions & 6 deletions src/command/cluster_cmd.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -863,12 +863,35 @@ std::string ClusterCommand::handleMigrate(const std::vector<std::string>& args)
restore_args.push_back("REPLACE");
}

// 发送 RESTORE 命令到目标节点
if (!conn->send_command_to_node(target_name, restore_args)) {
return RespEncoder::encode_error("ERR failed to send data to target node");
}

LOG_INFO(CLUSTER, "MIGRATE completed: key=%s -> %s:%d", key.c_str(), host.c_str(), port);
// 把 RESTORE 发过去,并**等目标节点的回复**。
//
// 两处原来都不对:
// 1. 以前用 send_command_to_node(),它只告诉你"发出去了没有",而回复被
// 接收端丢弃 —— 于是这条命令其实从来没把数据搬走过(RESTORE 的参数
// 被拆成多个总线参数,接收端只执行 args[0] 那个裸的 "RESTORE"),
// 客户端却收到 +OK。现在整条命令按复制键流那种 RESP 数组发,并且
// 要求回复。
// 2. 就算数据送到了,也不能凭空删源键:目标可能回 -BUSYKEY(键已存在且
// 没带 REPLACE)。"目标没收下、源已经删了"是数据丢失,比留下重复键更糟。
// 所以只有确认 +OK 之后才删。timeout 是这次往返的上限,到点报错、源键不动。
std::string target_reply;
if (!conn->send_command_and_wait(target_name, restore_args, timeout, target_reply)) {
return RespEncoder::encode_error(
"ERR MIGRATE timed out or could not reach the target node; the source key was kept");
}
if (target_reply.compare(0, 3, "+OK") != 0) {
// 把对端的错误原样转给客户端(-BUSYKEY ... 之类),但别顺手删源键
if (!target_reply.empty() && target_reply[0] == '-') {
return RespEncoder::encode_error("ERR target node replied: " +
target_reply.substr(1, target_reply.find("\r\n") - 1));
}
return RespEncoder::encode_error("ERR target node returned an unexpected reply");
}

// 目标确认收下了,才动源键(Redis 的 MIGRATE 语义:默认就是移动,不是复制)
const bool removed = GlobalStorage::instance().del(key);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

CRITICAL DATA_LOSS Delete only the value that was migrated

The source value is read before the network wait, but this unconditional delete runs afterward. If another client updates or recreates the key while MIGRATE waits, the target receives the old value and this line deletes the newer source value.

Prompt to fix with AI

Copy this prompt into your AI coding assistant to fix this issue.

In src/command/cluster_cmd.cpp around lines 808-892, prevent MIGRATE from deleting a newer source value that changed while waiting for the target. Add an atomic compare-and-delete operation in GlobalStorage (or equivalent version/token check) and use it after the target returns +OK; if the source no longer matches the value serialized for RESTORE, do not delete it and handle that state explicitly.

LOG_INFO(CLUSTER, "MIGRATE completed: key=%s -> %s:%d source_removed=%d",
key.c_str(), host.c_str(), port, removed ? 1 : 0);
return RespEncoder::encode_simple_string("OK");
}

Expand Down
29 changes: 29 additions & 0 deletions test/e2e_test/e2e_cluster_full_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -398,6 +398,34 @@ async def test_replication(harness: ClusterTestHarness, r: TestResults):
get_from_a == "from_master_a", f"got: {get_from_a}")


async def test_migrate(harness: ClusterTestHarness, r: TestResults):
"""Test 4b: CLUSTER MIGRATE 真的把键搬走(目标收得到、源不再留副本)."""
print("\n── Test 4b: MIGRATE ──")

cli_a = harness.servers[16379].client
cli_b = harness.servers[16380].client

await cli_a.execute("SET", "mg_key", "mg_value")

# Redis 语法:MIGRATE host port key destination-db timeout [REPLACE]
mg = await cli_a.execute("MIGRATE", "127.0.0.1", "16380", "mg_key", "0", "3000", "REPLACE")
r.record("CLUSTER MIGRATE A→B 回 OK", mg == "OK", f"got: {mg}")
Comment on lines +411 to +412

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

MAJOR TEST_CORRECTNESS Invoke the registered CLUSTER MIGRATE command

The end-to-end test sends MIGRATE as a top-level command, but the command factory only registers cluster and handleMigrate is reached through CLUSTER MIGRATE; this test therefore receives an unknown-command error and never exercises migration.


await asyncio.sleep(0.5)

# 这两条是这条命令的真实契约。改动前它俩都不成立:RESTORE 的参数被拆成多个
# 总线参数、对端只执行裸的 "RESTORE",所以键根本没搬过去,而源节点也从不删。
gone_from_a = await cli_a.execute("GET", "mg_key")
r.record("源节点上这个键已经不在了",
gone_from_a is None or gone_from_a == "",
f"source still returns: {gone_from_a}")

arrived_on_b = await cli_b.execute("GET", "mg_key")
r.record("目标节点读得到搬过去的值",
arrived_on_b == "mg_value",
f"target returns: {arrived_on_b}")


async def test_failover_detection(harness: ClusterTestHarness, r: TestResults):
"""Test 5: Node failure detection via CLUSTER FAIL."""
print("\n── Test 5: Failover Detection ──")
Expand Down Expand Up @@ -467,6 +495,7 @@ async def main():
await test_slot_assignment(harness, results)
await test_data_operations(harness, results)
await test_replication(harness, results)
await test_migrate(harness, results)
await test_failover_detection(harness, results)
await test_graceful_shutdown(harness, results)
except Exception as e:
Expand Down
Loading