diff --git a/src/snmalloc/ds/pool.h b/src/snmalloc/ds/pool.h index 2a94ec440..61a9ebff6 100644 --- a/src/snmalloc/ds/pool.h +++ b/src/snmalloc/ds/pool.h @@ -80,8 +80,7 @@ namespace snmalloc * * The third template argument is a method to retrieve the actual PoolState. * - * For the pool of allocators, refer to the AllocPool alias defined in - * corealloc.h. + * The allocator pool provides its own construction helper and PoolState. * * For a pool of another type, it is recommended to leave the * third template argument with its default value. The SingletonPoolState @@ -94,8 +93,53 @@ namespace snmalloc PoolState& get_state() = SingletonPoolState::pool> class Pool { + template + static auto call_reinit(U* p, int) -> decltype(p->reinit()) + { + return p->reinit(); + } + + template + static void call_reinit(U*, long) + {} + + template + static auto call_deinit(U* p, int) -> decltype(p->deinit()) + { + return p->deinit(); + } + + template + static void call_deinit(U*, long) + {} + + static void reinit(T* first) + { + T* item = first; + while (item != nullptr) + { + T* next = item->next.unsafe_ptr(); + call_reinit(item, 0); + item = next; + } + } + + static void deinit(T* first, T* last) + { + T* item = first; + while (true) + { + SNMALLOC_ASSERT(item != nullptr); + T* next = item->next.unsafe_ptr(); + call_deinit(item, 0); + if (item == last) + break; + item = next; + } + } + public: - static T* acquire() + static T* try_acquire_front() { PoolState& pool = get_state(); @@ -115,10 +159,19 @@ namespace snmalloc } }); + if (result != nullptr) + call_reinit(result, 0); + return result; + } + + static T* acquire() + { + T* result = try_acquire_front(); if (result != nullptr) return result; auto p = ConstructT::make(); + PoolState& pool = get_state(); with(pool.lock, [&]() { p->list_next = pool.list; @@ -155,6 +208,7 @@ namespace snmalloc pool.front = nullptr; pool.back = nullptr; }); + reinit(result); return result; } @@ -169,6 +223,7 @@ namespace snmalloc static void restore(T* first, T* last) { PoolState& pool = get_state(); + deinit(first, last); last->next = nullptr; with(pool.lock, [&]() { if (pool.front == nullptr) @@ -192,6 +247,7 @@ namespace snmalloc static void restore_front(T* first, T* last) { PoolState& pool = get_state(); + deinit(first, last); last->next = nullptr; with(pool.lock, [&]() { diff --git a/src/snmalloc/ds_core/ptrwrap.h b/src/snmalloc/ds_core/ptrwrap.h index 843eb8d73..b022ff586 100644 --- a/src/snmalloc/ds_core/ptrwrap.h +++ b/src/snmalloc/ds_core/ptrwrap.h @@ -517,6 +517,20 @@ namespace snmalloc return CapPtr::unsafe_from( this->unsafe_capptr.exchange(desired.unsafe_ptr(), order)); } + + SNMALLOC_FAST_PATH bool compare_exchange_strong( + CapPtr& expected, + CapPtr desired, + stl::MemoryOrder success_order, + stl::MemoryOrder failure_order) noexcept + { + auto raw_expected = expected.unsafe_ptr(); + bool result = this->unsafe_capptr.compare_exchange_strong( + raw_expected, desired.unsafe_ptr(), success_order, failure_order); + if (!result) + expected = CapPtr::unsafe_from(raw_expected); + return result; + } }; namespace capptr diff --git a/src/snmalloc/global/globalalloc.h b/src/snmalloc/global/globalalloc.h index 7607e582a..fc22fa478 100644 --- a/src/snmalloc/global/globalalloc.h +++ b/src/snmalloc/global/globalalloc.h @@ -11,10 +11,10 @@ namespace snmalloc static_assert( Config_::Options.AllocIsPoolAllocated, "Global cleanup is available only for pool-allocated configurations"); + // Call this periodically to free and coalesce memory allocated by // allocators that are not currently in use by any thread. - // One atomic operation to extract the stack, another to restore it. - // Handling the message queue for each stack is non-atomic. + // Handling the message queue for each allocator is non-atomic. auto* first = AllocPool::extract(); auto* alloc = first; @@ -43,6 +43,23 @@ namespace snmalloc static_assert( Config_::Options.AllocIsPoolAllocated, "Global status is available only for pool-allocated configurations"); + + // debug_is_empty() calls flush(), which requires an active message queue. + // Extract unused allocators to claim their queues for the whole check. + auto* first = AllocPool::extract(); + auto* last = first; + OnDestruct restore_free_allocators([&first, &last]() { + if (first != nullptr) + AllocPool::restore(first, last); + }); + while (last != nullptr) + { + auto* next = AllocPool::extract(last); + if (next == nullptr) + break; + last = next; + } + // This is a debugging function. It checks that all memory from all // allocators has been freed. auto* alloc = AllocPool::iterate(); diff --git a/src/snmalloc/mem/allocpool_assistance.h b/src/snmalloc/mem/allocpool_assistance.h new file mode 100644 index 000000000..a7475bd16 --- /dev/null +++ b/src/snmalloc/mem/allocpool_assistance.h @@ -0,0 +1,144 @@ +#pragma once + +#include "../pal/pal_consts.h" +#include "snmalloc/stl/atomic.h" + +#include + +#ifndef SNMALLOC_ASSIST_IDLE_MS +# define SNMALLOC_ASSIST_IDLE_MS 1000 +#endif + +namespace snmalloc +{ + template + inline constexpr bool uses_inactive_queue_marker = + Config::Options.AllocIsPoolAllocated && + pal_supports; + + /** + * State and policy for deciding when to assist a disused allocator. + */ + template + class AllocPoolAssistance + { + /** + * Best-effort count of inactive allocators with pending remote frees. + * Queue-state transitions are authoritative; publication ordering may make + * this count temporarily non-positive. + */ + SNMALLOC_REQUIRE_CONSTINIT + inline static stl::Atomic inactive_pending_count{0}; + + /** + * Time from which to measure the next assistance delay. Zero is an + * ordinary timestamp, not a sentinel. + */ + SNMALLOC_REQUIRE_CONSTINIT + inline static stl::Atomic idle_start_ms{0}; + + /** + * Incremented when pool acquisition claims a pending inactive allocator, + * which restarts the assistance delay. + */ + SNMALLOC_REQUIRE_CONSTINIT + inline static stl::Atomic deadline_reset_generation{0}; + + /** + * Last reset generation observed by the scheduling policy. + */ + SNMALLOC_REQUIRE_CONSTINIT + inline static stl::Atomic processed_deadline_reset_generation{0}; + + static bool process_deadline_reset(uint64_t sampled_time) + { + uint64_t processed = + processed_deadline_reset_generation.load(stl::memory_order_acquire); + uint64_t reset = + deadline_reset_generation.load(stl::memory_order_acquire); + if (processed == reset) + return false; + + // This sample only restarts the delay; it does not reserve assistance. + idle_start_ms.store(sampled_time, stl::memory_order_relaxed); + + // A failed CAS means another thread processed this generation. A newer + // generation remains different and will be processed by a later sample. + processed_deadline_reset_generation.compare_exchange_strong( + processed, reset, stl::memory_order_release, stl::memory_order_relaxed); + return true; + } + + public: + /** + * Record an inactive allocator becoming responsible for pending remote + * frees. + */ + static void pending_queue_added() + { + inactive_pending_count.fetch_add(1, stl::memory_order_relaxed); + } + + /** + * Record a pending inactive allocator being claimed from the pool. + */ + static void pending_queue_claimed() + { + inactive_pending_count.fetch_sub(1, stl::memory_order_relaxed); + deadline_reset_generation.fetch_add(1, stl::memory_order_release); + } + + /** + * Debug-only observation of the best-effort number of inactive allocators + * with pending remote frees. + */ + static int64_t debug_pending_count() + { + return inactive_pending_count.load(stl::memory_order_relaxed); + } + + /** + * Return whether this sample should perform one assistance attempt. + * + * A successful return advances the shared pacing origin, ensuring that + * competing threads do not assist for the same interval. + */ + [[nodiscard]] static bool should_assist(uint64_t sampled_time) + { + if constexpr (uses_inactive_queue_marker) + { + if (process_deadline_reset(sampled_time)) + return false; + + int64_t pending = + inactive_pending_count.load(stl::memory_order_acquire); + if (pending <= 0) + return false; + + // More pending queues shorten the interval between assistance attempts. + uint64_t delay = + uint64_t{SNMALLOC_ASSIST_IDLE_MS} / static_cast(pending); + if (delay == 0) + delay = 1; + uint64_t idle_start = idle_start_ms.load(stl::memory_order_relaxed); + + // Unsigned subtraction handles clock wraparound. An older sample may + // cause an extra best-effort attempt, which is harmless. + if ((sampled_time - idle_start) < delay) + return false; + + // Advancing the origin reserves this interval for the winning thread. + return idle_start_ms.compare_exchange_strong( + idle_start, + sampled_time, + stl::memory_order_relaxed, + stl::memory_order_relaxed); + } + else + { + UNUSED(sampled_time); + return false; + } + } + }; +} // namespace snmalloc diff --git a/src/snmalloc/mem/corealloc.h b/src/snmalloc/mem/corealloc.h index a28357cd9..c04b6d39b 100644 --- a/src/snmalloc/mem/corealloc.h +++ b/src/snmalloc/mem/corealloc.h @@ -2,6 +2,7 @@ #include "../ds/ds.h" #include "../ds/pool.h" +#include "allocpool_assistance.h" #include "check_init.h" #include "freelist.h" #include "metadata.h" @@ -118,6 +119,8 @@ namespace snmalloc using BackendSlabMetadata = typename Config::Backend::SlabMetadata; using PagemapEntry = typename Config::PagemapEntry; + static void assist_disused_allocators(uint64_t now_ms); + /// }@ /** @@ -237,6 +240,30 @@ namespace snmalloc return *public_state(); } + public: + void reinit() + { + if constexpr (uses_inactive_queue_marker) + { + if (message_queue().claim_from_inactive()) + { + AllocPoolAssistance::pending_queue_claimed(); + } + } + } + + void deinit() + { + if constexpr (uses_inactive_queue_marker) + { + if (message_queue().release_to_inactive()) + { + AllocPoolAssistance::pending_queue_added(); + } + } + } + + private: /** * Check if this allocator has messages to deallocate blocks from another * thread @@ -298,6 +325,16 @@ namespace snmalloc constexpr Allocator(bool){}; public: + bool debug_is_in_use() + { + return this->in_use.load(stl::memory_order_acquire); + } + + bool debug_has_pending_remote() + { + return !message_queue().is_empty(); + } + /** * Constructor for the case that the core allocator owns the local state. * SFINAE disabled if the allocator does not own the local state. @@ -429,10 +466,13 @@ namespace snmalloc }; auto cb = [this, domesticate, &need_post, &bytes_freed]( capptr::Alloc msg) SNMALLOC_FAST_PATH_LAMBDA { + if (bytes_freed >= static_cast(REMOTE_BATCH_LIMIT)) + return false; + auto& entry = Config::Backend::get_metaentry(snmalloc::address_cast(msg)); handle_dealloc_remote(entry, msg, need_post, domesticate, bytes_freed); - return bytes_freed < REMOTE_BATCH_LIMIT; + return true; }; #ifdef SNMALLOC_TRACING @@ -828,7 +868,8 @@ namespace snmalloc } auto r = finish_alloc(p, size); - return ticker.check_tick(r); + return ticker.check_tick( + r, [](uint64_t now_ms) { assist_disused_allocators(now_ms); }); } return small_refill_slow( sizeclass, fast_free_list, size); @@ -889,7 +930,8 @@ namespace snmalloc } auto r = finish_alloc(p, size); - return ticker.check_tick(r); + return ticker.check_tick( + r, [](uint64_t now_ms) { assist_disused_allocators(now_ms); }); }, [](Allocator* a, size_t size) SNMALLOC_FAST_PATH_LAMBDA { return a->template small_alloc(size); @@ -1386,14 +1428,13 @@ namespace snmalloc } /** - * Flush one message-queue snapshot, cached state, and delayed - * deallocations. Enqueues whose back exchange observes the reset remain for - * a later flush. - * - * Returns true if messages are sent to other threads. + * Flush queued messages, cached state, and delayed deallocations. */ - bool flush() + template + void flush_impl(bool& posted) { + message_queue().assert_not_inactive(); + auto local_state = backend_state_ptr(); auto domesticate = [local_state](freelist::QueuePtr p) SNMALLOC_FAST_PATH_LAMBDA { @@ -1410,6 +1451,7 @@ namespace snmalloc const PagemapEntry& entry = Config::Backend::get_metaentry(snmalloc::address_cast(m)); handle_dealloc_remote(entry, m, need_post, domesticate, bytes_flushed); + return true; }; if constexpr (Config::Options.QueueHeadsAreTame) @@ -1418,11 +1460,26 @@ namespace snmalloc [](freelist::QueuePtr p) SNMALLOC_FAST_PATH_LAMBDA { return freelist::HeadPtr::unsafe_from(p.unsafe_ptr()); }; - message_queue().drain_and_reset(domesticate_first, domesticate, cb); + if constexpr (Dequeue) + { + message_queue().template dequeue( + domesticate_first, domesticate, cb); + } + else + { + message_queue().drain_and_reset(domesticate_first, domesticate, cb); + } } else { - message_queue().drain_and_reset(domesticate, domesticate, cb); + if constexpr (Dequeue) + { + message_queue().template dequeue(domesticate, domesticate, cb); + } + else + { + message_queue().drain_and_reset(domesticate, domesticate, cb); + } } auto& key = freelist::Object::key_root; @@ -1444,7 +1501,7 @@ namespace snmalloc } while (!small_fast_free_lists[i].empty()); } - auto posted = remote_dealloc_cache.template post( + posted = remote_dealloc_cache.template post( local_state, get_trunc_id()); // We may now have unused slabs, return to the global allocator. @@ -1467,10 +1524,24 @@ namespace snmalloc } // Set the remote_dealloc_cache to immediately slow path. remote_dealloc_cache.capacity = 0; + } + bool flush() + { + bool posted = false; + flush_impl(posted); return posted; } + void try_flush() + { + if (message_queue().is_empty()) + return; + + bool posted = false; + flush_impl(posted); + } + /** * If result parameter is non-null, then false is assigned into the * the location pointed to by result if this allocator is non-empty. @@ -1601,10 +1672,28 @@ namespace snmalloc } }; - /** - * Use this alias to access the pool of allocators throughout snmalloc. - */ template using AllocPool = Pool, ConstructAllocator, Config::pool>; + + template + void Allocator::assist_disused_allocators(uint64_t now_ms) + { + if constexpr (uses_inactive_queue_marker) + { + if (!AllocPoolAssistance::should_assist(now_ms)) + return; + + Allocator* alloc = AllocPool::try_acquire_front(); + if (alloc == nullptr) + return; + + OnDestruct restore([alloc]() { AllocPool::release(alloc); }); + alloc->try_flush(); + } + else + { + UNUSED(now_ms); + } + } } // namespace snmalloc diff --git a/src/snmalloc/mem/freelist_queue.h b/src/snmalloc/mem/freelist_queue.h index 33f7fa489..bc1615c75 100644 --- a/src/snmalloc/mem/freelist_queue.h +++ b/src/snmalloc/mem/freelist_queue.h @@ -6,6 +6,13 @@ namespace snmalloc { + enum class EnqueueResult + { + StartedActive, + StartedInactive, + Appended + }; + /** * A FreeListMPSCQ is a chain of freed objects exposed as a MPSC append-only * atomic queue that uses one xchg per append. @@ -55,6 +62,9 @@ namespace snmalloc Domesticator_queue& domesticate, Cb& cb) { + // Read the successor before invoking the callback. If the callback + // rejects curr, curr remains owned by the queue with a published + // successor. while (address_cast(curr) != address_cast(target)) { auto next = curr->atomic_read_next(Key, Key_tweak, domesticate); @@ -63,15 +73,70 @@ namespace snmalloc Aal::prefetch(next.unsafe_ptr()); if (SNMALLOC_UNLIKELY(!cb(curr))) - return next; + return curr; curr = next; } + // The returned node is unprocessed and remains in the queue. return curr; } public: + static freelist::QueuePtr inactive_marker() + { + return freelist::QueuePtr::unsafe_from( + unsafe_from_uintptr>(1)); + } + + bool release_to_inactive() + { + freelist::QueuePtr expected = nullptr; + if (back.compare_exchange_strong( + expected, + inactive_marker(), + stl::memory_order_acq_rel, + stl::memory_order_acquire)) + { + return false; + } + + SNMALLOC_ASSERT(expected != inactive_marker()); + return true; + } + + bool claim_from_inactive() + { + auto expected = back.load(stl::memory_order_acquire); + if (expected == inactive_marker()) + { + if (back.compare_exchange_strong( + expected, + nullptr, + stl::memory_order_acq_rel, + stl::memory_order_acquire)) + { + return false; + } + } + + SNMALLOC_ASSERT(expected != nullptr); + SNMALLOC_ASSERT(expected != inactive_marker()); + return expected != nullptr; + } + + void assert_not_inactive() + { + SNMALLOC_ASSERT( + back.load(stl::memory_order_relaxed) != inactive_marker()); + } + + bool is_empty() + { + auto value = back.load(stl::memory_order_acquire); + return value == nullptr || value == inactive_marker(); + } + /** * Exactly one consumer role may execute this operation. No dequeue, * owning-allocator queue processing, or second drain may run concurrently. @@ -95,6 +160,8 @@ namespace snmalloc Domesticator_queue domesticate_queue, Cb cb) { + assert_not_inactive(); + // After reuse, acquire the release sequence headed by the preceding // reset, so front cannot observe an earlier queue generation. if (back.load(stl::memory_order_acquire) == nullptr) @@ -115,6 +182,7 @@ namespace snmalloc front.store(nullptr, stl::memory_order_relaxed); auto target = back.exchange(nullptr, stl::memory_order_acq_rel); SNMALLOC_ASSERT(target != nullptr); + SNMALLOC_ASSERT(target != inactive_marker()); auto process = [&cb](freelist::HeadPtr p) { cb(p); @@ -133,6 +201,7 @@ namespace snmalloc inline bool can_dequeue() { + assert_not_inactive(); return front.load(stl::memory_order_relaxed) != back.load(stl::memory_order_relaxed); } @@ -144,12 +213,11 @@ namespace snmalloc * The Domesticator here is used only on pointers read from the head. See * the commentary on the class. * - * Returns true if this enqueue observed an empty back and started a new - * queue generation by publishing front. Returns false if it appended to - * an existing chain. + * Reports whether this enqueue started an active or inactive queue + * generation, or appended to an existing chain. */ template - bool enqueue( + EnqueueResult enqueue( freelist::HeadPtr first, freelist::HeadPtr last, Domesticator_head domesticate_head) @@ -176,30 +244,37 @@ namespace snmalloc freelist::QueuePtr prev = back.exchange(capptr_rewild(last), stl::memory_order_acq_rel); - if (SNMALLOC_LIKELY(prev != nullptr)) + if (SNMALLOC_UNLIKELY(prev == nullptr || prev == inactive_marker())) { - // Once this store publishes first, a drain may observe it and release - // prev; this must therefore be this producer's final access to prev. - freelist::Object::atomic_store_next( - domesticate_head(prev), first, Key, Key_tweak); - return false; + // drain_and_reset clears front before resetting back, so only a + // producer whose exchange observed an empty state may publish a + // replacement front. + front.store(capptr_rewild(first)); + return prev == nullptr ? EnqueueResult::StartedActive : + EnqueueResult::StartedInactive; } - // drain_and_reset clears front before resetting back, so only a producer - // whose exchange observed null may publish a replacement front. - front.store(capptr_rewild(first)); - return true; + // Once this store publishes first, a drain may observe it and release + // prev; this must therefore be this producer's final access to prev. + freelist::Object::atomic_store_next( + domesticate_head(prev), first, Key, Key_tweak); + return EnqueueResult::Appended; } /** * Destructively iterate the queue. Each queue element is removed and fed - * to the callback in turn. The callback may return false to stop iteration - * early (but must have processed the element it was given!). + * to the callback in turn. The callback may return false to reject its + * argument and stop iteration early. A rejected object must not have been + * consumed, freed, or re-enqueued; it remains owned by this queue. + * + * Closing dequeue requires the callback to process the retained final + * object after the queue has been closed. * * Takes a domestication callback for each of "pointers read from head" and * "pointers read from queue". See the commentary on the class. */ template< + bool Close = false, typename Domesticator_head, typename Domesticator_queue, typename Cb> @@ -226,8 +301,34 @@ namespace snmalloc * Publishing it to client-accessible front requires it to be considered * Wild again in !QueueHeadsAreTame builds. */ - curr = process_chain(curr, b, domesticate_queue, cb); - front = capptr_rewild(curr); + if constexpr (Close) + { + curr = process_chain(curr, b, domesticate_queue, cb); + + auto current_back = back.load(stl::memory_order_acquire); + if (address_cast(curr) == address_cast(current_back)) + { + front.store(nullptr, stl::memory_order_relaxed); + if (back.compare_exchange_strong( + current_back, + nullptr, + stl::memory_order_acq_rel, + stl::memory_order_acquire)) + { + // A producer that could still access curr would have changed back, + // making the compare-exchange fail. + SNMALLOC_CHECK(cb(curr)); + return; + } + } + + front.store(capptr_rewild(curr), stl::memory_order_release); + } + else + { + curr = process_chain(curr, b, domesticate_queue, cb); + front = capptr_rewild(curr); + } invariant(); } }; diff --git a/src/snmalloc/mem/remoteallocator.h b/src/snmalloc/mem/remoteallocator.h index efd4ba942..1f40026e2 100644 --- a/src/snmalloc/mem/remoteallocator.h +++ b/src/snmalloc/mem/remoteallocator.h @@ -314,6 +314,26 @@ namespace snmalloc using alloc_id_t = address_t; + bool is_empty() + { + return list.is_empty(); + } + + bool release_to_inactive() + { + return list.release_to_inactive(); + } + + bool claim_from_inactive() + { + return list.claim_from_inactive(); + } + + void assert_not_inactive() + { + list.assert_not_inactive(); + } + constexpr RemoteAllocator() = default; void invariant() @@ -352,11 +372,11 @@ namespace snmalloc * The Domesticator here is used only on pointers read from the head. See * the commentary on the class. * - * Returns true if this enqueue started a new queue generation, or false if - * it appended to an existing chain. + * Reports whether this enqueue started an active or inactive queue + * generation, or appended to an existing chain. */ template - bool enqueue( + EnqueueResult enqueue( capptr::Alloc first, capptr::Alloc last, Domesticator_head domesticate_head) @@ -369,13 +389,18 @@ namespace snmalloc /** * Destructively iterate the queue. Each queue element is removed and fed - * to the callback in turn. The callback may return false to stop iteration - * early (but must have processed the element it was given!). + * to the callback in turn. The callback may return false to reject its + * argument and stop iteration early. A rejected message must not have been + * consumed, freed, or re-enqueued; it remains owned by the queue. + * + * Closing dequeue requires the callback to process the retained final + * message after the queue has been closed. * * Takes a domestication callback for each of "pointers read from head" and * "pointers read from queue". See the commentary on the class. */ template< + bool Close = false, typename Domesticator_head, typename Domesticator_queue, typename Cb> @@ -387,7 +412,7 @@ namespace snmalloc auto cbwrap = [cb](freelist::HeadPtr p) SNMALLOC_FAST_PATH_LAMBDA { return cb(RemoteMessage::from_message_link(p)); }; - list.dequeue(domesticate_head, domesticate_queue, cbwrap); + list.template dequeue(domesticate_head, domesticate_queue, cbwrap); } alloc_id_t trunc_id() diff --git a/src/snmalloc/mem/remotecache.h b/src/snmalloc/mem/remotecache.h index c92c99a65..7ffcc2c97 100644 --- a/src/snmalloc/mem/remotecache.h +++ b/src/snmalloc/mem/remotecache.h @@ -1,16 +1,15 @@ #pragma once #include "../ds/ds.h" +#include "allocpool_assistance.h" #include "backend_wrappers.h" #include "freelist.h" #include "metadata.h" #include "remoteallocator.h" #include "snmalloc/stl/array.h" -#include "snmalloc/stl/atomic.h" namespace snmalloc { - /** * Same-destination message batching. * @@ -317,16 +316,29 @@ namespace snmalloc mitigations(sanity_checks), !entry.is_backend_owned(), "Delayed detection of attempt to free internal structure."); + EnqueueResult result; if constexpr (Config::Options.QueueHeadsAreTame) { auto domesticate_nop = [](freelist::QueuePtr p) { return freelist::HeadPtr::unsafe_from(p.unsafe_ptr()); }; - remote->enqueue(first, last, domesticate_nop); + result = remote->enqueue(first, last, domesticate_nop); + } + else + { + result = remote->enqueue(first, last, domesticate); + } + + if constexpr (uses_inactive_queue_marker) + { + if (result == EnqueueResult::StartedInactive) + { + AllocPoolAssistance::pending_queue_added(); + } } else { - remote->enqueue(first, last, domesticate); + UNUSED(result); } sent_something = true; } diff --git a/src/snmalloc/mem/ticker.h b/src/snmalloc/mem/ticker.h index 33e3819fb..21be4bd29 100644 --- a/src/snmalloc/mem/ticker.h +++ b/src/snmalloc/mem/ticker.h @@ -1,6 +1,7 @@ #pragma once #include "../ds_core/ds_core.h" +#include "../pal/pal_consts.h" #include @@ -40,10 +41,11 @@ namespace snmalloc * Slow path that actually queries clock and sets up * how many calls for the next time we hit the slow path. */ - template - SNMALLOC_SLOW_PATH T check_tick_slow(T p = nullptr) noexcept + template + SNMALLOC_SLOW_PATH T check_tick_slow(T p, Callback callback) noexcept { uint64_t now_ms = PAL::time_in_ms(); + callback(now_ms); // Set up clock. if (last_query_ms == 0) @@ -73,6 +75,8 @@ namespace snmalloc auto new_deadline_in_ticks = ((1 + counted) * deadline_in_ms) / duration_ms; + if (new_deadline_in_ticks == 0) + new_deadline_in_ticks = 1; counted = new_deadline_in_ticks; count_down = new_deadline_in_ticks; @@ -82,6 +86,12 @@ namespace snmalloc public: template SNMALLOC_FAST_PATH T check_tick(T p = nullptr) + { + return check_tick(p, [](uint64_t) {}); + } + + template + SNMALLOC_FAST_PATH T check_tick(T p, Callback callback) { if constexpr (pal_supports) { @@ -91,7 +101,7 @@ namespace snmalloc // heart beat. if (--count_down == 0) { - return check_tick_slow(p); + return check_tick_slow(p, callback); } } return p; diff --git a/src/snmalloc/mitigations/allocconfig.h b/src/snmalloc/mitigations/allocconfig.h index 3f326a570..27f13f4aa 100644 --- a/src/snmalloc/mitigations/allocconfig.h +++ b/src/snmalloc/mitigations/allocconfig.h @@ -93,6 +93,9 @@ namespace snmalloc 1 * 1024 * 1024 #endif ; + static_assert( + REMOTE_BATCH_LIMIT > 0, + "SNMALLOC_REMOTE_BATCH_PROCESS_SIZE must be positive"); // Used to configure when the backend should use thread local buddies. // This only basically is used to disable some buddy allocators on small diff --git a/src/test/func/allocpool/allocpool.cc b/src/test/func/allocpool/allocpool.cc new file mode 100644 index 000000000..5e3084944 --- /dev/null +++ b/src/test/func/allocpool/allocpool.cc @@ -0,0 +1,174 @@ +#include +#include +#include +#include + +using namespace snmalloc; + +namespace +{ + constexpr uint64_t idle_constant_ms = uint64_t{SNMALLOC_ASSIST_IDLE_MS}; + static_assert(idle_constant_ms >= 2); + + struct AlternateConfig + { + using Pal = DefaultPal; + + static constexpr Flags Options = []() constexpr { + Flags opts = {}; + opts.IsQueueInline = false; + opts.QueueHeadsAreTame = false; + opts.AllocOwnsLocalState = false; + return opts; + }(); + }; + + struct NonPooledConfig + { + using Pal = DefaultPal; + + static constexpr Flags Options = []() constexpr { + Flags opts = {}; + opts.AllocIsPoolAllocated = false; + return opts; + }(); + }; + + struct SchedulingConfig + { + using Pal = DefaultPal; + + static constexpr Flags Options = []() constexpr { + Flags opts = {}; + opts.AllocIsPoolAllocated = true; + return opts; + }(); + }; + + static_assert(uses_inactive_queue_marker); + static_assert(uses_inactive_queue_marker>); + static_assert( + !uses_inactive_queue_marker>>); + static_assert(!uses_inactive_queue_marker); + static_assert(uses_inactive_queue_marker); + + template + Allocator* last_in_chain(Allocator* first) + { + auto* last = first; + while (last != nullptr) + { + auto* next = AllocPool::extract(last); + if (next == nullptr) + return last; + last = next; + } + return nullptr; + } + + void test_reservation_and_pacing() + { + using Assistance = AllocPoolAssistance; + + SNMALLOC_CHECK(Assistance::debug_pending_count() == 0); + SNMALLOC_CHECK(!Assistance::should_assist(1)); + + Assistance::pending_queue_claimed(); + SNMALLOC_CHECK(Assistance::debug_pending_count() == -1); + SNMALLOC_CHECK(!Assistance::should_assist(1)); + Assistance::pending_queue_added(); + SNMALLOC_CHECK(!Assistance::should_assist(1 + idle_constant_ms)); + + Assistance::pending_queue_added(); + SNMALLOC_CHECK(Assistance::should_assist(1 + idle_constant_ms)); + SNMALLOC_CHECK(!Assistance::should_assist(1 + idle_constant_ms)); + SNMALLOC_CHECK(Assistance::should_assist(idle_constant_ms)); + SNMALLOC_CHECK(!Assistance::should_assist(1 + idle_constant_ms)); + + Assistance::pending_queue_claimed(); + Assistance::pending_queue_added(); + std::atomic start{false}; + std::atomic assisted{0}; + auto assist_after_start = [&start]() { + while (!start.load(std::memory_order_acquire)) + {} + SNMALLOC_CHECK(!Assistance::should_assist(2 * idle_constant_ms)); + }; + std::thread first(assist_after_start); + std::thread second(assist_after_start); + start.store(true, std::memory_order_release); + first.join(); + second.join(); + + start.store(false, std::memory_order_relaxed); + auto reserve_after_start = [&start, &assisted]() { + while (!start.load(std::memory_order_acquire)) + {} + if (Assistance::should_assist(3 * idle_constant_ms)) + assisted.fetch_add(1, std::memory_order_relaxed); + }; + std::thread third(reserve_after_start); + std::thread fourth(reserve_after_start); + start.store(true, std::memory_order_release); + third.join(); + fourth.join(); + SNMALLOC_CHECK(assisted.load(std::memory_order_relaxed) == 1); + + Assistance::pending_queue_claimed(); + SNMALLOC_CHECK(Assistance::debug_pending_count() == 0); + } + + void add_forwarding_work() + { + using Pool = AllocPool; + + auto* owner = Pool::acquire(); + auto* sender = Pool::acquire(); + void* p = owner->alloc(64); + + owner->flush(); + Pool::release(owner); + sender->dealloc(p); + Pool::release(sender); + } + + void test_single_pass_cleanup() + { + using TestConfig = snmalloc::Config; + using Pool = AllocPool; + using Assistance = AllocPoolAssistance; + + auto* saved = Pool::extract(); + auto* saved_last = last_in_chain(saved); + + add_forwarding_work(); + SNMALLOC_CHECK(Assistance::debug_pending_count() == 0); + + cleanup_unused(); + SNMALLOC_CHECK(Assistance::debug_pending_count() == 1); + + cleanup_unused(); + SNMALLOC_CHECK(Assistance::debug_pending_count() == 0); + + auto* assisting = Pool::try_acquire_front(); + SNMALLOC_CHECK(assisting != nullptr); + SNMALLOC_CHECK(assisting->debug_is_in_use()); + cleanup_unused(); + SNMALLOC_CHECK(assisting->debug_is_in_use()); + Pool::release(assisting); + + bool empty = false; + debug_check_empty(&empty); + SNMALLOC_CHECK(empty); + + if (saved != nullptr) + Pool::restore(saved, saved_last); + } +} + +int main() +{ + setup(); + test_reservation_and_pacing(); + test_single_pass_cleanup(); +} diff --git a/src/test/func/domestication/domestication.cc b/src/test/func/domestication/domestication.cc index 1c2eb9fef..3e1069d42 100644 --- a/src/test/func/domestication/domestication.cc +++ b/src/test/func/domestication/domestication.cc @@ -1,3 +1,4 @@ +#include #include // # define SNMALLOC_TRACING @@ -116,6 +117,69 @@ namespace snmalloc #define SNMALLOC_NAME_MANGLE(a) test_##a #include +namespace +{ + using TestAlloc = Allocator; + using TestPool = AllocPool; + using TestAssistance = AllocPoolAssistance; + + TestAlloc* last_in_chain(TestAlloc* first) + { + auto* last = first; + while (last != nullptr) + { + auto* next = TestPool::extract(last); + if (next == nullptr) + return last; + last = next; + } + return nullptr; + } + + void test_disused_allocator_assistance() + { + auto* saved = TestPool::extract(); + auto* saved_last = last_in_chain(saved); + snmalloc::CustomConfig::domesticate_patch_location = nullptr; + + { + ScopedAllocator sender; + void* p; + { + ScopedAllocator owner; + p = owner->alloc(64); + } + + sender->dealloc(p); + sender->flush(); + SNMALLOC_CHECK(TestAssistance::debug_pending_count() == 1); + + const auto deadline = std::chrono::steady_clock::now() + + std::chrono::milliseconds(3 * uint64_t{SNMALLOC_ASSIST_IDLE_MS}); + while (TestAssistance::debug_pending_count() != 0) + { + if (std::chrono::steady_clock::now() >= deadline) + { + std::cerr << "Allocator assistance did not complete" << std::endl; + abort(); + } + + void* churn = sender->alloc(64); + sender->dealloc(churn); + } + + auto* assisted = TestPool::extract(); + SNMALLOC_CHECK(assisted != nullptr); + SNMALLOC_CHECK(!assisted->debug_has_pending_remote()); + auto* assisted_last = last_in_chain(assisted); + TestPool::restore(assisted, assisted_last); + } + + if (saved != nullptr) + TestPool::restore(saved, saved_last); + } +} + int main() { static constexpr bool pagemap_randomize = @@ -130,6 +194,8 @@ int main() entropy.make_free_list_key(RemoteAllocator::key_global); entropy.make_free_list_key(freelist::Object::key_root); + test_disused_allocator_assistance(); + ScopedAllocator alloc1; // Allocate from alloc1; the size doesn't matter a whole lot, it just needs to diff --git a/src/test/func/fixed_region_alloc/fixed_region_alloc.cc b/src/test/func/fixed_region_alloc/fixed_region_alloc.cc index a4d7bc60f..ed354378b 100644 --- a/src/test/func/fixed_region_alloc/fixed_region_alloc.cc +++ b/src/test/func/fixed_region_alloc/fixed_region_alloc.cc @@ -1,5 +1,6 @@ #include "test/setup.h" +#include #include #include #include @@ -31,6 +32,67 @@ int main() << pointer_offset(oe_base, size) << std::endl; CustomGlobals::init(nullptr, oe_base, size); + + { + using Pool = AllocPool; + using State = AllocPoolAssistance; + + auto count_allocators = []() { + size_t count = 0; + for (auto* alloc = Pool::iterate(); alloc != nullptr; + alloc = Pool::iterate(alloc)) + { + count++; + } + return count; + }; + + auto sender = get_scoped_allocator(); + void* remote; + { + auto owner = get_scoped_allocator(); + remote = owner->alloc(128); + SNMALLOC_CHECK(remote != nullptr); + } + + sender->dealloc(remote); + sender->flush(); + SNMALLOC_CHECK(State::debug_pending_count() == 1); + + const size_t allocator_count = count_allocators(); + constexpr size_t batch_size = 256; + void* batch[batch_size]; + size_t iterations = 0; + constexpr size_t iteration_limit = 1 << 24; + const auto deadline = + std::chrono::steady_clock::now() + + std::chrono::milliseconds( + 3 * allocator_count * uint64_t{SNMALLOC_ASSIST_IDLE_MS}); + + while (State::debug_pending_count() != 0) + { + if ( + (std::chrono::steady_clock::now() >= deadline) || + (iterations == iteration_limit)) + { + std::cerr << "Fixed-range assistance did not complete: pending=" + << State::debug_pending_count() << std::endl; + abort(); + } + + for (auto& p : batch) + { + p = sender->alloc(128); + SNMALLOC_CHECK(p != nullptr); + } + for (auto p : batch) + sender->dealloc(p); + iterations++; + } + + SNMALLOC_CHECK(count_allocators() == allocator_count); + } + auto a = get_scoped_allocator(); size_t object_size = 128; diff --git a/src/test/func/freelist_mpscq/freelist_mpscq.cc b/src/test/func/freelist_mpscq/freelist_mpscq.cc index 3c1af0b24..124abfd1e 100644 --- a/src/test/func/freelist_mpscq/freelist_mpscq.cc +++ b/src/test/func/freelist_mpscq/freelist_mpscq.cc @@ -1,5 +1,7 @@ +#include #include #include +#include using namespace snmalloc; @@ -75,7 +77,7 @@ namespace /** * Enqueue reports whether it starts a queue generation and preserves FIFO * order for single-element and pre-linked multi-element enqueues. After a - * callback stops dequeue, the next dequeue resumes at the successor. + * callback rejects an object, the next dequeue resumes at that object. * Dequeue retains the object at back; drain delivers it and resets the queue * for reuse. */ @@ -89,28 +91,32 @@ namespace Object replacement; CallbackOrder order; - SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(first), domesticate)); SNMALLOC_CHECK( - !queue.enqueue(as_head(second), as_head(second), domesticate)); + queue.enqueue(as_head(first), as_head(first), domesticate) == + EnqueueResult::StartedActive); + SNMALLOC_CHECK( + queue.enqueue(as_head(second), as_head(second), domesticate) == + EnqueueResult::Appended); freelist::Object::atomic_store_next( as_head(batch_first), as_head(batch_last), queue_key, NO_KEY_TWEAK); SNMALLOC_CHECK( - !queue.enqueue(as_head(batch_first), as_head(batch_last), domesticate)); + queue.enqueue(as_head(batch_first), as_head(batch_last), domesticate) == + EnqueueResult::Appended); - queue.dequeue(domesticate, domesticate, [&order](freelist::HeadPtr value) { - order.add(value); - return false; - }); - SNMALLOC_CHECK(order.count == 1); + queue.dequeue( + domesticate, domesticate, [](freelist::HeadPtr) { return false; }); + SNMALLOC_CHECK(order.count == 0); SNMALLOC_CHECK( - address_cast(order.values[0]) == address_cast(as_head(first))); + address_cast(queue.front.load()) == address_cast(as_head(first))); queue.dequeue(domesticate, domesticate, [&order](freelist::HeadPtr value) { order.add(value); return true; }); SNMALLOC_CHECK(order.count == 3); + SNMALLOC_CHECK( + address_cast(order.values[0]) == address_cast(as_head(first))); SNMALLOC_CHECK( address_cast(order.values[1]) == address_cast(as_head(second))); SNMALLOC_CHECK( @@ -125,7 +131,8 @@ namespace address_cast(order.values[3]) == address_cast(as_head(batch_last))); SNMALLOC_CHECK( - queue.enqueue(as_head(replacement), as_head(replacement), domesticate)); + queue.enqueue(as_head(replacement), as_head(replacement), domesticate) == + EnqueueResult::StartedActive); queue.drain_and_reset( domesticate, domesticate, [&order](freelist::HeadPtr value) { order.add(value); @@ -150,7 +157,9 @@ namespace size_t dequeue_callbacks = 0; size_t drain_callbacks = 0; - SNMALLOC_CHECK(queue.enqueue(as_head(tail), as_head(tail), domesticate)); + SNMALLOC_CHECK( + queue.enqueue(as_head(tail), as_head(tail), domesticate) == + EnqueueResult::StartedActive); freelist::Object::atomic_store_next( as_head(tail), as_head(outside), queue_key, NO_KEY_TWEAK); @@ -197,7 +206,8 @@ namespace freelist::Object::atomic_store_next( as_head(old_first), as_head(old_last), queue_key, NO_KEY_TWEAK); SNMALLOC_CHECK( - queue.enqueue(as_head(old_first), as_head(old_last), domesticate)); + queue.enqueue(as_head(old_first), as_head(old_last), domesticate) == + EnqueueResult::StartedActive); queue.drain_and_reset( domesticate, domesticate, [&](freelist::HeadPtr value) { @@ -210,7 +220,9 @@ namespace queue_key, NO_KEY_TWEAK); replacement_started = queue.enqueue( - as_head(replacement_first), as_head(replacement_last), domesticate); + as_head(replacement_first), + as_head(replacement_last), + domesticate) == EnqueueResult::StartedActive; } }); @@ -252,9 +264,12 @@ namespace size_t queue_domesticates = 0; CallbackOrder, 2> order; - SNMALLOC_CHECK(remote.enqueue(first_message, first_message, domesticate)); SNMALLOC_CHECK( - !remote.enqueue(second_message, second_message, domesticate)); + remote.enqueue(first_message, first_message, domesticate) == + EnqueueResult::StartedActive); + SNMALLOC_CHECK( + remote.enqueue(second_message, second_message, domesticate) == + EnqueueResult::Appended); auto domesticate_head = [&](freelist::QueuePtr value) -> freelist::HeadPtr { head_domesticates++; @@ -281,6 +296,203 @@ namespace SNMALLOC_CHECK( address_cast(order.values[1]) == address_cast(second_message)); } + + void test_inactive_marker_transitions() + { + Queue queue; + Object first; + size_t callbacks = 0; + + SNMALLOC_CHECK(!queue.release_to_inactive()); + SNMALLOC_CHECK(queue.is_empty()); + + SNMALLOC_CHECK( + queue.enqueue(as_head(first), as_head(first), domesticate) == + EnqueueResult::StartedInactive); + SNMALLOC_CHECK(queue.claim_from_inactive()); + queue.drain_and_reset( + domesticate, domesticate, [&callbacks, &first](freelist::HeadPtr value) { + SNMALLOC_CHECK(address_cast(value) == address_cast(as_head(first))); + callbacks++; + }); + SNMALLOC_CHECK(callbacks == 1); + + SNMALLOC_CHECK(!queue.release_to_inactive()); + SNMALLOC_CHECK(!queue.claim_from_inactive()); + } + + void test_closing_dequeue_contract() + { + Queue queue; + Object first; + Object last; + Object replacement; + CallbackOrder order; + + freelist::Object::atomic_store_next( + as_head(first), as_head(last), queue_key, NO_KEY_TWEAK); + SNMALLOC_CHECK( + queue.enqueue(as_head(first), as_head(last), domesticate) == + EnqueueResult::StartedActive); + queue.dequeue( + domesticate, domesticate, [&first](freelist::HeadPtr value) { + SNMALLOC_CHECK(address_cast(value) == address_cast(as_head(first))); + return false; + }); + SNMALLOC_CHECK(order.count == 0); + SNMALLOC_CHECK( + address_cast(queue.front.load()) == address_cast(as_head(first))); + SNMALLOC_CHECK( + address_cast(queue.back.load()) == address_cast(as_head(last))); + + queue.dequeue( + domesticate, + domesticate, + [&order, &last, &replacement, &queue](freelist::HeadPtr value) { + order.add(value); + if (address_cast(value) == address_cast(as_head(last))) + { + SNMALLOC_CHECK( + queue.enqueue( + as_head(replacement), as_head(replacement), domesticate) == + EnqueueResult::StartedActive); + } + return true; + }); + SNMALLOC_CHECK(order.count == 2); + SNMALLOC_CHECK( + address_cast(order.values[0]) == address_cast(as_head(first))); + SNMALLOC_CHECK( + address_cast(order.values[1]) == address_cast(as_head(last))); + + queue.dequeue( + domesticate, + domesticate, + [&order, &replacement](freelist::HeadPtr value) { + SNMALLOC_CHECK( + address_cast(value) == address_cast(as_head(replacement))); + order.add(value); + return true; + }); + SNMALLOC_CHECK(order.count == 3); + SNMALLOC_CHECK(queue.front.load() == nullptr); + SNMALLOC_CHECK(queue.back.load() == nullptr); + } + + void test_closing_dequeue_unpublished_states() + { + Queue queue; + Object first; + Object last; + size_t callbacks = 0; + + freelist::Object::atomic_store_null(as_head(last), queue_key, NO_KEY_TWEAK); + queue.back.store(capptr_rewild(as_head(last)), stl::memory_order_release); + queue.dequeue( + domesticate, domesticate, [&callbacks](freelist::HeadPtr) { + callbacks++; + return true; + }); + SNMALLOC_CHECK(callbacks == 0); + queue.front.store(capptr_rewild(as_head(last)), stl::memory_order_release); + queue.dequeue( + domesticate, domesticate, [&callbacks](freelist::HeadPtr) { + callbacks++; + return true; + }); + SNMALLOC_CHECK(callbacks == 1); + SNMALLOC_CHECK(queue.front.load() == nullptr); + SNMALLOC_CHECK(queue.back.load() == nullptr); + + Queue linked_queue; + freelist::Object::atomic_store_null( + as_head(first), queue_key, NO_KEY_TWEAK); + freelist::Object::atomic_store_null(as_head(last), queue_key, NO_KEY_TWEAK); + linked_queue.front.store( + capptr_rewild(as_head(first)), stl::memory_order_release); + linked_queue.back.store( + capptr_rewild(as_head(last)), stl::memory_order_release); + linked_queue.dequeue( + domesticate, domesticate, [&callbacks](freelist::HeadPtr) { + callbacks++; + return true; + }); + SNMALLOC_CHECK(callbacks == 1); + SNMALLOC_CHECK( + address_cast(linked_queue.front.load()) == address_cast(as_head(first))); + SNMALLOC_CHECK( + address_cast(linked_queue.back.load()) == address_cast(as_head(last))); + + freelist::Object::atomic_store_next( + as_head(first), as_head(last), queue_key, NO_KEY_TWEAK); + linked_queue.dequeue( + domesticate, domesticate, [&callbacks](freelist::HeadPtr) { + callbacks++; + return true; + }); + SNMALLOC_CHECK(callbacks == 3); + SNMALLOC_CHECK(linked_queue.front.load() == nullptr); + SNMALLOC_CHECK(linked_queue.back.load() == nullptr); + } + + void test_closing_dequeue_append_race() + { + Queue queue; + Object first; + Object old_tail; + Object appended; + std::atomic preflight_reached{false}; + std::atomic append_complete{false}; + CallbackOrder order; + + freelist::Object::atomic_store_next( + as_head(first), as_head(old_tail), queue_key, NO_KEY_TWEAK); + SNMALLOC_CHECK( + queue.enqueue(as_head(first), as_head(old_tail), domesticate) == + EnqueueResult::StartedActive); + + std::thread producer([&]() { + while (!preflight_reached.load(std::memory_order_acquire)) + {} + SNMALLOC_CHECK( + queue.enqueue(as_head(appended), as_head(appended), domesticate) == + EnqueueResult::Appended); + append_complete.store(true, std::memory_order_release); + }); + + auto block_during_preflight = + [&](freelist::QueuePtr value) -> freelist::HeadPtr { + preflight_reached.store(true, std::memory_order_release); + while (!append_complete.load(std::memory_order_acquire)) + {} + return domesticate(value); + }; + + queue.dequeue( + domesticate, block_during_preflight, [&order](freelist::HeadPtr value) { + order.add(value); + return true; + }); + producer.join(); + SNMALLOC_CHECK(order.count == 1); + SNMALLOC_CHECK( + address_cast(order.values[0]) == address_cast(as_head(first))); + SNMALLOC_CHECK( + address_cast(queue.front.load()) == address_cast(as_head(old_tail))); + SNMALLOC_CHECK( + address_cast(queue.back.load()) == address_cast(as_head(appended))); + + queue.dequeue( + domesticate, domesticate, [&order](freelist::HeadPtr value) { + order.add(value); + return true; + }); + SNMALLOC_CHECK(order.count == 3); + SNMALLOC_CHECK( + address_cast(order.values[1]) == address_cast(as_head(old_tail))); + SNMALLOC_CHECK( + address_cast(order.values[2]) == address_cast(as_head(appended))); + } } int main() @@ -291,4 +503,8 @@ int main() test_dequeue_checks_bound_first(); test_callback_starts_replacement(); test_remote_allocator_uses_distinct_domesticators(); + test_inactive_marker_transitions(); + test_closing_dequeue_contract(); + test_closing_dequeue_unpublished_states(); + test_closing_dequeue_append_race(); } diff --git a/src/test/func/multi_setspecific/multi_setspecific.cc b/src/test/func/multi_setspecific/multi_setspecific.cc index 486116fe5..181c62aa4 100644 --- a/src/test/func/multi_setspecific/multi_setspecific.cc +++ b/src/test/func/multi_setspecific/multi_setspecific.cc @@ -68,7 +68,8 @@ int main() std::thread(thread_setspecific).join(); // There should be a single allocator that can be extracted. - if (snmalloc::AllocPool::extract() == nullptr) + auto* first = snmalloc::AllocPool::extract(); + if (first == nullptr) { // The thread has not torn down its allocator. snmalloc::report_fatal_error( @@ -76,6 +77,11 @@ int main() return 1; } + auto* last = first; + while (auto* next = snmalloc::AllocPool::extract(last)) + last = next; + snmalloc::AllocPool::restore(first, last); + return 0; } #else diff --git a/src/test/func/pool/pool.cc b/src/test/func/pool/pool.cc index 8f11ff689..44f6e97b7 100644 --- a/src/test/func/pool/pool.cc +++ b/src/test/func/pool/pool.cc @@ -51,6 +51,24 @@ struct PoolSortEntry : Pooled using PoolSort = Pool; +struct PoolReportEntry : Pooled +{ + size_t reinit_count = 0; + size_t deinit_count = 0; + + void reinit() + { + reinit_count++; + } + + void deinit() + { + deinit_count++; + } +}; + +using PoolReport = Pool; + void test_alloc() { auto ptr = PoolA::acquire(); @@ -74,6 +92,53 @@ void test_constructor() PoolB::release(ptr2); } +void test_lifecycle_hooks() +{ + auto* first = PoolReport::acquire(); + auto* second = PoolReport::acquire(); + SNMALLOC_CHECK(first->reinit_count == 0); + SNMALLOC_CHECK(first->deinit_count == 0); + + PoolReport::release(first); + PoolReport::release(second); + SNMALLOC_CHECK(first->deinit_count == 1); + SNMALLOC_CHECK(second->deinit_count == 1); + + auto* removed = PoolReport::try_acquire_front(); + SNMALLOC_CHECK(removed == first); + SNMALLOC_CHECK(removed->reinit_count == 1); + PoolReport::release(removed); + + auto* reused = PoolReport::acquire(); + SNMALLOC_CHECK(reused == second); + SNMALLOC_CHECK(reused->reinit_count == 1); + PoolReport::release(reused); + + first = PoolReport::extract(); + SNMALLOC_CHECK(first != nullptr); + auto* last = first; + size_t count = 0; + while (last != nullptr) + { + SNMALLOC_CHECK(last->reinit_count == last->deinit_count); + count++; + auto* next = PoolReport::extract(last); + if (next == nullptr) + break; + last = next; + } + SNMALLOC_CHECK(count == 2); + PoolReport::restore(first, last); + + PoolReport::sort(); + auto* item = PoolReport::iterate(); + while (item != nullptr) + { + SNMALLOC_CHECK(item->reinit_count + 1 == item->deinit_count); + item = PoolReport::iterate(item); + } +} + void test_alloc_many() { constexpr size_t count = 16'000'000 / MIN_CHUNK_SIZE; @@ -250,6 +315,8 @@ int main(int argc, char** argv) std::cout << "test_alloc passed" << std::endl; test_constructor(); std::cout << "test_constructor passed" << std::endl; + test_lifecycle_hooks(); + std::cout << "test_lifecycle_hooks passed" << std::endl; test_alloc_many(); std::cout << "test_alloc_many passed" << std::endl; test_different_alloc(); diff --git a/src/test/func/teardown/teardown.cc b/src/test/func/teardown/teardown.cc index 02129374d..9df12bc13 100644 --- a/src/test/func/teardown/teardown.cc +++ b/src/test/func/teardown/teardown.cc @@ -142,4 +142,5 @@ int main(int, char**) f(shifted(7) - 1); printf("\n"); } + snmalloc::debug_in_use(0); } diff --git a/src/test/func/ticker/ticker.cc b/src/test/func/ticker/ticker.cc new file mode 100644 index 000000000..f88a1844a --- /dev/null +++ b/src/test/func/ticker/ticker.cc @@ -0,0 +1,71 @@ +#include +#include + +using namespace snmalloc; + +namespace +{ + struct TestPal + { + static constexpr uint64_t pal_features = PalFeatures::Time; + inline static uint64_t now = 1; + + static uint64_t time_in_ms() + { + return now; + } + }; + + void test_sample_callback_and_long_gap() + { + Ticker ticker; + size_t callbacks = 0; + uint64_t last_sample = 0; + + ticker.check_tick(nullptr, [&](uint64_t sample) { + callbacks++; + last_sample = sample; + }); + SNMALLOC_CHECK(callbacks == 1); + SNMALLOC_CHECK(last_sample == 1); + + TestPal::now = 1001; + ticker.check_tick(nullptr, [&](uint64_t sample) { + callbacks++; + last_sample = sample; + }); + SNMALLOC_CHECK(callbacks == 2); + SNMALLOC_CHECK(last_sample == 1001); + + TestPal::now = 1002; + ticker.check_tick(nullptr, [&](uint64_t sample) { + callbacks++; + last_sample = sample; + }); + SNMALLOC_CHECK(callbacks == 3); + SNMALLOC_CHECK(last_sample == 1002); + } + + void test_zero_duration_callback() + { + Ticker ticker; + size_t callbacks = 0; + TestPal::now = 10; + + ticker.check_tick(nullptr, [&](uint64_t) { callbacks++; }); + ticker.check_tick(nullptr, [&](uint64_t) { callbacks++; }); + SNMALLOC_CHECK(callbacks == 2); + + ticker.check_tick(nullptr, [&](uint64_t) { callbacks++; }); + SNMALLOC_CHECK(callbacks == 3); + ticker.check_tick(nullptr, [&](uint64_t) { callbacks++; }); + SNMALLOC_CHECK(callbacks == 3); + } +} + +int main() +{ + setup(); + test_sample_callback_and_long_gap(); + test_zero_duration_callback(); +} diff --git a/src/test/perf/disused_remote/disused_remote.cc b/src/test/perf/disused_remote/disused_remote.cc new file mode 100644 index 000000000..ac54580ee --- /dev/null +++ b/src/test/perf/disused_remote/disused_remote.cc @@ -0,0 +1,342 @@ +#include "test/opt.h" +#include "test/setup.h" +#include "test/xoroshiro.h" + +#include +#include +#include +#include +#include +#include +#include + +using namespace snmalloc; + +namespace +{ + struct Slot + { + void* allocation; + size_t owner; + bool replaced; + }; + + struct PoolStats + { + size_t count = 0; + size_t in_use = 0; + size_t pending = 0; + }; + + PoolStats pool_stats() + { + PoolStats result; + auto* alloc = AllocPool::iterate(); + while (alloc != nullptr) + { + result.count++; + result.in_use += alloc->debug_is_in_use() ? 1 : 0; + result.pending += alloc->debug_has_pending_remote() ? 1 : 0; + alloc = AllocPool::iterate(alloc); + } + return result; + } + + void wait_until(const std::atomic& value, size_t target) + { + while (value.load(std::memory_order_acquire) < target) + std::this_thread::yield(); + } + + size_t latency_bucket(uint64_t nanoseconds) + { + size_t bucket = 0; + while ((nanoseconds > 1) && (bucket < 63)) + { + nanoseconds >>= 1; + bucket++; + } + return bucket; + } +} + +int main(int argc, char** argv) +{ + setup(); + opt::Opt opt(argc, argv); + + const bool smoke = opt.has("--smoke"); + const size_t builder_count = opt.is("--builders", smoke ? 4 : 8); + const size_t table_size = + opt.is("--table", smoke ? (1 << 15) : (1 << 18)); + + if ((builder_count == 0) || (table_size < builder_count)) + { + std::cerr << "builders must be non-zero and no larger than table size" + << std::endl; + return 1; + } + + const size_t requested_size = opt.is("--size", 1024); + const size_t overlap = opt.is("--overlap", table_size / 32); + const size_t termination_gap = + opt.is("--termination-gap", table_size / (8 * builder_count)); + const size_t survivor_churn = opt.is("--survivor-churn", table_size); + const size_t minimum_churn_ms = opt.is("--minimum-churn-ms", 0); + using SchedulingState = AllocPoolAssistance; + + if (requested_size > MAX_SMALL_SIZECLASS_SIZE) + { + std::cerr << "size must use a small sizeclass" << std::endl; + return 1; + } + + std::vector table(table_size); + std::vector> builder_exit(builder_count); + std::vector> builder_released(builder_count); + std::vector replaced_before_release(builder_count); + std::vector replaced_after_release(builder_count); + std::vector builder_queues(builder_count); + + std::atomic builders_ready{0}; + std::atomic churn_ready{false}; + std::atomic start_churn{false}; + std::atomic churn_operations{0}; + std::atomic first_replacements{0}; + std::atomic stop_churn{false}; + std::array allocation_latency_histogram{}; + uint64_t maximum_allocation_latency_ns = 0; + uint64_t churn_elapsed_ns = 0; + uint64_t pending_cleared_ns = 0; + + for (size_t i = 0; i < builder_count; i++) + { + builder_exit[i] = false; + builder_released[i] = false; + } + + std::thread churn_thread([&]() { + void* initial = snmalloc::alloc(1); + snmalloc::dealloc(initial); + churn_ready.store(true, std::memory_order_release); + + while (!start_churn.load(std::memory_order_acquire)) + std::this_thread::yield(); + + xoroshiro::p128r32 random(0x5eed, 0x1234); + size_t remaining = table_size; + size_t survivor_count = 0; + bool observed_pending = false; + bool minimum_time_complete = minimum_churn_ms == 0; + const auto churn_start = std::chrono::steady_clock::now(); + + while ((remaining != 0) || (survivor_count < survivor_churn) || + !stop_churn.load(std::memory_order_acquire) || + !minimum_time_complete) + { + const size_t index = random.next() % table_size; + auto& slot = table[index]; + const auto allocation_start = std::chrono::steady_clock::now(); + void* replacement = snmalloc::alloc(requested_size); + const auto allocation_end = std::chrono::steady_clock::now(); + const auto allocation_latency_ns = static_cast( + std::chrono::duration_cast( + allocation_end - allocation_start) + .count()); + allocation_latency_histogram[latency_bucket(allocation_latency_ns)]++; + if (allocation_latency_ns > maximum_allocation_latency_ns) + maximum_allocation_latency_ns = allocation_latency_ns; + snmalloc::dealloc(slot.allocation); + slot.allocation = replacement; + + if (!slot.replaced) + { + slot.replaced = true; + remaining--; + first_replacements.fetch_add(1, std::memory_order_release); + if (builder_released[slot.owner].load(std::memory_order_acquire)) + replaced_after_release[slot.owner]++; + else + replaced_before_release[slot.owner]++; + } + else if (remaining == 0) + { + survivor_count++; + } + + const size_t operation = + churn_operations.fetch_add(1, std::memory_order_release) + 1; + if ((operation & 1023) == 0) + { + const auto pending = SchedulingState::debug_pending_count(); + observed_pending |= pending > 0; + if (observed_pending && (pending <= 0) && (pending_cleared_ns == 0)) + { + pending_cleared_ns = static_cast( + std::chrono::duration_cast( + allocation_end - churn_start) + .count()); + } + } + if (!minimum_time_complete && ((operation & 1023) == 0)) + { + const auto elapsed = + std::chrono::duration_cast( + allocation_end - churn_start); + minimum_time_complete = + elapsed.count() >= static_cast(minimum_churn_ms); + } + } + + churn_elapsed_ns = static_cast( + std::chrono::duration_cast( + std::chrono::steady_clock::now() - churn_start) + .count()); + }); + + std::vector builders; + builders.reserve(builder_count); + for (size_t builder = 0; builder < builder_count; builder++) + { + builders.emplace_back([&, builder]() { + for (size_t index = builder; index < table_size; index += builder_count) + { + void* allocation = snmalloc::alloc(requested_size); + table[index] = {allocation, builder, false}; + if (builder_queues[builder] == nullptr) + { + const auto& entry = + Config::Backend::get_metaentry(address_cast(allocation)); + builder_queues[builder] = entry.get_remote(); + } + } + + builders_ready.fetch_add(1, std::memory_order_acq_rel); + while (!builder_exit[builder].load(std::memory_order_acquire)) + std::this_thread::yield(); + }); + } + + wait_until(builders_ready, builder_count); + while (!churn_ready.load(std::memory_order_acquire)) + std::this_thread::yield(); + + for (size_t i = 0; i < builder_count; i++) + { + if (builder_queues[i] == nullptr) + { + std::cerr << "builder did not allocate" << std::endl; + return 1; + } + for (size_t j = 0; j < i; j++) + { + if (builder_queues[i] == builder_queues[j]) + { + std::cerr << "builders did not retain distinct allocators" << std::endl; + return 1; + } + } + } + + const size_t usable_size = snmalloc::alloc_size(table[0].allocation); + const size_t slab_size = + sizeclass_to_slab_size(size_to_sizeclass(requested_size)); + const size_t allocator_count = builder_count + 1; + const size_t retained_slab_tolerance = allocator_count * slab_size; + const size_t remote_batch_tolerance = + allocator_count * static_cast(REMOTE_CACHE); + const size_t metadata_tolerance = allocator_count * MIN_CHUNK_SIZE; + const size_t retention_tolerance = + retained_slab_tolerance + remote_batch_tolerance + metadata_tolerance; + const size_t phase1_current = Config::Backend::get_current_usage(); + const size_t phase1_peak = Config::Backend::get_peak_usage(); + + start_churn.store(true, std::memory_order_release); + wait_until(churn_operations, overlap); + + for (size_t builder = 0; builder < builder_count; builder++) + { + builder_exit[builder].store(true, std::memory_order_release); + builders[builder].join(); + builder_released[builder].store(true, std::memory_order_release); + wait_until(churn_operations, overlap + ((builder + 1) * termination_gap)); + } + + wait_until(first_replacements, table_size); + wait_until(churn_operations, overlap + table_size + survivor_churn); + stop_churn.store(true, std::memory_order_release); + churn_thread.join(); + + const auto stats = pool_stats(); + const size_t final_current = Config::Backend::get_current_usage(); + const size_t final_peak = Config::Backend::get_peak_usage(); + const auto inactive_pending_count = SchedulingState::debug_pending_count(); + + const size_t measured_operations = churn_operations.load(); + const size_t p99_target = ((measured_operations * 99) + 99) / 100; + size_t p99_bucket = 0; + size_t cumulative_latency_samples = 0; + for (; p99_bucket < allocation_latency_histogram.size(); p99_bucket++) + { + cumulative_latency_samples += allocation_latency_histogram[p99_bucket]; + if (cumulative_latency_samples >= p99_target) + break; + } + const uint64_t allocation_p99_upper_ns = + p99_bucket == 63 ? UINT64_MAX : (uint64_t{1} << (p99_bucket + 1)); + const uint64_t churn_operations_per_second = churn_elapsed_ns == 0 ? + 0 : + static_cast( + static_cast(measured_operations) * 1000000000.0 / + static_cast(churn_elapsed_ns)); + + size_t post_release_replacements = 0; + for (size_t builder = 0; builder < builder_count; builder++) + { + post_release_replacements += replaced_after_release[builder]; + std::cout << "builder[" << builder << "] first_replacements_before_release=" + << replaced_before_release[builder] + << " first_replacements_after_release=" + << replaced_after_release[builder] << std::endl; + } + + const size_t expected_retained_payload = + post_release_replacements * usable_size; + if (expected_retained_payload <= (4 * retention_tolerance)) + { + std::cerr << "retained-payload signal is too small for tolerance" + << std::endl; + return 1; + } + + std::cout << "builders=" << builder_count << " table=" << table_size + << " requested_size=" << requested_size + << " usable_size=" << usable_size + << " live_payload=" << (table_size * requested_size) + << " expected_retained_payload=" << expected_retained_payload + << " retained_slab_tolerance=" << retained_slab_tolerance + << " remote_batch_tolerance=" << remote_batch_tolerance + << " metadata_tolerance=" << metadata_tolerance + << " retention_tolerance=" << retention_tolerance + << " phase1_current=" << phase1_current + << " phase1_peak=" << phase1_peak + << " final_current=" << final_current + << " final_peak=" << final_peak << " pool_size=" << stats.count + << " pool_in_use=" << stats.in_use + << " pooled_pending_queues=" << stats.pending + << " inactive_pending_count=" << inactive_pending_count + << " pending_cleared_ns=" << pending_cleared_ns + << " churn_operations=" << measured_operations + << " churn_elapsed_ns=" << churn_elapsed_ns + << " churn_operations_per_second=" << churn_operations_per_second + << " allocation_p99_upper_ns=" << allocation_p99_upper_ns + << " maximum_allocation_latency_ns=" + << maximum_allocation_latency_ns << std::endl; + + for (auto& slot : table) + snmalloc::dealloc(slot.allocation); + cleanup_unused(); + cleanup_unused(); + debug_check_empty(); + return 0; +}