Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Fixed

- **Custom autoscaler trigger packs**: `autoscaler<Scheduler, Triggers...>` now
constructs and dispatches each configured trigger correctly instead of
failing template instantiation because the trigger argument was omitted.
Unit coverage now exercises all three documented custom trigger forms
together (#1020).
- **Non-blocking epoll connect precondition**: the epoll backend now rejects
`async_connect()` on a blocking socket with `-EINVAL` before `connect(2)` can
block a scheduler worker. Direct callers must set `O_NONBLOCK` first (#997).
Expand Down
3 changes: 2 additions & 1 deletion include/elio/runtime/autoscaler.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,8 @@ class autoscaler_impl {
// Handle each trigger
if constexpr (sizeof...(Triggers) > 0) {
// Custom triggers provided
((handle_trigger<Triggers>(pending, num_workers, now, last_idle_time, cfg)), ...);
(handle_trigger(Triggers{}, pending, num_workers, now,
last_idle_time, cfg), ...);
} else {
// Default behavior: scale up on overload, scale down on idle
handle_default_overload(pending, num_workers, cfg);
Expand Down
120 changes: 120 additions & 0 deletions tests/unit/test_autoscaler.cpp
Original file line number Diff line number Diff line change
@@ -1,11 +1,84 @@
#include <catch2/catch_test_macros.hpp>
#include <elio/elio.hpp>
#include <atomic>
#include <thread>
#include <chrono>

using namespace elio;
using namespace elio::runtime;

namespace {

template<typename Predicate>
bool wait_for_condition(Predicate&& predicate,
std::chrono::milliseconds timeout) {
const auto deadline = std::chrono::steady_clock::now() + timeout;
while (std::chrono::steady_clock::now() < deadline) {
if (predicate()) return true;
std::this_thread::sleep_for(std::chrono::milliseconds(1));
}
return predicate();
}

struct trigger_test_worker {
[[nodiscard]] bool is_idle() const noexcept { return false; }

[[nodiscard]] std::chrono::steady_clock::time_point
last_task_time() const noexcept {
return {};
}
};

struct trigger_test_scheduler {
[[nodiscard]] size_t pending_tasks() const noexcept {
return pending.load(std::memory_order_relaxed);
}

[[nodiscard]] size_t num_threads() const noexcept {
return workers.load(std::memory_order_relaxed);
}

void set_thread_count(size_t count) noexcept {
workers.store(count, std::memory_order_relaxed);
}

[[nodiscard]] trigger_test_worker* get_worker(size_t) noexcept {
return &worker;
}

std::atomic<size_t> pending{2};
std::atomic<size_t> workers{2};
trigger_test_worker worker;
};

struct overload_capture_action {
inline static std::atomic<size_t> calls{0};

void operator()(trigger_test_scheduler*, size_t) const noexcept {
calls.fetch_add(1, std::memory_order_release);
}
};

struct idle_capture_action {
inline static std::atomic<size_t> calls{0};

void operator()(trigger_test_scheduler*, size_t,
std::chrono::seconds) const noexcept {
calls.fetch_add(1, std::memory_order_release);
}
};

struct block_capture_action {
inline static std::atomic<size_t> calls{0};

void operator()(trigger_test_scheduler*, size_t,
std::chrono::milliseconds) const noexcept {
calls.fetch_add(1, std::memory_order_release);
}
};

} // namespace

TEST_CASE("autoscaler config defaults", "[autoscaler]") {
autoscaler_config config;

Expand All @@ -30,6 +103,53 @@ TEST_CASE("autoscaler start/stop", "[autoscaler]") {
sched.shutdown();
}

TEST_CASE("autoscaler custom trigger pack", "[autoscaler]") {
overload_capture_action::calls.store(0, std::memory_order_relaxed);
idle_capture_action::calls.store(0, std::memory_order_relaxed);
block_capture_action::calls.store(0, std::memory_order_relaxed);

autoscaler_config config;
config.tick_interval = std::chrono::milliseconds(1);
config.overload_threshold = 1;
config.idle_threshold = 1;
config.idle_delay = std::chrono::seconds(0);
config.min_workers = 1;
config.max_workers = 3;
config.block_threshold = std::chrono::milliseconds(0);

trigger_test_scheduler sched;
autoscaler<trigger_test_scheduler,
on_overload<overload_capture_action>,
on_idle<idle_capture_action>,
on_block<block_capture_action>> scaler{config};
scaler.start(&sched);

const bool overload_and_block_observed = wait_for_condition(
[] {
return overload_capture_action::calls.load(
std::memory_order_acquire) > 0 &&
block_capture_action::calls.load(
std::memory_order_acquire) > 0;
},
std::chrono::seconds(2));

sched.pending.store(0, std::memory_order_relaxed);
const bool idle_observed = wait_for_condition(
[] {
return idle_capture_action::calls.load(
std::memory_order_acquire) > 0;
},
std::chrono::seconds(2));

scaler.stop();

CHECK(overload_and_block_observed);
CHECK(idle_observed);
CHECK(overload_capture_action::calls.load(std::memory_order_relaxed) > 0);
CHECK(idle_capture_action::calls.load(std::memory_order_relaxed) > 0);
CHECK(block_capture_action::calls.load(std::memory_order_relaxed) > 0);
}

TEST_CASE("worker_thread last_task_time", "[autoscaler]") {
scheduler sched{1};
sched.start();
Expand Down