diff --git a/.gitignore b/.gitignore index 391be7f9..cc5365a2 100644 --- a/.gitignore +++ b/.gitignore @@ -69,4 +69,5 @@ TAGS ########################### src/rtorrent test/rtorrent_Test* +test/test-suite.log.tmp test-driver diff --git a/src/command_download.cc b/src/command_download.cc index 0e64bdbc..a8bc128f 100644 --- a/src/command_download.cc +++ b/src/command_download.cc @@ -18,6 +18,7 @@ #include #include #include +#include #include #include #include @@ -305,8 +306,10 @@ apply_d_add_peer(core::Download* download, const std::string& arg) { assert(std::this_thread::get_id() == torrent::main_thread::thread_id()); + auto callback_id = torrent::system::make_callback_id(); + // Currently discarding SOCK_STREAM. - torrent::this_thread::resolver()->resolve_preferred(NULL, host, AF_UNSPEC, AF_INET, [download, port](torrent::c_sa_shared_ptr sa, int err) { + torrent::this_thread::resolver()->resolve_preferred(callback_id, host, AF_UNSPEC, AF_INET, [download, port](torrent::c_sa_shared_ptr sa, int err) { if (sa == nullptr) { lt_log_print(torrent::LOG_TORRENT_WARN, "could not resolve hostname for added peer: %s", gai_strerror(err)); return; diff --git a/src/command_tracker.cc b/src/command_tracker.cc index f61b47fe..5543715f 100644 --- a/src/command_tracker.cc +++ b/src/command_tracker.cc @@ -6,6 +6,7 @@ #include #include #include +#include #include #include #include @@ -48,8 +49,10 @@ apply_dht_add_node(const std::string& arg) { // TODO: Move this lookup to DhtController. + auto callback_id = torrent::system::make_callback_id(); + // Currently discarding SOCK_STREAM. - torrent::this_thread::resolver()->resolve_specific(nullptr, host_str, PF_INET, [host_str, port](torrent::c_sa_shared_ptr sa, int err) { + torrent::this_thread::resolver()->resolve_specific(callback_id, host_str, PF_INET, [host_str, port](torrent::c_sa_shared_ptr sa, int err) { if (sa == nullptr) { lt_log_print(torrent::LOG_DHT_ERROR, "dht.add_node : could not resolve host : %s (%s)", gai_strerror(err), host_str.c_str()); return; diff --git a/src/globals.h b/src/globals.h index dca7ae58..7d5cc9e5 100644 --- a/src/globals.h +++ b/src/globals.h @@ -25,13 +25,11 @@ namespace scgi_thread { torrent::system::Thread* thread(); std::thread::id thread_id(); -void callback(void* target, std::function&& fn); -void cancel_callback(void* target); -void cancel_callback_and_wait(void* target); +void callback_interrupt(torrent::system::callback_id& id, std::function&& fn); -rpc::SCgi* scgi(); -void set_scgi(rpc::SCgi* scgi); -void set_rpc_log(const std::string& filename); +rpc::SCgi* scgi(); +void set_scgi(rpc::SCgi* scgi); +void set_rpc_log(const std::string& filename); } // namespace torrent::scgi_thread @@ -41,9 +39,9 @@ namespace session_thread { torrent::system::Thread* thread(); std::thread::id thread_id(); -void callback(void* target, std::function&& fn); -void cancel_callback(void* target); -void cancel_callback_and_wait(void* target); +void callback(std::function&& fn); +void callback(torrent::system::callback_id& id, std::function&& fn); +void cancel_callback(torrent::system::callback_id& id); session::SessionManager* manager(); std::string session_path(); diff --git a/src/rpc/scgi_task.cc b/src/rpc/scgi_task.cc index aa93a8c7..0af85c65 100644 --- a/src/rpc/scgi_task.cc +++ b/src/rpc/scgi_task.cc @@ -14,7 +14,7 @@ #include #include #include -#include +#include #include "control.h" #include "globals.h" @@ -24,6 +24,12 @@ namespace rpc { +SCgiTask::SCgiTask() + : m_callback_id(torrent::system::make_callback_id()) { + + set_file_descriptor(-1); +} + void SCgiTask::open(SCgi* parent, int fd) { set_file_descriptor(fd); @@ -67,8 +73,7 @@ SCgiTask::close() { if (!is_open()) return; - torrent::main_thread::thread()->cancel_callback_and_wait(this); - torrent::system::Thread::self()->cancel_callback(this); + torrent::system::cancel_callback_and_wait(m_callback_id, scgi_thread::thread(), torrent::main_thread::thread()); torrent::runtime::socket_manager()->close_event_or_throw(this, [this]() { torrent::this_thread::poll()->remove_and_close(this); @@ -86,13 +91,13 @@ SCgiTask::close() { void SCgiTask::event_read() { int read_length = m_buffer.size() - m_position; - + if (m_content_length == 0) read_length--; if (read_length <= 0) throw torrent::internal_error("SCgiTask::event_read() no space in buffer for event_read."); - + int bytes = ::recv(m_fileDesc, m_buffer.data() + m_position, read_length, 0); if (bytes <= 0) { @@ -321,7 +326,7 @@ SCgiTask::detect_content_type(const std::string& content_type) { void SCgiTask::receive_call(const char* buffer, uint32_t length) { - assert(torrent::system::Thread::self() == scgi_thread::thread()); + assert(torrent::this_thread::thread() == scgi_thread::thread()); RpcManager::RPCType rpc_type; @@ -348,7 +353,7 @@ SCgiTask::receive_call(const char* buffer, uint32_t length) { m_result_mutex.lock(); m_result_mutex.unlock(); - scgi_thread::thread()->callback_interrupt_polling(this, [this]() { + scgi_thread::callback_interrupt(m_callback_id, [this]() { if (!is_open()) return; @@ -368,7 +373,7 @@ SCgiTask::receive_call(const char* buffer, uint32_t length) { m_result_mutex.lock(); m_result_mutex.unlock(); - torrent::main_thread::thread()->callback_interrupt_polling(this, [this, rpc_type, buffer, length, result_callback]() { + torrent::main_thread::callback_interrupt(m_callback_id, [this, rpc_type, buffer, length, result_callback]() { // Memory barrier for the input data. // std::atomic_thread_fence(std::memory_order_acquire); m_result_mutex.lock(); @@ -383,7 +388,7 @@ SCgiTask::receive_call(const char* buffer, uint32_t length) { void SCgiTask::receive_write(const char* buffer, uint32_t length) { - assert(torrent::system::Thread::self() == torrent::main_thread::thread()); + assert(torrent::this_thread::thread() == torrent::main_thread::thread()); if (buffer == nullptr || length > (100 << 20)) throw torrent::internal_error("SCgiTask::receive_write(...) received bad input."); diff --git a/src/rpc/scgi_task.h b/src/rpc/scgi_task.h index e012e9be..1b3961b1 100644 --- a/src/rpc/scgi_task.h +++ b/src/rpc/scgi_task.h @@ -18,7 +18,7 @@ public: enum ContentType { XML, JSON }; - SCgiTask() { m_fileDesc = -1; } + SCgiTask(); const char* type_name() const override { return "scgi-task"; } @@ -54,7 +54,8 @@ private: void plaintext_response(const char* buffer, uint32_t content_length); void gzip_response(const char* buffer, uint32_t content_length); - SCgi* m_parent{}; + SCgi* m_parent{}; + torrent::system::callback_id m_callback_id; std::mutex m_result_mutex; diff --git a/src/scgi/thread_scgi.cc b/src/scgi/thread_scgi.cc index ada7b7ad..3ff0d5f6 100644 --- a/src/scgi/thread_scgi.cc +++ b/src/scgi/thread_scgi.cc @@ -12,11 +12,6 @@ namespace scgi { -class ThreadScgiInternal { -public: - static ThreadScgi* thread_scgi() { return ThreadScgi::internal_thread_scgi(); } -}; - ThreadScgi* ThreadScgi::m_thread_scgi{}; void @@ -60,7 +55,7 @@ ThreadScgi::set_scgi(rpc::SCgi* scgi) { change_rpc_log(); - callback(nullptr, [this]() { + callback([this]() { if (m_scgi == nullptr) throw torrent::internal_error("Tried to start SCGI but object was not present."); @@ -72,7 +67,7 @@ ThreadScgi::set_scgi(rpc::SCgi* scgi) { void ThreadScgi::set_rpc_log(const std::string& filename) { - callback(nullptr, [this, filename]() { + callback([this, filename]() { m_rpc_log_filename = filename; change_rpc_log(); }); @@ -127,15 +122,13 @@ ThreadScgi::next_timeout() { namespace scgi_thread { -torrent::system::Thread* thread() { return scgi::ThreadScgiInternal::thread_scgi(); } -std::thread::id thread_id() { return scgi::ThreadScgiInternal::thread_scgi()->thread_id(); } +torrent::system::Thread* thread() { return scgi::ThreadScgi::thread_scgi(); } +std::thread::id thread_id() { return scgi::ThreadScgi::thread_scgi()->thread_id(); } -void callback(void* target, std::function&& fn) { scgi::ThreadScgiInternal::thread_scgi()->callback(target, std::move(fn)); } -void cancel_callback(void* target) { scgi::ThreadScgiInternal::thread_scgi()->cancel_callback(target); } -void cancel_callback_and_wait(void* target) { scgi::ThreadScgiInternal::thread_scgi()->cancel_callback_and_wait(target); } +void callback_interrupt(torrent::system::callback_id& id, std::function&& fn) { scgi::ThreadScgi::thread_scgi()->callback_interrupt(id, std::move(fn)); } -rpc::SCgi* scgi() { return scgi::ThreadScgiInternal::thread_scgi()->scgi(); } -void set_scgi(rpc::SCgi* scgi) { scgi::ThreadScgiInternal::thread_scgi()->set_scgi(scgi); } -void set_rpc_log(const std::string& filename) { scgi::ThreadScgiInternal::thread_scgi()->set_rpc_log(filename); } +rpc::SCgi* scgi() { return scgi::ThreadScgi::thread_scgi()->scgi(); } +void set_scgi(rpc::SCgi* scgi) { scgi::ThreadScgi::thread_scgi()->set_scgi(scgi); } +void set_rpc_log(const std::string& filename) { scgi::ThreadScgi::thread_scgi()->set_rpc_log(filename); } } // namespace scgi_thread diff --git a/src/scgi/thread_scgi.h b/src/scgi/thread_scgi.h index 824ffa98..38f5db27 100644 --- a/src/scgi/thread_scgi.h +++ b/src/scgi/thread_scgi.h @@ -11,8 +11,6 @@ class SCgi; namespace scgi { -class ThreadScgiInternal; - class ThreadScgi : public torrent::system::Thread { public: @@ -28,8 +26,6 @@ public: void set_rpc_log(const std::string& filename); protected: - friend class ThreadScgiInternal; - ThreadScgi() = default; static auto internal_thread_scgi() { return m_thread_scgi; } @@ -43,9 +39,9 @@ private: void task_touch_log(); void change_rpc_log(); - static ThreadScgi* m_thread_scgi; + static ThreadScgi* m_thread_scgi; - std::atomic m_scgi{nullptr}; + std::atomic m_scgi{}; std::string m_rpc_log_filename; }; diff --git a/src/session/session_manager.cc b/src/session/session_manager.cc index 3d2468be..933ad5fa 100644 --- a/src/session/session_manager.cc +++ b/src/session/session_manager.cc @@ -5,6 +5,7 @@ #include #include #include +#include #include #include "globals.h" @@ -23,6 +24,7 @@ namespace session { SessionManager::SessionManager(torrent::system::Thread* thread) : m_thread(thread), + m_callback_id(torrent::system::make_callback_id()), m_lockfile(std::make_unique()) { } @@ -221,8 +223,7 @@ SessionManager::cleanup() { LT_LOG("session manager cleaned up", 0); - session_thread::cancel_callback(this); - torrent::main_thread::cancel_callback(this); + torrent::system::cancel_callback_and_wait(m_callback_id, session_thread::thread(), torrent::main_thread::thread()); } void @@ -233,7 +234,9 @@ SessionManager::callback_pending_builds() { if (m_callback_scheduled_process_pending_builds.exchange(true)) return; - torrent::main_thread::callback(this, [this]() { process_pending_builds(false); }); + torrent::main_thread::callback(m_callback_id, [this]() { + process_pending_builds(false); + }); } void @@ -247,7 +250,9 @@ SessionManager::callback_save_request() { if (m_callback_scheduled_process_saves_request.exchange(true)) return; - session_thread::callback(this, [this]() { process_save_request(); }); + session_thread::callback(m_callback_id, [this]() { + process_save_request(); + }); } void @@ -255,7 +260,9 @@ SessionManager::callback_finished_saves() { if (m_callback_scheduled_process_finished_saves.exchange(true)) return; - session_thread::callback(this, [this]() { process_finished_saves(); }); + session_thread::callback(m_callback_id, [this]() { + process_finished_saves(); + }); } void diff --git a/src/session/session_manager.h b/src/session/session_manager.h index 4fa76843..eff135ac 100644 --- a/src/session/session_manager.h +++ b/src/session/session_manager.h @@ -89,37 +89,39 @@ private: bool replace_save_request_unsafe(SaveRequest& download); bool remove_completely_unsafe(core::Download* download, std::unique_lock& lock); - torrent::system::Thread* m_thread; + using ProcessingSave = std::pair, SaveRequest>; - bool m_freeze_info{}; - std::string m_path; - bool m_use_fsyncdisk{true}; - bool m_use_lock{true}; + torrent::system::Thread* m_thread; + torrent::system::callback_id m_callback_id; - std::mutex m_mutex; - bool m_active{}; + bool m_freeze_info{}; + std::string m_path; + bool m_use_fsyncdisk{true}; + bool m_use_lock{true}; - typedef std::pair, SaveRequest> ProcessingSave; + align_cacheline std::mutex m_mutex; - std::deque m_save_requests; - std::atomic m_save_request_counter{}; - std::list m_processing_saves; - std::atomic m_processing_save_counter{}; - std::condition_variable m_finished_condition; - std::vector m_finished_saves; + bool m_active{}; - std::atomic m_callback_scheduled_process_pending_builds{}; - std::atomic m_callback_scheduled_process_saves_request{}; - std::atomic m_callback_scheduled_process_finished_saves{}; + std::deque m_save_requests; + std::atomic m_save_request_counter{}; + std::list m_processing_saves; + std::atomic m_processing_save_counter{}; + std::condition_variable m_finished_condition; + std::vector m_finished_saves; + + std::atomic m_callback_scheduled_process_pending_builds{}; + std::atomic m_callback_scheduled_process_saves_request{}; + std::atomic m_callback_scheduled_process_finished_saves{}; std::unique_ptr m_lockfile; - std::chrono::microseconds m_last_storage_error_message{}; - unsigned int m_ignored_storage_error_count{}; + std::chrono::microseconds m_last_storage_error_message{}; + unsigned int m_ignored_storage_error_count{}; // Pending builds are only ever locked by main thread. - std::mutex m_pending_builds_mutex; - std::deque m_pending_builds; + std::mutex m_pending_builds_mutex; + std::deque m_pending_builds; }; inline bool SessionManager::is_used() const { return !m_path.empty(); } diff --git a/src/session/thread_session.cc b/src/session/thread_session.cc index bf698beb..e9015cc9 100644 --- a/src/session/thread_session.cc +++ b/src/session/thread_session.cc @@ -4,15 +4,11 @@ #include +#include "globals.h" #include "session/session_manager.h" namespace session { -class ThreadSessionInternal { -public: - static ThreadSession* thread_session() { return ThreadSession::internal_thread_session(); } -}; - ThreadSession* ThreadSession::m_thread_session{nullptr}; void @@ -74,14 +70,14 @@ ThreadSession::next_timeout() { namespace session_thread { -torrent::system::Thread* thread() { return session::ThreadSessionInternal::thread_session(); } -std::thread::id thread_id() { return session::ThreadSessionInternal::thread_session()->thread_id(); } +torrent::system::Thread* thread() { return session::ThreadSession::thread_session(); } +std::thread::id thread_id() { return thread()->thread_id(); } -void callback(void* target, std::function&& fn) { session::ThreadSessionInternal::thread_session()->callback(target, std::move(fn)); } -void cancel_callback(void* target) { session::ThreadSessionInternal::thread_session()->cancel_callback(target); } -void cancel_callback_and_wait(void* target) { session::ThreadSessionInternal::thread_session()->cancel_callback_and_wait(target); } +void callback(std::function&& fn) { thread()->callback(std::move(fn)); } +void callback(torrent::system::callback_id& id, std::function&& fn) { thread()->callback(id, std::move(fn)); } +void cancel_callback(torrent::system::callback_id& id) { thread()->cancel_callback(id); } -session::SessionManager* manager() { return session::ThreadSessionInternal::thread_session()->manager(); } -std::string session_path() { return session::ThreadSessionInternal::thread_session()->manager()->path(); } +session::SessionManager* manager() { return session::ThreadSession::thread_session()->manager(); } +std::string session_path() { return manager()->path(); } } // namespace session_thread