Skip to content

fix(cluster): CLUSTER MIGRATE 真的把键搬走,并且只在目标确认后才删源键 - #92

Merged
zimingttkx merged 1 commit into
mainfrom
fix-cluster-migrate-roundtrip
Oct 5, 2026
Merged

zimingttkx merged 1 commit into
mainfrom
fix-cluster-migrate-roundtrip

Conversation

@zimingttkx

Copy link
Copy Markdown
Owner

两个独立的问题凑在一起,才解释得通这条命令为什么"看起来能用"

1. 数据从来没被搬走过。 MIGRATE 把 RESTORE 的 token 逐个塞进 msg.args,而接收端的 kRepData 分支只执行 msg.args[0] —— 也就是裸的一个词 "RESTORE",参数一个都没送到。handle_replication_command 用 CommandFactory 执行它 → 少参数报错 → 回复被丢弃。客户端那边收到的是 MIGRATE 自己回的 +OK。复制键流走的是另一种打包(RespEncoder::encode_array(args) 整段当一个参数),所以复制能通、迁移不能通。

2. 就算送到了也不能删源键。 send_command_to_node() 的返回值只说明"发出去了没有",总线上没有请求/回复关联(没有 correlation id、没有 promise/future),所以发送方永远不知道目标是 +OK 还是 -BUSYKEY。在那个前提下补删除,等于可能"目标没收下、源已经删了"—— 那是数据丢失,比留下重复键更糟。

改法

  • ReplicationMgr::execute_bus_command_line():同一条解析+执行路径,但把 RESP 回复带回来;handle_replication_command 变成"执行并丢回复"的薄壳,副本路径行为不变。
  • ClusterConnection::send_command_and_wait():请求打 CCREQ <id> 、回复打 CCRESP <id> ,仍然只是普通的 kRepData 参数 —— 帧结构没动,ClusterMsgHeader 长度那几条测试不受影响。等待有硬超时;超时/发送失败都不碰源键。下一次请求时顺手清掉"回复比超时晚到"留下的条目,避免无界增长。
  • MIGRATE 按复制键流那种 RESP 数组发命令,只有目标回 +OK 才 del 源键;其它情况把对端的错误透传给客户端并保留源键。等待用的是命令自带的 timeout 参数 —— 这也是 Redis MIGRATE 本身的语义(它会阻塞到 timeout)。
  • 接收端收到 CCRESP 时不把它当写命令执行("+OK\r\n" 会被空格切成一条假命令)。

测试(Daily 的 e2e 档)

e2e_cluster_full_test.py 新增 Test 4b:A 上 SET mg_key → MIGRATE 到 B → 断言 ① 命令回 OK ② A 上读不到了 ③ B 上读得到原值。改动前 ② 和 ③ 都不成立。该脚本的退出码是 failed > 0 → 1,所以这三条真能红。

本机做实与限制

五个文件的括号净差与 origin/main 逐文件对比一致(防 rebase 压坏结构);py_compile 过。其余是 sys/epoll.h 依赖,MinGW 编不了,由 CI 学。

