Fix SCGI threading and added missing header.

This commit is contained in:
Jari Sundell
2025-03-31 18:26:07 +02:00
committed by GitHub
parent 2a998df4f1
commit 2ab7460cbc
9 changed files with 135 additions and 258 deletions
+97 -59
View File
@@ -8,26 +8,22 @@
#include <sys/socket.h>
#include <torrent/exceptions.h>
#include <torrent/poll.h>
#include <torrent/torrent.h>
#include <torrent/utils/log.h>
#include "utils/socket_fd.h"
#include <torrent/utils/thread.h>
#include "control.h"
#include "globals.h"
#include "scgi.h"
// Test:
// #include "core/manager.h"
// #include <rak/timer.h>
// static rak::timer scgiTimer;
#include "rpc/parse_commands.h"
#include "utils/socket_fd.h"
namespace rpc {
// If bufferSize is zero then memcpy won't do anything.
inline void
SCgiTask::realloc_buffer(uint32_t size, const char* buffer, uint32_t bufferSize) {
char* tmp = rak::cacheline_allocator<char>::alloc_size(size);
char* tmp = new char[size];
std::memcpy(tmp, buffer, bufferSize);
::free(m_buffer);
@@ -36,17 +32,16 @@ SCgiTask::realloc_buffer(uint32_t size, const char* buffer, uint32_t bufferSize)
void
SCgiTask::open(SCgi* parent, int fd) {
m_parent = parent;
m_fileDesc = fd;
m_buffer = rak::cacheline_allocator<char>::alloc_size((m_bufferSize = default_buffer_size) + 1);
m_position = m_buffer;
m_body = NULL;
m_parent = parent;
m_fileDesc = fd;
m_buffer = new char[default_buffer_size + 1];
m_buffer_size = default_buffer_size;
m_position = m_buffer;
m_body = NULL;
worker_thread->poll()->open(this);
worker_thread->poll()->insert_read(this);
worker_thread->poll()->insert_error(this);
// scgiTimer = rak::timer::current();
torrent::thread_self->poll()->open(this);
torrent::thread_self->poll()->insert_read(this);
torrent::thread_self->poll()->insert_error(this);
}
void
@@ -54,26 +49,26 @@ SCgiTask::close() {
if (!get_fd().is_valid())
return;
worker_thread->poll()->remove_read(this);
worker_thread->poll()->remove_write(this);
worker_thread->poll()->remove_error(this);
worker_thread->poll()->close(this);
torrent::main_thread()->cancel_callback_and_wait(this);
torrent::thread_self->cancel_callback(this);
torrent::thread_self->poll()->remove_read(this);
torrent::thread_self->poll()->remove_write(this);
torrent::thread_self->poll()->remove_error(this);
torrent::thread_self->poll()->close(this);
get_fd().close();
get_fd().clear();
::free(m_buffer);
m_buffer = NULL;
auto lock = std::lock_guard<std::mutex>(m_result_mutex);
// Test
// char buffer[512];
// sprintf(buffer, "SCgi system call processed: %i", (int)(rak::timer::current() - scgiTimer).usec());
// control->core()->push_log(std::string(buffer));
delete[] m_buffer;
m_buffer = NULL;
}
void
SCgiTask::event_read() {
int bytes = ::recv(m_fileDesc, m_position, m_bufferSize - (m_position - m_buffer), 0);
int bytes = ::recv(m_fileDesc, m_position, m_buffer_size - (m_position - m_buffer), 0);
if (bytes <= 0) {
if (bytes == 0 || !rak::error_number::current().is_blocked_momentary())
@@ -119,17 +114,21 @@ SCgiTask::event_read() {
while (current < header_end) {
char* key = current;
char* key_end = static_cast<char*>(std::memchr(current, '\0', header_end - current));
if (!key_end)
goto event_read_failed;
current = key_end + 1;
if (current >= header_end)
goto event_read_failed;
char* value = current;
char* value_end = static_cast<char*>(std::memchr(current, '\0', header_end - current));
if (!value_end)
goto event_read_failed;
current = value_end + 1;
if (strcmp(key, "CONTENT_LENGTH") == 0) {
@@ -164,45 +163,41 @@ SCgiTask::event_read() {
goto event_read_failed;
}
if ((unsigned int)(content_length + header_size) < m_bufferSize) {
m_bufferSize = content_length + header_size;
if ((unsigned int)(content_length + header_size) < m_buffer_size) {
m_buffer_size = content_length + header_size;
} else if ((unsigned int)content_length <= default_buffer_size) {
m_bufferSize = content_length;
m_buffer_size = content_length;
std::memmove(m_buffer, m_body, std::distance(m_body, m_position));
m_position = m_buffer + std::distance(m_body, m_position);
m_body = m_buffer;
} else {
realloc_buffer((m_bufferSize = content_length) + 1, m_body, std::distance(m_body, m_position));
realloc_buffer((m_buffer_size = content_length) + 1, m_body, std::distance(m_body, m_position));
m_position = m_buffer + std::distance(m_body, m_position);
m_body = m_buffer;
}
}
if ((unsigned int)std::distance(m_buffer, m_position) != m_bufferSize)
if ((unsigned int)std::distance(m_buffer, m_position) != m_buffer_size)
return;
worker_thread->poll()->remove_read(this);
worker_thread->poll()->insert_write(this);
torrent::thread_self->poll()->remove_read(this);
if (m_parent->log_fd() >= 0) {
int __UNUSED result;
// Clean up logging, this is just plain ugly...
// write(m_logFd, "\n---\n", sizeof("\n---\n"));
result = write(m_parent->log_fd(), m_buffer, m_bufferSize);
result = write(m_parent->log_fd(), m_buffer, m_buffer_size);
result = write(m_parent->log_fd(), "\n---\n", sizeof("\n---\n"));
}
lt_log_print_dump(torrent::LOG_RPC_DUMP, m_body, m_bufferSize - std::distance(m_buffer, m_body), "scgi", "RPC read.", 0);
// Close if the call failed, else stay open to write back data.
if (!m_parent->receive_call(this, m_body, m_bufferSize - std::distance(m_buffer, m_body)))
close();
lt_log_print_dump(torrent::LOG_RPC_DUMP, m_body, m_buffer_size - std::distance(m_buffer, m_body), "scgi", "RPC read.", 0);
receive_call(m_body, m_buffer_size - std::distance(m_buffer, m_body));
return;
event_read_failed:
@@ -215,9 +210,9 @@ SCgiTask::event_write() {
// Apple and Solaris do not support MSG_NOSIGNAL,
// so disable this fix until we find a better solution
#if defined(__APPLE__) || defined(__sun__)
int bytes = ::send(m_fileDesc, m_position, m_bufferSize, 0);
int bytes = ::send(m_fileDesc, m_position, m_buffer_size, 0);
#else
int bytes = ::send(m_fileDesc, m_position, m_bufferSize, MSG_NOSIGNAL);
int bytes = ::send(m_fileDesc, m_position, m_buffer_size, MSG_NOSIGNAL);
#endif
if (bytes == -1) {
@@ -228,9 +223,9 @@ SCgiTask::event_write() {
}
m_position += bytes;
m_bufferSize -= bytes;
m_buffer_size -= bytes;
if (bytes == 0 || m_bufferSize == 0)
if (bytes == 0 || m_buffer_size == 0)
return close();
}
@@ -239,13 +234,61 @@ SCgiTask::event_error() {
close();
}
bool
void
SCgiTask::receive_call(const char* buffer, uint32_t length) {
// TODO: Rewrite RpcManager.process to pass the result buffer instead of having to copy it.
auto scgi_thread = torrent::thread_self;
auto result_callback = [this, scgi_thread](const char* b, uint32_t l) {
receive_write(b, l);
scgi_thread->callback(this, [this]() {
// Only need to lock once here as a memory barrier.
m_result_mutex.lock();
m_result_mutex.unlock();
torrent::thread_self->poll()->insert_write(this);
});
};
auto lock = std::lock_guard<std::mutex>(m_result_mutex);
switch (content_type()) {
case rpc::SCgiTask::ContentType::JSON:
torrent::main_thread()->callback(this, [buffer, length, result_callback]() {
rpc.process(RpcManager::RPCType::JSON, buffer, length,
[result_callback](const char* b, uint32_t l) {
result_callback(b, l);
return true;
});
});
break;
case rpc::SCgiTask::ContentType::XML:
torrent::main_thread()->callback(this, [buffer, length, result_callback]() {
rpc.process(RpcManager::RPCType::XML, buffer, length,
[result_callback](const char* b, uint32_t l) {
result_callback(b, l);
return true;
});
});
break;
default:
throw torrent::internal_error("SCgiTask::receive_call(...) received bad input.");
}
}
void
SCgiTask::receive_write(const char* buffer, uint32_t length) {
if (buffer == NULL || length > (100 << 20))
throw torrent::internal_error("SCgiTask::receive_write(...) received bad input.");
auto lock = std::lock_guard<std::mutex>(m_result_mutex);
// Need to cast due to a bug in MacOSX gcc-4.0.1.
if (length + 256 > std::max(m_bufferSize, (unsigned int)default_buffer_size))
if (length + 256 > std::max(m_buffer_size, (unsigned int)default_buffer_size))
realloc_buffer(length + 256, NULL, 0);
const auto header = m_content_type == ContentType::JSON
@@ -253,25 +296,20 @@ SCgiTask::receive_write(const char* buffer, uint32_t length) {
: "Status: 200 OK\r\nContent-Type: text/xml\r\nContent-Length: %i\r\n\r\n";
// Who ever bothers to check the return value?
int headerSize = sprintf(m_buffer, header, length);
int headerSize = snprintf(m_buffer, m_buffer_size, header, length);
m_position = m_buffer;
m_bufferSize = length + headerSize;
m_buffer_size = length + headerSize;
std::memcpy(m_buffer + headerSize, buffer, length);
if (m_parent->log_fd() >= 0) {
int __UNUSED result;
// Clean up logging, this is just plain ugly...
// write(m_logFd, "\n---\n", sizeof("\n---\n"));
result = write(m_parent->log_fd(), m_buffer, m_bufferSize);
int result [[maybe_unused]];
result = write(m_parent->log_fd(), m_buffer, m_buffer_size);
result = write(m_parent->log_fd(), "\n---\n", sizeof("\n---\n"));
}
lt_log_print_dump(torrent::LOG_RPC_DUMP, m_buffer, m_bufferSize, "scgi", "RPC write.", 0);
event_write();
return true;
lt_log_print_dump(torrent::LOG_RPC_DUMP, m_buffer, m_buffer_size, "scgi", "RPC write.", 0);
}
} // namespace rpc