From 6a1a0b183147c8631ab19b9e031801a0b84e86b4 Mon Sep 17 00:00:00 2001 From: Matthew Parkinson Date: Thu, 17 Sep 2026 11:48:09 +0100 Subject: [PATCH 1/4] Refactor freelist MPSC queue draining Expose empty-to-nonempty enqueue transitions and add a concurrent drain-and-reset operation. Share bounded chain processing with ordinary dequeue while preserving its retained tail semantics. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 75ccdc6c-2024-4d81-b22f-f1b52fbe29cf --- src/snmalloc/mem/corealloc.h | 37 +- src/snmalloc/mem/freelist_queue.h | 142 ++-- src/snmalloc/mem/remoteallocator.h | 19 +- .../func/freelist_mpscq/freelist_mpscq.cc | 658 ++++++++++++++++++ src/test/perf/msgpass/msgpass.cc | 13 +- 5 files changed, 770 insertions(+), 99 deletions(-) create mode 100644 src/test/func/freelist_mpscq/freelist_mpscq.cc diff --git a/src/snmalloc/mem/corealloc.h b/src/snmalloc/mem/corealloc.h index b9d63c957..4e5b2872a 100644 --- a/src/snmalloc/mem/corealloc.h +++ b/src/snmalloc/mem/corealloc.h @@ -1386,11 +1386,13 @@ namespace snmalloc } /** - * Flush the cached state and delayed deallocations + * 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. */ - bool flush(bool destroy_queue = false) + bool flush() { auto local_state = backend_state_ptr(); auto domesticate = [local_state](freelist::QueuePtr p) @@ -1400,26 +1402,17 @@ namespace snmalloc size_t bytes_flushed = 0; // Not currently used. - if (destroy_queue) - { - auto cb = - [this, domesticate, &bytes_flushed](capptr::Alloc m) { - bool need_post = true; // Always going to post, so ignore. - const PagemapEntry& entry = - Config::Backend::get_metaentry(snmalloc::address_cast(m)); - handle_dealloc_remote( - entry, m, need_post, domesticate, bytes_flushed); - }; + auto cb = [this, domesticate, &bytes_flushed]( + capptr::Alloc m) { + // Forwarded messages are posted together after the drain, so skip + // per-message capacity checks and intermediate-post signalling. + bool need_post = true; + const PagemapEntry& entry = + Config::Backend::get_metaentry(snmalloc::address_cast(m)); + handle_dealloc_remote(entry, m, need_post, domesticate, bytes_flushed); + }; - message_queue().destroy_and_iterate(domesticate, cb); - } - else - { - // Process incoming message queue - // Loop as normally only processes a batch - while (has_messages()) - handle_message_queue([]() {}); - } + message_queue().drain_and_reset(domesticate, cb); auto& key = freelist::Object::key_root; @@ -1514,7 +1507,7 @@ namespace snmalloc }); }; - bool sent_something = flush(true); + bool sent_something = flush(); for (auto& alloc_class : alloc_classes) { diff --git a/src/snmalloc/mem/freelist_queue.h b/src/snmalloc/mem/freelist_queue.h index 452033249..f79d6be54 100644 --- a/src/snmalloc/mem/freelist_queue.h +++ b/src/snmalloc/mem/freelist_queue.h @@ -16,18 +16,15 @@ namespace snmalloc * for the client to reach as the Pagemap, which we trust to store not just * Tame CapPtr<>s but raw C++ pointers. * - * Where necessary, methods expose two domesticator callbacks at the - * interface and are careful to use one for the front and back values and the - * other for pointers read from the queue itself. That's not ideal, but it - * lets the client condition its behavior appropriately and prevents us from - * accidentally following either of these pointers in generic code. - * Specifically, + * Where necessary, dequeue exposes two domesticator callbacks and is careful + * to use one for the front value and the other for pointers read from the + * queue itself. Draining uses its queue domesticator for both the front + * value and links in the chain. Specifically, * * * `domesticate_head` is used for the MPSCQ pointers used to reach into * the chain of objects * - * * `domesticate_queue` is used to traverse links in that chain (and in - * fact, we traverse only the first). + * * `domesticate_queue` is used to traverse links in that chain. * * In the case that the MPSCQ is not easily accessible to the client, * `domesticate_head` can just be a type coersion, and `domesticate_queue` @@ -51,35 +48,79 @@ namespace snmalloc SNMALLOC_ASSERT(pointer_align_up(this, REMOTE_MIN_ALIGN) == this); } - void init() + private: + template + freelist::HeadPtr process_chain( + freelist::HeadPtr curr, + freelist::QueuePtr target, + Domesticator_queue& domesticate, + Cb& cb) { - back.store(nullptr); - front.store(nullptr); - invariant(); - } + while (address_cast(curr) != address_cast(target)) + { + auto next = curr->atomic_read_next(Key, Key_tweak, domesticate); + if (SNMALLOC_UNLIKELY(next == nullptr)) + return curr; - freelist::QueuePtr destroy() - { - if (back.load(stl::memory_order_relaxed) == nullptr) - return nullptr; + Aal::prefetch(next.unsafe_ptr()); + if (SNMALLOC_UNLIKELY(!cb(curr))) + return next; - freelist::QueuePtr fnt = front.load(); - back.store(nullptr, stl::memory_order_relaxed); - front.store(nullptr, stl::memory_order_relaxed); - return fnt; + curr = next; + } + + return curr; } + public: + /** + * Exactly one consumer role may execute this operation. No dequeue, + * owning-allocator queue processing, or second drain may run concurrently. + * + * A producer that has exchanged back must eventually publish front or its + * predecessor link; otherwise this operation waits indefinitely. + * + * The queue is reset before the first callback. The callback may therefore + * release or re-enqueue an object; any re-enqueue belongs to the + * replacement chain and is not consumed by this invocation. + */ template - void destroy_and_iterate(Domesticator_queue domesticate, Cb cb) + void drain_and_reset(Domesticator_queue domesticate, Cb cb) { - auto p = domesticate(destroy()); + // 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) + return; - while (p != nullptr) + freelist::HeadPtr curr = nullptr; + do { - auto n = p->atomic_read_next(Key, Key_tweak, domesticate); + auto raw = front.load(stl::memory_order_acquire); + if (raw != nullptr) + curr = domesticate(raw); + if (curr == nullptr) + Aal::pause(); + } while (curr == nullptr); + + // A producer that observes null back may immediately publish a new front, + // so the old front must be cleared before resetting back. + front.store(nullptr, stl::memory_order_relaxed); + auto target = back.exchange(nullptr, stl::memory_order_acq_rel); + SNMALLOC_ASSERT(target != nullptr); + + auto process = [&cb](freelist::HeadPtr p) { cb(p); - p = n; + return true; + }; + + while (true) + { + curr = process_chain(curr, target, domesticate, process); + if (address_cast(curr) == address_cast(target)) + break; + Aal::pause(); } + cb(curr); } inline bool can_dequeue() @@ -94,9 +135,13 @@ 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. */ template - void enqueue( + bool enqueue( freelist::HeadPtr first, freelist::HeadPtr last, Domesticator_head domesticate_head) @@ -125,12 +170,17 @@ namespace snmalloc if (SNMALLOC_LIKELY(prev != nullptr)) { + // 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; + return false; } + // 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; } /** @@ -163,40 +213,12 @@ namespace snmalloc // Use back to bound, so we don't handle new entries. auto b = back.load(stl::memory_order_relaxed); - while (address_cast(curr) != address_cast(b)) - { - freelist::HeadPtr next = - curr->atomic_read_next(Key, Key_tweak, domesticate_queue); - // We have observed a non-linearisable effect of the queue. - // Just go back to allocating normally. - if (SNMALLOC_UNLIKELY(next == nullptr)) - break; - // We want this element next, so start it loading. - Aal::prefetch(next.unsafe_ptr()); - if (SNMALLOC_UNLIKELY(!cb(curr))) - { - /* - * We've domesticate_queue-d next so that we can read through it, but - * we're storing it back into client-accessible memory in - * !QueueHeadsAreTame builds, so go ahead and consider it Wild again. - * On QueueHeadsAreTame builds, the subsequent domesticate_head call - * above will also be a type-level sleight of hand, but we can still - * justify it by the domesticate_queue that happened in this - * dequeue(). - */ - front = capptr_rewild(next); - invariant(); - return; - } - - curr = next; - } - /* - * Here, we've hit the end of the queue: next is nullptr and curr has not - * been handed to the callback. The same considerations about Wildness - * above hold here. + * process_chain may return a pointer domesticated from a queue link. + * 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); invariant(); } diff --git a/src/snmalloc/mem/remoteallocator.h b/src/snmalloc/mem/remoteallocator.h index 1b02381a9..1aab15af2 100644 --- a/src/snmalloc/mem/remoteallocator.h +++ b/src/snmalloc/mem/remoteallocator.h @@ -2,6 +2,7 @@ #include "freelist_queue.h" #include "snmalloc/stl/new.h" +#include "snmalloc/stl/utility.h" namespace snmalloc { @@ -320,19 +321,14 @@ namespace snmalloc list.invariant(); } - void init() - { - list.init(); - } - template - void destroy_and_iterate(Domesticator_queue domesticate, Cb cb) + void drain_and_reset(Domesticator_queue domesticate, Cb cb) { - auto cbwrap = [cb](freelist::HeadPtr p) SNMALLOC_FAST_PATH_LAMBDA { + auto cbwrap = [&cb](freelist::HeadPtr p) SNMALLOC_FAST_PATH_LAMBDA { cb(RemoteMessage::from_message_link(p)); }; - return list.destroy_and_iterate(domesticate, cbwrap); + return list.drain_and_reset(stl::move(domesticate), stl::move(cbwrap)); } inline bool can_dequeue() @@ -346,14 +342,17 @@ 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. */ template - void enqueue( + bool enqueue( capptr::Alloc first, capptr::Alloc last, Domesticator_head domesticate_head) { - list.enqueue( + return list.enqueue( RemoteMessage::to_message_link(first), RemoteMessage::to_message_link(last), domesticate_head); diff --git a/src/test/func/freelist_mpscq/freelist_mpscq.cc b/src/test/func/freelist_mpscq/freelist_mpscq.cc new file mode 100644 index 000000000..06b893b8f --- /dev/null +++ b/src/test/func/freelist_mpscq/freelist_mpscq.cc @@ -0,0 +1,658 @@ +#include +#include +#include +#include +#include +#include +#include +#include + +using namespace snmalloc; + +namespace +{ + FreeListKey key{0x1234, 0x5678, 0x9abc}; + using Queue = FreeListMPSCQ; + using Object = freelist::Object::T<>; + + freelist::HeadPtr domesticate(freelist::QueuePtr p) + { + return freelist::HeadPtr::unsafe_from(p.unsafe_ptr()); + } + + freelist::HeadPtr as_head(Object& object) + { + return freelist::HeadPtr::unsafe_from(&object); + } + + template + void wait_until(Predicate predicate) + { + auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(10); + while (!predicate()) + { + SNMALLOC_CHECK(std::chrono::steady_clock::now() < deadline); + std::this_thread::yield(); + } + } + + struct CallbackOrder + { + freelist::HeadPtr values[8]; + size_t count = 0; + + void add(freelist::HeadPtr value) + { + values[count++] = value; + } + }; + + void check_empty(Queue& queue) + { + SNMALLOC_CHECK(queue.back.load(stl::memory_order_relaxed) == nullptr); + SNMALLOC_CHECK(queue.front.load(stl::memory_order_relaxed) == nullptr); + } + + void test_empty() + { + Queue queue; + size_t callbacks = 0; + + queue.drain_and_reset( + domesticate, [&callbacks](freelist::HeadPtr) { callbacks++; }); + + SNMALLOC_CHECK(callbacks == 0); + check_empty(queue); + } + + void test_enqueue_and_iterate() + { + Queue queue; + Object first; + Object second; + Object batch_first; + Object batch_last; + 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)); + + freelist::Object::atomic_store_next( + as_head(batch_first), as_head(batch_last), key, NO_KEY_TWEAK); + SNMALLOC_CHECK( + !queue.enqueue(as_head(batch_first), as_head(batch_last), domesticate)); + + queue.drain_and_reset( + domesticate, [&order](freelist::HeadPtr value) { order.add(value); }); + + SNMALLOC_CHECK(order.count == 4); + 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( + address_cast(order.values[2]) == address_cast(as_head(batch_first))); + SNMALLOC_CHECK( + address_cast(order.values[3]) == address_cast(as_head(batch_last))); + check_empty(queue); + + SNMALLOC_CHECK( + queue.enqueue(as_head(replacement), as_head(replacement), domesticate)); + queue.enqueue(as_head(second), as_head(second), domesticate); + + order.count = 0; + queue.drain_and_reset( + domesticate, [&order](freelist::HeadPtr value) { order.add(value); }); + SNMALLOC_CHECK(order.count == 2); + SNMALLOC_CHECK( + address_cast(order.values[0]) == address_cast(as_head(replacement))); + SNMALLOC_CHECK( + address_cast(order.values[1]) == address_cast(as_head(second))); + check_empty(queue); + } + + void test_dequeue_empty() + { + Queue queue; + Object pending; + size_t callbacks = 0; + auto cb = [&callbacks](freelist::HeadPtr) { + callbacks++; + return true; + }; + + queue.dequeue(domesticate, domesticate, cb); + SNMALLOC_CHECK(callbacks == 0); + check_empty(queue); + + freelist::Object::atomic_store_null(as_head(pending), key, NO_KEY_TWEAK); + auto prev = queue.back.exchange( + capptr_rewild(as_head(pending)), stl::memory_order_acq_rel); + SNMALLOC_CHECK(prev == nullptr); + + queue.dequeue(domesticate, domesticate, cb); + SNMALLOC_CHECK(callbacks == 0); + SNMALLOC_CHECK(queue.front.load(stl::memory_order_relaxed) == nullptr); + SNMALLOC_CHECK( + address_cast(queue.back.load(stl::memory_order_relaxed)) == + address_cast(as_head(pending))); + + queue.back.store(nullptr, stl::memory_order_relaxed); + check_empty(queue); + } + + void test_dequeue_positions() + { + Queue queue; + Object objects[3]; + CallbackOrder order; + + freelist::Object::atomic_store_next( + as_head(objects[0]), as_head(objects[1]), key, NO_KEY_TWEAK); + freelist::Object::atomic_store_next( + as_head(objects[1]), as_head(objects[2]), key, NO_KEY_TWEAK); + SNMALLOC_CHECK( + queue.enqueue(as_head(objects[0]), as_head(objects[2]), domesticate)); + + queue.dequeue(domesticate, domesticate, [&order](freelist::HeadPtr value) { + order.add(value); + return false; + }); + + SNMALLOC_CHECK(order.count == 1); + SNMALLOC_CHECK( + address_cast(order.values[0]) == address_cast(as_head(objects[0]))); + SNMALLOC_CHECK( + address_cast(queue.front.load(stl::memory_order_relaxed)) == + address_cast(as_head(objects[1]))); + + queue.dequeue(domesticate, domesticate, [&order](freelist::HeadPtr value) { + order.add(value); + return true; + }); + + SNMALLOC_CHECK(order.count == 2); + SNMALLOC_CHECK( + address_cast(order.values[1]) == address_cast(as_head(objects[1]))); + SNMALLOC_CHECK( + address_cast(queue.front.load(stl::memory_order_relaxed)) == + address_cast(as_head(objects[2]))); + SNMALLOC_CHECK( + address_cast(queue.back.load(stl::memory_order_relaxed)) == + address_cast(as_head(objects[2]))); + + order.count = 0; + queue.drain_and_reset( + domesticate, [&order](freelist::HeadPtr value) { order.add(value); }); + SNMALLOC_CHECK(order.count == 1); + SNMALLOC_CHECK( + address_cast(order.values[0]) == address_cast(as_head(objects[2]))); + check_empty(queue); + } + + void test_dequeue_publication_gap() + { + Queue queue; + Object first; + Object second; + size_t callbacks = 0; + + SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(first), domesticate)); + + freelist::Object::atomic_store_null(as_head(second), key, NO_KEY_TWEAK); + auto prev = queue.back.exchange( + capptr_rewild(as_head(second)), stl::memory_order_acq_rel); + SNMALLOC_CHECK(address_cast(prev) == address_cast(as_head(first))); + + queue.dequeue(domesticate, domesticate, [&callbacks](freelist::HeadPtr) { + callbacks++; + return true; + }); + + SNMALLOC_CHECK(callbacks == 0); + SNMALLOC_CHECK( + address_cast(queue.front.load(stl::memory_order_relaxed)) == + address_cast(as_head(first))); + + freelist::Object::atomic_store_next( + as_head(first), as_head(second), key, NO_KEY_TWEAK); + queue.drain_and_reset( + domesticate, [&callbacks](freelist::HeadPtr) { callbacks++; }); + SNMALLOC_CHECK(callbacks == 2); + check_empty(queue); + } + + void test_dequeue_checks_bound_first() + { + Queue queue; + Object first; + Object outside; + size_t callbacks = 0; + size_t domesticates = 0; + + SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(first), domesticate)); + freelist::Object::atomic_store_next( + as_head(first), as_head(outside), key, NO_KEY_TWEAK); + + auto counting_domesticate = + [&domesticates](freelist::QueuePtr value) -> freelist::HeadPtr { + domesticates++; + return domesticate(value); + }; + + queue.dequeue( + domesticate, counting_domesticate, [&callbacks](freelist::HeadPtr) { + callbacks++; + return true; + }); + + SNMALLOC_CHECK(callbacks == 0); + SNMALLOC_CHECK(domesticates == 0); + SNMALLOC_CHECK( + address_cast(queue.front.load(stl::memory_order_relaxed)) == + address_cast(as_head(first))); + SNMALLOC_CHECK( + address_cast(queue.back.load(stl::memory_order_relaxed)) == + address_cast(as_head(first))); + + freelist::Object::atomic_store_null(as_head(first), key, NO_KEY_TWEAK); + queue.drain_and_reset(domesticate, [](freelist::HeadPtr) {}); + check_empty(queue); + } + + void test_remote_allocator_move_only_callback() + { + RemoteAllocator remote; + RemoteMessage message{}; + size_t callbacks = 0; + auto message_ptr = capptr::Alloc::unsafe_from(&message); + + SNMALLOC_CHECK(remote.enqueue(message_ptr, message_ptr, domesticate)); + + auto cb = [token = std::make_unique(1), + &callbacks](capptr::Alloc) mutable { + callbacks += (*token)++; + }; + remote.drain_and_reset(domesticate, std::move(cb)); + + SNMALLOC_CHECK(callbacks == 1); + SNMALLOC_CHECK(remote.list.back.load(stl::memory_order_relaxed) == nullptr); + SNMALLOC_CHECK( + remote.list.front.load(stl::memory_order_relaxed) == nullptr); + } + + void test_front_publication_gap() + { + Queue queue; + Object object; + std::atomic entered{false}; + std::atomic completed{false}; + std::atomic callbacks{0}; + + freelist::Object::atomic_store_null(as_head(object), key, NO_KEY_TWEAK); + auto prev = queue.back.exchange( + capptr_rewild(as_head(object)), stl::memory_order_acq_rel); + SNMALLOC_CHECK(prev == nullptr); + + std::thread consumer([&]() { + entered.store(true, std::memory_order_release); + queue.drain_and_reset(domesticate, [&](freelist::HeadPtr value) { + SNMALLOC_CHECK(address_cast(value) == address_cast(as_head(object))); + callbacks.fetch_add(1, std::memory_order_relaxed); + }); + completed.store(true, std::memory_order_release); + }); + + wait_until([&]() { return entered.load(std::memory_order_acquire); }); + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + SNMALLOC_CHECK(!completed.load(std::memory_order_acquire)); + SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 0); + + queue.front.store(capptr_rewild(as_head(object))); + consumer.join(); + + SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 1); + SNMALLOC_CHECK(completed.load(std::memory_order_acquire)); + check_empty(queue); + } + + void test_successor_publication_gap() + { + Queue queue; + Object first; + Object second; + std::atomic entered{false}; + std::atomic completed{false}; + std::atomic callbacks{0}; + + SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(first), domesticate)); + + freelist::Object::atomic_store_null(as_head(second), key, NO_KEY_TWEAK); + auto prev = queue.back.exchange( + capptr_rewild(as_head(second)), stl::memory_order_acq_rel); + SNMALLOC_CHECK(address_cast(prev) == address_cast(as_head(first))); + + std::thread consumer([&]() { + entered.store(true, std::memory_order_release); + queue.drain_and_reset(domesticate, [&callbacks](freelist::HeadPtr) { + callbacks.fetch_add(1, std::memory_order_relaxed); + }); + completed.store(true, std::memory_order_release); + }); + + wait_until([&]() { return entered.load(std::memory_order_acquire); }); + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + SNMALLOC_CHECK(!completed.load(std::memory_order_acquire)); + SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 0); + + freelist::Object::atomic_store_next( + as_head(first), as_head(second), key, NO_KEY_TWEAK); + consumer.join(); + + SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 2); + SNMALLOC_CHECK(completed.load(std::memory_order_acquire)); + check_empty(queue); + } + + void test_rejected_front() + { + Queue queue; + Object object; + std::atomic reject{true}; + std::atomic attempted{false}; + std::atomic completed{false}; + std::atomic callbacks{0}; + + SNMALLOC_CHECK( + queue.enqueue(as_head(object), as_head(object), domesticate)); + + auto validating_domesticate = [&reject, + &attempted](freelist::QueuePtr value) { + attempted.store(true, std::memory_order_release); + if (reject.load(std::memory_order_acquire)) + return freelist::HeadPtr(nullptr); + return domesticate(value); + }; + + std::thread consumer([&]() { + queue.drain_and_reset( + validating_domesticate, [&callbacks](freelist::HeadPtr) { + callbacks.fetch_add(1, std::memory_order_relaxed); + }); + completed.store(true, std::memory_order_release); + }); + + wait_until([&]() { return attempted.load(std::memory_order_acquire); }); + SNMALLOC_CHECK(!completed.load(std::memory_order_acquire)); + SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 0); + + reject.store(false, std::memory_order_release); + consumer.join(); + + SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 1); + SNMALLOC_CHECK(completed.load(std::memory_order_acquire)); + check_empty(queue); + } + + void test_callback_starts_replacement() + { + Queue queue; + Object objects[4]; + CallbackOrder order; + std::atomic first_callback{false}; + std::atomic release_callback{false}; + + freelist::Object::atomic_store_next( + as_head(objects[0]), as_head(objects[1]), key, NO_KEY_TWEAK); + SNMALLOC_CHECK( + queue.enqueue(as_head(objects[0]), as_head(objects[1]), domesticate)); + + std::thread consumer([&]() { + queue.drain_and_reset(domesticate, [&](freelist::HeadPtr value) { + order.add(value); + if (address_cast(value) == address_cast(as_head(objects[0]))) + { + first_callback.store(true, std::memory_order_release); + wait_until( + [&]() { return release_callback.load(std::memory_order_acquire); }); + } + }); + }); + + wait_until( + [&]() { return first_callback.load(std::memory_order_acquire); }); + + freelist::Object::atomic_store_next( + as_head(objects[2]), as_head(objects[3]), key, NO_KEY_TWEAK); + SNMALLOC_CHECK( + queue.enqueue(as_head(objects[2]), as_head(objects[3]), domesticate)); + release_callback.store(true, std::memory_order_release); + consumer.join(); + + SNMALLOC_CHECK(order.count == 2); + for (size_t i = 0; i < 2; i++) + { + SNMALLOC_CHECK( + address_cast(order.values[i]) == address_cast(as_head(objects[i]))); + } + SNMALLOC_CHECK(queue.back.load(stl::memory_order_relaxed) != nullptr); + + order.count = 0; + queue.drain_and_reset( + domesticate, [&order](freelist::HeadPtr value) { order.add(value); }); + SNMALLOC_CHECK(order.count == 2); + for (size_t i = 0; i < 2; i++) + { + SNMALLOC_CHECK( + address_cast(order.values[i]) == address_cast(as_head(objects[i + 2]))); + } + check_empty(queue); + } + + size_t find_object(Object* objects, size_t count, freelist::HeadPtr value) + { + for (size_t i = 0; i < count; i++) + { + if (address_cast(as_head(objects[i])) == address_cast(value)) + return i; + } + SNMALLOC_CHECK(false); + return count; + } + + void test_concurrent_producers(size_t object_count) + { + constexpr size_t producer_count = 4; + Queue queue; + auto objects = std::make_unique(object_count); + auto seen = std::make_unique[]>(object_count); + std::atomic next_object{0}; + std::atomic producers_live{producer_count}; + std::atomic true_results{0}; + + for (size_t i = 0; i < object_count; i++) + seen[i].store(0, std::memory_order_relaxed); + + std::vector producers; + for (size_t producer = 0; producer < producer_count; producer++) + { + producers.emplace_back([&]() { + while (true) + { + size_t first_index = + next_object.fetch_add(4, std::memory_order_relaxed); + if (first_index >= object_count) + break; + + size_t last_index = first_index + 3; + if (last_index >= object_count) + last_index = object_count - 1; + for (size_t i = first_index; i < last_index; i++) + { + freelist::Object::atomic_store_next( + as_head(objects[i]), as_head(objects[i + 1]), key, NO_KEY_TWEAK); + } + + if (queue.enqueue( + as_head(objects[first_index]), + as_head(objects[last_index]), + domesticate)) + { + true_results.fetch_add(1, std::memory_order_relaxed); + } + } + producers_live.fetch_sub(1, std::memory_order_release); + }); + } + + size_t callbacks = 0; + while (producers_live.load(std::memory_order_acquire) != 0 || + queue.back.load(stl::memory_order_relaxed) != nullptr) + { + queue.drain_and_reset(domesticate, [&](freelist::HeadPtr value) { + size_t index = find_object(objects.get(), object_count, value); + SNMALLOC_CHECK( + seen[index].fetch_add(1, std::memory_order_relaxed) == 0); + callbacks++; + }); + std::this_thread::yield(); + } + + for (auto& producer : producers) + producer.join(); + + SNMALLOC_CHECK(callbacks == object_count); + for (size_t i = 0; i < object_count; i++) + SNMALLOC_CHECK(seen[i].load(std::memory_order_relaxed) == 1); + SNMALLOC_CHECK(true_results.load(std::memory_order_relaxed) >= 1); + SNMALLOC_CHECK( + true_results.load(std::memory_order_relaxed) <= object_count); + check_empty(queue); + + Object replacement; + SNMALLOC_CHECK( + queue.enqueue(as_head(replacement), as_head(replacement), domesticate)); + SNMALLOC_CHECK(queue.back.load(stl::memory_order_relaxed) != nullptr); + queue.drain_and_reset(domesticate, [](freelist::HeadPtr) {}); + check_empty(queue); + } + + void test_reuse(size_t generations) + { + constexpr size_t object_count = 8; + constexpr size_t reserved = 2; + Queue queue; + Object objects[object_count]; + std::atomic states[object_count]; + std::atomic generation[object_count]; + auto seen = + std::make_unique[]>(object_count * generations); + std::atomic producer_done{false}; + + for (size_t i = 0; i < object_count; i++) + { + states[i].store(0, std::memory_order_relaxed); + generation[i].store(0, std::memory_order_relaxed); + } + for (size_t i = 0; i < object_count * generations; i++) + seen[i].store(0, std::memory_order_relaxed); + + std::thread producer([&]() { + for (size_t round = 0; round < generations; round++) + { + size_t claimed[object_count]; + size_t claim_count = 0; + while (claim_count < object_count) + { + for (size_t i = 0; i < object_count; i++) + { + size_t available = 0; + if (states[i].compare_exchange_strong( + available, + reserved, + std::memory_order_acq_rel, + std::memory_order_relaxed)) + { + generation[i].store(round, std::memory_order_relaxed); + claimed[claim_count++] = i; + } + } + if (claim_count < object_count) + std::this_thread::yield(); + } + + for (size_t first = 0; first < object_count; first += 3) + { + size_t last = first + 2; + if (last >= object_count) + last = object_count - 1; + for (size_t i = first; i < last; i++) + { + freelist::Object::atomic_store_next( + as_head(objects[claimed[i]]), + as_head(objects[claimed[i + 1]]), + key, + NO_KEY_TWEAK); + } + for (size_t i = first; i <= last; i++) + states[claimed[i]].store(1, std::memory_order_release); + queue.enqueue( + as_head(objects[claimed[first]]), + as_head(objects[claimed[last]]), + domesticate); + } + } + producer_done.store(true, std::memory_order_release); + }); + + size_t callbacks = 0; + while (!producer_done.load(std::memory_order_acquire) || + queue.back.load(stl::memory_order_relaxed) != nullptr) + { + queue.drain_and_reset(domesticate, [&](freelist::HeadPtr value) { + size_t index = find_object(objects, object_count, value); + SNMALLOC_CHECK(states[index].load(std::memory_order_acquire) == 1); + size_t round = generation[index].load(std::memory_order_relaxed); + SNMALLOC_CHECK( + seen[(round * object_count) + index].fetch_add( + 1, std::memory_order_relaxed) == 0); + callbacks++; + states[index].store(0, std::memory_order_release); + std::this_thread::yield(); + }); + } + producer.join(); + + SNMALLOC_CHECK(callbacks == object_count * generations); + for (size_t i = 0; i < object_count * generations; i++) + SNMALLOC_CHECK(seen[i].load(std::memory_order_relaxed) == 1); + check_empty(queue); + } +} + +int main(int argc, char** argv) +{ + bool heavy = false; + for (int i = 1; i < argc; i++) + { + if (std::strcmp(argv[i], "--heavy") == 0) + heavy = true; + } + + setup(); + test_empty(); + test_enqueue_and_iterate(); + test_dequeue_empty(); + test_dequeue_positions(); + test_dequeue_publication_gap(); + test_dequeue_checks_bound_first(); + test_remote_allocator_move_only_callback(); + test_front_publication_gap(); + test_successor_publication_gap(); + test_rejected_front(); + test_callback_starts_replacement(); + test_concurrent_producers(heavy ? 4096 : 256); + test_reuse(heavy ? 256 : 16); +} diff --git a/src/test/perf/msgpass/msgpass.cc b/src/test/perf/msgpass/msgpass.cc index b8c0d9d2b..cf312aeef 100644 --- a/src/test/perf/msgpass/msgpass.cc +++ b/src/test/perf/msgpass/msgpass.cc @@ -102,7 +102,9 @@ void consumer(const struct params* param, size_t qix) (queue_gate > param->N_CONSUMER)); chatty("Cl %zu fini\n", qix); - snmalloc::dealloc(myq.destroy().unsafe_ptr()); + myq.drain_and_reset(domesticate_nop, [](freelist::HeadPtr o) { + snmalloc::dealloc(o.as_void().unsafe_ptr()); + }); } void proxy(const struct params* param, size_t qix) @@ -134,7 +136,9 @@ void proxy(const struct params* param, size_t qix) chatty("Px %zu fini\n", qix); - snmalloc::dealloc(myq.destroy().unsafe_ptr()); + myq.drain_and_reset(domesticate_nop, [](freelist::HeadPtr o) { + snmalloc::dealloc(o.as_void().unsafe_ptr()); + }); queue_gate--; } @@ -218,11 +222,6 @@ int main(int argc, char** argv) auto* producer_threads = new std::thread[param.N_PRODUCER]; auto* queue_threads = new std::thread[param.N_QUEUE]; - for (size_t i = 0; i < param.N_QUEUE; i++) - { - param.msgqueue[i].init(); - } - producers_live = true; queue_gate = param.N_QUEUE; messages_outstanding = 0; From 41af95597ec9a3998cc182e34e3de04b7aabbd86 Mon Sep 17 00:00:00 2001 From: Matthew Parkinson Date: Thu, 17 Sep 2026 13:30:18 +0100 Subject: [PATCH 2/4] Avoid MSVC shadow warning in queue test Rename the focused test's namespace-scope freelist key so instantiated allocator parameters named key do not trigger C4459 under warnings-as-errors. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 75ccdc6c-2024-4d81-b22f-f1b52fbe29cf --- .../func/freelist_mpscq/freelist_mpscq.cc | 42 +++++++++++-------- 1 file changed, 25 insertions(+), 17 deletions(-) diff --git a/src/test/func/freelist_mpscq/freelist_mpscq.cc b/src/test/func/freelist_mpscq/freelist_mpscq.cc index 06b893b8f..5fb9a4896 100644 --- a/src/test/func/freelist_mpscq/freelist_mpscq.cc +++ b/src/test/func/freelist_mpscq/freelist_mpscq.cc @@ -11,8 +11,8 @@ using namespace snmalloc; namespace { - FreeListKey key{0x1234, 0x5678, 0x9abc}; - using Queue = FreeListMPSCQ; + FreeListKey queue_key{0x1234, 0x5678, 0x9abc}; + using Queue = FreeListMPSCQ; using Object = freelist::Object::T<>; freelist::HeadPtr domesticate(freelist::QueuePtr p) @@ -80,7 +80,7 @@ namespace !queue.enqueue(as_head(second), as_head(second), domesticate)); freelist::Object::atomic_store_next( - as_head(batch_first), as_head(batch_last), key, NO_KEY_TWEAK); + 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)); @@ -127,7 +127,8 @@ namespace SNMALLOC_CHECK(callbacks == 0); check_empty(queue); - freelist::Object::atomic_store_null(as_head(pending), key, NO_KEY_TWEAK); + freelist::Object::atomic_store_null( + as_head(pending), queue_key, NO_KEY_TWEAK); auto prev = queue.back.exchange( capptr_rewild(as_head(pending)), stl::memory_order_acq_rel); SNMALLOC_CHECK(prev == nullptr); @@ -150,9 +151,9 @@ namespace CallbackOrder order; freelist::Object::atomic_store_next( - as_head(objects[0]), as_head(objects[1]), key, NO_KEY_TWEAK); + as_head(objects[0]), as_head(objects[1]), queue_key, NO_KEY_TWEAK); freelist::Object::atomic_store_next( - as_head(objects[1]), as_head(objects[2]), key, NO_KEY_TWEAK); + as_head(objects[1]), as_head(objects[2]), queue_key, NO_KEY_TWEAK); SNMALLOC_CHECK( queue.enqueue(as_head(objects[0]), as_head(objects[2]), domesticate)); @@ -201,7 +202,8 @@ namespace SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(first), domesticate)); - freelist::Object::atomic_store_null(as_head(second), key, NO_KEY_TWEAK); + freelist::Object::atomic_store_null( + as_head(second), queue_key, NO_KEY_TWEAK); auto prev = queue.back.exchange( capptr_rewild(as_head(second)), stl::memory_order_acq_rel); SNMALLOC_CHECK(address_cast(prev) == address_cast(as_head(first))); @@ -217,7 +219,7 @@ namespace address_cast(as_head(first))); freelist::Object::atomic_store_next( - as_head(first), as_head(second), key, NO_KEY_TWEAK); + as_head(first), as_head(second), queue_key, NO_KEY_TWEAK); queue.drain_and_reset( domesticate, [&callbacks](freelist::HeadPtr) { callbacks++; }); SNMALLOC_CHECK(callbacks == 2); @@ -234,7 +236,7 @@ namespace SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(first), domesticate)); freelist::Object::atomic_store_next( - as_head(first), as_head(outside), key, NO_KEY_TWEAK); + as_head(first), as_head(outside), queue_key, NO_KEY_TWEAK); auto counting_domesticate = [&domesticates](freelist::QueuePtr value) -> freelist::HeadPtr { @@ -257,7 +259,8 @@ namespace address_cast(queue.back.load(stl::memory_order_relaxed)) == address_cast(as_head(first))); - freelist::Object::atomic_store_null(as_head(first), key, NO_KEY_TWEAK); + freelist::Object::atomic_store_null( + as_head(first), queue_key, NO_KEY_TWEAK); queue.drain_and_reset(domesticate, [](freelist::HeadPtr) {}); check_empty(queue); } @@ -291,7 +294,8 @@ namespace std::atomic completed{false}; std::atomic callbacks{0}; - freelist::Object::atomic_store_null(as_head(object), key, NO_KEY_TWEAK); + freelist::Object::atomic_store_null( + as_head(object), queue_key, NO_KEY_TWEAK); auto prev = queue.back.exchange( capptr_rewild(as_head(object)), stl::memory_order_acq_rel); SNMALLOC_CHECK(prev == nullptr); @@ -329,7 +333,8 @@ namespace SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(first), domesticate)); - freelist::Object::atomic_store_null(as_head(second), key, NO_KEY_TWEAK); + freelist::Object::atomic_store_null( + as_head(second), queue_key, NO_KEY_TWEAK); auto prev = queue.back.exchange( capptr_rewild(as_head(second)), stl::memory_order_acq_rel); SNMALLOC_CHECK(address_cast(prev) == address_cast(as_head(first))); @@ -348,7 +353,7 @@ namespace SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 0); freelist::Object::atomic_store_next( - as_head(first), as_head(second), key, NO_KEY_TWEAK); + as_head(first), as_head(second), queue_key, NO_KEY_TWEAK); consumer.join(); SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 2); @@ -405,7 +410,7 @@ namespace std::atomic release_callback{false}; freelist::Object::atomic_store_next( - as_head(objects[0]), as_head(objects[1]), key, NO_KEY_TWEAK); + as_head(objects[0]), as_head(objects[1]), queue_key, NO_KEY_TWEAK); SNMALLOC_CHECK( queue.enqueue(as_head(objects[0]), as_head(objects[1]), domesticate)); @@ -425,7 +430,7 @@ namespace [&]() { return first_callback.load(std::memory_order_acquire); }); freelist::Object::atomic_store_next( - as_head(objects[2]), as_head(objects[3]), key, NO_KEY_TWEAK); + as_head(objects[2]), as_head(objects[3]), queue_key, NO_KEY_TWEAK); SNMALLOC_CHECK( queue.enqueue(as_head(objects[2]), as_head(objects[3]), domesticate)); release_callback.store(true, std::memory_order_release); @@ -492,7 +497,10 @@ namespace for (size_t i = first_index; i < last_index; i++) { freelist::Object::atomic_store_next( - as_head(objects[i]), as_head(objects[i + 1]), key, NO_KEY_TWEAK); + as_head(objects[i]), + as_head(objects[i + 1]), + queue_key, + NO_KEY_TWEAK); } if (queue.enqueue( @@ -593,7 +601,7 @@ namespace freelist::Object::atomic_store_next( as_head(objects[claimed[i]]), as_head(objects[claimed[i + 1]]), - key, + queue_key, NO_KEY_TWEAK); } for (size_t i = first; i <= last; i++) From d7f262e835ac5fea8631daa8a53b1a1974ae9fbc Mon Sep 17 00:00:00 2001 From: Matthew Parkinson Date: Fri, 18 Sep 2026 09:09:12 +0100 Subject: [PATCH 3/4] Separate queue head domestication Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 75ccdc6c-2024-4d81-b22f-f1b52fbe29cf --- src/snmalloc/mem/corealloc.h | 13 +- src/snmalloc/mem/freelist_queue.h | 24 ++- src/snmalloc/mem/remoteallocator.h | 15 +- .../func/freelist_mpscq/freelist_mpscq.cc | 194 ++++++++++++++---- src/test/perf/msgpass/msgpass.cc | 14 +- 5 files changed, 199 insertions(+), 61 deletions(-) diff --git a/src/snmalloc/mem/corealloc.h b/src/snmalloc/mem/corealloc.h index 4e5b2872a..a28357cd9 100644 --- a/src/snmalloc/mem/corealloc.h +++ b/src/snmalloc/mem/corealloc.h @@ -1412,7 +1412,18 @@ namespace snmalloc handle_dealloc_remote(entry, m, need_post, domesticate, bytes_flushed); }; - message_queue().drain_and_reset(domesticate, cb); + if constexpr (Config::Options.QueueHeadsAreTame) + { + auto domesticate_first = + [](freelist::QueuePtr p) SNMALLOC_FAST_PATH_LAMBDA { + return freelist::HeadPtr::unsafe_from(p.unsafe_ptr()); + }; + message_queue().drain_and_reset(domesticate_first, domesticate, cb); + } + else + { + message_queue().drain_and_reset(domesticate, domesticate, cb); + } auto& key = freelist::Object::key_root; diff --git a/src/snmalloc/mem/freelist_queue.h b/src/snmalloc/mem/freelist_queue.h index f79d6be54..33f7fa489 100644 --- a/src/snmalloc/mem/freelist_queue.h +++ b/src/snmalloc/mem/freelist_queue.h @@ -16,10 +16,9 @@ namespace snmalloc * for the client to reach as the Pagemap, which we trust to store not just * Tame CapPtr<>s but raw C++ pointers. * - * Where necessary, dequeue exposes two domesticator callbacks and is careful - * to use one for the front value and the other for pointers read from the - * queue itself. Draining uses its queue domesticator for both the front - * value and links in the chain. Specifically, + * Where necessary, dequeue and draining expose two domesticator callbacks + * and are careful to use one for the front value and the other for pointers + * read from the queue itself. Specifically, * * * `domesticate_head` is used for the MPSCQ pointers used to reach into * the chain of objects @@ -80,12 +79,21 @@ namespace snmalloc * A producer that has exchanged back must eventually publish front or its * predecessor link; otherwise this operation waits indefinitely. * + * domesticate_head applies only to values loaded from front. + * domesticate_queue applies to successors decoded from message links. + * * The queue is reset before the first callback. The callback may therefore * release or re-enqueue an object; any re-enqueue belongs to the * replacement chain and is not consumed by this invocation. */ - template - void drain_and_reset(Domesticator_queue domesticate, Cb cb) + template< + typename Domesticator_head, + typename Domesticator_queue, + typename Cb> + void drain_and_reset( + Domesticator_head domesticate_head, + Domesticator_queue domesticate_queue, + Cb cb) { // After reuse, acquire the release sequence headed by the preceding // reset, so front cannot observe an earlier queue generation. @@ -97,7 +105,7 @@ namespace snmalloc { auto raw = front.load(stl::memory_order_acquire); if (raw != nullptr) - curr = domesticate(raw); + curr = domesticate_head(raw); if (curr == nullptr) Aal::pause(); } while (curr == nullptr); @@ -115,7 +123,7 @@ namespace snmalloc while (true) { - curr = process_chain(curr, target, domesticate, process); + curr = process_chain(curr, target, domesticate_queue, process); if (address_cast(curr) == address_cast(target)) break; Aal::pause(); diff --git a/src/snmalloc/mem/remoteallocator.h b/src/snmalloc/mem/remoteallocator.h index 1aab15af2..efd4ba942 100644 --- a/src/snmalloc/mem/remoteallocator.h +++ b/src/snmalloc/mem/remoteallocator.h @@ -321,14 +321,23 @@ namespace snmalloc list.invariant(); } - template - void drain_and_reset(Domesticator_queue domesticate, Cb cb) + template< + typename Domesticator_head, + typename Domesticator_queue, + typename Cb> + void drain_and_reset( + Domesticator_head domesticate_head, + Domesticator_queue domesticate_queue, + Cb cb) { auto cbwrap = [&cb](freelist::HeadPtr p) SNMALLOC_FAST_PATH_LAMBDA { cb(RemoteMessage::from_message_link(p)); }; - return list.drain_and_reset(stl::move(domesticate), stl::move(cbwrap)); + return list.drain_and_reset( + stl::move(domesticate_head), + stl::move(domesticate_queue), + stl::move(cbwrap)); } inline bool can_dequeue() diff --git a/src/test/func/freelist_mpscq/freelist_mpscq.cc b/src/test/func/freelist_mpscq/freelist_mpscq.cc index 5fb9a4896..7375f8f2f 100644 --- a/src/test/func/freelist_mpscq/freelist_mpscq.cc +++ b/src/test/func/freelist_mpscq/freelist_mpscq.cc @@ -59,7 +59,9 @@ namespace size_t callbacks = 0; queue.drain_and_reset( - domesticate, [&callbacks](freelist::HeadPtr) { callbacks++; }); + domesticate, domesticate, [&callbacks](freelist::HeadPtr) { + callbacks++; + }); SNMALLOC_CHECK(callbacks == 0); check_empty(queue); @@ -85,7 +87,9 @@ namespace !queue.enqueue(as_head(batch_first), as_head(batch_last), domesticate)); queue.drain_and_reset( - domesticate, [&order](freelist::HeadPtr value) { order.add(value); }); + domesticate, domesticate, [&order](freelist::HeadPtr value) { + order.add(value); + }); SNMALLOC_CHECK(order.count == 4); SNMALLOC_CHECK( @@ -104,7 +108,9 @@ namespace order.count = 0; queue.drain_and_reset( - domesticate, [&order](freelist::HeadPtr value) { order.add(value); }); + domesticate, domesticate, [&order](freelist::HeadPtr value) { + order.add(value); + }); SNMALLOC_CHECK(order.count == 2); SNMALLOC_CHECK( address_cast(order.values[0]) == address_cast(as_head(replacement))); @@ -186,7 +192,9 @@ namespace order.count = 0; queue.drain_and_reset( - domesticate, [&order](freelist::HeadPtr value) { order.add(value); }); + domesticate, domesticate, [&order](freelist::HeadPtr value) { + order.add(value); + }); SNMALLOC_CHECK(order.count == 1); SNMALLOC_CHECK( address_cast(order.values[0]) == address_cast(as_head(objects[2]))); @@ -221,7 +229,9 @@ namespace freelist::Object::atomic_store_next( as_head(first), as_head(second), queue_key, NO_KEY_TWEAK); queue.drain_and_reset( - domesticate, [&callbacks](freelist::HeadPtr) { callbacks++; }); + domesticate, domesticate, [&callbacks](freelist::HeadPtr) { + callbacks++; + }); SNMALLOC_CHECK(callbacks == 2); check_empty(queue); } @@ -261,7 +271,53 @@ namespace freelist::Object::atomic_store_null( as_head(first), queue_key, NO_KEY_TWEAK); - queue.drain_and_reset(domesticate, [](freelist::HeadPtr) {}); + queue.drain_and_reset(domesticate, domesticate, [](freelist::HeadPtr) {}); + check_empty(queue); + } + + void test_drain_uses_distinct_domesticators() + { + Queue queue; + Object objects[3]; + CallbackOrder order; + size_t head_domesticates = 0; + size_t queue_domesticates = 0; + + freelist::Object::atomic_store_next( + as_head(objects[0]), as_head(objects[1]), queue_key, NO_KEY_TWEAK); + freelist::Object::atomic_store_next( + as_head(objects[1]), as_head(objects[2]), queue_key, NO_KEY_TWEAK); + SNMALLOC_CHECK( + queue.enqueue(as_head(objects[0]), as_head(objects[2]), domesticate)); + + auto domesticate_head = [&](freelist::QueuePtr value) -> freelist::HeadPtr { + head_domesticates++; + SNMALLOC_CHECK(address_cast(value) == address_cast(as_head(objects[0]))); + return domesticate(value); + }; + auto domesticate_queue = + [&](freelist::QueuePtr value) -> freelist::HeadPtr { + queue_domesticates++; + SNMALLOC_CHECK(address_cast(value) != address_cast(as_head(objects[0]))); + SNMALLOC_CHECK( + (address_cast(value) == address_cast(as_head(objects[1]))) || + (address_cast(value) == address_cast(as_head(objects[2])))); + return domesticate(value); + }; + + queue.drain_and_reset( + domesticate_head, domesticate_queue, [&order](freelist::HeadPtr value) { + order.add(value); + }); + + SNMALLOC_CHECK(head_domesticates == 1); + SNMALLOC_CHECK(queue_domesticates == 2); + SNMALLOC_CHECK(order.count == 3); + for (size_t i = 0; i < 3; i++) + { + SNMALLOC_CHECK( + address_cast(order.values[i]) == address_cast(as_head(objects[i]))); + } check_empty(queue); } @@ -278,7 +334,7 @@ namespace &callbacks](capptr::Alloc) mutable { callbacks += (*token)++; }; - remote.drain_and_reset(domesticate, std::move(cb)); + remote.drain_and_reset(domesticate, domesticate, std::move(cb)); SNMALLOC_CHECK(callbacks == 1); SNMALLOC_CHECK(remote.list.back.load(stl::memory_order_relaxed) == nullptr); @@ -302,10 +358,11 @@ namespace std::thread consumer([&]() { entered.store(true, std::memory_order_release); - queue.drain_and_reset(domesticate, [&](freelist::HeadPtr value) { - SNMALLOC_CHECK(address_cast(value) == address_cast(as_head(object))); - callbacks.fetch_add(1, std::memory_order_relaxed); - }); + queue.drain_and_reset( + domesticate, domesticate, [&](freelist::HeadPtr value) { + SNMALLOC_CHECK(address_cast(value) == address_cast(as_head(object))); + callbacks.fetch_add(1, std::memory_order_relaxed); + }); completed.store(true, std::memory_order_release); }); @@ -341,9 +398,10 @@ namespace std::thread consumer([&]() { entered.store(true, std::memory_order_release); - queue.drain_and_reset(domesticate, [&callbacks](freelist::HeadPtr) { - callbacks.fetch_add(1, std::memory_order_relaxed); - }); + queue.drain_and_reset( + domesticate, domesticate, [&callbacks](freelist::HeadPtr) { + callbacks.fetch_add(1, std::memory_order_relaxed); + }); completed.store(true, std::memory_order_release); }); @@ -383,7 +441,7 @@ namespace std::thread consumer([&]() { queue.drain_and_reset( - validating_domesticate, [&callbacks](freelist::HeadPtr) { + validating_domesticate, domesticate, [&callbacks](freelist::HeadPtr) { callbacks.fetch_add(1, std::memory_order_relaxed); }); completed.store(true, std::memory_order_release); @@ -401,6 +459,48 @@ namespace check_empty(queue); } + void test_rejected_successor() + { + Queue queue; + Object first; + Object second; + std::atomic reject{true}; + std::atomic attempted{false}; + std::atomic completed{false}; + std::atomic callbacks{0}; + + freelist::Object::atomic_store_next( + as_head(first), as_head(second), queue_key, NO_KEY_TWEAK); + SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(second), domesticate)); + + auto validating_domesticate = [&reject, + &attempted](freelist::QueuePtr value) { + attempted.store(true, std::memory_order_release); + if (reject.load(std::memory_order_acquire)) + return freelist::HeadPtr(nullptr); + return domesticate(value); + }; + + std::thread consumer([&]() { + queue.drain_and_reset( + domesticate, validating_domesticate, [&callbacks](freelist::HeadPtr) { + callbacks.fetch_add(1, std::memory_order_relaxed); + }); + completed.store(true, std::memory_order_release); + }); + + wait_until([&]() { return attempted.load(std::memory_order_acquire); }); + SNMALLOC_CHECK(!completed.load(std::memory_order_acquire)); + SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 0); + + reject.store(false, std::memory_order_release); + consumer.join(); + + SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 2); + SNMALLOC_CHECK(completed.load(std::memory_order_acquire)); + check_empty(queue); + } + void test_callback_starts_replacement() { Queue queue; @@ -415,15 +515,17 @@ namespace queue.enqueue(as_head(objects[0]), as_head(objects[1]), domesticate)); std::thread consumer([&]() { - queue.drain_and_reset(domesticate, [&](freelist::HeadPtr value) { - order.add(value); - if (address_cast(value) == address_cast(as_head(objects[0]))) - { - first_callback.store(true, std::memory_order_release); - wait_until( - [&]() { return release_callback.load(std::memory_order_acquire); }); - } - }); + queue.drain_and_reset( + domesticate, domesticate, [&](freelist::HeadPtr value) { + order.add(value); + if (address_cast(value) == address_cast(as_head(objects[0]))) + { + first_callback.store(true, std::memory_order_release); + wait_until([&]() { + return release_callback.load(std::memory_order_acquire); + }); + } + }); }); wait_until( @@ -446,7 +548,9 @@ namespace order.count = 0; queue.drain_and_reset( - domesticate, [&order](freelist::HeadPtr value) { order.add(value); }); + domesticate, domesticate, [&order](freelist::HeadPtr value) { + order.add(value); + }); SNMALLOC_CHECK(order.count == 2); for (size_t i = 0; i < 2; i++) { @@ -519,12 +623,13 @@ namespace while (producers_live.load(std::memory_order_acquire) != 0 || queue.back.load(stl::memory_order_relaxed) != nullptr) { - queue.drain_and_reset(domesticate, [&](freelist::HeadPtr value) { - size_t index = find_object(objects.get(), object_count, value); - SNMALLOC_CHECK( - seen[index].fetch_add(1, std::memory_order_relaxed) == 0); - callbacks++; - }); + queue.drain_and_reset( + domesticate, domesticate, [&](freelist::HeadPtr value) { + size_t index = find_object(objects.get(), object_count, value); + SNMALLOC_CHECK( + seen[index].fetch_add(1, std::memory_order_relaxed) == 0); + callbacks++; + }); std::this_thread::yield(); } @@ -543,7 +648,7 @@ namespace SNMALLOC_CHECK( queue.enqueue(as_head(replacement), as_head(replacement), domesticate)); SNMALLOC_CHECK(queue.back.load(stl::memory_order_relaxed) != nullptr); - queue.drain_and_reset(domesticate, [](freelist::HeadPtr) {}); + queue.drain_and_reset(domesticate, domesticate, [](freelist::HeadPtr) {}); check_empty(queue); } @@ -619,17 +724,18 @@ namespace while (!producer_done.load(std::memory_order_acquire) || queue.back.load(stl::memory_order_relaxed) != nullptr) { - queue.drain_and_reset(domesticate, [&](freelist::HeadPtr value) { - size_t index = find_object(objects, object_count, value); - SNMALLOC_CHECK(states[index].load(std::memory_order_acquire) == 1); - size_t round = generation[index].load(std::memory_order_relaxed); - SNMALLOC_CHECK( - seen[(round * object_count) + index].fetch_add( - 1, std::memory_order_relaxed) == 0); - callbacks++; - states[index].store(0, std::memory_order_release); - std::this_thread::yield(); - }); + queue.drain_and_reset( + domesticate, domesticate, [&](freelist::HeadPtr value) { + size_t index = find_object(objects, object_count, value); + SNMALLOC_CHECK(states[index].load(std::memory_order_acquire) == 1); + size_t round = generation[index].load(std::memory_order_relaxed); + SNMALLOC_CHECK( + seen[(round * object_count) + index].fetch_add( + 1, std::memory_order_relaxed) == 0); + callbacks++; + states[index].store(0, std::memory_order_release); + std::this_thread::yield(); + }); } producer.join(); @@ -656,10 +762,12 @@ int main(int argc, char** argv) test_dequeue_positions(); test_dequeue_publication_gap(); test_dequeue_checks_bound_first(); + test_drain_uses_distinct_domesticators(); test_remote_allocator_move_only_callback(); test_front_publication_gap(); test_successor_publication_gap(); test_rejected_front(); + test_rejected_successor(); test_callback_starts_replacement(); test_concurrent_producers(heavy ? 4096 : 256); test_reuse(heavy ? 256 : 16); diff --git a/src/test/perf/msgpass/msgpass.cc b/src/test/perf/msgpass/msgpass.cc index cf312aeef..7344161b7 100644 --- a/src/test/perf/msgpass/msgpass.cc +++ b/src/test/perf/msgpass/msgpass.cc @@ -102,9 +102,10 @@ void consumer(const struct params* param, size_t qix) (queue_gate > param->N_CONSUMER)); chatty("Cl %zu fini\n", qix); - myq.drain_and_reset(domesticate_nop, [](freelist::HeadPtr o) { - snmalloc::dealloc(o.as_void().unsafe_ptr()); - }); + myq.drain_and_reset( + domesticate_nop, domesticate_nop, [](freelist::HeadPtr o) { + snmalloc::dealloc(o.as_void().unsafe_ptr()); + }); } void proxy(const struct params* param, size_t qix) @@ -136,9 +137,10 @@ void proxy(const struct params* param, size_t qix) chatty("Px %zu fini\n", qix); - myq.drain_and_reset(domesticate_nop, [](freelist::HeadPtr o) { - snmalloc::dealloc(o.as_void().unsafe_ptr()); - }); + myq.drain_and_reset( + domesticate_nop, domesticate_nop, [](freelist::HeadPtr o) { + snmalloc::dealloc(o.as_void().unsafe_ptr()); + }); queue_gate--; } From d4b55dadf5891a6f2f1561936d45263dd7b25e15 Mon Sep 17 00:00:00 2001 From: Matthew Parkinson Date: Fri, 18 Sep 2026 14:29:20 +0100 Subject: [PATCH 4/4] Simplify focused MPSC queue tests Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 75ccdc6c-2024-4d81-b22f-f1b52fbe29cf --- .../func/freelist_mpscq/freelist_mpscq.cc | 790 ++++-------------- 1 file changed, 155 insertions(+), 635 deletions(-) diff --git a/src/test/func/freelist_mpscq/freelist_mpscq.cc b/src/test/func/freelist_mpscq/freelist_mpscq.cc index 7375f8f2f..3c1af0b24 100644 --- a/src/test/func/freelist_mpscq/freelist_mpscq.cc +++ b/src/test/func/freelist_mpscq/freelist_mpscq.cc @@ -1,11 +1,5 @@ -#include -#include -#include -#include #include #include -#include -#include using namespace snmalloc; @@ -25,49 +19,67 @@ namespace return freelist::HeadPtr::unsafe_from(&object); } - template - void wait_until(Predicate predicate) - { - auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(10); - while (!predicate()) - { - SNMALLOC_CHECK(std::chrono::steady_clock::now() < deadline); - std::this_thread::yield(); - } - } - + template struct CallbackOrder { - freelist::HeadPtr values[8]; + T values[Size]; size_t count = 0; - void add(freelist::HeadPtr value) + void add(T value) { + SNMALLOC_CHECK(count < Size); values[count++] = value; } }; - void check_empty(Queue& queue) - { - SNMALLOC_CHECK(queue.back.load(stl::memory_order_relaxed) == nullptr); - SNMALLOC_CHECK(queue.front.load(stl::memory_order_relaxed) == nullptr); - } - + /** + * An empty drain invokes no callback. An empty dequeue domesticates the null + * head, but applies neither the queue domesticator nor the callback. + */ void test_empty() { Queue queue; + size_t head_domesticates = 0; + size_t queue_domesticates = 0; size_t callbacks = 0; queue.drain_and_reset( domesticate, domesticate, [&callbacks](freelist::HeadPtr) { callbacks++; }); + SNMALLOC_CHECK(callbacks == 0); + + auto domesticate_head = + [&head_domesticates](freelist::QueuePtr value) -> freelist::HeadPtr { + head_domesticates++; + SNMALLOC_CHECK(value == nullptr); + return nullptr; + }; + auto domesticate_queue = + [&queue_domesticates](freelist::QueuePtr) -> freelist::HeadPtr { + queue_domesticates++; + return nullptr; + }; + + queue.dequeue( + domesticate_head, domesticate_queue, [&callbacks](freelist::HeadPtr) { + callbacks++; + return true; + }); + SNMALLOC_CHECK(head_domesticates == 1); + SNMALLOC_CHECK(queue_domesticates == 0); SNMALLOC_CHECK(callbacks == 0); - check_empty(queue); } - void test_enqueue_and_iterate() + /** + * 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. + * Dequeue retains the object at back; drain delivers it and resets the queue + * for reuse. + */ + void test_queue_contract() { Queue queue; Object first; @@ -75,7 +87,7 @@ namespace Object batch_first; Object batch_last; Object replacement; - CallbackOrder order; + CallbackOrder order; SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(first), domesticate)); SNMALLOC_CHECK( @@ -86,689 +98,197 @@ namespace SNMALLOC_CHECK( !queue.enqueue(as_head(batch_first), as_head(batch_last), domesticate)); - queue.drain_and_reset( - domesticate, domesticate, [&order](freelist::HeadPtr value) { - order.add(value); - }); - - SNMALLOC_CHECK(order.count == 4); - 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( - address_cast(order.values[2]) == address_cast(as_head(batch_first))); - SNMALLOC_CHECK( - address_cast(order.values[3]) == address_cast(as_head(batch_last))); - check_empty(queue); - - SNMALLOC_CHECK( - queue.enqueue(as_head(replacement), as_head(replacement), domesticate)); - queue.enqueue(as_head(second), as_head(second), domesticate); - - order.count = 0; - queue.drain_and_reset( - domesticate, domesticate, [&order](freelist::HeadPtr value) { - order.add(value); - }); - SNMALLOC_CHECK(order.count == 2); - SNMALLOC_CHECK( - address_cast(order.values[0]) == address_cast(as_head(replacement))); - SNMALLOC_CHECK( - address_cast(order.values[1]) == address_cast(as_head(second))); - check_empty(queue); - } - - void test_dequeue_empty() - { - Queue queue; - Object pending; - size_t callbacks = 0; - auto cb = [&callbacks](freelist::HeadPtr) { - callbacks++; - return true; - }; - - queue.dequeue(domesticate, domesticate, cb); - SNMALLOC_CHECK(callbacks == 0); - check_empty(queue); - - freelist::Object::atomic_store_null( - as_head(pending), queue_key, NO_KEY_TWEAK); - auto prev = queue.back.exchange( - capptr_rewild(as_head(pending)), stl::memory_order_acq_rel); - SNMALLOC_CHECK(prev == nullptr); - - queue.dequeue(domesticate, domesticate, cb); - SNMALLOC_CHECK(callbacks == 0); - SNMALLOC_CHECK(queue.front.load(stl::memory_order_relaxed) == nullptr); - SNMALLOC_CHECK( - address_cast(queue.back.load(stl::memory_order_relaxed)) == - address_cast(as_head(pending))); - - queue.back.store(nullptr, stl::memory_order_relaxed); - check_empty(queue); - } - - void test_dequeue_positions() - { - Queue queue; - Object objects[3]; - CallbackOrder order; - - freelist::Object::atomic_store_next( - as_head(objects[0]), as_head(objects[1]), queue_key, NO_KEY_TWEAK); - freelist::Object::atomic_store_next( - as_head(objects[1]), as_head(objects[2]), queue_key, NO_KEY_TWEAK); - SNMALLOC_CHECK( - queue.enqueue(as_head(objects[0]), as_head(objects[2]), domesticate)); - queue.dequeue(domesticate, domesticate, [&order](freelist::HeadPtr value) { order.add(value); return false; }); - SNMALLOC_CHECK(order.count == 1); SNMALLOC_CHECK( - address_cast(order.values[0]) == address_cast(as_head(objects[0]))); - SNMALLOC_CHECK( - address_cast(queue.front.load(stl::memory_order_relaxed)) == - address_cast(as_head(objects[1]))); + address_cast(order.values[0]) == address_cast(as_head(first))); queue.dequeue(domesticate, domesticate, [&order](freelist::HeadPtr value) { order.add(value); return true; }); - - SNMALLOC_CHECK(order.count == 2); - SNMALLOC_CHECK( - address_cast(order.values[1]) == address_cast(as_head(objects[1]))); + SNMALLOC_CHECK(order.count == 3); SNMALLOC_CHECK( - address_cast(queue.front.load(stl::memory_order_relaxed)) == - address_cast(as_head(objects[2]))); + address_cast(order.values[1]) == address_cast(as_head(second))); SNMALLOC_CHECK( - address_cast(queue.back.load(stl::memory_order_relaxed)) == - address_cast(as_head(objects[2]))); + address_cast(order.values[2]) == address_cast(as_head(batch_first))); - order.count = 0; queue.drain_and_reset( domesticate, domesticate, [&order](freelist::HeadPtr value) { order.add(value); }); - SNMALLOC_CHECK(order.count == 1); + SNMALLOC_CHECK(order.count == 4); SNMALLOC_CHECK( - address_cast(order.values[0]) == address_cast(as_head(objects[2]))); - check_empty(queue); - } - - void test_dequeue_publication_gap() - { - Queue queue; - Object first; - Object second; - size_t callbacks = 0; - - SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(first), domesticate)); - - freelist::Object::atomic_store_null( - as_head(second), queue_key, NO_KEY_TWEAK); - auto prev = queue.back.exchange( - capptr_rewild(as_head(second)), stl::memory_order_acq_rel); - SNMALLOC_CHECK(address_cast(prev) == address_cast(as_head(first))); - - queue.dequeue(domesticate, domesticate, [&callbacks](freelist::HeadPtr) { - callbacks++; - return true; - }); + address_cast(order.values[3]) == address_cast(as_head(batch_last))); - SNMALLOC_CHECK(callbacks == 0); SNMALLOC_CHECK( - address_cast(queue.front.load(stl::memory_order_relaxed)) == - address_cast(as_head(first))); - - freelist::Object::atomic_store_next( - as_head(first), as_head(second), queue_key, NO_KEY_TWEAK); + queue.enqueue(as_head(replacement), as_head(replacement), domesticate)); queue.drain_and_reset( - domesticate, domesticate, [&callbacks](freelist::HeadPtr) { - callbacks++; + domesticate, domesticate, [&order](freelist::HeadPtr value) { + order.add(value); }); - SNMALLOC_CHECK(callbacks == 2); - check_empty(queue); + SNMALLOC_CHECK(order.count == 5); + SNMALLOC_CHECK( + address_cast(order.values[4]) == address_cast(as_head(replacement))); } + /** + * The captured back bounds a dequeue before any successor link is read. A + * link beyond that back belongs to an enqueue this dequeue does not cover, so + * it is neither domesticated nor followed. The subsequent drain confirms + * that the retained object remains queued. + */ void test_dequeue_checks_bound_first() { Queue queue; - Object first; + Object tail; Object outside; - size_t callbacks = 0; - size_t domesticates = 0; + size_t queue_domesticates = 0; + size_t dequeue_callbacks = 0; + size_t drain_callbacks = 0; - SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(first), domesticate)); + SNMALLOC_CHECK(queue.enqueue(as_head(tail), as_head(tail), domesticate)); freelist::Object::atomic_store_next( - as_head(first), as_head(outside), queue_key, NO_KEY_TWEAK); + as_head(tail), as_head(outside), queue_key, NO_KEY_TWEAK); auto counting_domesticate = - [&domesticates](freelist::QueuePtr value) -> freelist::HeadPtr { - domesticates++; + [&queue_domesticates](freelist::QueuePtr value) -> freelist::HeadPtr { + queue_domesticates++; return domesticate(value); }; queue.dequeue( - domesticate, counting_domesticate, [&callbacks](freelist::HeadPtr) { - callbacks++; + domesticate, + counting_domesticate, + [&dequeue_callbacks](freelist::HeadPtr) { + dequeue_callbacks++; return true; }); - SNMALLOC_CHECK(callbacks == 0); - SNMALLOC_CHECK(domesticates == 0); - SNMALLOC_CHECK( - address_cast(queue.front.load(stl::memory_order_relaxed)) == - address_cast(as_head(first))); - SNMALLOC_CHECK( - address_cast(queue.back.load(stl::memory_order_relaxed)) == - address_cast(as_head(first))); + SNMALLOC_CHECK(queue_domesticates == 0); + SNMALLOC_CHECK(dequeue_callbacks == 0); - freelist::Object::atomic_store_null( - as_head(first), queue_key, NO_KEY_TWEAK); - queue.drain_and_reset(domesticate, domesticate, [](freelist::HeadPtr) {}); - check_empty(queue); + queue.drain_and_reset( + domesticate, domesticate, [&](freelist::HeadPtr value) { + SNMALLOC_CHECK(address_cast(value) == address_cast(as_head(tail))); + drain_callbacks++; + }); + SNMALLOC_CHECK(drain_callbacks == 1); } - void test_drain_uses_distinct_domesticators() + /** + * Drain resets the queue before invoking callbacks. An enqueue from a + * callback therefore starts a replacement chain that is not consumed until + * the next drain. + */ + void test_callback_starts_replacement() { Queue queue; - Object objects[3]; - CallbackOrder order; - size_t head_domesticates = 0; - size_t queue_domesticates = 0; + Object old_first; + Object old_last; + Object replacement_first; + Object replacement_last; + CallbackOrder order; + bool replacement_started = false; freelist::Object::atomic_store_next( - as_head(objects[0]), as_head(objects[1]), queue_key, NO_KEY_TWEAK); - freelist::Object::atomic_store_next( - as_head(objects[1]), as_head(objects[2]), queue_key, NO_KEY_TWEAK); + as_head(old_first), as_head(old_last), queue_key, NO_KEY_TWEAK); SNMALLOC_CHECK( - queue.enqueue(as_head(objects[0]), as_head(objects[2]), domesticate)); - - auto domesticate_head = [&](freelist::QueuePtr value) -> freelist::HeadPtr { - head_domesticates++; - SNMALLOC_CHECK(address_cast(value) == address_cast(as_head(objects[0]))); - return domesticate(value); - }; - auto domesticate_queue = - [&](freelist::QueuePtr value) -> freelist::HeadPtr { - queue_domesticates++; - SNMALLOC_CHECK(address_cast(value) != address_cast(as_head(objects[0]))); - SNMALLOC_CHECK( - (address_cast(value) == address_cast(as_head(objects[1]))) || - (address_cast(value) == address_cast(as_head(objects[2])))); - return domesticate(value); - }; + queue.enqueue(as_head(old_first), as_head(old_last), domesticate)); queue.drain_and_reset( - domesticate_head, domesticate_queue, [&order](freelist::HeadPtr value) { + domesticate, domesticate, [&](freelist::HeadPtr value) { order.add(value); + if (address_cast(value) == address_cast(as_head(old_first))) + { + freelist::Object::atomic_store_next( + as_head(replacement_first), + as_head(replacement_last), + queue_key, + NO_KEY_TWEAK); + replacement_started = queue.enqueue( + as_head(replacement_first), as_head(replacement_last), domesticate); + } }); - SNMALLOC_CHECK(head_domesticates == 1); - SNMALLOC_CHECK(queue_domesticates == 2); - SNMALLOC_CHECK(order.count == 3); - for (size_t i = 0; i < 3; i++) - { - SNMALLOC_CHECK( - address_cast(order.values[i]) == address_cast(as_head(objects[i]))); - } - check_empty(queue); - } - - void test_remote_allocator_move_only_callback() - { - RemoteAllocator remote; - RemoteMessage message{}; - size_t callbacks = 0; - auto message_ptr = capptr::Alloc::unsafe_from(&message); - - SNMALLOC_CHECK(remote.enqueue(message_ptr, message_ptr, domesticate)); - - auto cb = [token = std::make_unique(1), - &callbacks](capptr::Alloc) mutable { - callbacks += (*token)++; - }; - remote.drain_and_reset(domesticate, domesticate, std::move(cb)); - - SNMALLOC_CHECK(callbacks == 1); - SNMALLOC_CHECK(remote.list.back.load(stl::memory_order_relaxed) == nullptr); + SNMALLOC_CHECK(replacement_started); + SNMALLOC_CHECK(order.count == 2); SNMALLOC_CHECK( - remote.list.front.load(stl::memory_order_relaxed) == nullptr); - } - - void test_front_publication_gap() - { - Queue queue; - Object object; - std::atomic entered{false}; - std::atomic completed{false}; - std::atomic callbacks{0}; - - freelist::Object::atomic_store_null( - as_head(object), queue_key, NO_KEY_TWEAK); - auto prev = queue.back.exchange( - capptr_rewild(as_head(object)), stl::memory_order_acq_rel); - SNMALLOC_CHECK(prev == nullptr); - - std::thread consumer([&]() { - entered.store(true, std::memory_order_release); - queue.drain_and_reset( - domesticate, domesticate, [&](freelist::HeadPtr value) { - SNMALLOC_CHECK(address_cast(value) == address_cast(as_head(object))); - callbacks.fetch_add(1, std::memory_order_relaxed); - }); - completed.store(true, std::memory_order_release); - }); - - wait_until([&]() { return entered.load(std::memory_order_acquire); }); - std::this_thread::sleep_for(std::chrono::milliseconds(10)); - SNMALLOC_CHECK(!completed.load(std::memory_order_acquire)); - SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 0); - - queue.front.store(capptr_rewild(as_head(object))); - consumer.join(); - - SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 1); - SNMALLOC_CHECK(completed.load(std::memory_order_acquire)); - check_empty(queue); - } - - void test_successor_publication_gap() - { - Queue queue; - Object first; - Object second; - std::atomic entered{false}; - std::atomic completed{false}; - std::atomic callbacks{0}; - - SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(first), domesticate)); - - freelist::Object::atomic_store_null( - as_head(second), queue_key, NO_KEY_TWEAK); - auto prev = queue.back.exchange( - capptr_rewild(as_head(second)), stl::memory_order_acq_rel); - SNMALLOC_CHECK(address_cast(prev) == address_cast(as_head(first))); - - std::thread consumer([&]() { - entered.store(true, std::memory_order_release); - queue.drain_and_reset( - domesticate, domesticate, [&callbacks](freelist::HeadPtr) { - callbacks.fetch_add(1, std::memory_order_relaxed); - }); - completed.store(true, std::memory_order_release); - }); - - wait_until([&]() { return entered.load(std::memory_order_acquire); }); - std::this_thread::sleep_for(std::chrono::milliseconds(10)); - SNMALLOC_CHECK(!completed.load(std::memory_order_acquire)); - SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 0); - - freelist::Object::atomic_store_next( - as_head(first), as_head(second), queue_key, NO_KEY_TWEAK); - consumer.join(); + address_cast(order.values[0]) == address_cast(as_head(old_first))); + SNMALLOC_CHECK( + address_cast(order.values[1]) == address_cast(as_head(old_last))); - SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 2); - SNMALLOC_CHECK(completed.load(std::memory_order_acquire)); - check_empty(queue); + queue.drain_and_reset( + domesticate, domesticate, [&order](freelist::HeadPtr value) { + order.add(value); + }); + SNMALLOC_CHECK(order.count == 4); + SNMALLOC_CHECK( + address_cast(order.values[2]) == + address_cast(as_head(replacement_first))); + SNMALLOC_CHECK( + address_cast(order.values[3]) == address_cast(as_head(replacement_last))); } - void test_rejected_front() + /** + * RemoteAllocator::drain_and_reset uses domesticate_head for the value read + * from front and domesticate_queue for successors read from message links. + * It recovers each RemoteMessage from its link, which requires a non-zero + * displacement when remote messages are batched. + */ + void test_remote_allocator_uses_distinct_domesticators() { - Queue queue; - Object object; - std::atomic reject{true}; - std::atomic attempted{false}; - std::atomic completed{false}; - std::atomic callbacks{0}; + RemoteAllocator remote; + RemoteMessage first{}; + RemoteMessage second{}; + auto first_message = capptr::Alloc::unsafe_from(&first); + auto second_message = capptr::Alloc::unsafe_from(&second); + auto first_link = RemoteMessage::to_message_link(first_message); + auto second_link = RemoteMessage::to_message_link(second_message); + size_t head_domesticates = 0; + size_t queue_domesticates = 0; + CallbackOrder, 2> order; + SNMALLOC_CHECK(remote.enqueue(first_message, first_message, domesticate)); SNMALLOC_CHECK( - queue.enqueue(as_head(object), as_head(object), domesticate)); + !remote.enqueue(second_message, second_message, domesticate)); - auto validating_domesticate = [&reject, - &attempted](freelist::QueuePtr value) { - attempted.store(true, std::memory_order_release); - if (reject.load(std::memory_order_acquire)) - return freelist::HeadPtr(nullptr); + auto domesticate_head = [&](freelist::QueuePtr value) -> freelist::HeadPtr { + head_domesticates++; + SNMALLOC_CHECK(address_cast(value) == address_cast(first_link)); return domesticate(value); }; - - std::thread consumer([&]() { - queue.drain_and_reset( - validating_domesticate, domesticate, [&callbacks](freelist::HeadPtr) { - callbacks.fetch_add(1, std::memory_order_relaxed); - }); - completed.store(true, std::memory_order_release); - }); - - wait_until([&]() { return attempted.load(std::memory_order_acquire); }); - SNMALLOC_CHECK(!completed.load(std::memory_order_acquire)); - SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 0); - - reject.store(false, std::memory_order_release); - consumer.join(); - - SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 1); - SNMALLOC_CHECK(completed.load(std::memory_order_acquire)); - check_empty(queue); - } - - void test_rejected_successor() - { - Queue queue; - Object first; - Object second; - std::atomic reject{true}; - std::atomic attempted{false}; - std::atomic completed{false}; - std::atomic callbacks{0}; - - freelist::Object::atomic_store_next( - as_head(first), as_head(second), queue_key, NO_KEY_TWEAK); - SNMALLOC_CHECK(queue.enqueue(as_head(first), as_head(second), domesticate)); - - auto validating_domesticate = [&reject, - &attempted](freelist::QueuePtr value) { - attempted.store(true, std::memory_order_release); - if (reject.load(std::memory_order_acquire)) - return freelist::HeadPtr(nullptr); + auto domesticate_queue = + [&](freelist::QueuePtr value) -> freelist::HeadPtr { + queue_domesticates++; + SNMALLOC_CHECK(address_cast(value) == address_cast(second_link)); return domesticate(value); }; - std::thread consumer([&]() { - queue.drain_and_reset( - domesticate, validating_domesticate, [&callbacks](freelist::HeadPtr) { - callbacks.fetch_add(1, std::memory_order_relaxed); - }); - completed.store(true, std::memory_order_release); - }); - - wait_until([&]() { return attempted.load(std::memory_order_acquire); }); - SNMALLOC_CHECK(!completed.load(std::memory_order_acquire)); - SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 0); + remote.drain_and_reset( + domesticate_head, + domesticate_queue, + [&order](capptr::Alloc value) { order.add(value); }); - reject.store(false, std::memory_order_release); - consumer.join(); - - SNMALLOC_CHECK(callbacks.load(std::memory_order_relaxed) == 2); - SNMALLOC_CHECK(completed.load(std::memory_order_acquire)); - check_empty(queue); - } - - void test_callback_starts_replacement() - { - Queue queue; - Object objects[4]; - CallbackOrder order; - std::atomic first_callback{false}; - std::atomic release_callback{false}; - - freelist::Object::atomic_store_next( - as_head(objects[0]), as_head(objects[1]), queue_key, NO_KEY_TWEAK); - SNMALLOC_CHECK( - queue.enqueue(as_head(objects[0]), as_head(objects[1]), domesticate)); - - std::thread consumer([&]() { - queue.drain_and_reset( - domesticate, domesticate, [&](freelist::HeadPtr value) { - order.add(value); - if (address_cast(value) == address_cast(as_head(objects[0]))) - { - first_callback.store(true, std::memory_order_release); - wait_until([&]() { - return release_callback.load(std::memory_order_acquire); - }); - } - }); - }); - - wait_until( - [&]() { return first_callback.load(std::memory_order_acquire); }); - - freelist::Object::atomic_store_next( - as_head(objects[2]), as_head(objects[3]), queue_key, NO_KEY_TWEAK); - SNMALLOC_CHECK( - queue.enqueue(as_head(objects[2]), as_head(objects[3]), domesticate)); - release_callback.store(true, std::memory_order_release); - consumer.join(); - - SNMALLOC_CHECK(order.count == 2); - for (size_t i = 0; i < 2; i++) - { - SNMALLOC_CHECK( - address_cast(order.values[i]) == address_cast(as_head(objects[i]))); - } - SNMALLOC_CHECK(queue.back.load(stl::memory_order_relaxed) != nullptr); - - order.count = 0; - queue.drain_and_reset( - domesticate, domesticate, [&order](freelist::HeadPtr value) { - order.add(value); - }); + SNMALLOC_CHECK(head_domesticates == 1); + SNMALLOC_CHECK(queue_domesticates == 1); SNMALLOC_CHECK(order.count == 2); - for (size_t i = 0; i < 2; i++) - { - SNMALLOC_CHECK( - address_cast(order.values[i]) == address_cast(as_head(objects[i + 2]))); - } - check_empty(queue); - } - - size_t find_object(Object* objects, size_t count, freelist::HeadPtr value) - { - for (size_t i = 0; i < count; i++) - { - if (address_cast(as_head(objects[i])) == address_cast(value)) - return i; - } - SNMALLOC_CHECK(false); - return count; - } - - void test_concurrent_producers(size_t object_count) - { - constexpr size_t producer_count = 4; - Queue queue; - auto objects = std::make_unique(object_count); - auto seen = std::make_unique[]>(object_count); - std::atomic next_object{0}; - std::atomic producers_live{producer_count}; - std::atomic true_results{0}; - - for (size_t i = 0; i < object_count; i++) - seen[i].store(0, std::memory_order_relaxed); - - std::vector producers; - for (size_t producer = 0; producer < producer_count; producer++) - { - producers.emplace_back([&]() { - while (true) - { - size_t first_index = - next_object.fetch_add(4, std::memory_order_relaxed); - if (first_index >= object_count) - break; - - size_t last_index = first_index + 3; - if (last_index >= object_count) - last_index = object_count - 1; - for (size_t i = first_index; i < last_index; i++) - { - freelist::Object::atomic_store_next( - as_head(objects[i]), - as_head(objects[i + 1]), - queue_key, - NO_KEY_TWEAK); - } - - if (queue.enqueue( - as_head(objects[first_index]), - as_head(objects[last_index]), - domesticate)) - { - true_results.fetch_add(1, std::memory_order_relaxed); - } - } - producers_live.fetch_sub(1, std::memory_order_release); - }); - } - - size_t callbacks = 0; - while (producers_live.load(std::memory_order_acquire) != 0 || - queue.back.load(stl::memory_order_relaxed) != nullptr) - { - queue.drain_and_reset( - domesticate, domesticate, [&](freelist::HeadPtr value) { - size_t index = find_object(objects.get(), object_count, value); - SNMALLOC_CHECK( - seen[index].fetch_add(1, std::memory_order_relaxed) == 0); - callbacks++; - }); - std::this_thread::yield(); - } - - for (auto& producer : producers) - producer.join(); - - SNMALLOC_CHECK(callbacks == object_count); - for (size_t i = 0; i < object_count; i++) - SNMALLOC_CHECK(seen[i].load(std::memory_order_relaxed) == 1); - SNMALLOC_CHECK(true_results.load(std::memory_order_relaxed) >= 1); SNMALLOC_CHECK( - true_results.load(std::memory_order_relaxed) <= object_count); - check_empty(queue); - - Object replacement; + address_cast(order.values[0]) == address_cast(first_message)); SNMALLOC_CHECK( - queue.enqueue(as_head(replacement), as_head(replacement), domesticate)); - SNMALLOC_CHECK(queue.back.load(stl::memory_order_relaxed) != nullptr); - queue.drain_and_reset(domesticate, domesticate, [](freelist::HeadPtr) {}); - check_empty(queue); - } - - void test_reuse(size_t generations) - { - constexpr size_t object_count = 8; - constexpr size_t reserved = 2; - Queue queue; - Object objects[object_count]; - std::atomic states[object_count]; - std::atomic generation[object_count]; - auto seen = - std::make_unique[]>(object_count * generations); - std::atomic producer_done{false}; - - for (size_t i = 0; i < object_count; i++) - { - states[i].store(0, std::memory_order_relaxed); - generation[i].store(0, std::memory_order_relaxed); - } - for (size_t i = 0; i < object_count * generations; i++) - seen[i].store(0, std::memory_order_relaxed); - - std::thread producer([&]() { - for (size_t round = 0; round < generations; round++) - { - size_t claimed[object_count]; - size_t claim_count = 0; - while (claim_count < object_count) - { - for (size_t i = 0; i < object_count; i++) - { - size_t available = 0; - if (states[i].compare_exchange_strong( - available, - reserved, - std::memory_order_acq_rel, - std::memory_order_relaxed)) - { - generation[i].store(round, std::memory_order_relaxed); - claimed[claim_count++] = i; - } - } - if (claim_count < object_count) - std::this_thread::yield(); - } - - for (size_t first = 0; first < object_count; first += 3) - { - size_t last = first + 2; - if (last >= object_count) - last = object_count - 1; - for (size_t i = first; i < last; i++) - { - freelist::Object::atomic_store_next( - as_head(objects[claimed[i]]), - as_head(objects[claimed[i + 1]]), - queue_key, - NO_KEY_TWEAK); - } - for (size_t i = first; i <= last; i++) - states[claimed[i]].store(1, std::memory_order_release); - queue.enqueue( - as_head(objects[claimed[first]]), - as_head(objects[claimed[last]]), - domesticate); - } - } - producer_done.store(true, std::memory_order_release); - }); - - size_t callbacks = 0; - while (!producer_done.load(std::memory_order_acquire) || - queue.back.load(stl::memory_order_relaxed) != nullptr) - { - queue.drain_and_reset( - domesticate, domesticate, [&](freelist::HeadPtr value) { - size_t index = find_object(objects, object_count, value); - SNMALLOC_CHECK(states[index].load(std::memory_order_acquire) == 1); - size_t round = generation[index].load(std::memory_order_relaxed); - SNMALLOC_CHECK( - seen[(round * object_count) + index].fetch_add( - 1, std::memory_order_relaxed) == 0); - callbacks++; - states[index].store(0, std::memory_order_release); - std::this_thread::yield(); - }); - } - producer.join(); - - SNMALLOC_CHECK(callbacks == object_count * generations); - for (size_t i = 0; i < object_count * generations; i++) - SNMALLOC_CHECK(seen[i].load(std::memory_order_relaxed) == 1); - check_empty(queue); + address_cast(order.values[1]) == address_cast(second_message)); } } -int main(int argc, char** argv) +int main() { - bool heavy = false; - for (int i = 1; i < argc; i++) - { - if (std::strcmp(argv[i], "--heavy") == 0) - heavy = true; - } - setup(); test_empty(); - test_enqueue_and_iterate(); - test_dequeue_empty(); - test_dequeue_positions(); - test_dequeue_publication_gap(); + test_queue_contract(); test_dequeue_checks_bound_first(); - test_drain_uses_distinct_domesticators(); - test_remote_allocator_move_only_callback(); - test_front_publication_gap(); - test_successor_publication_gap(); - test_rejected_front(); - test_rejected_successor(); test_callback_starts_replacement(); - test_concurrent_producers(heavy ? 4096 : 256); - test_reuse(heavy ? 256 : 16); + test_remote_allocator_uses_distinct_domesticators(); }