docs/api.md §14 顺手改了三件事:RESTORE 的错误文案(#89 之后)、MIGRATE 的"确认后删"语义、以及明确记下总线裸 \xC0 分隔没有转义这个仍然存在的独立缺陷(含 0xC0 的 value 在复制/迁移时会在总线层被切断)—— 不假装这条已经解决。

两个独立的问题,凑在一起才解释得通这条命令为什么"看起来能用":

1. 数据从来没被搬走。MIGRATE 把 RESTORE 的 token 逐个塞进 msg.args,而接收端的
   kRepData 分支只执行 `msg.args[0]` —— 也就是裸的一个词 "RESTORE",参数一个都没送到。
   走 handle_replication_command 之后是 CommandFactory 执行 RESTORE 少参数 → 回一条
   错误 → 回复被丢弃。客户端那边收到的是 MIGRATE 自己的 +OK。复制键流走的是另一种
   打包(RespEncoder::encode_array(args) 整段当一个参数),所以复制能通、迁移不能通。
2. 就算送到了也不能删源键。send_command_to_node() 的返回值只说明"发出去了没有",
   总线上没有请求/回复关联(没有 correlation id、没有 promise/future),所以发送方永远
   不知道目标是 +OK 还是 -BUSYKEY。在这个前提下补删除等于可能把数据删没。

改法:

- ReplicationMgr 拆出 execute_bus_command_line():同一条解析+执行路径,但把 RESP 回复
  带回来(handle_replication_command 变成"执行并丢回复"的薄壳,副本路径行为不变)。
- ClusterConnection::send_command_and_wait():请求打 "CCREQ <id> "、回复打 "CCRESP <id> ",
  仍然只是普通的 kRepData 参数,**帧结构没动**(header 长度那几条测试不受影响)。
  回复到达就把结果交给等待者并唤醒 cv;等待有硬超时,超时/发送失败都不删源键。
  顺手在下一次请求时清掉"回复比超时晚到"留下的条目,避免无界增长。
- MIGRATE 按复制键流那种 RESP 数组发命令,等目标回话;只有 +OK 才 del 源键,
  其它情况把对端的错误透传给客户端并且保留源键。等待用的是命令自带的 timeout 参数,
  这也是 Redis MIGRATE 本身的语义(它会阻塞到 timeout)。
- 接收端收到 CCRESP 时**不**把它当写命令执行("+OK\r\n" 会被空格切成一条假命令)。

测试(e2e,Daily 的 e2e 档跑):e2e_cluster_full_test.py 新增 Test 4b:A 上 SET mg_key →
MIGRATE 到 B → 断言 ①命令回 OK ②A 上读不到了 ③B 上读得到原值。改动前 ②③ 都不成立。
该脚本的退出码是 failed>0 即 1,所以这三条真能红。

本机做实:括号净差与 origin/main 逐文件对比一致;其余部分依赖 epoll/socket,MinGW 编不了,
由 CI 学。docs/api.md §14 顺手把 RESTORE 的错误文案与 MIGRATE 的"确认后删"写清楚,
并明确记下总线裸 \xC0 分隔没有转义这个仍存在的独立缺陷(不假装已经解决)。
Copilot AI balanced review requested due to automatic review settings October 5, 2026 06:20
@zimingttkx
zimingttkx enabled auto-merge (squash) October 5, 2026 06:20

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

@entelligence-ai-pr-reviews

Copy link
Copy Markdown

EntelligenceAI PR Summary

修复 CLUSTER MIGRATE 未正确传递 RESTORE 参数且无确认即返回成功的问题。新增带请求 ID 的总线请求/回复机制,只有目标返回 +OK 后才删除源键,并补充序列化边界说明与端到端迁移测试。


Review Scorecard

Dimension Rating Basis
Code Quality ●○○○○ 1/5 — Unacceptable 2 critical, 3 significant finding(s) — reviewer rated the code Unacceptable
Blast Radius Low changed symbols are referenced only within their own file(s); no high-impact surface touched, 6 file(s) / ~227 line(s) changed (size only — not a blast signal)
Merge Confidence ●●○○○ 2/5 — Changes Needed code quality 1/5 × Low blast radius

Issues found:

  • Critical src/cluster/cluster_connection.h — Bind pending replies to the requested node
  • Critical src/command/cluster_cmd.cpp — Delete only the value that was migrated
  • Significant docs/api.md — Document the actual RESTORE serialization format
  • Significant test/e2e_test/e2e_cluster_full_test.py — Invoke the registered CLUSTER MIGRATE command
  • Significant docs/api.md — Remove the unsupported destination-db argument

Fix before merge but low risk; code quality rated Unacceptable (1/5): 2 critical, 3 significant findings in 4 files, the most severe in src/cluster/cluster_connection.h (Bind pending replies to the requested node).

Need to merge before these are addressed? Anyone with write access can comment @entelligence /approve to approve it now. @entelligence help lists every command.

Evaluated against
  • 7/7 changed files reviewed
  • criteria: correctness, security & access control, robustness & error handling, concurrency & data integrity, repo conventions / steering docs
  • steering docs: none found in repo
Files requiring special attention
  • src/cluster/cluster_connection.h
  • src/command/cluster_cmd.cpp
  • docs/api.md
  • test/e2e_test/e2e_cluster_full_test.py

@entelligence-ai-pr-reviews

Copy link
Copy Markdown

Walkthrough

This PR adds synchronous command execution over cluster links and uses confirmed target responses to make CLUSTER MIGRATE transactional. It updates bus command handling, documents RESTORE and migration protocols and limitations, and adds end-to-end coverage verifying successful key movement and source deletion.

Changes

Files

  • src/cluster/cluster_connection.cpp
  • src/cluster/cluster_connection.h

Added tagged CCREQ/CCRESP bus-message handling, request IDs, timeout-based pending-reply tracking, condition-variable synchronization, and reply delivery. Cluster commands are encoded as RESP arrays, executed remotely through ReplicationMgr, and correlated responses are routed to waiting callers while unknown or late replies are discarded.

Files

  • src/cluster/replication_mgr.cpp
  • src/cluster/replication_mgr.h

Added synchronous execute_bus_command_line support for parsing and executing bus-delivered commands through CommandFactory, returning RESP-encoded success or error responses. The existing replication handler remains fire-and-forget, while malformed payloads, empty commands, unknown commands, and execution failures are reported explicitly.

Files

  • src/command/cluster_cmd.cpp

Updated ClusterCommand::handleMigrate to send the complete RESTORE command through send_command_and_wait, enforce migration timeouts, validate the target response, and delete the source key only after receiving +OK. Connection failures, timeouts, target errors such as -BUSYKEY, and unexpected replies preserve the source key and return errors; completion logging records deletion status.

Files

  • docs/api.md

Documented length-prefixed CacheObject serialization for RESTORE and CLUSTER MIGRATE, including embedded newline and carriage-return preservation, incomplete-frame atomic failure, and the remaining raw 0xC0 delimiter limitation. Added CLUSTER MIGRATE request/reply and transactional move semantics, including source preservation on target or connection errors and the updated deserialization error message.

Files

  • test/e2e_test/e2e_cluster_full_test.py

Added and integrated test_migrate to verify that MIGRATE returns OK, removes the source key, and makes the value available on the destination node, including coverage for RESTORE argument handling and source-key deletion.

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.

}

