Parse a RdmaMessage according to the pre-defined format Args: rm: the message structure where the parsed message will be saved buffer: the place where the raw message is stored Returns: None
| 1358 | // Returns: |
| 1359 | // None |
| 1360 | void RdmaMessage::ParseMessage(RdmaMessage& rm, void* buffer) { |
| 1361 | char* message = static_cast<char*>(buffer); |
| 1362 | // type |
| 1363 | rm.type_ = static_cast<RdmaMessageType>(message[kTypeStartIndex]); |
| 1364 | // request index |
| 1365 | memcpy(&rm.request_index_, &message[kRequestIndexStartIndex], |
| 1366 | sizeof(rm.request_index_)); |
| 1367 | // name, step_id, remote_addr, rkey |
| 1368 | if ((rm.type_ == RDMA_MESSAGE_TENSOR_REQUEST) || |
| 1369 | (rm.type_ == RDMA_MESSAGE_TENSOR_RE_REQUEST)) { |
| 1370 | memcpy(&rm.name_size_, &message[kNameSizeStartIndex], |
| 1371 | sizeof(rm.name_size_)); |
| 1372 | rm.name_ = string(&message[kNameStartIndex], rm.name_size_); |
| 1373 | memcpy(&rm.remote_addr_, &message[kRemoteAddrStartIndex], |
| 1374 | sizeof(rm.remote_addr_)); |
| 1375 | memcpy(&rm.rkey_, &message[kRkeyStartIndex], sizeof(rm.rkey_)); |
| 1376 | memcpy(&rm.step_id_, &message[kStepIdStartIndex], sizeof(rm.step_id_)); |
| 1377 | } |
| 1378 | // data_type, tensor_bytes, tensor_shape, is_dead |
| 1379 | if ((rm.type_ == RDMA_MESSAGE_TENSOR_REQUEST) || |
| 1380 | (rm.type_ == RDMA_MESSAGE_META_DATA_UPDATE) || |
| 1381 | (rm.type_ == RDMA_MESSAGE_TENSOR_RE_REQUEST)) { |
| 1382 | memcpy(&rm.is_dead_, &message[kIsDeadStartIndex], sizeof(rm.is_dead_)); |
| 1383 | memcpy(&rm.data_type_, &message[kDataTypeStartIndex], |
| 1384 | sizeof(rm.data_type_)); |
| 1385 | memcpy(&rm.tensor_shape_, &message[kTensorShapeStartIndex], |
| 1386 | sizeof(rm.tensor_shape_)); |
| 1387 | memcpy(&rm.tensor_bytes_, &message[kTensorBytesStartIndex], |
| 1388 | sizeof(rm.tensor_bytes_)); |
| 1389 | } |
| 1390 | // checksum |
| 1391 | #ifdef RDMA_DATA_VALIDATION |
| 1392 | memcpy(&rm.checksum_, &message[kChecksumStartIndex], sizeof(rm.checksum_)); |
| 1393 | #endif |
| 1394 | // error status |
| 1395 | if (rm.type_ == RDMA_MESSAGE_ERROR_STATUS) { |
| 1396 | ErrorStatusProto gsProto; |
| 1397 | uint32_t gsProtoSize = *(uint32_t*)&message[kErrorStatusStartIndex]; |
| 1398 | CHECK(ParseProtoUnlimited(&gsProto, &message[kErrorStatusStartIndex + 4], |
| 1399 | gsProtoSize)) |
| 1400 | << "Failed to parse error status proto from message. Aborting."; |
| 1401 | ::grpc::Status gs((::grpc::StatusCode)gsProto.error_code(), |
| 1402 | gsProto.error_message(), gsProto.error_details()); |
| 1403 | rm.status_ = FromGrpcStatus(gs); |
| 1404 | } |
| 1405 | } |
| 1406 | |
| 1407 | //***************************************************************************** |
| 1408 | // RdmaMemoryMgr |
nothing calls this directly
no test coverage detected