Compare commits

...

28 Commits

Author SHA1 Message Date
rakshasa 9507ac17b3 Stuff. 2026-08-04 11:08:52 +02:00
rakshasa 49d45c1088 Stuff. 2026-08-04 11:04:20 +02:00
rakshasa 5ab973504a Stuff. 2026-08-04 10:47:23 +02:00
rakshasa ebe28e07b1 Stuff. 2026-08-04 10:40:43 +02:00
rakshasa 06a391f9d8 Stuff. 2026-08-03 09:11:23 +02:00
rakshasa 247ae7621a Tagged release 0.16.19. 2026-08-02 10:44:30 +02:00
xirvik 05831942a7 commands: don't register min_alloc/max_alloc for the generic socket category
SocketManager throws internal_error for category_generic, aborting rtorrent over RPC.
2026-07-30 12:42:24 +02:00
rakshasa ee18a1f3fc Added commit message linter. 2026-07-28 10:04:51 +02:00
Jari Sundell f52f20f1b6 Cleaned up system headers. 2026-07-26 19:53:18 +09:00
xirvik 41eef969ea commands: fix group.<name>.ratio.disable calling a removed command
group_insert() builds the disable method as "schedule_remove=...", but that
command no longer exists: the current name is schedule.remove, and the only
old-style aliases registered are schedule2 and schedule_remove2 (and those only
when deprecated commands are enabled at registration time). Every
group.<name>.ratio.disable therefore faults with

  Command "schedule_remove" does not exist.

including the built-in group.seeding.ratio.disable, so a ratio group can be
enabled but never disabled. The enable half already uses the current name.
2026-07-25 10:18:49 +02:00
fffe 1def5e8b36 Add log.print and enum.log_group methods 2026-07-22 11:08:32 +02:00
simonc56 8355870586 Mark 'd.timestamp.started' as safe in command download initialization 2026-07-20 08:21:32 +02:00
Jari Sundell ba2bc7e64d Various release bug fixes. 2026-07-19 18:53:57 +09:00
rakshasa 0c11deac50 Tagged release 0.16.18. 2026-07-17 10:15:28 +02:00
Jari Sundell f6b3ad0efd Minor cleanup of pex and other commands. 2026-07-15 19:40:13 +09:00
Jari Sundell 5d350f0bcc Use new encryption modes. 2026-07-14 17:34:14 +09:00
Jari Sundell 86fa0d195f Changing listen/dht port changes the listening/dht ports. 2026-07-12 17:19:33 +09:00
Jari Sundell 2354c9cdb5 Moved ChunkManager out of public header directory. 2026-07-10 17:54:37 +09:00
Jari Sundell 1b25bc7f56 Moved all sync/diskspace related methods to MemoryManager. 2026-07-10 02:21:08 +09:00
Jari Sundell c6bde213b4 Added runtime::MemoryManager to split out unrelated features from ChunkManager. 2026-07-09 01:56:30 +09:00
Jari Sundell 959448acda Fixed command_base t_pod align static asserts. 2026-07-08 17:15:56 +09:00
Jari Sundell 3b6da6feac Reordered http queue slots to avoid race conditions. 2026-07-07 16:34:19 +09:00
Jari Sundell ea2d22cb87 Removed unused add/remove error-event code. 2026-07-07 03:05:20 +09:00
rakshasa fa351c017d Tagged release 0.16.17. 2026-07-06 14:06:52 +02:00
Pluto Yang 1e1bdd551a Add LoongArch CPU support 2026-07-06 13:38:08 +02:00
Jari Sundell a07f57a703 Fixed fd close order in ExecFile. 2026-07-06 18:06:51 +09:00
Jari Sundell 0494ce70c8 Add http done/failed slots before starting request. 2026-07-06 17:18:39 +09:00
Xeonacid cfad36bf12 Add RISC-V cacheline fallbacks
The Linux cacheline probe first tries <linux/cache.h>, but that header is not part of the installed uapi headers on Arch Linux. When the compile probe fails, configure falls back to a host_cpu mapping and currently aborts for riscv64 with:

  Unrecognized CPU architecture (riscv64) on Linux fallback path.

