| 87 | } |
| 88 | |
| 89 | void DelayedThreadPool::ProcessTasks() |
| 90 | { |
| 91 | ImmediateQueue pendingImmediate; |
| 92 | DelayedQueue pendingDelayed; |
| 93 | |
| 94 | while (true) |
| 95 | { |
| 96 | std::array<Task, QUEUE_TYPE_COUNT> tasks; |
| 97 | |
| 98 | { |
| 99 | std::unique_lock lk(m_mu); |
| 100 | if (!m_delayed.IsEmpty()) |
| 101 | { |
| 102 | // We need to wait until the moment when the earliest delayed |
| 103 | // task may be executed, given that an immediate task or a |
| 104 | // delayed task with an earlier execution time may arrive |
| 105 | // while we are waiting. |
| 106 | auto const when = m_delayed.GetFirstValue()->m_when; |
| 107 | m_cv.wait_until(lk, when, [this, when]() |
| 108 | { |
| 109 | if (m_shutdown || !m_immediate.IsEmpty() || m_delayed.IsEmpty()) |
| 110 | return true; |
| 111 | return m_delayed.GetFirstValue()->m_when < when; |
| 112 | }); |
| 113 | } |
| 114 | else |
| 115 | { |
| 116 | // When there are no delayed tasks in the queue, we need to |
| 117 | // wait until there is at least one immediate or delayed task. |
| 118 | m_cv.wait(lk, [this]() { return m_shutdown || !m_immediate.IsEmpty() || !m_delayed.IsEmpty(); }); |
| 119 | } |
| 120 | |
| 121 | if (m_shutdown) |
| 122 | { |
| 123 | switch (m_exit) |
| 124 | { |
| 125 | case Exit::ExecPending: |
| 126 | ASSERT(pendingImmediate.IsEmpty(), ()); |
| 127 | m_immediate.Swap(pendingImmediate); |
| 128 | |
| 129 | ASSERT(pendingDelayed.IsEmpty(), ()); |
| 130 | m_delayed.Swap(pendingDelayed); |
| 131 | break; |
| 132 | case Exit::SkipPending: break; |
| 133 | } |
| 134 | |
| 135 | break; |
| 136 | } |
| 137 | |
| 138 | auto const canExecImmediate = !m_immediate.IsEmpty(); |
| 139 | auto const canExecDelayed = !m_delayed.IsEmpty() && Now() >= m_delayed.GetFirstValue()->m_when; |
| 140 | |
| 141 | if (canExecImmediate) |
| 142 | { |
| 143 | tasks[QUEUE_TYPE_IMMEDIATE] = std::move(m_immediate.Front()); |
| 144 | m_immediate.PopFront(); |
| 145 | } |
| 146 | |