一、Socket编程TCP
服务端流程: socket() -> bind() -> listen() -> accept() -> read() -> (处理业务逻辑) -> write() -> close() 客户端流程: socket() -> connect() -> read()/write() -> close()
1.1 服务端流程
① socket()
-
作用:创建一个 socket 文件描述符,指定协议族(如
AF_INET)、类型(SOCK_STREAM表示 TCP)、协议(通常填 0)。 -
内核行为:分配一个
struct socket和struct sock,初始化状态为CLOSED。 -
注意:此时该 socket 还没有绑定地址,不能直接用于通信。
② bind()
-
作用:将 socket 绑定到本地 IP 和端口(例如
0.0.0.0:8080)。 -
内核行为:检查端口是否被占用,将本地地址记录到 socket 结构中。
-
注意:服务器必须显式绑定端口,客户端通常不需要(内核会自动分配临时端口)。
③ listen()
-
作用:将 socket 从主动套接字转变为被动套接字,进入
LISTEN状态,并设置内核维护的连接队列长度(backlog)。 -
内核行为:创建两个队列:
-
半连接队列(SYN Queue):存放收到 SYN 但尚未完成三次握手的连接。
-
已完成连接队列(Accept Queue):存放已完成三次握手、等待应用程序
accept()取走的连接。
-
-
注意:
listen()之后该 socket 只能用来接受连接,不能发送数据。
④ accept()
-
作用:从已完成连接队列中取出一个已建立的连接,返回新的 socket 文件描述符。
-
内核行为:若队列为空且 socket 为阻塞模式,则进程睡眠直到有新连接到达;否则立即返回。
-
返回的新 socket:它继承了监听 socket 的本地地址,但具有独立的四元组(加上客户端 IP 和端口),状态为
ESTABLISHED,可进行数据收发。 -
注意:监听 socket 本身不参与数据通信,只负责“接客”。
⑤ read() / recv()
-
作用:从连接 socket 读取客户端发来的数据。
-
内核行为:从该连接的接收缓冲区拷贝数据到用户空间;若缓冲区为空则阻塞(默认阻塞模式)。
-
注意:TCP 是字节流,
read()返回的数据长度可能小于请求长度,也可能包含多条消息,需要应用层自行处理粘包/拆包。
⑥ 处理业务逻辑
-
根据应用需求处理收到的数据,可能涉及数据库、计算等。
⑦ write() / send()
-
作用:将响应数据发送给客户端。
-
内核行为:将数据拷贝到该连接的发送缓冲区,内核负责 TCP 分段、重传、流量控制等。
-
注意:
write()成功返回只表示数据进入了内核缓冲区,并不保证对端已经收到(TCP 可靠性由内核保证,但应用层无法知道对方是否真正处理)。
⑧ close()
-
作用:关闭连接 socket,释放文件描述符和相关资源。
-
内核行为:若发送缓冲区还有数据则尝试发送完毕,然后发送 FIN 包,进入四次挥手流程。
-
注意:监听 socket 通常在整个服务运行期间保持打开,直到服务停止才关闭。
1.2 客户端流程
① socket()
-
同服务端,创建一个未连接的 TCP socket。
② connect()
-
作用:向服务器发起连接请求,指定服务器的 IP 和端口。
-
内核行为:
-
客户端发送 SYN 包,进入
SYN_SENT状态。 -
收到服务器的 SYN+ACK 后,回复 ACK,连接建立,状态变为
ESTABLISHED。 -
同时内核为客户端自动绑定一个临时端口(如果之前未
bind())。
-
-
注意:
connect()在阻塞模式下会一直等待直到连接成功或失败(如超时、拒绝)。
③ read() / write()
-
连接建立后,客户端可以收发数据。其行为和服务端的
read()/write()对称。
④ close()
-
关闭连接,发起四次挥手,释放资源。
二、Demo演示
2.1 EchoServer
2.1.1 头文件包含
A. TcpServer.hpp
#include "myLog.hpp"
#include "Common.hpp"
#include "InetAddr.hpp"
#include "myThread.hpp"
#include "ThreadPool.hpp"
#include <signal.h>
namespace TcpServerModule
{
using namespace LogModule;
using namespace ThreadModule;
using namespace ThreadPoolModule;
using task_t = std::function<void()>;
const int defaultbacklog = 128;
class TcpServer : public NoCopy
{
public:
TcpServer() = default;
TcpServer(int port)
: _port(port)
{
}
void Service(int connfd, const InetAddr &peer)
{
// 1.面向字节流接受消息
char buffer[1024];
for (;;)
{
//简单处理:未对字节流进行消息边界处理
ssize_t n = read(connfd, buffer, sizeof(buffer) - 1);
if (n > 0)
{
buffer[n] = 0; // 设置为C风格字符串, n<= sizeof(buffer)-1
LOG(LogLevel::DEBUG) << peer.GetStringAddr() << " #" << buffer;
// 2. 写回数据
std::string echo_string = "echo# " + std::string(buffer);
// write不保证一次写完, 需循环发送直到全部写入(对端断开时返回-1)
size_t total = 0;
while (total < echo_string.size())
{
ssize_t w = write(connfd, echo_string.c_str() + total,
echo_string.size() - total);
if (w == -1)
{
LOG(LogLevel::ERROR) << "write failed: "
<< peer.GetStringAddr() << " "
<< strerror(errno) << std::endl;
close(connfd);
return;
}
total += w;
}
}
else if (n == 0)
{
LOG(LogLevel::DEBUG) << peer.GetStringAddr() << " Exited...";
close(connfd);
break;
}
else
{
LOG(LogLevel::DEBUG) << peer.GetStringAddr() << " Abnormal...";
close(connfd);
break;
}
}
}
// 打开socket文件 + bind绑定端口号 + listen监听连接请求
void Init()
{
// 0.忽略SIGPIPE: 对端断开后write返回-1并由Service统一处理, 而不是进程被信号杀死
signal(SIGPIPE, SIG_IGN);
// 0.5 SIGCHLD置为SIG_IGN: 子进程退出时由内核自动回收, 不会产生僵尸进程
signal(SIGCHLD, SIG_IGN);
// 1.打开socket文件
_listentfd = socket(AF_INET, SOCK_STREAM, 0);
if (_listentfd == -1)
{
LOG(LogLevel::ERROR) << "The TcpServer failed to open the socket file" << std::endl;
exit(SOCKET_ERR);
}
LOG(LogLevel::INFO) << "socket success: " << _listentfd; // 3
// 1.5 允许服务端退出后立刻重启复用端口(否则TIME_WAIT残留会导致bind失败: Address already in use)
int opt = 1;
if (setsockopt(_listentfd, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt)) == -1)
{
LOG(LogLevel::WARNING) << "setsockopt SO_REUSEADDR failed" << std::endl;
}
// 2.绑定端口号
InetAddr local(_port);
int n = bind(_listentfd, local.GetAddrPtr(), local.GetLen());
if (n == -1)
{
LOG(LogLevel::ERROR) << "The TcpServer failed to bind the port" << std::endl;
close(_listentfd);
exit(BIND_ERR);
}
LOG(LogLevel::INFO) << "bind success: " << _listentfd; // 3
// 3.设置监听状态
n = listen(_listentfd, defaultbacklog);
if (n == -1)
{
LOG(LogLevel::ERROR) << "The TcpServer failed to set up listener" << std::endl;
close(_listentfd);
exit(LISTEN_ERR);
}
LOG(LogLevel::INFO) << "listen success: " << _listentfd; // 3
}
// 循环accept: 服务完一个连接后继续等待下一个, 常驻运行, 不会服务一个客户端就退出
// 注:
// v0: 串行服务(一个时刻只能服务一个连接),
// v1: 多进程 (父进程进行接受连接,关闭服务文件描述符; 子进程进行服务,关闭监听文件描述符)
// v2: 多线程
// v3: 线程池
void Processv1(int connfd, InetAddr &addr)
{
// v1多进程版本
// version1.0:利用多进程(每个连接fork一个子进程独立服务, 互不阻塞)
// 四个要点: 1.子进程处理完必须exit, 否则会回到for循环抢accept
// 2.父子都要close自己不再需要的fd(fd表是共享的, 引用计数)
// 3.fork失败只丢弃当前连接, 不应终止整个服务器
// 4.Init中已将SIGCHLD置为SIG_IGN, 子进程退出由内核回收
// version1.1:利用多进程(每个连接fork一个子进程独立服务, 互不阻塞)
// 四个要点: 1.子进程提前退出,由父进程进行回收,孙子进程执行服务。
// 孙子进程退出时变为孤儿进程由1号进程进行接管
// 2.父子都要close自己不再需要的fd(fd表是共享的, 引用计数)
// 3.fork失败只丢弃当前连接, 不应终止整个服务器
int pid = fork();
if (pid < 0)
{
LOG(LogLevel::ERROR) << "fork fail: " << strerror(errno) << std::endl;
close(connfd); // 无人处理该连接, 关闭释放
// continue; // 单个连接fork失败, 不终止服务器
}
else if (pid == 0)
{
// 子进程: 只服务, 不再需要监听套接字
close(_listentfd);
Service(connfd, addr); // 子进程执行服务
close(connfd);
exit(NormalExit); // 必须exit, 否则子进程会回到for(;;)循环继续accept
}
else
{
// 父进程: 不参与服务, 必须关闭connfd, 避免文件描述符泄漏
close(connfd);
}
}
void processv2(int connfd, InetAddr &addr)
{
// version2:利用多线程
// 三点: 1.必须调用start()真正创建线程, 否则线程从未启动, Service不执行
// 2.必须按值捕获connfd/addr: 引用捕获的是accept循环的栈局部量,
// 下一轮accept复用该栈内存后, 服务线程读到的将是新连接的数据
// 3.必须detach()分离: 主线程只负责accept, 不等待服务线程;
// Service内部已close(connfd), 主线程不得关闭(线程共享fd表)
std::unique_ptr<Thread> pthread = std::make_unique<Thread>(
"thread",
[this, connfd, addr]()
{
Service(connfd, addr);
});
pthread->start();
pthread->detach();
}
void Run()
{
// version3:利用线程池
// 四点: 1.池只创建一次且常驻: 绝不能放进accept循环, 否则每个连接都新建一个15线程的池
// 2.必须调用Start()才真正启动worker线程, 否则任务无人执行(worker等队列, 队列永远空)
// 3.任务必须按值捕获connfd/addr: 引用捕获会随下一轮accept悬空(与v2同因)
// 4.connfd的关闭仍由Service负责; 任务粒度= 一次 Service(connfd, addr) 调用
ThreadPool<task_t> pool(15);
pool.Start();
for (;;)
{
// 1.通过accept获取连接
struct sockaddr_in peer;
socklen_t len = sizeof(peer);
int connfd = accept(_listentfd, (struct sockaddr *)&peer, &len);
if (connfd == -1)
{
if (errno == EINTR)
continue; // 被信号打断, 属正常情况, 重新accept
LOG(LogLevel::ERROR) << "The TcpServer failed to obtain a new connection" << std::endl;
continue; // 单个accept失败不应终止整个服务
}
// 2.成功建立连接
InetAddr addr(peer);
LOG(LogLevel::INFO) << "accept success, peer addr : " << addr.GetStringAddr() << std::endl;
// 3.执行服务
// version3:利用线程池
pool.PushTask(
[this, connfd, addr]()
{
Service(connfd, addr);
});
}
}
private:
uint16_t _port;
int _listentfd = -1;
};
}
2.1.2 源代码
A. TcpClient.cc
#include "Common.hpp"
#include "TcpServer.hpp"
#include "myLog.hpp"
#include "InetAddr.hpp"
#include <signal.h>
using namespace LogModule;
// 服务端接收缓冲 1024, 回显时前缀 "echo# " 占 6 字节,
// 超过该长度的消息回声会被截断, 故客户端发送前先对齐限制
constexpr size_t kMaxMsgSize = 1024 - 1 - 6;
int main(int argc, char *argv[])
{
if (argc != 3)
{
std::cerr << "Usage: " << argv[0] << " " << "ip" << " " << "port" << std::endl;
exit(USAGE_ERR);
}
// 0.忽略SIGPIPE: 对端关闭后继续写会让进程默认退出, 改为由read/write正常返回错误
signal(SIGPIPE, SIG_IGN);
// 1.创建socket文件
int sfd = socket(AF_INET, SOCK_STREAM, 0);
if (sfd == -1)
{
LOG(LogLevel::ERROR) << "socket failed" << std::endl;
return 1;
}
// 2.构建服务端地址
InetAddr dest_addr(argv[1], argv[2]);
// 3.建立连接: TCP面向连接, 失败必须在此暴露(如服务端未启动), 而不是继续执行
if (connect(sfd, dest_addr.GetAddrPtr(), dest_addr.GetLen()) == -1)
{
LOG(LogLevel::ERROR) << "connect failed: " << strerror(errno) << std::endl;
close(sfd);
return 1;
}
LOG(LogLevel::INFO) << "connect success" << std::endl;
for (;;)
{
// 4.向服务端发送消息
std::string msg;
std::cout << "Please Enter# ";
if (!std::getline(std::cin, msg))
break; // stdin结束(Ctrl+D或管道输入完毕), 退出循环而不是用空消息刷屏
if (msg.size() > kMaxMsgSize)
msg.resize(kMaxMsgSize);
// TCP面向字节流: 使用通用fd接口write(与send在flags=0时等价, 无UDP遗留的目的地址参数)
ssize_t sent = write(sfd, msg.c_str(), msg.size());
if (sent == -1)
{
LOG(LogLevel::ERROR) << "send failed: " << strerror(errno) << std::endl;
break; // 连接可能已断开, 退出并统一收尾
}
// 5.接受服务端的回显消息
char buffer[1024];
ssize_t n = read(sfd, buffer, sizeof(buffer) - 1);
if (n > 0)
{
buffer[n] = 0;
std::cout << buffer << std::endl;
}
else if (n == 0)
{
// 对端正常关闭连接
LOG(LogLevel::WARNING) << "server closed the connection" << std::endl;
break;
}
else
{
LOG(LogLevel::ERROR) << "recv failed: " << strerror(errno) << std::endl;
break;
}
// 注: 阻塞socket下recv不会返回EAGAIN, 原"超时"分支是不可能触发的死代码, 已删除;
// TCP本身不丢消息, 收不到回显只会是对端断开(返回0)或出错(返回-1)
}
close(sfd);
return 0;
}
B. TcpServer.cc
#include "TcpServer.hpp"
#include "InetAddr.hpp"
#include <memory>
using namespace TcpServerModule;
int main(int argc, char *argv[])
{
if (argc != 2)
{
std::cerr << "Usage: " << argv[0] << " port" << std::endl;
exit(USAGE_ERR);
}
uint16_t port = std::stoi(argv[1]);
std::unique_ptr<TcpServer> tcp = std::make_unique<TcpServer>(port);
tcp->Init();
tcp->Run();
return 0;
}
2.2 CommandServer
2.2.1 头文件包含
A.TcpServer.hpp
#pragma once
#include "myLog.hpp"
#include "Common.hpp"
#include "InetAddr.hpp"
#include "ThreadPool.hpp"
#include "Command.hpp"
#include <signal.h>
namespace TcpServerModule
{
using namespace LogModule;
using namespace ThreadPoolModule;
using task_t = std::function<void()>;
const int defaultbacklog = 128;
class TcpServer : public NoCopy
{
public:
TcpServer() = default;
TcpServer(int port)
: _port(port)
{
}
// 打开socket文件 + bind绑定端口号 + listen监听连接请求
void Init()
{
// 0.忽略SIGPIPE: 对端断开后write返回-1并由Service统一处理, 而不是进程被信号杀死
signal(SIGPIPE, SIG_IGN);
// 0.5 SIGCHLD置为SIG_IGN: 子进程退出时由内核自动回收, 不会产生僵尸进程
signal(SIGCHLD, SIG_IGN);
// 1.打开socket文件
_listentfd = socket(AF_INET, SOCK_STREAM, 0);
if (_listentfd == -1)
{
LOG(LogLevel::ERROR) << "The TcpServer failed to open the socket file" << std::endl;
exit(SOCKET_ERR);
}
LOG(LogLevel::INFO) << "socket success: " << _listentfd; // 3
// 1.5 允许服务端退出后立刻重启复用端口(否则TIME_WAIT残留会导致bind失败: Address already in use)
int opt = 1;
if (setsockopt(_listentfd, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt)) == -1)
{
LOG(LogLevel::WARNING) << "setsockopt SO_REUSEADDR failed" << std::endl;
}
// 2.绑定端口号
InetAddr local(_port);
int n = bind(_listentfd, local.GetAddrPtr(), local.GetLen());
if (n == -1)
{
LOG(LogLevel::ERROR) << "The TcpServer failed to bind the port" << std::endl;
close(_listentfd);
exit(BIND_ERR);
}
LOG(LogLevel::INFO) << "bind success: " << _listentfd; // 3
// 3.设置监听状态
n = listen(_listentfd, defaultbacklog);
if (n == -1)
{
LOG(LogLevel::ERROR) << "The TcpServer failed to set up listener" << std::endl;
close(_listentfd);
exit(LISTEN_ERR);
}
LOG(LogLevel::INFO) << "listen success: " << _listentfd; // 3
}
// 循环accept: 服务完一个连接后继续等待下一个, 常驻运行, 不会服务一个客户端就退出
// 注: 利用线程池,并发执行
void Service(const InetAddr &addr, int connfd)
{
Command cmd(connfd);
std::string commandstr = cmd.RecvCommand();
std::string result = cmd.Execute(commandstr);
cmd.SendCommand(result);
}
void Run()
{
// version3:利用线程池
// 四点: 1.池只创建一次且常驻: 绝不能放进accept循环, 否则每个连接都新建一个15线程的池
// 2.必须调用Start()才真正启动worker线程, 否则任务无人执行(worker等队列, 队列永远空)
// 3.任务必须按值捕获connfd/addr: 引用捕获会随下一轮accept悬空(与v2同因)
// 4.connfd的关闭由Service内Command的析构(RAII)负责; 任务粒度=一次 Service(connfd, addr) 调用
ThreadPool<task_t> pool(15);
pool.Start();
for (;;)
{
// 1.通过accept获取连接
struct sockaddr_in peer;
socklen_t len = sizeof(peer);
int connfd = accept(_listentfd, (struct sockaddr *)&peer, &len);
if (connfd == -1)
{
if (errno == EINTR)
continue; // 被信号打断, 属正常情况, 重新accept
LOG(LogLevel::ERROR) << "The TcpServer failed to obtain a new connection" << std::endl;
continue; // 单个accept失败不应终止整个服务
}
// 2.成功建立连接
InetAddr addr(peer);
LOG(LogLevel::INFO) << "accept success, peer addr : " << addr.GetStringAddr() << std::endl;
// 3.执行服务
// 利用线程池
pool.PushTask(
[this, addr, connfd]()
{
Service(addr, connfd);
});
}
}
private:
uint16_t _port;
int _listentfd = -1;
};
}
B.Command.hpp
#pragma once
#include "myLog.hpp"
#include "Common.hpp"
#include <cstdio> // popen/pclose
#include <cstring> // strerror
#include <cerrno> // errno
#include <iostream>
#include <string>
#include <unistd.h>
#include <sys/types.h>
#include <sys/socket.h>
#include <set>
using namespace LogModule;
class Command
{
public:
Command() = default;
Command(int sockfd)
: _sockfd(sockfd)
{
}
// 连接结束必须回收 fd: 否则每个连接泄漏一个 fd, 长跑服务器最终 EMFILE 无法 accept
// 由析构统一负责, 成功/失败路径都能关闭, 避免在业务代码里漏掉 close
~Command()
{
if (_sockfd != -1)
{
close(_sockfd);
}
}
// 同一 fd 只能被一个对象管理: 拷贝会导致两个对象关闭同一 fd (double close)
Command(const Command &) = delete;
Command &operator=(const Command &) = delete;
bool IsSafe(const std::string &command)
{
return _command.find(command) != _command.end();
}
// 接受指令
std::string RecvCommand()
{
char line[1024];
ssize_t n = recv(_sockfd, line, sizeof(line) - 1, 0);
if (n > 0)
{
line[n] = 0;
std::string cmd(line);
// 去除尾部换行/回车: 客户端按行发送时通常带 '\n',
// 否则白名单精确匹配会失败返回 "command unsafe"
while (!cmd.empty() && (cmd.back() == '\n' || cmd.back() == '\r'))
cmd.pop_back();
return cmd;
}
return std::string(); // n==0(对端关闭) 或 n==-1(出错) 一律视为无命令
}
// 处理指令
std::string Execute(const std::string &command)
{
if (!IsSafe(command))
return "command unsafe";
FILE *fp = popen(command.c_str(), "r");
if (fp == nullptr)
{
// 注意: 绝不能在这里 exit(), 它运行在工作线程里, 会杀死整个服务器进程
LOG(LogLevel::ERROR) << "popen fail" << std::endl;
return "popen fail";
}
char buf[256];
std::string result;
while (fgets(buf, sizeof(buf) - 1, fp))
{
result += std::string(buf);
}
pclose(fp); // 必须回收: 否则每次命令泄漏一个 FILE* 和管道 fd
return result;
}
// 发送指令 (循环发送: 结果大于发送缓冲时单次 send 会部分发送, 客户端将拿到截断内容)
void SendCommand(const std::string &result)
{
std::string resp = result.empty() ? "none" : result;
size_t total = resp.size();
size_t sent = 0;
while (sent < total)
{
ssize_t n = send(_sockfd, resp.c_str() + sent, total - sent, 0);
if (n == -1)
{
LOG(LogLevel::ERROR) << "send failed: " << strerror(errno) << std::endl;
close(_sockfd);
_sockfd = -1; // 已关闭, 防止析构时重复 close
return;
}
sent += n;
}
}
private:
// 白名单只读共享: inline static 只在程序启动时构造一次, 无需每个连接重复插入
inline static const std::set<std::string> _command = {
"ls", "pwd", "ls -l", "ll", "touch", "who", "whoami"
};
int _sockfd = -1;
};
2.2.2 源代码
A.TcpClient.cc
#include "Common.hpp"
#include "TcpServer.hpp"
#include "myLog.hpp"
#include "InetAddr.hpp"
#include <signal.h>
using namespace LogModule;
// 服务端接收缓冲 1024, 回显时前缀 "echo# " 占 6 字节,
// 超过该长度的消息回声会被截断, 故客户端发送前先对齐限制
constexpr size_t kMaxMsgSize = 1024 - 1 - 6;
int main(int argc, char *argv[])
{
if (argc != 3)
{
std::cerr << "Usage: " << argv[0] << " " << "ip" << " " << "port" << std::endl;
exit(USAGE_ERR);
}
// 0.忽略SIGPIPE: 对端关闭后继续写会让进程默认退出, 改为由read/write正常返回错误
signal(SIGPIPE, SIG_IGN);
// 1.创建socket文件
int sfd = socket(AF_INET, SOCK_STREAM, 0);
if (sfd == -1)
{
LOG(LogLevel::ERROR) << "socket failed" << std::endl;
return 1;
}
// 2.构建服务端地址
InetAddr dest_addr(argv[1], argv[2]);
// 3.建立连接: TCP面向连接, 失败必须在此暴露(如服务端未启动), 而不是继续执行
if (connect(sfd, dest_addr.GetAddrPtr(), dest_addr.GetLen()) == -1)
{
LOG(LogLevel::ERROR) << "connect failed: " << strerror(errno) << std::endl;
close(sfd);
return 1;
}
LOG(LogLevel::INFO) << "connect success" << std::endl;
for (;;)
{
// 4.向服务端发送消息
std::string msg;
std::cout << "Please Enter# ";
if (!std::getline(std::cin, msg))
break; // stdin结束(Ctrl+D或管道输入完毕), 退出循环而不是用空消息刷屏
if (msg.size() > kMaxMsgSize)
msg.resize(kMaxMsgSize);
// TCP面向字节流: 使用通用fd接口write(与send在flags=0时等价, 无UDP遗留的目的地址参数)
ssize_t sent = write(sfd, msg.c_str(), msg.size());
if (sent == -1)
{
LOG(LogLevel::ERROR) << "send failed: " << strerror(errno) << std::endl;
break; // 连接可能已断开, 退出并统一收尾
}
// 5.接受服务端的回显消息
char buffer[1024];
ssize_t n = read(sfd, buffer, sizeof(buffer) - 1);
if (n > 0)
{
buffer[n] = 0;
std::cout << buffer << std::endl;
}
else if (n == 0)
{
// 对端正常关闭连接
LOG(LogLevel::WARNING) << "server closed the connection" << std::endl;
break;
}
else
{
LOG(LogLevel::ERROR) << "recv failed: " << strerror(errno) << std::endl;
break;
}
// 注: 阻塞socket下recv不会返回EAGAIN, 原"超时"分支是不可能触发的死代码, 已删除;
// TCP本身不丢消息, 收不到回显只会是对端断开(返回0)或出错(返回-1)
}
close(sfd);
return 0;
}
B. TcpServer.cc
#include "TcpServer.hpp"
#include <memory>
#include <cstdlib>
using namespace TcpServerModule;
int main(int argc, char *argv[])
{
if (argc != 2)
{
std::cerr << "Usage: " << argv[0] << " port" << std::endl;
exit(USAGE_ERR);
}
uint16_t port = std::stoi(argv[1]);
TcpServer server(port);
server.Init();
server.Run(); // 内部为 for(;;) 阻塞, 不会返回
return 0;
}
三、附录
3.1 头文件补充
3.1.1 Common.hpp
#pragma once
#include <iostream>
#include <string.h>
#include <sys/types.h>
#include <sys/socket.h>
#include <netinet/in.h>
#include <arpa/inet.h>
#include <unistd.h>
#include <string>
#include <sys/time.h>
#include <cerrno>
// 错误码
enum ExitCode
{
NormalExit,
USAGE_ERR,
SOCKET_ERR,
BIND_ERR,
LISTEN_ERR,
ACCEPT_ERR
};
class NoCopy
{
public:
NoCopy() = default; // 显式保留默认构造(声明了删除的拷贝构造会抑制隐式默认构造)
NoCopy(const NoCopy &) = delete;
NoCopy &operator=(const NoCopy &) = delete;
~NoCopy() = default;
private:
};
3.1.2 InetAddr.hpp
#pragma once
#include <string> // std::string / std::stoi
#include <cstring> // memset
#include <cstdint> // uint16_t
#include <netinet/in.h> // struct sockaddr_in / htons / INADDR_ANY
#include <arpa/inet.h> // inet_pton / inet_ntop
#include "Common.hpp"
class InetAddr
{
public:
InetAddr() = default;
//点分十进制ip 和 字符串型port
InetAddr(const std::string &ip, const std::string port)
: _ip(ip),
_port(std::stoi(port))
{
memset(&_addr, 0, sizeof(_addr));
_addr.sin_family = AF_INET;
_addr.sin_port = htons(_port);
inet_pton(AF_INET, _ip.c_str(), &_addr.sin_addr.s_addr);
}
//点分十进制ip 和 短整型port
InetAddr(const std::string &ip, uint16_t port)
: _ip(ip),
_port(port)
{
memset(&_addr, 0, sizeof(_addr));
_addr.sin_family = AF_INET;
_addr.sin_port = htons(_port);
inet_pton(AF_INET, _ip.c_str(), &_addr.sin_addr.s_addr);
}
//默认ip "0.0.0.0" 监听所有窗口 和 短整型port
InetAddr(uint16_t port)
: _ip("0.0.0.0"),
_port(port)
{
memset(&_addr, 0, sizeof(_addr));
_addr.sin_family = AF_INET;
_addr.sin_port = htons(port);
_addr.sin_addr.s_addr = INADDR_ANY;
}
// 网络套接字 -> 本地套接字
InetAddr(const struct sockaddr_in &addr)
: _addr(addr)
{
char buffer[64];
_ip = inet_ntop(AF_INET, &_addr.sin_addr.s_addr, buffer, sizeof(buffer));
_port = ntohs(_addr.sin_port);
}
std::string GetIp() const
{
return _ip;
}
uint16_t GetPort() const
{
return _port;
}
socklen_t GetLen() const
{
return sizeof(_addr);
}
std::string GetStringAddr() const
{
return _ip + ":" + std::to_string(_port);
}
const struct sockaddr_in &GetAddr() const
{
return _addr;
}
const struct sockaddr *GetAddrPtr() const
{
return (const struct sockaddr *)&_addr;
}
private:
struct sockaddr_in _addr;
std::string _ip;
uint16_t _port;
};
3.1.3 myLog.hpp
#pragma once
#include <sys/types.h>
#include <unistd.h>
#include <fstream>
#include <sstream>
#include <memory>
#include <filesystem>
#include "myMutex.hpp"
namespace LogModule
{
// 默认路径
const std::string defaultpath = "./log/";
const std::string defaultname = "log.txt";
// 获取当前时间
static std::string GetCurrTime()
{
time_t now = time(nullptr); // 获取当前时间
struct tm t;
localtime_r(&now, &t); // 线程安全,必须 _r
char buf[64] = {0};
strftime(buf, sizeof(buf), "%Y-%m-%d %H:%M:%S", &t);
return buf;
}
// 日志等级
enum class LogLevel
{
DEBUG,
INFO,
WARNING,
ERROR,
FATAL
};
// 等级枚举 → 字符串
std::string LogLevelToString(LogLevel level)
{
switch (level)
{
case LogLevel::DEBUG:
return "DEBUG";
case LogLevel::INFO:
return "INFO";
case LogLevel::WARNING:
return "WARNING";
case LogLevel::ERROR:
return "ERROR";
case LogLevel::FATAL:
return "FATAL";
default:
return "UNKNOWN";
}
}
// LevelTag: 流式日志的"参数包"
// LOG(INFO) << "x"; 中, LOG 宏把 等级/文件/行号 打包后经 << 传给 Logger
struct LevelTag
{
LogLevel level;
const char *file;
int line;
};
// LogStrategy: 抽象的日志刷新策略
// 通过多态实现: a.在显示器打印 b.向指定文件写入
class LogStrategy
{
public:
LogStrategy() = default;
virtual ~LogStrategy() = default; // 让编译器合成默认析构,不抑制移动语义
virtual void SyncLog(const std::string &message) = 0; // 纯虚函数
private:
};
// ConsoleLogStrategy:将日志输出到显示器
class ConsoleLogStrategy : public LogStrategy
{
public:
ConsoleLogStrategy() = default;
void SyncLog(const std::string &message) override
{
{
MutexModule::LockGuard lg(_mutex);
std::cout << message << std::endl;
}
}
private:
MutexModule::Mutex _mutex;
};
// FileLogStrategy : 将日志信息输出到文件
class FileLogStrategy : public LogStrategy
{
public:
FileLogStrategy(const std::string &path = defaultpath, // 默认参数已覆盖无参调用
const std::string &name = defaultname)
: _fullpath(path + name)
{
std::filesystem::create_directories(path);
_ofs.open(_fullpath, std::ios::app);
}
void SyncLog(const std::string &message) override
{
{
MutexModule::LockGuard lg(_mutex);
_ofs << message << '\n';
}
}
~FileLogStrategy()
{
if (_ofs.is_open())
_ofs.close();
}
private:
std::string _fullpath;
std::ofstream _ofs;
MutexModule::Mutex _mutex;
};
class Message; // 前置声明, 完整定义见 Logger 之后 (它需要调用 Logger 的公开接口)
class Logger
{
public:
static Logger &GetInstance() // 单例入口
{
static Logger instance;
return instance;
}
void Log(LogLevel level, const char *file, int line, std::string message); // 定义见 Message 之后
// 判断某等级是否会被记录 (供 Message 构造时提前短路)
bool ShouldLog(LogLevel level) const { return level >= _logLevel; }
// 统一落盘出口: 只做策略输出 (格式化由 Message 完成)
void SyncText(const std::string &line); // 定义见 Message 之后
void SetLogLevel(LogLevel level) { _logLevel = level; }
void SetStrategy(std::unique_ptr<LogStrategy> strategy) { _strategy = std::move(strategy); }
// 流式入口: LOG(INFO) << "x" << 43;
// 返回临时 Message, 整个表达式结束后由析构函数统一落盘
Message operator<<(const LevelTag &tag);
private:
Logger()
: _logLevel(LogLevel::DEBUG),
_strategy(std::make_unique<ConsoleLogStrategy>())
{
}
Logger(const Logger &) = delete;
Logger &operator=(const Logger &) = delete;
~Logger() = default;
LogLevel _logLevel; // 日志的等级
std::unique_ptr<LogStrategy> _strategy;
};
// Message: 一条正在构建的日志行 (流式收集 + 格式化 + 落盘 合一)
// 流式用法: LOG(INFO) << a << b; 内容先收集进 _oss,
// 整个表达式结束时临时对象析构 → 格式化并交给策略落盘。
// 函数式用法: Logger::Log() 内部也构造 Message 填入文本, 复用同一通道。
// 只通过 Logger 的公开接口 (ShouldLog / SyncText) 协作, 不触碰内部实现。
// 每条语句一个独立缓冲 → 线程安全, 也不需要写 std::endl。
class Message
{
public:
Message(LogLevel level, const char *file, int line, Logger &logger)
: _level(level),
_file(file),
_line(line),
_pid(getpid()),
_valid(logger.ShouldLog(level)), // 等级过滤: 不达标直接短路
_logger(logger)
{
}
~Message()
{
if (!_valid)
return;
std::string text = _oss.str();
if (text.empty()) // 空语句 LOG(INFO); 什么都不做
return;
_logger.SyncText(Format(text)); // 拼好完整一行, 交给策略落盘
}
template <typename T>
Message &operator<<(const T &value) // 数字/字符串/自定义类型均可用
{
if (_valid)
_oss << value;
return *this;
}
// 兼容 std::endl / std::flush 等流操纵符: 忽略即可, 析构时才真实落盘
Message &operator<<(std::ostream &(*)(std::ostream &))
{
return *this;
}
Message(Message &&) = default; // 允许临时对象移动 (C++17 返回值优化也兜底)
Message(const Message &) = delete; // 拷贝会重复落盘, 禁止
private:
// 拼完整行: [时间] [等级] [pid] [文件] [行号] - 内容
std::string Format(const std::string &content) const
{
std::ostringstream oss;
oss << '[' << GetCurrTime() << "] [" << LogLevelToString(_level) << "] ["
<< _pid << "] [" << _file << "] [" << _line << "] - " << content;
return oss.str();
}
LogLevel _level; // 日志的等级
const char *_file; // 在哪个文件打印
int _line; // 文件的行号
pid_t _pid; // 进程 id
bool _valid; // 等级过滤结果
Logger &_logger; // 落盘时回交给 Logger
std::ostringstream _oss; // 流式内容的缓冲
};
// ---- 以下成员定义必须放在 Message 之后 (需要 Message 的完整类型) ----
// 函数式入口: 构造一个 Message 填入文本, 由 m 析构时落盘 (过滤由 Message 内部短路)
inline void Logger::Log(LogLevel level, const char *file, int line, std::string message)
{
Message m(level, file, line, *this);
m << message;
}
inline void Logger::SyncText(const std::string &line)
{
_strategy->SyncLog(line); // 策略输出
}
// 流式入口: operator<< 按值返回也需要 Message 的完整类型
inline Message Logger::operator<<(const LevelTag &tag)
{
return Message(tag.level, tag.file, tag.line, *this);
}
}
#define LOG(level) \
LogModule::Logger::GetInstance() << LogModule::LevelTag { (level), __FILE__, __LINE__ }
3.1.4 myMutex.hpp
#pragma once
#include <iostream>
#include <pthread.h>
namespace MutexModule
{
class Mutex
{
public:
Mutex()
{
pthread_mutex_init(&_mutex, nullptr);
}
Mutex(const Mutex &) = delete;
Mutex &operator=(const Mutex &) = delete;
void lock()
{
pthread_mutex_lock(&_mutex);
}
void unlock()
{
pthread_mutex_unlock(&_mutex);
}
// 暴露底层句柄,供条件变量等需要原生 pthread_mutex_t* 的场景使用
pthread_mutex_t *native_handle()
{
return &_mutex;
}
~Mutex()
{
pthread_mutex_destroy(&_mutex);
}
private:
pthread_mutex_t _mutex;
};
class LockGuard
{
public:
LockGuard(Mutex &mutex)
: _mutex(mutex)
{
_mutex.lock();
}
LockGuard(const LockGuard &) = delete;
LockGuard &operator=(const LockGuard &) = delete;
~LockGuard()
{
_mutex.unlock();
}
private:
Mutex &_mutex;
};
}
3.1.5 myThread.hpp
#pragma once
#include <iostream>
#include <string>
#include <functional>
#include <pthread.h>
#include <cstring>
#include <cstdlib>
// 基于 C 接口 (pthread) 的 C++ 线程封装
namespace ThreadModule
{
// 已绑参的可调用对象类型: 无参无返回
using func_t = std::function<void()>;
// 内部上下文: 承载可调用对象 + 线程名
// 独立分配在堆上, 是为了把生命周期与 Thread 对象解耦
// (线程可能在 Thread 对象析构后仍在运行, 不能让回调引用悬空)
class ThreadData
{
friend class Thread; // 须带 class: 否则 Thread 尚未声明, 会被当成友元函数声明
public:
ThreadData(const func_t &func, const std::string &name)
: _func(func),
_name(name)
{
}
private:
func_t _func;
std::string _name;
};
// 线程状态: 防止非法操作 (重复 join / join 已 detach 的线程等)
enum class ThreadStatus
{
NEW, // 已构造, 未启动
RUNNING, // 运行中
DETACHED, // 已分离
JOINED // 已回收
};
// 线程自动编号: inline 函数 + static 局部变量
// C++11 下 inline 函数的 static 局部量跨编译单元共享, 头文件安全
inline int NextThreadId()
{
static int id = 1;
return id++;
}
class Thread
{
public:
// 接受任意可调用对象 + 参数, 用 std::bind 绑参后转为 std::function<void()>
template <typename F, typename... Args>
Thread(const std::string &name, F &&f, Args &&...args)
{
_name = name + "-" + std::to_string(NextThreadId());
_status = ThreadStatus::NEW;
_td = new ThreadData(
std::bind(std::forward<F>(f), std::forward<Args>(args)...),
_name);
}
~Thread()
{
// 析构策略 B : 仍在运行就自动 detach, 保证不泄漏
if (_status == ThreadStatus::RUNNING)
{
pthread_detach(_tid); // 不断言, 容错 (线程可能已结束)
}
// 从未启动: 跳板不会运行, 自己清理 _td
if (_status == ThreadStatus::NEW && _td != nullptr)
{
delete _td;
_td = nullptr;
}
// 已启动 (RUNNING/DETACHED/JOINED): _td 由跳板结束时释放, 此处不动
}
// 禁用拷贝: 一个 pthread_t 不能被两个对象管理 (double join / double free)
Thread(const Thread &) = delete;
Thread &operator=(const Thread &) = delete;
// 启动线程 (构造与启动分离, 便于统一管理一批线程)
void start()
{
if (_status != ThreadStatus::NEW)
{
std::cerr << "[Thread:" << _name << "] already started\n";
return;
}
int n = pthread_create(&_tid, nullptr, start_routine, _td);
if (n != 0)
{
std::cerr << "pthread_create error: " << strerror(n) << "\n";
std::abort();
}
_status = ThreadStatus::RUNNING;
}
// 等待回收
void join()
{
if (_status != ThreadStatus::RUNNING)
{
std::cerr << "[Thread:" << _name << "] cannot join (status not RUNNING)\n";
return;
}
int n = pthread_join(_tid, nullptr);
if (n != 0)
{
std::cerr << "pthread_join error: " << strerror(n) << "\n";
std::abort();
}
_status = ThreadStatus::JOINED;
_td = nullptr; // _td 已被跳板释放, 仅置空标记
}
// 分离
void detach()
{
if (_status != ThreadStatus::RUNNING)
{
std::cerr << "[Thread:" << _name << "] cannot detach (status not RUNNING)\n";
return;
}
int n = pthread_detach(_tid);
if (n != 0)
{
std::cerr << "pthread_detach error: " << strerror(n) << "\n";
std::abort();
}
_status = ThreadStatus::DETACHED;
_td = nullptr; // _td 由跳板结束时释放
}
const std::string &name() const { return _name; }
bool isRunning() const { return _status == ThreadStatus::RUNNING; }
private:
// 跳板: 签名必须匹配 pthread 要求的 void*(*)(void*)
// 由于非 static 成员函数隐含 this, 签名不匹配, 因此必须用 static
static void *start_routine(void *arg)
{
ThreadData *td = static_cast<ThreadData *>(arg);
td->_func(); // 执行用户回调
delete td; // 跳板作为最终消费者, 释放上下文
return nullptr;
}
private:
pthread_t _tid;
std::string _name; // 自己保留一份, name() 在 join/detach 后仍可用
ThreadStatus _status;
ThreadData *_td = nullptr;
};
} // namespace ThreadModule
3.1.6 ThreadPool.hpp
#pragma once
#include "myLog.hpp"
#include "myThread.hpp"
#include "myMutex.hpp"
#include "myCond.hpp"
#include <vector>
#include <queue>
#include <memory>
namespace ThreadPoolModule
{
using namespace LogModule;
const int defalutnum = 5;
template <typename T>
class ThreadPool
{
public:
ThreadPool(const int num = defalutnum)
{
for (int i = 0; i < num; i++)
{
_threads.emplace_back(std::make_unique<ThreadModule::Thread>(
"worker",
[this]()
{
HandleTask();
}));
}
}
// 工作线程调用:处理任务队列中的任务
void HandleTask()
{
while (true)
{
T task;
{
MutexModule::LockGuard lg(_mutex);
// 1.用while循环防止虚假唤醒 和 线程池是否退出
while (_task_q.empty() && _isrunning)
{
_cond.Wait(_mutex);
}
// 线程退出
if (_task_q.empty() && !_isrunning)
{
return;
}
// 2.线程获取任务
task = _task_q.front();
_task_q.pop();
}
// 3.线程处理任务
task();
}
}
// 用户线程调用:向任务队列中添加任务
void PushTask(const T &task)
{
{
MutexModule::LockGuard lg(_mutex);
_task_q.push(task);
}
_cond.Signal();
}
// 创建所有线程
void Start()
{
{
MutexModule::LockGuard lg(_mutex);
if (_isrunning)
return;
_isrunning = true;
}
for (const auto &t : _threads)
{
t->start();
}
}
// 广播唤醒所有阻塞的 worker
void Stop()
{
// 注意:需要加锁对_isrunning进行设置
{
MutexModule::LockGuard lg(_mutex);
if (!_isrunning)
return;
_isrunning = false;
}
_cond.Broadcast();
}
void Join()
{
for (auto &t : _threads)
{
if (t->isRunning()) // 只回收运行中的;未启动/已回收的跳过
t->join();
}
}
~ThreadPool()
{
Stop();
Join();
}
private:
std::vector<std::unique_ptr<ThreadModule::Thread>> _threads; // 管理线程
std::queue<T> _task_q;
MutexModule::Mutex _mutex;
CondModule::Cond _cond;
bool _isrunning = false;
};
}
转载自 CSDN-专业IT技术社区
原文链接:https://blog.csdn.net/hii_echo/article/details/166994043




