Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 45 additions & 0 deletions src/data/thread_disk.cc
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
#include "data/thread_disk.h"

#include <cassert>
#include <unistd.h>

#include "thread_main.h"
#include "data/hash_check_queue.h"
Expand Down Expand Up @@ -71,7 +72,50 @@ ThreadDisk::init_thread() {

void
ThreadDisk::cleanup_thread() {
// Drain pending closes so we do not leak fds at shutdown.
perform_close_fds();

assert(m_hash_check_queue->empty() && "ThreadDisk::cleanup_thread(): m_hash_check_queue not empty.");
assert(m_close_fds.empty() && "ThreadDisk::cleanup_thread(): m_close_fds not empty.");
}

void
ThreadDisk::queue_close_fd(int fd) {
if (fd < 0)
return;

{
auto lock = std::lock_guard(m_close_fds_lock);
m_close_fds.push_back(fd);
}

interrupt();
}

void
ThreadDisk::queue_close_fds(const std::vector<int>& fds) {
if (fds.empty())
return;

{
auto lock = std::lock_guard(m_close_fds_lock);
m_close_fds.insert(m_close_fds.end(), fds.begin(), fds.end());
}

interrupt();
}

void
ThreadDisk::perform_close_fds() {
std::deque<int> local;

{
auto lock = std::lock_guard(m_close_fds_lock);
local.swap(m_close_fds);
}

for (int fd : local)
::close(fd);
}

void
Expand All @@ -87,6 +131,7 @@ ThreadDisk::call_events() {
throw shutdown_exception();
}

perform_close_fds();
process_callbacks();
}

Expand Down
12 changes: 12 additions & 0 deletions src/data/thread_disk.h
Original file line number Diff line number Diff line change
@@ -1,7 +1,10 @@
#ifndef LIBTORRENT_DATA_THREAD_DISK_H
#define LIBTORRENT_DATA_THREAD_DISK_H

#include <deque>
#include <memory>
#include <mutex>
#include <vector>

#include "torrent/common.h"
#include "torrent/system/thread.h"
Expand All @@ -25,15 +28,24 @@ class LIBTORRENT_EXPORT ThreadDisk : public system::Thread {

HashCheckQueue* hash_check_queue() { return m_hash_check_queue.get(); }

// Detach on main first, then queue the raw fd here. Disk thread owns ::close.
void queue_close_fd(int fd);
void queue_close_fds(const std::vector<int>& fds);

private:
ThreadDisk() = default;

void call_events() override;
std::chrono::microseconds next_timeout() override;

void perform_close_fds();

static ThreadDisk* m_thread_disk;

std::unique_ptr<HashCheckQueue> m_hash_check_queue;

std::mutex m_close_fds_lock;
std::deque<int> m_close_fds;
};

} // namespace torrent
Expand Down
14 changes: 12 additions & 2 deletions src/torrent/data/file_list.cc
Original file line number Diff line number Diff line change
Expand Up @@ -484,11 +484,16 @@ FileList::close() {

LT_LOG_FL(INFO, "Closing.", 0);

std::vector<File*> files;
files.reserve(size());

for (auto& entry : *this) {
entry->unset_flags_protected(File::flag_active);
manager->file_manager()->close(entry.get());
files.push_back(entry.get());
}

manager->file_manager()->close_files(files);

m_is_open = false;
m_indirect_links.clear();

Expand All @@ -502,8 +507,13 @@ FileList::close_all_files() {

LT_LOG_FL(INFO, "Closing all files.", 0);

std::vector<File*> files;
files.reserve(size());

for (auto& entry : *this)
manager->file_manager()->close(entry.get());
files.push_back(entry.get());

manager->file_manager()->close_files(files);
}

void
Expand Down
83 changes: 73 additions & 10 deletions src/torrent/data/file_manager.cc
Original file line number Diff line number Diff line change
Expand Up @@ -6,14 +6,55 @@
#include <cassert>
#include <fcntl.h>
#include <limits>
#include <unistd.h>

#include "manager.h"
#include "data/socket_file.h"
#include "data/thread_disk.h"
#include "torrent/exceptions.h"
#include "torrent/data/file.h"

namespace torrent {

namespace {

// Synchronous close on the calling thread (main under pressure / tests).
void
close_fd_now(int fd) {
if (fd < 0)
return;

::close(fd);
}

// Prefer disk thread so main/UI is not blocked on slow FS close.
void
close_fd_deferred(int fd) {
if (fd < 0)
return;

if (ThreadDisk* disk = ThreadDisk::thread_disk(); disk != nullptr && disk->is_active())
disk->queue_close_fd(fd);
else
close_fd_now(fd);
}

void
close_fds_deferred(const std::vector<int>& fds) {
if (fds.empty())
return;

if (ThreadDisk* disk = ThreadDisk::thread_disk(); disk != nullptr && disk->is_active()) {
disk->queue_close_fds(fds);
return;
}

for (int fd : fds)
close_fd_now(fd);
}

} // namespace

FileManager::~FileManager() {
assert(empty() && "FileManager::~FileManager() called but empty() != true.");
}
Expand Down Expand Up @@ -71,28 +112,49 @@ FileManager::open(value_type file, [[maybe_unused]] bool hashing, int prot, int
return true;
}

void
FileManager::close(value_type file) {
if (!file->is_open())
return;

if (file->is_padding())
return;
int
FileManager::detach(value_type file) {
if (file == nullptr || !file->is_open() || file->is_padding())
return -1;

SocketFile(file->file_descriptor()).close();
int fd = file->file_descriptor();

file->set_protection(0);
file->reset_file_descriptor();

auto itr = std::find(begin(), end(), file);

if (itr == end())
throw internal_error("FileManager::close_file(...) itr == end().");
throw internal_error("FileManager::detach(...) itr == end().");

*itr = back();
base_type::pop_back();

m_files_closed_counter++;
return fd;
}

void
FileManager::close(value_type file) {
close_fd_deferred(detach(file));
}

void
FileManager::close_files(const std::vector<value_type>& files) {
if (files.empty())
return;

std::vector<int> fds;
fds.reserve(files.size());

for (value_type file : files) {
int fd = detach(file);

if (fd >= 0)
fds.push_back(fd);
}

close_fds_deferred(fds);
}

void
Expand All @@ -107,8 +169,9 @@ FileManager::close_least_active() {
}
}

// Free a kernel slot immediately so open-at-cap does not soft-overshoot.
if (least)
close(least);
close_fd_now(detach(least));
}

} // namespace torrent
6 changes: 6 additions & 0 deletions src/torrent/data/file_manager.h
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,9 @@ class LIBTORRENT_EXPORT FileManager : private std::vector<File*> {
bool open(value_type file, bool hashing, int prot, int flags);
void close(value_type file);

// Detach files and queue their fds for close on the disk thread.
void close_files(const std::vector<value_type>& files);

// TODO: Close all files held by a download after hashing. Also flush all memory chunks.

void close_least_active();
Expand All @@ -52,6 +55,9 @@ class LIBTORRENT_EXPORT FileManager : private std::vector<File*> {
FileManager(const FileManager&) = delete;
FileManager& operator=(const FileManager&) = delete;

// Detach bookkeeping and return the raw fd (or -1). FileManager close paths own ::close.
int detach(value_type file);

size_type m_max_open_files{0};
bool m_advise_random{false};
bool m_advise_random_hashing{false};
Expand Down
Loading