thread_manager.cpp 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624
  1. #include "thread_manager.h"
  2. #include "../utils/utils_packet_queue.h"
  3. #include "../utils/utils_frame_queue.h"
  4. #include "thread_manager.h"
  5. #include "../base/logger.h"
  6. #include "../base/media_common.h"
  7. #include "../codec/codec_video_decoder.h"
  8. #include "../codec/codec_audio_decoder.h"
  9. #include "../utils/utils_packet_queue.h"
  10. #include "../utils/utils_frame_queue.h"
  11. #include "../utils/utils_synchronizer.h"
  12. #include <thread>
  13. #include <chrono>
  14. #include <vector>
  15. namespace av {
  16. namespace player {
  17. // ThreadBase 实现
  18. ThreadBase::ThreadBase(const std::string& name) : m_name(name)
  19. {
  20. av::Logger::instance().debug("Thread created: " + m_name);
  21. }
  22. ThreadBase::~ThreadBase()
  23. {
  24. stop();
  25. join();
  26. av::Logger::instance().debug("Thread destroyed: " + m_name);
  27. }
  28. bool ThreadBase::start()
  29. {
  30. if (m_running.load()) {
  31. av::Logger::instance().warning("Thread already running: " + m_name);
  32. return false;
  33. }
  34. m_shouldStop = false;
  35. m_thread = std::make_unique<std::thread>(&ThreadBase::threadEntry, this);
  36. av::Logger::instance().info("Thread started: " + m_name);
  37. return true;
  38. }
  39. void ThreadBase::stop()
  40. {
  41. if (!m_running.load()) {
  42. return;
  43. }
  44. m_shouldStop = true;
  45. av::Logger::instance().info("Thread stop requested: " + m_name);
  46. }
  47. void ThreadBase::join()
  48. {
  49. if (m_thread && m_thread->joinable()) {
  50. m_thread->join();
  51. m_thread.reset();
  52. av::Logger::instance().info("Thread joined: " + m_name);
  53. }
  54. }
  55. bool ThreadBase::isRunning() const
  56. {
  57. return m_running.load();
  58. }
  59. void ThreadBase::threadEntry()
  60. {
  61. m_running = true;
  62. try {
  63. av::Logger::instance().info("Thread running: " + m_name);
  64. run();
  65. }
  66. catch (const std::exception& e) {
  67. av::Logger::instance().error("Thread exception in " + m_name + ": " + e.what());
  68. }
  69. catch (...) {
  70. av::Logger::instance().error("Unknown exception in thread: " + m_name);
  71. }
  72. m_running = false;
  73. av::Logger::instance().info("Thread finished: " + m_name);
  74. }
  75. // ReadThread 实现
  76. ReadThread::ReadThread(AVFormatContext* formatContext,
  77. av::utils::PacketQueue* packetQueue,
  78. int videoStreamIndex,
  79. int audioStreamIndex)
  80. : ThreadBase("ReadThread")
  81. , m_formatContext(formatContext)
  82. , m_packetQueue(packetQueue)
  83. , m_videoStreamIndex(videoStreamIndex)
  84. , m_audioStreamIndex(audioStreamIndex)
  85. {
  86. }
  87. void ReadThread::seek(int64_t timestamp)
  88. {
  89. m_seekTarget = timestamp;
  90. m_seeking = true;
  91. av::Logger::instance().info("Seek requested to: " + std::to_string(timestamp));
  92. }
  93. void ReadThread::run()
  94. {
  95. if (!m_formatContext || !m_packetQueue) {
  96. av::Logger::instance().error("Invalid parameters for ReadThread");
  97. return;
  98. }
  99. while (!shouldStop()) {
  100. // 处理跳转
  101. if (m_seeking.load()) {
  102. handleSeek();
  103. continue;
  104. }
  105. // 检查队列是否已满
  106. if (m_packetQueue->full()) {
  107. std::this_thread::sleep_for(std::chrono::milliseconds(10));
  108. continue;
  109. }
  110. // 读取数据包
  111. if (!readPacket()) {
  112. // 读取结束或出错
  113. // EOF处理由队列内部管理
  114. break;
  115. }
  116. }
  117. av::Logger::instance().info("ReadThread finished");
  118. }
  119. bool ReadThread::readPacket()
  120. {
  121. AVPacket* packet = av_packet_alloc();
  122. if (!packet) {
  123. av::Logger::instance().error("Failed to allocate packet");
  124. return false;
  125. }
  126. int ret = av_read_frame(m_formatContext, packet);
  127. if (ret < 0) {
  128. av_packet_free(&packet);
  129. if (ret == AVERROR_EOF) {
  130. av::Logger::instance().info("End of file reached");
  131. } else {
  132. av::Logger::instance().error("Error reading frame: " + std::to_string(ret));
  133. }
  134. return false;
  135. }
  136. // 只处理我们关心的流
  137. if (packet->stream_index == m_videoStreamIndex ||
  138. packet->stream_index == m_audioStreamIndex) {
  139. if (m_packetQueue->enqueue(packet) != av::ErrorCode::SUCCESS) {
  140. av_packet_free(&packet);
  141. return false;
  142. }
  143. } else {
  144. av_packet_free(&packet);
  145. }
  146. return true;
  147. }
  148. void ReadThread::handleSeek()
  149. {
  150. int64_t seekTarget = m_seekTarget.load();
  151. // 执行跳转
  152. int ret = av_seek_frame(m_formatContext, -1, seekTarget, AVSEEK_FLAG_BACKWARD);
  153. if (ret < 0) {
  154. av::Logger::instance().error("Seek failed: " + std::to_string(ret));
  155. } else {
  156. av::Logger::instance().info("Seek completed to: " + std::to_string(seekTarget));
  157. // 清空队列
  158. m_packetQueue->clear();
  159. }
  160. m_seeking = false;
  161. }
  162. // VideoDecodeThread 实现
  163. VideoDecodeThread::VideoDecodeThread(av::utils::PacketQueue* packetQueue,
  164. av::utils::FrameQueue* frameQueue,
  165. VideoDecoder* decoder,
  166. av::utils::Synchronizer* synchronizer,
  167. int streamIndex,
  168. AVCodecParameters* codecParams,
  169. AVRational timeBase)
  170. : ThreadBase("VideoDecodeThread")
  171. , m_packetQueue(packetQueue)
  172. , m_frameQueue(frameQueue)
  173. , m_decoder(decoder)
  174. , m_synchronizer(synchronizer)
  175. , m_streamIndex(streamIndex)
  176. , m_bsfContext(nullptr)
  177. , m_frameOutputCallback(nullptr)
  178. , m_timeBase(timeBase)
  179. {
  180. if (codecParams) {
  181. initBitStreamFilter(codecParams);
  182. }
  183. }
  184. void VideoDecodeThread::setFrameOutputCallback(FrameOutputCallback callback)
  185. {
  186. m_frameOutputCallback = callback;
  187. }
  188. VideoDecodeThread::~VideoDecodeThread()
  189. {
  190. cleanupBitStreamFilter();
  191. }
  192. void VideoDecodeThread::run()
  193. {
  194. if (!m_packetQueue || !m_frameQueue || !m_decoder) {
  195. av::Logger::instance().error("Invalid parameters for VideoDecodeThread");
  196. return;
  197. }
  198. while (!shouldStop()) {
  199. // 检查帧队列是否已满
  200. if (m_frameQueue->full()) {
  201. std::this_thread::sleep_for(std::chrono::milliseconds(10));
  202. continue;
  203. }
  204. // 解码帧
  205. if (!decodeFrame()) {
  206. if (m_packetQueue->empty()) {
  207. break;
  208. }
  209. std::this_thread::sleep_for(std::chrono::milliseconds(1));
  210. }
  211. }
  212. av::Logger::instance().info("VideoDecodeThread finished");
  213. }
  214. bool VideoDecodeThread::decodeFrame()
  215. {
  216. // 获取数据包
  217. AVPacket* packet = m_packetQueue->dequeuePacket(10); // 10ms超时
  218. if (!packet) {
  219. return false;
  220. }
  221. // 检查是否是视频包
  222. if (packet->stream_index != m_streamIndex) {
  223. av_packet_free(&packet);
  224. return true;
  225. }
  226. // 如果有比特流过滤器,先过滤数据包
  227. AVPacket* filteredPacket = packet;
  228. if (m_bsfContext) {
  229. int ret = av_bsf_send_packet(m_bsfContext, packet);
  230. if (ret < 0) {
  231. av::Logger::instance().errorf("Failed to send packet to bitstream filter: {}",
  232. av::ffmpeg_utils::errorToString(ret));
  233. av_packet_free(&packet);
  234. return false;
  235. }
  236. AVPacket* outputPacket = av_packet_alloc();
  237. ret = av_bsf_receive_packet(m_bsfContext, outputPacket);
  238. if (ret < 0) {
  239. if (ret != AVERROR(EAGAIN) && ret != AVERROR_EOF) {
  240. av::Logger::instance().errorf("Failed to receive packet from bitstream filter: {}",
  241. av::ffmpeg_utils::errorToString(ret));
  242. }
  243. av_packet_free(&packet);
  244. av_packet_free(&outputPacket);
  245. return false;
  246. }
  247. av_packet_free(&packet);
  248. filteredPacket = outputPacket;
  249. }
  250. // 解码
  251. std::vector<AVFramePtr> frames;
  252. AVPacketPtr packetPtr(filteredPacket);
  253. ErrorCode result = m_decoder->decode(packetPtr, frames);
  254. if (result != ErrorCode::SUCCESS) {
  255. av::Logger::instance().errorf("Video decode failed with error: {}", static_cast<int>(result));
  256. return false;
  257. }
  258. if (frames.empty()) {
  259. // 这是正常情况,解码器可能需要更多输入
  260. return true;
  261. }
  262. av::Logger::instance().debugf("Successfully decoded {} video frames", frames.size());
  263. // 处理解码后的帧
  264. for (auto& frame : frames) {
  265. // 更新视频时钟
  266. if (m_synchronizer && frame->pts != AV_NOPTS_VALUE) {
  267. // 使用正确的时间基计算PTS
  268. double pts = frame->pts * av_q2d(m_timeBase);
  269. m_synchronizer->setVideoClock(pts);
  270. }
  271. // 调用输出回调函数(用于视频渲染)
  272. if (m_frameOutputCallback && frame) {
  273. m_frameOutputCallback(frame);
  274. }
  275. // 添加到帧队列
  276. if (m_frameQueue->enqueue(frame.release()) != av::ErrorCode::SUCCESS) {
  277. return false;
  278. }
  279. }
  280. return true;
  281. }
  282. bool VideoDecodeThread::initBitStreamFilter(AVCodecParameters* codecParams)
  283. {
  284. if (!codecParams || codecParams->codec_id != AV_CODEC_ID_H264) {
  285. return true; // 不需要比特流过滤器
  286. }
  287. const AVBitStreamFilter* bsf = av_bsf_get_by_name("h264_mp4toannexb");
  288. if (!bsf) {
  289. av::Logger::instance().error("Failed to find h264_mp4toannexb bitstream filter");
  290. return false;
  291. }
  292. int ret = av_bsf_alloc(bsf, &m_bsfContext);
  293. if (ret < 0) {
  294. av::Logger::instance().errorf("Failed to allocate bitstream filter context: {}",
  295. av::ffmpeg_utils::errorToString(ret));
  296. return false;
  297. }
  298. ret = avcodec_parameters_copy(m_bsfContext->par_in, codecParams);
  299. if (ret < 0) {
  300. av::Logger::instance().errorf("Failed to copy codec parameters to bitstream filter: {}",
  301. av::ffmpeg_utils::errorToString(ret));
  302. av_bsf_free(&m_bsfContext);
  303. return false;
  304. }
  305. ret = av_bsf_init(m_bsfContext);
  306. if (ret < 0) {
  307. av::Logger::instance().errorf("Failed to initialize bitstream filter: {}",
  308. av::ffmpeg_utils::errorToString(ret));
  309. av_bsf_free(&m_bsfContext);
  310. return false;
  311. }
  312. av::Logger::instance().info("H.264 bitstream filter initialized successfully");
  313. return true;
  314. }
  315. void VideoDecodeThread::cleanupBitStreamFilter()
  316. {
  317. if (m_bsfContext) {
  318. av_bsf_free(&m_bsfContext);
  319. m_bsfContext = nullptr;
  320. }
  321. }
  322. // AudioDecodeThread 实现
  323. AudioDecodeThread::AudioDecodeThread(av::utils::PacketQueue* packetQueue,
  324. av::utils::FrameQueue* frameQueue,
  325. AudioDecoder* decoder,
  326. av::utils::Synchronizer* synchronizer,
  327. int streamIndex,
  328. AVRational timeBase)
  329. : ThreadBase("AudioDecodeThread")
  330. , m_packetQueue(packetQueue)
  331. , m_frameQueue(frameQueue)
  332. , m_decoder(decoder)
  333. , m_synchronizer(synchronizer)
  334. , m_streamIndex(streamIndex)
  335. , m_frameOutputCallback(nullptr)
  336. , m_timeBase(timeBase)
  337. {
  338. }
  339. void AudioDecodeThread::setFrameOutputCallback(FrameOutputCallback callback)
  340. {
  341. m_frameOutputCallback = callback;
  342. }
  343. void AudioDecodeThread::run()
  344. {
  345. if (!m_packetQueue || !m_frameQueue || !m_decoder) {
  346. av::Logger::instance().error("Invalid parameters for AudioDecodeThread");
  347. return;
  348. }
  349. while (!shouldStop()) {
  350. // 检查帧队列是否已满
  351. if (m_frameQueue->full()) {
  352. std::this_thread::sleep_for(std::chrono::milliseconds(10));
  353. continue;
  354. }
  355. // 解码帧
  356. if (!decodeFrame()) {
  357. if (m_packetQueue->empty()) {
  358. break;
  359. }
  360. std::this_thread::sleep_for(std::chrono::milliseconds(1));
  361. }
  362. }
  363. av::Logger::instance().info("AudioDecodeThread finished");
  364. }
  365. bool AudioDecodeThread::decodeFrame()
  366. {
  367. // 获取数据包
  368. AVPacket* packet = m_packetQueue->dequeuePacket(10); // 10ms超时
  369. if (!packet) {
  370. return false;
  371. }
  372. // 检查是否是音频包
  373. if (packet->stream_index != m_streamIndex) {
  374. av_packet_free(&packet);
  375. return true;
  376. }
  377. // 解码
  378. std::vector<AVFramePtr> frames;
  379. AVPacketPtr packetPtr(packet);
  380. ErrorCode result = m_decoder->decode(packetPtr, frames);
  381. if (result != ErrorCode::SUCCESS) {
  382. av::Logger::instance().errorf("Audio decode failed with error: {}", static_cast<int>(result));
  383. return false;
  384. }
  385. if (frames.empty()) {
  386. // 这是正常情况,解码器可能需要更多输入
  387. return true;
  388. }
  389. av::Logger::instance().debugf("Successfully decoded {} audio frames", frames.size());
  390. // 处理解码后的帧
  391. for (auto& frame : frames) {
  392. // 更新音频时钟
  393. if (m_synchronizer && frame->pts != AV_NOPTS_VALUE) {
  394. // 使用正确的时间基计算PTS
  395. double pts = frame->pts * av_q2d(m_timeBase);
  396. m_synchronizer->setAudioClock(pts);
  397. }
  398. // 调用输出回调函数(用于音频输出)
  399. if (m_frameOutputCallback && frame) {
  400. try {
  401. m_frameOutputCallback(frame);
  402. } catch (const std::exception& e) {
  403. av::Logger::instance().error("Exception in audio frame output callback: " + std::string(e.what()));
  404. } catch (...) {
  405. av::Logger::instance().error("Unknown exception in audio frame output callback");
  406. }
  407. }
  408. // 添加到帧队列
  409. if (m_frameQueue->enqueue(frame.release()) != av::ErrorCode::SUCCESS) {
  410. return false;
  411. }
  412. }
  413. return true;
  414. }
  415. // ThreadManager 实现
  416. ThreadManager::ThreadManager()
  417. {
  418. av::Logger::instance().debug("ThreadManager created");
  419. }
  420. ThreadManager::~ThreadManager()
  421. {
  422. stopAll();
  423. joinAll();
  424. av::Logger::instance().debug("ThreadManager destroyed");
  425. }
  426. bool ThreadManager::createReadThread(AVFormatContext* formatContext,
  427. av::utils::PacketQueue* packetQueue,
  428. int videoStreamIndex,
  429. int audioStreamIndex)
  430. {
  431. if (m_readThread) {
  432. av::Logger::instance().warning("ReadThread already exists");
  433. return false;
  434. }
  435. m_readThread = std::make_unique<ReadThread>(formatContext, packetQueue,
  436. videoStreamIndex, audioStreamIndex);
  437. av::Logger::instance().info("ReadThread created");
  438. return true;
  439. }
  440. bool ThreadManager::createVideoDecodeThread(av::utils::PacketQueue* packetQueue,
  441. av::utils::FrameQueue* frameQueue,
  442. VideoDecoder* decoder,
  443. av::utils::Synchronizer* synchronizer,
  444. int streamIndex,
  445. AVCodecParameters* codecParams,
  446. AVRational timeBase)
  447. {
  448. if (m_videoDecodeThread) {
  449. av::Logger::instance().warning("VideoDecodeThread already exists");
  450. return false;
  451. }
  452. m_videoDecodeThread = std::make_unique<VideoDecodeThread>(packetQueue, frameQueue,
  453. decoder, synchronizer, streamIndex, codecParams, timeBase);
  454. av::Logger::instance().info("VideoDecodeThread created");
  455. return true;
  456. }
  457. bool ThreadManager::createAudioDecodeThread(av::utils::PacketQueue* packetQueue,
  458. av::utils::FrameQueue* frameQueue,
  459. AudioDecoder* decoder,
  460. av::utils::Synchronizer* synchronizer,
  461. int streamIndex,
  462. AVRational timeBase)
  463. {
  464. if (m_audioDecodeThread) {
  465. av::Logger::instance().warning("AudioDecodeThread already exists");
  466. return false;
  467. }
  468. m_audioDecodeThread = std::make_unique<AudioDecodeThread>(packetQueue, frameQueue,
  469. decoder, synchronizer, streamIndex, timeBase);
  470. av::Logger::instance().info("AudioDecodeThread created");
  471. return true;
  472. }
  473. bool ThreadManager::startAll()
  474. {
  475. bool success = true;
  476. if (m_readThread && !m_readThread->start()) {
  477. av::Logger::instance().error("Failed to start ReadThread");
  478. success = false;
  479. }
  480. if (m_videoDecodeThread && !m_videoDecodeThread->start()) {
  481. av::Logger::instance().error("Failed to start VideoDecodeThread");
  482. success = false;
  483. }
  484. if (m_audioDecodeThread && !m_audioDecodeThread->start()) {
  485. av::Logger::instance().error("Failed to start AudioDecodeThread");
  486. success = false;
  487. }
  488. if (success) {
  489. av::Logger::instance().info("All threads started successfully");
  490. }
  491. return success;
  492. }
  493. void ThreadManager::stopAll()
  494. {
  495. if (m_readThread) {
  496. m_readThread->stop();
  497. }
  498. if (m_videoDecodeThread) {
  499. m_videoDecodeThread->stop();
  500. }
  501. if (m_audioDecodeThread) {
  502. m_audioDecodeThread->stop();
  503. }
  504. av::Logger::instance().info("All threads stop requested");
  505. }
  506. void ThreadManager::joinAll()
  507. {
  508. if (m_readThread) {
  509. m_readThread->join();
  510. m_readThread.reset();
  511. }
  512. if (m_videoDecodeThread) {
  513. m_videoDecodeThread->join();
  514. m_videoDecodeThread.reset();
  515. }
  516. if (m_audioDecodeThread) {
  517. m_audioDecodeThread->join();
  518. m_audioDecodeThread.reset();
  519. }
  520. av::Logger::instance().info("All threads joined");
  521. }
  522. bool ThreadManager::hasRunningThreads() const
  523. {
  524. return (m_readThread && m_readThread->isRunning()) ||
  525. (m_videoDecodeThread && m_videoDecodeThread->isRunning()) ||
  526. (m_audioDecodeThread && m_audioDecodeThread->isRunning());
  527. }
  528. } // namespace player
  529. } // namespace av