Skip to content

Commit 0beba7a

Browse files
authored
[CICD] Add P2P engine perf benchmark (one-sided read/write) (flagos-ai#498)
1 parent fe583c3 commit 0beba7a

2 files changed

Lines changed: 310 additions & 0 deletions

File tree

.github/workflows/test.yml

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -133,3 +133,9 @@ jobs:
133133
# export LD_LIBRARY_PATH=/__w/FlagCX/FlagCX/build/lib:$LD_LIBRARY_PATH
134134
# mpirun -np 8 --allow-run-as-root \
135135
# $DEVICE_API_BIN/perf_internode_onesided -b 1M -e 64M -f 2 -R 2
136+
137+
- name: "P2P Engine perf (one-sided read/write)"
138+
run: |
139+
export PATH=$MPI_HOME/bin:$PATH
140+
export LD_LIBRARY_PATH=/__w/FlagCX/FlagCX/build/lib:$LD_LIBRARY_PATH
141+
mpirun -np 2 --allow-run-as-root $PERF_BIN/perf_p2p_engine -b 4K -e 64M -f 2 -n 10
Lines changed: 304 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,304 @@
1+
/*************************************************************************
2+
* Copyright (c) 2026 BAAI. All rights reserved.
3+
*
4+
* P2P Engine performance benchmark.
5+
*
6+
* Exercises the FlagCX P2P Engine one-sided RDMA APIs directly:
7+
* - flagcxP2pEngineRead (RDMA GET)
8+
* - flagcxP2pEngineWrite (RDMA PUT)
9+
*
10+
* Uses two MPI ranks with the RPC control-plane path:
11+
* Rank 0 = server (target), Rank 1 = client (initiator).
12+
* Both register GPU buffers, start RPC servers, connect, then the
13+
* client initiates reads/writes against the server's buffer.
14+
*
15+
* Usage:
16+
* mpirun -np 2 perf_p2p_engine -b 4K -e 256M -f 2 -n 20
17+
*
18+
* Environment:
19+
* FLAGCX_P2P_PERF_OP=read|write|both (default: both)
20+
************************************************************************/
21+
22+
#include "flagcx.h"
23+
#include "flagcx_p2p.h"
24+
#include "tools.h"
25+
26+
#include <cassert>
27+
#include <chrono>
28+
#include <cinttypes>
29+
#include <climits>
30+
#include <cstdio>
31+
#include <cstdlib>
32+
#include <cstring>
33+
#include <string>
34+
#include <thread>
35+
#include <unistd.h>
36+
37+
static void fatal(const char *msg, int rank) {
38+
fprintf(stderr, "[rank %d] FATAL: %s\n", rank, msg);
39+
MPI_Abort(MPI_COMM_WORLD, 1);
40+
}
41+
42+
static void fatalIf(bool cond, const char *msg, int rank) {
43+
if (cond)
44+
fatal(msg, rank);
45+
}
46+
47+
enum PerfOp { OP_READ = 1, OP_WRITE = 2, OP_BOTH = 3 };
48+
49+
static PerfOp getOpMode() {
50+
const char *env = getenv("FLAGCX_P2P_PERF_OP");
51+
if (env == nullptr)
52+
return OP_BOTH;
53+
if (strcmp(env, "read") == 0)
54+
return OP_READ;
55+
if (strcmp(env, "write") == 0)
56+
return OP_WRITE;
57+
return OP_BOTH;
58+
}
59+
60+
static bool pollTransferDone(FlagcxP2pConn *conn, uint64_t transferId,
61+
int timeoutMs) {
62+
if (transferId == 0)
63+
return true;
64+
auto deadline =
65+
std::chrono::steady_clock::now() + std::chrono::milliseconds(timeoutMs);
66+
while (std::chrono::steady_clock::now() < deadline) {
67+
if (flagcxP2pEngineXferStatus(conn, transferId))
68+
return true;
69+
std::this_thread::yield();
70+
}
71+
return flagcxP2pEngineXferStatus(conn, transferId);
72+
}
73+
74+
static void printHeader(int rank, const char *opName) {
75+
if (rank == 0) {
76+
printf("\n");
77+
printf("# P2P Engine: %s\n", opName);
78+
printf("#%14s %12s %12s\n", "Size (bytes)", "Latency (us)", "BW (GB/s)");
79+
}
80+
}
81+
82+
static void printResult(int rank, size_t size, double latencyUs,
83+
double bandwidth) {
84+
if (rank == 0) {
85+
printf("%15zu %12.2f %12.3f\n", size, latencyUs, bandwidth);
86+
}
87+
}
88+
89+
static void benchmarkOp(FlagcxP2pConn *conn, FlagcxP2pMr localMr,
90+
void *localBuf, uint64_t remoteVa, size_t maxBytes,
91+
size_t minBytes, int stepFactor, int numWarmupIters,
92+
int numIters, int worldRank, bool isRead) {
93+
const char *opName = isRead ? "READ (RDMA GET)" : "WRITE (RDMA PUT)";
94+
printHeader(worldRank, opName);
95+
96+
// Only rank 1 (client) initiates transfers
97+
const bool isInitiator = (worldRank == 1);
98+
99+
for (size_t size = minBytes; size <= maxBytes; size *= stepFactor) {
100+
if (size > UINT32_MAX) {
101+
if (worldRank == 0)
102+
printf("# Skipping size %zu (exceeds uint32_t desc limit)\n", size);
103+
MPI_Barrier(MPI_COMM_WORLD);
104+
MPI_Barrier(MPI_COMM_WORLD);
105+
continue;
106+
}
107+
108+
FlagcxP2pRdmaDesc desc = {};
109+
if (isInitiator) {
110+
int ret = flagcxP2pEngineMakeDesc(conn, remoteVa, (uint32_t)size, &desc);
111+
fatalIf(ret != 0, "flagcxP2pEngineMakeDesc failed", worldRank);
112+
113+
// Warmup
114+
for (int i = 0; i < numWarmupIters; i++) {
115+
uint64_t transferId = 0;
116+
if (isRead) {
117+
ret = flagcxP2pEngineRead(conn, localMr, localBuf, size, desc,
118+
&transferId);
119+
} else {
120+
ret = flagcxP2pEngineWrite(conn, localMr, localBuf, size, desc,
121+
&transferId);
122+
}
123+
fatalIf(ret != 0, "warmup transfer failed", worldRank);
124+
fatalIf(!pollTransferDone(conn, transferId, 10000),
125+
"warmup transfer timed out", worldRank);
126+
}
127+
}
128+
129+
MPI_Barrier(MPI_COMM_WORLD);
130+
131+
// Timed iterations (only initiator)
132+
double elapsed = 0.0;
133+
if (isInitiator) {
134+
timer tim;
135+
for (int i = 0; i < numIters; i++) {
136+
uint64_t transferId = 0;
137+
int ret;
138+
if (isRead) {
139+
ret = flagcxP2pEngineRead(conn, localMr, localBuf, size, desc,
140+
&transferId);
141+
} else {
142+
ret = flagcxP2pEngineWrite(conn, localMr, localBuf, size, desc,
143+
&transferId);
144+
}
145+
fatalIf(ret != 0, "timed transfer failed", worldRank);
146+
fatalIf(!pollTransferDone(conn, transferId, 10000),
147+
"timed transfer timed out", worldRank);
148+
}
149+
elapsed = tim.elapsed();
150+
}
151+
152+
// Sync timing across ranks
153+
double avgLatency = elapsed / (numIters > 0 ? numIters : 1);
154+
MPI_Allreduce(MPI_IN_PLACE, &avgLatency, 1, MPI_DOUBLE, MPI_MAX,
155+
MPI_COMM_WORLD);
156+
157+
double latencyUs = avgLatency * 1e6;
158+
double bandwidth = (double)size / avgLatency / 1e9;
159+
printResult(worldRank, size, latencyUs, bandwidth);
160+
161+
MPI_Barrier(MPI_COMM_WORLD);
162+
}
163+
}
164+
165+
int main(int argc, char *argv[]) {
166+
parser args(argc, argv);
167+
size_t minBytes = args.getMinBytes();
168+
size_t maxBytes = args.getMaxBytes();
169+
int stepFactor = args.getStepFactor();
170+
int numWarmupIters = args.getWarmupIters();
171+
int numIters = args.getTestIters();
172+
173+
int worldRank = 0, worldSize = 0;
174+
MPI_Init(&argc, &argv);
175+
MPI_Comm_rank(MPI_COMM_WORLD, &worldRank);
176+
MPI_Comm_size(MPI_COMM_WORLD, &worldSize);
177+
178+
if (worldSize != 2) {
179+
if (worldRank == 0)
180+
fprintf(stderr, "perf_p2p_engine requires exactly 2 MPI ranks.\n");
181+
MPI_Finalize();
182+
return 1;
183+
}
184+
185+
PerfOp opMode = getOpMode();
186+
187+
if (stepFactor <= 1) {
188+
if (worldRank == 0)
189+
fprintf(stderr, "perf_p2p_engine: step factor (-f) must be >= 2.\n");
190+
MPI_Finalize();
191+
return 1;
192+
}
193+
194+
// Initialize device handle and set GPU
195+
flagcxDeviceHandle_t devHandle = nullptr;
196+
fatalIf(flagcxDeviceHandleInit(&devHandle) != flagcxSuccess,
197+
"flagcxDeviceHandleInit failed", worldRank);
198+
199+
int nGpu = 0;
200+
devHandle->getDeviceCount(&nGpu);
201+
fatalIf(nGpu <= 0, "No GPU devices found", worldRank);
202+
devHandle->setDevice(worldRank % nGpu);
203+
204+
// Allocate GPU buffer using flagcxMemAlloc (GDR-capable)
205+
void *gpuBuf = nullptr;
206+
fatalIf(flagcxMemAlloc(&gpuBuf, maxBytes) != flagcxSuccess,
207+
"flagcxMemAlloc failed", worldRank);
208+
209+
// Create P2P engine
210+
FlagcxP2pEngine *engine = flagcxP2pEngineCreate();
211+
fatalIf(engine == nullptr, "flagcxP2pEngineCreate failed", worldRank);
212+
213+
// Register the GPU buffer
214+
FlagcxP2pMr mr = 0;
215+
fatalIf(flagcxP2pEngineReg(engine, reinterpret_cast<uintptr_t>(gpuBuf),
216+
maxBytes, mr) != 0,
217+
"flagcxP2pEngineReg failed", worldRank);
218+
219+
// Start RPC server
220+
fatalIf(flagcxP2pEngineStartRpcServer(engine) != 0,
221+
"flagcxP2pEngineStartRpcServer failed", worldRank);
222+
223+
// Get RPC port
224+
int rpcPort = flagcxP2pEngineGetRpcPort(engine);
225+
fatalIf(rpcPort < 0, "flagcxP2pEngineGetRpcPort failed", worldRank);
226+
227+
// Get metadata to extract local IP
228+
char *metadataRaw = nullptr;
229+
fatalIf(flagcxP2pEngineGetMetadata(engine, &metadataRaw) != 0,
230+
"flagcxP2pEngineGetMetadata failed", worldRank);
231+
232+
// Parse IP from metadata format "ip:rdma_port?gpu_index?notif_port"
233+
std::string metaStr(metadataRaw);
234+
delete[] metadataRaw;
235+
size_t firstSep = metaStr.find('?');
236+
std::string endpoint = metaStr.substr(0, firstSep);
237+
size_t lastColon = endpoint.rfind(':');
238+
std::string localIp = endpoint.substr(0, lastColon);
239+
240+
// Build session string "ip:rpc_port"
241+
char localSession[256];
242+
snprintf(localSession, sizeof(localSession), "%s:%d", localIp.c_str(),
243+
rpcPort);
244+
245+
// Exchange session strings via MPI
246+
char remoteSession[256] = {};
247+
if (worldRank == 0) {
248+
MPI_Send(localSession, 256, MPI_CHAR, 1, 0, MPI_COMM_WORLD);
249+
MPI_Recv(remoteSession, 256, MPI_CHAR, 1, 0, MPI_COMM_WORLD,
250+
MPI_STATUS_IGNORE);
251+
} else {
252+
MPI_Recv(remoteSession, 256, MPI_CHAR, 0, 0, MPI_COMM_WORLD,
253+
MPI_STATUS_IGNORE);
254+
MPI_Send(localSession, 256, MPI_CHAR, 0, 0, MPI_COMM_WORLD);
255+
}
256+
257+
// Connect to peer
258+
FlagcxP2pConn *conn = flagcxP2pEngineGetConn(engine, remoteSession);
259+
fatalIf(conn == nullptr, "flagcxP2pEngineGetConn failed", worldRank);
260+
261+
// Exchange remote buffer VA
262+
uint64_t localVa = reinterpret_cast<uint64_t>(gpuBuf);
263+
uint64_t remoteVa = 0;
264+
if (worldRank == 0) {
265+
MPI_Send(&localVa, 1, MPI_UINT64_T, 1, 1, MPI_COMM_WORLD);
266+
MPI_Recv(&remoteVa, 1, MPI_UINT64_T, 1, 1, MPI_COMM_WORLD,
267+
MPI_STATUS_IGNORE);
268+
} else {
269+
MPI_Recv(&remoteVa, 1, MPI_UINT64_T, 0, 1, MPI_COMM_WORLD,
270+
MPI_STATUS_IGNORE);
271+
MPI_Send(&localVa, 1, MPI_UINT64_T, 0, 1, MPI_COMM_WORLD);
272+
}
273+
274+
MPI_Barrier(MPI_COMM_WORLD);
275+
276+
if (worldRank == 0) {
277+
printf("P2P Engine Perf Benchmark\n");
278+
printf(" Local session: %s\n", localSession);
279+
printf(" Remote session: %s\n", remoteSession);
280+
printf(" Buffer size: %zu bytes\n", maxBytes);
281+
printf(" Iterations: %d (warmup: %d)\n", numIters, numWarmupIters);
282+
}
283+
284+
// Both ranks run benchmarkOp; only rank 1 (client) initiates transfers,
285+
// rank 0 (server) participates in MPI collectives for timing sync.
286+
if (opMode & OP_READ) {
287+
benchmarkOp(conn, mr, gpuBuf, remoteVa, maxBytes, minBytes, stepFactor,
288+
numWarmupIters, numIters, worldRank, true);
289+
}
290+
if (opMode & OP_WRITE) {
291+
benchmarkOp(conn, mr, gpuBuf, remoteVa, maxBytes, minBytes, stepFactor,
292+
numWarmupIters, numIters, worldRank, false);
293+
}
294+
295+
// Cleanup
296+
MPI_Barrier(MPI_COMM_WORLD);
297+
flagcxP2pEngineMrDestroy(engine, mr);
298+
flagcxP2pEngineDestroy(engine);
299+
flagcxMemFree(gpuBuf);
300+
flagcxDeviceHandleFree(devHandle);
301+
MPI_Finalize();
302+
303+
return 0;
304+
}

0 commit comments

Comments
 (0)