Handle riscv* in both cacheline fallback maps and use a 64-byte cacheline. That matches the Linux RISC-V kernel default L1_CACHE_SHIFT value of 6 and the value reported on a riscv64 machine through getconf LEVEL1_DCACHE_LINESIZE and sysfs cache coherency_line_size.
2026-07-05 15:40:50 +02:00
42 changed files with 729 additions and 564 deletions
+65
View File
@@ -0,0 +1,65 @@
name: "Lint Commit Message Size"
on:
pull_request:
types: [opened, synchronize, reopened]
jobs:
check-commit-bounds:
runs-on: ubuntu-latest
steps:
- name: Check out code
uses: actions/checkout@v4
with:
fetch-depth: 0
- name: Validate Line Count, Width, and Spacing
run: |
# Fetch commit hashes unique to this PR branch
COMMITS=$(git log --no-merges --pretty=format:"%H" origin/${{ github.base_ref }}..HEAD)
MAX_LINES=3
MAX_CHARS=90
FAILED=0
for commit in $COMMITS; do
SUBJECT=$(git log --format="%s" -n 1 $commit)
# Extract clean commit message, trimming trailing blank lines
RAW_MSG=$(git log --format="%B" -n 1 $commit)
CLEAN_MSG=$(echo "$RAW_MSG" | awk '{msg[NR]=$0} END {while(NR>0 && msg[NR]=="") NR--; for(i=1;i<=NR;i++) print msg[i]}')
# 1. Check total line count
LINE_COUNT=$(echo "$CLEAN_MSG" | wc -l)
if [ "$LINE_COUNT" -gt "$MAX_LINES" ]; then
echo "❌ Error: Commit message has too many lines ($LINE_COUNT/$MAX_LINES)."
echo " Commit: '$SUBJECT'"
FAILED=1
fi
# 2. Check for multiple consecutive blank lines
# This regex looks for 2 or more empty lines anywhere in the message
if echo "$CLEAN_MSG" | grep -pz '(\r?\n){3,}'; then
echo "❌ Error: Commit message contains multiple consecutive line breaks."
echo " Commit: '$SUBJECT'"
FAILED=1
fi
# 3. Check maximum width of any individual line
while IFS= read -r line; do
LINE_LENGTH=${#line}
if [ "$LINE_LENGTH" -gt "$MAX_CHARS" ]; then
echo "❌ Error: Line length exceeds limit ($LINE_LENGTH / $MAX_CHARS chars)."
echo " Offending line: '$line'"
FAILED=1
fi
done <<< "$CLEAN_MSG"
done
# Fail the job if any check failed
if [ "$FAILED" -ne 0 ]; then
exit 1
fi
echo "✅ All commit messages passed style rules!"
+3 -3
View File
@@ -1,6 +1,6 @@
m4_pattern_allow([PKG_CHECK_EXISTS])
AC_INIT([rtorrent],[0.16.16],[sundell.software@gmail.com])
AC_INIT([rtorrent],[0.16.19],[sundell.software@gmail.com])
AC_CONFIG_HEADERS([config.h])
AC_CONFIG_MACRO_DIRS([scripts])
@@ -14,7 +14,7 @@ AX_CXX_COMPILE_STDCXX(20, noext, mandatory)
PKG_PROG_PKG_CONFIG
AC_DEFINE([API_VERSION], [23], [api version])
AC_DEFINE([API_VERSION], [24], [api version])
RAK_CHECK_CFLAGS
RAK_CHECK_CXXFLAGS
@@ -47,7 +47,7 @@ fi
PKG_CHECK_MODULES([CPPUNIT], [cppunit],, [no_cppunit="yes"])
PKG_CHECK_MODULES([ZLIB], [zlib])
PKG_CHECK_MODULES([DEPENDENCIES], [libtorrent >= 0.16.16])
PKG_CHECK_MODULES([DEPENDENCIES], [libtorrent >= 0.16.19])
AC_LANG_PUSH(C++)
TORRENT_WITH_XMLRPC_C
+8
View File
@@ -118,6 +118,14 @@ AC_DEFUN([TORRENT_CHECK_CACHELINE], [
AC_MSG_RESULT([linux fallback enterprise 128 bytes])
AC_DEFINE([LT_SMP_CACHE_BYTES], 128, [Fallback 128-byte alignment for Linux enterprise hardware.])
;;
riscv32*|riscv64*)
AC_MSG_RESULT([linux fallback RISC-V 64 bytes])
AC_DEFINE([LT_SMP_CACHE_BYTES], 64, [Fallback 64-byte alignment for Linux RISC-V hardware.])
;;
loongarch32*|loongarch64*)
AC_MSG_RESULT([linux fallback LoongArch 64 bytes])
AC_DEFINE([LT_SMP_CACHE_BYTES], 64, [Fallback 64-byte alignment for Linux LoongArch hardware.])
;;
*)
AC_MSG_RESULT([unrecognized CPU arch on Linux header fallback])
AC_MSG_FAILURE([Unrecognized CPU architecture ($host_cpu) on Linux fallback path. Aborting build.])
+2
View File
@@ -174,6 +174,8 @@ libsub_root_a_SOURCES = \
utils/list_focus.h \
utils/lockfile.cc \
utils/lockfile.h \
utils/waitpid_queue.cc \
utils/waitpid_queue.h \
utils/watch_ready_queue.cc \
utils/watch_ready_queue.h \
\
+1
View File
@@ -974,6 +974,7 @@ initialize_command_download() {
rpc::rpc.mark_safe("d.size_pex");
rpc::rpc.mark_safe("d.completed_bytes");
rpc::rpc.mark_safe("d.complete");
rpc::rpc.mark_safe("d.timestamp.started");
rpc::rpc.mark_safe("d.timestamp.finished");
rpc::rpc.mark_safe("d.bytes_done");
rpc::rpc.mark_safe("d.peers_accounted");
+16 -10
View File
@@ -444,16 +444,20 @@ initialize_command_dynamic() {
CMD2_ANY ("catch", std::bind(&cmd_catch, std::placeholders::_1, std::placeholders::_2));
CMD2_ANY ("strings.choke_heuristics", std::bind(&torrent::option_list_strings, torrent::OPTION_CHOKE_HEURISTICS));
CMD2_ANY ("strings.choke_heuristics.upload", std::bind(&torrent::option_list_strings, torrent::OPTION_CHOKE_HEURISTICS_UPLOAD));
CMD2_ANY ("strings.choke_heuristics.download", std::bind(&torrent::option_list_strings, torrent::OPTION_CHOKE_HEURISTICS_DOWNLOAD));
CMD2_ANY ("strings.connection_type", std::bind(&torrent::option_list_strings, torrent::OPTION_CONNECTION_TYPE));
CMD2_ANY ("strings.encryption", std::bind(&torrent::option_list_strings, torrent::OPTION_ENCRYPTION));
CMD2_ANY ("strings.ip_filter", std::bind(&torrent::option_list_strings, torrent::OPTION_IP_FILTER));
CMD2_ANY ("strings.ip_tos", std::bind(&torrent::option_list_strings, torrent::OPTION_IP_TOS));
CMD2_ANY ("strings.log_group", std::bind(&torrent::option_list_strings, torrent::OPTION_LOG_GROUP));
CMD2_ANY ("strings.tracker_event", std::bind(&torrent::option_list_strings, torrent::OPTION_TRACKER_EVENT));
CMD2_ANY ("strings.tracker_mode", std::bind(&torrent::option_list_strings, torrent::OPTION_TRACKER_MODE));
CMD2_ANY_STRING ("enum.log_group", [](auto, const auto& str) { return torrent::option_find_string_str(torrent::OPTION_LOG_GROUP, str); });
CMD2_ANY ("strings.choke_heuristics", [](auto, auto) { return torrent::option_list_strings(torrent::OPTION_CHOKE_HEURISTICS); });
CMD2_ANY ("strings.choke_heuristics.upload", [](auto, auto) { return torrent::option_list_strings(torrent::OPTION_CHOKE_HEURISTICS_UPLOAD); });
CMD2_ANY ("strings.choke_heuristics.download", [](auto, auto) { return torrent::option_list_strings(torrent::OPTION_CHOKE_HEURISTICS_DOWNLOAD); });
CMD2_ANY ("strings.connection_type", [](auto, auto) { return torrent::option_list_strings(torrent::OPTION_CONNECTION_TYPE); });
CMD2_ANY ("strings.encryption", [](auto, auto) { return torrent::Object::create_list(); });
CMD2_ANY ("strings.encryption.handshake", [](auto, auto) { return torrent::option_list_strings(torrent::OPTION_ENCRYPTION_HANDSHAKE); });
CMD2_ANY ("strings.encryption.stream", [](auto, auto) { return torrent::option_list_strings(torrent::OPTION_ENCRYPTION_STREAM); });
CMD2_ANY ("strings.ip_filter", [](auto, auto) { return torrent::option_list_strings(torrent::OPTION_IP_FILTER); });
CMD2_ANY ("strings.ip_tos", [](auto, auto) { return torrent::option_list_strings(torrent::OPTION_IP_TOS); });
CMD2_ANY ("strings.log_group", [](auto, auto) { return torrent::option_list_strings(torrent::OPTION_LOG_GROUP); });
CMD2_ANY ("strings.tracker_event", [](auto, auto) { return torrent::option_list_strings(torrent::OPTION_TRACKER_EVENT); });
CMD2_ANY ("strings.tracker_mode", [](auto, auto) { return torrent::option_list_strings(torrent::OPTION_TRACKER_MODE); });
// clang-format on
#ifdef HAVE_XMLRPC_TINYXML2
@@ -468,6 +472,8 @@ initialize_command_dynamic() {
rpc::rpc.mark_safe("method.rlookup");
rpc::rpc.mark_safe("catch");
rpc::rpc.mark_safe("enum.log_group");
rpc::rpc.mark_safe("strings.choke_heuristics");
rpc::rpc.mark_safe("strings.choke_heuristics.upload");
rpc::rpc.mark_safe("strings.choke_heuristics.download");
+32 -31
View File
@@ -8,10 +8,10 @@
#include <sys/types.h>
#include <sys/stat.h>
#include <torrent/torrent.h>
#include <torrent/chunk_manager.h>
#include <torrent/data/file_manager.h>
#include <torrent/data/chunk_utils.h>
#include <torrent/runtime/runtime.h>
#include <torrent/runtime/memory_manager.h>
#include <torrent/runtime/socket_manager.h>
#include <torrent/utils/chrono.h>
#include <torrent/utils/option_strings.h>
@@ -117,7 +117,7 @@ group_insert(const torrent::Object::list_type& args) {
rpc::commands.call("method.insert", rpc::create_object_list("group." + name + ".ratio.enable", "simple",
"schedule=group." + name + ".ratio,5,60,on_ratio=" + name));
rpc::commands.call("method.insert", rpc::create_object_list("group." + name + ".ratio.disable", "simple",
"schedule_remove=group." + name + ".ratio"));
"schedule.remove=group." + name + ".ratio"));
rpc::commands.call("method.insert", rpc::create_object_list("group." + name + ".ratio.command", "simple",
"d.try_close= ;d.ignore_commands.set=1"));
rpc::commands.call("method.insert", rpc::create_object_list("group." + name + ".view", "string", view));
@@ -188,13 +188,8 @@ cmd_file_append(const torrent::Object::list_type& args) {
void
initialize_command_local() {
core::DownloadList* dList = control->core()->download_list();
torrent::ChunkManager* chunkManager = torrent::chunk_manager();
torrent::FileManager* fileManager = torrent::file_manager();
if (rpc::call_command_value("method.use_deprecated") == 1) {
CMD_ANY_LIST ("file.append", std::bind(&cmd_file_append, std::placeholders::_2));
}
CMD_ANY ("system.hostname", std::bind(&system_hostname));
CMD_ANY ("system.pid", std::bind(&getpid));
@@ -255,42 +250,44 @@ initialize_command_local() {
CMD_ANY (category_name + ".size", [category](auto, auto) { return torrent::runtime::socket_manager()->category_managed_size(category); });
CMD_ANY (category_name + ".max_size", [category](auto, auto) { return torrent::runtime::socket_manager()->category_max_size(category); });
CMD_ANY (category_name + ".min_alloc", [category](auto, auto) { return torrent::runtime::socket_manager()->category_min_allocation(category); });
CMD_ANY (category_name + ".max_alloc", [category](auto, auto) { return torrent::runtime::socket_manager()->category_max_allocation(category); });
if (i == 0)
continue;
CMD_ANY (category_name + ".min_alloc", [category](auto, auto) { return torrent::runtime::socket_manager()->category_min_allocation(category); });
CMD_ANY (category_name + ".max_alloc", [category](auto, auto) { return torrent::runtime::socket_manager()->category_max_allocation(category); });
CMD_ANY_VALUE_V(category_name + ".min_alloc.set", [category](auto, auto& value) { torrent::runtime::socket_manager()->set_category_min_allocation(category, value); });
CMD_ANY_VALUE_V(category_name + ".max_alloc.set", [category](auto, auto& value) { torrent::runtime::socket_manager()->set_category_max_allocation(category, value); });
}
CMD_ANY ("pieces.sync.always_safe", std::bind(&CM_t::safe_sync, chunkManager));
CMD_ANY_VALUE_V ("pieces.sync.always_safe.set", std::bind(&CM_t::set_safe_sync, chunkManager, std::placeholders::_2));
CMD_ANY ("pieces.sync.safe_free_diskspace", std::bind(&CM_t::safe_free_diskspace, chunkManager));
CMD_ANY ("pieces.sync.timeout", std::bind(&CM_t::timeout_sync, chunkManager));
CMD_ANY_VALUE_V ("pieces.sync.timeout.set", std::bind(&CM_t::set_timeout_sync, chunkManager, std::placeholders::_2));
CMD_ANY ("pieces.sync.timeout_safe", std::bind(&CM_t::timeout_safe_sync, chunkManager));
CMD_ANY_VALUE_V ("pieces.sync.timeout_safe.set", std::bind(&CM_t::set_timeout_safe_sync, chunkManager, std::placeholders::_2));
CMD_ANY ("pieces.sync.queue_size", std::bind(&CM_t::sync_queue_size, chunkManager));
CMD_ANY ("pieces.sync.always_safe", [](auto, auto) { return torrent::runtime::memory_manager()->safe_sync(); });
CMD_ANY_VALUE_V ("pieces.sync.always_safe.set", [](auto, auto& value) { return torrent::runtime::memory_manager()->set_safe_sync(value); });
CMD_ANY ("pieces.sync.safe_free_diskspace", [](auto, auto) { return torrent::runtime::memory_manager()->sync_safe_free_diskspace(); });
CMD_ANY ("pieces.sync.timeout", [](auto, auto) { return torrent::runtime::memory_manager()->timeout_sync().count(); });
CMD_ANY_VALUE_V ("pieces.sync.timeout.set", [](auto, auto& value) { return torrent::runtime::memory_manager()->set_timeout_sync(value); });
// CMD_ANY ("pieces.sync.timeout_safe", [](auto, auto) { return torrent::runtime::memory_manager()->timeout_safe_sync(); });
// CMD_ANY_VALUE_V ("pieces.sync.timeout_safe.set", [](auto, auto& value) { return torrent::runtime::memory_manager()->set_timeout_safe_sync(value); });
CMD_ANY ("pieces.sync.timeout_safe", [](auto, auto) { return 0; });
CMD_ANY_VALUE_V ("pieces.sync.timeout_safe.set", [](auto, auto) { });
CMD_ANY ("pieces.sync.queue_size", [](auto, auto) { return torrent::runtime::memory_manager()->sync_queue_block_count(); });
CMD_ANY ("pieces.preload.type", std::bind(&CM_t::preload_type, chunkManager));
CMD_ANY_VALUE_V ("pieces.preload.type.set", std::bind(&CM_t::set_preload_type, chunkManager, std::placeholders::_2));
CMD_ANY ("pieces.preload.min_size", std::bind(&CM_t::preload_min_size, chunkManager));
CMD_ANY_VALUE_V ("pieces.preload.min_size.set", std::bind(&CM_t::set_preload_min_size, chunkManager, std::placeholders::_2));
CMD_ANY ("pieces.preload.min_rate", std::bind(&CM_t::preload_required_rate, chunkManager));
CMD_ANY_VALUE_V ("pieces.preload.min_rate.set", std::bind(&CM_t::set_preload_required_rate, chunkManager, std::placeholders::_2));
CMD_ANY ("pieces.memory.current", std::bind(&CM_t::memory_usage, chunkManager));
CMD_ANY ("pieces.memory.sync_queue", std::bind(&CM_t::sync_queue_memory_usage, chunkManager));
CMD_ANY ("pieces.memory.block_count", std::bind(&CM_t::memory_block_count, chunkManager));
CMD_ANY ("pieces.memory.max", std::bind(&CM_t::max_memory_usage, chunkManager));
CMD_ANY_VALUE_V ("pieces.memory.max.set", std::bind(&CM_t::set_max_memory_usage, chunkManager, std::placeholders::_2));
CMD_ANY ("pieces.stats_preloaded", std::bind(&CM_t::stats_preloaded, chunkManager));
CMD_ANY ("pieces.stats_not_preloaded", std::bind(&CM_t::stats_not_preloaded, chunkManager));
CMD_ANY ("pieces.preload.type", [](auto, auto) { return torrent::runtime::memory_manager()->preload_type(); });
CMD_ANY_VALUE_V ("pieces.preload.type.set", [](auto, auto& value) { return torrent::runtime::memory_manager()->set_preload_type(value); });
CMD_ANY ("pieces.preload.min_size", [](auto, auto) { return torrent::runtime::memory_manager()->preload_min_size(); });
CMD_ANY_VALUE_V ("pieces.preload.min_size.set", [](auto, auto& value) { return torrent::runtime::memory_manager()->set_preload_min_size(value); });
CMD_ANY ("pieces.preload.min_rate", [](auto, auto) { return torrent::runtime::memory_manager()->preload_required_rate(); });
CMD_ANY_VALUE_V ("pieces.preload.min_rate.set", [](auto, auto& value) { return torrent::runtime::memory_manager()->set_preload_required_rate(value); });
CMD_ANY ("pieces.stats_preloaded", [](auto, auto) { return torrent::runtime::memory_manager()->stats_preloaded(); });
CMD_ANY ("pieces.stats_not_preloaded", [](auto, auto) { return torrent::runtime::memory_manager()->stats_not_preloaded(); });
CMD_ANY ("pieces.stats.total_size", std::bind(&apply_pieces_stats_total_size));
CMD_ANY ("pieces.memory.current", [](auto, auto) { return torrent::runtime::memory_manager()->memory_usage(); });
CMD_ANY ("pieces.memory.sync_queue", [](auto, auto) { return torrent::runtime::memory_manager()->sync_queue_memory_usage(); });
CMD_ANY ("pieces.memory.block_count", [](auto, auto) { return torrent::runtime::memory_manager()->memory_block_count(); });
CMD_ANY ("pieces.memory.max", [](auto, auto) { return torrent::runtime::memory_manager()->max_memory_usage(); });
CMD_ANY_VALUE_V ("pieces.memory.max.set", [](auto, auto& value) { return torrent::runtime::memory_manager()->set_max_memory_usage(value); });
CMD_ANY ("pieces.hash.queue_size", std::bind(&torrent::main_thread::hash_queue_size));
CMD_VAR_BOOL ("pieces.hash.on_completion", true);
@@ -357,6 +354,10 @@ initialize_command_local() {
rpc::rpc.mark_safe(category_name + ".size");
rpc::rpc.mark_safe(category_name + ".max_size");
if (i == 0)
continue;
rpc::rpc.mark_safe(category_name + ".min_alloc");
rpc::rpc.mark_safe(category_name + ".max_alloc");
}
+29
View File
@@ -1,6 +1,7 @@
#include "config.h"
#include <fcntl.h>
#include <iterator>
#include <stdio.h>
#include <unistd.h>
#include <torrent/data/chunk_utils.h>
@@ -15,6 +16,7 @@
#include "core/download.h"
#include "core/download_list.h"
#include "core/manager.h"
#include "rpc/parse.h"
#include "rpc/parse_commands.h"
torrent::Object
@@ -38,6 +40,32 @@ apply_log_add_output(const torrent::Object::list_type& args) {
return torrent::Object();
}
torrent::Object
apply_log_print(const torrent::Object::list_type& args) {
if (args.size() < 2)
throw torrent::input_error("Invalid number of arguments.");
torrent::Object::value_type group;
if (args.front().is_value())
group = args.front().as_value();
else if (args.front().is_string())
group = torrent::option_find_string_str(torrent::OPTION_LOG_GROUP, args.front().as_string());
else
throw torrent::input_error("Invalid log group.");
if (group < 0 || group >= torrent::LOG_GROUP_MAX_SIZE)
throw torrent::input_error("Invalid log group.");
std::string message;
for (auto itr = std::next(args.begin()); itr != args.end(); ++itr)
rpc::print_object_std(&message, &*itr, 0);
torrent::log_groups[group].internal_print(message);
return torrent::Object();
}
// TODO: Deprecated.
torrent::Object
apply_log(const torrent::Object::string_type& arg, int logType) {
@@ -112,6 +140,7 @@ initialize_command_logging() {
CMD2_ANY_STRING_V("log.close", std::bind(&torrent::log_close_output_str, std::placeholders::_2));
CMD2_ANY_LIST ("log.add_output", std::bind(&apply_log_add_output, std::placeholders::_2));
CMD2_ANY_LIST ("log.print", std::bind(&apply_log_print, std::placeholders::_2));
CMD2_ANY_STRING ("log.execute", std::bind(&apply_log, std::placeholders::_2, 0));
CMD2_ANY_STRING ("log.vmmap.dump", std::bind(&log_vmmap_dump, std::placeholders::_2));
+124 -21
View File
@@ -9,7 +9,9 @@
#include <torrent/download/resource_manager.h>
#include <torrent/net/http_stack.h>
#include <torrent/net/socket_address.h>
#include <torrent/runtime/client_config.h>
#include <torrent/runtime/network_config.h>
#include <torrent/runtime/network_manager.h>
#include <torrent/runtime/proxy_manager.h>
#include <torrent/runtime/runtime.h>
#include <torrent/runtime/socket_manager.h>
@@ -33,21 +35,115 @@
#endif
torrent::Object
apply_encryption(const torrent::Object::list_type& args) {
uint32_t options_mask = torrent::runtime::NetworkConfig::encryption_none;
listen_port_range() {
auto port_range = torrent::runtime::client_config()->listen_port_range();
for (const auto& arg : args) {
uint32_t opt = torrent::option_find_string(torrent::OPTION_ENCRYPTION, arg.as_string().c_str());
return std::to_string(port_range.first) + "-" + std::to_string(port_range.second);
}
if (opt == torrent::runtime::NetworkConfig::encryption_none)
options_mask = torrent::runtime::NetworkConfig::encryption_none;
else
options_mask |= opt;
void
set_listen_port_range(const std::string& arg) {
unsigned int port_first{}, port_last{};
if (std::sscanf(arg.c_str(), "%i-%i", &port_first, &port_last) != 2)
throw torrent::input_error("Invalid port_range argument.");
if (port_first >= (1 << 16) || port_last >= (1 << 16))
throw torrent::input_error("Port range out-of-bounds.");
torrent::runtime::client_config()->set_listen_port_range(port_first, port_last);
}
torrent::Object
get_encryption() {
auto encryption_modes = torrent::runtime::network_config()->encryption_modes();
return torrent::option_to_str_or_throw(torrent::OPTION_ENCRYPTION_HANDSHAKE, encryption_modes.first) + "," +
torrent::option_to_str_or_throw(torrent::OPTION_ENCRYPTION_STREAM, encryption_modes.second);
}
torrent::Object
get_handshake_encryption() {
auto encryption_modes = torrent::runtime::network_config()->encryption_modes();
return torrent::option_to_str_or_throw(torrent::OPTION_ENCRYPTION_MODE, encryption_modes.first);
}
torrent::Object
get_stream_encryption() {
auto encryption_modes = torrent::runtime::network_config()->encryption_modes();
return torrent::option_to_str_or_throw(torrent::OPTION_ENCRYPTION_MODE, encryption_modes.second);
}
torrent::Object
apply_obsolete_encryption(const torrent::Object::list_type& args) {
torrent::encryption_mode handshake_mode{torrent::ENCRYPTION_MODE_ALLOW};
torrent::encryption_mode stream_mode{torrent::ENCRYPTION_MODE_ALLOW};
for (auto& itr : args) {
auto arg = itr.as_string();
if (arg == "none") {
handshake_mode = torrent::ENCRYPTION_MODE_DENY;
stream_mode = torrent::ENCRYPTION_MODE_DENY;
break;
} else if (arg == "allow_incoming") {
} else if (arg == "try_outgoing") {
} else if (arg == "require") {
handshake_mode = torrent::ENCRYPTION_MODE_REQUIRE;
} else if (arg == "require_RC4" || arg == "require_rc4") {
handshake_mode = torrent::ENCRYPTION_MODE_REQUIRE;
stream_mode = torrent::ENCRYPTION_MODE_REQUIRE;
break;
} else if (arg == "enable_retry") {
} else if (arg == "prefer_plaintext") {
} else {
throw torrent::input_error("Invalid encryption option: '" + arg + "'");
}
}
torrent::runtime::network_config()->set_encryption_options(options_mask);
lt_log_print(torrent::LOG_WARN, "Obsolete encryption options used, use 'handshake_{deny,allow,prefer,require}, stream_{deny,allow,prefer,require}' instead.");
return torrent::Object();
torrent::runtime::network_config()->set_encryption_modes(handshake_mode, stream_mode);
return {};
}
torrent::Object
apply_encryption(const torrent::Object::list_type& args) {
if (args.empty())
throw torrent::input_error("No encryption options specified.");
torrent::encryption_mode encryption_mode, handshake_mode, stream_mode;
if (args.size() == 1) {
try {
encryption_mode = static_cast<torrent::encryption_mode>(torrent::option_find_string_str(torrent::OPTION_ENCRYPTION_MODE, args.front().as_string()));
} catch (torrent::input_error& e) {
return apply_obsolete_encryption(args);
}
torrent::runtime::network_config()->set_encryption_modes(encryption_mode, encryption_mode);
return {};
}
if (args.size() != 2)
return apply_obsolete_encryption(args);
try {
handshake_mode = static_cast<torrent::encryption_mode>(torrent::option_find_string_str(torrent::OPTION_ENCRYPTION_HANDSHAKE, args.front().as_string()));
stream_mode = static_cast<torrent::encryption_mode>(torrent::option_find_string_str(torrent::OPTION_ENCRYPTION_STREAM, args.back().as_string()));
} catch (torrent::input_error& e) {
return apply_obsolete_encryption(args);
}
torrent::runtime::network_config()->set_encryption_modes(handshake_mode, stream_mode);
return {};
}
torrent::Object
@@ -220,17 +316,22 @@ initialize_command_network() {
auto http_stack = torrent::net_thread::http_stack();
auto nw_config = torrent::runtime::network_config();
// Isn't port_open used?
CMD_VAR_BOOL ("network.port_open", true);
CMD_VAR_BOOL ("network.port_random", true);
CMD_VAR_STRING ("network.port_range", "6881-6999");
CMD_ANY ("network.listen.port", [](auto, auto) { return torrent::runtime::network_manager()->listen_port(); });
CMD_ANY_VALUE_V ("network.listen.port.set", [](auto, auto& value) { return torrent::runtime::network_manager()->set_listen_port(value); });
CMD_ANY ("network.listen.port.random", [](auto, auto) { return torrent::runtime::client_config()->listen_port_random(); });
CMD_ANY_VALUE_V ("network.listen.port.random.set", [](auto, auto& value) { return torrent::runtime::client_config()->set_listen_port_random(value); });
CMD_ANY ("network.listen.port.range", [](auto, auto) { return listen_port_range(); });
CMD_ANY_STRING_V("network.listen.port.range.set", [](auto, auto& value) { return set_listen_port_range(value); });
CMD_ANY ("network.listen.backlog", [](auto, auto) { return torrent::runtime::network_config()->listen_backlog(); });
CMD_ANY_VALUE_V ("network.listen.backlog.set", [](auto, auto& value) { return torrent::runtime::network_config()->set_listen_backlog(value); });
CMD_ANY ("network.listen.port", [](auto, auto) { return torrent::runtime::listen_port(); });
CMD_ANY ("network.listen.backlog", [nw_config](auto, auto) { return nw_config->listen_backlog(); });
CMD_ANY_VALUE_V ("network.listen.backlog.set", [nw_config](auto, auto& value) { return nw_config->set_listen_backlog(value); });
CMD_ANY ("protocol.pex", [](auto, auto) { return torrent::runtime::client_config()->is_pex_enabled(); });
CMD_ANY_VALUE_V ("protocol.pex.set", [](auto, auto& value) { return torrent::runtime::client_config()->set_pex_enabled(value); });
CMD_VAR_BOOL ("protocol.pex", true);
CMD_ANY_LIST ("protocol.encryption.set", [](auto, auto& args) { return apply_encryption(args); });
CMD_ANY_LIST ("protocol.encryption", [](auto, auto) { return get_encryption(); });
CMD_ANY_LIST ("protocol.encryption.set", [](auto, auto& args) { return apply_encryption(args); });
CMD_ANY_LIST ("protocol.encryption.handshake", [](auto, auto) { return get_handshake_encryption(); });
CMD_ANY_LIST ("protocol.encryption.stream", [](auto, auto) { return get_stream_encryption(); });
CMD_VAR_STRING ("protocol.connection.leech", "leech");
CMD_VAR_STRING ("protocol.connection.seed", "seed");
@@ -301,8 +402,10 @@ initialize_command_network() {
CMD_ANY ("network.xmlrpc.size_limit", [](auto, auto) { return rpc::rpc.size_limit(); });
CMD_ANY_VALUE_V ("network.xmlrpc.size_limit.set", [](auto, auto& arg) { return rpc::rpc.set_size_limit(arg); });
CMD_VAR_BOOL ("network.rpc.use_xmlrpc", true);
CMD_VAR_BOOL ("network.rpc.use_jsonrpc", true);
CMD_ANY ("network.rpc.use_xmlrpc", [](auto, auto) { return rpc::rpc.use_xmlrpc(); });
CMD_ANY_VALUE_V ("network.rpc.use_xmlrpc.set", [](auto, auto& arg) { return rpc::rpc.set_use_xmlrpc(arg); });
CMD_ANY ("network.rpc.use_jsonrpc", [](auto, auto) { return rpc::rpc.use_jsonrpc(); });
CMD_ANY_VALUE_V ("network.rpc.use_jsonrpc.set", [](auto, auto& arg) { return rpc::rpc.set_use_jsonrpc(arg); });
CMD_ANY ("network.block.ipv4", [nw_config](auto, auto) { return nw_config->is_block_ipv4(); });
CMD_ANY_VALUE_V ("network.block.ipv4.set", [nw_config](auto, auto& value) { return nw_config->set_block_ipv4(value); });
+1 -1
View File
@@ -142,7 +142,7 @@ initialize_command_tracker() {
lt_log_print(torrent::LOG_DHT_ERROR, "dht.port.set is no longer supported, use dht.override_port.set", 0);
});
CMD2_ANY ("dht.override_port", [](auto, auto) { return torrent::runtime::network_config()->override_dht_port(); });
CMD2_ANY_VALUE_V ("dht.override_port.set", [](auto, auto& value) { return torrent::runtime::network_config()->set_override_dht_port(value); });
CMD2_ANY_VALUE_V ("dht.override_port.set", [](auto, auto& value) { return torrent::runtime::network_manager()->set_dht_port(value); });
CMD2_ANY_STRING ("dht.add_node", [](auto, auto& str) { return apply_dht_add_node(str); });
CMD2_ANY ("dht.statistics", [](auto, auto) { return control->dht_manager()->dht_statistics(); });
-3
View File
@@ -70,9 +70,6 @@ Control::initialize() {
display::Window::slot_unschedule([this](display::Window* w) { m_display->unschedule(w); });
display::Window::slot_adjust([this]() { m_display->adjust_layout(); });
torrent::net_thread::http_stack()->set_user_agent(USER_AGENT);
m_core->listen_open();
m_core->set_hashing_view(*m_view_manager->find_throw("hashing"));
m_ui->init(this);
+21 -7
View File
@@ -8,6 +8,7 @@
#include <torrent/object_stream.h>
#include <torrent/rate.h>
#include <torrent/runtime/network_manager.h>
#include <torrent/runtime/runtime.h>
#include <torrent/tracker/dht_controller.h>
#include <torrent/utils/log.h>
@@ -141,18 +142,29 @@ DhtManager::save_dht_cache() {
void
DhtManager::set_mode_by_user(const std::string& arg) {
for (int i = 0; i < dht_settings_num; i++) {
if (arg == dht_settings[i]) {
m_set_by_user = true;
return set_mode_directly(i);
}
unsigned int mode = [arg]() {
for (int i = 0; i < dht_settings_num; i++) {
if (arg == dht_settings[i])
return i;
}
throw torrent::input_error("Invalid dht mode: " + arg);
}();
m_set_by_user = true;
if (!torrent::runtime::is_network_initialized()) {
m_start = mode;
return;
}
set_mode_directly(mode);
}
void
DhtManager::set_mode_directly(unsigned int mode) {
if (mode >= dht_settings_num)
throw torrent::input_error("Invalid argument.");
throw torrent::input_error("Invalid dht mode.");
m_start = mode;
@@ -164,8 +176,10 @@ DhtManager::set_mode_directly(unsigned int mode) {
void
DhtManager::set_auto_if_untouched_and_has_session() {
if (m_set_by_user)
if (m_set_by_user) {
set_mode_directly(m_start);
return;
}
if (rpc::call_command_string("session.path").empty()) {
LT_LOG("DHT auto-start disabled, session path not set.", 0);
+24 -21
View File
@@ -15,6 +15,7 @@
#include <torrent/rate.h>
#include <torrent/data/file_utils.h>
#include <torrent/net/http_stack.h>
#include <torrent/runtime/client_config.h>
#include <torrent/utils/string_manip.h>
#include "control.h"
@@ -104,40 +105,42 @@ DownloadFactory::receive_load() {
throw torrent::internal_error("DownloadFactory::load*() called on an object with m_stream != NULL");
if (is_network_uri(m_uri)) {
// Http handling here.
m_stream.reset(new std::stringstream);
HttpQueue::iterator itr = m_manager->http_queue()->insert(m_uri, m_stream);
auto done_fn = [this]() { receive_loaded(); };
auto failed_fn = [this](const std::string& error) { receive_failed(error); };
itr->add_done_slot(torrent::this_thread::thread(), [this]() { receive_loaded(); });
itr->add_failed_slot(torrent::this_thread::thread(), [this](const std::string& error) { receive_failed(error); });
m_manager->http_queue()->insert(m_uri, m_stream, done_fn, failed_fn);
m_variables["tied_to_file"] = (int64_t)false;
return;
}
} else if (is_magnet_uri(m_uri)) {
if (is_magnet_uri(m_uri)) {
// DEBUG: Use m_object.
m_stream.reset(new std::stringstream());
*m_stream << "d10:magnet-uri" << m_uri.length() << ":" << m_uri << "e";
m_variables["tied_to_file"] = (int64_t)false;
receive_loaded();
} else {
std::fstream stream(expand_path(m_uri).c_str(), std::ios::in | std::ios::binary);
if (!stream.is_open())
return receive_failed("Could not open file");
m_object = new torrent::Object;
stream >> *m_object;
if (!stream.good())
return receive_failed("Reading torrent file failed");
m_isFile = true;
receive_loaded();
return;
}
std::fstream stream(expand_path(m_uri).c_str(), std::ios::in | std::ios::binary);
if (!stream.is_open())
return receive_failed("Could not open file");
m_object = new torrent::Object;
stream >> *m_object;
if (!stream.good())
return receive_failed("Reading torrent file failed");
m_isFile = true;
receive_loaded();
}
void
@@ -269,7 +272,7 @@ DownloadFactory::receive_success() {
if (!m_session && m_variables["tied_to_file"].as_value())
rpc::call_command("d.tied_to_file.set", m_uri.empty() ? m_variables["tied_file"] : m_uri, rpc::make_target(download));
rpc::call_command("d.peer_exchange.set", rpc::call_command_value("protocol.pex"), rpc::make_target(download));
rpc::call_command("d.peer_exchange.set", torrent::runtime::client_config()->is_pex_enabled(), rpc::make_target(download));
torrent::resume_load_addresses(*download->download(), resumeObject);
torrent::resume_load_file_priorities(*download->download(), resumeObject);
+6 -5
View File
@@ -9,7 +9,8 @@
namespace core {
HttpQueue::iterator
HttpQueue::insert(const std::string& url, std::shared_ptr<std::ostream> stream) {
HttpQueue::insert(const std::string& url, std::shared_ptr<std::ostream> stream,
std::function<void()> done_fn, std::function<void(const std::string&)> failed_fn) {
auto itr = base_type::insert(end(), torrent::net::HttpGet(url, stream));
itr->set_max_file_size(15 << 20);
@@ -18,11 +19,11 @@ HttpQueue::insert(const std::string& url, std::shared_ptr<std::ostream> stream)
for (auto& slot : m_signal_insert)
slot(*itr);
itr->add_done_slot(torrent::this_thread::thread(), [this, itr]() { erase(itr); });
itr->add_failed_slot(torrent::this_thread::thread(), [this, itr](auto) { erase(itr); });
itr->add_done_slot(torrent::this_thread::thread(), std::move(done_fn));
itr->add_done_slot(torrent::this_thread::thread(), [this, itr]() { erase(itr); });
// TODO: Downloading http torrents doesn't seem to work.
// TODO: Quitting no longer works.
itr->add_failed_slot(torrent::this_thread::thread(), std::move(failed_fn));
itr->add_failed_slot(torrent::this_thread::thread(), [this, itr](auto) { erase(itr); });
torrent::net_thread::http_stack()->start_get(*itr);
+2 -1
View File
@@ -39,7 +39,8 @@ public:
//
// Consider adding a flag to indicate whetever HttpQueue should
// delete the stream.
iterator insert(const std::string& url, std::shared_ptr<std::ostream> stream);
iterator insert(const std::string& url, std::shared_ptr<std::ostream> stream,
std::function<void()> done_fn, std::function<void(const std::string&)> failed_fn);
void erase(iterator itr);
void clear();
+2 -34
View File
@@ -33,6 +33,8 @@
#include "core/http_queue.h"
#include "core/view.h"
#include <torrent/runtime/client_config.h>
namespace core {
const int Manager::create_start;
@@ -160,40 +162,6 @@ Manager::shutdown(bool force) {
}
}
void
Manager::listen_open() {
// This stuff really should be moved outside of manager, make it
// part of the init script.
if (!rpc::call_command_value("network.port_open"))
return;
int portFirst, portLast;
torrent::Object portRange = rpc::call_command("network.port_range");
if (!portRange.is_string())
throw torrent::input_error("Invalid port_range argument type.");
if (std::sscanf(portRange.as_string().c_str(), "%i-%i", &portFirst, &portLast) != 2)
throw torrent::input_error("Invalid port_range argument.");
if (portFirst > portLast || portLast >= (1 << 16))
throw torrent::input_error("Invalid port range.");
if (rpc::call_command_value("network.port_random")) {
int boundary = portFirst + random() % (portLast - portFirst + 1);
if (torrent::runtime::network_manager()->listen_open(boundary, portLast) ||
torrent::runtime::network_manager()->listen_open(portFirst, boundary))
return;
} else {
if (torrent::runtime::network_manager()->listen_open(portFirst, portLast))
return;
}
throw torrent::input_error("Could not open/bind port for listening: " + std::string(std::strerror(errno)));
}
void
Manager::receive_http_failed(std::string msg) {
push_log_std("Http download error: \"" + msg + "\"");
-2
View File
@@ -55,8 +55,6 @@ public:
void cleanup();
void listen_open();
const std::string& magnet_path();
void set_magnet_path(const std::string& path);
-1
View File
@@ -13,7 +13,6 @@ void
InputEvent::insert() {
torrent::this_thread::poll()->open(this);
torrent::this_thread::poll()->insert_read(this);
torrent::this_thread::poll()->insert_error(this);
}
void
+3 -4
View File
@@ -2,16 +2,15 @@
#define RTORRENT_INPUT_INPUT_EVENT_H
#include <functional>
#include <torrent/event.h>
#include <torrent/system/event.h>
namespace input {
class InputEvent : public torrent::Event {
class InputEvent : public torrent::system::Event {
public:
typedef std::function<void (int)> slot_int;
InputEvent(int fd) { m_fileDesc = fd; }
InputEvent(int fd) { set_file_descriptor(fd); }
const char* type_name() const override { return "input"; }
+26 -8
View File
@@ -11,6 +11,9 @@
#include <torrent/exceptions.h>
#include <torrent/data/chunk_utils.h>
#include <torrent/net/fd.h>
#include <torrent/net/http_stack.h>
#include <torrent/runtime/memory_manager.h>
#include <torrent/runtime/runtime.h>
#include <torrent/utils/chrono.h>
#include <torrent/utils/log.h>
@@ -149,8 +152,8 @@ main(int argc, char** argv) {
SignalHandler::set_sigaction_handler(SIGBUS, &handle_sigbus);
torrent::log_add_group_output(torrent::LOG_NOTICE, "important");
torrent::log_add_group_output(torrent::LOG_DHT_ERROR, "important");
torrent::log_add_group_output(torrent::LOG_NOTICE, "important");
torrent::log_add_group_output(torrent::LOG_DHT_ERROR, "important");
torrent::log_add_group_output(torrent::LOG_INFO, "complete");
torrent::log_add_group_output(torrent::LOG_DHT_ERROR, "complete");
@@ -298,8 +301,6 @@ main(int argc, char** argv) {
"schedule = low_diskspace,5,60,((close_low_diskspace,500M))\n"
"schedule = prune_file_status,3600,86400,((system.file_status_cache.prune))\n"
"protocol.encryption.set=allow_incoming,prefer_plaintext,enable_retry\n"
"ui.color.focus.set=reverse\n"
);
@@ -343,8 +344,9 @@ main(int argc, char** argv) {
CMD_REDIRECT("directory", "directory.default.set");
CMD_REDIRECT("session", "session.path.set");
CMD_REDIRECT("scgi_port", "network.scgi.open_port");
CMD_REDIRECT("scgi_local", "network.scgi.open_local");
CMD_REDIRECT_NO_EXPORT("port_range", "network.listen.port.range.set");
CMD_REDIRECT_NO_EXPORT("scgi_port", "network.scgi.open_port");
CMD_REDIRECT_NO_EXPORT("scgi_local", "network.scgi.open_local");
CMD_REDIRECT("to_gm_time", "convert.gm_time");
CMD_REDIRECT("to_gm_date", "convert.gm_date");
@@ -388,8 +390,6 @@ main(int argc, char** argv) {
CMD_REDIRECT("schedule2", "schedule");
CMD_REDIRECT("schedule_remove2", "schedule.remove");
// TODO: Remove file.append when cleaning these up.
CMD_REDIRECT("bind", "network.bind_address.set");
CMD_REDIRECT("ip", "network.local_address.set");
CMD_REDIRECT("port_range", "network.port_range.set");
@@ -417,6 +417,20 @@ main(int argc, char** argv) {
lt_log_print(torrent::LOG_WARN, "The 'throttle.ip' command is deprecated and does nothing.");
return torrent::Object();
});
CMD_ANY("network.port_open", [](auto, auto) {
lt_log_print(torrent::LOG_WARN, "The 'network.port_open' command is deprecated and does nothing.");
return torrent::Object();
});
CMD_ANY("network.port_open.set", [](auto, auto) {
lt_log_print(torrent::LOG_WARN, "The 'network.port_open.set' command is deprecated and does nothing.");
return torrent::Object();
});
CMD_REDIRECT("network.port_random", "network.listen.port.random");
CMD_REDIRECT("network.port_random.set", "network.listen.port.random.set");
CMD_REDIRECT("network.port_range", "network.listen.port.range");
CMD_REDIRECT("network.port_range.set", "network.listen.port.range.set");
}
{
@@ -444,10 +458,14 @@ main(int argc, char** argv) {
});
LT_LOG("seeded srandom and srand48 (seed:%u)", random_seed);
LT_LOG("max memory usage: %" PRIu64, torrent::runtime::memory_manager()->max_memory_usage());
control->initialize();
control->ui()->load_input_history();
torrent::net_thread::http_stack()->set_user_agent(USER_AGENT);
torrent::runtime::initialize_network();
// Load session torrents and perform scheduled tasks to ensure session torrents are loaded
// before arg torrents.
control->dht_manager()->set_auto_if_untouched_and_has_session();
+15 -16
View File
@@ -55,12 +55,11 @@ struct rt_triple : private std::pair<T1, T2> {
base_type(src.first, src.second), third(src.third) {}
};
typedef rt_triple<int, void*, void*> target_type;
class command_base;
typedef const torrent::Object (*command_base_call_type)(command_base*, target_type, const torrent::Object&);
typedef std::function<torrent::Object (target_type, const torrent::Object&)> base_function;
using target_type = rt_triple<int, void*, void*>;
using base_function = std::function<torrent::Object (target_type, const torrent::Object&)>;
using command_base_call_type = const torrent::Object (command_base*, target_type, const torrent::Object&);
template <typename tmpl> struct command_base_is_valid {};
template <command_base_call_type tmpl_func> struct command_base_is_type {};
@@ -84,17 +83,18 @@ public:
typedef const torrent::Object (*download_pair_slot) (command_base*, core::Download*, core::Download*, const torrent::Object&);
static const int target_generic = 0;
static const int target_any = 1;
static const int target_download = 2;
static const int target_peer = 3;
static const int target_tracker = 4;
static const int target_file = 5;
static const int target_file_itr = 6;
static constexpr int target_generic = 0;
static constexpr int target_any = 1;
static constexpr int target_download = 2;
static constexpr int target_peer = 3;
static constexpr int target_tracker = 4;
static constexpr int target_file = 5;
static constexpr int target_file_itr = 6;
static constexpr int target_download_pair = 7;
static const int target_download_pair = 7;
static constexpr unsigned int max_arguments = 10;
static const unsigned int max_arguments = 10;
static constexpr std::size_t optimal_alignment = std::max(alignof(std::max_align_t), alignof(base_function));
struct stack_type {
torrent::Object* begin() { return reinterpret_cast<torrent::Object*>(buffer); }
@@ -169,8 +169,7 @@ public:
template <typename T>
void set_function(T s, [[maybe_unused]] int value = command_base_is_valid<T>::value) {
static_assert(sizeof(T) <= sizeof(t_pod), "t_pod storage overflow");
static_assert(alignof(std::max_align_t) % alignof(T) == 0, "t_pod alignment insufficient for type");
static_assert(alignof(std::max_align_t) >= alignof(T), "t_pod structural capacity mismatch");
static_assert(optimal_alignment >= alignof(T), "t_pod alignment insufficient for type");
if (m_dest_helper)
m_dest_helper(t_pod);
@@ -200,7 +199,7 @@ protected:
using copy_fn_t = void (*)(void* dest, const void* src);
using dest_fn_t = void (*)(void* ptr);
alignas(std::max_align_t) char t_pod[sizeof(base_function)];
alignas(optimal_alignment) char t_pod[sizeof(base_function)];
copy_fn_t m_copy_helper;
dest_fn_t m_dest_helper;
+46 -202
View File
@@ -1,20 +1,22 @@
#include "config.h"
#include <cerrno>
#include <cstring>
#include <fcntl.h>
#include <spawn.h>
#include <string>
#include "rpc/exec_file.h"
// #include <cassert>
// #include <cerrno>
// #include <cstring>
// #include <fcntl.h>
// #include <spawn.h>
// #include <string>
#include <unistd.h>
#include <sys/types.h>
#include <sys/wait.h>
// #include <sys/types.h>
#include <sys/uio.h>
// #include <torrent/net/fd.h>
#include <torrent/system/thread.h>
#include <torrent/system/spawn_process.h>
#include <torrent/system/types.h>
#include "exec_file.h"
#include "parse.h"
// Standard POSIX environment pointer
extern char** environ;
#include "rpc/parse.h"
namespace rpc {
@@ -22,214 +24,56 @@ namespace rpc {
int
ExecFile::execute(const char* file, char* const* argv, int flags) {
// Write the executed command and its parameters to the log fd.
[[maybe_unused]] int result;
torrent::system::SpawnProcess spawn_process;
spawn_process.set_log_fd(m_log_fd);
spawn_process.set_background(flags & flag_background);
spawn_process.set_capture_output(flags & flag_capture);
if (m_log_fd != -1) {
for (char* const* itr = argv; *itr != NULL; itr++) {
if (itr == argv)
result = write(m_log_fd, "\n---\n", sizeof("\n---\n"));
else
result = write(m_log_fd, " ", 1);
std::vector<struct iovec> iovecs;
iovecs.reserve(32);
iovecs.push_back({const_cast<char*>("\n---\n"), 5});
result = write(m_log_fd, *itr, std::strlen(*itr));
for (auto* itr = argv; *itr != nullptr; itr++) {
if (itr != argv)
iovecs.push_back({const_cast<char*>(" "), 1});
iovecs.push_back({*itr, std::strlen(*itr)});
}
result = write(m_log_fd, "\n---\n", sizeof("\n---\n"));
iovecs.push_back({const_cast<char*>("\n---\n"), 5});
[[maybe_unused]] int result = ::writev(m_log_fd, iovecs.data(), iovecs.size());
}
posix_spawn_file_actions_t actions{};
if (posix_spawn_file_actions_init(&actions) != 0)
throw torrent::internal_error("ExecFile::execute(...) posix_spawn_file_actions_init failed.");
posix_spawnattr_t attr;
posix_spawnattr_init(&attr);
short spawn_flags = 0;
// Try to avoid leaking open fds to the spawned process. Prefer POSIX_SPAWN_CLOEXEC_DEFAULT
// (macOS-only) or posix_spawn_file_actions_addclosefrom_np (glibc >= 2.34, FreeBSD >= 13.1).
//
// Other platforms like musl libc, OpenBSD and NetBSD must rely on explicit O_CLOEXEC.
#if defined(POSIX_SPAWN_CLOEXEC_DEFAULT)
spawn_flags |= POSIX_SPAWN_CLOEXEC_DEFAULT;
#elif defined(HAVE_POSIX_SPAWN_FILE_ACTIONS_ADDCLOSEFROM_NP)
posix_spawn_file_actions_addclosefrom_np(&actions, 3);
#endif
// Handle standard input redirection (/dev/null), posix_spawn_file_actions_addopen handles opening
// and dup2 natively
if (posix_spawn_file_actions_addopen(&actions, 0, "/dev/null", O_RDWR, 0) != 0) {
// Fallback if open fails inside action setup
posix_spawn_file_actions_addclose(&actions, 0);
}
int pipe_fd[2] = {-1, -1};
if ((flags & flag_capture) && pipe(pipe_fd))
throw torrent::input_error("ExecFile::execute(...) Pipe creation failed.");
// Handle standard output redirection
if (flags & flag_capture) {
posix_spawn_file_actions_adddup2(&actions, pipe_fd[1], 1);
// Ensure the write end of the pipe is closed in the child after duplicating.
posix_spawn_file_actions_addclose(&actions, pipe_fd[1]);
posix_spawn_file_actions_addclose(&actions, pipe_fd[0]);
} else if (m_log_fd != -1) {
posix_spawn_file_actions_adddup2(&actions, m_log_fd, 1);
} else {
posix_spawn_file_actions_addopen(&actions, 1, "/dev/null", O_RDWR, 0);
}
if (m_log_fd != -1) {
posix_spawn_file_actions_adddup2(&actions, m_log_fd, 2);
} else {
posix_spawn_file_actions_addopen(&actions, 2, "/dev/null", O_RDWR, 0);
}
if (flags & flag_background) {
#ifdef POSIX_SPAWN_SETSID
spawn_flags |= POSIX_SPAWN_SETSID;
#else
spawn_flags |= POSIX_SPAWN_SETPGROUP;
posix_spawnattr_setpgroup(&attr, 0);
#endif
}
posix_spawnattr_setflags(&attr, spawn_flags);
pid_t child_pid{};
int spawn_status = posix_spawnp(&child_pid, file, &actions, &attr, argv, environ);
posix_spawn_file_actions_destroy(&actions);
posix_spawnattr_destroy(&attr);
int spawn_status = spawn_process.execute(file, argv);
if (spawn_status != 0) {
if (pipe_fd[0] != -1)
::close(pipe_fd[0]);
if (pipe_fd[1] != -1)
::close(pipe_fd[1]);
throw torrent::input_error("ExecFile::execute(...) posix_spawn failed: " + std::string(std::strerror(spawn_status)));
}
if (flags & flag_capture) {
m_capture = std::string();
::close(pipe_fd[1]);
char buffer[4096];
ssize_t length;
do {
length = read(pipe_fd[0], buffer, sizeof(buffer));
if (length > 0)
m_capture += std::string(buffer, length);
} while (length > 0);
::close(pipe_fd[0]);
if (m_log_fd != -1) {
result = write(m_log_fd, "Captured output:\n", sizeof("Captured output:\n"));
result = write(m_log_fd, m_capture.data(), m_capture.length());
auto prefix = "\n--- posix_spawn failed: ";
auto errno_str = torrent::system::errno_enum_str(spawn_status) + " ---\n";
struct iovec iovecs[2] = {
{const_cast<char*>(prefix), std::strlen(prefix)},
{const_cast<char*>(errno_str.c_str()), errno_str.size()},
};
[[maybe_unused]] int result = ::writev(m_log_fd, iovecs, 2);
}
throw torrent::input_error("ExecFile::execute() posix_spawn failed: " + torrent::system::errno_enum_str(spawn_status));
}
if (flags & flag_background) {
if (m_log_fd != -1)
result = write(m_log_fd, "\n--- Running in Background ---\n", sizeof("\n--- Running in Background ---\n"));
m_waitpid_queue.close_pid(spawn_process.child_pid());
return 0;
}
int status;
while (waitpid(child_pid, &status, 0) == -1) {
switch (errno) {
case EINTR:
continue;
case ECHILD:
throw torrent::internal_error("ExecFile::execute(...) waitpid failed with ECHILD, child process not found.");
case EINVAL:
throw torrent::internal_error("ExecFile::execute(...) waitpid failed with EINVAL.");
default:
throw torrent::internal_error("ExecFile::execute(...) waitpid failed with unexpected error: " + std::string(std::strerror(errno)));
}
};
// Check return value?
if (m_log_fd != -1) {
if (WIFEXITED(status) && WEXITSTATUS(status) == 0)
result = write(m_log_fd, "\n--- Success ---\n", sizeof("\n--- Success ---\n"));
else
result = write(m_log_fd, "\n--- Error ---\n", sizeof("\n--- Error ---\n"));
}
return status;
}
torrent::Object
ExecFile::execute_object(const torrent::Object& rawArgs, int flags) {
char* argsBuffer[max_args];
char** argsCurrent = argsBuffer;
// Size of value strings are less than 24.
char valueBuffer[buffer_size+1];
char* valueCurrent = valueBuffer;
if (rawArgs.is_list()) {
const torrent::Object::list_type& args = rawArgs.as_list();
if (args.empty())
throw torrent::input_error("Too few arguments.");
for (torrent::Object::list_const_iterator itr = args.begin(), last = args.end(); itr != last; itr++, argsCurrent++) {
if (argsCurrent == argsBuffer + max_args - 1)
throw torrent::input_error("Too many arguments.");
if (itr->is_string() && (!(flags & flag_expand_tilde) || *itr->as_string().c_str() != '~')) {
*argsCurrent = const_cast<char*>(itr->as_string().c_str());
} else {
*argsCurrent = valueCurrent;
valueCurrent = print_object(valueCurrent, valueBuffer + buffer_size, &*itr, flags) + 1;
if (valueCurrent >= valueBuffer + buffer_size)
throw torrent::input_error("Overflowed execute arg buffer.");
}
}
} else {
const torrent::Object::string_type& args = rawArgs.as_string();
if ((flags & flag_expand_tilde) && args.c_str()[0] == '~') {
*argsCurrent = valueCurrent;
valueCurrent = print_object(valueCurrent, valueBuffer + buffer_size, &rawArgs, flags) + 1;
} else {
*argsCurrent = const_cast<char*>(args.c_str());
}
argsCurrent++;
}
*argsCurrent = NULL;
int status = execute(argsBuffer[0], argsBuffer, flags);
if ((flags & flag_throw) && status != 0)
throw torrent::input_error("Bad return code.");
if (flags & flag_capture)
return m_capture;
m_capture = spawn_process.capture_child_output();
return torrent::Object((int64_t)status);
return spawn_process.wait_for_child();
}
}
} // namespace rpc
+4
View File
@@ -3,6 +3,8 @@
#include <torrent/object.h>
#include "utils/waitpid_queue.h"
namespace rpc {
class ExecFile {
@@ -24,6 +26,8 @@ public:
private:
int m_log_fd{-1};
std::string m_capture;
utils::WaitpidQueue m_waitpid_queue;
};
}
+6 -37
View File
@@ -113,25 +113,20 @@ bool
RpcManager::process(RPCType type, const char* in_buffer, uint32_t length, slot_response_callback callback) {
switch (type) {
case RPCType::XML:
// TODO: 'network.rpc.use_xmlrpc' should be a bool in RpcManager, not a command variable.
if (m_xmlrpc.is_valid() && rpc::call_command_value("network.rpc.use_xmlrpc")) {
return m_xmlrpc.process(in_buffer, length, callback);
} else {
if (!m_xmlrpc.is_valid() || !m_use_xmlrpc) {
const std::string response = "<?xml version=\"1.0\"?><methodResponse><fault><value><struct><member><name>faultCode</name><value><i8>-501</i8></value></member><member><name>faultString</name><value><string>XML-RPC not supported</string></value></member></struct></value></fault></methodResponse>";
return callback(response.c_str(), response.size());
}
break;
return m_xmlrpc.process(in_buffer, length, callback);
case RPCType::JSON:
if (rpc::call_command_value("network.rpc.use_jsonrpc")) {
return m_jsonrpc.process(in_buffer, length, callback);
} else {
if (!m_use_jsonrpc) {
const std::string response = "{\"jsonrpc\":\"2.0\",\"error\":{\"code\":-32601,\"message\":\"JSON-RPC not supported\"},\"id\":null}";
return callback(response.c_str(), response.size());
}
break;
return m_jsonrpc.process(in_buffer, length, callback);
default:
throw torrent::input_error("invalid parameters: unknown RPC type");
@@ -172,32 +167,6 @@ RpcManager::cleanup() {
m_jsonrpc.cleanup();
}
bool
RpcManager::is_type_enabled(RPCType type) const {
switch (type) {
case RPCType::XML:
return m_is_xmlrpc_enabled;
case RPCType::JSON:
return m_is_jsonrpc_enabled;
default:
throw torrent::input_error("invalid parameters: unknown RPC type");
}
}
void
RpcManager::set_type_enabled(RPCType type, bool enabled) {
switch (type) {
case RPCType::XML:
m_is_xmlrpc_enabled = enabled;
break;
case RPCType::JSON:
m_is_jsonrpc_enabled = enabled;
break;
default:
throw torrent::input_error("invalid parameters: unknown RPC type");
}
}
void
RpcManager::insert_command(const char* name, const char* parm, const char* doc) {
m_xmlrpc.insert_command(name, parm, doc);
+13 -4
View File
@@ -67,8 +67,11 @@ public:
int dialect() { return m_xmlrpc.dialect(); }
void set_dialect(int dialect) { m_xmlrpc.set_dialect(dialect); }
bool is_type_enabled(RPCType type) const;
void set_type_enabled(RPCType type, bool enabled);
bool use_xmlrpc() const;
void set_use_xmlrpc(bool v);
bool use_jsonrpc() const;
void set_use_jsonrpc(bool v);
bool process(RPCType type, const char* in_buffer, uint32_t length, slot_response_callback callback);
bool process_untrusted(RPCType type, const char* in_buffer, uint32_t length, slot_response_callback callback);
@@ -102,8 +105,9 @@ private:
JsonRpc m_jsonrpc;
bool m_handlers_initialized{};
bool m_is_jsonrpc_enabled{true};
bool m_is_xmlrpc_enabled{true};
bool m_use_jsonrpc{true};
bool m_use_xmlrpc{true};
std::atomic<bool> m_scgi_allow_compression{true};
std::atomic<unsigned int> m_scgi_min_compress_size{1000};
@@ -116,6 +120,11 @@ private:
extern RpcManager rpc;
inline bool RpcManager::use_xmlrpc() const { return m_use_xmlrpc; }
inline void RpcManager::set_use_xmlrpc(bool v) { m_use_xmlrpc = v; }
inline bool RpcManager::use_jsonrpc() const { return m_use_jsonrpc; }
inline void RpcManager::set_use_jsonrpc(bool v) { m_use_jsonrpc = v; }
} // namespace rpc
#endif
-1
View File
@@ -111,7 +111,6 @@ SCgi::activate() {
torrent::this_thread::poll()->open(this);
torrent::this_thread::poll()->insert_read(this);
torrent::this_thread::poll()->insert_error(this);
}
// TODO: This should close the fd to avoid reuse.
+2 -2
View File
@@ -3,13 +3,13 @@
#include <array>
#include <memory>
#include <torrent/event.h>
#include <torrent/system/event.h>
#include "rpc/scgi_task.h"
namespace rpc {
class SCgi : public torrent::Event {
class SCgi : public torrent::system::Event {
public:
static const int max_tasks = 100;
+5 -6
View File
@@ -27,7 +27,7 @@ namespace rpc {
SCgiTask::SCgiTask()
: m_callback_id(torrent::system::make_callback_id()) {
set_file_descriptor(-1);
reset_file_descriptor();
}
void
@@ -49,7 +49,6 @@ SCgiTask::open(SCgi* parent, int fd) {
torrent::this_thread::poll()->open(this);
torrent::this_thread::poll()->insert_read(this);
torrent::this_thread::poll()->insert_error(this);
auto lock = std::lock_guard<std::mutex>(m_result_mutex);
@@ -65,7 +64,7 @@ SCgiTask::cancel_open() {
torrent::this_thread::poll()->remove_and_close(this);
torrent::fd_close(file_descriptor());
set_file_descriptor(-1);
reset_file_descriptor();
};
void
@@ -79,7 +78,7 @@ SCgiTask::close() {
torrent::this_thread::poll()->remove_and_close(this);
torrent::fd_close(file_descriptor());
set_file_descriptor(-1);
reset_file_descriptor();
});
// The callbacks are guaranteed to be finished/canceled at this point.
@@ -98,7 +97,7 @@ SCgiTask::event_read() {
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);
int bytes = ::recv(file_descriptor(), m_buffer.data() + m_position, read_length, 0);
if (bytes <= 0) {
if (bytes == 0 || !(errno == EAGAIN || errno == EINTR))
@@ -181,7 +180,7 @@ event_read_failed:
void
SCgiTask::event_write() {
int bytes = ::send(m_fileDesc, m_buffer.data() + m_position, m_buffer.size() - m_position, 0);
int bytes = ::send(file_descriptor(), m_buffer.data() + m_position, m_buffer.size() - m_position, 0);
if (bytes == -1) {
if (!(errno == EAGAIN || errno == EINTR))
+4 -4
View File
@@ -4,13 +4,13 @@
#include <memory>
#include <mutex>
#include <vector>
#include <torrent/event.h>
#include <torrent/system/event.h>
namespace rpc {
class SCgi;
class SCgiTask : public torrent::Event {
class SCgiTask : public torrent::system::Event {
public:
static constexpr int default_buffer_size = 8191;
static constexpr int max_header_size = 2000;
@@ -22,8 +22,8 @@ public:
const char* type_name() const override { return "scgi-task"; }
bool is_open() const { return m_fileDesc != -1; }
bool is_available() const { return m_fileDesc == -1; }
bool is_open() const { return file_descriptor() != -1; }
bool is_available() const { return file_descriptor() == -1; }
void open(SCgi* parent, int fd);
void cancel_open();
+1 -1
View File
@@ -38,7 +38,7 @@ parse_main_options(int argc, char** argv) {
optionParser.insert_option('b', [](auto& arg) { rpc::call_command_set_string("network.bind_address.set", arg); });
optionParser.insert_option('d', [](auto& arg) { rpc::call_command_set_string("directory.default.set", arg); });
optionParser.insert_option('i', [](auto& arg) { rpc::call_command_set_string("ip", arg); });
optionParser.insert_option('p', [](auto& arg) { rpc::call_command_set_string("network.port_range.set", arg); });
optionParser.insert_option('p', [](auto& arg) { rpc::call_command_set_string("network.listen.port.range.set", arg); });
optionParser.insert_option('s', [](auto& arg) { rpc::call_command_set_string("session", arg); });
optionParser.insert_option('O', [](auto& arg) { rpc::parse_command_single_std(arg); });
+4 -7
View File
@@ -3,7 +3,6 @@
#include <cassert>
#include <torrent/exceptions.h>
#include <torrent/chunk_manager.h>
#include <torrent/throttle.h>
#include <torrent/torrent.h>
#include <torrent/data/file_list.h>
@@ -31,10 +30,10 @@
namespace ui {
Download::Download(core::Download* d) :
m_download(d) {
Download::Download(core::Download* d)
: m_download(d) {
m_windowDownloadStatus = new WDownloadStatus(d);
m_windowDownloadStatus = std::make_unique<WDownloadStatus>(d);
m_windowDownloadStatus->set_bottom(true);
m_uiArray[DISPLAY_MENU] = create_menu();
@@ -60,8 +59,6 @@ Download::~Download() {
assert(!is_active() && "ui::Download::~Download() called on an active object.");
std::for_each(m_uiArray, m_uiArray + DISPLAY_MAX_SIZE, [](ElementBase* eb) { delete eb; });
delete m_windowDownloadStatus;
}
inline ElementBase*
@@ -171,7 +168,7 @@ Download::activate(display::Frame* frame, [[maybe_unused]] bool focus) {
m_frame = frame;
m_frame->initialize_row(2);
m_frame->frame(1)->initialize_window(m_windowDownloadStatus);
m_frame->frame(1)->initialize_window(m_windowDownloadStatus.get());
m_windowDownloadStatus->set_active(true);
activate_display_menu(DISPLAY_PEER_LIST);
+1 -35
View File
@@ -1,37 +1,3 @@
// rTorrent - BitTorrent client
// 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 <sundell.software@gmail.com>
#ifndef RTORRENT_UI_DOWNLOAD_H
#define RTORRENT_UI_DOWNLOAD_H
@@ -113,7 +79,7 @@ private:
bool m_focusDisplay{};
WDownloadStatus* m_windowDownloadStatus;
std::unique_ptr<WDownloadStatus> m_windowDownloadStatus;
};
}
+6 -4
View File
@@ -15,13 +15,13 @@ namespace ui {
class ElementLogComplete : public ElementBase {
public:
typedef display::WindowLogComplete WLogComplete;
using WLogComplete = display::WindowLogComplete;
ElementLogComplete(torrent::log_buffer* l);
~ElementLogComplete() override;
void activate(display::Frame* frame, bool focus = true);
void disable();
void activate(display::Frame* frame, bool focus = true) override;
void disable() override;
display::Window* window();
@@ -32,7 +32,9 @@ private:
torrent::log_buffer* m_log;
align_cacheline std::atomic<bool> m_log_updating{};
align_cacheline
std::atomic<bool> m_log_updating{};
};
}
+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
+2
View File
@@ -49,6 +49,8 @@ rtorrent_Test_Rpc_SOURCES = $(rtorrent_Test_Common) \
rtorrent_Test_Src_SOURCES = $(rtorrent_Test_Common) \
src/test_command_dynamic.cc \
src/test_command_dynamic.h \
src/test_command_local.cc \
src/test_command_local.h \
src/test_watch_ready_queue.cc \
src/test_watch_ready_queue.h
+22 -8
View File
@@ -2,9 +2,10 @@
#define LIBTORRENT_HELPERS_MOCK_COMPARE_H
#include <algorithm>
#include <map>
#include <type_traits>
#include <torrent/event.h>
#include <torrent/net/socket_address.h>
#include <torrent/system/event.h>
// Compare arguments to mock functions with what is expected. The lhs
// are the expected arguments, rhs are the ones called with.
@@ -13,13 +14,13 @@ template <typename Arg>
inline bool mock_compare_arg(Arg lhs, Arg rhs) { return lhs == rhs; }
template <int I, typename A, typename... Args>
typename std::enable_if<I == 1, int>::type
std::enable_if_t<I == 1, int>
mock_compare_tuple(const std::tuple<A, Args...>& lhs, const std::tuple<Args...>& rhs) {
return mock_compare_arg(std::get<I>(lhs), std::get<I - 1>(rhs)) ? 0 : 1;
}
template <int I, typename A, typename... Args>
typename std::enable_if<1 < I, int>::type
std::enable_if_t<1 < I, int>
mock_compare_tuple(const std::tuple<A, Args...>& lhs, const std::tuple<Args...>& rhs) {
auto res = mock_compare_tuple<I - 1>(lhs, rhs);
@@ -70,24 +71,37 @@ void mock_compare_add(T* v) {
// Specialize:
//
constexpr int mock_compare_gt_two_int = -0xFD30;
template <>
inline bool mock_compare_arg<int>(int lhs, int rhs) {
if (lhs == mock_compare_gt_two_int)
return rhs > 2;
if (rhs == mock_compare_gt_two_int)
return lhs > 2;
return lhs == rhs;
}
template <>
inline bool mock_compare_arg<sockaddr*>(sockaddr* lhs, sockaddr* rhs) {
return lhs != nullptr && rhs != nullptr && torrent::sa_equal(lhs, rhs);
}
template <>
inline bool mock_compare_arg<const sockaddr*>(const sockaddr* lhs, const sockaddr* rhs) {
return lhs != nullptr && rhs != nullptr && torrent::sa_equal(lhs, rhs);
}
template <>
inline bool mock_compare_arg<torrent::Event*>(torrent::Event* lhs, torrent::Event* rhs) {
if (mock_compare_map<torrent::Event>::is_key(lhs)) {
if (!mock_compare_map<torrent::Event>::has_value(rhs)) {
mock_compare_map<torrent::Event>::values[lhs] = rhs;
inline bool mock_compare_arg<torrent::system::Event*>(torrent::system::Event* lhs, torrent::system::Event* rhs) {
if (mock_compare_map<torrent::system::Event>::is_key(lhs)) {
if (!mock_compare_map<torrent::system::Event>::has_value(rhs)) {
mock_compare_map<torrent::system::Event>::values[lhs] = rhs;
return true;
}
return mock_compare_map<torrent::Event>::has_key(lhs) && mock_compare_map<torrent::Event>::get(lhs) == rhs;
return mock_compare_map<torrent::system::Event>::has_key(lhs) && mock_compare_map<torrent::system::Event>::get(lhs) == rhs;
}
return lhs == rhs;
+5 -66
View File
@@ -6,9 +6,9 @@
#include <iostream>
#include <unistd.h>
#include "torrent/event.h"
#include "torrent/net/socket_address.h"
#include "torrent/net/fd.h"
#include "torrent/system/event.h"
#include "torrent/utils/log.h"
#include "torrent/utils/random.h"
@@ -32,7 +32,7 @@ mock_clear(bool ignore_assert) {
MOCK_CLEANUP_MAP(torrent::random_uniform_uint16);
MOCK_CLEANUP_MAP(torrent::random_uniform_uint32);
mock_compare_map<torrent::Event>::values.clear();
mock_compare_map<torrent::system::Event>::values.clear();
}
} // namespace
@@ -55,7 +55,9 @@ mock_redirect_defaults([[maybe_unused]] mock_redirect_flags flags) {
mock_redirect(torrent::fd__close, std::function<int(int fildes)>([](int fildes) { return ::close(fildes); }));
mock_redirect(torrent::fd__fcntl_int, std::function<int(int fildes, int cmd, int arg)>([](int fildes, int cmd, int arg) { return ::fcntl(fildes, cmd, arg); }));
mock_redirect(torrent::fd__setsockopt_int, std::function<int(int socket, int level, int option_name, int option_value)>([](int socket, int level, int option_name, int option_value) { return ::setsockopt(socket, level, option_name, &option_value, sizeof(int)); }));
mock_redirect(torrent::fd__setsockopt_int, std::function<int(int socket, int level, int option_name, int option_value)>([](int socket, int level, int option_name, int option_value) {
return ::setsockopt(socket, level, option_name, &option_value, sizeof(int));
}));
mock_redirect(torrent::fd__socket, std::function<int(int domain, int type, int protocol)>([](int domain, int type, int protocol) { return ::socket(domain, type, protocol); }));
}
@@ -112,69 +114,6 @@ int fd__socket(int domain, int type, int protocol) {
return mock_call<int>(__func__, &torrent::fd__socket, domain, type, protocol);
}
//
// Mock functions for 'torrent/common.h':
//
namespace this_thread {
void event_open(Event* event) {
MOCK_LOG("fd:%i type_name:%s", event->file_descriptor(), event->type_name());
return mock_call<void>(__func__, &torrent::this_thread::event_open, event);
}
void event_open_and_count(Event* event) {
MOCK_LOG("fd:%i type_name:%s", event->file_descriptor(), event->type_name());
return mock_call<void>(__func__, &torrent::this_thread::event_open_and_count, event);
}
void event_close_and_count(Event* event) {
MOCK_LOG("fd:%i type_name:%s", event->file_descriptor(), event->type_name());
return mock_call<void>(__func__, &torrent::this_thread::event_close_and_count, event);
}
void event_closed_and_count(Event* event) {
MOCK_LOG("fd:%i type_name:%s", event->file_descriptor(), event->type_name());
return mock_call<void>(__func__, &torrent::this_thread::event_closed_and_count, event);
}
void event_insert_read(Event* event) {
MOCK_LOG("fd:%i type_name:%s", event->file_descriptor(), event->type_name());
return mock_call<void>(__func__, &torrent::this_thread::event_insert_read, event);
}
void event_insert_write(Event* event) {
MOCK_LOG("fd:%i type_name:%s", event->file_descriptor(), event->type_name());
return mock_call<void>(__func__, &torrent::this_thread::event_insert_write, event);
}
void event_insert_error(Event* event) {
MOCK_LOG("fd:%i type_name:%s", event->file_descriptor(), event->type_name());
return mock_call<void>(__func__, &torrent::this_thread::event_insert_error, event);
}
void event_remove_read(Event* event) {
MOCK_LOG("fd:%i type_name:%s", event->file_descriptor(), event->type_name());
return mock_call<void>(__func__, &torrent::this_thread::event_remove_read, event);
}
void event_remove_write(Event* event) {
MOCK_LOG("fd:%i type_name:%s", event->file_descriptor(), event->type_name());
return mock_call<void>(__func__, &torrent::this_thread::event_remove_write, event);
}
void event_remove_error(Event* event) {
MOCK_LOG("fd:%i type_name:%s", event->file_descriptor(), event->type_name());
return mock_call<void>(__func__, &torrent::this_thread::event_remove_error, event);
}
void event_remove_and_close(Event* event) {
MOCK_LOG("fd:%i type_name:%s", event->file_descriptor(), event->type_name());
return mock_call<void>(__func__, &torrent::this_thread::event_remove_and_close, event);
}
}
//
// Mock functions for 'torrent/utils/random.h':
//
-17
View File
@@ -169,20 +169,3 @@ TestParseOptions::test_flag_libtorrent() {
FLAG_LT_LOG_ASSERT_ERROR("resume_data|rpc_dump");
}
#define FLAGS_LT_ENCRYPTION_ASSERT(flags, result) \
CPPUNIT_ASSERT(rpc::parse_option_flags(flags, std::bind(&torrent::option_find_string_str, torrent::OPTION_ENCRYPTION, std::placeholders::_1)) == (result))
#define FLAGS_LT_ENCRYPTION_ASSERT_ERROR(flags) \
ASSERT_CATCH_INPUT_ERROR(rpc::parse_option_flags(flags, std::bind(&torrent::option_find_string_str, torrent::OPTION_ENCRYPTION, std::placeholders::_1)))
void
TestParseOptions::test_flags_libtorrent() {
FLAGS_LT_ENCRYPTION_ASSERT("", torrent::runtime::NetworkConfig::encryption_none);
FLAGS_LT_ENCRYPTION_ASSERT("none", torrent::runtime::NetworkConfig::encryption_none);
FLAGS_LT_ENCRYPTION_ASSERT("require_rc4", torrent::runtime::NetworkConfig::encryption_require_RC4);
FLAGS_LT_ENCRYPTION_ASSERT("require_RC4", torrent::runtime::NetworkConfig::encryption_require_RC4);
FLAGS_LT_ENCRYPTION_ASSERT("require_RC4 | enable_retry", torrent::runtime::NetworkConfig::encryption_require_RC4 | torrent::runtime::NetworkConfig::encryption_enable_retry);
FLAGS_LT_ENCRYPTION_ASSERT_ERROR("require_");
}
-2
View File
@@ -14,7 +14,6 @@ class TestParseOptions : public test_fixture {
CPPUNIT_TEST(test_flags_print_flags);
CPPUNIT_TEST(test_flag_libtorrent);
CPPUNIT_TEST(test_flags_libtorrent);
CPPUNIT_TEST_SUITE_END();
@@ -30,5 +29,4 @@ public:
void test_flags_print_flags();
void test_flag_libtorrent();
void test_flags_libtorrent();
};
+54
View File
@@ -0,0 +1,54 @@
#include "config.h"
#include "test/src/test_command_local.h"
#include <torrent/torrent.h>
#include <torrent/runtime/socket_manager.h>
#include <torrent/utils/option_strings.h>
#include "control.h"
#include "globals.h"
#include "rpc/parse_commands.h"
CPPUNIT_TEST_SUITE_REGISTRATION(TestCommandLocal);
void initialize_command_local();
void
TestCommandLocal::setUp() {
torrent::initialize_main_thread();
torrent::initialize();
if (control == nullptr)
control = new Control;
if (!rpc::commands.has("system.sockets.size"))
initialize_command_local();
}
void
TestCommandLocal::tearDown() {
torrent::cleanup();
}
void
TestCommandLocal::test_socket_category_commands() {
for (uint32_t i = 0; i < torrent::runtime::SocketManager::category_count; ++i) {
auto category = static_cast<torrent::runtime::socket_manager_category_t>(i);
auto name = "system.sockets." + torrent::option_to_str_or_throw(torrent::OPTION_SOCKET_CATEGORY, i);
for (const auto suffix : {".size", ".max_size", ".min_alloc", ".max_alloc"})
if (rpc::commands.has(name + suffix))
CPPUNIT_ASSERT_NO_THROW(rpc::commands.call(name + suffix));
const bool has_allocation = category != torrent::runtime::category_generic;
CPPUNIT_ASSERT(rpc::commands.has(name + ".size"));
CPPUNIT_ASSERT(rpc::commands.has(name + ".max_size"));
CPPUNIT_ASSERT_EQUAL(has_allocation, rpc::commands.has(name + ".min_alloc"));
CPPUNIT_ASSERT_EQUAL(has_allocation, rpc::commands.has(name + ".max_alloc"));
CPPUNIT_ASSERT_EQUAL(has_allocation, rpc::commands.has(name + ".min_alloc.set"));
CPPUNIT_ASSERT_EQUAL(has_allocation, rpc::commands.has(name + ".max_alloc.set"));
}
}
+15
View File
@@ -0,0 +1,15 @@
#include "test/helpers/test_fixture.h"
class TestCommandLocal : public test_fixture {
CPPUNIT_TEST_SUITE(TestCommandLocal);
CPPUNIT_TEST(test_socket_category_commands);
CPPUNIT_TEST_SUITE_END();
public:
void setUp();
void tearDown();
void test_socket_category_commands();
};