Skip to content

Commit 3cd7d77

Browse files
authored
Urma event mode (#5)
* feat(coro_rpc): add URMA RDMA transport support
1 parent c1cef74 commit 3cd7d77

61 files changed

Lines changed: 31404 additions & 75 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.
Lines changed: 142 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,142 @@
1+
# URMA CTP 实现计划
2+
3+
## 当前状态
4+
5+
### 已完成
6+
1. `circle_buffer` 提取到 `coro_io/detail/circle_buffer.hpp`
7+
2. `urma_device.hpp` 类型冲突修复(`urma_device_wrapper_t`
8+
3. `urma_socket.hpp` 基础重构
9+
4. JFC/JFR/Jetty 创建流程修正
10+
11+
### 已知问题
12+
1. `urma_query_jetty` 签名不正确
13+
2. `urma_import_jetty` 参数结构错误
14+
3. ASIO 协程兼容性问题
15+
4. Socket wrapper 接口不匹配
16+
17+
---
18+
19+
## URMA CTP API 正确用法
20+
21+
### 1. JFC 创建 (Completion Channel)
22+
```c
23+
urma_jfc_cfg_t jfc_cfg = {};
24+
jfc_cfg.depth = 64;
25+
jfc_cfg.flag.value = 0;
26+
jfc_cfg.jfce = nullptr; // polling mode
27+
jfc_cfg.user_ctx = 0;
28+
urma_jfc_t* jfc = urma_create_jfc(ctx, &jfc_cfg);
29+
```
30+
31+
### 2. JFR 创建 (Receive Queue)
32+
```c
33+
urma_jfr_cfg_t jfr_cfg = {};
34+
jfr_cfg.depth = recv_cnt;
35+
jfr_cfg.flag.value = 0;
36+
jfr_cfg.trans_mode = URMA_TM_RM;
37+
jfr_cfg.max_sge = 1;
38+
jfr_cfg.min_rnr_timer = 12;
39+
jfr_cfg.jfc = jfc;
40+
jfr_cfg.token_value = {};
41+
urma_jfr_t* jfr = urma_create_jfr(ctx, &jfr_cfg);
42+
```
43+
44+
### 3. Jetty 创建 (CTP Mode)
45+
```c
46+
urma_jetty_cfg_t jetty_cfg = {};
47+
jetty_cfg.flag.bs.share_jfr = 1;
48+
jetty_cfg.jfs_cfg.depth = send_cnt + 1;
49+
jetty_cfg.jfs_cfg.flag.value = 0;
50+
jetty_cfg.jfs_cfg.trans_mode = URMA_TM_RM;
51+
jetty_cfg.jfs_cfg.jfc = jfc;
52+
jetty_cfg.jfs_cfg.user_ctx = 0;
53+
jetty_cfg.shared.jfr = jfr;
54+
urma_jetty_t* jetty = urma_create_jetty(ctx, &jetty_cfg);
55+
```
56+
57+
### 4. Jetty ID 获取
58+
```c
59+
// 通过 jetty->jfs_id.id 获取
60+
uint32_t jetty_id = jetty->jfs_id.id;
61+
```
62+
63+
### 5. Segment 注册
64+
```c
65+
urma_seg_cfg_t seg_cfg = {};
66+
seg_cfg.va = (uint64_t)buffer;
67+
seg_cfg.len = buffer_size;
68+
seg_cfg.flag.bs.access = URMA_ACCESS_READ | URMA_ACCESS_WRITE;
69+
seg_cfg.flag.bs.token_policy = URMA_TOKEN_NONE;
70+
urma_target_seg_t* tseg = urma_register_seg(ctx, &seg_cfg);
71+
```
72+
73+
### 6. 导入远端 Jetty
74+
```c
75+
urma_rjetty_t remote = {};
76+
remote.jetty_id.eid = peer_eid;
77+
remote.jetty_id.id = peer_jetty_id;
78+
remote.trans_mode = URMA_TM_RM;
79+
remote.type = URMA_JETTY;
80+
remote.tp_type = URMA_CTP; // or URMA_RTP
81+
urma_target_jetty_t* remote_tjetty = urma_import_jetty(ctx, &remote, nullptr);
82+
```
83+
84+
### 7. 发送数据 (SEND)
85+
```c
86+
urma_sge_t sge = {
87+
.addr = (uint64_t)buffer,
88+
.len = data_len,
89+
.tseg = local_tseg
90+
};
91+
urma_sg_t src = {
92+
.sge = &sge,
93+
.num_sge = 1
94+
};
95+
urma_send_wr_t send_wr = {
96+
.src = src
97+
};
98+
urma_jfs_wr_t jfs_wr = {
99+
.opcode = URMA_OPC_SEND,
100+
.send = send_wr,
101+
.tjetty = remote_tjetty
102+
};
103+
urma_post_jetty_send_wr(jetty, &jfs_wr, &bad_wr);
104+
```
105+
106+
### 8. 轮询完成
107+
```c
108+
urma_cr_t cr[8];
109+
int cnt = urma_poll_jfc(jfc, 8, cr);
110+
for (int i = 0; i < cnt; ++i) {
111+
// cr[i].completion_len - 传输字节数
112+
// cr[i].status - 状态
113+
// cr[i].flag.bs.s_r - 0=send, 1=recv
114+
}
115+
```
116+
117+
---
118+
119+
## 待修复清单
120+
121+
### 高优先级
122+
- [ ] `urma_socket.hpp`: 修复 `urma_query_jetty` 调用
123+
- [ ] `urma_socket.hpp`: 修复 `urma_import_jetty` 参数
124+
- [ ] `urma_socket.hpp`: 正确获取 Jetty ID
125+
- [ ] `socket_wrapper.hpp`: 添加缺失接口
126+
127+
### 中优先级
128+
- [ ] ASIO 协程兼容性问题
129+
- [ ] `await_ready` Future 使用错误
130+
- [ ] `async_read/write` 调用方式
131+
132+
### 低优先级
133+
- [ ] `urma_socket_info` 序列化格式
134+
- [ ] 连接握手协议
135+
136+
---
137+
138+
## 参考文档
139+
140+
- URMA API Guide: `.claude/skills/query-urma-docs/URMA API Guide.ch.md`
141+
- URMA User Guide: `.claude/skills/query-urma-docs/URMA User Guide.ch.md`
142+
- URMA 头文件: `include/ylt/urma/urma_api.h`, `urma_types.h`
Lines changed: 169 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,169 @@
1+
# URMA CTP 实现计划 - 下一步
2+
3+
## 当前状态
4+
5+
### 已完成 ✅
6+
1. `urma_socket.hpp` 重构修复:
7+
- `urma_query_jetty` 调用已移除,直接通过 `jetty_->jetty_id.id` 访问
8+
- `urma_import_jetty` 参数已正确初始化 (`trans_mode`, `type`, `tp_type`, `flag`)
9+
- ASIO 协程兼容性问题已修复 (`async_simple::CurrentExecutor{}`)
10+
- 添加缺失方法: `prepare_accept`, `get_remote_address`, `get_local_address`, `get_remote_qp_num`, `get_local_qp_num`
11+
- 添加 `urma_md5_header`, `urma_md5_first_header` 常量
12+
13+
### 待解决问题 🔴
14+
15+
#### 1. urma_example.cpp 集成问题
16+
17+
**错误1**: `get_global_urma_device()` 不接受参数
18+
```cpp
19+
// 当前错误用法:
20+
coro_io::get_global_urma_device({
21+
.dev_name = "",
22+
.buffer_pool_config = {...},
23+
.eid_index = 0
24+
});
25+
```
26+
27+
**错误2**: `coro_rpc_client::config::socket_config` 是 variant 类型,不支持 URMA
28+
```cpp
29+
// 当前错误用法:
30+
conf.socket_config = urma_config; // urma_socket_t::config_t 不能直接赋值
31+
```
32+
33+
#### 2. 缺失的 API
34+
- `urma_device_wrapper_t::init(const urma_config&)` - 需要支持完整配置
35+
- `coro_rpc_client::config` variant 需要 URMA 类型支持
36+
37+
---
38+
39+
## 参考实现: ib_socket 集成模式
40+
41+
### ib_socket 的结构
42+
43+
```cpp
44+
// ib_socket.hpp
45+
class ib_socket_t {
46+
public:
47+
struct config_t {
48+
uint32_t cq_size = 128;
49+
uint16_t recv_buffer_cnt = 8;
50+
uint16_t send_buffer_cnt = 4;
51+
ibv_qp_type qp_type = IBV_QPT_RC;
52+
ibv_qp_cap cap = {...};
53+
std::shared_ptr<ib_device_t> device;
54+
};
55+
56+
struct ib_socket_info {
57+
uint8_t gid[16];
58+
uint16_t lid;
59+
uint32_t buffer_size;
60+
uint32_t qp_num;
61+
};
62+
63+
ib_socket_t(coro_io::ExecutorWrapper<>* executor, const config_t& config);
64+
async_simple::coro::Lazy<std::error_code> connect(const std::string& addr, const std::string& port);
65+
async_simple::coro::Lazy<std::error_code> accept() noexcept;
66+
void prepare_accpet(asio::ip::tcp::socket soc) noexcept;
67+
// ... 其他方法
68+
};
69+
```
70+
71+
### socket_wrapper 的 variant 支持
72+
73+
```cpp
74+
// socket_wrapper.hpp
75+
class socket_wrapper_t {
76+
#ifdef YLT_ENABLE_IBV
77+
std::unique_ptr<ib_socket_t> ib_socket_;
78+
#endif
79+
#ifdef YLT_ENABLE_URMA
80+
std::unique_ptr<urma_socket_t> urma_socket_;
81+
#endif
82+
83+
public:
84+
template<typename T>
85+
auto visit(T&& op) {
86+
#ifdef YLT_ENABLE_IBV
87+
if (ib_socket_) return op(*ib_socket_);
88+
#endif
89+
#ifdef YLT_ENABLE_URMA
90+
if (urma_socket_) return op(*urma_socket_);
91+
#endif
92+
return op(*socket_);
93+
}
94+
};
95+
```
96+
97+
---
98+
99+
## 下一步修改计划
100+
101+
### 阶段1: 修复 urma_example.cpp
102+
103+
#### 1.1 修改 get_global_urma_device() 支持配置参数
104+
105+
参考 ib_device 的初始化模式:
106+
107+
```cpp
108+
// urma_device.hpp
109+
inline bool urma_device_wrapper_t::init(const urma_init_config_t& config) {
110+
// 支持配置参数初始化
111+
}
112+
```
113+
114+
#### 1.2 修改 coro_rpc_client::config 支持 URMA
115+
116+
需要在 variant 中添加 URMA 类型:
117+
118+
```cpp
119+
// coro_rpc_client.hpp
120+
struct config_t {
121+
using socket_config_t = std::variant<
122+
tcp_config,
123+
tcp_with_ssl_config,
124+
tcp_with_ntiles_config
125+
#ifdef YLT_ENABLE_URMA
126+
, urma_socket_t::config_t // 添加 URMA 支持
127+
#endif
128+
>;
129+
socket_config_t socket_config;
130+
};
131+
```
132+
133+
### 阶段2: 实现 URMA 连接的建立流程
134+
135+
#### 2.1 参考 ib_socket 的连接建立
136+
137+
```cpp
138+
// ib_socket 连接流程:
139+
async_simple::coro::Lazy<std::error_code> connect(const std::string& addr, const std::string& port) {
140+
// 1. TCP 连接
141+
// 2. 交换 socket_info (gid, lid, buffer_size, qp_num)
142+
// 3. 初始化 QP
143+
// 4. 返回
144+
}
145+
```
146+
147+
#### 2.2 URMA CTP 模式连接
148+
149+
根据 URMA API Guide,CTP 模式流程:
150+
1. 创建 JFC (Completion Channel)
151+
2. 创建 JFR (Receive Queue)
152+
3. 创建 Jetty (Combined send/recv)
153+
4. 通过 TCP 交换 Jetty ID 和 EID
154+
5. 导入远端 Jetty
155+
6. 建立连接
156+
157+
### 阶段3: 测试验证
158+
159+
1. 单元测试: urma_socket 基本功能
160+
2. 集成测试: urma_example 与 coro_rpc 集成
161+
162+
---
163+
164+
## 参考文档
165+
166+
- URMA API Guide: `.claude/skills/query-urma-docs/URMA API Guide.ch.md`
167+
- URMA User Guide: `.claude/skills/query-urma-docs/URMA User Guide.ch.md`
168+
- ib_socket 实现: `include/ylt/coro_io/ibverbs/ib_socket.hpp`
169+
- socket_wrapper: `include/ylt/coro_io/socket_wrapper.hpp`

.claude/rules/common/agents.md

Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
1+
# Agent Orchestration
2+
3+
## Available Agents
4+
5+
Located in `~/.claude/agents/`:
6+
7+
| Agent | Purpose | When to Use |
8+
|-------|---------|-------------|
9+
| planner | Implementation planning | Complex features, refactoring |
10+
| architect | System design | Architectural decisions |
11+
| tdd-guide | Test-driven development | New features, bug fixes |
12+
| code-reviewer | Code review | After writing code |
13+
| security-reviewer | Security analysis | Before commits |
14+
| build-error-resolver | Fix build errors | When build fails |
15+
| e2e-runner | E2E testing | Critical user flows |
16+
| refactor-cleaner | Dead code cleanup | Code maintenance |
17+
| doc-updater | Documentation | Updating docs |
18+
| rust-reviewer | Rust code review | Rust projects |
19+
20+
## Immediate Agent Usage
21+
22+
No user prompt needed:
23+
1. Complex feature requests - Use **planner** agent
24+
2. Code just written/modified - Use **code-reviewer** agent
25+
3. Bug fix or new feature - Use **tdd-guide** agent
26+
4. Architectural decision - Use **architect** agent
27+
28+
## Parallel Task Execution
29+
30+
ALWAYS use parallel Task execution for independent operations:
31+
32+
```markdown
33+
# GOOD: Parallel execution
34+
Launch 3 agents in parallel:
35+
1. Agent 1: Security analysis of auth module
36+
2. Agent 2: Performance review of cache system
37+
3. Agent 3: Type checking of utilities
38+
39+
# BAD: Sequential when unnecessary
40+
First agent 1, then agent 2, then agent 3
41+
```
42+
43+
## Multi-Perspective Analysis
44+
45+
For complex problems, use split role sub-agents:
46+
- Factual reviewer
47+
- Senior engineer
48+
- Security expert
49+
- Consistency reviewer
50+
- Redundancy checker

0 commit comments

Comments
 (0)