thread_manager.cpp 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590
  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. : ThreadBase("VideoDecodeThread")
  170. , m_packetQueue(packetQueue)
  171. , m_frameQueue(frameQueue)
  172. , m_decoder(decoder)
  173. , m_synchronizer(synchronizer)
  174. , m_streamIndex(streamIndex)
  175. , m_bsfContext(nullptr)
  176. {
  177. if (codecParams) {
  178. initBitStreamFilter(codecParams);
  179. }
  180. }
  181. VideoDecodeThread::~VideoDecodeThread()
  182. {
  183. cleanupBitStreamFilter();
  184. }
  185. void VideoDecodeThread::run()
  186. {
  187. if (!m_packetQueue || !m_frameQueue || !m_decoder) {
  188. av::Logger::instance().error("Invalid parameters for VideoDecodeThread");
  189. return;
  190. }
  191. while (!shouldStop()) {
  192. // 检查帧队列是否已满
  193. if (m_frameQueue->full()) {
  194. std::this_thread::sleep_for(std::chrono::milliseconds(10));
  195. continue;
  196. }
  197. // 解码帧
  198. if (!decodeFrame()) {
  199. if (m_packetQueue->empty()) {
  200. break;
  201. }
  202. std::this_thread::sleep_for(std::chrono::milliseconds(1));
  203. }
  204. }
  205. av::Logger::instance().info("VideoDecodeThread finished");
  206. }
  207. bool VideoDecodeThread::decodeFrame()
  208. {
  209. // 获取数据包
  210. AVPacket* packet = m_packetQueue->dequeuePacket(10); // 10ms超时
  211. if (!packet) {
  212. return false;
  213. }
  214. // 检查是否是视频包
  215. if (packet->stream_index != m_streamIndex) {
  216. av_packet_free(&packet);
  217. return true;
  218. }
  219. // 如果有比特流过滤器,先过滤数据包
  220. AVPacket* filteredPacket = packet;
  221. if (m_bsfContext) {
  222. int ret = av_bsf_send_packet(m_bsfContext, packet);
  223. if (ret < 0) {
  224. av::Logger::instance().errorf("Failed to send packet to bitstream filter: {}",
  225. av::ffmpeg_utils::errorToString(ret));
  226. av_packet_free(&packet);
  227. return false;
  228. }
  229. AVPacket* outputPacket = av_packet_alloc();
  230. ret = av_bsf_receive_packet(m_bsfContext, outputPacket);
  231. if (ret < 0) {
  232. if (ret != AVERROR(EAGAIN) && ret != AVERROR_EOF) {
  233. av::Logger::instance().errorf("Failed to receive packet from bitstream filter: {}",
  234. av::ffmpeg_utils::errorToString(ret));
  235. }
  236. av_packet_free(&packet);
  237. av_packet_free(&outputPacket);
  238. return false;
  239. }
  240. av_packet_free(&packet);
  241. filteredPacket = outputPacket;
  242. }
  243. // 解码
  244. std::vector<AVFramePtr> frames;
  245. AVPacketPtr packetPtr(filteredPacket);
  246. ErrorCode result = m_decoder->decode(packetPtr, frames);
  247. if (result != ErrorCode::SUCCESS) {
  248. av::Logger::instance().errorf("Video decode failed with error: {}", static_cast<int>(result));
  249. return false;
  250. }
  251. if (frames.empty()) {
  252. // 这是正常情况,解码器可能需要更多输入
  253. return true;
  254. }
  255. av::Logger::instance().debugf("Successfully decoded {} video frames", frames.size());
  256. // 处理解码后的帧
  257. for (auto& frame : frames) {
  258. // 更新视频时钟
  259. if (m_synchronizer && frame->pts != AV_NOPTS_VALUE) {
  260. // 需要获取时间基,这里暂时使用默认值
  261. double pts = frame->pts * 0.000040; // 临时值,需要从解码器获取正确的时间基
  262. m_synchronizer->setVideoClock(pts);
  263. }
  264. // 添加到帧队列
  265. if (m_frameQueue->enqueue(frame.release()) != av::ErrorCode::SUCCESS) {
  266. return false;
  267. }
  268. }
  269. return true;
  270. }
  271. bool VideoDecodeThread::initBitStreamFilter(AVCodecParameters* codecParams)
  272. {
  273. if (!codecParams || codecParams->codec_id != AV_CODEC_ID_H264) {
  274. return true; // 不需要比特流过滤器
  275. }
  276. const AVBitStreamFilter* bsf = av_bsf_get_by_name("h264_mp4toannexb");
  277. if (!bsf) {
  278. av::Logger::instance().error("Failed to find h264_mp4toannexb bitstream filter");
  279. return false;
  280. }
  281. int ret = av_bsf_alloc(bsf, &m_bsfContext);
  282. if (ret < 0) {
  283. av::Logger::instance().errorf("Failed to allocate bitstream filter context: {}",
  284. av::ffmpeg_utils::errorToString(ret));
  285. return false;
  286. }
  287. ret = avcodec_parameters_copy(m_bsfContext->par_in, codecParams);
  288. if (ret < 0) {
  289. av::Logger::instance().errorf("Failed to copy codec parameters to bitstream filter: {}",
  290. av::ffmpeg_utils::errorToString(ret));
  291. av_bsf_free(&m_bsfContext);
  292. return false;
  293. }
  294. ret = av_bsf_init(m_bsfContext);
  295. if (ret < 0) {
  296. av::Logger::instance().errorf("Failed to initialize bitstream filter: {}",
  297. av::ffmpeg_utils::errorToString(ret));
  298. av_bsf_free(&m_bsfContext);
  299. return false;
  300. }
  301. av::Logger::instance().info("H.264 bitstream filter initialized successfully");
  302. return true;
  303. }
  304. void VideoDecodeThread::cleanupBitStreamFilter()
  305. {
  306. if (m_bsfContext) {
  307. av_bsf_free(&m_bsfContext);
  308. m_bsfContext = nullptr;
  309. }
  310. }
  311. // AudioDecodeThread 实现
  312. AudioDecodeThread::AudioDecodeThread(av::utils::PacketQueue* packetQueue,
  313. av::utils::FrameQueue* frameQueue,
  314. AudioDecoder* decoder,
  315. av::utils::Synchronizer* synchronizer,
  316. int streamIndex)
  317. : ThreadBase("AudioDecodeThread")
  318. , m_packetQueue(packetQueue)
  319. , m_frameQueue(frameQueue)
  320. , m_decoder(decoder)
  321. , m_synchronizer(synchronizer)
  322. , m_streamIndex(streamIndex)
  323. {
  324. }
  325. void AudioDecodeThread::run()
  326. {
  327. if (!m_packetQueue || !m_frameQueue || !m_decoder) {
  328. av::Logger::instance().error("Invalid parameters for AudioDecodeThread");
  329. return;
  330. }
  331. while (!shouldStop()) {
  332. // 检查帧队列是否已满
  333. if (m_frameQueue->full()) {
  334. std::this_thread::sleep_for(std::chrono::milliseconds(10));
  335. continue;
  336. }
  337. // 解码帧
  338. if (!decodeFrame()) {
  339. if (m_packetQueue->empty()) {
  340. break;
  341. }
  342. std::this_thread::sleep_for(std::chrono::milliseconds(1));
  343. }
  344. }
  345. av::Logger::instance().info("AudioDecodeThread finished");
  346. }
  347. bool AudioDecodeThread::decodeFrame()
  348. {
  349. // 获取数据包
  350. AVPacket* packet = m_packetQueue->dequeuePacket(10); // 10ms超时
  351. if (!packet) {
  352. return false;
  353. }
  354. // 检查是否是音频包
  355. if (packet->stream_index != m_streamIndex) {
  356. av_packet_free(&packet);
  357. return true;
  358. }
  359. // 解码
  360. std::vector<AVFramePtr> frames;
  361. AVPacketPtr packetPtr(packet);
  362. ErrorCode result = m_decoder->decode(packetPtr, frames);
  363. if (result != ErrorCode::SUCCESS) {
  364. av::Logger::instance().errorf("Audio decode failed with error: {}", static_cast<int>(result));
  365. return false;
  366. }
  367. if (frames.empty()) {
  368. // 这是正常情况,解码器可能需要更多输入
  369. return true;
  370. }
  371. av::Logger::instance().debugf("Successfully decoded {} audio frames", frames.size());
  372. // 处理解码后的帧
  373. for (auto& frame : frames) {
  374. // 更新音频时钟
  375. if (m_synchronizer && frame->pts != AV_NOPTS_VALUE) {
  376. // 需要获取时间基,这里暂时使用默认值
  377. double pts = frame->pts * 0.000023; // 临时值,需要从解码器获取正确的时间基
  378. m_synchronizer->setAudioClock(pts);
  379. }
  380. // 添加到帧队列
  381. if (m_frameQueue->enqueue(frame.release()) != av::ErrorCode::SUCCESS) {
  382. return false;
  383. }
  384. }
  385. return true;
  386. }
  387. // ThreadManager 实现
  388. ThreadManager::ThreadManager()
  389. {
  390. av::Logger::instance().debug("ThreadManager created");
  391. }
  392. ThreadManager::~ThreadManager()
  393. {
  394. stopAll();
  395. joinAll();
  396. av::Logger::instance().debug("ThreadManager destroyed");
  397. }
  398. bool ThreadManager::createReadThread(AVFormatContext* formatContext,
  399. av::utils::PacketQueue* packetQueue,
  400. int videoStreamIndex,
  401. int audioStreamIndex)
  402. {
  403. if (m_readThread) {
  404. av::Logger::instance().warning("ReadThread already exists");
  405. return false;
  406. }
  407. m_readThread = std::make_unique<ReadThread>(formatContext, packetQueue,
  408. videoStreamIndex, audioStreamIndex);
  409. av::Logger::instance().info("ReadThread created");
  410. return true;
  411. }
  412. bool ThreadManager::createVideoDecodeThread(av::utils::PacketQueue* packetQueue,
  413. av::utils::FrameQueue* frameQueue,
  414. VideoDecoder* decoder,
  415. av::utils::Synchronizer* synchronizer,
  416. int streamIndex,
  417. AVCodecParameters* codecParams)
  418. {
  419. if (m_videoDecodeThread) {
  420. av::Logger::instance().warning("VideoDecodeThread already exists");
  421. return false;
  422. }
  423. m_videoDecodeThread = std::make_unique<VideoDecodeThread>(packetQueue, frameQueue,
  424. decoder, synchronizer, streamIndex, codecParams);
  425. av::Logger::instance().info("VideoDecodeThread created");
  426. return true;
  427. }
  428. bool ThreadManager::createAudioDecodeThread(av::utils::PacketQueue* packetQueue,
  429. av::utils::FrameQueue* frameQueue,
  430. AudioDecoder* decoder,
  431. av::utils::Synchronizer* synchronizer,
  432. int streamIndex)
  433. {
  434. if (m_audioDecodeThread) {
  435. av::Logger::instance().warning("AudioDecodeThread already exists");
  436. return false;
  437. }
  438. m_audioDecodeThread = std::make_unique<AudioDecodeThread>(packetQueue, frameQueue,
  439. decoder, synchronizer, streamIndex);
  440. av::Logger::instance().info("AudioDecodeThread created");
  441. return true;
  442. }
  443. bool ThreadManager::startAll()
  444. {
  445. bool success = true;
  446. if (m_readThread && !m_readThread->start()) {
  447. av::Logger::instance().error("Failed to start ReadThread");
  448. success = false;
  449. }
  450. if (m_videoDecodeThread && !m_videoDecodeThread->start()) {
  451. av::Logger::instance().error("Failed to start VideoDecodeThread");
  452. success = false;
  453. }
  454. if (m_audioDecodeThread && !m_audioDecodeThread->start()) {
  455. av::Logger::instance().error("Failed to start AudioDecodeThread");
  456. success = false;
  457. }
  458. if (success) {
  459. av::Logger::instance().info("All threads started successfully");
  460. }
  461. return success;
  462. }
  463. void ThreadManager::stopAll()
  464. {
  465. if (m_readThread) {
  466. m_readThread->stop();
  467. }
  468. if (m_videoDecodeThread) {
  469. m_videoDecodeThread->stop();
  470. }
  471. if (m_audioDecodeThread) {
  472. m_audioDecodeThread->stop();
  473. }
  474. av::Logger::instance().info("All threads stop requested");
  475. }
  476. void ThreadManager::joinAll()
  477. {
  478. if (m_readThread) {
  479. m_readThread->join();
  480. m_readThread.reset();
  481. }
  482. if (m_videoDecodeThread) {
  483. m_videoDecodeThread->join();
  484. m_videoDecodeThread.reset();
  485. }
  486. if (m_audioDecodeThread) {
  487. m_audioDecodeThread->join();
  488. m_audioDecodeThread.reset();
  489. }
  490. av::Logger::instance().info("All threads joined");
  491. }
  492. bool ThreadManager::hasRunningThreads() const
  493. {
  494. return (m_readThread && m_readThread->isRunning()) ||
  495. (m_videoDecodeThread && m_videoDecodeThread->isRunning()) ||
  496. (m_audioDecodeThread && m_audioDecodeThread->isRunning());
  497. }
  498. } // namespace player
  499. } // namespace av