Skip to content

struct ​

c++
#ifndef STRUCT_H
#define STRUCT_H
#include <QUdpSocket>
#include <QDebug>
#include <QString>
#include <QHash>
#include <QList>
#include <QJsonDocument>
#include <QJsonObject>
#include <QJsonArray>
#include <QObject>
#include <QByteArray>
#include <QObject>
#include <QThread>
#include <QtCore/qglobal.h>

#if defined(PUBSUB_LIB_LIBRARY)
#  define PUBSUB_LIB_EXPORT Q_DECL_EXPORT
#else
#  define PUBSUB_LIB_EXPORT Q_DECL_IMPORT
#endif

class ClientInfo{
public:
    QString ip;
    unsigned short port;
    ClientInfo(){};
    ClientInfo(QString ip_,unsigned short port_){
        ip = ip_;
        port = port_;
    }
    bool operator==(const ClientInfo& other) const{
        if(this->ip == other.ip && this->port == other.port){
            return true;
        }
        return false;
    }

    bool operator!=(const ClientInfo& other) const{
        if(this->ip != other.ip || this->port != other.port){
            return true;
        }
        return false;
    }
};

//主题 里边有订阅者
class Topic{
public:
    Topic(){};
    QString name;
    ClientInfo publisher;
    bool is_publisher_private = false;
    QList<ClientInfo> subscriber_list;
};

enum OrderType{
    PublishTopic = 1,
    SubscribeTopic,
    UnPublishTopic,
    UnSubscribeTopic,
    PublishDataToTopic,
    SetTopicPrivate,
    SetTopicPublic
};

class Order{
public:
    short type;
    QString topic_name;
    QByteArray data;
};


static QByteArray OrderToJson(const Order& order){
    QJsonObject obj;
    obj["data"] = QString(order.data);
    obj["type"] = order.type;
    obj["topic_name"] = order.topic_name;
    QJsonDocument doc(obj);
    return doc.toJson();


};

static Order JsonToOrder(QByteArray order_json){
    QJsonDocument doc = QJsonDocument::fromJson(order_json);
    QJsonObject obj = doc.object();
    Order order;
    order.data = obj["data"].toString().toUtf8();
    order.type = obj["type"].toInt();
    order.topic_name = obj["topic_name"].toString();
    return order;
};
//QByteArray ReplyToJson(const Reply& reply){

//};
//Reply QByteArrayToReply(QByteArray reply_json){

//};
#endif // STRUCT_H

PubCenter ​

c++
#ifndef PUBCENTER_H
#define PUBCENTER_H
#include "struct.h"
class PUBSUB_LIB_EXPORT PubCenter : public QObject
{
    Q_OBJECT
public:
    PubCenter(QString ip_,unsigned short port_);

    bool Bind();

    void PublishTopic(QString topic_name,const ClientInfo& client);
    void SubscribeTopic(QString topic_name,const ClientInfo& client);
    void UnPublishTopic(QString topic_name,const ClientInfo& client);
    void UnSubscribeTopic(QString topic_name,const ClientInfo& client);
    void PublishDataToTopic(QString topic_name, QString data ,const ClientInfo& client);
    void SetTopicPrivate(QString topic_name,const ClientInfo& client);
    void SetTopicPublic(QString topic_name,const ClientInfo& client);

    //处理客户端的指令
    virtual void HandleOrder(QByteArray data,QString ip,unsigned short port) = 0;


private:
    void Send(QString ip, unsigned short port, QByteArray data);
private:

    QString ip;
    unsigned short port;


    QUdpSocket* udp_socket;
    QHash<QString,Topic> topic_list;
};
#endif // PUBCENTER_H
c++
#include "pubcenter.h"



PubCenter::PubCenter(QString ip_,unsigned short port_):
    ip(ip_),port(port_)
{
    qRegisterMetaType<ClientInfo>("ClientInfo");

}

