From f29ab8aa139909ff8ef503b8636a2ddfc1896ccf Mon Sep 17 00:00:00 2001 From: Dmytro Bogovych Date: Fri, 14 Aug 2026 11:32:17 +0300 Subject: [PATCH] - memory optimization for vq-core project --- src/engine/helper/HL_PoolAllocator.h | 193 +++++++++++++++++++++++---- 1 file changed, 166 insertions(+), 27 deletions(-) diff --git a/src/engine/helper/HL_PoolAllocator.h b/src/engine/helper/HL_PoolAllocator.h index f2ba70d8..02d1c765 100644 --- a/src/engine/helper/HL_PoolAllocator.h +++ b/src/engine/helper/HL_PoolAllocator.h @@ -26,12 +26,26 @@ namespace hl { #if HL_RTP_POOL /// @class FixedBlockPool - /// A process-wide, fixed-block pool with a lock-free thread-local fast path. Identical in - /// design to the pcpp Layer pool: uniform 256-byte blocks carved from 64 KB chunks, an - /// intrusive thread-local free list, and a per-block header tag so deallocate() is O(1) and - /// lock-free for any block (and can tell pooled blocks from the global-allocator fallback used - /// for oversized requests) regardless of the freeing thread. Uniform block size makes a block - /// allocated on one thread safe to free on another (it joins the freeing thread's free list). + /// A process-wide, fixed-block pool with a lock-free thread-local fast path: uniform 256-byte + /// blocks carved from 64 KB chunks, an intrusive thread-local free list, and a per-block header + /// tag so deallocate() is O(1) and lock-free for any block (and can tell pooled blocks from the + /// global-allocator fallback used for oversized requests) regardless of the freeing thread. + /// + /// CROSS-THREAD RECLAIM (2026-07-26) — fixes an unbounded capacity ratchet. A uniform block size + /// makes a block allocated on one thread safe to free on another, but "safe" is not "reclaimed". + /// With a purely thread-local free list the block joins the FREEING thread's list, so a + /// producer/consumer split strands it: vq-core allocates RTPPacket / jitter-buffer nodes on the + /// capture threads and frees them on the interval/reap side, so the capture threads' lists were + /// perpetually empty and refill() carved a fresh chunk every BlocksPerChunk allocations, forever. + /// Observed on vqmonitor (the only box with >1 capture thread): capacityBlocks() 256 -> 11.8M + /// = 2.7 GiB over 4 days at ~25 MiB/h, while activeBlocks() drained to 0 every idle interval. + /// + /// The fix is a process-wide depot of free-block batches. A thread whose free list grows past + /// ThreadCacheHighWater hands a batch down to the depot; a thread that runs dry takes a batch + /// from the depot before carving a new chunk; a thread flushes its remainder to the depot on + /// exit. The depot is touched once per BlocksPerChunk operations and the hysteresis between + /// "empty" and the high-water keeps it off the per-packet path, so the fast path stays lock-free. + /// capacityBlocks() now plateaus at the true concurrent high-water instead of ratcheting. class FixedBlockPool { public: @@ -39,6 +53,10 @@ namespace hl /// shared_ptr nodes we pool (control block + RTPPacket / RtpBuffer::Packet, ~90-130 bytes). static constexpr std::size_t PayloadSize = 240; static constexpr std::size_t BlocksPerChunk = 256; + /// Free blocks a thread keeps to itself before handing a BlocksPerChunk batch to the depot. + /// Must exceed BlocksPerChunk so a thread that churns around the boundary does not bounce + /// on the depot mutex; 2x leaves a full batch in hand after every release. + static constexpr std::size_t ThreadCacheHighWater = 2 * BlocksPerChunk; static void* allocate(std::size_t size) { @@ -50,12 +68,24 @@ namespace hl return raw + HeaderSize; } - void*& head = freeListHead(); - if (head == nullptr) - head = registry().refill(); + ThreadCache& tc = cache(); + if (tc.head == nullptr) + { + // Reclaim from the depot before asking the OS for more memory. Only when the + // depot is dry as well has the pool genuinely outgrown its current capacity. + std::size_t taken = 0; + tc.head = depot().takeBatch(taken); + if (tc.head == nullptr) + { + tc.head = registry().refill(); + taken = BlocksPerChunk; + } + tc.count = taken; + } - uint8_t* block = static_cast(head); - head = nextOf(block); + uint8_t* block = static_cast(tc.head); + tc.head = nextOf(block); + --tc.count; s_activeBlocks.fetch_add(1, std::memory_order_relaxed); return block + HeaderSize; } @@ -68,10 +98,16 @@ namespace hl uint8_t* block = static_cast(ptr) - HeaderSize; if (tagOf(block) == TagPool) { - void*& head = freeListHead(); - nextOf(block) = head; - head = block; + ThreadCache& tc = cache(); + nextOf(block) = tc.head; + tc.head = block; + ++tc.count; s_activeBlocks.fetch_sub(1, std::memory_order_relaxed); + + // The one line that closes the ratchet: surplus goes back to the depot, where + // the thread that actually allocates can reach it. + if (tc.count >= ThreadCacheHighWater) + releaseBatch(tc, BlocksPerChunk); } else { @@ -98,10 +134,19 @@ namespace hl return s_activeGlobal.load(std::memory_order_relaxed); } /// Total blocks ever carved (chunks x BlocksPerChunk) = high-water capacity. + /// With cross-thread reclaim this PLATEAUS once the pool has covered the + /// concurrent high-water; sustained linear growth means reclaim is not working. static std::int64_t capacityBlocks() noexcept { return s_chunks.load(std::memory_order_relaxed) * static_cast(BlocksPerChunk); } + /// Free blocks parked in the process-wide depot, i.e. reclaimed from a freeing + /// thread and available to any allocating thread. Healthy steady state is a + /// non-zero, bounded value: it is the surplus that used to be stranded. + static std::int64_t depotBlocks() noexcept + { + return s_depotBlocks.load(std::memory_order_relaxed); + } /// Approximate live pooled payload bytes (excludes per-block header). static std::int64_t activeBytes() noexcept { @@ -113,6 +158,7 @@ namespace hl inline static std::atomic s_activeBlocks{0}; inline static std::atomic s_activeGlobal{0}; inline static std::atomic s_chunks{0}; + inline static std::atomic s_depotBlocks{0}; static constexpr std::size_t HeaderSize = alignof(std::max_align_t) >= sizeof(uint64_t) ? alignof(std::max_align_t) : sizeof(uint64_t); @@ -121,6 +167,13 @@ namespace hl static constexpr uint64_t TagPool = 0x504F4F4C52545008ULL; // "POOLRTP\b" static constexpr uint64_t TagGlobal = 0x474C4F42524C0808ULL; // "GLOBRL\b\b" + /// A free block's payload is dead space, so the list links live there. Slot 0 chains + /// blocks within a list/batch; slots 1-2 are used only on a batch's head block, to chain + /// batches inside the depot and to carry the batch length (batches from a thread-exit + /// flush are shorter than BlocksPerChunk). + static_assert(PayloadSize >= 2 * sizeof(void*) + sizeof(std::size_t), + "pooled block payload must hold the free-list, depot and length slots"); + static uint64_t& tagOf(void* block) noexcept { return *reinterpret_cast(block); @@ -131,23 +184,95 @@ namespace hl return *reinterpret_cast(static_cast(block) + HeaderSize); } - static void*& freeListHead() noexcept + static void*& batchNextOf(void* block) noexcept { - static thread_local void* head = nullptr; - return head; + return *reinterpret_cast(static_cast(block) + HeaderSize + sizeof(void*)); } + static std::size_t& batchLenOf(void* block) noexcept + { + return *reinterpret_cast(static_cast(block) + HeaderSize + 2 * sizeof(void*)); + } + + /// Per-thread free list. The destructor hands the remainder to the depot so a thread that + /// exits does not strand its blocks — the same ratchet as the cross-thread case, just + /// triggered once per thread instead of continuously. + struct ThreadCache + { + void* head = nullptr; + std::size_t count = 0; + + ~ThreadCache() + { + while (count > 0 && head != nullptr) + releaseBatch(*this, count); + } + }; + + static ThreadCache& cache() noexcept + { + static thread_local ThreadCache tc; + return tc; + } + + /// Detach the first `blocks` entries of the thread list and park them in the depot. + static void releaseBatch(ThreadCache& tc, std::size_t blocks) noexcept + { + if (blocks == 0 || tc.head == nullptr) + return; + + void* batchHead = tc.head; + void* batchTail = batchHead; + std::size_t length = 1; + for (; length < blocks && nextOf(batchTail) != nullptr; ++length) + batchTail = nextOf(batchTail); + + tc.head = nextOf(batchTail); + tc.count -= length; + nextOf(batchTail) = nullptr; + + batchLenOf(batchHead) = length; + depot().giveBatch(batchHead, length); + } + + /// Process-wide store of reclaimed batches. Intrusive (the links live in the free blocks + /// themselves) so neither path can allocate or throw — deallocate() is noexcept. The mutex + /// is taken once per BlocksPerChunk operations, not per packet. + class Depot + { + public: + void* takeBatch(std::size_t& lengthOut) noexcept + { + std::lock_guard lock(m_Mutex); + void* batch = m_Head; + if (batch == nullptr) + { + lengthOut = 0; + return nullptr; + } + m_Head = batchNextOf(batch); + lengthOut = batchLenOf(batch); + batchNextOf(batch) = nullptr; + s_depotBlocks.fetch_sub(static_cast(lengthOut), std::memory_order_relaxed); + return batch; + } + + void giveBatch(void* batch, std::size_t length) noexcept + { + std::lock_guard lock(m_Mutex); + batchNextOf(batch) = m_Head; + m_Head = batch; + s_depotBlocks.fetch_add(static_cast(length), std::memory_order_relaxed); + } + + private: + std::mutex m_Mutex; + void* m_Head = nullptr; + }; + class ChunkRegistry { public: - ~ChunkRegistry() - { - std::lock_guard lock(m_Mutex); - for (uint8_t* chunk : m_Chunks) - ::operator delete(chunk); - m_Chunks.clear(); - } - void* refill() { const std::size_t chunkBytes = BlockSize * BlocksPerChunk; @@ -175,10 +300,22 @@ namespace hl std::vector m_Chunks; }; + /// Both singletons are deliberately immortal (leaked at exit, reclaimed by the OS). + /// A thread_local ThreadCache is destroyed at thread exit — for the main thread that is + /// during static destruction — and its destructor touches the depot, so a depot with a + /// destructor could be used after being destroyed. The same applies to the chunk memory: + /// any pooled object outliving the registry would free into deleted chunks. Never + /// destroying them removes both hazards; the process is exiting either way. + static Depot& depot() noexcept + { + static Depot* instance = new Depot(); + return *instance; + } + static ChunkRegistry& registry() { - static ChunkRegistry instance; - return instance; + static ChunkRegistry* instance = new ChunkRegistry(); + return *instance; } }; #endif // HL_RTP_POOL @@ -189,11 +326,13 @@ namespace hl #if HL_RTP_POOL inline std::int64_t poolActiveBlocks() noexcept { return FixedBlockPool::activeBlocks(); } inline std::int64_t poolCapacityBlocks() noexcept { return FixedBlockPool::capacityBlocks(); } + inline std::int64_t poolDepotBlocks() noexcept { return FixedBlockPool::depotBlocks(); } inline std::int64_t poolActiveGlobal() noexcept { return FixedBlockPool::activeGlobalAllocations(); } inline std::int64_t poolActiveBytes() noexcept { return FixedBlockPool::activeBytes(); } #else inline std::int64_t poolActiveBlocks() noexcept { return -1; } inline std::int64_t poolCapacityBlocks() noexcept { return -1; } + inline std::int64_t poolDepotBlocks() noexcept { return -1; } inline std::int64_t poolActiveGlobal() noexcept { return -1; } inline std::int64_t poolActiveBytes() noexcept { return -1; } #endif