ws_socket
c++
#include <boost/beast.hpp>
#include "boost/asio.hpp"
#include "boost/beast/core.hpp"
#include "boost/beast/version.hpp"
#include "boost/beast/websocket.hpp"
//接收数据缓冲区
boost::beast::flat_buffer _buffer;
//升级成的websocket 升级以后 这个会话就用_ws_ptr
std::unique_ptr<websocket::stream<boost::beast::tcp_stream>> _ws_ptr;握手
c++
void Session::HandShake() {
auto self = shared_from_this();
//挂起async_accept 异步等待客户端的握手包
_ws_ptr->async_accept([self](boost::system::error_code err){
std::cout<<"_ws_ptr->async_accept\n";
try {
if (!err) {
emit self->_server->HandShakeSuccess();
//握手成功 把session添加进入map 管理 并且挂起异步读
self->_server->AddSession(self->GetUUid(), self);
self->AsyncRead();
}
else {
emit self->_server->HandShakeFail();
std::cout << "握手失败: " << err.what() << std::endl;
}
} catch (std::exception& e) {
std::cout<<"连接异常: "<<e.what()<<"\n";
return;
}
});
/**
* 如果异步操作可能很久还没有返回回调, 也就是客户端建立了tcp连接 但是一直没有发握手包
* 可以设置一个定时器 在定时器中调用被挂起异步操作对象的cancel()方法 取消其在ioc中挂起的异步操作
* // 设置超时时间为5秒
asio::steady_timer timer(io_context);
timer.expires_after(std::chrono::seconds(5));
timer.async_wait([&_ws_ptr](const boost::system::error_code& error) {
handle_timeout(error, _ws_ptr);
});
*/
}异步读
c++
void Session::AsyncRead(){
auto self = shared_from_this();
_ws_ptr->async_read(_buffer,
[self](beast::error_code err, std::size_t buffer_bytes) {
try {
if (err) {
//客户端断开连接 触发读事件 清理客户端连接
std::cout << "websocket async read error is " << err.what() << std::endl;
self->_server->ClearSession(self->GetUUid());
return;
}
//获取数据
self->_ws_ptr->text(self->_ws_ptr->got_text());
std::string recv_data = boost::beast::buffers_to_string(self->_buffer.data());
//清理缓冲区为下个数据接收做准备
self->_buffer.consume(self->_buffer.size());
std::cout << "websocket receive msg is " << recv_data << std::endl;
emit self->_server->GetMsg(QString::fromStdString(self->_ip),self->_port, QString::fromStdString(recv_data));
//读完上一个数据 重新挂起异步读 继续读
self->AsyncRead();
}
catch (std::exception& exp) {
std::cout << "exception is " << exp.what() << std::endl;
self->_server->ClearSession(self->GetUUid());
}
});
}异步写
c++
void Session::SendCallBack(std::string msg){
auto self = shared_from_this();
_ws_ptr->async_write(asio::buffer(msg.c_str(), msg.length()),
[self](beast::error_code err, std::size_t nsize) {
try {
if (err) {
std::cout << "async send err is " << err.what() << std::endl;
self->_server->ClearSession(self->GetUUid());
return;
}
std::string send_msg;
{
std::lock_guard<std::mutex> lck_gurad(self->_send_lock);
//因为上一条已经写完 就把上一条出队
self->_send_queue.pop();
//出队完为空 说明没数据了 就不挂起异步写
if (self->_send_queue.empty()) {
return;
}
send_msg = self->_send_queue.front();
}
//继续挂起异步写
self->SendCallBack(std::move(send_msg));
}
catch (std::exception& exp) {
std::cout << "async send exception is " << exp.what() << std::endl;
self->_server->ClearSession(self->GetUUid());
}
});
}关闭
c++
void Session::Close()
{
beast::get_lowest_layer(*_ws_ptr.get()).socket().close();
}
void Session::ShutDownReceive()
{
auto& sock = beast::get_lowest_layer(*_ws_ptr.get()).socket();
sock.shutdown(sock.shutdown_receive);
}