bool PubCenter::Bind()
{
    udp_socket = new QUdpSocket;
    if(!udp_socket->bind(QHostAddress(ip),port)){
        return false;
    }

    QObject::connect(udp_socket,&QUdpSocket::readyRead,this,[=](){

        while(udp_socket->hasPendingDatagrams()){
            QByteArray buffer;
            buffer.resize(udp_socket->pendingDatagramSize());
            QHostAddress client_ip;
            unsigned short client_port;
            udp_socket->readDatagram(buffer.data(),buffer.size(),&client_ip,&client_port);
             qDebug()<<client_ip<<client_port <<" : "<< QString(buffer).toLocal8Bit();
            //收到客户端数据 执行HandleData处理
            HandleOrder(buffer,client_ip.toString(),client_port);
        }
    });

    return true;
}

void PubCenter::Send(QString ip, unsigned short port, QByteArray data)
{
    udp_socket->writeDatagram(data,QHostAddress(ip),port);
}

void PubCenter::PublishTopic(QString topic_name,const ClientInfo& client)
{

    const auto& flag = topic_list.find(topic_name);
    if(flag != topic_list.end()){

        Send(client.ip,client.port,QString("主题%1已经存在 无需发布").arg(topic_name).toUtf8());
        return ;
    }


    Topic topic;
    topic.name = topic_name;
    topic.publisher = client;
    topic.is_publisher_private = false;
    topic_list[topic_name] = topic;
    Send(client.ip,client.port,QString("主题%1发布").arg(topic_name).toUtf8());
}

void PubCenter::SubscribeTopic(QString topic_name,const ClientInfo& client)
{
    const auto& flag = topic_list.find(topic_name);
    if(flag == topic_list.end()){
        Send(client.ip,client.port,QString("主题%1不存在 无法订阅").arg(topic_name).toUtf8());
        return ;
    }

    flag.value().subscriber_list.append(client);
    Send(client.ip,client.port,QString("订阅主题%1").arg(topic_name).toUtf8());
}

void PubCenter::UnPublishTopic(QString topic_name, const ClientInfo& client)
{
    const auto& flag = topic_list.find(topic_name);
    if(flag == topic_list.end()){

        Send(client.ip,client.port,QString("主题%1不存在 无需取消发布").arg(topic_name).toUtf8());
        return ;
    }
    //如果是发布者独有
    if(flag.value().is_publisher_private){
        //如果取消发布的不是该主题的发布者
        if(flag.value().publisher != client){
            Send(client.ip,client.port,QString("主题%1为发布者独有, 您不是该主题的发布者,无法取消发布").arg(topic_name).toUtf8());
            return;
        }

    }

    //删除之前 通知该主题的所有订阅者
    QList<ClientInfo>& subscriber_list = flag.value().subscriber_list;

    for(const auto& subscriber: flag.value().subscriber_list){
        Send(subscriber.ip,subscriber.port,QString("主题%1已经被删除").arg(topic_name).toUtf8());
    }


    topic_list.remove(topic_name);
    Send(client.ip,client.port,QString("删除主题%1").arg(topic_name).toUtf8());
}

void PubCenter::UnSubscribeTopic(QString topic_name,const ClientInfo& client)
{
    const auto& flag = topic_list.find(topic_name);
    if(flag == topic_list.end()){
        Send(client.ip,client.port,QString("主题%1不存在 无需取消订阅").arg(topic_name).toUtf8());
        return ;
    }

    //如果没有订阅该主题 则不需要取消订阅
    //如果删除不成功 说明client没有订阅
    if(!flag.value().subscriber_list.removeOne(client)){
        Send(client.ip,client.port,QString("您尚未订阅主题%1,不能取消订阅").arg(topic_name).toUtf8());
        return;
    }

}

void PubCenter::PublishDataToTopic(QString topic_name, QString data ,const ClientInfo& client)
{
    const auto& flag = topic_list.find(topic_name);
    if(flag == topic_list.end()){
        Send(client.ip,client.port,QString("主题%1不存在 无法推送消息").arg(topic_name).toUtf8());
        return ;
    }
    //可以设置为只允许发布者往该主题推送消息
    if(flag.value().is_publisher_private){//是发布者独有
        if(flag.value().publisher != client){//不是发布者
            Send(client.ip,client.port,QString("您没有权限往主题%1推送消息").arg(topic_name).toUtf8());
            return;
        }
    }

    //逐个给订阅者推送消息即可
    for(const auto& subscriber : flag.value().subscriber_list){
        Send(subscriber.ip,subscriber.port, data.toUtf8());
    }
    Send(client.ip,client.port,QString("向主题%1订阅者推送%2").arg(topic_name).arg(data).toUtf8());

}

void PubCenter::SetTopicPrivate(QString topic_name, const ClientInfo &client)
{
    const auto& flag = topic_list.find(topic_name);
    if(flag == topic_list.end()){
        Send(client.ip,client.port,QString("主题%1不存在 设置为私有").arg(topic_name).toUtf8());
        return ;
    }
    //是主题发布者
    if(client == flag.value().publisher){
        flag.value().is_publisher_private = true;
        Send(client.ip,client.port,QString("主题%1设置为私有").arg(topic_name).toUtf8());
        return;
    }
}

void PubCenter::SetTopicPublic(QString topic_name, const ClientInfo &client)
{
    const auto& flag = topic_list.find(topic_name);
    if(flag == topic_list.end()){

        Send(client.ip,client.port,QString("主题%1设置为公有").arg(topic_name).toUtf8());
        return ;
    }

    //是主题发布者
    if(client == flag.value().publisher){
        flag.value().is_publisher_private = false;
    }
}

Client ​

c++
#ifndef CLIENT_H
#define CLIENT_H
#include "struct.h"

class PUBSUB_LIB_EXPORT Client : public QObject
{
    Q_OBJECT
public:
    Client(QString ip_,unsigned short port_,
         QString des_ip_,unsigned short des_port_);


    void PublishTopic(QString topic_name);
    void SubscribeTopic(QString topic_name);
    void UnPublishTopic(QString topic_name);
    void UnSubscribeTopic(QString topic_name);
    void PublishDataToTopic(QString topic_name, QString data);
    void SetTopicPrivate(QString topic_name);
    void SetTopicPublic(QString topic_name);

signals:
    void StartThread();
private:

    void Send(short type,QString topic_name,QByteArray data);
private:
    QString ip;
    unsigned short port;
    QString des_ip;
    unsigned short des_port;

    QThread* thread;
    QUdpSocket* udp_socket;
};

#endif // CLIENT_H
c++
#include "client.h"


Client::Client(QString ip_, unsigned short port_, QString des_ip_, unsigned short des_port_):
    ip(ip_),port(port_),des_ip(des_ip_),des_port(des_port_)
{

//    thread = new QThread;
//    moveToThread(thread);
//    thread->start();

    udp_socket = new QUdpSocket;
    qDebug()<< udp_socket->bind(QHostAddress(ip),port);
    QObject::connect(udp_socket,&QUdpSocket::readyRead,this,[=](){
        qDebug()<<"asd\n";
        while(udp_socket->hasPendingDatagrams()){
            QByteArray buffer;
            buffer.resize(udp_socket->pendingDatagramSize());
            QHostAddress client_ip;
            unsigned short client_port;
            udp_socket->readDatagram(buffer.data(),buffer.size(),&client_ip,&client_port);

            //收到udp报文
             qDebug()<<client_ip<<client_port <<" : "<< QString(buffer);
        }
    });

}

void Client::PublishTopic(QString topic_name)
{
    Send(OrderType::PublishTopic,topic_name,"");
}

void Client::SubscribeTopic(QString topic_name)
{
    Send(OrderType::SubscribeTopic,topic_name,"");
}

void Client::UnPublishTopic(QString topic_name)
{
    Send(OrderType::UnPublishTopic,topic_name,"");
}

void Client::UnSubscribeTopic(QString topic_name)
{
    Send(OrderType::UnSubscribeTopic,topic_name,"");
}

void Client::PublishDataToTopic(QString topic_name, QString data)
{
    Send(OrderType::PublishDataToTopic,topic_name, data.toUtf8());
}

void Client::SetTopicPrivate(QString topic_name)
{
    Send(OrderType::SetTopicPrivate,topic_name,"");
}

void Client::SetTopicPublic(QString topic_name)
{
    Send(OrderType::SetTopicPublic,topic_name,"");
}

void Client::Send(short type,QString topic_name,QByteArray data)
{
    Order order;
    order.type = type;
    order.topic_name = topic_name;
    order.data = data;


    udp_socket->writeDatagram(OrderToJson(order),QHostAddress(des_ip),des_port);
}

学 习 记 录