Initial commit

This commit is contained in:
2026-06-16 10:56:40 +08:00
commit d2f95e0e27
2047 changed files with 619063 additions and 0 deletions
+248
View File
@@ -0,0 +1,248 @@
#include "TCP_Client.h"
#include "Core/Statistics/Frequency_Limit.h"
#include "Socket_p.h"
#include <chrono>
#include <string>
#ifdef _WIN32
#include <winsock2.h>
#include <ws2tcpip.h>
#else
#include <sys/types.h>
#include <sys/socket.h>
#include <sys/select.h>
#include <errno.h>
#endif
namespace Psc::socket {
static bool fd_writable_now(int fd) {
#ifdef _WIN32
fd_set wfds;
FD_ZERO(&wfds);
FD_SET((SOCKET)fd, &wfds);
timeval tv{};
tv.tv_sec = 0; tv.tv_usec = 0;
int r = select(0, nullptr, &wfds, nullptr, &tv);
return r > 0 && FD_ISSET((SOCKET)fd, &wfds);
#else
fd_set wfds;
FD_ZERO(&wfds);
FD_SET(fd, &wfds);
timeval tv{};
tv.tv_sec = 0; tv.tv_usec = 0;
int r = select(fd + 1, nullptr, &wfds, nullptr, &tv);
return r > 0 && FD_ISSET(fd, &wfds);
#endif
}
void TCP_Client::create() {
send_buffer.init(max_buffer_size);
auto rsf = socket::create_socket_fd(Socket_Type::TCP);
if (!rsf) {
last_err = rsf.error();
socket_logger->c_debug("Tcp_Client::create", {}, LOG_POS + last_err.to_string());
return;
}
socket_fd = rsf.value();
set_state(Not_Set_Field);
auto r = socket::set_reuse(socket_fd, true)
.and_then([this]() { return socket::set_block(socket_fd, false); })
.and_then([this]() { return socket::TCP::set_keep_alive(socket_fd, true); })
.and_then([this]() { return socket::TCP::set_no_delay(socket_fd, true); })
.and_then([this]() { return socket::set_reuse_addr(socket_fd, true); })
;
if (!r) {
last_err = r.error();
socket_logger->error("Tcp_Client::create", {}, Psc::to_string(last_err) + " " + LOG_POS);
(void)socket::close(socket_fd);
socket_fd = -1;
set_state(Not_Created);
return;
}
// 这里最好检查 dest_address 是否已设置(你自己加一个标志位也行)
set_state(Reconnecting);
next_reconnect_tp = std::chrono::steady_clock::now();
tick(); // 立即尝试一次
}
void TCP_Client::schedule_recreate(Enum_Err<NetError> e) {
last_err = e;
(void)socket::close(socket_fd);
set_state(Not_Created);
next_reconnect_tp = std::chrono::steady_clock::now() + std::chrono::milliseconds(reconnect_backoff_ms);
}
bool TCP_Client::start_connect() {
if (socket_fd == -1) return false;
auto r = socket::connect(socket_fd, dest_address);
if (!r) {
last_err = r.error();
if (last_err.nerr == NetError::in_progress) {
set_state(Connecting);
return false; // 还没完成
}
if (last_err.nerr == NetError::already_connected) {
set_state(Wait_Check);
return true;
}
socket_logger->c_debug({}, {}, "连接失败:"+ VAR_STR_2(last_err, *this) + "\n");
// 其它错误:进入退避重连
schedule_recreate(last_err);
return false;
}
return false;
}
bool TCP_Client::check_writeable() {
return true;
}
void TCP_Client::tick() {
if (socket_fd == -1) return;
auto st = state.load();
auto now = std::chrono::steady_clock::now();
if (st == Not_Created) {
create();
return;
}
if (st == Reconnecting) {
if (now >= next_reconnect_tp) {
(void) start_connect();
}
return;
}
if (st == Connecting || st == Wait_Check) {
set_state(Connected);
return;
}
if (st == Connected) {
flush_send_buffer();
return;
}
}
std::string TCP_Client::read() {
if (socket_fd == -1) return "";
// 如果还没连上,先推进连接状态
tick();
if (state.load(std::memory_order_relaxed) != Connected) return "";
auto ret = socket::recv_all(socket_fd);
if (!ret) {
Enum_Err e = ret.error();
// 下面这几个名字按你 NetError 实际枚举改:
// - would_block: 非阻塞没数据,不算错
// - in_progress/not_connected: 连接没完成
if (e.nerr == NetError::would_block) return "";
static Frequency_Limit drop_log_fl;
if (drop_log_fl.test()) {
socket_logger->c_debug("Tcp_Client::read recv失败", {},
to_string() + " " + Ret_To_String(ret) + " " + LOG_POS);
}
// 断线类错误:进入退避重连(必要时你也可以 close+重建fd)
schedule_recreate(e);
return "";
}
auto r = ret.value();
return r;
}
void TCP_Client::flush_send_buffer() {
if (socket_fd == -1) return;
if (state.load(std::memory_order_relaxed) != Connected) return;
if (send_buffer.empty()) return;
constexpr uint32_t patch_len = 1500;
std::string buf;
buf.resize(patch_len);
while (true) {
uint32_t n = send_buffer.peek_best_effort(buf.data(), patch_len);
if (n == 0) break;
auto r = socket::send(socket_fd, buf.data(), n);
if (!r) {
last_err = r.error();
if (last_err.nerr == NetError::would_block) break;
if (last_err.nerr == NetError::connection_aborted) {
schedule_recreate(last_err);
}
else if (last_err.nerr == NetError::connection_reset) {
schedule_recreate(last_err);
} else {
socket_logger->c_debug("Tcp_Client::send", {}, "触发重连 " + VAR_STR_2(last_err, *this));
schedule_recreate(last_err);
}
return;
}
send_buffer.skip_best_effort((uint32_t)r.value());
}
}
void TCP_Client::send(const std::string& data) {
if (data.empty()) return;
(void)send_buffer.write_best_effort((const uint8_t*)data.data(), (uint32_t)data.size());
if (state != Connected) return;
tick();
flush_send_buffer();
}
std::string TCP_Client::to_string() {
return VAR_STR_4(socket_fd, address, dest_address, state);
}
void TCP_Client::close() {
if (socket_fd != -1) {
auto r = socket::close(socket_fd);
(void)r;
socket_fd = -1;
set_state(Not_Created);
}
}
void TCP_Client::set_state(State s) {
auto old_state = state.load();
std::ostringstream oss;
oss << to_string() << "状态变化:" << Psc::to_string(s) << "===>" << Psc::to_string(old_state) << std::endl;
socket_logger->debug({}, {}, oss.str());
if (on_state_change) on_state_change(state, s);
state = s;
}
TCP_Client::~TCP_Client() {
TCP_Client::close();
}
} // namespace Psc::socket