From 0ea7615d0abe799894e4575d6df5b36cc299d50f Mon Sep 17 00:00:00 2001 From: kannibalox Date: Sun, 1 Dec 2024 13:30:05 -0500 Subject: [PATCH] Convert __sync* built-ins to std::atomic Also converts thread_queue_hack to a simpler reserved vector --- src/control.h | 9 ++-- src/thread_base.cc | 121 ++++++++++--------------------------------- src/thread_base.h | 3 +- src/thread_worker.cc | 18 +++---- src/thread_worker.h | 10 +--- 5 files changed, 44 insertions(+), 117 deletions(-) diff --git a/src/control.h b/src/control.h index 023ecffc..e40b4dac 100644 --- a/src/control.h +++ b/src/control.h @@ -37,6 +37,7 @@ #ifndef RTORRENT_CONTROL_H #define RTORRENT_CONTROL_H +#include #include #include #include @@ -87,8 +88,8 @@ public: void handle_shutdown(); - void receive_normal_shutdown() { m_shutdownReceived = true; __sync_synchronize(); } - void receive_quick_shutdown() { m_shutdownReceived = true; m_shutdownQuick = true; __sync_synchronize(); } + void receive_normal_shutdown() { m_shutdownReceived = true; } + void receive_quick_shutdown() { m_shutdownReceived = true; m_shutdownQuick = true; } core::Manager* core() { return m_core; } core::ViewManager* view_manager() { return m_viewManager; } @@ -134,8 +135,8 @@ private: rak::priority_item m_taskShutdown; - bool m_shutdownReceived lt_cacheline_aligned; - bool m_shutdownQuick lt_cacheline_aligned; + std::atomic m_shutdownReceived lt_cacheline_aligned; + std::atomic m_shutdownQuick lt_cacheline_aligned; }; #endif diff --git a/src/thread_base.cc b/src/thread_base.cc index 081791d4..6564133a 100644 --- a/src/thread_base.cc +++ b/src/thread_base.cc @@ -1,39 +1,3 @@ -// libTorrent - BitTorrent library -// Copyright (C) 2005-2011, Jari Sundell -// -// This program is free software; you can redistribute it and/or modify -// it under the terms of the GNU General Public License as published by -// the Free Software Foundation; either version 2 of the License, or -// (at your option) any later version. -// -// This program is distributed in the hope that it will be useful, -// but WITHOUT ANY WARRANTY; without even the implied warranty of -// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the -// GNU General Public License for more details. -// -// You should have received a copy of the GNU General Public License -// along with this program; if not, write to the Free Software -// Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA 02111-1307 USA -// -// In addition, as a special exception, the copyright holders give -// permission to link the code of portions of this program with the -// OpenSSL library under certain conditions as described in each -// individual source file, and distribute linked combinations -// including the two. -// -// You must obey the GNU General Public License in all respects for -// all of the code used other than OpenSSL. If you modify file(s) -// with this exception, you may extend this exception to your version -// of the file(s), but you are not obligated to do so. If you do not -// wish to do so, delete this exception statement from your version. -// If you delete this exception statement from all source files in the -// program, then also delete it here. -// -// Contact: Jari Sundell -// -// Skomakerveien 33 -// 3185 Skoppum, NORWAY - #include "config.h" #include "thread_base.h" @@ -41,82 +5,54 @@ #include #include #include -#include -#include -#include #include #include #include #include +#include #include "globals.h" -#include "control.h" -#include "core/manager.h" -// Temporarly injected into config.h. -/* temp hack */ -//#define lt_cacheline_aligned __attribute__((__aligned__(128))) - -class lt_cacheline_aligned thread_queue_hack { +class lt_cacheline_aligned thread_queue_hack : private std::vector { public: - typedef ThreadBase::thread_base_func value_type; - typedef ThreadBase::thread_base_func* iterator; + static constexpr unsigned int max_size = 32; - static const unsigned int max_size = 32; + using value_type = ThreadBase::thread_base_func; + using base_type = std::vector; - thread_queue_hack() { std::memset(this, 0, sizeof(thread_queue_hack)); } + using base_type::empty; - void lock() { while (!__sync_bool_compare_and_swap(&m_lock, 0, 1)) usleep(0); } - void unlock() { __sync_bool_compare_and_swap(&m_lock, 1, 0); } - - iterator begin() { return m_queue; } - iterator max_capacity() { return m_queue + max_size; } - - iterator end_and_lock() { lock(); return std::find(begin(), max_capacity(), (value_type)NULL); } - - bool empty() const { return m_queue[0] == NULL; } + thread_queue_hack() { reserve(max_size); }; void push_back(value_type v) { - iterator itr = end_and_lock(); - - if (itr == max_capacity()) - throw torrent::internal_error("Overflowed thread_queue."); - - __sync_bool_compare_and_swap(itr, NULL, v); - unlock(); + if (size() >= max_size) + throw torrent::internal_error("Overflowed thread_queue max size of " + std::to_string(max_size) + "."); + const std::lock_guard lguard(m_lock); + base_type::push_back(v); } - value_type* copy_and_clear(value_type* dest) { - iterator itr = begin(); - lock(); - - while (*itr != NULL) *dest++ = *itr++; - - clear_and_unlock(); - return dest; + void copy_and_clear(base_type& target) { + std::lock_guard lguard(m_lock); + target.assign(begin(), end()); + clear(); } - void clear_and_unlock() { - std::memset(this, 0, sizeof(thread_queue_hack)); - __sync_synchronize(); - } - - private: - int m_lock; - value_type m_queue[max_size + 1]; +private: + std::mutex m_lock; }; -void throw_shutdown_exception() { throw torrent::shutdown_exception(); } +void +throw_shutdown_exception() { throw torrent::shutdown_exception(); } ThreadBase::ThreadBase() { m_taskShutdown.slot() = std::bind(&throw_shutdown_exception); - m_threadQueue = new thread_queue_hack; + m_threadQueue = std::make_unique(); } -ThreadBase::~ThreadBase() { - delete m_threadQueue; -} +// Defined in here and not the header so that the default destructor +// can properly deduce how to destruct the unique_ptr +ThreadBase::~ThreadBase() = default; // Move to libtorrent... void @@ -137,12 +73,11 @@ ThreadBase::next_timeout_usec() { void ThreadBase::call_queued_items() { - thread_base_func result[thread_queue_hack::max_size]; - thread_base_func* first = result; - thread_base_func* last = m_threadQueue->copy_and_clear((thread_base_func*)result); - - while (first != last && *first) - (*first++)(this); + std::vector queue; + m_threadQueue->copy_and_clear(queue); + for (auto itr : queue) { + itr(this); + } } void diff --git a/src/thread_base.h b/src/thread_base.h index b125ea8b..f8d7e855 100644 --- a/src/thread_base.h +++ b/src/thread_base.h @@ -37,6 +37,7 @@ #ifndef RTORRENT_UTILS_THREAD_BASE_H #define RTORRENT_UTILS_THREAD_BASE_H +#include #include #include #include @@ -82,7 +83,7 @@ protected: // Temporary hack to pass messages to a thread. This really needs to // be cleaned up and/or integrated into the priority queue itself. - thread_queue_hack* m_threadQueue; + std::unique_ptr m_threadQueue; }; #endif diff --git a/src/thread_worker.cc b/src/thread_worker.cc index 939d7026..61fc4083 100644 --- a/src/thread_worker.cc +++ b/src/thread_worker.cc @@ -55,8 +55,8 @@ ThreadWorker::ThreadWorker() { } ThreadWorker::~ThreadWorker() { - if (m_safe.scgi) - m_safe.scgi->deactivate(); + if (m_scgi) + m_scgi.load()->deactivate(); } void @@ -67,17 +67,13 @@ ThreadWorker::init_thread() { bool ThreadWorker::set_scgi(rpc::SCgi* scgi) { - if (!__sync_bool_compare_and_swap(&m_safe.scgi, NULL, scgi)) + if (m_scgi != nullptr) return false; + m_scgi = scgi; + change_xmlrpc_log(); - // The xmlrpc process call requires a global lock. -// m_safe.scgi->set_slot_process(rak::mem_fn(&rpc::xmlrpc, &rpc::XmlRpc::process)); - - // Synchronize in order to ensure the worker thread sees the updated - // SCgi object. - __sync_synchronize(); queue_item((thread_base_func)&start_scgi); return true; } @@ -93,10 +89,10 @@ void ThreadWorker::start_scgi(ThreadBase* baseThread) { ThreadWorker* thread = (ThreadWorker*)baseThread; - if (thread->m_safe.scgi == NULL) + if (thread->scgi() == NULL) throw torrent::internal_error("Tried to start SCGI but object was not present."); - thread->m_safe.scgi->activate(); + thread->scgi()->activate(); } void diff --git a/src/thread_worker.h b/src/thread_worker.h index 49414b89..ffb7c855 100644 --- a/src/thread_worker.h +++ b/src/thread_worker.h @@ -57,7 +57,7 @@ public: virtual void init_thread(); - rpc::SCgi* scgi() { return m_safe.scgi; } + rpc::SCgi* scgi() { return m_scgi; } bool set_scgi(rpc::SCgi* scgi); void set_xmlrpc_log(const std::string& filename); @@ -70,13 +70,7 @@ private: void change_xmlrpc_log(); - struct lt_cacheline_aligned safe_type { - safe_type() : scgi(NULL) {} - - rpc::SCgi* scgi; - }; - - safe_type m_safe; + std::atomic lt_cacheline_aligned m_scgi{ nullptr }; // The following types shall only be modified while holding the // global lock.