#include "DataStream.h" DataStream::DataStream(QObject *parent) : QObject(parent) { running_flag = false; thread = new QThread(); this->moveToThread(thread); thread->setPriority(QThread::NormalPriority); connect(thread, &QThread::started, this, &DataStream::process); setRunFrq(80); initbuff(); start(); Parser2_init(&parser); Packer2_init(&packer); } DataStream::~DataStream() { stop(); if(thread) { thread->deleteLater(); } } void DataStream::setRunFrq(qreal frq) { if((frq != 0)||(frq <= 1000)) { running_frq = frq; qDebug() << "set running frquency:" <isRunning()) { running_flag = true; thread->start(); qDebug() << "thread start" << thread->isRunning(); } else { qDebug() << "thread has started"; } } void DataStream::stop() { if(thread) { if(!thread->isFinished()) { running_flag = false; thread->requestInterruption(); thread->setPriority(QThread::HighPriority);//设置成最高,让线程优先退出 bool flag = disconnect(thread, nullptr, nullptr, nullptr);//没有完全断开 qDebug() << "disconnect" << flag << "isInterruptionRequested" << thread->isInterruptionRequested(); thread->quit(); thread->wait(1000);//这个地方有问题、没办法退出线程,断开之前,先断开串口连接,这样就能断开线程,否则一直在占用 } else { qDebug() << "thread is not running"; } } } void DataStream::process()//线程函数 { uint8_t count = 0; while (true) { count ++; QThread::msleep(1000.0/frq()); switch(status.m_Mode) { default: case Nop_Mode : QThread::yieldCurrentThread();break; case RecieveMode : break; case TransmitMode : sendAFrame();status.m_Mode = Nop_Mode;break; } if(isInterruptionRequested())//退出 { break; } } } void DataStream::setGCSID(int m_sysid, int m_compid) { GCS_SysID = m_sysid; GCS_CompID = m_compid; } void DataStream::setID(int m_sysid,int m_compid) { sysid = (uint8_t)m_sysid; compid = (uint8_t)m_compid; } void DataStream::Send(mavlink_message_t msg) { if(isActive()) { uint8_t buff[MAVLINK_MAX_PACKET_LEN+sizeof(quint64)]; uint16_t len = mavlink_msg_to_send_buffer(buff, &msg); if(plogFileName.size() > 0) { setLogData(plogFileName,1,msg); } emit SendMessageTo(0,buff, len);//使用信号和槽 } } bool DataStream::setLogData(QString name,uint8_t type,mavlink_message_t msg) { mutex.lock(); QFile file(name); if(file.open(QIODevice::Append)) { uint8_t buff[MAVLINK_MAX_PACKET_LEN+sizeof(quint64)]; qToLittleEndian(type, buff); quint64 currentTimestamp = (quint64)QDateTime::currentMSecsSinceEpoch(); qToLittleEndian(currentTimestamp, buff+1); uint16_t len = mavlink_msg_to_send_buffer(buff+1+sizeof(quint64), &msg); auto size = Packer2_pack(&packer,1,buff,len+1+sizeof(quint64)); QByteArray data; data.clear(); data.append((const char *)packer.buff,size); QDataStream stream(&file); stream << data; /* if (file.write((const char *)packer.buff, size) == size) { } */ file.flush(); file.close(); mutex.unlock(); return true; } mutex.unlock(); return false; } void DataStream::initbuff(void) { client_buff.max_size = 10 * 1024 *1024;//10M client_buff.buff[0].clear(); client_buff.buff[1].clear(); client_buff.select = 0; serial_buff.max_size = 10 * 1024 *1024; serial_buff.buff[0].clear(); serial_buff.buff[1].clear(); serial_buff.select = 0; } void DataStream::setbuff(quint32 src,QByteArray data) { switch (src) { default: case SourceType::c_sock: //当前的buff超过10M字节之后就清除,防爆机制,不然buff太大后容易卡死 if(client_buff.buff[client_buff.select].size() >= client_buff.max_size) { client_buff.buff[client_buff.select].clear(); qDebug() << "client_buff.buff " << client_buff.select <<" OVERFLOW"; } client_buff.buff[client_buff.select].append(data); break; case SourceType::s_port: //当前的buff超过10M字节之后就清除,防爆机制,不然buff太大后容易卡死 if(serial_buff.buff[serial_buff.select].size() >= serial_buff.max_size) { serial_buff.buff[serial_buff.select].clear(); qDebug() << "serial_buff.buff " << serial_buff.select <<" OVERFLOW"; } serial_buff.buff[serial_buff.select].append(data); break; } } QByteArray DataStream::readbuff(quint32 src) { QByteArray datagram; switch (src) { default: case SourceType::c_sock: if(client_buff.select == 0) { datagram.clear(); datagram.append(client_buff.buff[1]); //清除这个未选择的buff client_buff.buff[1].clear(); //读取完成,可以往这个内存里面写数了 client_buff.select = 1; } else if(client_buff.select == 1) { datagram.clear(); datagram.append(client_buff.buff[0]); //清除这个未选择的buff client_buff.buff[0].clear(); //读取完成,可以往这个内存里面写数了 client_buff.select = 0; } break; case SourceType::s_port: if(serial_buff.select == 0) { datagram.clear(); datagram.append(serial_buff.buff[1]); //读取完成,可以往这个内存里面写数了 serial_buff.buff[1].clear(); serial_buff.select = 1; } else if(serial_buff.select == 1) { datagram.clear(); datagram.append(serial_buff.buff[0]); //读取完成,可以往这个内存里面写数了 serial_buff.buff[0].clear(); serial_buff.select = 0; } break; } return datagram; } //发送函数 void DataStream::sendData(const int &id, const QByteArray &data) { idCurrent = id; databuff.clear(); databuff.append(data); status.m_Mode = TransmitMode; } void DataStream::sendAFrame(void) { //收到指令 mavlink_message_t msg; mavlink_encapsulated_data_t encapsulated_data; encapsulated_data.seqnr = idCurrent; size_t len = databuff.size(); for(uint16_t i = 0;i < len;i++) { encapsulated_data.data[i] = (uint8_t)databuff[i]; } if(len < 253) { for(uint16_t i = len;i < 253;i++) { encapsulated_data.data[i] = 0; } } qDebug() << "send encap" << databuff; mavlink_msg_encapsulated_data_encode(GCS_SysID,GCS_CompID, &msg,&encapsulated_data); Send(msg); }