MultithreadedCompressor.h 6.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225
  1. // Copyright 2020 Dolphin Emulator Project
  2. // SPDX-License-Identifier: GPL-2.0-or-later
  3. #pragma once
  4. #include <atomic>
  5. #include <functional>
  6. #include <memory>
  7. #include <thread>
  8. #include <utility>
  9. #include <variant>
  10. #include <vector>
  11. #include "Common/Assert.h"
  12. #include "Common/Event.h"
  13. #include "Common/Result.h"
  14. namespace DiscIO
  15. {
  16. enum class ConversionResultCode
  17. {
  18. Success,
  19. Canceled,
  20. ReadFailed,
  21. WriteFailed,
  22. InternalError,
  23. };
  24. template <typename T>
  25. using ConversionResult = Common::Result<ConversionResultCode, T>;
  26. // This class starts a number of compression threads and one output thread.
  27. // The set_up_compress_thread_state function is called at the start of each compression thread.
  28. // When CompressAndWrite is called, the compress function will be called on one of the
  29. // compression threads, and then the output function will be called on the output thread.
  30. // The output thread handles data in the order that data was submitted using CompressAndWrite,
  31. // but the compression threads are not guaranteed to handle data in a predictable order.
  32. // Remember to check GetStatus regularly and cancel if it doesn't return Success,
  33. // and call Shutdown when you want to ensure that everything finishes.
  34. template <typename CompressThreadState, typename CompressParameters, typename OutputParameters>
  35. class MultithreadedCompressor
  36. {
  37. public:
  38. MultithreadedCompressor(
  39. std::function<ConversionResultCode(CompressThreadState*)> set_up_compress_thread_state,
  40. std::function<ConversionResult<OutputParameters>(CompressThreadState*, CompressParameters)>
  41. compress,
  42. std::function<ConversionResultCode(OutputParameters)> output)
  43. : m_set_up_compress_thread_state(std::move(set_up_compress_thread_state)),
  44. m_compress(std::move(compress)), m_output(std::move(output)),
  45. m_threads(std::max<unsigned int>(1, std::thread::hardware_concurrency()))
  46. {
  47. m_compress_threads = std::make_unique<CompressThread[]>(m_threads);
  48. for (size_t i = 0; i < m_threads; ++i)
  49. {
  50. m_compress_threads[i].thread =
  51. std::thread(std::mem_fn(&MultithreadedCompressor::CompressThreadFunction), this,
  52. &m_compress_threads[i]);
  53. }
  54. m_output_thread =
  55. std::thread(std::mem_fn(&MultithreadedCompressor::OutputThreadFunction), this);
  56. }
  57. ~MultithreadedCompressor()
  58. {
  59. if (!m_shutting_down.load())
  60. Shutdown();
  61. }
  62. void CompressAndWrite(CompressParameters parameters)
  63. {
  64. if (GetStatus() != ConversionResultCode::Success)
  65. return;
  66. CompressThread& compress_thread = m_compress_threads[m_current_index];
  67. compress_thread.compress_ready_event.Wait();
  68. compress_thread.compress_parameters = std::move(parameters);
  69. compress_thread.compress_event.Set();
  70. ++m_current_index;
  71. if (m_current_index >= m_threads)
  72. m_current_index -= m_threads;
  73. }
  74. void SetError(ConversionResultCode result)
  75. {
  76. ASSERT(result != ConversionResultCode::Success);
  77. // If we already have an error, don't overwrite it
  78. ConversionResultCode expected = ConversionResultCode::Success;
  79. m_result.compare_exchange_strong(expected, result);
  80. }
  81. ConversionResultCode GetStatus() const { return m_result.load(); }
  82. void Shutdown()
  83. {
  84. for (size_t i = 0; i < m_threads; ++i)
  85. m_compress_threads[i].compress_ready_event.Wait();
  86. for (size_t i = 0; i < m_threads; ++i)
  87. m_compress_threads[i].compress_done_event.Wait();
  88. for (size_t i = 0; i < m_threads; ++i)
  89. m_compress_threads[i].output_ready_event.Wait();
  90. m_shutting_down.store(true);
  91. for (size_t i = 0; i < m_threads; ++i)
  92. m_compress_threads[i].compress_event.Set();
  93. for (size_t i = 0; i < m_threads; ++i)
  94. m_compress_threads[i].output_event.Set();
  95. for (size_t i = 0; i < m_threads; ++i)
  96. m_compress_threads[i].thread.join();
  97. m_output_thread.join();
  98. }
  99. private:
  100. struct CompressThread
  101. {
  102. std::thread thread;
  103. Common::Event compress_ready_event;
  104. Common::Event compress_event;
  105. Common::Event compress_done_event;
  106. Common::Event output_ready_event;
  107. Common::Event output_event;
  108. CompressParameters compress_parameters;
  109. OutputParameters output_parameters;
  110. };
  111. void CompressThreadFunction(CompressThread* state)
  112. {
  113. CompressThreadState compress_thread_state;
  114. ConversionResultCode setup_result = m_set_up_compress_thread_state(&compress_thread_state);
  115. if (setup_result != ConversionResultCode::Success)
  116. SetError(setup_result);
  117. state->compress_ready_event.Set();
  118. state->compress_done_event.Set();
  119. while (true)
  120. {
  121. state->compress_event.Wait();
  122. if (m_shutting_down.load())
  123. return;
  124. CompressParameters parameters = std::move(state->compress_parameters);
  125. state->compress_done_event.Reset();
  126. state->compress_ready_event.Set();
  127. ConversionResult<OutputParameters> result =
  128. m_compress(&compress_thread_state, std::move(parameters));
  129. if (result)
  130. {
  131. state->output_ready_event.Wait();
  132. state->output_parameters = std::move(*result);
  133. state->output_event.Set();
  134. }
  135. else
  136. {
  137. SetError(result.Error());
  138. }
  139. state->compress_done_event.Set();
  140. }
  141. }
  142. void OutputThreadFunction()
  143. {
  144. for (size_t i = 0; i < m_threads; ++i)
  145. m_compress_threads[i].output_ready_event.Set();
  146. size_t index = 0;
  147. while (true)
  148. {
  149. CompressThread& compress_thread = m_compress_threads[index];
  150. compress_thread.output_event.Wait();
  151. if (m_shutting_down.load())
  152. return;
  153. OutputParameters parameters = std::move(compress_thread.output_parameters);
  154. compress_thread.output_ready_event.Set();
  155. const ConversionResultCode result = m_output(std::move(parameters));
  156. if (result != ConversionResultCode::Success)
  157. SetError(result);
  158. ++index;
  159. if (index >= m_threads)
  160. index -= m_threads;
  161. }
  162. }
  163. std::function<ConversionResultCode(CompressThreadState*)> m_set_up_compress_thread_state;
  164. std::function<ConversionResult<OutputParameters>(CompressThreadState*, CompressParameters)>
  165. m_compress;
  166. std::function<ConversionResultCode(OutputParameters)> m_output;
  167. // We can't use std::vector for this, because Common::Event is not movable
  168. std::unique_ptr<CompressThread[]> m_compress_threads;
  169. std::thread m_output_thread;
  170. const size_t m_threads;
  171. size_t m_current_index = 0;
  172. std::atomic<ConversionResultCode> m_result = ConversionResultCode::Success;
  173. std::atomic<bool> m_shutting_down = false;
  174. };
  175. } // namespace DiscIO