| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 2de845b commit 340b983
3 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -157,16 +157,9 @@ class WorkerThreadsTaskRunner::DelayedTaskScheduler { | |||
| 157 | 157 | DelayedTaskScheduler* scheduler = | |
| 158 | 158 | ContainerOf(&DelayedTaskScheduler::loop_, flush_tasks->loop); | |
| 159 | 159 | ||
| 160 | - auto tasks_to_run = scheduler->tasks_.Lock().PopAll(); | ||
| 161 | - while (!tasks_to_run.empty()) { | ||
| 162 | - // We have to use const_cast because std::priority_queue::top() does not | ||
| 163 | - // return a movable item. | ||
| 164 | - std::unique_ptr<Task> task = | ||
| 165 | - std::move(const_cast<std::unique_ptr<Task>&>(tasks_to_run.top())); | ||
| 166 | - tasks_to_run.pop(); | ||
| 167 | - // This runs either the ScheduleTasks that scheduels the timers to | ||
| 168 | - // pop the tasks back into the worker task runner queue, or the | ||
| 169 | - // or the StopTasks to stop the timers and drop all the pending tasks. | ||
| 160 | + // ScheduleTasks (start a timer that pops the task into the worker queue) | ||
| 161 | + // in posting order, then, once Stop() was called, the StopTask. | ||
| 162 | + for (std::unique_ptr<Task>& task : scheduler->tasks_.Lock().PopAll()) { | ||
| 170 | 163 | task->Run(); | |
| 171 | 164 | } | |
| 172 | 165 | } | |
@@ -611,15 +604,8 @@ void NodePlatform::DrainTasks(Isolate* isolate) { | |||
| 611 | 604 | bool PerIsolatePlatformData::FlushForegroundTasksInternal() { | |
| 612 | 605 | bool did_work = false; | |
| 613 | 606 | ||
| 614 | - auto delayed_tasks_to_schedule = foreground_delayed_tasks_.Lock().PopAll(); | ||
| 615 | - while (!delayed_tasks_to_schedule.empty()) { | ||
| 616 | - // We have to use const_cast because std::priority_queue::top() does not | ||
| 617 | - // return a movable item. | ||
| 618 | - std::unique_ptr<DelayedTask> delayed = | ||
| 619 | - std::move(const_cast<std::unique_ptr<DelayedTask>&>( | ||
| 620 | - delayed_tasks_to_schedule.top())); | ||
| 621 | - delayed_tasks_to_schedule.pop(); | ||
| 622 | - | ||
| 607 | + for (std::unique_ptr<DelayedTask>& delayed : | ||
| 608 | + foreground_delayed_tasks_.Lock().PopAll()) { | ||
| 623 | 609 | did_work = true; | |
| 624 | 610 | uint64_t delay_millis = llround(delayed->timeout * 1000); | |
| 625 | 611 | ||
@@ -642,18 +628,8 @@ bool PerIsolatePlatformData::FlushForegroundTasksInternal() { | |||
| 642 | 628 | }); | |
| 643 | 629 | } | |
| 644 | 630 | ||
| 645 | - TaskQueue<TaskQueueEntry>::PriorityQueue tasks; | ||
| 646 | - { | ||
| 647 | - auto locked = foreground_tasks_.Lock(); | ||
| 648 | - tasks = locked.PopAll(); | ||
| 649 | - } | ||
| 650 | - | ||
| 651 | - while (!tasks.empty()) { | ||
| 652 | - // We have to use const_cast because std::priority_queue::top() does not | ||
| 653 | - // return a movable item. | ||
| 654 | - std::unique_ptr<TaskQueueEntry> entry = | ||
| 655 | - std::move(const_cast<std::unique_ptr<TaskQueueEntry>&>(tasks.top())); | ||
| 656 | - tasks.pop(); | ||
| 631 | + for (std::unique_ptr<TaskQueueEntry>& entry : | ||
| 632 | + foreground_tasks_.Lock().PopAll()) { | ||
| 657 | 633 | did_work = true; | |
| 658 | 634 | RunForegroundTask(std::move(entry->task)); | |
| 659 | 635 | } | |
@@ -788,12 +764,21 @@ template <class T> | |||
| 788 | 764 | TaskQueue<T>::Locked::Locked(TaskQueue* queue) | |
| 789 | 765 | : queue_(queue), lock_(queue->lock_) {} | |
| 790 | 766 | ||
| 767 | + template <class T> | ||
| 768 | + std::unique_ptr<T> TaskQueue<T>::PopTask() { | ||
| 769 | + // std::priority_queue::top() only hands out a const reference. | ||
| 770 | + Item& top = const_cast<Item&>(task_queue_.top()); | ||
| 771 | + std::unique_ptr<T> task = std::move(top.task); | ||
| 772 | + task_queue_.pop(); | ||
| 773 | + return task; | ||
| 774 | + } | ||
| 775 | + | ||
| 791 | 776 | template <class T> | |
| 792 | 777 | void TaskQueue<T>::Locked::Push(std::unique_ptr<T> task, bool outstanding) { | |
| 793 | 778 | if (outstanding) { | |
| 794 | 779 | queue_->outstanding_tasks_++; | |
| 795 | 780 | } | |
| 796 | - queue_->task_queue_.push(std::move(task)); | ||
| 781 | + queue_->task_queue_.push({std::move(task), queue_->next_sequence_++}); | ||
| 797 | 782 | queue_->tasks_available_.Signal(lock_); | |
| 798 | 783 | } | |
| 799 | 784 | ||
@@ -802,10 +787,7 @@ std::unique_ptr<T> TaskQueue<T>::Locked::Pop() { | |||
| 802 | 787 | if (queue_->task_queue_.empty()) { | |
| 803 | 788 | return std::unique_ptr<T>(nullptr); | |
| 804 | 789 | } | |
| 805 | - std::unique_ptr<T> result = std::move( | ||
| 806 | - std::move(const_cast<std::unique_ptr<T>&>(queue_->task_queue_.top()))); | ||
| 807 | - queue_->task_queue_.pop(); | ||
| 808 | - return result; | ||
| 790 | + return queue_->PopTask(); | ||
| 809 | 791 | } | |
| 810 | 792 | ||
| 811 | 793 | template <class T> | |
@@ -816,10 +798,7 @@ std::unique_ptr<T> TaskQueue<T>::Locked::BlockingPop() { | |||
| 816 | 798 | if (queue_->stopped_) { | |
| 817 | 799 | return std::unique_ptr<T>(nullptr); | |
| 818 | 800 | } | |
| 819 | - std::unique_ptr<T> result = std::move( | ||
| 820 | - std::move(const_cast<std::unique_ptr<T>&>(queue_->task_queue_.top()))); | ||
| 821 | - queue_->task_queue_.pop(); | ||
| 822 | - return result; | ||
| 801 | + return queue_->PopTask(); | ||
| 823 | 802 | } | |
| 824 | 803 | ||
| 825 | 804 | template <class T> | |
@@ -843,12 +822,19 @@ void TaskQueue<T>::Locked::Stop() { | |||
| 843 | 822 | } | |
| 844 | 823 | ||
| 845 | 824 | template <class T> | |
| 846 | - TaskQueue<T>::PriorityQueue TaskQueue<T>::Locked::PopAll() { | ||
| 847 | - TaskQueue<T>::PriorityQueue result; | ||
| 848 | - result.swap(queue_->task_queue_); | ||
| 825 | + std::vector<std::unique_ptr<T>> TaskQueue<T>::Locked::PopAll() { | ||
| 826 | + std::vector<std::unique_ptr<T>> result; | ||
| 827 | + result.reserve(queue_->task_queue_.size()); | ||
| 828 | + while (!queue_->task_queue_.empty()) { | ||
| 829 | + result.push_back(queue_->PopTask()); | ||
| 830 | + } | ||
| 849 | 831 | return result; | |
| 850 | 832 | } | |
| 851 | 833 | ||
| 834 | + template class TaskQueue<Task>; | ||
| 835 | + template class TaskQueue<TaskQueueEntry>; | ||
| 836 | + template class TaskQueue<DelayedTask>; | ||
| 837 | + | ||
| 852 | 838 | void MultiIsolatePlatform::DisposeIsolate(Isolate* isolate) { | |
| 853 | 839 | // The order of these calls is important. When the Isolate is disposed, | |
| 854 | 840 | // it may still post tasks to the platform, so it must still be registered | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -27,22 +27,6 @@ concept has_priority = requires(T t) { t.priority; }; | |||
| 27 | 27 | template <class T> | |
| 28 | 28 | class TaskQueue { | |
| 29 | 29 | public: | |
| 30 | - // If the entry type has a priority member, order the priority queue by | ||
| 31 | - // that - higher priority first. Otherwise, maintain insertion order. | ||
| 32 | - struct EntryCompare { | ||
| 33 | - bool operator()(const std::unique_ptr<T>& a, | ||
| 34 | - const std::unique_ptr<T>& b) const { | ||
| 35 | - if constexpr (has_priority<T>) { | ||
| 36 | - return a->priority < b->priority; | ||
| 37 | - } else { | ||
| 38 | - return false; | ||
| 39 | - } | ||
| 40 | - } | ||
| 41 | - }; | ||
| 42 | - | ||
| 43 | - using PriorityQueue = std::priority_queue<std::unique_ptr<T>, | ||
| 44 | - std::vector<std::unique_ptr<T>>, | ||
| 45 | - EntryCompare>; | ||
| 46 | 30 | class Locked { | |
| 47 | 31 | public: | |
| 48 | 32 | void Push(std::unique_ptr<T> task, bool outstanding = false); | |
@@ -51,7 +35,8 @@ class TaskQueue { | |||
| 51 | 35 | void NotifyOfOutstandingCompletion(); | |
| 52 | 36 | void BlockingDrain(); | |
| 53 | 37 | void Stop(); | |
| 54 | - PriorityQueue PopAll(); | ||
| 38 | + // All queued tasks, in the order Pop() would have returned them. | ||
| 39 | + std::vector<std::unique_ptr<T>> PopAll(); | ||
| 55 | 40 | ||
| 56 | 41 | private: | |
| 57 | 42 | friend class TaskQueue; | |
@@ -67,11 +52,33 @@ class TaskQueue { | |||
| 67 | 52 | Locked Lock() { return Locked(this); } | |
| 68 | 53 | ||
| 69 | 54 | private: | |
| 55 | + struct Item { | ||
| 56 | + std::unique_ptr<T> task; | ||
| 57 | + uint64_t sequence; | ||
| 58 | + }; | ||
| 59 | + // Higher priority first if the entry type has one; posting order otherwise | ||
| 60 | + // and among equal priorities (a sequence number breaks the tie). | ||
| 61 | + struct ItemCompare { | ||
| 62 | + bool operator()(const Item& a, const Item& b) const { | ||
| 63 | + if constexpr (has_priority<T>) { | ||
| 64 | + if (a.task->priority != b.task->priority) { | ||
| 65 | + return a.task->priority < b.task->priority; | ||
| 66 | + } | ||
| 67 | + } | ||
| 68 | + return a.sequence > b.sequence; | ||
| 69 | + } | ||
| 70 | + }; | ||
| 71 | + using PriorityQueue = | ||
| 72 | + std::priority_queue<Item, std::vector<Item>, ItemCompare>; | ||
| 73 | + | ||
| 74 | + std::unique_ptr<T> PopTask(); | ||
| 75 | + | ||
| 70 | 76 | Mutex lock_; | |
| 71 | 77 | ConditionVariable tasks_available_; | |
| 72 | 78 | ConditionVariable outstanding_tasks_drained_; | |
| 73 | 79 | int outstanding_tasks_; | |
| 74 | 80 | bool stopped_; | |
| 81 | + uint64_t next_sequence_ = 0; | ||
| 75 | 82 | PriorityQueue task_queue_; | |
| 76 | 83 | }; | |
| 77 | 84 | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -128,3 +128,53 @@ TEST_F(PlatformTest, TracingControllerNullptr) { | |||
| 128 | 128 | node::SetTracingController(orig_controller); | |
| 129 | 129 | EXPECT_EQ(node::GetTracingController(), orig_controller); | |
| 130 | 130 | } | |
| 131 | + | ||
| 132 | + class RecordingTask : public v8::Task { | ||
| 133 | + public: | ||
| 134 | + RecordingTask(std::vector<int>* log, int id) : log_(log), id_(id) {} | ||
| 135 | + void Run() override { log_->push_back(id_); } | ||
| 136 | + | ||
| 137 | + private: | ||
| 138 | + std::vector<int>* log_; | ||
| 139 | + int id_; | ||
| 140 | + }; | ||
| 141 | + | ||
| 142 | + TEST(TaskQueueTest, HigherPriorityFirstThenPostingOrder) { | ||
| 143 | + std::vector<int> log; | ||
| 144 | + { | ||
| 145 | + node::TaskQueue<v8::Task> queue; | ||
| 146 | + for (int i = 0; i < 64; i++) { | ||
| 147 | + queue.Lock().Push(std::make_unique<RecordingTask>(&log, i)); | ||
| 148 | + } | ||
| 149 | + for (std::unique_ptr<v8::Task>& task : queue.Lock().PopAll()) task->Run(); | ||
| 150 | + for (int i = 64; i < 96; i++) { | ||
| 151 | + queue.Lock().Push(std::make_unique<RecordingTask>(&log, i)); | ||
| 152 | + } | ||
| 153 | + while (std::unique_ptr<v8::Task> task = queue.Lock().Pop()) task->Run(); | ||
| 154 | + } | ||
| 155 | + ASSERT_EQ(log.size(), 96u); | ||
| 156 | + for (int i = 0; i < 96; i++) EXPECT_EQ(log[i], i); | ||
| 157 | + | ||
| 158 | + log.clear(); | ||
| 159 | + { | ||
| 160 | + using v8::TaskPriority; | ||
| 161 | + node::TaskQueue<node::TaskQueueEntry> queue; | ||
| 162 | + const TaskPriority priorities[] = {TaskPriority::kUserVisible, | ||
| 163 | + TaskPriority::kBestEffort, | ||
| 164 | + TaskPriority::kUserBlocking, | ||
| 165 | + TaskPriority::kUserVisible, | ||
| 166 | + TaskPriority::kUserBlocking, | ||
| 167 | + TaskPriority::kBestEffort, | ||
| 168 | + TaskPriority::kUserVisible}; | ||
| 169 | + int id = 0; | ||
| 170 | + for (TaskPriority priority : priorities) { | ||
| 171 | + queue.Lock().Push(std::make_unique<node::TaskQueueEntry>( | ||
| 172 | + std::make_unique<RecordingTask>(&log, id++), priority)); | ||
| 173 | + } | ||
| 174 | + for (std::unique_ptr<node::TaskQueueEntry>& entry : queue.Lock().PopAll()) { | ||
| 175 | + entry->task->Run(); | ||
| 176 | + } | ||
| 177 | + } | ||
| 178 | + const std::vector<int> expected = {2, 4, 0, 3, 6, 1, 5}; | ||
| 179 | + EXPECT_EQ(log, expected); | ||
| 180 | + } | ||
| Back | FazBrowse Home | New Git URL |
0 commit comments