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;
};
}
#endifc++
#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_ ;
#endifTcpServer
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;
};
}
// #endifc++
#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();//清理所有连接
}
}