Added WaitpidQueue to handle background process reaping.

This commit is contained in:
Jari Sundell
2026-08-04 12:03:26 +02:00
committed by GitHub
parent 247ae7621a
commit f90a96d57a
16 changed files with 244 additions and 180 deletions
+115
View File
@@ -0,0 +1,115 @@
#include "config.h"
#include "utils/waitpid_queue.h"
#include <sys/wait.h>
#include <torrent/exceptions.h>
namespace utils {
WaitpidQueue::WaitpidQueue() {
m_worker = std::async(std::launch::async, [this]() {
bool is_running = true;
auto wait_time = 50ms;
while (is_running) {
if (!m_queue.empty()) {
auto start_time = std::chrono::steady_clock::now();
while (!m_wakeup_worker.load(std::memory_order_acquire) && !m_should_shutdown) {
auto elapsed = std::chrono::steady_clock::now() - start_time;
if (elapsed >= wait_time)
break;
std::this_thread::sleep_for(50ms);
}
} else {
m_wakeup_worker.wait(false, std::memory_order_acquire);
std::this_thread::sleep_for(50ms);
}
std::set<pid_t> queue;
{
std::lock_guard<std::mutex> guard(m_mutex);
if (m_should_shutdown) {
if (m_queue.empty())
return;
is_running = false;
}
if (m_queue.empty())
throw torrent::internal_error("WaitpidQueue worker thread woke up but queue is empty.");
queue = m_queue;
m_wakeup_worker.store(false, std::memory_order_release);
}
wait_time = std::min(10 * 1000ms, wait_time * 2);
for (auto pid : queue) {
if (::waitpid(pid, nullptr, WNOHANG) == 0)
continue;
{
std::lock_guard<std::mutex> guard(m_mutex);
if (m_queue.erase(pid) != 1)
throw torrent::internal_error("WaitpidQueue worker thread could not find pid in queue.");
}
wait_time = std::max(50ms, wait_time / 2);
m_remaining.fetch_sub(1, std::memory_order_release);
m_remaining.notify_all();
}
}
});
}
WaitpidQueue::~WaitpidQueue() {
{
std::lock_guard<std::mutex> guard(m_mutex);
m_should_shutdown = true;
}
m_wakeup_worker.store(true, std::memory_order_release);
m_wakeup_worker.notify_all();
// m_worker.wait();
}
void
WaitpidQueue::close_pid(pid_t pid) {
if (pid < 0)
throw torrent::internal_error("WaitpidQueue::close_pid() called with invalid pid.");
m_remaining.fetch_add(1, std::memory_order_acquire);
{
std::lock_guard<std::mutex> guard(m_mutex);
// if (!m_queue.empty()) {
// m_queue.push_back(pid);
// return;
// }
m_queue.insert(pid);
}
m_wakeup_worker.store(true, std::memory_order_release);
m_wakeup_worker.notify_all();
}
void
WaitpidQueue::wait_for(uint32_t max_remaining) {
while (m_remaining.load(std::memory_order_acquire) > max_remaining)
m_remaining.wait(max_remaining, std::memory_order_acquire);
}
} // namespace torrent::utils
+44
View File
@@ -0,0 +1,44 @@
#ifndef RTORRENT_UTILS_WAITPID_QUEUE_H
#define RTORRENT_UTILS_WAITPID_QUEUE_H
#include <future>
#include <set>
#include <torrent/system/common.h>
namespace utils {
class WaitpidQueue {
public:
WaitpidQueue();
~WaitpidQueue();
uint32_t size() const;
void close_pid(int pid);
void wait_for(uint32_t max_remaining);
private:
WaitpidQueue(const WaitpidQueue&) = delete;
WaitpidQueue& operator=(const WaitpidQueue&) = delete;
std::future<void> m_worker;
align_cacheline
std::mutex m_mutex;
std::set<pid_t> m_queue;
bool m_should_shutdown{};
align_cacheline
std::atomic<bool> m_wakeup_worker{};
std::atomic<uint32_t> m_remaining{};
};
inline uint32_t WaitpidQueue::size() const { return m_remaining.load(std::memory_order_acquire); }
} // namespace utils
#endif
+7 -5
View File
@@ -9,7 +9,7 @@
#include <utility>
#include <vector>
#include <torrent/utils/scheduler.h>
#include <torrent/system/scheduler.h>
namespace utils {
@@ -49,10 +49,12 @@ private:
void update_status(Entry* entry);
void schedule();
std::map<std::string, Entry> m_entries;
std::vector<Entry*> m_entry_queue;
torrent::utils::SchedulerEntry m_task_process;
bool m_active{true};
std::map<std::string, Entry> m_entries;
std::vector<Entry*> m_entry_queue;
bool m_active{true};
torrent::system::SchedulerEntry m_task_process;
};
}