thread_manager.cpp 18 KB

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