316 lines
8.3 KiB
C++
316 lines
8.3 KiB
C++
#include "TCP_Server.h"
|
|
#include "Core/Statistics/Frequency_Limit.h"
|
|
#include "Core/Statistics/Statistics.h"
|
|
#include "Core/transmit_protocol/core/statistics.h"
|
|
#include "Socket_p.h"
|
|
|
|
#include <unordered_set>
|
|
|
|
namespace Psc::socket {
|
|
|
|
TCP_Server::~TCP_Server() {
|
|
if (socket_fd != -1) {
|
|
socket_logger->c_debug({},{}, "析构关闭的Server_Socket" + TCP_Server::to_string());
|
|
TCP_Server::close();
|
|
}
|
|
}
|
|
|
|
TCP_Server &TCP_Server::set_tcp_no_delay(bool _no_delay) {
|
|
this->no_delay = _no_delay; return *this;
|
|
}
|
|
|
|
void TCP_Server::close() {
|
|
if (socket_fd == -1) {
|
|
socket_logger->c_debug({},{}, "警告 重复关闭的Server_Socket,提前退出!" + LOG_POS);
|
|
return;
|
|
}
|
|
|
|
std::ostringstream oss;
|
|
oss << to_string() << "开始关闭, 剩余连接数:" << tcp_clients.size() << " ";
|
|
|
|
for (auto &kv : tcp_clients) {
|
|
const Socket_FD fd = kv.first;
|
|
auto &conn = kv.second;
|
|
auto &info = conn->info;
|
|
|
|
oss << info.to_string();
|
|
auto r = Psc::socket::close(fd);
|
|
if (!r) {
|
|
oss << info.to_string() + "连接关闭异常" + Psc::to_string(r.error())
|
|
<< " ";
|
|
}
|
|
}
|
|
|
|
auto ret = Psc::socket::close(socket_fd);
|
|
if (!ret) {
|
|
oss << to_string() + "关闭异常" + Psc::to_string(ret.error()) << " ";
|
|
}
|
|
|
|
oss << " 关闭完毕!";
|
|
socket_fd = -1;
|
|
tcp_clients.clear();
|
|
socket_logger->c_debug("Server_Socket", {}, oss.str());
|
|
state = Not_Created;
|
|
}
|
|
|
|
TCP_Server& TCP_Server::create() {
|
|
auto ret = socket::create_socket_fd(Socket_Type::TCP);
|
|
if (!ret) {
|
|
socket_logger->c_debug("创建socket失败!", {}, to_string() + VAR_STR_1(state));
|
|
return *this;
|
|
}
|
|
state = Not_Set_Field;
|
|
socket_fd = ret.value();
|
|
|
|
auto r2 =
|
|
socket::set_block(socket_fd, false)
|
|
.and_then(
|
|
[this]() { return socket::set_reuse_addr(socket_fd, true); })
|
|
.and_then([this]() {
|
|
return socket::TCP::set_no_delay(socket_fd, no_delay);
|
|
})
|
|
// .and_then([this]() { return socket::set_debug(socket_fd, true); })
|
|
.and_then([this]() {
|
|
return socket::set_buffer_size(socket_fd, recv_system_buffer_size);
|
|
});
|
|
|
|
if (!r2) {
|
|
socket_logger->c_debug("设置属性失败!", {}, to_string() + VAR_STR_2(socket_fd, state));
|
|
}
|
|
state = Not_bind;
|
|
return *this;
|
|
}
|
|
TCP_Server &TCP_Server::set_connect_system_buffer_size(size_t size) {
|
|
this->connect_system_buffer_size = size;
|
|
return *this;
|
|
}
|
|
TCP_Server &TCP_Server::set_connect_user_buffer_size(size_t size) {
|
|
this->connect_user_buffer_size = size;
|
|
return *this;
|
|
}
|
|
TCP_Server &TCP_Server::set_recv_system_buffer_size(size_t size) {
|
|
this->recv_system_buffer_size = size;
|
|
return *this;
|
|
}
|
|
|
|
bool TCP_Server::listen(const std::string &ip, std::uint32_t port) {
|
|
return listen(Psc::socket::Sockaddr_In(ip, port));
|
|
}
|
|
|
|
std::string TCP_Server::to_string() {
|
|
return "Server_Socket:[" + address->to_string() + "]";
|
|
}
|
|
|
|
void TCP_Server::flush_clients() {
|
|
if (state != Working) {
|
|
return;
|
|
}
|
|
if (socket_fd == -1)
|
|
return;
|
|
|
|
auto cur = accept();
|
|
while (cur) {
|
|
const Socket_FD fd = cur->info.fd;
|
|
|
|
// insert_or_assign 避免 operator[] 的二次构造/二次查找
|
|
auto [it, inserted] = tcp_clients.insert_or_assign(fd, cur);
|
|
auto &t = it->second;
|
|
|
|
socket_logger->c_debug(
|
|
"", {}, to_string() + " 添加了连接 [" + t->to_string() + "]");
|
|
cur = accept();
|
|
}
|
|
}
|
|
|
|
std::shared_ptr<TCP_Connect> TCP_Server::accept() const {
|
|
|
|
if (socket_fd == -1)
|
|
return nullptr;
|
|
|
|
auto rc = socket::TCP::accept(socket_fd);
|
|
if (!rc) {
|
|
socket_logger->error("accept", {}, Ret_To_String(rc) + " " + LOG_POS);
|
|
return nullptr;
|
|
}
|
|
|
|
auto cur = rc.value(); // std::optional<Accept_Info>
|
|
if (!cur)
|
|
return nullptr;
|
|
|
|
auto cfg =
|
|
socket::set_block(cur->fd, false)
|
|
.and_then([&]() { return socket::set_reuse_addr(cur->fd, true); })
|
|
.and_then(
|
|
[&]() { return socket::set_buffer_size(cur->fd, connect_system_buffer_size); });
|
|
|
|
if (!cfg) {
|
|
// 先记录“配置失败原因”
|
|
socket_logger->error("accept post-config failed", {},
|
|
Ret_To_String(cfg) + " " + LOG_POS);
|
|
|
|
// 再尝试关闭,并记录“关闭是否成功”
|
|
auto cr = socket::close(cur->fd);
|
|
if (!cr) {
|
|
socket_logger->error("accept post-config close failed", {},
|
|
Ret_To_String(cr) + " " + LOG_POS);
|
|
}
|
|
return nullptr;
|
|
}
|
|
|
|
auto ret = std::make_shared<TCP_Connect>();
|
|
ret->info = *cur;
|
|
ret->send_buffer.init(connect_user_buffer_size);
|
|
return ret;
|
|
}
|
|
|
|
bool TCP_Server::listen(const Sockaddr_In &address) {
|
|
this->address = address;
|
|
auto ret = socket::bind(socket_fd, address);
|
|
if (!ret) {
|
|
state = Not_bind;
|
|
return false;
|
|
}
|
|
auto r = socket::TCP::listen(socket_fd, 100);
|
|
if (!r) {
|
|
state = Not_Listen;
|
|
return false;
|
|
}
|
|
state = Working;
|
|
return true;
|
|
}
|
|
|
|
void TCP_Server::write_to_all_clients(const std::string &data) {
|
|
if (socket_fd == -1)
|
|
return;
|
|
// 去重,避免同一 fd 被多次 close/erase
|
|
std::unordered_set<Socket_FD> need_close_fds;
|
|
for (auto &kv : tcp_clients) {
|
|
const Socket_FD fd = kv.first;
|
|
auto &conn = kv.second;
|
|
Accept_Info &info = conn->info;
|
|
|
|
bool ok = conn->send_buffer.write(data.c_str(), data.size());
|
|
if (!ok) {
|
|
conn->lose_speed.update(data.size());
|
|
static Frequency_Limit fl;
|
|
if (fl.test()) {
|
|
socket_logger->c_debug(
|
|
"lose", {},
|
|
"警告:" + info.sockaddr.to_string() +
|
|
" 发送缓冲已满,已丢弃本次数据。 pending=" +
|
|
std::to_string(data.size()) +
|
|
" add=" + std::to_string(data.size()) + " max_pending=" +
|
|
std::to_string(max_pending_send_bytes) + get_error_message() + " " + LOG_POS);
|
|
}
|
|
// 不return 仍然尝试发送
|
|
} else {
|
|
conn->push_speed.update(data.size());
|
|
}
|
|
|
|
std::string buf;
|
|
buf.resize(2000);
|
|
|
|
while (true) {
|
|
auto peek = conn->send_buffer.peek_best_effort(buf.data(), buf.size());
|
|
//if (peek < 500) break;
|
|
|
|
if (peek == 0)
|
|
break;
|
|
auto r_send = socket::send(fd, buf.data(), peek);
|
|
|
|
if (r_send.has_value()) {
|
|
const size_t n = static_cast<size_t>(r_send.value());
|
|
conn->send_num.update(n);
|
|
conn->send_buffer.skip_best_effort(n);
|
|
if (n != peek) {
|
|
break;
|
|
}
|
|
continue;
|
|
}
|
|
|
|
const int sys_ec = r_send.error().native.value();
|
|
|
|
if (socket::is_client_need_close_ec(sys_ec)) {
|
|
need_close_fds.insert(fd);
|
|
break;
|
|
}
|
|
break;
|
|
}
|
|
}
|
|
|
|
// 3) 统一关闭(不在遍历 map 时 erase)
|
|
for (auto fd : need_close_fds) {
|
|
auto it = tcp_clients.find(fd);
|
|
if (it == tcp_clients.end())
|
|
continue;
|
|
|
|
auto info = it->second->info;
|
|
std::string log = to_string() + " 清除了连接 [" + info.to_string() + "]";
|
|
|
|
auto r = socket::close(fd);
|
|
log += r ? " 成功!" : (" 失败!" + Psc::to_string(r.error()));
|
|
|
|
tcp_clients.erase(it);
|
|
socket_logger->c_debug("", {}, log);
|
|
}
|
|
}
|
|
|
|
std::vector<TCP_Server::Read_Info> TCP_Server::read_from_all_clients() {
|
|
if (socket_fd == -1)
|
|
return {};
|
|
|
|
std::vector<Read_Info> ret;
|
|
|
|
auto r_readable = socket::select_read(client_fds(), 0);
|
|
if (!r_readable)
|
|
return {};
|
|
|
|
auto readable = r_readable.value();
|
|
|
|
for (auto fd : readable) {
|
|
auto it = tcp_clients.find(fd);
|
|
if (it == tcp_clients.end())
|
|
continue;
|
|
|
|
auto &conn = it->second;
|
|
auto &info = conn->info;
|
|
|
|
auto r = socket::recv_all(fd);
|
|
if (r.has_value()) {
|
|
ret.push_back(Read_Info{info, r.value()});
|
|
continue;
|
|
}
|
|
|
|
const int ec = r.error().native.value();
|
|
if (is_client_need_close_ec(ec)) {
|
|
auto r2 = socket::close(fd);
|
|
if (!r2) {
|
|
socket_logger->c_debug("tcp服务端", {},
|
|
to_string() + VAR_STR_1(fd) + info.to_string() +
|
|
"关闭连接失败!!!!" + LOG_POS);
|
|
}
|
|
tcp_clients.erase(it);
|
|
}
|
|
}
|
|
|
|
return ret;
|
|
}
|
|
|
|
std::vector<Socket_FD> TCP_Server::client_fds() {
|
|
std::vector<Socket_FD> ret;
|
|
ret.reserve(tcp_clients.size());
|
|
for (auto &c : tcp_clients)
|
|
ret.push_back(c.first);
|
|
return ret;
|
|
}
|
|
|
|
std::vector<std::shared_ptr<TCP_Connect>> TCP_Server::get_all_clients() {
|
|
std::vector<std::shared_ptr<TCP_Connect>> ret;
|
|
ret.reserve(tcp_clients.size());
|
|
for (auto &c : tcp_clients)
|
|
ret.push_back(c.second);
|
|
return ret;
|
|
}
|
|
|
|
} // namespace Psc::socket
|