// 目标确认收下了,才动源键(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.

Comment thread docs/api.md
Comment on lines +214 to +215
- 配套序列化由 `CacheObject::serialize()` 提供:一行类型标签 + 若干条
`<字节数>\n<原始字节>` 记录,因此成员/字段里含 `\n`、`\r` 不会再被截断

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.

Comment on lines +411 to +412
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}")

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.

Comment thread docs/api.md
- 错误返回:`-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.

@zimingttkx
zimingttkx merged commit fff14ed into main Oct 5, 2026
8 checks passed
zimingttkx added a commit that referenced this pull request Oct 5, 2026
最终审查时抓到的三处,都是 #92 带进来的:

1. timeout 没有上限。#92 之前它被完全忽略,收下 999999999 无害;现在它是
   send_command_and_wait 的等待时长,而命令跑在 SubReactor 的事件循环线程上 ——
   等待期间这条 reactor 上所有连接都停着。一个整数就能把整条 reactor 冻十几天。
   加上 1..60000 的硬范围,并在校验顺序上把它放到"集群没启用"之前,这样参数错误
   在单机模式下也判得出来(契约用例才有地方钉)。
2. std::stoi("5000abc") 返回 5000 且不抛 —— 一个决定阻塞多久的参数不该这么猜。
   换成 from_chars 严格解析,跟 #86/#88 同一套判据。
3. e2e 里那条 MIGRATE 用成了 Redis 顶层语法(多了个 destination-db 的 "0"),
   于是实际传进去的 timeout 是 0:send_command_and_wait 对 <=0 直接失败,三条断言
   会在 Daily 的 e2e 档全红。本项目的签名是
   `CLUSTER MIGRATE host port key timeout [REPLACE]`,改成它。
   docs/api.md §14 与内部命令表里我也照 Redis 语法写过,一并改正,并把
   "等待会占住这条 reactor、所以必须有上限"写进文档而不是只留在代码里。

测试:contract_test.cpp 新增 4 组(gate 标签)—— 超过 60000 / 0 / 负数 / 带尾巴的
"5000abc" 全部报错,且错误文案区分"超出范围"与"不是整数";合法值仍然走到
"cluster mode is not enabled",这条顺序探针证明上面三条是被边界拦下的,
不是被 disabled 抢先返回的。e2e 那三条断言(回 OK、源读不到、目标读得到)保持不变。

本机做实:两个改动文件的括号净差与 origin/main 一致;e2e 脚本 py_compile 过。
cluster_cmd.cpp 依赖 sys/epoll.h,MinGW 编不了,由 CI 学。
zimingttkx added a commit that referenced this pull request Oct 5, 2026
最终审查时抓到的三处,都是 #92 带进来的:

1. timeout 没有上限。#92 之前它被完全忽略,收下 999999999 无害;现在它是
   send_command_and_wait 的等待时长,而命令跑在 SubReactor 的事件循环线程上 ——
   等待期间这条 reactor 上所有连接都停着。一个整数就能把整条 reactor 冻十几天。
   加上 1..60000 的硬范围,并在校验顺序上把它放到"集群没启用"之前,这样参数错误
   在单机模式下也判得出来(契约用例才有地方钉)。
2. std::stoi("5000abc") 返回 5000 且不抛 —— 一个决定阻塞多久的参数不该这么猜。
   换成 from_chars 严格解析,跟 #86/#88 同一套判据。
3. e2e 里那条 MIGRATE 用成了 Redis 顶层语法(多了个 destination-db 的 "0"),
   于是实际传进去的 timeout 是 0:send_command_and_wait 对 <=0 直接失败,三条断言
   会在 Daily 的 e2e 档全红。本项目的签名是
   `CLUSTER MIGRATE host port key timeout [REPLACE]`,改成它。
   docs/api.md §14 与内部命令表里我也照 Redis 语法写过,一并改正,并把
   "等待会占住这条 reactor、所以必须有上限"写进文档而不是只留在代码里。

测试:contract_test.cpp 新增 4 组(gate 标签)—— 超过 60000 / 0 / 负数 / 带尾巴的
"5000abc" 全部报错,且错误文案区分"超出范围"与"不是整数";合法值仍然走到
"cluster mode is not enabled",这条顺序探针证明上面三条是被边界拦下的,
不是被 disabled 抢先返回的。e2e 那三条断言(回 OK、源读不到、目标读得到)保持不变。

本机做实:两个改动文件的括号净差与 origin/main 一致;e2e 脚本 py_compile 过。
cluster_cmd.cpp 依赖 sys/epoll.h,MinGW 编不了,由 CI 学。
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants