From 9e83a378733a5211b655fbdb83b51dae6e9a75f0 Mon Sep 17 00:00:00 2001 From: antoniofilipovic Date: Tue, 26 Sep 2023 14:50:01 +0200 Subject: [PATCH] add working tracker --- src/memgraph.cpp | 1 + src/memory/memory_control.cpp | 143 +++++++++++++++------------- src/memory/memory_control.hpp | 59 ++++++++++++ src/memory/new_delete.cpp | 10 +- src/query/interpreter.cpp | 2 + src/storage/v2/inmemory/storage.cpp | 1 + src/utils/memory_tracker.cpp | 20 ++-- 7 files changed, 156 insertions(+), 80 deletions(-) diff --git a/src/memgraph.cpp b/src/memgraph.cpp index 60ebbbebd..33bbde55f 100644 --- a/src/memgraph.cpp +++ b/src/memgraph.cpp @@ -89,6 +89,7 @@ void InitSignalHandlers(const std::function &shutdown_fun) { auto shutdown = [shutdown_fun]() { if (is_shutting_down) return; is_shutting_down = 1; + // memgraph::memory::UnSetHooks(); shutdown_fun(); }; diff --git a/src/memory/memory_control.cpp b/src/memory/memory_control.cpp index 689982070..929449aab 100644 --- a/src/memory/memory_control.cpp +++ b/src/memory/memory_control.cpp @@ -39,66 +39,30 @@ extent_hooks_t *old_hooks = nullptr; extent_alloc_t *old_alloc = nullptr; -// TODO: think how to solve issue of updating memory status if negative, probably we shouldn't care as it will jump to -// positive pretty quickly -// std::mutex m; - -struct ExtentHooksStats { - struct Alloc { - std::atomic commited{0}; - std::atomic uncommited{0}; - }; - - struct Dalloc { - std::atomic commited{0}; - std::atomic uncommited{0}; - std::atomic error{0}; - }; - - struct Destroy { - std::atomic commited{0}; - std::atomic uncommited{0}; - }; - - struct PurgeForced { - std::atomic counter{0}; - }; - - struct PurgeLazy { - std::atomic counter{0}; - }; - - struct Merge { - std::atomic commited{0}; - std::atomic uncommited{0}; - }; - - struct Split { - std::atomic commited{0}; - std::atomic uncommited{0}; - }; - - struct Commit { - std::atomic counter{0}; - }; - - struct Decommit { - std::atomic counter{0}; - }; - - Alloc alloc; - Dalloc dalloc; - Destroy destroy; - PurgeForced purge_forced; - PurgeLazy purge_lazy; - Merge merge; - Split split; - Commit commit; - Decommit decommit; -}; +std::vector arena_allocations{}; +std::vector arena_upper_limit{}; +std::vector tracking_arenas{}; ExtentHooksStats extent_hook_stats; +int GetArenaForThread() { +#if USE_JEMALLOC + unsigned thread_arena{0}; + size_t size_thread_arena = sizeof(thread_arena); + int err = mallctl("thread.arena", &thread_arena, &size_thread_arena, nullptr, 0); + if (err) { + return -1; + } + return static_cast(thread_arena); +#endif + return -1; +} + +void TrackMemoryForThread(int arena_id, size_t size) { + tracking_arenas[arena_id] = true; + arena_upper_limit[arena_id] = arena_allocations[arena_id] + size; +} + void PrintJemallocInternalStats() { bool config_stats{false}; size_t size_of_config_stats = sizeof(config_stats); @@ -200,6 +164,8 @@ void PrintJemallocInternalStats() { } void PrintStats() { + /* + std::cout << "[TOTAL] RAM:" << utils::GetReadableSize(utils::total_memory_tracker.Amount()) << ", VIRT: " << utils::GetReadableSize(utils::total_memory_tracker.AmountVirt()); @@ -247,6 +213,7 @@ void PrintStats() { if (merge_commited || merge_uncommited) { std::cout << "[merge]commited: " << merge_commited << ", uncommited: " << merge_uncommited << std::endl; } + */ } void *my_alloc(extent_hooks_t *extent_hooks, void *new_addr, size_t size, size_t alignment, bool *zero, bool *commit, @@ -260,6 +227,16 @@ void *my_alloc(extent_hooks_t *extent_hooks, void *new_addr, size_t size, size_t // This needs to be before, to trow exception in case of too big alloc if (*commit) { + arena_allocations[arena_ind] += size; + if (tracking_arenas[arena_ind]) { + if (arena_allocations[arena_ind] > arena_upper_limit[arena_ind]) { + throw utils::OutOfMemoryException( + fmt::format("Memory limit exceeded! Attempting to allocate a chunk of {} which would put the current " + "use to {}, while the maximum allowed size for allocation is set to {}.", + utils::GetReadableSize(size), utils::GetReadableSize(arena_allocations[arena_ind]), + utils::GetReadableSize(arena_upper_limit[arena_ind]))); + } + } memgraph::utils::total_memory_tracker.Alloc(static_cast(size)); } else { memgraph::utils::total_memory_tracker.AllocVirt(static_cast(size)); @@ -268,6 +245,7 @@ void *my_alloc(extent_hooks_t *extent_hooks, void *new_addr, size_t size, size_t auto *ptr = original_hooks_vec[arena_ind]->alloc(extent_hooks, new_addr, size, alignment, zero, commit, arena_ind); if (ptr == nullptr) { if (*commit) { + arena_allocations[arena_ind] -= size; memgraph::utils::total_memory_tracker.Free(static_cast(size)); } else { memgraph::utils::total_memory_tracker.FreeVirt(static_cast(size)); @@ -296,7 +274,7 @@ static bool my_dalloc(extent_hooks_t *extent_hooks, void *addr, size_t size, boo // dangerous assert, useful for testing MG_ASSERT(size % 4096 == 0, "Dalloc size not multiple of page size!"); - MG_ASSERT(original_hooks_vec[arena_ind] && original_hooks_vec[arena_ind]->dalloc); + // MG_ASSERT(original_hooks_vec[arena_ind] && original_hooks_vec[arena_ind]->dalloc); auto err = old_hooks->dalloc(extent_hooks, addr, size, committed, arena_ind); if (err) { @@ -306,6 +284,7 @@ static bool my_dalloc(extent_hooks_t *extent_hooks, void *addr, size_t size, boo if (committed) { memgraph::utils::total_memory_tracker.Free(static_cast(size)); + arena_allocations[arena_ind] -= size; extent_hook_stats.dalloc.commited.fetch_add(1, std::memory_order_relaxed); // spdlog::trace(fmt::format("[DALLOC][RAM] memory pages: {} ,equals to: {}", size / 4096UL, // utils::GetReadableSize(size))); @@ -327,6 +306,7 @@ static void my_destroy(extent_hooks_t *extent_hooks, void *addr, size_t size, bo // utils::GetReadableSize(size))); memgraph::utils::total_memory_tracker.Free(static_cast(size)); + arena_allocations[arena_ind] -= size; extent_hook_stats.destroy.commited.fetch_add(1, std::memory_order_relaxed); } else { // spdlog::trace(fmt::format("[DESTROY][VIRT] memory pages: {}, size: {} ", size / 4096UL, @@ -368,16 +348,16 @@ static bool my_decommit(extent_hooks_t *extent_hooks, void *addr, size_t size, s MG_ASSERT(original_hooks_vec[arena_ind] && original_hooks_vec[arena_ind]->alloc); auto err = old_hooks->decommit(extent_hooks, addr, size, offset, length, arena_ind); + extent_hook_stats.decommit.counter.fetch_add(1, std::memory_order_relaxed); if (err) { return err; } + + memgraph::utils::total_memory_tracker.AllocVirt(static_cast(length)); + memgraph::utils::total_memory_tracker.Free(static_cast(length)); // spdlog::trace(fmt::format("[DECOMMIT][RAM] memory pages: {}, size: {} ", length / 4096UL, // utils::GetReadableSize(length))); // TODO: check is this correct behavior - memgraph::utils::total_memory_tracker.AllocVirt(static_cast(length)); - memgraph::utils::total_memory_tracker.Free(static_cast(length)); - - extent_hook_stats.decommit.counter.fetch_add(1, std::memory_order_relaxed); return false; } @@ -389,9 +369,9 @@ static bool my_purge_forced(extent_hooks_t *extent_hooks, void *addr, size_t siz MG_ASSERT(original_hooks_vec[arena_ind] && original_hooks_vec[arena_ind]->purge_forced); auto err = original_hooks_vec[arena_ind]->purge_forced(extent_hooks, addr, size, offset, length, arena_ind); - if (err) { - return err; - } + // if (err) { + // return err; + // } // spdlog::trace(fmt::format("[PURGE F][RAM] memory pages: {}, size: {} ", length / 4096UL, // utils::GetReadableSize(length))); memgraph::utils::total_memory_tracker.Free(static_cast(length)); @@ -476,6 +456,9 @@ void SetHooks() { original_hooks_vec.reserve(narenas + 5); for (int i = 0; i < narenas; i++) { + arena_allocations.emplace_back(0); + arena_upper_limit.emplace_back(0); + tracking_arenas.emplace_back(0); std::string func_name = "arena." + std::to_string(i) + ".extent_hooks"; size_t hooks_len = sizeof(old_hooks); @@ -507,6 +490,36 @@ void SetHooks() { #endif } +// TODO this can be designed if we fail setting hooks to rollback to classic jemalloc tracker +void UnSetHooks() { +#if USE_JEMALLOC + + uint64_t allocated{0}; + uint64_t sz{sizeof(allocated)}; + + sz = sizeof(unsigned); + unsigned narenas{0}; + int err = mallctl("opt.narenas", (void *)&narenas, &sz, nullptr, 0); + + if (err) { + return; + } + + std::cout << narenas << " : n arenas" << std::endl; + + for (int i = 0; i < narenas; i++) { + std::string func_name = "arena." + std::to_string(i) + ".extent_hooks"; + + err = mallctl(func_name.c_str(), nullptr, nullptr, &old_hooks, sizeof(old_hooks)); + + if (err) { + LOG_FATAL("Error setting jemalloc hooks for jemalloc arena {}", i); + } + } + +#endif +} + void PurgeUnusedMemory() { #if USE_JEMALLOC mallctl("arena." STRINGIFY(MALLCTL_ARENAS_ALL) ".purge", nullptr, nullptr, nullptr, 0); diff --git a/src/memory/memory_control.hpp b/src/memory/memory_control.hpp index 34e20ebe9..69e2a0ade 100644 --- a/src/memory/memory_control.hpp +++ b/src/memory/memory_control.hpp @@ -36,9 +36,68 @@ int mallctlWrite(const char *cmd, T in) { void PurgeUnusedMemory(); void SetHooks(); +void UnSetHooks(); void PrintStats(); +int GetArenaForThread(); +void TrackMemoryForThread(int arena_ind, size_t size); +void SetGlobalLimit(size_t size); inline std::atomic allocated_memory{0}; inline std::atomic virtual_allocated_memory{0}; +inline size_t global_limit{0}; + +struct ExtentHooksStats { + struct Alloc { + std::atomic commited{0}; + std::atomic uncommited{0}; + }; + + struct Dalloc { + std::atomic commited{0}; + std::atomic uncommited{0}; + std::atomic error{0}; + }; + + struct Destroy { + std::atomic commited{0}; + std::atomic uncommited{0}; + }; + + struct PurgeForced { + std::atomic counter{0}; + }; + + struct PurgeLazy { + std::atomic counter{0}; + }; + + struct Merge { + std::atomic commited{0}; + std::atomic uncommited{0}; + }; + + struct Split { + std::atomic commited{0}; + std::atomic uncommited{0}; + }; + + struct Commit { + std::atomic counter{0}; + }; + + struct Decommit { + std::atomic counter{0}; + }; + + Alloc alloc; + Dalloc dalloc; + Destroy destroy; + PurgeForced purge_forced; + PurgeLazy purge_lazy; + Merge merge; + Split split; + Commit commit; + Decommit decommit; +}; } // namespace memgraph::memory diff --git a/src/memory/new_delete.cpp b/src/memory/new_delete.cpp index 92f5c6726..a2b251627 100644 --- a/src/memory/new_delete.cpp +++ b/src/memory/new_delete.cpp @@ -91,7 +91,7 @@ void TrackMemory(std::size_t size) { if (size != 0) [[likely]] { size = nallocx(size, 0); } - memgraph::utils::old_jemalloc_total_memory_tracker.Alloc(static_cast(size)); + // memgraph::utils::old_jemalloc_total_memory_tracker.Alloc(static_cast(size)); #else memgraph::utils::total_memory_tracker.Alloc(static_cast(size)); #endif @@ -102,7 +102,7 @@ void TrackMemory(std::size_t size, const std::align_val_t align) { if (size != 0) [[likely]] { size = nallocx(size, MALLOCX_ALIGN(align)); // NOLINT(hicpp-signed-bitwise) } - memgraph::utils::old_jemalloc_total_memory_tracker.Alloc(static_cast(size)); + // memgraph::utils::old_jemalloc_total_memory_tracker.Alloc(static_cast(size)); #else memgraph::utils::total_memory_tracker.Alloc(static_cast(size)); #endif @@ -132,7 +132,7 @@ void UntrackMemory([[maybe_unused]] void *ptr, [[maybe_unused]] std::size_t size try { #if USE_JEMALLOC if (ptr != nullptr) [[likely]] { - memgraph::utils::old_jemalloc_total_memory_tracker.Free(sallocx(ptr, 0)); + // memgraph::utils::old_jemalloc_total_memory_tracker.Free(sallocx(ptr, 0)); } #else if (size) { @@ -150,8 +150,8 @@ void UntrackMemory(void *ptr, const std::align_val_t align, [[maybe_unused]] std try { #if USE_JEMALLOC if (ptr != nullptr) [[likely]] { - memgraph::utils::old_jemalloc_total_memory_tracker.Free( - sallocx(ptr, MALLOCX_ALIGN(align))); // NOLINT(hicpp-signed-bitwise) + // memgraph::utils::old_jemalloc_total_memory_tracker.Free( + // sallocx(ptr, MALLOCX_ALIGN(align))); // NOLINT(hicpp-signed-bitwise) } #else if (size) { diff --git a/src/query/interpreter.cpp b/src/query/interpreter.cpp index 9216f1e5b..090604dff 100644 --- a/src/query/interpreter.cpp +++ b/src/query/interpreter.cpp @@ -1250,6 +1250,8 @@ std::optional PullPlan::Pull(AnyStream *strea std::optional maybe_limited_resource; if (memory_limit_) { maybe_limited_resource.emplace(&*pool_memory, *memory_limit_); + + memgraph::memory::TrackMemoryForThread(memgraph::memory::GetArenaForThread(), *memory_limit_); ctx_.evaluation_context.memory = &*maybe_limited_resource; } else { ctx_.evaluation_context.memory = &*pool_memory; diff --git a/src/storage/v2/inmemory/storage.cpp b/src/storage/v2/inmemory/storage.cpp index d26b75b9b..1322b832a 100644 --- a/src/storage/v2/inmemory/storage.cpp +++ b/src/storage/v2/inmemory/storage.cpp @@ -278,6 +278,7 @@ Result> InMemoryStorage::InMemoryAccessor::DeleteV Result>>> InMemoryStorage::InMemoryAccessor::DetachDeleteVertex(VertexAccessor *vertex) { + utils::MemoryTracker::OutOfMemoryExceptionEnabler oom_exception; using ReturnType = std::pair>; MG_ASSERT(vertex->transaction_ == &transaction_, diff --git a/src/utils/memory_tracker.cpp b/src/utils/memory_tracker.cpp index a5aa09a4d..bdbc4081e 100644 --- a/src/utils/memory_tracker.cpp +++ b/src/utils/memory_tracker.cpp @@ -104,20 +104,20 @@ void MemoryTracker::Alloc(const int64_t size) { const int64_t will_be = size + amount_.fetch_add(size, std::memory_order_relaxed); - // const auto current_hard_limit = hard_limit_.load(std::memory_order_relaxed); + const auto current_hard_limit = hard_limit_.load(std::memory_order_relaxed); - // if (UNLIKELY(current_hard_limit && will_be > current_hard_limit && MemoryTrackerCanThrow())) { - // MemoryTracker::OutOfMemoryExceptionBlocker exception_blocker; + if (UNLIKELY(current_hard_limit && will_be > current_hard_limit && MemoryTrackerCanThrow())) { + MemoryTracker::OutOfMemoryExceptionBlocker exception_blocker; - // amount_.fetch_sub(size, std::memory_order_relaxed); + amount_.fetch_sub(size, std::memory_order_relaxed); - // throw OutOfMemoryException( - // fmt::format("Memory limit exceeded! Attempting to allocate a chunk of {} which would put the current " - // "use to {}, while the maximum allowed size for allocation is set to {}.", - // GetReadableSize(size), GetReadableSize(will_be), GetReadableSize(current_hard_limit))); - // } + throw OutOfMemoryException( + fmt::format("Memory limit exceeded! Attempting to allocate a chunk of {} which would put the current " + "use to {}, while the maximum allowed size for allocation is set to {}.", + GetReadableSize(size), GetReadableSize(will_be), GetReadableSize(current_hard_limit))); + } - // UpdatePeak(will_be); + UpdatePeak(will_be); } void MemoryTracker::Free(const int64_t size) { amount_.fetch_sub(size, std::memory_order_relaxed); }