diff --git a/src/snmalloc/mem/corealloc.h b/src/snmalloc/mem/corealloc.h index b9d63c957..a28357cd9 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,25 +1402,27 @@ namespace snmalloc size_t bytes_flushed = 0; // Not currently used. - if (destroy_queue) + 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); + }; + + if constexpr (Config::Options.QueueHeadsAreTame) { - 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 domesticate_first = + [](freelist::QueuePtr p) SNMALLOC_FAST_PATH_LAMBDA { + return freelist::HeadPtr::unsafe_from(p.unsafe_ptr()); }; - - message_queue().destroy_and_iterate(domesticate, cb); + message_queue().drain_and_reset(domesticate_first, 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, domesticate, cb); } auto& key = freelist::Object::key_root; @@ -1514,7 +1518,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..33f7fa489 100644 --- a/src/snmalloc/mem/freelist_queue.h +++ b/src/snmalloc/mem/freelist_queue.h @@ -16,18 +16,14 @@ 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 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 * - * * `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 +47,88 @@ 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; } - template - void destroy_and_iterate(Domesticator_queue domesticate, Cb cb) + 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. + * + * 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< + typename Domesticator_head, + typename Domesticator_queue, + typename Cb> + void drain_and_reset( + Domesticator_head domesticate_head, + Domesticator_queue domesticate_queue, + 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_head(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_queue, process); + if (address_cast(curr) == address_cast(target)) + break; + Aal::pause(); } + cb(curr); } inline bool can_dequeue() @@ -94,9 +143,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 +178,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 +221,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..efd4ba942 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,23 @@ namespace snmalloc list.invariant(); } - void init() - { - list.init(); - } - - template - void destroy_and_iterate(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 { + 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_head), + stl::move(domesticate_queue), + stl::move(cbwrap)); } inline bool can_dequeue() @@ -346,14 +351,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..3c1af0b24 --- /dev/null +++ b/src/test/func/freelist_mpscq/freelist_mpscq.cc @@ -0,0 +1,294 @@ +#include +#include + +using namespace snmalloc; + +namespace +{ + FreeListKey queue_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 + struct CallbackOrder + { + T values[Size]; + size_t count = 0; + + void add(T value) + { + SNMALLOC_CHECK(count < Size); + values[count++] = value; + } + }; + + /** + * 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); + } + + /** + * 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; + 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), queue_key, NO_KEY_TWEAK); + SNMALLOC_CHECK( + !queue.enqueue(as_head(batch_first), as_head(batch_last), 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(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[1]) == address_cast(as_head(second))); + SNMALLOC_CHECK( + address_cast(order.values[2]) == address_cast(as_head(batch_first))); + + queue.drain_and_reset( + domesticate, domesticate, [&order](freelist::HeadPtr value) { + order.add(value); + }); + SNMALLOC_CHECK(order.count == 4); + SNMALLOC_CHECK( + address_cast(order.values[3]) == address_cast(as_head(batch_last))); + + SNMALLOC_CHECK( + queue.enqueue(as_head(replacement), as_head(replacement), domesticate)); + queue.drain_and_reset( + domesticate, domesticate, [&order](freelist::HeadPtr value) { + order.add(value); + }); + 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 tail; + Object outside; + size_t queue_domesticates = 0; + size_t dequeue_callbacks = 0; + size_t drain_callbacks = 0; + + SNMALLOC_CHECK(queue.enqueue(as_head(tail), as_head(tail), domesticate)); + freelist::Object::atomic_store_next( + as_head(tail), as_head(outside), queue_key, NO_KEY_TWEAK); + + auto counting_domesticate = + [&queue_domesticates](freelist::QueuePtr value) -> freelist::HeadPtr { + queue_domesticates++; + return domesticate(value); + }; + + queue.dequeue( + domesticate, + counting_domesticate, + [&dequeue_callbacks](freelist::HeadPtr) { + dequeue_callbacks++; + return true; + }); + + SNMALLOC_CHECK(queue_domesticates == 0); + SNMALLOC_CHECK(dequeue_callbacks == 0); + + 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); + } + + /** + * 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 old_first; + Object old_last; + Object replacement_first; + Object replacement_last; + CallbackOrder order; + bool replacement_started = false; + + 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.drain_and_reset( + 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(replacement_started); + SNMALLOC_CHECK(order.count == 2); + SNMALLOC_CHECK( + address_cast(order.values[0]) == address_cast(as_head(old_first))); + SNMALLOC_CHECK( + address_cast(order.values[1]) == address_cast(as_head(old_last))); + + 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))); + } + + /** + * 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() + { + 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( + !remote.enqueue(second_message, second_message, domesticate)); + + auto domesticate_head = [&](freelist::QueuePtr value) -> freelist::HeadPtr { + head_domesticates++; + SNMALLOC_CHECK(address_cast(value) == address_cast(first_link)); + return domesticate(value); + }; + auto domesticate_queue = + [&](freelist::QueuePtr value) -> freelist::HeadPtr { + queue_domesticates++; + SNMALLOC_CHECK(address_cast(value) == address_cast(second_link)); + return domesticate(value); + }; + + remote.drain_and_reset( + domesticate_head, + domesticate_queue, + [&order](capptr::Alloc value) { order.add(value); }); + + SNMALLOC_CHECK(head_domesticates == 1); + SNMALLOC_CHECK(queue_domesticates == 1); + SNMALLOC_CHECK(order.count == 2); + SNMALLOC_CHECK( + address_cast(order.values[0]) == address_cast(first_message)); + SNMALLOC_CHECK( + address_cast(order.values[1]) == address_cast(second_message)); + } +} + +int main() +{ + setup(); + test_empty(); + test_queue_contract(); + test_dequeue_checks_bound_first(); + test_callback_starts_replacement(); + test_remote_allocator_uses_distinct_domesticators(); +} diff --git a/src/test/perf/msgpass/msgpass.cc b/src/test/perf/msgpass/msgpass.cc index b8c0d9d2b..7344161b7 100644 --- a/src/test/perf/msgpass/msgpass.cc +++ b/src/test/perf/msgpass/msgpass.cc @@ -102,7 +102,10 @@ 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, domesticate_nop, [](freelist::HeadPtr o) { + snmalloc::dealloc(o.as_void().unsafe_ptr()); + }); } void proxy(const struct params* param, size_t qix) @@ -134,7 +137,10 @@ 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, domesticate_nop, [](freelist::HeadPtr o) { + snmalloc::dealloc(o.as_void().unsafe_ptr()); + }); queue_gate--; } @@ -218,11 +224,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;