A lot of file transfer updates

This commit is contained in:
WolverinDEV
2020-05-13 11:32:08 +02:00
parent 1a2dd4a008
commit 90b1646876
45 changed files with 1989 additions and 2806 deletions
+167 -13
View File
@@ -8,7 +8,9 @@
#include <random>
#include "./LocalFileProvider.h"
#include "LocalFileProvider.h"
#include <experimental/filesystem>
namespace fs = std::experimental::filesystem;
using namespace ts::server::file;
using namespace ts::server::file::transfer;
@@ -101,29 +103,59 @@ void LocalFileTransfer::stop() {
this->shutdown_client_worker();
}
std::shared_ptr<ExecuteResponse<TransferInitError, std::shared_ptr<Transfer>>> LocalFileTransfer::initialize_icon_transfer(Transfer::Direction direction, ServerId sid, const TransferInfo &info) {
return this->initialize_transfer(direction, sid, 0, Transfer::TARGET_TYPE_ICON, info);
std::shared_ptr<ExecuteResponse<TransferInitError, std::shared_ptr<Transfer>>>
LocalFileTransfer::initialize_icon_transfer(Transfer::Direction direction, const std::shared_ptr<VirtualFileServer> &server, const TransferInfo &info) {
return this->initialize_transfer(direction, server, 0, Transfer::TARGET_TYPE_ICON, info);
}
std::shared_ptr<ExecuteResponse<TransferInitError, std::shared_ptr<Transfer>>> LocalFileTransfer::initialize_avatar_transfer(Transfer::Direction direction, ServerId sid, const TransferInfo &info) {
return this->initialize_transfer(direction, sid, 0, Transfer::TARGET_TYPE_AVATAR, info);
std::shared_ptr<ExecuteResponse<TransferInitError, std::shared_ptr<Transfer>>>
LocalFileTransfer::initialize_avatar_transfer(Transfer::Direction direction, const std::shared_ptr<VirtualFileServer> &server, const TransferInfo &info) {
return this->initialize_transfer(direction, server, 0, Transfer::TARGET_TYPE_AVATAR, info);
}
std::shared_ptr<ExecuteResponse<TransferInitError, std::shared_ptr<Transfer>>> LocalFileTransfer::initialize_channel_transfer(Transfer::Direction direction, ServerId sid, ChannelId cid, const TransferInfo &info) {
return this->initialize_transfer(direction, sid, cid, Transfer::TARGET_TYPE_CHANNEL_FILE, info);
std::shared_ptr<ExecuteResponse<TransferInitError, std::shared_ptr<Transfer>>>
LocalFileTransfer::initialize_channel_transfer(Transfer::Direction direction, const std::shared_ptr<VirtualFileServer> &server, ChannelId cid, const TransferInfo &info) {
return this->initialize_transfer(direction, server, cid, Transfer::TARGET_TYPE_CHANNEL_FILE, info);
}
std::shared_ptr<ExecuteResponse<TransferInitError, std::shared_ptr<Transfer>>> LocalFileTransfer::initialize_transfer(
Transfer::Direction direction, ServerId sid, ChannelId cid,
Transfer::Direction direction, const std::shared_ptr<VirtualFileServer> &server, ChannelId cid,
Transfer::TargetType ttype,
const TransferInfo &info) {
auto response = this->create_execute_response<TransferInitError, std::shared_ptr<Transfer>>();
/* TODO: test for a transfer limit */
std::lock_guard clock{this->transfer_create_mutex};
if(info.max_concurrent_transfers > 0) {
std::unique_lock tlock{this->transfers_mutex};
{
auto transfers = std::count_if(this->transfers_.begin(), this->transfers_.end(), [&](const std::shared_ptr<FileClient>& client) {
return client->transfer && client->transfer->client_unique_id == info.client_unique_id;
});
transfers += std::count_if(this->pending_transfers.begin(), this->pending_transfers.end(), [&](const std::shared_ptr<Transfer>& transfer) {
return transfer->client_unique_id == info.client_unique_id;
});
if(transfers >= info.max_concurrent_transfers) {
response->emplace_fail(TransferInitError::CLIENT_TOO_MANY_TRANSFERS, std::to_string(transfers));
return response;
}
}
{
auto server_transfers = this->pending_transfers.size();
server_transfers += std::count_if(this->transfers_.begin(), this->transfers_.end(), [&](const std::shared_ptr<FileClient>& client) {
return client->transfer;
});
if(server_transfers >= this->max_concurrent_transfers) {
response->emplace_fail(TransferInitError::SERVER_TOO_MANY_TRANSFERS, std::to_string(server_transfers));
return response;
}
}
}
auto transfer = std::make_shared<Transfer>();
transfer->server_transfer_id = ++this->current_transfer_id;
transfer->server_id = sid;
transfer->server_transfer_id = server->generate_transfer_id();
transfer->server = server;
transfer->channel_id = cid;
transfer->target_type = ttype;
transfer->direction = direction;
@@ -142,6 +174,8 @@ std::shared_ptr<ExecuteResponse<TransferInitError, std::shared_ptr<Transfer>>> L
transfer->file_offset = info.file_offset;
transfer->expected_file_size = info.expected_file_size;
transfer->max_bandwidth = info.max_bandwidth;
transfer->client_unique_id = info.client_unique_id;
transfer->client_id = info.client_id;
constexpr static std::string_view kTokenCharacters{"abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789"};
for(auto& c : transfer->transfer_key)
@@ -152,6 +186,70 @@ std::shared_ptr<ExecuteResponse<TransferInitError, std::shared_ptr<Transfer>>> L
transfer->initialized_timestamp = std::chrono::system_clock::now();
{
std::string absolute_path{};
switch (transfer->target_type) {
case Transfer::TARGET_TYPE_AVATAR:
absolute_path = this->file_system_.absolute_avatar_path(transfer->server, transfer->target_file_path);
logMessage(LOG_FT, "Initialized avatar transfer for avatar \"{}\" ({} bytes, transferring {} bytes).", transfer->target_file_path, transfer->expected_file_size, transfer->expected_file_size - transfer->file_offset);
break;
case Transfer::TARGET_TYPE_ICON:
absolute_path = this->file_system_.absolute_icon_path(transfer->server, transfer->target_file_path);
logMessage(LOG_FT, "Initialized icon transfer for icon \"{}\" ({} bytes, transferring {} bytes).",
transfer->target_file_path, transfer->expected_file_size, transfer->expected_file_size - transfer->file_offset);
break;
case Transfer::TARGET_TYPE_CHANNEL_FILE:
absolute_path = this->file_system_.absolute_channel_path(transfer->server, transfer->channel_id, transfer->target_file_path);
logMessage(LOG_FT, "Initialized channel transfer for file \"{}/{}\" ({} bytes, transferring {} bytes).",
transfer->channel_id, transfer->target_file_path, transfer->expected_file_size, transfer->expected_file_size - transfer->file_offset);
break;
case Transfer::TARGET_TYPE_UNKNOWN:
default:
response->emplace_fail(TransferInitError::INVALID_FILE_TYPE, "");
return response;
}
transfer->absolute_file_path = absolute_path;
const auto root_path_length = this->file_system_.root_path().size();
if(root_path_length < absolute_path.size())
transfer->relative_file_path = absolute_path.substr(root_path_length);
else
transfer->relative_file_path = "error";
transfer->file_name = fs::u8path(absolute_path).filename();
}
if(direction == Transfer::DIRECTION_DOWNLOAD) {
auto path = fs::u8path(transfer->absolute_file_path);
std::error_code error{};
if(!fs::exists(path, error)) {
response->emplace_fail(TransferInitError::FILE_DOES_NOT_EXISTS, "");
return response;
} else if(error) {
logWarning(LOG_FT, "Failed to check for file at {}: {}. Assuming it does not exists.", transfer->absolute_file_path, error.value(), error.message());
response->emplace_fail(TransferInitError::FILE_DOES_NOT_EXISTS, "");
return response;
}
auto status = fs::status(path, error);
if(error) {
logWarning(LOG_FT, "Failed to status for file at {}: {}. Ignoring file transfer.", transfer->absolute_file_path, error.value(), error.message());
response->emplace_fail(TransferInitError::IO_ERROR, "stat");
return response;
}
if(status.type() != fs::file_type::regular) {
response->emplace_fail(TransferInitError::FILE_IS_NOT_A_FILE, "");
return response;
}
transfer->expected_file_size = fs::file_size(path, error);
if(error) {
logWarning(LOG_FT, "Failed to get file size for file at {}: {}. Ignoring file transfer.", transfer->absolute_file_path, error.value(), error.message());
response->emplace_fail(TransferInitError::IO_ERROR, "file_size");
return response;
}
}
{
std::lock_guard tlock{this->transfers_mutex};
this->pending_transfers.push_back(transfer);
@@ -164,7 +262,7 @@ std::shared_ptr<ExecuteResponse<TransferInitError, std::shared_ptr<Transfer>>> L
return response;
}
std::shared_ptr<ExecuteResponse<TransferActionError>> LocalFileTransfer::stop_transfer(transfer_id id, bool flush) {
std::shared_ptr<ExecuteResponse<TransferActionError>> LocalFileTransfer::stop_transfer(const std::shared_ptr<VirtualFileServer>& server, transfer_id id, bool flush) {
auto response = this->create_execute_response<TransferActionError>();
std::shared_ptr<Transfer> transfer{};
@@ -174,13 +272,13 @@ std::shared_ptr<ExecuteResponse<TransferActionError>> LocalFileTransfer::stop_tr
std::lock_guard tlock{this->transfers_mutex};
auto ct_it = std::find_if(this->transfers_.begin(), this->transfers_.end(), [&](const std::shared_ptr<FileClient>& t) {
return t->transfer && t->transfer->server_transfer_id == id;
return t->transfer && t->transfer->server_transfer_id == id && t->transfer->server == server;
});
if(ct_it != this->transfers_.end())
connected_transfer = *ct_it;
else {
auto t_it = std::find_if(this->pending_transfers.begin(), this->pending_transfers.end(), [&](const std::shared_ptr<Transfer>& t) {
return t->server_transfer_id == id;
return t->server_transfer_id == id && t->server == server;
});
if(t_it != this->pending_transfers.end()) {
transfer = *t_it;
@@ -209,4 +307,60 @@ std::shared_ptr<ExecuteResponse<TransferActionError>> LocalFileTransfer::stop_tr
response->emplace_success();
return response;
}
inline void apply_transfer_info(const std::shared_ptr<Transfer>& transfer, ActiveFileTransfer& info) {
info.server_transfer_id = transfer->server_transfer_id;
info.client_transfer_id = transfer->client_transfer_id;
info.direction = transfer->direction;
info.client_id = transfer->client_id;
info.client_unique_id = transfer->client_unique_id;
info.file_path = transfer->relative_file_path;
info.file_name = transfer->file_name;
info.expected_size = transfer->expected_file_size;
}
std::shared_ptr<ExecuteResponse<TransferListError, std::vector<ActiveFileTransfer>>> LocalFileTransfer::list_transfer() {
std::vector<ActiveFileTransfer> transfer_infos{};
auto response = this->create_execute_response<TransferListError, std::vector<ActiveFileTransfer>>();
std::unique_lock tlock{this->transfers_mutex};
auto awaiting_transfers = this->pending_transfers;
auto running_transfers = this->transfers_;
tlock.unlock();
transfer_infos.reserve(awaiting_transfers.size() + running_transfers.size());
for(const auto& transfer : awaiting_transfers) {
ActiveFileTransfer info{};
apply_transfer_info(transfer, info);
info.size_done = transfer->file_offset;
info.status = ActiveFileTransfer::NOT_STARTED;
info.runtime = std::chrono::milliseconds{0};
info.average_speed = 0;
info.current_speed = 0;
transfer_infos.push_back(info);
}
for(const auto& client : running_transfers) {
auto transfer = client->transfer;
if(!transfer) continue;
ActiveFileTransfer info{};
apply_transfer_info(transfer, info);
info.size_done = transfer->file_offset + client->statistics.file_transferred.total_bytes;
info.status = ActiveFileTransfer::RUNNING;
info.runtime = std::chrono::floor<std::chrono::milliseconds>(std::chrono::system_clock::now() - client->timings.key_received);
info.average_speed = client->statistics.file_transferred.average_bandwidth();
info.current_speed = client->statistics.file_transferred.current_bandwidth();
transfer_infos.push_back(info);
}
response->emplace_success(std::move(transfer_infos));
return response;
}