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_HPubCenter
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_Hc++
#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_Hc++
#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);
}