优化了UDP通信代码,增加了消息的时间戳

This commit is contained in:
zjk
2024-11-21 18:13:26 +08:00
parent eb2877f03f
commit abd91d0e88
6 changed files with 219 additions and 48 deletions
+58 -13
View File
@@ -83,19 +83,16 @@ bool PowerManger::OnConnectToServer()
bool PowerManger::Iterate()
{
AppCastingMOOSApp::Iterate();
// Do your thing here!
if(!m_udpComm.m_qReceiveCcuStateBuffer.empty())
{
m_CcuCurrentState = m_udpComm.m_qReceiveCcuStateBuffer.front();
cout << m_CcuCurrentState.checkCode << endl;
m_udpComm.m_qReceiveCcuStateBuffer.pop();
}
// 更新当前状态
updatePmState();
updateTime();
//输出错误信息
if(!m_udpComm.sError.empty())
{
for(int i=0;i<m_udpComm.sError.size();i++)
{
cout<<MOOS::ConsoleColours::red()<< m_udpComm.sError[i] << MOOS::ConsoleColours::reset() << endl;
cout<<MOOS::ConsoleColours::red() << m_udpComm.sError[i] << MOOS::ConsoleColours::reset() << endl;
}
m_udpComm.sError.clear();
}
@@ -223,11 +220,59 @@ bool PowerManger::ListenLoopUDP()
bool PowerManger::updatePmState()
{
m_pmCurrentState.ccustate = m_CcuCurrentState.data;
m_pmCurrentState.bus1State = m_disHighVolBusCurrentState.data;
m_pmCurrentState.bus2State = m_disHighAVolBusCurrentState.data;
m_pmCurrentState.bus3State = m_disHighBVolBusCurrentState.data;
m_pmCurrentState.bus4State = m_disLowBusState.data;
if(!m_udpComm.m_qReceiveCcuStateBuffer.empty())
{
double time;
m_udpComm.getCcuStateMsgFromQueue(time,m_CcuCurrentState);
m_pmCurrentState.ccustate = m_CcuCurrentState.data;
double df;
IsSkewed(time, &df);
cout << "CcuState Skew:" << df << endl;
}
if(!m_udpComm.m_qReceiveDisHighVolBusBuffer.empty())
{
double time;
m_udpComm.getDisSysHVBusFbMsgFromQueue(time,m_disHighVolBusCurrentState);
m_pmCurrentState.bus1State = m_disHighVolBusCurrentState.data;
}
if(!m_udpComm.m_qReceiveDisHighAVolBusBuffer.empty())
{
double time;
m_udpComm.getDisSysHVABusFbMsgFromQueue(time,m_disHighAVolBusCurrentState);
m_pmCurrentState.bus2State = m_disHighAVolBusCurrentState.data;
}
if(!m_udpComm.m_qReceiveDisHighBVolBusBuffer.empty())
{
double time;
m_udpComm.getDisSysHVBBusFbMsgFromQueue(time,m_disHighBVolBusCurrentState);
m_pmCurrentState.bus3State = m_disHighBVolBusCurrentState.data;
}
if(!m_udpComm.m_qReceiveDisLowBusBuffer.empty())
{
double time;
m_udpComm.getDisSysLVBusFbMsgFromQueue(time,m_disLowBusState);
m_pmCurrentState.bus4State = m_disLowBusState.data;
}
return true;
}
bool PowerManger::IsSkewed(double dfTimeNow, double * pdfSkew)
{
double dfSkew = fabs(dfTimeNow - m_current_time);
if(pdfSkew != NULL)
{
*pdfSkew = dfSkew;
}
return (dfSkew > SKEW_TOLERANCE) ? true : false;
}
bool PowerManger::buildPowerSysReport()
{
//TODO: 构建能源系统报告
// string suffix;
// Json::Value json;
return true;
}
+6 -2
View File
@@ -23,6 +23,7 @@
#include "MOOS/libMOOS/Thirdparty/AppCasting/AppCastingMOOSApp.h"
#include "MOOS/libMOOS/Utils/MOOSThread.h"
#include "MOOS/libMOOS/Comms/XPCUdpSocket.h"
#include "MOOS/libMOOS/Utils/MOOSUtilityFunctions.h"
#include "ccuUdpMsg.h"
#include "json/json.h"
#include "udpComm.h"
@@ -30,7 +31,7 @@
#define IPORT 5000
#define MAX_UDP_PKT_SIZE 65535
#define SKEW_TOLERANCE 5
class PowerManger : public AppCastingMOOSApp
{
public:
@@ -52,8 +53,10 @@ class PowerManger : public AppCastingMOOSApp
XPCUdpSocket* setUdpScoket(long lPort);
bool updatePmState();
bool updateTime(){m_current_time = MOOSTime();}
bool IsSkewed(double dfTimeNow, double * pdfSkew = NULL);
bool buildEquipmentReport();
bool buildPowerSysReport();
bool buildDisributionReport();
bool buildSubDistributionReport();
bool buildPredictionReport();
@@ -73,6 +76,7 @@ class PowerManger : public AppCastingMOOSApp
long m_lPort;
unsigned int m_nReceiveBufferSizeKB;
unsigned int m_nSendBufferSizeKB;
double m_current_time;
// CMOOSThread m_SendThread;
private: // State variables
+129 -20
View File
@@ -34,7 +34,6 @@ bool udpComm::getCcuState(char *m, msg_CcuStateFbMsg &s) {
Checksum sum;
msg_CcuStateFbMsg *p;
p = reinterpret_cast<msg_CcuStateFbMsg *>(m);
h = p->header;
// 判断消息是否有效
@@ -209,11 +208,15 @@ bool udpComm::pushMsgToQueue(unsigned char *m)
switch (id)
{
case CCU_STATE_FB_ID: //CCU_STATE_FB 0x0051
msg_CcuStateFbMsg s;
{ msg_CcuStateFbMsg s;
if(getCcuState((char *)Buff,s))
{
if(m_qReceiveCcuStateBuffer.size()<QUENUE_SIZE)
m_qReceiveCcuStateBuffer.push(s);
{
// m_qReceiveCcuStateBuffer.push(s);
msgWithTime* ps = stampTimeStamp(&s);
m_qReceiveCcuStateBuffer.push(*ps);
}
else
{
sError.push_back("pushMsgToQueue faile,CcuState buffer is full");
@@ -226,12 +229,17 @@ bool udpComm::pushMsgToQueue(unsigned char *m)
return false;
}
break;
}
case CCU_SET_PARM_FB_ID: //CCU_SET_PARM_FB 0x0040
msg_CcuSetParmFbMsg p;
if(getCcuSetParmFb((char *)Buff,p))
{
msg_CcuSetParmFbMsg s;
if(getCcuSetParmFb((char *)Buff,s))
{
if(m_qReceiveCcuSetParmBuffer.size()<QUENUE_SIZE)
m_qReceiveCcuSetParmBuffer.push(p);
{
msgWithTime* ps = stampTimeStamp(&s);
m_qReceiveCcuSetParmBuffer.push(*ps);
}
else
{
sError.push_back("pushMsgToQueue faile,CcuSetParm buffer is full");
@@ -243,12 +251,17 @@ bool udpComm::pushMsgToQueue(unsigned char *m)
return false;
}
break;
}
case DIS_HVBUS_FB_ID: //DIS_HVBUS_FB 0x0020
msg_disHighVolBusFbMsg h;
if(getDisSysHVBusFb((char *)Buff,h))
{
msg_disHighVolBusFbMsg s;
if(getDisSysHVBusFb((char *)Buff,s))
{
if(m_qReceiveDisHighVolBusBuffer.size()<QUENUE_SIZE)
m_qReceiveDisHighVolBusBuffer.push(h);
{
msgWithTime* ps = stampTimeStamp(&s);
m_qReceiveDisHighVolBusBuffer.push(*ps);
}
else
{
sError.push_back("pushMsgToQueue faile,DisSysHVBus buffer is full");
@@ -260,12 +273,17 @@ bool udpComm::pushMsgToQueue(unsigned char *m)
return false;
}
break;
case DIS_HVBUSA_FB_ID: //DIS_HVBUSA_FB 0x0006
msg_disHighAVolBusFbMsg a;
if(getDisSysHVABusFb((char *)Buff,a))
}
case DIS_HVBUSA_FB_ID: //DIS_HVBUSA_FB 0x0006
{
msg_disHighAVolBusFbMsg s;
if(getDisSysHVABusFb((char *)Buff,s))
{
if(m_qReceiveDisHighAVolBusBuffer.size()<QUENUE_SIZE)
m_qReceiveDisHighAVolBusBuffer.push(a);
{
msgWithTime* ps = stampTimeStamp(&s);
m_qReceiveDisHighAVolBusBuffer.push(*ps);
}
else
{
sError.push_back("pushMsgToQueue faile,DisSysHVABus buffer is full");
@@ -277,12 +295,17 @@ bool udpComm::pushMsgToQueue(unsigned char *m)
return false;
}
break;
case DIS_HVBUSB_FB_ID: //DIS_HVBUSB_FB 0x0007
msg_disHighBVolBusFbMsg b;
if(getDisSysHVBBusFb((char *)Buff,b))
}
case DIS_HVBUSB_FB_ID: //DIS_HVBUSB_FB 0x0007
{
msg_disHighBVolBusFbMsg s;
if(getDisSysHVBBusFb((char *)Buff,s))
{
if(m_qReceiveDisHighBVolBusBuffer.size()<QUENUE_SIZE)
m_qReceiveDisHighBVolBusBuffer.push(b);
{
msgWithTime* ps = stampTimeStamp(&s);
m_qReceiveDisHighBVolBusBuffer.push(*ps);
}
else
{
sError.push_back("pushMsgToQueue faile,DisSysHVBBus buffer is full");
@@ -294,12 +317,17 @@ bool udpComm::pushMsgToQueue(unsigned char *m)
return false;
}
break;
}
case DIS_LVBUS_FB_ID: //DIS_LVBUS_FB 0x0008
msg_disLowBusFbMsg l;
if(getDisSysLVBusFb((char *)Buff,l))
{
msg_disLowBusFbMsg s;
if(getDisSysLVBusFb((char *)Buff,s))
{
if(m_qReceiveDisLowBusBuffer.size()<QUENUE_SIZE)
m_qReceiveDisLowBusBuffer.push(l);
{
msgWithTime* ps = stampTimeStamp(&s);
m_qReceiveDisLowBusBuffer.push(*ps);
}
else
{
sError.push_back("pushMsgToQueue faile,DisSysLVBus buffer is full");
@@ -310,6 +338,8 @@ bool udpComm::pushMsgToQueue(unsigned char *m)
{
return false;
}
break;
}
default:
break;
}
@@ -335,4 +365,83 @@ bool udpComm::clearMsgQueue()
while(!m_qReceiveDisLowBusBuffer.empty())
m_qReceiveDisLowBusBuffer.pop();
return true;
}
msgWithTime* udpComm::stampTimeStamp(void* msg)
{
msgWithTime* p;
p->data = msg;
p->timeStamp = MOOSTime();
return p;
}
int udpComm::getMsgFormQueue(queue<msgWithTime> &q,double &time,void* data)
{
if (!q.empty())
{
if(q.front().data!= NULL)
{
data = q.front().data;
time = q.front().timeStamp;
q.pop();
return 0;
}
else
return 2;
}
else
return 3;
}
int udpComm::getCcuStateMsgFromQueue(double &time,msg_CcuStateFbMsg &m)
{
void* p;
int ret;
ret = getMsgFormQueue(m_qReceiveCcuStateBuffer,time,p);
m = *reinterpret_cast<msg_CcuStateFbMsg *>(p);
return ret;
}
int udpComm::getCcuSetParmMsgFromQueue(double &time,msg_CcuSetParmFbMsg &m)
{
void* p;
int ret;
ret = getMsgFormQueue(m_qReceiveCcuSetParmBuffer,time,p);
m = *reinterpret_cast<msg_CcuSetParmFbMsg *>(p);
return ret;
}
int udpComm::getDisSysHVBusFbMsgFromQueue(double &time,msg_disHighVolBusFbMsg &m)
{
void* p;
int ret;
ret = getMsgFormQueue(m_qReceiveDisHighVolBusBuffer,time,p);
m = *reinterpret_cast<msg_disHighVolBusFbMsg *>(p);
return ret;
}
int udpComm::getDisSysHVABusFbMsgFromQueue(double &time,msg_disHighAVolBusFbMsg &m)
{
void* p;
int ret;
ret = getMsgFormQueue(m_qReceiveDisHighAVolBusBuffer,time,p);
m = *reinterpret_cast<msg_disHighAVolBusFbMsg *>(p);
return ret;
}
int udpComm::getDisSysHVBBusFbMsgFromQueue(double &time,msg_disHighBVolBusFbMsg &m)
{
void* p;
int ret;
ret = getMsgFormQueue(m_qReceiveDisHighBVolBusBuffer,time,p);
m = *reinterpret_cast<msg_disHighBVolBusFbMsg *>(p);
return ret;
}
int udpComm::getDisSysLVBusFbMsgFromQueue(double &time,msg_disLowBusFbMsg &m)
{
void* p;
int ret;
ret = getMsgFormQueue(m_qReceiveDisLowBusBuffer,time,p);
m = *reinterpret_cast<msg_disLowBusFbMsg *>(p);
return ret;
}
+24 -8
View File
@@ -5,8 +5,15 @@
#include <vector>
#include <queue>
#include <string.h>
#include "MOOS/libMOOS/Utils/MOOSUtilityFunctions.h"
#include <iostream>
using namespace std;
#define QUENUE_SIZE 100
typedef struct {
double timeStamp;
void* data;
} msgWithTime;
class udpComm {
private:
/* data */
@@ -15,8 +22,15 @@ public:
~udpComm();
bool pushMsgToQueue(unsigned char *m);
int getCcuStateMsgFromQueue(double &time,msg_CcuStateFbMsg &m);
int getCcuSetParmMsgFromQueue(double &time,msg_CcuSetParmFbMsg &m);
int getDisSysHVBusFbMsgFromQueue(double &time,msg_disHighVolBusFbMsg &m);
int getDisSysHVABusFbMsgFromQueue(double &time,msg_disHighAVolBusFbMsg &m);
int getDisSysHVBBusFbMsgFromQueue(double &time,msg_disHighBVolBusFbMsg &m);
int getDisSysLVBusFbMsgFromQueue(double &time,msg_disLowBusFbMsg &m);
bool clearMsgQueue();
bool getMsgId(char *m,unsigned short &id);
bool getCcuState(char *m, msg_CcuStateFbMsg &s);
@@ -30,15 +44,17 @@ public:
bool sendCcuColCmd(const msg_CcuColCmdMsg m);
bool sendCcuSetParmCmd(const msg_CcuSetParmMsg m);
msgWithTime* stampTimeStamp(void* msg);
Checksum calculateChecksum(const unsigned char* data, int size);
int getMsgFormQueue(queue<msgWithTime> &q,double &time,void* m);
std::vector<std::string> sError;
std::queue<msg_CcuStateFbMsg> m_qReceiveCcuStateBuffer;
std::queue<msg_CcuSetParmFbMsg> m_qReceiveCcuSetParmBuffer;
std::queue<msg_disHighVolBusFbMsg> m_qReceiveDisHighVolBusBuffer;
std::queue<msg_disHighAVolBusFbMsg> m_qReceiveDisHighAVolBusBuffer;
std::queue<msg_disHighBVolBusFbMsg> m_qReceiveDisHighBVolBusBuffer;
std::queue<msg_disLowBusFbMsg> m_qReceiveDisLowBusBuffer;
std::queue<msgWithTime> m_qReceiveCcuStateBuffer;
std::queue<msgWithTime> m_qReceiveCcuSetParmBuffer;
std::queue<msgWithTime> m_qReceiveDisHighVolBusBuffer;
std::queue<msgWithTime> m_qReceiveDisHighAVolBusBuffer;
std::queue<msgWithTime> m_qReceiveDisHighBVolBusBuffer;
std::queue<msgWithTime> m_qReceiveDisLowBusBuffer;
private:
msg_CcuColCmdMsg m_ccuColCmd;
+1 -1
View File
@@ -25,7 +25,7 @@ sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
# message = "Hello, this is a test message."
message = valid_msg3
# 发送信息的次数
num_messages = 10
num_messages = 100
# 发送信息的时间间隔(秒)
interval = 0.1
+1 -4
View File
@@ -13,7 +13,7 @@ class msg_CcuStateFbMsg:
self.totalPowerGeneration = 0x00
self.ch4o = 0x00
self.o2 = 0x00
self.fcuVoltage = 0x00
self.fcuVoltage = 0x1230
self.fcuCurrent = 0x00
self.o2Level = 0x00
self.placeholder1 = 0x00
@@ -83,9 +83,6 @@ class msg_CcuStateFbMsg:
self.placeholder5 = 0x00
self.checkCode = 0x00
def pack(self):
# 初始化数据
self.length = 0x0000