Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
59 changes: 35 additions & 24 deletions src/core/nixl_agent.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1222,44 +1222,49 @@ nixlAgent::getLocalPartialMD(const nixl_reg_dlist_t &descs,
nixl_status_t
nixlAgent::loadRemoteMD (const nixl_blob_t &remote_metadata,
std::string &agent_name) {
int count = 0;
nixlSerDes sd;
size_t conn_cnt;
nixl_blob_t conn_info;
nixl_backend_t nixl_backend;
nixlBackendEngine* eng;
nixl_status_t ret;

NIXL_LOCK_GUARD(data->lock);
ret = sd.importStr(remote_metadata);
if(ret)
if (ret != NIXL_SUCCESS) {
return ret;
}

std::string remote_agent = sd.getStr("Agent");
if (remote_agent.size() == 0)
if (remote_agent.empty()) {
return NIXL_ERR_MISMATCH;
}

if (remote_agent == data->name)
if (remote_agent == data->name) {
return NIXL_ERR_INVALID_PARAM;
}

NIXL_DEBUG << "Loading remote metadata for agent: " << remote_agent;

size_t conn_cnt;
ret = sd.getBuf("Conns", &conn_cnt, sizeof(conn_cnt));
if(ret) {
NIXL_ERROR << "Error getting connection count: " << nixlEnumStrings::statusStr(ret);
return ret;
}

for (size_t i=0; i<conn_cnt; ++i) {
int count = 0;
for (size_t i = 0; i < conn_cnt; ++i) {
nixl_backend = sd.getStr("t");
if (nixl_backend.size() == 0)
if (nixl_backend.empty()) {
return NIXL_ERR_MISMATCH;
}

conn_info = sd.getStr("c");
if (conn_info.size() == 0)
if (conn_info.empty()) {
return NIXL_ERR_MISMATCH;
}

// Current agent might not support a remote backend
if (data->backendEngines.count(nixl_backend)!=0) {
if (data->backendEngines.count(nixl_backend) != 0) {

// No need to reload same conn info, error if it changed
if (data->remoteBackends.count(remote_agent) != 0 &&
Expand All @@ -1270,11 +1275,13 @@ nixlAgent::loadRemoteMD (const nixl_blob_t &remote_metadata,
continue;
}

eng = data->backendEngines[nixl_backend];
nixlBackendEngine *eng = data->backendEngines[nixl_backend];
if (eng->supportsRemote()) {
ret = eng->loadRemoteConnInfo(remote_agent, conn_info);
if (ret)
if (ret != NIXL_SUCCESS) {
return ret; // Error in load
}

count++;
data->remoteBackends[remote_agent].emplace(nixl_backend, conn_info);
} else {
Expand All @@ -1286,21 +1293,22 @@ nixlAgent::loadRemoteMD (const nixl_blob_t &remote_metadata,
}

// No common backend, no point in loading the rest, unexpected
if (count == 0 && conn_cnt > 0)
if ((count == 0) && (conn_cnt > 0)) {
return NIXL_ERR_BACKEND;
}

if (sd.getStr("") != "MemSection")
if (sd.getStr("") != "MemSection") {
return NIXL_ERR_MISMATCH;
}

if (data->remoteSections.count(remote_agent) == 0)
data->remoteSections[remote_agent] = new nixlRemoteSection(
remote_agent);
if (data->remoteSections.count(remote_agent) == 0) {
data->remoteSections[remote_agent] = new nixlRemoteSection(remote_agent);
}

ret = data->remoteSections[remote_agent]->loadRemoteData(&sd,
data->backendEngines);
ret = data->remoteSections[remote_agent]->loadRemoteData(&sd, data->backendEngines);

// TODO: can be more graceful, if just the new MD blob was improper
if (ret) {
if (ret != NIXL_SUCCESS) {
delete data->remoteSections[remote_agent];
data->remoteSections.erase(remote_agent);
data->remoteBackends.erase(remote_agent);
Expand All @@ -1315,19 +1323,22 @@ nixl_status_t
nixlAgent::invalidateRemoteMD(const std::string &remote_agent) {
NIXL_LOCK_GUARD(data->lock);

if (remote_agent == data->name)
if (remote_agent == data->name) {
return NIXL_ERR_INVALID_PARAM;
}

nixl_status_t ret = NIXL_ERR_NOT_FOUND;
if (data->remoteSections.count(remote_agent)!=0) {
if (data->remoteSections.count(remote_agent) != 0) {
delete data->remoteSections[remote_agent];
data->remoteSections.erase(remote_agent);
ret = NIXL_SUCCESS;
}

if (data->remoteBackends.count(remote_agent)!=0) {
for (auto & it: data->remoteBackends[remote_agent])
if (data->remoteBackends.count(remote_agent) != 0) {
for (auto &it : data->remoteBackends[remote_agent]) {
data->backendEngines[it.first]->disconnect(remote_agent);
}

data->remoteBackends.erase(remote_agent);
ret = NIXL_SUCCESS;
}
Expand Down