diff --git a/docs/api.md b/docs/api.md index d618bfa..fe8eb2c 100644 --- a/docs/api.md +++ b/docs/api.md @@ -211,8 +211,18 @@ RESTORE - 从 `CLUSTER MIGRATE` 流程接收已序列化的 `CacheObject` - `ttl` 单位为毫秒;0 表示永不过期 -- 配套序列化由 `CacheObject::serialize()` 提供 -- 错误返回:`-ERR invalid TTL`(ttl 非法)/ `-ERR invalid serialized data for `(反序列化失败)/ `-BUSYKEY Target key name already exists`(key 已存在且未带 REPLACE) +- 配套序列化由 `CacheObject::serialize()` 提供:一行类型标签 + 若干条 + `<字节数>\n<原始字节>` 记录,因此成员/字段里含 `\n`、`\r` 不会再被截断 +- 注意:总线本身用裸 `\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]` 走的是"发到目标 + 等它回复": +只有目标回 `+OK` 之后才删源键(默认语义是移动,不是复制)。超时、连不上、或目标回了 +`-BUSYKEY` 之类的错误时,**源键保持不动**,错误转给客户端。总线的请求/回复用 +`CCREQ ` / `CCRESP ` 两种标记,帧结构本身没有变。 ## 15. 错误码 diff --git a/src/cluster/cluster_connection.cpp b/src/cluster/cluster_connection.cpp index 6ff1cb2..5c8ba87 100644 --- a/src/cluster/cluster_connection.cpp +++ b/src/cluster/cluster_connection.cpp @@ -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::steady_clock::now().time_since_epoch()).count(); + } + + // "CCREQ <其余全部>" / "CCRESP <其余全部>"。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& 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 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 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 lock(pending_replies_mutex_); + pending_replies_.erase(request_id); + return false; + } + + std::unique_lock 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 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); @@ -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(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:" std::string replica_name = cmd_line.substr(9); LOG_INFO(CLUSTER, "Received replication sync request from %s (replica=%s)", diff --git a/src/cluster/cluster_connection.h b/src/cluster/cluster_connection.h index 7916756..60dca75 100644 --- a/src/cluster/cluster_connection.h +++ b/src/cluster/cluster_connection.h @@ -14,6 +14,7 @@ #include #include #include +#include namespace cc_server { @@ -109,6 +110,24 @@ class ClusterConnection { // 向节点发送 RESP 命令(用于 MIGRATE 等场景) bool send_command_to_node(const std::string& node_name, const std::vector& args); + /** + * @brief 给节点发一条命令,并等它的 RESP 回复(带超时) + * + * CLUSTER MIGRATE 需要这个:只有确认目标节点收下(+OK)才能删源键。总线原本 + * 只有 kRepData 的单向推送、没有请求/回复关联,所以发送方永远不知道对端是 + * 接受了还是回了 -BUSYKEY —— 在那个前提下"补上删源键"等于可能把数据删没, + * 比留下重复键更糟。 + * + * 帧格式不动(header 保持原样):请求把命令行前加一个 "CCREQ ",回复用 + * "CCRESP ",两者仍是普通的 kRepData 参数,所以老的结构体长度测试不受影响。 + * + * @param timeout_ms 最长等待;到点返回 false,调用方因此不会去删源键 + * @return 拿到回复为 true(回复本身可能是错误回复,语义由调用方判断) + */ + bool send_command_and_wait(const std::string& node_name, + const std::vector& args, + int timeout_ms, std::string& reply); + // 向节点发送原始字符串数据(用于复制命令推送) bool send_raw_to_node(const std::string& node_name, const std::string& data); @@ -172,6 +191,19 @@ class ClusterConnection { std::unordered_map 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 不认识(已超时)就直接丢 + 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 pending_replies_; + std::atomic next_request_id_{1}; + NodeCallback node_connected_callback_; NodeCallback node_disconnected_callback_; ClusterLink::MsgCallback msg_callback_; diff --git a/src/cluster/replication_mgr.cpp b/src/cluster/replication_mgr.cpp index 26c5492..09db395 100644 --- a/src/cluster/replication_mgr.cpp +++ b/src/cluster/replication_mgr.cpp @@ -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 repl_seq{0}; int64_t seq = repl_seq.fetch_add(1); @@ -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()) { @@ -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]; @@ -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) { diff --git a/src/cluster/replication_mgr.h b/src/cluster/replication_mgr.h index 1153ba8..ecb62eb 100644 --- a/src/cluster/replication_mgr.h +++ b/src/cluster/replication_mgr.h @@ -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); diff --git a/src/command/cluster_cmd.cpp b/src/command/cluster_cmd.cpp index 038fa91..78daacf 100644 --- a/src/command/cluster_cmd.cpp +++ b/src/command/cluster_cmd.cpp @@ -863,12 +863,35 @@ std::string ClusterCommand::handleMigrate(const std::vector& 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); + 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"); } diff --git a/test/e2e_test/e2e_cluster_full_test.py b/test/e2e_test/e2e_cluster_full_test.py index 02a3c03..dd82798 100755 --- a/test/e2e_test/e2e_cluster_full_test.py +++ b/test/e2e_test/e2e_cluster_full_test.py @@ -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}") + + 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 ──") @@ -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: