Skip to content

IOContextPool ​

c++
#include <iostream>
#include <thread>
#include <vector>
#include <memory>
#include <boost/asio.hpp>
#include "SingleClass.hpp"
namespace tcpserver{
class IOContextPool : public Singleton<IOContextPool>
{
    friend class Singleton<IOContextPool>;
private:
    /**
     * 函数名: IOContextPool
     * 参数: threadCount , 设置为当前电脑的线程数量
     * 作用: 初始化io_context池 ,
     *      当连接到来时从池里取出一个io_context,
     *      把新连接的socket放入其中,
     *      调用run()运行io_context事件循环
     *      让其处理其中的socket的挂起的异步操作
     */
    IOContextPool(size_t threadCount = std::thread::hardware_concurrency());
    IOContextPool(const IOContextPool&) = delete;
    IOContextPool& operator=(const IOContextPool&) = delete;
    size_t _threadCount;
    size_t _next;
    std::vector<std::thread> _threads;
    std::vector<asio::io_context> _contexts;
    //用来给io_context注册一个初始事件 防止它一调用run()了就直接退出
    std::vector<std::unique_ptr<asio::io_context::work>> _io_works;
public:
    asio::io_context& GetIoContext();
    ~IOContextPool();
    void Stop();
//    void Reset();
};

}
c++
#include "IOContextPool.h"
namespace tcpserver{
IOContextPool::IOContextPool(size_t threadCount ):_threadCount(threadCount),_next(0),_contexts(threadCount),_io_works(threadCount)
{
    for(size_t i = 0 ; i < threadCount;i++){
        _io_works[i] = std::unique_ptr<asio::io_context::work>(
            new asio::io_context::work(_contexts[i]));

    }
    for(size_t i = 0 ; i < threadCount;i++){
        std::cout<<"IOContext"<<i<<" 启动\n";
        _threads.emplace_back(std::thread([=](){
            _contexts[i].run();
        })); //emplace_back可以插入 然后执行构造函数
//        _threads[i].detach();
    }
}

IOContextPool::~IOContextPool()
{
    for(size_t i = 0 ; i < _threadCount;i++){
        _threads[i].join();
    }
}

void IOContextPool::Stop()
{
    for(size_t i = 0 ; i < _threadCount;i++){
        _contexts[i].stop();
    }
}

asio::io_context& IOContextPool::GetIoContext(){
    auto& io_context = _contexts[_next++];

    if(_next == _threadCount){
        _next = 0;
        
    }
    return io_context;
}

}

LogicSystem ​

c++

#pragma once
#include "MsgNode.h"
#include <queue>
#include <vector>
#include <map>
#include <functional>
#include <memory>
#include <string>
#include <thread>
#include <condition_variable>
#include "SingleClass.hpp"
namespace tcpserver{
class Session; //前置声明
typedef std::function<void(std::shared_ptr<Session> session, 
                           short msg_id, 
                           std::string msg_data)> Callback;

//跑在_worker单线程中,用于处理tcp层已经解析好的数据包
class LogicSystem : public Singleton<LogicSystem>
{
    friend class Singleton<LogicSystem> ;
public:
    /**
     * 函数名: AddToQueue
     * 作用: 把RecvNode添加到_recvNodeQueue
     */
    void AddToQueue(std::shared_ptr<RecvNode> node);
    /**
     * 函数名: Stop
     * 作用: 停止线程
     */
    void Stop();
    void SayHello(std::shared_ptr<Session> session, 
                  short msg_id, 
                  std::string msg_data);

    ~LogicSystem();
private:
    //用于注册逻辑层对数据包的处理逻辑
    std::map<short, Callback> _callbackMap;
    bool _stop;
    std::queue<std::shared_ptr<RecvNode>> _recvNodeQueue;
    std::mutex _mtx;
    std::condition_variable _cv;
    std::thread _worker;
    enum ID_TYPE{
        REPLAY = 1001
    };
    LogicSystem();
    /**
     * 函数名: RegisterCallback
     * 作用: 在其中注册ID为xxx的包对应的操作,如某个包的ID为1001则执行SayHello
     */
    void RegisterCallback();
};
}
c++
#include "LogicSystem.h"
#include "Session.h"
namespace tcpserver{

LogicSystem::LogicSystem():_stop(false)
{
    RegisterCallback();
    _worker = std::thread([=](){
        std::cout<<"逻辑层线程启动"<<"\n";
        while(true){
            std::unique_lock<std::mutex> lock(_mtx);
            while(_recvNodeQueue.empty() && !_stop){
                _cv.wait(lock);
            }
            if(_stop){
                while(!_recvNodeQueue.empty()){
                    std::shared_ptr<RecvNode> node = _recvNodeQueue.front();
                    auto func = _callbackMap.find(node->_id);
                    if(func != _callbackMap.end()){
                        func->second(node->_session,node->_id,node->_msg);
                    }
                    _recvNodeQueue.pop();
                }
                break;
            }

            std::shared_ptr<RecvNode> node = _recvNodeQueue.front();
            auto func = _callbackMap.find(node->_id);
            std::cout<<"id: "<<node->_id<<" msg:"<<node->_msg<<"\n";
            if(func != _callbackMap.end()){
                func->second(node->_session,node->_id,node->_msg);
                
            }
            _recvNodeQueue.pop();
        }
         std::cout<<"逻辑层线程关闭"<<"\n";
    });
}

//在这里注册包ID对应的处理函数
void LogicSystem::RegisterCallback(){
    _callbackMap.insert({ID_TYPE::REPLAY,
                         std::bind(
                         &LogicSystem::SayHello,
                         this,std::placeholders::_1,
                         std::placeholders::_2,std::placeholders::_3)});
}

void LogicSystem::SayHello(std::shared_ptr<Session> session, short msg_id, std::string msg_data){

    std::cout<<"SayHello:"<< msg_data<<" "<<msg_data.size()<<"\n";
    session->AsyncSend(msg_data.c_str(),msg_data.size());
}

void LogicSystem::Stop(){
    _stop = true;
}

LogicSystem::~LogicSystem()
{
    _mtx.unlock();
    _cv.notify_one();
    _worker.join();
}

void LogicSystem::AddToQueue(std::shared_ptr<RecvNode> node){
    
    std::unique_lock<std::mutex> lock(_mtx);
    _recvNodeQueue.push(node);
    if(_recvNodeQueue.size() == 1){
        _cv.notify_one();
    }
}
}

MsgNode ​

c++
#ifndef MsgNode_H
#define MsgNode_H
#define ID_SIZE 2
#define HEAD_SIZE 2
#define MAX_LENGTH  1024
#include <iostream>
#include <string.h> 
#include <memory>
#include <boost/asio.hpp>
namespace tcpserver{
class Session;

class MsgNode
{
public:
    MsgNode();
    ~MsgNode();
};

class SendNode
{
public:
    SendNode(const char* data, size_t data_len);
    ~SendNode(); 
    char* data_;
    size_t total_len_; // 数据包总大小
    size_t cur_len_;   // 当前已发送
};

class RecvNode
{
public:
    
    RecvNode(std::shared_ptr<Session> session, short id, std::string msg);
    RecvNode(std::shared_ptr<Session> session);
    ~RecvNode();
    std::shared_ptr<Session> _session;
    short _id;
    std::string _msg;
};

}
#endif
c++
#include "MsgNode.h"

#include "Session.h"
namespace tcpserver{


MsgNode::MsgNode(){}
MsgNode::~MsgNode(){}
SendNode::SendNode(const char* data,size_t data_len):total_len_(data_len + HEAD_SIZE),cur_len_(0){
    data_ = new char[total_len_];
    int net_data_len = asio::detail::socket_ops::host_to_network_short(data_len);
    memcpy(data_ , &net_data_len, HEAD_SIZE);
    memcpy(data_ + HEAD_SIZE, data, data_len);

}

SendNode::~SendNode()
{
    delete[] data_;
}

RecvNode::RecvNode(std::shared_ptr<Session> session,short id, std::string msg): _session(session),_id(id),_msg(msg){}

RecvNode::RecvNode(std::shared_ptr<Session> session): _session(session){}

RecvNode::~RecvNode(){}
}

Session ​

c++

#pragma once
#include <iostream>
#include <map>
#include <queue>
#include <memory>
#include <mutex>
#include <QObject>
#include "MsgNode.h"
#include <QString>
#include "LogicSystem.h"
namespace tcpserver{
class TcpServer;
// 会话
class Session : public std::enable_shared_from_this<Session>
{

public:
    /**
     * 函数名: Session
     * 参数: ioc , 该Session的socket挂起的异步操作 由这个ioc事件循环处理
     * 参数: server ,  指向服务服务器类的指针
     * 作用: 封装了客户端socket的异步读写、关闭等操作
     */
    Session(asio::io_context &ioc, TcpServer* server);
    ~Session();

    /**
     * 函数名: AsyncSend
     * 参数: buf , 服务端要给客户端连接发送的消息,传指针即可 内部进行拷贝
     * 参数: len , 消息长度
     * 作用: 异步发送接口,数据先入队,再发送,用于服务端给该socket发消息。 异步发送的回调由SendCallBack处理
     *       由于在ioc中挂起异步写后,数据发送的顺序不能确定,得通过发送队列保证数据的发送顺序
     */
    void AsyncSend(const char* buf, int len);


    /**
     * 函数名: AsyncRead
     * 作用: 挂起异步读。读事件触发,读完一次后内部会再次自己调用
     * AsyncRead 给socket注册读事件挂起异步读。所以只需要在连接建立时调用一次
     */
    void AsyncRead();

    /**
     * 函数名: Close
     * 作用: 关闭socket
     */
    void Close();

    asio::ip::tcp::socket &Socket();
    std::string GetUUid();
    int _port;
    std::string _ip;
    std::string _uuid;
private:
    //处理写事件的回调函数
    void HandleWrite(const boost::system::error_code &error, std::shared_ptr<Session> session);
    
    //以下三个回调函数是拆包过程,客户端的包是由 ID+包头+包体构成
    void HandleReadId(const boost::system::error_code &error, size_t bytes_transfered, std::shared_ptr<Session> session);
    void HandleReadHead(const boost::system::error_code &error, size_t bytes_transfered, std::shared_ptr<Session> session);
    void HandleReadData(const boost::system::error_code &error, size_t bytes_transfered, std::shared_ptr<Session> session);
    char _id[ID_SIZE]; // 消息id缓冲区
    char _head[HEAD_SIZE]; //包头缓冲区
    char _data[MAX_LENGTH];//读包体缓冲区
    
    asio::ip::tcp::socket _socket;

    std::queue<std::shared_ptr<SendNode>> _send_queue;//发送队列
    std::mutex _send_lock;//由于可能多线程操作 给发送队列加锁

    TcpServer* _server;
};

}
c++

#include "Session.h"
#include "TcpServer.h"
#include "boost/asio/socket_base.hpp"
namespace tcpserver{
Session::Session(asio::io_context &ioc, TcpServer *server)
    : _socket(ioc)
{
    _server = server;
    //初始化读包缓冲区  
    memset(_id, 0, ID_SIZE);
    memset(_head, 0, HEAD_SIZE);
    memset(_data,0,MAX_LENGTH);
}

Session::~Session()
{
    if(_socket.is_open()){
        _socket.close();
    }
    std::cout << "bye  " << _ip << ":" << _port << "\n";
}

asio::ip::tcp::socket &Session::Socket()
{
    return _socket;
}


//异步发送接口
void Session::AsyncSend(const char *buf, int len)
{
    std::lock_guard<std::mutex> lock(_send_lock);
    // 先入队
    _send_queue.push(std::make_shared<SendNode>(buf, len));
    if (_send_queue.size() == 1)
    {
        std::shared_ptr<SendNode> node = _send_queue.front();
        asio::async_write(_socket,
                                 asio::buffer(node->data_, node->total_len_),
                                 std::bind(&Session::HandleWrite,
                                           this,
                                           std::placeholders::_1,
                                           shared_from_this()));
    }
}

// 对于HandleWrite回调  他负责把发送队列里面的所有node清空
// async_write在底层多次调用async_write_some 如果async_write返回回调函数 说明已经全部写完
void Session::HandleWrite(const boost::system::error_code &error, std::shared_ptr<Session> session)
{
    if (!error)
    {
        std::lock_guard<std::mutex> lock(_send_lock);
        // 走到这一步 说明上一个node被发完了 就出队
        _send_queue.pop();
        if (!session->_send_queue.empty())
        {
            std::shared_ptr<SendNode> node = session->_send_queue.front();

            asio::async_write(_socket,
                                     asio::buffer(node->data_, node->total_len_),
                                     std::bind(&Session::HandleWrite,
                                               this, std::placeholders::_1, session));
        }
    }
    else
    {
        //出问题就清理这个会话,但不是马上就会析构掉这个session,因为回调函数中持指向这个Session对象的智能指针 只有这个session的所有操作都跑完以后才会被析构掉 是安全的
        this->_server->ClearSession(this->GetUUid());
    }
}

std::string Session::GetUUid()
{
    return _uuid;
}

void Session::AsyncRead()
{
    // 挂起异步读
    asio::async_read(_socket, asio::buffer(_id, ID_SIZE),
                            std::bind(&Session::HandleReadId, this,
                                      std::placeholders::_1, std::placeholders::_2, shared_from_this()));
}

void Session::Close()
{
    emit _server->DisConnect(QString::fromStdString(_ip),_port);
    _socket.close();
}


void Session::HandleReadId(const boost::system::error_code &error, size_t bytes_transfered, std::shared_ptr<Session> session)
{
    if (!error)
    {
     
        asio::async_read(_socket, asio::buffer(_head, HEAD_SIZE),
                                std::bind(
                                    &Session::HandleReadHead, this, std::placeholders::_1, std::placeholders::_2, session));
    }
    else
    {
        std::cout<<"HandleReadId ClearSession";
        this->_server->ClearSession(this->GetUUid());
    }
}


void Session::HandleReadHead(const boost::system::error_code &error, size_t bytes_transfered, std::shared_ptr<Session> session)
{
    if (!error)
    {

        //解析包头 读包体
        int data_len = 0;
        memcpy(&data_len, _head, HEAD_SIZE);
        data_len = asio::detail::socket_ops::network_to_host_short(data_len);
     
        asio::async_read(_socket, asio::buffer(_data, data_len),
                                std::bind(
                                    &Session::HandleReadData, this, std::placeholders::_1, std::placeholders::_2, session));
    }
    else
    {
        std::cout<<"HandleReadHead ClearSession";
        this->_server->ClearSession(this->GetUUid());
    }
}

void Session::HandleReadData(const boost::system::error_code &error, size_t bytes_transfered, std::shared_ptr<Session> session)
{
    if (!error)
    {
        // std::cout << session->_ip << ":" << session->_port << ":  " << _data << std::endl;
        //走到这里  一个包的头已经读到_head里  包体已经读到_data里
        short id = 0;
        memcpy(&id, _id, ID_SIZE);
        // 网络字节序转化为本地字节序
        id = asio::detail::socket_ops::network_to_host_short(id);
        emit _server->GetMsg(QString::fromStdString(_ip),_port,QString::fromStdString(std::string(_data,bytes_transfered)));
//        emit GetMsg(QString::fromStdString(_ip),_port,QString::fromStdString(std::string(_data,bytes_transfered)));
        // //封装成逻辑结点 传给逻辑层处理
        LogicSystem::getInstance()->AddToQueue(
            std::make_shared<RecvNode>(shared_from_this(),id,std::string(_data,bytes_transfered))
            );

        //清空一下读包用的缓冲区 继续挂起读id读下一个包
        memset(_id, 0, ID_SIZE);
        memset(_head, 0, HEAD_SIZE);
        memset(_data, 0, bytes_transfered);

        asio::async_read(_socket, asio::buffer(_id, ID_SIZE),
                                std::bind(
                                    &Session::HandleReadId, this, std::placeholders::_1, std::placeholders::_2, session));
    }
    else
    {
        //出问题就移除该链接
        std::cout<<"HandleReadData ClearSession";
        this->_server->ClearSession(this->GetUUid());
    }
}

}

Singleton ​

c++

#ifndef SINGLECLASS_HPP 
#define SINGLECLASS_HPP

#include <iostream>
#include <memory>
#include <mutex>


//制作一个单例类T,需要
/*
    1.继承Singleton<T>
    2.把Singleton<T>作为T的友元
    3.禁用T的拷贝构造,赋值重载运算符
*/

template <typename T>
class Singleton
{
protected:
    // 使用默认构造
    Singleton() = default;
    // 删除拷贝构造
    Singleton(const Singleton<T> &) = delete;
    // 禁用 重载赋值运算符
    Singleton<T> &operator=(const Singleton<T> &st) = delete;
    // 静态实例
    static std::shared_ptr<T> m_instance;

    // 加锁
    static std::mutex instance_mtx_;

public:
    static std::shared_ptr<T> getInstance()
    {
        //获取单例之前先加锁
        instance_mtx_.lock();

        if (m_instance == nullptr)
        {

            m_instance = std::shared_ptr<T>(new T());
        }
        instance_mtx_.unlock();
        return m_instance;
    }
};

// 初始化静态成员
template <typename T>
std::shared_ptr<T> Singleton<T>::m_instance = nullptr;

template <typename T>
std::mutex Singleton<T>::instance_mtx_ ;
#endif

TcpServer ​

c++
// #ifndef TcpServer_H
// #define TcpServer_H
#pragma once
#include <iostream>
#include <map>
#include <queue>
#include <memory>
#include <mutex>
#include "MsgNode.h"
#include <boost/asio.hpp>
#include "IOContextPool.h"
#include "LogicSystem.h"
#include <QObject>
#include <QString>

namespace tcpserver{
class Session;

class TcpServer : public QObject
{
    Q_OBJECT
public:
    /**
     * 函数名: TcpServer
     * 参数: ioc , asio上下文对象 在某个线程中调用run()启动。
     * 参数: port , 绑定本地端口
     * 作用: 主要用于初始化_acceptor
     */
    TcpServer(asio::io_context &ioc, short port);
    ~TcpServer();

    /**
     * 函数名: Bind
     * 作用: 绑定在构造函数中传入的port
     */
    void Bind();

    /**
     * 函数名: StartAccept
     * 作用: 在io_context中挂起异步accept
     */
    void StartAccept();
    void HandleAccept(std::shared_ptr<Session> new_session, const boost::system::error_code &error);

    /**
     * 函数名: ClearSession
     * 作用: 清理连接
     */
    void ClearSession(std::string uuid);

    /**
     * 函数名: Close
     * 作用: 关闭服务器
     *    1.停止接收新连接
     *    2.遍历_sessions,处理完毕每个session的读写事件后调用close关闭连接
     *    3.清空_sessions 释放所有连接对象
     */
    void Close();

signals:
    void Connect(QString ip, int port);//新tcp连接成功信号
    void DisConnect(QString ip, int port);//连接断开信号
    void GetMsg(QString ip, int port, QString msg);//收到消息信号

public:
    asio::io_context &_ioc;
    asio::ip::tcp::acceptor _acceptor;
    std::map<std::string, std::shared_ptr<Session>>  _sessions;

    short _local_port;
};

}
// #endif
c++

#include "TcpServer.h"
#include <QDebug>
#include "Session.h"
namespace tcpserver{
TcpServer::TcpServer(asio::io_context &ioc, short port) : _ioc(ioc),
    _acceptor(ioc, asio::ip::tcp::endpoint(asio::ip::tcp::v4(), port)),
    _local_port(port)
{
    _acceptor.set_option(asio::ip::tcp::acceptor::reuse_address(true));
}

TcpServer::~TcpServer()
{

    std::cout<<"服务被销毁\n";

}

void TcpServer::Bind()
{
    if(!_acceptor.is_open()){
        boost::system::error_code code;
        _acceptor.open(asio::ip::tcp::v4(),code);
        std::cout<<"open: "<<code.what()<<'\n';
        _acceptor.bind(asio::ip::tcp::endpoint(asio::ip::tcp::v4(),_local_port),code);
        std::cout<<"bind: "<<code.what()<<'\n';
        std::cout << "开始监听"<< "\n";

    }else{
        StartAccept();
    }


}

void TcpServer::StartAccept()
{

    if(_acceptor.is_open()){
        //从ioc池里取一个ioc上下文 跟该session关联起来
        std::shared_ptr<Session> new_session = std::make_shared<Session>(IOContextPool::getInstance()->GetIoContext(), this);

        //因为这个acceptor给io_context注册了一个读事件 监听对端连接 所以io_context.run()不会退出
        _acceptor.async_accept(new_session->Socket(),
                               std::bind(&TcpServer::HandleAccept, this, new_session, std::placeholders::_1));
//        _acceptor.async_accept([](){
//            qDebug()<<"asdsad";
//        });
    }
}

void TcpServer::ClearSession(std::string uuid){
    const auto& session = _sessions.find(uuid);
    if(session != _sessions.end()){
        emit DisConnect(QString::fromStdString(session->second->_ip), session->second->_port);
        _sessions.erase(uuid);

    }
}


//执行到HandleAccept 时 说明async_accept已经处理完毕
void TcpServer::HandleAccept(std::shared_ptr<Session> new_session, const boost::system::error_code &error)
{
    if (!error)
    {
        //开始异步读这个socket
        new_session->_ip = new_session->Socket().remote_endpoint().address().to_string();
        new_session->_port = new_session->Socket().remote_endpoint().port();
        std::cout<<new_session->_ip<<":"<<new_session->_port<<"连接 \n";

        //设置连接id
        new_session->_uuid = new_session->_ip+std::to_string(new_session->_port);

        emit Connect(QString::fromStdString(new_session->_ip), new_session->_port);
        //设置客户端端口复用
        new_session->Socket().set_option(asio::socket_base::reuse_address(true));
        new_session->AsyncRead();
        //把会话放到map里管理
        this->_sessions.insert(std::make_pair(new_session->GetUUid(), new_session));
    }
    else
    {
        std::cout <<"HandleAccept: "<< error.what() << "\n";
    }
    //继续监听
    StartAccept();
}

void TcpServer::Close()
{
    _acceptor.close();//停止接收新连接
    for(const auto& session: _sessions){
         //断开连接
         session.second->Close();
    }
    _sessions.clear();//清理所有连接

}
}

学 习 记 录