Skip to content

Commit 7f05565

Browse files
authored
Merge pull request #295 from elbeno/urgent-triggers
✨ Add `urgent_trigger_scheduler`
2 parents 3e3c8aa + e006e68 commit 7f05565

5 files changed

Lines changed: 83 additions & 9 deletions

File tree

docs/schedulers.adoc

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -407,3 +407,30 @@ NOTE: It is possible to use a `trigger_scheduler` that takes arguments as a
407407
"normal" scheduler, i.e. functions like `start_on` will work; however the
408408
arguments passed to `run_triggers` will be discarded when used with constructs
409409
like `start_on(trigger_scheduler<"name", int>{}, just(42))`.
410+
411+
Normally, calling `run_triggers` runs all the work in order -- and there can be
412+
multiple senders waiting to run on a given trigger. Occasionally it is useful to
413+
have a sender jump to the head of the queue. To do this use an `urgent_trigger_scheduler`:
414+
415+
[source,cpp]
416+
----
417+
int x{};
418+
async::trigger_scheduler<"name", int>{}.schedule()
419+
| async::then([&] (auto i) { x = i; })
420+
| async::start_detached();
421+
422+
// this sender will be queued before the one above
423+
async::urgent_trigger_scheduler<"name", int>{}.schedule()
424+
| async::then([&] (auto i) { x = i*2; })
425+
| async::start_detached();
426+
427+
// when event occurs
428+
async::run_one_trigger<"name">(42);
429+
// x is now 84
430+
async::run_triggers<"name">(12);
431+
// x is now 12
432+
----
433+
434+
This is not a general-purpose priority management mechanism, but it can be
435+
useful to prevent priority inversion when integrating triggers with other work
436+
that uses fixed interrupt priorities.

include/async/incite_on.hpp

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -178,7 +178,8 @@ template <sender S, typename Sched>
178178
namespace _incite_on {
179179
template <typename Uniq>
180180
class trigger_scheduler
181-
: public trigger_mgr::scheduler<trigger_scheduler<Uniq>, Uniq> {
181+
: public trigger_mgr::scheduler<trigger_scheduler<Uniq>, Uniq,
182+
trigger_mgr::queue_at_back> {
182183
[[nodiscard]] friend constexpr auto operator==(trigger_scheduler,
183184
trigger_scheduler)
184185
-> bool = default;

include/async/schedulers/trigger_manager.hpp

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,20 @@ struct all {
7777
};
7878
} // namespace run_policy
7979

80+
namespace trigger_mgr {
81+
struct queue_at_front {
82+
constexpr static auto push(auto &list, auto &task) -> void {
83+
list.push_front(std::addressof(task));
84+
}
85+
};
86+
87+
struct queue_at_back {
88+
constexpr static auto push(auto &list, auto &task) -> void {
89+
list.push_back(std::addressof(task));
90+
}
91+
};
92+
} // namespace trigger_mgr
93+
8094
template <typename Name, typename... Args> struct trigger_manager {
8195
using task_t = trigger_task<Args...>;
8296

@@ -86,12 +100,12 @@ template <typename Name, typename... Args> struct trigger_manager {
86100
stdx::atomic<int> task_count;
87101

88102
public:
89-
auto enqueue(task_t &t) -> bool {
103+
template <typename QueuePolicy> auto enqueue(task_t &t) -> bool {
90104
return conc::call_in_critical_section<mutex>([&]() -> bool {
91105
auto const added = not std::exchange(t.pending, true);
92106
if (added) {
93107
++task_count;
94-
tasks[0].push_back(std::addressof(t));
108+
QueuePolicy::push(tasks[0], t);
95109
}
96110
return added;
97111
});

include/async/schedulers/trigger_scheduler.hpp

Lines changed: 19 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -67,9 +67,10 @@ struct op_state_base<Rcvr, Ops> {
6767
auto clear_stop_cb() -> void {}
6868
};
6969

70-
template <typename Name, typename Rcvr, typename... Args>
71-
struct op_state final : op_state_base<Rcvr, op_state<Name, Rcvr, Args...>>,
72-
trigger_task<Args...> {
70+
template <typename Name, typename Rcvr, typename QueuePolicy, typename... Args>
71+
struct op_state final
72+
: op_state_base<Rcvr, op_state<Name, Rcvr, QueuePolicy, Args...>>,
73+
trigger_task<Args...> {
7374
template <stdx::same_as_unqualified<Rcvr> R>
7475
// NOLINTNEXTLINE(bugprone-forwarding-reference-overload)
7576
constexpr explicit(true) op_state(R &&r) : rcvr{std::forward<R>(r)} {}
@@ -93,7 +94,7 @@ struct op_state final : op_state_base<Rcvr, op_state<Name, Rcvr, Args...>>,
9394
complete_stopped();
9495
return;
9596
}
96-
triggers<Name, Args...>.enqueue(*this);
97+
triggers<Name, Args...>.template enqueue<QueuePolicy>(*this);
9798
this->emplace_stop_cb();
9899
}
99100

@@ -111,7 +112,8 @@ struct op_state final : op_state_base<Rcvr, op_state<Name, Rcvr, Args...>>,
111112
[[no_unique_address]] Rcvr rcvr;
112113
};
113114

114-
template <typename S, typename Name, typename... Args> class scheduler {
115+
template <typename S, typename Name, typename QueuePolicy, typename... Args>
116+
class scheduler {
115117
struct sender {
116118
using is_sender = void;
117119

@@ -127,7 +129,8 @@ template <typename S, typename Name, typename... Args> class scheduler {
127129
template <receiver R>
128130
[[nodiscard]] constexpr auto connect(R &&r) const {
129131
check_connect<sender, R>();
130-
return trigger_mgr::op_state<Name, std::remove_cvref_t<R>, Args...>{
132+
return trigger_mgr::op_state<Name, std::remove_cvref_t<R>,
133+
QueuePolicy, Args...>{
131134
std::forward<R>(r)};
132135
}
133136
};
@@ -140,6 +143,7 @@ template <typename S, typename Name, typename... Args> class scheduler {
140143
};
141144
} // namespace trigger_mgr
142145

146+
namespace detail {
143147
template <stdx::ct_string Name, typename... Args>
144148
class trigger_scheduler
145149
: public trigger_mgr::scheduler<trigger_scheduler<Name, Args...>,
@@ -148,6 +152,15 @@ class trigger_scheduler
148152
trigger_scheduler)
149153
-> bool = default;
150154
};
155+
} // namespace detail
156+
157+
template <stdx::ct_string Name, typename... Args>
158+
using trigger_scheduler =
159+
detail::trigger_scheduler<Name, trigger_mgr::queue_at_back, Args...>;
160+
161+
template <stdx::ct_string Name, typename... Args>
162+
using urgent_trigger_scheduler =
163+
detail::trigger_scheduler<Name, trigger_mgr::queue_at_front, Args...>;
151164

152165
struct trigger_scheduler_sender_t;
153166

test/schedulers/trigger_scheduler.cpp

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -355,3 +355,22 @@ TEST_CASE("thread safety for immediate execution", "[trigger_scheduler]") {
355355
CHECK(var2 == 42);
356356
CHECK(async::triggers<stdx::cts_t<"rqp_imm">>.empty());
357357
}
358+
359+
TEMPLATE_TEST_CASE(
360+
"urgent_trigger_scheduler puts a task at the front of the queue",
361+
"[trigger_scheduler]", decltype([] {})) {
362+
constexpr auto name = type_string<TestType>;
363+
auto s = async::urgent_trigger_scheduler<name>{};
364+
365+
int var1{};
366+
async::sender auto sndr1 =
367+
async::start_on(s, async::just_result_of([&] { var1 *= 2; }));
368+
CHECK(async::start_detached(sndr1));
369+
async::sender auto sndr2 =
370+
async::start_on(s, async::just_result_of([&] { var1 += 42; }));
371+
CHECK(async::start_detached(sndr2));
372+
373+
async::run_triggers<name>();
374+
CHECK(var1 == 84);
375+
CHECK(async::triggers<stdx::cts_t<name>>.empty());
376+
}

0 commit comments

Comments
 (0)