3#include <tess/core/assert.h>
4#include <tess/core/config.h>
5#include <tess/core/fail_fast.h>
6#include <tess/diagnostics/diagnostics.h>
17namespace tess::experimental::maintenance {
20class MaintenanceBudget {
22 explicit constexpr MaintenanceBudget(
23 std::uint64_t units = std::numeric_limits<std::uint64_t>::max()) noexcept
24 : remaining_(units) {}
26 [[nodiscard]]
constexpr auto consume(std::uint64_t units = 1)
noexcept
28 if (units > remaining_) {
35 [[nodiscard]]
constexpr auto remaining()
const noexcept -> std::uint64_t {
40 std::uint64_t remaining_;
44class MaintenanceTask {
46 MaintenanceTask() =
default;
47 MaintenanceTask(
const MaintenanceTask&) =
delete;
48 auto operator=(
const MaintenanceTask&) -> MaintenanceTask& =
delete;
49 MaintenanceTask(MaintenanceTask&&) =
delete;
50 auto operator=(MaintenanceTask&&) -> MaintenanceTask& =
delete;
52 virtual ~MaintenanceTask() {
53 if (registration_epoch_.load(std::memory_order_relaxed) != 0) {
54 ::tess::detail::fail_fast(
55 "MaintenanceTask destroyed while registered; release it or "
56 "destroy its registered scheduler first");
62 template <
typename Backend>
65 std::atomic<std::uint64_t> registration_epoch_ = 0;
70 std::uint64_t schedule_calls = 0;
71 std::uint64_t coalesced_calls = 0;
72 std::uint64_t executions = 0;
73 std::uint64_t capacity_failures = 0;
101 [[nodiscard]]
virtual auto flush() ->
bool = 0;
106 [[nodiscard]] virtual auto
has_pending() const noexcept ->
bool = 0;
113 void record_schedule()
noexcept {
114 schedule_calls_.fetch_add(1, std::memory_order_relaxed);
116 void record_coalesced()
noexcept {
117 coalesced_calls_.fetch_add(1, std::memory_order_relaxed);
119 void record_execution()
noexcept {
120 executions_.fetch_add(1, std::memory_order_relaxed);
122 void record_capacity_failure()
noexcept {
123 capacity_failures_.fetch_add(1, std::memory_order_relaxed);
128 schedule_calls_.load(std::memory_order_relaxed),
129 coalesced_calls_.load(std::memory_order_relaxed),
130 executions_.load(std::memory_order_relaxed),
131 capacity_failures_.load(std::memory_order_relaxed)};
135 std::atomic<std::uint64_t> schedule_calls_ = 0;
136 std::atomic<std::uint64_t> coalesced_calls_ = 0;
137 std::atomic<std::uint64_t> executions_ = 0;
138 std::atomic<std::uint64_t> capacity_failures_ = 0;
142struct QueuedMaintenanceEntry {
144 std::uint64_t admitted_tick = 0;
145 std::size_t queue_slot = std::numeric_limits<std::size_t>::max();
148class BoundedTaskQueue {
150 explicit BoundedTaskQueue(std::size_t capacity) : entries_(capacity) {}
152 [[nodiscard]]
auto push(
MaintenanceTask& task, std::uint64_t admitted_tick,
153 std::size_t* queue_slot =
nullptr)
noexcept ->
bool {
154 if (size_ == entries_.size()) {
157 const auto tail = (head_ + size_) % entries_.size();
158 entries_[tail] = QueuedMaintenanceEntry{&task, admitted_tick, tail};
159 if (queue_slot !=
nullptr) {
166 [[nodiscard]]
auto pop()
noexcept -> QueuedMaintenanceEntry {
170 auto entry = entries_[head_];
171 entries_[head_] = {};
172 head_ = (head_ + 1) % entries_.size();
177 [[nodiscard]]
auto empty()
const noexcept ->
bool {
return size_ == 0; }
179 [[nodiscard]]
auto size()
const noexcept -> std::size_t {
return size_; }
181 [[nodiscard]]
auto oldest_admitted_tick()
const noexcept -> std::uint64_t {
182 return size_ == 0 ? 0 : entries_[head_].admitted_tick;
186 std::vector<QueuedMaintenanceEntry> entries_;
187 std::size_t head_ = 0;
188 std::size_t size_ = 0;
192class PendingTaskIndex {
194 explicit PendingTaskIndex(std::size_t capacity)
195 : buckets_(bucket_count(capacity), npos), nodes_(capacity) {}
197 [[nodiscard]]
auto contains(
const MaintenanceTask& task)
const noexcept
199 if (buckets_.empty()) {
202 for (
auto index = buckets_[home(task)]; index != npos;
203 index = nodes_[index].next) {
204 if (nodes_[index].task == &task) {
212 std::size_t queue_slot)
noexcept ->
bool {
213 if (buckets_.empty() || queue_slot >= nodes_.size() ||
214 nodes_[queue_slot].task !=
nullptr) {
217 const auto bucket = home(task);
218 nodes_[queue_slot] = Node{&task, buckets_[bucket]};
219 buckets_[bucket] = queue_slot;
224 std::size_t queue_slot)
noexcept ->
bool {
225 if (buckets_.empty() || queue_slot >= nodes_.size()) {
228 auto* link = &buckets_[home(task)];
229 while (*link != npos) {
230 if (*link == queue_slot) {
231 *link = nodes_[queue_slot].next;
232 nodes_[queue_slot] = {};
235 link = &nodes_[*link].next;
241 static constexpr auto npos = std::numeric_limits<std::size_t>::max();
245 std::size_t next = npos;
248 [[nodiscard]]
static auto bucket_count(std::size_t capacity)
noexcept
253 const auto maximum = std::numeric_limits<std::size_t>::max();
254 return capacity > maximum - capacity ? capacity : capacity * 2;
259 auto value =
reinterpret_cast<std::uintptr_t
>(&task);
261 value ^= value >> 17u;
262 if constexpr (
sizeof(value) >=
sizeof(std::uint64_t)) {
263 value *= std::uintptr_t{0x9e3779b97f4a7c15ULL};
265 value *= std::uintptr_t{0x9e3779b9U};
267 return static_cast<std::size_t
>(value) % buckets_.size();
270 std::vector<std::size_t> buckets_;
271 std::vector<Node> nodes_;
274template <
bool Coalescing>
277 explicit QueuedScheduler(std::size_t capacity)
278 : queue_(capacity), pending_(Coalescing ? capacity : 0) {}
281 metrics_.record_schedule();
282 const auto lock = std::scoped_lock{queue_mutex_};
285 const auto called_from_task = running_thread_ == std::this_thread::get_id();
286 if constexpr (Coalescing) {
287 if (pending_.contains(task)) {
288 if (called_from_task) {
289 running_task_scheduled_ =
true;
291 metrics_.record_coalesced();
292 if (accounting_ !=
nullptr) {
293 ++accounting_->counters.offered;
294 ++accounting_->counters.coalesced_into_pending;
299 const auto admitted_tick =
300 accounting_ !=
nullptr ? accounting_->last_observed_tick : 0;
301 auto queue_slot = std::size_t{0};
302 if (!queue_.push(task, admitted_tick, &queue_slot)) {
303 metrics_.record_capacity_failure();
304 if (accounting_ !=
nullptr) {
305 ++accounting_->counters.offered;
306 ++accounting_->counters.rejected;
310 if constexpr (Coalescing) {
311 const auto inserted = pending_.insert(task, queue_slot);
312 TESS_ASSERT(inserted);
313 static_cast<void>(inserted);
315 if (called_from_task) {
316 running_task_scheduled_ =
true;
318 if (accounting_ !=
nullptr) {
319 ++accounting_->counters.offered;
320 accounting_->record_admitted();
326 const auto run_lock = std::scoped_lock{run_mutex_};
327 if (accounting_ !=
nullptr &&
328 budget.remaining() != std::numeric_limits<std::uint64_t>::max()) {
329 const auto lock = std::scoped_lock{queue_mutex_};
330 accounting_->counters.offered_work_units += budget.remaining();
332 while (budget.remaining() != 0) {
333 auto entry = detail::QueuedMaintenanceEntry{};
335 const auto queue_lock = std::scoped_lock{queue_mutex_};
336 entry = pop_pending();
338 if (entry.task ==
nullptr) {
341 if (!run_task(entry, budget)) {
348 [[nodiscard]]
auto flush() ->
bool override {
349 const auto run_lock = std::scoped_lock{run_mutex_};
354 auto entry = detail::QueuedMaintenanceEntry{};
356 const auto queue_lock = std::scoped_lock{queue_mutex_};
357 entry = pop_pending();
359 if (entry.task ==
nullptr) {
362 if (!run_task(entry, budget)) {
369 return metrics_.snapshot();
372 [[nodiscard]]
auto has_pending()
const noexcept ->
bool override {
373 const auto lock = std::scoped_lock{queue_mutex_};
374 return !queue_.empty() || running_active_;
385 const auto run_lock = std::scoped_lock{run_mutex_};
386 const auto lock = std::scoped_lock{queue_mutex_};
387 TESS_ASSERT(queue_.empty());
388 accounting_ = accounting;
393 void observe_flow_tick(std::uint64_t tick) {
394 const auto lock = std::scoped_lock{queue_mutex_};
395 if (accounting_ ==
nullptr) {
399 const auto now = accounting_->last_observed_tick;
402 if (!queue_.empty()) {
404 oldest = queue_.oldest_admitted_tick();
406 if (running_active_ && running_admitted_tick_ < oldest) {
408 oldest = running_admitted_tick_;
410 if (running_active_ && !any) {
413 accounting_->counters.oldest_outstanding_age_ticks = any ? now - oldest : 0;
417 [[nodiscard]]
auto run_task(detail::QueuedMaintenanceEntry entry,
419 auto& task = *entry.task;
420 metrics_.record_execution();
421 const auto before = budget.remaining();
423 const auto lock = std::scoped_lock{queue_mutex_};
424 running_thread_ = std::this_thread::get_id();
425 running_task_scheduled_ =
false;
426 running_active_ =
true;
427 running_admitted_tick_ = entry.admitted_tick;
429#if TESS_HAS_EXCEPTIONS
438 const auto lock = std::scoped_lock{queue_mutex_};
439 running_thread_ = {};
440 running_task_scheduled_ =
false;
441 running_active_ =
false;
442 account_terminal(entry, before, budget.remaining(),
false);
448 auto scheduled_follow_up =
false;
450 const auto lock = std::scoped_lock{queue_mutex_};
451 scheduled_follow_up = running_task_scheduled_;
452 running_thread_ = {};
453 running_task_scheduled_ =
false;
454 running_active_ =
false;
455 account_terminal(entry, before, budget.remaining(),
true);
461 return budget.remaining() != before || !scheduled_follow_up;
466 [[nodiscard]]
auto pop_pending()
noexcept -> QueuedMaintenanceEntry {
467 auto entry = queue_.pop();
468 if constexpr (Coalescing) {
469 if (entry.task !=
nullptr) {
470 const auto erased = pending_.erase(*entry.task, entry.queue_slot);
472 static_cast<void>(erased);
481 void account_terminal(
const detail::QueuedMaintenanceEntry& entry,
482 std::uint64_t budget_before, std::uint64_t budget_after,
483 bool completed)
noexcept {
484 if (accounting_ ==
nullptr) {
487 auto& counters = accounting_->counters;
488 counters.consumed_work_units += budget_before - budget_after;
489 ++(completed ? counters.completed : counters.failed);
490 accounting_->record_left_outstanding();
491 counters.residence_ticks_accumulated +=
492 accounting_->last_observed_tick - entry.admitted_tick;
495 mutable std::mutex queue_mutex_;
496 std::mutex run_mutex_;
497 BoundedTaskQueue queue_;
498 PendingTaskIndex pending_;
503 ~QueuedScheduler()
override {
504 if (accounting_ ==
nullptr) {
508 const auto entry = pop_pending();
509 if (entry.task ==
nullptr) {
512 ++accounting_->counters.dropped_after_admission;
513 accounting_->record_left_outstanding();
514 accounting_->counters.residence_ticks_accumulated +=
515 accounting_->last_observed_tick - entry.admitted_tick;
520 MetricsStore metrics_;
521 std::thread::id running_thread_;
522 bool running_task_scheduled_ =
false;
523 bool running_active_ =
false;
524 std::uint64_t running_admitted_tick_ = 0;
533 explicit ImmediateScheduler(std::size_t = 0) {}
539 const auto run_lock = std::scoped_lock{run_mutex_};
540 metrics_.record_schedule();
541 for (
auto* active = active_run_; active !=
nullptr;
542 active = active->parent) {
543 if (active->task != &task) {
546 if (active->pending == std::numeric_limits<std::uint64_t>::max()) {
547 metrics_.record_capacity_failure();
548 account_offer(Offer::Rejected);
552 account_offer(Offer::Coalesced);
559 auto active = ActiveRun{&task, 1, active_run_};
560 struct ActiveRunGuard {
563 ~ActiveRunGuard() { current = previous; }
565 active_run_ = &active;
569 const auto guard = ActiveRunGuard{active_run_, active.parent};
570 account_offer(Offer::Admitted);
572 auto completed_ok =
true;
573 auto consumed_total = std::uint64_t{0};
574 while (active.pending != 0) {
576 const auto before = budget.remaining();
577 metrics_.record_execution();
578#if TESS_HAS_EXCEPTIONS
582 consumed_total += before - budget.remaining();
583 account_run_terminal(consumed_total,
true);
589 consumed_total += before - budget.remaining();
590 if (active.pending != 0 && budget.remaining() == before) {
591 completed_ok =
false;
597 account_run_terminal(consumed_total,
false);
604 [[nodiscard]]
auto flush() ->
bool override {
return true; }
606 return metrics_.snapshot();
609 [[nodiscard]]
auto has_pending() const noexcept ->
bool override {
610 const auto run_lock = std::scoped_lock{run_mutex_};
611 return active_run_ !=
nullptr;
621 const auto run_lock = std::scoped_lock{run_mutex_};
622 TESS_ASSERT(active_run_ ==
nullptr);
623 accounting_ = accounting;
629 const auto run_lock = std::scoped_lock{run_mutex_};
630 if (accounting_ !=
nullptr) {
631 accounting_->observe_tick(tick);
632 accounting_->counters.oldest_outstanding_age_ticks = 0;
639 std::uint64_t pending = 0;
640 ActiveRun* parent =
nullptr;
643 enum class Offer : std::uint8_t { Admitted, Rejected, Coalesced };
645 void account_offer(Offer offer)
noexcept {
646 if (accounting_ ==
nullptr) {
649 ++accounting_->counters.offered;
651 case Offer::Admitted:
652 accounting_->record_admitted();
654 case Offer::Rejected:
655 ++accounting_->counters.rejected;
657 case Offer::Coalesced:
658 ++accounting_->counters.coalesced_into_pending;
663 void account_run_terminal(std::uint64_t consumed,
bool failed)
noexcept {
664 if (accounting_ ==
nullptr) {
667 accounting_->counters.consumed_work_units += consumed;
668 ++(failed ? accounting_->counters.failed : accounting_->counters.completed);
669 accounting_->record_left_outstanding();
672 detail::MetricsStore metrics_;
673 mutable std::recursive_mutex run_mutex_;
674 ActiveRun* active_run_ =
nullptr;
675 diagnostics::FlowAccounting* accounting_ =
nullptr;
693 explicit DirtyBitScheduler(std::size_t capacity)
694 : tasks_(capacity,
nullptr),
695 registrations_(index_capacity(capacity)),
696 pending_(word_count(capacity)) {
697 for (
auto& word : pending_) {
698 word.store(0, std::memory_order_relaxed);
704 const auto setup_lock = std::scoped_lock{setup_mutex_};
705 if (sealed_.load(std::memory_order_relaxed)) {
708 if (find_index(task) != npos) {
711 if (registered_ == tasks_.size()) {
712 metrics_.record_capacity_failure();
715 const auto index = registered_++;
716 tasks_[index] = &task;
717 auto slot = home(task);
718 while (registrations_[slot].task !=
nullptr) {
719 slot = next_registration(slot);
721 registrations_[slot] = Registration{&task, index};
727 const auto setup_lock = std::scoped_lock{setup_mutex_};
728 sealed_.store(
true, std::memory_order_release);
732 metrics_.record_schedule();
733 if (!sealed_.load(std::memory_order_acquire)) {
736 const auto index = find_index(task);
740 const auto mask = std::uint64_t{1} << (index % 64u);
741 const auto previous =
742 pending_[index / 64u].fetch_or(mask, std::memory_order_release);
743 if ((previous & mask) != 0) {
744 metrics_.record_coalesced();
746 if (active_scheduler_ ==
this) {
747 running_task_scheduled_ =
true;
753 if (!sealed_.load(std::memory_order_acquire)) {
756 const auto run_lock = std::scoped_lock{run_mutex_};
757 return drain(budget);
760 [[nodiscard]]
auto flush() ->
bool override {
761 if (!sealed_.load(std::memory_order_acquire)) {
764 const auto run_lock = std::scoped_lock{run_mutex_};
766 return drain(budget);
770 return metrics_.snapshot();
773 [[nodiscard]]
auto has_pending() const noexcept ->
bool override {
774 for (
const auto& word : pending_) {
775 if (word.load(std::memory_order_acquire) != 0) {
783 struct Registration {
785 std::size_t index = 0;
788 static constexpr auto npos = std::numeric_limits<std::size_t>::max();
790 [[nodiscard]]
static auto index_capacity(std::size_t capacity)
noexcept
795 const auto maximum = std::numeric_limits<std::size_t>::max();
796 return capacity > maximum - capacity ? capacity : capacity * 2;
799 [[nodiscard]]
static auto word_count(std::size_t capacity)
noexcept
801 return capacity / 64u + (capacity % 64u == 0 ? 0u : 1u);
804 [[nodiscard]]
auto home(
const MaintenanceTask& task)
const noexcept
806 auto value =
reinterpret_cast<std::uintptr_t
>(&task);
808 value ^= value >> 17u;
809 if constexpr (
sizeof(value) >=
sizeof(std::uint64_t)) {
810 value *= std::uintptr_t{0x9e3779b97f4a7c15ULL};
812 value *= std::uintptr_t{0x9e3779b9U};
814 return static_cast<std::size_t
>(value) % registrations_.size();
817 [[nodiscard]]
auto next_registration(std::size_t slot)
const noexcept
819 return slot + 1 == registrations_.size() ? 0 : slot + 1;
822 [[nodiscard]]
auto find_index(
const MaintenanceTask& task)
const noexcept
824 if (registrations_.empty()) {
827 auto slot = home(task);
828 for (std::size_t probed = 0; probed < registrations_.size(); ++probed) {
829 if (registrations_[slot].task ==
nullptr) {
832 if (registrations_[slot].task == &task) {
833 return registrations_[slot].index;
835 slot = next_registration(slot);
840 [[nodiscard]]
auto drain(MaintenanceBudget& budget) ->
bool {
842 auto executed =
false;
843 for (std::size_t word_index = 0; word_index < pending_.size();
845 auto word = pending_[word_index].exchange(0, std::memory_order_acquire);
847 if (budget.remaining() == 0) {
848 pending_[word_index].fetch_or(word, std::memory_order_release);
851 const auto bit =
static_cast<std::size_t
>(std::countr_zero(word));
852 const auto index = word_index * 64u + bit;
853 const auto mask = std::uint64_t{1} << bit;
856#if TESS_HAS_EXCEPTIONS
858 if (!run_task(*tasks_[index], budget)) {
860 pending_[word_index].fetch_or(word, std::memory_order_release);
866 pending_[word_index].fetch_or(word, std::memory_order_release);
871 if (!run_task(*tasks_[index], budget)) {
873 pending_[word_index].fetch_or(word, std::memory_order_release);
886 [[nodiscard]]
auto run_task(MaintenanceTask& task, MaintenanceBudget& budget)
888 metrics_.record_execution();
889 const auto before = budget.remaining();
890 struct ActiveRunGuard {
891 DirtyBitScheduler*& active;
892 DirtyBitScheduler* previous;
893 ~ActiveRunGuard() { active = previous; }
895 const auto previous = active_scheduler_;
896 active_scheduler_ =
this;
897 const auto guard = ActiveRunGuard{active_scheduler_, previous};
898 running_task_scheduled_ =
false;
900 return budget.remaining() != before || !running_task_scheduled_;
903 detail::MetricsStore metrics_;
904 std::mutex setup_mutex_;
905 std::mutex run_mutex_;
906 std::vector<MaintenanceTask*> tasks_;
907 std::vector<Registration> registrations_;
908 std::vector<std::atomic<std::uint64_t>> pending_;
909 std::size_t registered_ = 0;
910 std::atomic<bool> sealed_ =
false;
911 bool running_task_scheduled_ =
false;
912 inline static thread_local DirtyBitScheduler* active_scheduler_ =
nullptr;
916using FifoScheduler = detail::QueuedScheduler<false>;
919using CoalescingScheduler = detail::QueuedScheduler<true>;
auto flush() -> bool override
Completes all reachable work; see the non-reentrant drain contract above.
Definition maintenance.h:760
void seal()
Definition maintenance.h:726
auto schedule(MaintenanceTask &task) -> bool override
Definition maintenance.h:731
auto run_some(MaintenanceBudget budget) -> bool override
Definition maintenance.h:752
auto has_pending() const noexcept -> bool override
Definition maintenance.h:773
auto register_task(MaintenanceTask &task) -> bool
Definition maintenance.h:703
auto schedule(MaintenanceTask &task) -> bool override
Definition maintenance.h:535
auto run_some(MaintenanceBudget) -> bool override
Definition maintenance.h:601
auto flush() -> bool override
Completes all reachable work; see the non-reentrant drain contract above.
Definition maintenance.h:604
void observe_flow_tick(std::uint64_t tick)
Definition maintenance.h:628
auto has_pending() const noexcept -> bool override
Definition maintenance.h:609
void set_flow_accounting(diagnostics::FlowAccounting *accounting)
Definition maintenance.h:620
Shared unit budget passed through one maintenance drain.
Definition maintenance.h:20
Backend-neutral experimental maintenance scheduler interface.
Definition maintenance.h:77
virtual auto has_pending() const noexcept -> bool=0
virtual auto flush() -> bool=0
Completes all reachable work; see the non-reentrant drain contract above.
virtual auto run_some(MaintenanceBudget budget) -> bool=0
virtual auto schedule(MaintenanceTask &task) -> bool=0
Long-lived derived-state maintenance operation.
Definition maintenance.h:44
friend class RegisteredScheduler
Definition maintenance.h:63
Definition diagnostics.h:505
void observe_tick(std::uint64_t tick) noexcept
Definition diagnostics.h:513
Scheduler observations used by experiments and diagnostics.
Definition maintenance.h:69