mirror of
https://github.com/kunkundi/crossdesk.git
synced 2025-10-27 04:35:34 +08:00
Fix crash caused by multi threads during program termination
This commit is contained in:
@@ -6,19 +6,9 @@
|
||||
|
||||
RtpVideoReceiver::RtpVideoReceiver() {}
|
||||
|
||||
RtpVideoReceiver::~RtpVideoReceiver() {
|
||||
if (jitter_thread_ && jitter_thread_->joinable()) {
|
||||
jitter_thread_->join();
|
||||
delete jitter_thread_;
|
||||
jitter_thread_ = nullptr;
|
||||
}
|
||||
}
|
||||
RtpVideoReceiver::~RtpVideoReceiver() {}
|
||||
|
||||
void RtpVideoReceiver::InsertRtpPacket(RtpPacket& rtp_packet) {
|
||||
if (!jitter_thread_) {
|
||||
jitter_thread_ = new std::thread(&RtpVideoReceiver::Process, this);
|
||||
}
|
||||
|
||||
if (NAL_UNIT_TYPE::NALU == rtp_packet.NalUnitType()) {
|
||||
compelete_video_frame_queue_.push(
|
||||
VideoFrame(rtp_packet.Payload(), rtp_packet.Size()));
|
||||
@@ -98,23 +88,36 @@ bool RtpVideoReceiver::CheckIsFrameCompleted(RtpPacket& rtp_packet) {
|
||||
return false;
|
||||
}
|
||||
|
||||
void RtpVideoReceiver::Process() {
|
||||
while (1) {
|
||||
if (!compelete_video_frame_queue_.isEmpty()) {
|
||||
VideoFrame video_frame;
|
||||
compelete_video_frame_queue_.pop(video_frame);
|
||||
if (on_receive_complete_frame_) {
|
||||
auto now_complete_frame_ts = std::chrono::high_resolution_clock::now()
|
||||
.time_since_epoch()
|
||||
.count() /
|
||||
1000000;
|
||||
uint32_t duration = now_complete_frame_ts - last_complete_frame_ts_;
|
||||
LOG_ERROR("Duration {}", 1000 / duration);
|
||||
last_complete_frame_ts_ = now_complete_frame_ts;
|
||||
on_receive_complete_frame_(video_frame);
|
||||
}
|
||||
}
|
||||
void RtpVideoReceiver::Start() {
|
||||
std::lock_guard<std::mutex> lock_guard(mutex_);
|
||||
stop_ = false;
|
||||
}
|
||||
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(13));
|
||||
void RtpVideoReceiver::Stop() {
|
||||
std::lock_guard<std::mutex> lock_guard(mutex_);
|
||||
stop_ = true;
|
||||
}
|
||||
|
||||
bool RtpVideoReceiver::Process() {
|
||||
std::lock_guard<std::mutex> lock_guard(mutex_);
|
||||
if (stop_) {
|
||||
return false;
|
||||
}
|
||||
|
||||
if (!compelete_video_frame_queue_.isEmpty()) {
|
||||
VideoFrame video_frame;
|
||||
compelete_video_frame_queue_.pop(video_frame);
|
||||
if (on_receive_complete_frame_) {
|
||||
auto now_complete_frame_ts =
|
||||
std::chrono::high_resolution_clock::now().time_since_epoch().count() /
|
||||
1000000;
|
||||
uint32_t duration = now_complete_frame_ts - last_complete_frame_ts_;
|
||||
LOG_ERROR("Duration {}", 1000 / duration);
|
||||
last_complete_frame_ts_ = now_complete_frame_ts;
|
||||
on_receive_complete_frame_(video_frame);
|
||||
}
|
||||
}
|
||||
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(13));
|
||||
return true;
|
||||
}
|
||||
@@ -3,14 +3,15 @@
|
||||
|
||||
#include <functional>
|
||||
#include <map>
|
||||
#include <mutex>
|
||||
#include <queue>
|
||||
#include <thread>
|
||||
|
||||
#include "frame.h"
|
||||
#include "ringbuffer.h"
|
||||
#include "rtp_video_session.h"
|
||||
#include "thread_base.h"
|
||||
|
||||
class RtpVideoReceiver {
|
||||
class RtpVideoReceiver : public ThreadBase {
|
||||
public:
|
||||
RtpVideoReceiver();
|
||||
~RtpVideoReceiver();
|
||||
@@ -23,13 +24,16 @@ class RtpVideoReceiver {
|
||||
on_receive_complete_frame_ = on_receive_complete_frame;
|
||||
}
|
||||
|
||||
void Start();
|
||||
void Stop();
|
||||
|
||||
private:
|
||||
bool CheckIsFrameCompleted(RtpPacket& rtp_packet);
|
||||
void Process();
|
||||
|
||||
// private:
|
||||
// void OnReceiveFrame(uint8_t* payload) {}
|
||||
|
||||
private:
|
||||
bool Process() override;
|
||||
|
||||
private:
|
||||
std::map<uint16_t, RtpPacket> incomplete_frame_list_;
|
||||
uint8_t* nv12_data_ = nullptr;
|
||||
@@ -37,8 +41,9 @@ class RtpVideoReceiver {
|
||||
uint32_t last_complete_frame_ts_ = 0;
|
||||
|
||||
RingBuffer<VideoFrame> compelete_video_frame_queue_;
|
||||
std::thread* jitter_thread_ = nullptr;
|
||||
bool start_ = false;
|
||||
|
||||
bool stop_ = true;
|
||||
std::mutex mutex_;
|
||||
};
|
||||
|
||||
#endif
|
||||
|
||||
@@ -2,29 +2,35 @@
|
||||
|
||||
#include <chrono>
|
||||
|
||||
#include "log.h"
|
||||
|
||||
RtpVideoSender::RtpVideoSender() {}
|
||||
|
||||
RtpVideoSender::~RtpVideoSender() {
|
||||
if (send_thread_ && send_thread_->joinable()) {
|
||||
send_thread_->join();
|
||||
delete send_thread_;
|
||||
send_thread_ = nullptr;
|
||||
}
|
||||
}
|
||||
RtpVideoSender::~RtpVideoSender() {}
|
||||
|
||||
void RtpVideoSender::Enqueue(std::vector<RtpPacket>& rtp_packets) {
|
||||
if (!send_thread_) {
|
||||
send_thread_ = new std::thread(&RtpVideoSender::Process, this);
|
||||
}
|
||||
|
||||
for (auto& rtp_packet : rtp_packets) {
|
||||
start_ = true;
|
||||
rtp_packe_queue_.push(rtp_packet);
|
||||
}
|
||||
}
|
||||
|
||||
void RtpVideoSender::Process() {
|
||||
while (1) {
|
||||
void RtpVideoSender::Start() {
|
||||
std::lock_guard<std::mutex> lock_guard(mutex_);
|
||||
stop_ = false;
|
||||
}
|
||||
|
||||
void RtpVideoSender::Stop() {
|
||||
std::lock_guard<std::mutex> lock_guard(mutex_);
|
||||
stop_ = true;
|
||||
}
|
||||
|
||||
bool RtpVideoSender::Process() {
|
||||
std::lock_guard<std::mutex> lock_guard(mutex_);
|
||||
if (stop_) {
|
||||
return false;
|
||||
}
|
||||
|
||||
for (size_t i = 0; i < 50; i++)
|
||||
if (!rtp_packe_queue_.isEmpty()) {
|
||||
RtpPacket rtp_packet;
|
||||
rtp_packe_queue_.pop(rtp_packet);
|
||||
@@ -33,6 +39,5 @@ void RtpVideoSender::Process() {
|
||||
}
|
||||
}
|
||||
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(1));
|
||||
}
|
||||
return true;
|
||||
}
|
||||
@@ -2,12 +2,14 @@
|
||||
#define _RTP_VIDEO_SENDER_H_
|
||||
|
||||
#include <functional>
|
||||
#include <mutex>
|
||||
#include <thread>
|
||||
|
||||
#include "ringbuffer.h"
|
||||
#include "rtp_packet.h"
|
||||
#include "thread_base.h"
|
||||
|
||||
class RtpVideoSender {
|
||||
class RtpVideoSender : public ThreadBase {
|
||||
public:
|
||||
RtpVideoSender();
|
||||
~RtpVideoSender();
|
||||
@@ -21,14 +23,18 @@ class RtpVideoSender {
|
||||
rtp_packet_send_func_ = rtp_packet_send_func;
|
||||
}
|
||||
|
||||
void Start();
|
||||
void Stop();
|
||||
|
||||
private:
|
||||
void Process();
|
||||
bool Process() override;
|
||||
|
||||
private:
|
||||
std::function<void(RtpPacket &)> rtp_packet_send_func_ = nullptr;
|
||||
RingBuffer<RtpPacket> rtp_packe_queue_;
|
||||
std::thread *send_thread_ = nullptr;
|
||||
bool start_ = false;
|
||||
|
||||
bool stop_ = true;
|
||||
std::mutex mutex_;
|
||||
};
|
||||
|
||||
#endif
|
||||
Reference in New Issue
Block a user