threadsafe_queue.h 3.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147
  1. // Copyright 2010 Dolphin Emulator Project
  2. // Licensed under GPLv2+
  3. // Refer to the license.txt file included.
  4. #pragma once
  5. // a simple lockless thread-safe,
  6. // single reader, single writer queue
  7. #include <atomic>
  8. #include <cstddef>
  9. #include <mutex>
  10. #include <utility>
  11. namespace Common {
  12. template <typename T>
  13. class SPSCQueue {
  14. public:
  15. SPSCQueue() {
  16. write_ptr = read_ptr = new ElementPtr();
  17. }
  18. ~SPSCQueue() {
  19. // this will empty out the whole queue
  20. delete read_ptr;
  21. }
  22. std::size_t Size() const {
  23. return size.load();
  24. }
  25. bool Empty() const {
  26. return Size() == 0;
  27. }
  28. T& Front() const {
  29. return read_ptr->current;
  30. }
  31. template <typename Arg>
  32. void Push(Arg&& t) {
  33. // create the element, add it to the queue
  34. write_ptr->current = std::forward<Arg>(t);
  35. // set the next pointer to a new element ptr
  36. // then advance the write pointer
  37. ElementPtr* new_ptr = new ElementPtr();
  38. write_ptr->next.store(new_ptr, std::memory_order_release);
  39. write_ptr = new_ptr;
  40. ++size;
  41. }
  42. void Pop() {
  43. --size;
  44. ElementPtr* tmpptr = read_ptr;
  45. // advance the read pointer
  46. read_ptr = tmpptr->next.load();
  47. // set the next element to nullptr to stop the recursive deletion
  48. tmpptr->next.store(nullptr);
  49. delete tmpptr; // this also deletes the element
  50. }
  51. bool Pop(T& t) {
  52. if (Empty())
  53. return false;
  54. --size;
  55. ElementPtr* tmpptr = read_ptr;
  56. read_ptr = tmpptr->next.load(std::memory_order_acquire);
  57. t = std::move(tmpptr->current);
  58. tmpptr->next.store(nullptr);
  59. delete tmpptr;
  60. return true;
  61. }
  62. // not thread-safe
  63. void Clear() {
  64. size.store(0);
  65. delete read_ptr;
  66. write_ptr = read_ptr = new ElementPtr();
  67. }
  68. private:
  69. // stores a pointer to element
  70. // and a pointer to the next ElementPtr
  71. class ElementPtr {
  72. public:
  73. ElementPtr() {}
  74. ~ElementPtr() {
  75. ElementPtr* next_ptr = next.load();
  76. if (next_ptr)
  77. delete next_ptr;
  78. }
  79. T current;
  80. std::atomic<ElementPtr*> next{nullptr};
  81. };
  82. ElementPtr* write_ptr;
  83. ElementPtr* read_ptr;
  84. std::atomic_size_t size{0};
  85. };
  86. // a simple thread-safe,
  87. // single reader, multiple writer queue
  88. template <typename T>
  89. class MPSCQueue {
  90. public:
  91. std::size_t Size() const {
  92. return spsc_queue.Size();
  93. }
  94. bool Empty() const {
  95. return spsc_queue.Empty();
  96. }
  97. T& Front() const {
  98. return spsc_queue.Front();
  99. }
  100. template <typename Arg>
  101. void Push(Arg&& t) {
  102. std::lock_guard<std::mutex> lock(write_lock);
  103. spsc_queue.Push(t);
  104. }
  105. void Pop() {
  106. return spsc_queue.Pop();
  107. }
  108. bool Pop(T& t) {
  109. return spsc_queue.Pop(t);
  110. }
  111. // not thread-safe
  112. void Clear() {
  113. spsc_queue.Clear();
  114. }
  115. private:
  116. SPSCQueue<T> spsc_queue;
  117. std::mutex write_lock;
  118. };
  119. } // namespace Common