Skip to content

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);

}

学 习 记 录