Skip to content
Merged
Show file tree
Hide file tree
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
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@
#define URMA_ENDPOINT_H
#include <atomic>
#include <memory>
#include <mutex>
#include <shared_mutex>
#include <string>
#include <thread>
#include <utility>
Expand Down Expand Up @@ -131,6 +133,7 @@ class UrmaContext : public UbContext {
std::vector<urma_target_seg_t*> local_tseg_list_;
std::vector<urma_seg_t*> remote_seg_list_;
std::vector<urma_target_seg_t*> imported_seg_list_;
std::shared_mutex import_tseg_mutex_;

std::vector<UrmaJFR> jfr_list_;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -192,20 +192,23 @@ int UrmaContext::deconstruct() {
}
seg_region_list_.clear();

for (auto& seg : imported_seg_list_) {
int ret = urma_unimport_seg(seg);
if (ret) {
PLOG(ERROR) << "Failed to unimport segment";
{
std::unique_lock<std::shared_mutex> lock(import_tseg_mutex_);
for (auto& seg : imported_seg_list_) {
int ret = urma_unimport_seg(seg);
if (ret) {
PLOG(ERROR) << "Failed to unimport segment";
}
}
}
imported_seg_list_.clear();
imported_seg_list_.clear();

for (auto& seg : remote_seg_list_) {
free(seg);
}
remote_seg_list_.clear();
for (auto& seg : remote_seg_list_) {
free(seg);
}
remote_seg_list_.clear();

import_tseg_map.clear();
import_tseg_map.clear();
}

for (size_t i = 0; i < jfr_list_.size(); i++) {
if (!jfr_list_[i].native) continue;
Expand Down Expand Up @@ -373,21 +376,34 @@ int UrmaContext::doProcessContextEvents() {
}

void* UrmaContext::retrieveRemoteSeg(const std::string& remoteSegmentStr) {
{
std::shared_lock<std::shared_mutex> lock(import_tseg_mutex_);
auto ret = import_tseg_map.find(remoteSegmentStr);
if (ret != import_tseg_map.end()) return ret->second;
}

std::unique_lock<std::shared_mutex> lock(import_tseg_mutex_);
auto ret = import_tseg_map.find(remoteSegmentStr);
if (ret != import_tseg_map.end()) return ret->second;

std::vector<unsigned char> output_buffer;
deserializeBinaryData(remoteSegmentStr, output_buffer);
urma_seg_t* handle;
handle = (urma_seg_t*)malloc(sizeof(urma_seg_t));
auto* handle = static_cast<urma_seg_t*>(malloc(sizeof(urma_seg_t)));
if (!handle) {
LOG(ERROR) << "Allocate remote segment handle failed";
return nullptr;
}
memcpy(handle, output_buffer.data(), sizeof(urma_seg_t));
remote_seg_list_.push_back(handle);

auto import_tseg =
urma_import_seg(urma_context_, handle, &urma_token, 0, import_flag_);
if (import_tseg == NULL) {
LOG(ERROR) << "Import segment Failed With " << remoteSegmentStr;
free(handle);
return nullptr;
}

remote_seg_list_.push_back(handle);
imported_seg_list_.push_back(import_tseg);
import_tseg_map[remoteSegmentStr] = import_tseg;
return import_tseg;
Expand Down
Loading