Skip to content

Commit 2ea5bb8

Browse files
authored
Reland "TaskSources register tasks with MessageLoopTaskQueues dispatcher" (#25692)
1 parent 7e69744 commit 2ea5bb8

14 files changed

Lines changed: 547 additions & 57 deletions

‎ci/licenses_golden/licenses_flutter‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -269,8 +269,13 @@ FILE: ../../../flutter/fml/synchronization/sync_switch_unittest.cc
269269
FILE: ../../../flutter/fml/synchronization/waitable_event.cc
270270
FILE: ../../../flutter/fml/synchronization/waitable_event.h
271271
FILE: ../../../flutter/fml/synchronization/waitable_event_unittest.cc
272+
FILE: ../../../flutter/fml/task_queue_id.h
272273
FILE: ../../../flutter/fml/task_runner.cc
273274
FILE: ../../../flutter/fml/task_runner.h
275+
FILE: ../../../flutter/fml/task_source.cc
276+
FILE: ../../../flutter/fml/task_source.h
277+
FILE: ../../../flutter/fml/task_source_grade.h
278+
FILE: ../../../flutter/fml/task_source_unittests.cc
274279
FILE: ../../../flutter/fml/thread.cc
275280
FILE: ../../../flutter/fml/thread.h
276281
FILE: ../../../flutter/fml/thread_local.cc

‎fml/BUILD.gn‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,8 +70,11 @@ source_set("fml") {
7070
"synchronization/sync_switch.h",
7171
"synchronization/waitable_event.cc",
7272
"synchronization/waitable_event.h",
73+
"task_queue_id.h",
7374
"task_runner.cc",
7475
"task_runner.h",
76+
"task_source.cc",
77+
"task_source.h",
7578
"thread.cc",
7679
"thread.h",
7780
"thread_local.cc",
@@ -261,6 +264,7 @@ if (enable_unittests) {
261264
"synchronization/semaphore_unittest.cc",
262265
"synchronization/sync_switch_unittest.cc",
263266
"synchronization/waitable_event_unittest.cc",
267+
"task_source_unittests.cc",
264268
"thread_local_unittests.cc",
265269
"thread_unittests.cc",
266270
"time/time_delta_unittest.cc",

‎fml/delayed_task.cc‎

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -10,13 +10,17 @@ namespace fml {
1010

1111
DelayedTask::DelayedTask(size_t order,
1212
const fml::closure& task,
13-
fml::TimePoint target_time)
14-
: order_(order), task_(task), target_time_(target_time) {}
15-
16-
DelayedTask::DelayedTask(const DelayedTask& other) = default;
13+
fml::TimePoint target_time,
14+
fml::TaskSourceGrade task_source_grade)
15+
: order_(order),
16+
task_(task),
17+
target_time_(target_time),
18+
task_source_grade_(task_source_grade) {}
1719

1820
DelayedTask::~DelayedTask() = default;
1921

22+
DelayedTask::DelayedTask(const DelayedTask& other) = default;
23+
2024
const fml::closure& DelayedTask::GetTask() const {
2125
return task_;
2226
}
@@ -25,6 +29,10 @@ fml::TimePoint DelayedTask::GetTargetTime() const {
2529
return target_time_;
2630
}
2731

32+
fml::TaskSourceGrade DelayedTask::GetTaskSourceGrade() const {
33+
return task_source_grade_;
34+
}
35+
2836
bool DelayedTask::operator>(const DelayedTask& other) const {
2937
if (target_time_ == other.target_time_) {
3038
return order_ > other.order_;

‎fml/delayed_task.h‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
#include <queue>
99

1010
#include "flutter/fml/closure.h"
11+
#include "flutter/fml/task_source_grade.h"
1112
#include "flutter/fml/time/time_point.h"
1213

1314
namespace fml {
@@ -16,7 +17,8 @@ class DelayedTask {
1617
public:
1718
DelayedTask(size_t order,
1819
const fml::closure& task,
19-
fml::TimePoint target_time);
20+
fml::TimePoint target_time,
21+
fml::TaskSourceGrade task_source_grade);
2022

2123
DelayedTask(const DelayedTask& other);
2224

@@ -26,12 +28,15 @@ class DelayedTask {
2628

2729
fml::TimePoint GetTargetTime() const;
2830

31+
fml::TaskSourceGrade GetTaskSourceGrade() const;
32+
2933
bool operator>(const DelayedTask& other) const;
3034

3135
private:
3236
size_t order_;
3337
fml::closure task_;
3438
fml::TimePoint target_time_;
39+
fml::TaskSourceGrade task_source_grade_;
3540
};
3641

3742
using DelayedTaskQueue = std::priority_queue<DelayedTask,

‎fml/message_loop_task_queues.cc‎

Lines changed: 89 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -7,29 +7,54 @@
77
#include "flutter/fml/message_loop_task_queues.h"
88

99
#include <iostream>
10+
#include <memory>
1011

1112
#include "flutter/fml/make_copyable.h"
1213
#include "flutter/fml/message_loop_impl.h"
14+
#include "flutter/fml/task_source.h"
15+
#include "flutter/fml/thread_local.h"
1316

1417
namespace fml {
1518

1619
std::mutex MessageLoopTaskQueues::creation_mutex_;
1720

1821
const size_t TaskQueueId::kUnmerged = ULONG_MAX;
1922

23+
// Guarded by creation_mutex_.
2024
fml::RefPtr<MessageLoopTaskQueues> MessageLoopTaskQueues::instance_;
2125

22-
TaskQueueEntry::TaskQueueEntry()
23-
: owner_of(_kUnmerged), subsumed_by(_kUnmerged) {
26+
namespace {
27+
28+
// iOS prior to version 9 prevents c++11 thread_local and __thread specefier,
29+
// having us resort to boxed enum containers.
30+
class TaskSourceGradeHolder {
31+
public:
32+
TaskSourceGrade task_source_grade;
33+
34+
explicit TaskSourceGradeHolder(TaskSourceGrade task_source_grade_arg)
35+
: task_source_grade(task_source_grade_arg) {}
36+
};
37+
} // namespace
38+
39+
// Guarded by creation_mutex_.
40+
FML_THREAD_LOCAL ThreadLocalUniquePtr<TaskSourceGradeHolder>
41+
tls_task_source_grade;
42+
43+
TaskQueueEntry::TaskQueueEntry(TaskQueueId created_for_arg)
44+
: owner_of(_kUnmerged),
45+
subsumed_by(_kUnmerged),
46+
created_for(created_for_arg) {
2447
wakeable = NULL;
2548
task_observers = TaskObservers();
26-
delayed_tasks = DelayedTaskQueue();
49+
task_source = std::make_unique<TaskSource>(created_for);
2750
}
2851

2952
fml::RefPtr<MessageLoopTaskQueues> MessageLoopTaskQueues::GetInstance() {
3053
std::scoped_lock creation(creation_mutex_);
3154
if (!instance_) {
3255
instance_ = fml::MakeRefCounted<MessageLoopTaskQueues>();
56+
tls_task_source_grade.reset(
57+
new TaskSourceGradeHolder{TaskSourceGrade::kUnspecified});
3358
}
3459
return instance_;
3560
}
@@ -38,7 +63,7 @@ TaskQueueId MessageLoopTaskQueues::CreateTaskQueue() {
3863
std::lock_guard guard(queue_mutex_);
3964
TaskQueueId loop_id = TaskQueueId(task_queue_id_counter_);
4065
++task_queue_id_counter_;
41-
queue_entries_[loop_id] = std::make_unique<TaskQueueEntry>();
66+
queue_entries_[loop_id] = std::make_unique<TaskQueueEntry>(loop_id);
4267
return loop_id;
4368
}
4469

@@ -63,24 +88,36 @@ void MessageLoopTaskQueues::DisposeTasks(TaskQueueId queue_id) {
6388
const auto& queue_entry = queue_entries_.at(queue_id);
6489
FML_DCHECK(queue_entry->subsumed_by == _kUnmerged);
6590
TaskQueueId subsumed = queue_entry->owner_of;
66-
queue_entry->delayed_tasks = {};
91+
queue_entry->task_source->ShutDown();
6792
if (subsumed != _kUnmerged) {
68-
queue_entries_.at(subsumed)->delayed_tasks = {};
93+
queue_entries_.at(subsumed)->task_source->ShutDown();
6994
}
7095
}
7196

72-
void MessageLoopTaskQueues::RegisterTask(TaskQueueId queue_id,
73-
const fml::closure& task,
74-
fml::TimePoint target_time) {
97+
TaskSourceGrade MessageLoopTaskQueues::GetCurrentTaskSourceGrade() {
98+
std::scoped_lock creation(creation_mutex_);
99+
return tls_task_source_grade.get()->task_source_grade;
100+
}
101+
102+
void MessageLoopTaskQueues::RegisterTask(
103+
TaskQueueId queue_id,
104+
const fml::closure& task,
105+
fml::TimePoint target_time,
106+
fml::TaskSourceGrade task_source_grade) {
75107
std::lock_guard guard(queue_mutex_);
76108
size_t order = order_++;
77109
const auto& queue_entry = queue_entries_.at(queue_id);
78-
queue_entry->delayed_tasks.push({order, task, target_time});
110+
queue_entry->task_source->RegisterTask(
111+
{order, task, target_time, task_source_grade});
79112
TaskQueueId loop_to_wake = queue_id;
80113
if (queue_entry->subsumed_by != _kUnmerged) {
81114
loop_to_wake = queue_entry->subsumed_by;
82115
}
83-
WakeUpUnlocked(loop_to_wake, GetNextWakeTimeUnlocked(loop_to_wake));
116+
117+
// This can happen when the secondary tasks are paused.
118+
if (HasPendingTasksUnlocked(loop_to_wake)) {
119+
WakeUpUnlocked(loop_to_wake, GetNextWakeTimeUnlocked(loop_to_wake));
120+
}
84121
}
85122

86123
bool MessageLoopTaskQueues::HasPendingTasks(TaskQueueId queue_id) const {
@@ -94,20 +131,25 @@ fml::closure MessageLoopTaskQueues::GetNextTaskToRun(TaskQueueId queue_id,
94131
if (!HasPendingTasksUnlocked(queue_id)) {
95132
return nullptr;
96133
}
97-
TaskQueueId top_queue = _kUnmerged;
98-
const auto& top = PeekNextTaskUnlocked(queue_id, top_queue);
134+
TaskSource::TopTask top = PeekNextTaskUnlocked(queue_id);
99135

100136
if (!HasPendingTasksUnlocked(queue_id)) {
101137
WakeUpUnlocked(queue_id, fml::TimePoint::Max());
102138
} else {
103139
WakeUpUnlocked(queue_id, GetNextWakeTimeUnlocked(queue_id));
104140
}
105141

106-
if (top.GetTargetTime() > from_time) {
142+
if (top.task.GetTargetTime() > from_time) {
107143
return nullptr;
108144
}
109-
fml::closure invocation = top.GetTask();
110-
queue_entries_.at(top_queue)->delayed_tasks.pop();
145+
fml::closure invocation = top.task.GetTask();
146+
queue_entries_.at(top.task_queue_id)
147+
->task_source->PopTask(top.task.GetTaskSourceGrade());
148+
{
149+
std::scoped_lock creation(creation_mutex_);
150+
const auto task_source_grade = top.task.GetTaskSourceGrade();
151+
tls_task_source_grade.reset(new TaskSourceGradeHolder{task_source_grade});
152+
}
111153
return invocation;
112154
}
113155

@@ -126,12 +168,12 @@ size_t MessageLoopTaskQueues::GetNumPendingTasks(TaskQueueId queue_id) const {
126168
}
127169

128170
size_t total_tasks = 0;
129-
total_tasks += queue_entry->delayed_tasks.size();
171+
total_tasks += queue_entry->task_source->GetNumPendingTasks();
130172

131173
TaskQueueId subsumed = queue_entry->owner_of;
132174
if (subsumed != _kUnmerged) {
133175
const auto& subsumed_entry = queue_entries_.at(subsumed);
134-
total_tasks += subsumed_entry->delayed_tasks.size();
176+
total_tasks += subsumed_entry->task_source->GetNumPendingTasks();
135177
}
136178
return total_tasks;
137179
}
@@ -248,6 +290,20 @@ TaskQueueId MessageLoopTaskQueues::GetSubsumedTaskQueueId(
248290
return queue_entries_.at(owner)->owner_of;
249291
}
250292

293+
void MessageLoopTaskQueues::PauseSecondarySource(TaskQueueId queue_id) {
294+
std::lock_guard guard(queue_mutex_);
295+
queue_entries_.at(queue_id)->task_source->PauseSecondary();
296+
}
297+
298+
void MessageLoopTaskQueues::ResumeSecondarySource(TaskQueueId queue_id) {
299+
std::lock_guard guard(queue_mutex_);
300+
queue_entries_.at(queue_id)->task_source->ResumeSecondary();
301+
// Schedule a wake as needed.
302+
if (HasPendingTasksUnlocked(queue_id)) {
303+
WakeUpUnlocked(queue_id, GetNextWakeTimeUnlocked(queue_id));
304+
}
305+
}
306+
251307
// Subsumed queues will never have pending tasks.
252308
// Owning queues will consider both their and their subsumed tasks.
253309
bool MessageLoopTaskQueues::HasPendingTasksUnlocked(
@@ -258,7 +314,7 @@ bool MessageLoopTaskQueues::HasPendingTasksUnlocked(
258314
return false;
259315
}
260316

261-
if (!entry->delayed_tasks.empty()) {
317+
if (!entry->task_source->IsEmpty()) {
262318
return true;
263319
}
264320

@@ -267,37 +323,35 @@ bool MessageLoopTaskQueues::HasPendingTasksUnlocked(
267323
// this is not an owner and queue is empty.
268324
return false;
269325
} else {
270-
return !queue_entries_.at(subsumed)->delayed_tasks.empty();
326+
return !queue_entries_.at(subsumed)->task_source->IsEmpty();
271327
}
272328
}
273329

274330
fml::TimePoint MessageLoopTaskQueues::GetNextWakeTimeUnlocked(
275331
TaskQueueId queue_id) const {
276-
TaskQueueId tmp = _kUnmerged;
277-
return PeekNextTaskUnlocked(queue_id, tmp).GetTargetTime();
332+
return PeekNextTaskUnlocked(queue_id).task.GetTargetTime();
278333
}
279334

280-
const DelayedTask& MessageLoopTaskQueues::PeekNextTaskUnlocked(
281-
TaskQueueId owner,
282-
TaskQueueId& top_queue_id) const {
335+
TaskSource::TopTask MessageLoopTaskQueues::PeekNextTaskUnlocked(
336+
TaskQueueId owner) const {
283337
FML_DCHECK(HasPendingTasksUnlocked(owner));
284338
const auto& entry = queue_entries_.at(owner);
285339
const TaskQueueId subsumed = entry->owner_of;
286340
if (subsumed == _kUnmerged) {
287-
top_queue_id = owner;
288-
return entry->delayed_tasks.top();
341+
return entry->task_source->Top();
289342
}
290343

291-
const auto& owner_tasks = entry->delayed_tasks;
292-
const auto& subsumed_tasks = queue_entries_.at(subsumed)->delayed_tasks;
344+
TaskSource* owner_tasks = entry->task_source.get();
345+
TaskSource* subsumed_tasks = queue_entries_.at(subsumed)->task_source.get();
293346

294347
// we are owning another task queue
295-
const bool subsumed_has_task = !subsumed_tasks.empty();
296-
const bool owner_has_task = !owner_tasks.empty();
348+
const bool subsumed_has_task = !subsumed_tasks->IsEmpty();
349+
const bool owner_has_task = !owner_tasks->IsEmpty();
350+
fml::TaskQueueId top_queue_id = owner;
297351
if (owner_has_task && subsumed_has_task) {
298-
const auto owner_task = owner_tasks.top();
299-
const auto subsumed_task = subsumed_tasks.top();
300-
if (owner_task > subsumed_task) {
352+
const auto owner_task = owner_tasks->Top();
353+
const auto subsumed_task = subsumed_tasks->Top();
354+
if (owner_task.task > subsumed_task.task) {
301355
top_queue_id = subsumed;
302356
} else {
303357
top_queue_id = owner;
@@ -307,7 +361,7 @@ const DelayedTask& MessageLoopTaskQueues::PeekNextTaskUnlocked(
307361
} else {
308362
top_queue_id = subsumed;
309363
}
310-
return queue_entries_.at(top_queue_id)->delayed_tasks.top();
364+
return queue_entries_.at(top_queue_id)->task_source->Top();
311365
}
312366

313367
} // namespace fml

0 commit comments

Comments
 (0)