初始化流程
boost::container::static_vector, max_scheduling_groups()> _task_queues;
task_queue_list _active_task_queues;
task_queue_list _activating_task_queues;
add_task
void add_task(task* t) noexcept {
auto sg = t->group();
auto* q = _task_queues[sg._id].get();
bool was_empty = q->_q.empty();
q->_q.push_back(std::move(t));
#ifdef SEASTAR_SHUFFLE_TASK_QUEUE
shuffle(q->_q.back(), *q);
#endif
if (was_empty) {
activate(*q);
}
}
往task队列里添加任务:
void schedule(task* t) noexcept {
engine().add_task(t);
}
void schedule_urgent(task* t) noexcept {
engine().add_urgent_task(t);
}
用户层api使用:
future<>
yield() noexcept {
memory::scoped_critical_alloc_section _;
auto tsk = make_task([] {});
schedule(tsk);
return tsk->get_future();
}
在future函数中,make_ready()方法会调用schedule函数添加_task到task队列中:
template
void promise_base::make_ready() noexcept {
if (_task) {
if (Urgent == urgent::yes) {
::seastar::schedule_urgent(std::exchange(_task, nullptr));
} else {
::seastar::schedule(std::exchange(_task, nullptr));
}
}
}
make_ready是怎么使用的呢?
set_value或者set_exception调用时,future会变得ready,然后continuation会被attach到该future上,该future就被run起来; /// \brief Gets the promise's associated future.
///
/// The future and promise will be remember each other, even if either or
/// both are moved. When \c set_value() or \c set_exception() are called
/// on the promise, the future will be become ready, and if a continuation
/// was attached to the future, it will run.
future get_future() noexcept;
template
template
SEASTAR_CONCEPT( requires std::invocable )
void futurize::satisfy_with_result_of(promise_base_with_type&& pr, Func&& func) {
using ret_t = decltype(func());
if constexpr (std::is_void_v) {
func();
pr.set_value();
} else if constexpr (is_future::value) {
func().forward_to(std::move(pr));
} else {
pr.set_value(func());
}
}
void
reactor::run_some_tasks() {
if (!have_more_tasks()) {
return;
}
sched_print("run_some_tasks: start");
reset_preemption_monitor();
update_lowres_clocks();
sched_clock::time_point t_run_completed = now();
STAP_PROBE(seastar, reactor_run_tasks_start);
_cpu_stall_detector->start_task_run(t_run_completed);
do {
auto t_run_started = t_run_completed;
insert_activating_task_queues();
// 从r
task_queue* tq = pop_active_task_queue(t_run_started);
sched_print("running tq {} {}", (void*)tq, tq->_name);
tq->_current = true;
_last_vruntime = std::max(tq->_vruntime, _last_vruntime);
// 真正的开始run task
run_tasks(*tq);
tq->_current = false;
t_run_completed = now();
auto delta = t_run_completed - t_run_started;
account_runtime(*tq, delta);
sched_print("run complete ({} {}); time consumed {} usec; final vruntime {} empty {}",
(void*)tq, tq->_name, delta / 1us, tq->_vruntime, tq->_q.empty());
tq->_ts = t_run_completed;
if (!tq->_q.empty()) {
insert_active_task_queue(tq);
} else {
tq->_active = false;
}
} while (have_more_tasks() && !need_preempt());
_cpu_stall_detector->end_task_run(t_run_completed);
STAP_PROBE(seastar, reactor_run_tasks_end);
*internal::current_scheduling_group_ptr() = default_scheduling_group(); // Prevent inheritance from last group run
sched_print("run_some_tasks: end");
}