199class RegisteredScheduler final {
201 detail::is_raw_backend_v<Backend>,
202 "RegisteredScheduler requires a MaintenanceBackend or a "
203 "built-in experimental raw backend");
205 static constexpr bool has_backend_register =
requires(
206 Backend& backend,
MaintenanceTask& task) { backend.register_task(task); };
207 static constexpr bool has_backend_seal =
208 requires(Backend& backend) { backend.seal(); };
210 detail::is_dirty_bit_v<Backend> ||
211 ((!has_backend_register && !has_backend_seal) ||
213 "a custom maintenance backend must provide both no-throw "
214 "register_task(MaintenanceTask&) -> bool and seal() -> void hooks, or "
221 if (target ==
nullptr) {
222 ::tess::detail::fail_fast(
223 "RegisteredScheduler invoked a released maintenance slot");
231 std::uint64_t generation = 1;
234 class OperationGuard {
236 explicit OperationGuard(RegisteredScheduler& scheduler)
237 : previous_(detail::active_registered_scheduler_epoch) {
238 if (previous_ == scheduler.owner_epoch_) {
241 if (previous_ != 0) {
242 ::tess::detail::fail_fast(
243 "RegisteredScheduler nested cross-scheduler operation called "
244 "from a running task");
246 lock_ = std::shared_lock<std::shared_mutex>{scheduler.lifecycle_mutex_};
247 detail::active_registered_scheduler_epoch = scheduler.owner_epoch_;
251 OperationGuard(
const OperationGuard&) =
delete;
252 auto operator=(
const OperationGuard&) -> OperationGuard& =
delete;
256 detail::active_registered_scheduler_epoch = previous_;
261 std::uint64_t previous_ = 0;
262 std::shared_lock<std::shared_mutex> lock_{};
266 class ScheduleActivityGuard {
268 explicit ScheduleActivityGuard(std::atomic<std::uint64_t>& in_flight,
269 std::mutex& observation_mutex)
270 : in_flight_(in_flight), observation_mutex_(observation_mutex) {
271 const auto lock = std::scoped_lock{observation_mutex_};
272 const auto previous = in_flight_.fetch_add(1, std::memory_order_acq_rel);
273 if (previous == std::numeric_limits<std::uint64_t>::max()) {
274 ::tess::detail::fail_fast(
275 "RegisteredScheduler exhausted in-flight schedule count");
279 ScheduleActivityGuard(
const ScheduleActivityGuard&) =
delete;
280 auto operator=(
const ScheduleActivityGuard&)
281 -> ScheduleActivityGuard& =
delete;
283 ~ScheduleActivityGuard() {
284 const auto lock = std::scoped_lock{observation_mutex_};
285 const auto previous = in_flight_.fetch_sub(1, std::memory_order_acq_rel);
287 ::tess::detail::fail_fast(
288 "RegisteredScheduler in-flight schedule count underflowed");
293 std::atomic<std::uint64_t>& in_flight_;
294 std::mutex& observation_mutex_;
298 explicit RegisteredScheduler(std::size_t capacity)
299 : RegisteredScheduler(capacity, capacity) {}
301 RegisteredScheduler(std::size_t registry_capacity,
302 std::size_t backend_capacity)
303 : owner_epoch_(detail::claim_owner_epoch()),
304 capacity_(registry_capacity),
305 slots_(std::make_unique<Slot[]>(registry_capacity)),
306 backend_(backend_capacity) {}
308 RegisteredScheduler(
const RegisteredScheduler&) =
delete;
309 auto operator=(
const RegisteredScheduler&) -> RegisteredScheduler& =
delete;
310 RegisteredScheduler(RegisteredScheduler&&) =
delete;
311 auto operator=(RegisteredScheduler&&) -> RegisteredScheduler& =
delete;
313 ~RegisteredScheduler() {
314 for (std::size_t index = 0; index < capacity_; ++index) {
315 auto* task = slots_[index].task.target;
316 if (task ==
nullptr) {
319 auto expected = owner_epoch_;
320 if (!task->registration_epoch_.compare_exchange_strong(
321 expected, 0, std::memory_order_relaxed)) {
322 ::tess::detail::fail_fast(
323 "RegisteredScheduler task ownership changed unexpectedly");
325 slots_[index].task.target =
nullptr;
331 -> std::optional<MaintenanceHandle> {
332 reject_reentrant_lifecycle(
"register_task");
333 const auto lock = std::unique_lock<std::shared_mutex>{lifecycle_mutex_};
335 ::tess::detail::fail_fast(
336 "RegisteredScheduler::register_task called after seal");
338 for (std::size_t index = 0; index < capacity_; ++index) {
339 if (slots_[index].task.target == &task) {
340 return make_handle(index);
343 auto index = capacity_;
344 for (std::size_t candidate = 0; candidate < capacity_; ++candidate) {
345 if (slots_[candidate].task.target ==
nullptr) {
350 if (index == capacity_) {
353 auto expected = std::uint64_t{0};
354 if (!task.registration_epoch_.compare_exchange_strong(
355 expected, owner_epoch_, std::memory_order_relaxed)) {
356 ::tess::detail::fail_fast(
357 "RegisteredScheduler::register_task task belongs to another "
358 "scheduler or registration epoch");
360 slots_[index].task.target = &task;
361 return make_handle(index);
366 reject_reentrant_lifecycle(
"seal");
367 const auto lock = std::unique_lock<std::shared_mutex>{lifecycle_mutex_};
369 ::tess::detail::fail_fast(
"RegisteredScheduler::seal called twice");
372 detail::is_dirty_bit_v<Backend>) {
373 for (std::size_t index = 0; index < capacity_; ++index) {
374 if (slots_[index].task.target !=
nullptr &&
375 !backend_.register_task(slots_[index].task)) {
376 ::tess::detail::fail_fast(
377 "RegisteredScheduler fixed backend registry capacity mismatch");
387 auto& self =
const_cast<RegisteredScheduler&
>(*this);
388 const auto operation = OperationGuard{self};
389 return self.resolve(handle) !=
nullptr;
394 -> std::optional<ScheduleResult> {
395 const auto operation = OperationGuard{*
this};
396 auto* slot = resolve(handle);
397 if (slot ==
nullptr) {
400 require_sealed(
"try_schedule");
401 const auto schedule_activity =
402 ScheduleActivityGuard{in_flight_schedules_, observation_mutex_};
403#if TESS_HAS_EXCEPTIONS
405 const auto result = schedule_slot(*slot);
406 if (result != ScheduleResult::CapacityExhausted) {
407 invalidate_idle_observation();
411 invalidate_idle_observation();
415 const auto result = schedule_slot(*slot);
416 if (result != ScheduleResult::CapacityExhausted) {
417 invalidate_idle_observation();
426 if (!result.has_value()) {
427 ::tess::detail::fail_fast(
428 "RegisteredScheduler::schedule received a stale handle from the "
429 "wrong scheduler or registration epoch; use try_schedule for "
430 "expected uncertainty");
437 reject_reentrant_drain(
"run_some");
438 const auto operation = OperationGuard{*
this};
439 require_sealed(
"run_some");
440 const auto drain_lock = std::scoped_lock{drain_mutex_};
441 return observe_drain([&] {
return run_backend_some(budget); });
445 [[nodiscard]]
auto flush() -> DrainResult {
446 reject_reentrant_drain(
"flush");
447 const auto operation = OperationGuard{*
this};
448 require_sealed(
"flush");
449 const auto drain_lock = std::scoped_lock{drain_mutex_};
450 return observe_drain([&] {
return flush_backend(); });
458 reject_reentrant_lifecycle(
"try_release");
459 const auto lock = std::unique_lock<std::shared_mutex>{lifecycle_mutex_};
460 auto* slot = resolve(handle);
461 if (slot ==
nullptr) {
462 return ReleaseResult::InvalidHandle;
464 if (sealed_ && (idle_epoch_.load(std::memory_order_acquire) !=
465 activity_epoch_.load(std::memory_order_acquire) ||
466 backend_.has_pending())) {
467 return ReleaseResult::NotIdle;
469 auto* task = slot->task.target;
470 auto expected = owner_epoch_;
471 if (!task->registration_epoch_.compare_exchange_strong(
472 expected, 0, std::memory_order_relaxed)) {
473 ::tess::detail::fail_fast(
474 "RegisteredScheduler::try_release task ownership changed");
476 slot->task.target =
nullptr;
478 if (slot->generation == 0) {
481 return ReleaseResult::Released;
487 case ReleaseResult::Released:
489 case ReleaseResult::InvalidHandle:
490 ::tess::detail::fail_fast(
491 "RegisteredScheduler::release received a stale maintenance "
492 "handle; use try_release for expected uncertainty");
493 case ReleaseResult::NotIdle:
494 ::tess::detail::fail_fast(
495 "RegisteredScheduler::release requires a positive Idle result "
496 "after the last successful schedule; release never cancels");
498 ::tess::detail::fail_fast(
"RegisteredScheduler::release invalid result");
503 return backend_.metrics();
507 [[nodiscard]]
auto make_handle(std::size_t index)
const noexcept
512 [[nodiscard]]
auto resolve(MaintenanceHandle handle)
noexcept -> Slot* {
513 if (handle.owner_epoch_ != owner_epoch_ || handle.slot_ >= capacity_) {
516 auto& slot = slots_[handle.slot_];
517 if (slot.task.target ==
nullptr ||
518 slot.generation != handle.slot_generation_) {
524 void require_sealed(
const char* operation)
const {
528 if (operation ==
nullptr) {
529 ::tess::detail::fail_fast(
"RegisteredScheduler used before seal");
531 ::tess::detail::fail_fast(
532 "RegisteredScheduler operation called before seal");
535 void reject_reentrant_lifecycle(
const char* operation)
const {
536 if (detail::active_registered_scheduler_epoch == 0) {
539 static_cast<void>(operation);
540 ::tess::detail::fail_fast(
541 "RegisteredScheduler lifecycle mutation called from a running task");
544 void reject_reentrant_drain(
const char* operation)
const {
545 if (detail::active_registered_scheduler_epoch != owner_epoch_) {
548 static_cast<void>(operation);
549 ::tess::detail::fail_fast(
550 "RegisteredScheduler drain called from a running task");
553 void invalidate_idle_observation() noexcept {
554 const auto previous =
555 activity_epoch_.fetch_add(1, std::memory_order_acq_rel);
556 if (previous == std::numeric_limits<std::uint64_t>::max()) {
557 ::tess::detail::fail_fast(
558 "RegisteredScheduler exhausted activity epochs");
560 idle_epoch_.store(std::numeric_limits<std::uint64_t>::max(),
561 std::memory_order_release);
564 [[nodiscard]]
auto schedule_slot(Slot& slot) -> ScheduleResult {
565 if constexpr (MaintenanceBackend<Backend>) {
566 return backend_.schedule(slot.task);
568 const auto accepted = backend_.schedule(slot.task);
570 return ScheduleResult::Accepted;
572 if constexpr (detail::is_immediate_v<Backend>) {
573 return ScheduleResult::Stalled;
575 if constexpr (detail::is_dirty_bit_v<Backend>) {
576 ::tess::detail::fail_fast(
577 "RegisteredScheduler dirty-bit backend rejected a live sealed "
580 return ScheduleResult::CapacityExhausted;
584 [[nodiscard]]
auto run_backend_some(MaintenanceBudget budget)
585 -> BackendDrainResult {
586 if constexpr (MaintenanceBackend<Backend>) {
587 return backend_.run_some(budget);
589 return backend_.run_some(budget) ? BackendDrainResult::Completed
590 : BackendDrainResult::Stalled;
594 [[nodiscard]]
auto flush_backend() -> BackendDrainResult {
595 if constexpr (MaintenanceBackend<Backend>) {
596 return backend_.flush();
598 return backend_.flush() ? BackendDrainResult::Completed
599 : BackendDrainResult::Stalled;
603 template <
typename Drain>
604 [[nodiscard]]
auto observe_drain(Drain&& drain) -> DrainResult {
605 idle_epoch_.store(std::numeric_limits<std::uint64_t>::max(),
606 std::memory_order_release);
607 const auto activity_before =
608 activity_epoch_.load(std::memory_order_acquire);
609 const auto in_flight_before =
610 in_flight_schedules_.load(std::memory_order_acquire);
611 const auto pending_before = backend_.has_pending();
612 const auto completed = drain();
613 const auto pending_after = backend_.has_pending();
614 const auto observation_lock = std::scoped_lock{observation_mutex_};
615 const auto in_flight_after =
616 in_flight_schedules_.load(std::memory_order_acquire);
617 const auto activity_after = activity_epoch_.load(std::memory_order_acquire);
618 if (completed == BackendDrainResult::Stalled) {
619 return DrainResult::Stalled;
622 return DrainResult::BudgetExhausted;
624 if (!pending_before && in_flight_before == 0 && in_flight_after == 0 &&
625 activity_before == activity_after) {
626 idle_epoch_.store(activity_after, std::memory_order_release);
627 return DrainResult::Idle;
629 return DrainResult::Drained;
632 const std::uint64_t owner_epoch_;
633 const std::size_t capacity_;
634 std::unique_ptr<Slot[]> slots_;
636 mutable std::shared_mutex lifecycle_mutex_;
637 std::mutex observation_mutex_;
638 std::mutex drain_mutex_;
639 std::atomic<std::uint64_t> activity_epoch_ = 0;
640 std::atomic<std::uint64_t> in_flight_schedules_ = 0;
641 std::atomic<std::uint64_t> idle_epoch_ =
642 std::numeric_limits<std::uint64_t>::max();
643 bool sealed_ =
false;