Skip to content

Commit 121482b

Browse files
committed
src: run same-priority platform tasks in posting order
TaskQueue became a std::priority_queue when worker tasks started to honor v8::TaskPriority. Its comparator returns false for entry types without a priority member, and for entries of equal priority, on the assumption that the heap then keeps insertion order. It does not: three tasks pushed A, B, C pop as A, C, B, and larger batches come out in heap order. That affects the per-isolate foreground task queue (tasks of one priority no longer run in the order they were posted), the foreground delayed task queue, and the delayed task scheduler of the worker thread task runner, whose local queue is drained in one batch: when v8 posts a delayed worker task shortly before the platform shuts down, the StopTask pushed by Stop() can run before a ScheduleTask that was pushed earlier, that ScheduleTask then starts a timer on the scheduler's loop after all timers were supposed to be stopped, and Shutdown() blocks in uv_thread_join() until the delay (e.g. the 8 s of the memory reducer) expires. Give every queued item a sequence number and use it as the tie breaker, so that tasks of equal priority, and tasks without one, come out in FIFO order again; higher priorities still come first. PopAll() now returns the tasks in that order instead of handing out the heap. Signed-off-by: Shelley Vohr <shelley.vohr@gmail.com>
1 parent 30bff4a commit 121482b

3 files changed

Lines changed: 103 additions & 60 deletions

File tree

src/node_platform.cc

Lines changed: 29 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -157,16 +157,9 @@ class WorkerThreadsTaskRunner::DelayedTaskScheduler {
157157
DelayedTaskScheduler* scheduler =
158158
ContainerOf(&DelayedTaskScheduler::loop_, flush_tasks->loop);
159159

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()) {
170163
task->Run();
171164
}
172165
}
@@ -611,15 +604,8 @@ void NodePlatform::DrainTasks(Isolate* isolate) {
611604
bool PerIsolatePlatformData::FlushForegroundTasksInternal() {
612605
bool did_work = false;
613606

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()) {
623609
did_work = true;
624610
uint64_t delay_millis = llround(delayed->timeout * 1000);
625611

@@ -642,18 +628,8 @@ bool PerIsolatePlatformData::FlushForegroundTasksInternal() {
642628
});
643629
}
644630

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()) {
657633
did_work = true;
658634
RunForegroundTask(std::move(entry->task));
659635
}
@@ -788,12 +764,21 @@ template <class T>
788764
TaskQueue<T>::Locked::Locked(TaskQueue* queue)
789765
: queue_(queue), lock_(queue->lock_) {}
790766

767+
template <class T>
768+
std::unique_ptr<T> TaskQueue<T>::PopTask() {
769+
// std::priority_queue::top() only hands out a const reference.
770+
std::unique_ptr<T> task =
771+
std::move(const_cast<Item&>(task_queue_.top()).task);
772+
task_queue_.pop();
773+
return task;
774+
}
775+
791776
template <class T>
792777
void TaskQueue<T>::Locked::Push(std::unique_ptr<T> task, bool outstanding) {
793778
if (outstanding) {
794779
queue_->outstanding_tasks_++;
795780
}
796-
queue_->task_queue_.push(std::move(task));
781+
queue_->task_queue_.push({std::move(task), queue_->next_sequence_++});
797782
queue_->tasks_available_.Signal(lock_);
798783
}
799784

@@ -802,10 +787,7 @@ std::unique_ptr<T> TaskQueue<T>::Locked::Pop() {
802787
if (queue_->task_queue_.empty()) {
803788
return std::unique_ptr<T>(nullptr);
804789
}
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();
809791
}
810792

811793
template <class T>
@@ -816,10 +798,7 @@ std::unique_ptr<T> TaskQueue<T>::Locked::BlockingPop() {
816798
if (queue_->stopped_) {
817799
return std::unique_ptr<T>(nullptr);
818800
}
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();
823802
}
824803

825804
template <class T>
@@ -843,12 +822,19 @@ void TaskQueue<T>::Locked::Stop() {
843822
}
844823

845824
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+
}
849831
return result;
850832
}
851833

834+
template class TaskQueue<Task>;
835+
template class TaskQueue<TaskQueueEntry>;
836+
template class TaskQueue<DelayedTask>;
837+
852838
void MultiIsolatePlatform::DisposeIsolate(Isolate* isolate) {
853839
// The order of these calls is important. When the Isolate is disposed,
854840
// it may still post tasks to the platform, so it must still be registered

src/node_platform.h

Lines changed: 24 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -27,22 +27,6 @@ concept has_priority = requires(T t) { t.priority; };
2727
template <class T>
2828
class TaskQueue {
2929
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>;
4630
class Locked {
4731
public:
4832
void Push(std::unique_ptr<T> task, bool outstanding = false);
@@ -51,7 +35,8 @@ class TaskQueue {
5135
void NotifyOfOutstandingCompletion();
5236
void BlockingDrain();
5337
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();
5540

5641
private:
5742
friend class TaskQueue;
@@ -67,11 +52,33 @@ class TaskQueue {
6752
Locked Lock() { return Locked(this); }
6853

6954
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+
7076
Mutex lock_;
7177
ConditionVariable tasks_available_;
7278
ConditionVariable outstanding_tasks_drained_;
7379
int outstanding_tasks_;
7480
bool stopped_;
81+
uint64_t next_sequence_ = 0;
7582
PriorityQueue task_queue_;
7683
};
7784

test/cctest/test_platform.cc

Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -128,3 +128,53 @@ TEST_F(PlatformTest, TracingControllerNullptr) {
128128
node::SetTracingController(orig_controller);
129129
EXPECT_EQ(node::GetTracingController(), orig_controller);
130130
}
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+
}

0 commit comments

Comments
 (0)