80 static absl::NoDestructor<UvExecutor> executor;
83 executor->EnsureStarted();
96 absl::Status
Post(std::function<
void()> work,
97 const void* order_key =
nullptr) {
99 return absl::InvalidArgumentError(
"uv work must be callable");
102 thread::MutexLock lock(&mu_);
104 return absl::FailedPreconditionError(
"The A11 libuv loop is stopped");
106 work_.push_back(Item{.key = order_key, .work = std::move(work)});
108 const int result = async_->send();
110 return UvError(result,
"uv_async_send");
112 return absl::OkStatus();
115 [[nodiscard]] std::shared_ptr<uvw::loop>
loop()
const {
return loop_; }
118 thread::MutexLock lock(&mu_);
119 return loop_thread_id_.has_value() &&
120 *loop_thread_id_ == std::this_thread::get_id();
139 static void IgnoreSigPipeIfDefaulted() {
144 struct sigaction current{};
145 if (sigaction(SIGPIPE,
nullptr, ¤t) != 0) {
148 if (current.sa_handler != SIG_DFL) {
151 struct sigaction ignore{};
152 ignore.sa_handler = SIG_IGN;
153 sigemptyset(&ignore.sa_mask);
154 if (sigaction(SIGPIPE, &ignore,
nullptr) != 0) {
155 LOG(WARNING) <<
"Could not ignore SIGPIPE; a peer that closes mid-write "
156 "will terminate this process";
163 IgnoreSigPipeIfDefaulted();
164 loop_ = uvw::loop::create();
165 async_ = loop_->resource<uvw::async_handle>();
166 const int initialized = async_->init();
167 if (initialized != 0) {
168 LOG(FATAL) <<
"Could not initialize the A11 libuv executor: "
169 << uv_strerror(initialized);
171 async_->on<uvw::async_event>(
172 [
this](
const uvw::async_event&, uvw::async_handle&) { Drain(); });
173 thread_ = std::thread([
this]() {
175 thread::MutexLock lock(&mu_);
176 loop_thread_id_ = std::this_thread::get_id();
180 thread::MutexLock lock(&mu_);
187 void EnsureStarted() {
188 thread::MutexLock lock(&mu_);
189 while (!loop_thread_id_.has_value()) {
201 std::deque<Item> batch;
203 thread::MutexLock lock(&mu_);
206 if (DrainStatsEnabled()) {
207 RecordDrainBatch(batch);
209 for (Item& item : batch) {
210 const absl::Time started = absl::Now();
212 RecordItemDuration(absl::ToInt64Nanoseconds(absl::Now() - started));
216 if (batch.size() <= 1 || !FairDraining()) {
217 for (Item& item : batch) {
223 std::vector<const void*> keys;
224 absl::flat_hash_map<
const void*, std::deque<std::function<void()>>> lanes;
225 for (Item& item : batch) {
226 auto [lane, fresh] = lanes.try_emplace(item.key);
228 keys.push_back(item.key);
230 lane->second.push_back(std::move(item.work));
232 if (keys.size() == 1) {
233 for (std::function<
void()>& work : lanes.begin()->second) {
238 size_t remaining = batch.size();
239 while (remaining > 0) {
240 for (
const void*
key : keys) {
241 std::deque<std::function<void()>>& lane = lanes.at(
key);
245 std::function<void()> work = std::move(lane.front());
255 const void*
key =
nullptr;
256 std::function<void()> work;
260 static bool DrainStatsEnabled() {
261 static const bool on = [] {
262 const char* setting = std::getenv(
"A11_UV_DRAIN_STATS");
264 return setting !=
nullptr && absl::SimpleAtoi(setting, &
value) &&
271 static void RecordItemDuration(std::int64_t nanos) {
273 std::atomic<std::uint64_t> under_10us{0};
274 std::atomic<std::uint64_t> under_100us{0};
275 std::atomic<std::uint64_t> under_1ms{0};
276 std::atomic<std::uint64_t> over_1ms{0};
277 std::atomic<std::uint64_t> total_nanos{0};
278 std::atomic<std::int64_t> worst_nanos{0};
281 static absl::NoDestructor<Buckets> buckets;
282 static const bool registered = [] {
286 "uv item: <10us %llu, <100us %llu, <1ms %llu, >=1ms %llu, "
287 "busy %.1f ms total, worst %.0f us\n",
288 static_cast<unsigned long long>(buckets->under_10us.load()),
289 static_cast<unsigned long long>(buckets->under_100us.load()),
290 static_cast<unsigned long long>(buckets->under_1ms.load()),
291 static_cast<unsigned long long>(buckets->over_1ms.load()),
292 static_cast<double>(buckets->total_nanos.load()) / 1e6,
293 static_cast<double>(buckets->worst_nanos.load()) / 1e3);
300 constexpr std::int64_t kTenMicros = 10000;
301 constexpr std::int64_t kHundredMicros = 100000;
302 constexpr std::int64_t kMilli = 1000000;
303 buckets->total_nanos.fetch_add(
static_cast<std::uint64_t
>(nanos),
304 std::memory_order_relaxed);
305 if (nanos < kTenMicros) {
306 buckets->under_10us.fetch_add(1, std::memory_order_relaxed);
307 }
else if (nanos < kHundredMicros) {
308 buckets->under_100us.fetch_add(1, std::memory_order_relaxed);
309 }
else if (nanos < kMilli) {
310 buckets->under_1ms.fetch_add(1, std::memory_order_relaxed);
312 buckets->over_1ms.fetch_add(1, std::memory_order_relaxed);
314 std::int64_t seen = buckets->worst_nanos.load(std::memory_order_relaxed);
315 while (nanos > seen && !buckets->worst_nanos.compare_exchange_weak(
316 seen, nanos, std::memory_order_relaxed)) {}
319 static void RecordDrainBatch(
const std::deque<Item>& batch) {
321 std::atomic<std::uint64_t> drains{0};
322 std::atomic<std::uint64_t>
items{0};
323 std::atomic<std::uint64_t> multi_item{0};
324 std::atomic<std::uint64_t> multi_key{0};
325 std::atomic<std::uint64_t> largest{0};
328 static absl::NoDestructor<Stats> stats;
329 static const bool registered = [] {
331 const std::uint64_t drains = stats->drains.load();
332 const double per = drains == 0 ? 1.0 :
static_cast<double>(drains);
335 "uv drain: %llu drains, %llu items (%.2f/drain), "
336 "%llu with >1 item (%.1f%%), %llu with >1 key (%.1f%%), "
338 static_cast<unsigned long long>(drains),
339 static_cast<unsigned long long>(stats->items.load()),
340 drains == 0 ? 0.0 : static_cast<double>(stats->
items.load()) / per,
341 static_cast<unsigned long long>(stats->multi_item.load()),
344 : 100.0 * static_cast<double>(stats->multi_item.load()) / per,
345 static_cast<unsigned long long>(stats->multi_key.load()),
348 : 100.0 * static_cast<double>(stats->multi_key.load()) / per,
349 static_cast<unsigned long long>(stats->largest.load()));
354 stats->drains.fetch_add(1, std::memory_order_relaxed);
355 stats->items.fetch_add(batch.size(), std::memory_order_relaxed);
356 if (batch.size() > 1) {
357 stats->multi_item.fetch_add(1, std::memory_order_relaxed);
358 absl::flat_hash_set<const void*> keys;
359 for (
const Item& item : batch) {
360 keys.insert(item.key);
362 if (keys.size() > 1) {
363 stats->multi_key.fetch_add(1, std::memory_order_relaxed);
366 std::uint64_t seen = stats->largest.load(std::memory_order_relaxed);
367 while (batch.size() > seen &&
368 !stats->largest.compare_exchange_weak(seen, batch.size(),
369 std::memory_order_relaxed)) {}
373 static bool FairDraining() {
374 static const bool fair = [] {
375 const char* setting = std::getenv(
"A11_UV_FAIR");
379 return setting ==
nullptr || !absl::SimpleAtoi(setting, &
value) ||
385 mutable thread::Mutex mu_;
387 bool running_ ABSL_GUARDED_BY(mu_) =
true;
388 std::optional<std::thread::id> loop_thread_id_ ABSL_GUARDED_BY(mu_);
389 std::deque<Item> work_ ABSL_GUARDED_BY(mu_);
390 std::shared_ptr<uvw::loop> loop_;
391 std::shared_ptr<uvw::async_handle> async_;