高山有多高头像
关注
【Linux笔记】TCPSocket封面图

【Linux笔记】TCPSocket

一、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

文章来源转载

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--