A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
loop.h
Go to the documentation of this file.
1/*
2 * Copyright 2026 The A11 Authors
3 *
4 * Licensed under the Apache License, Version 2.0 (the "License");
5 * you may not use this file except in compliance with the License.
6 * You may obtain a copy of the License at
7 *
8 * http://www.apache.org/licenses/LICENSE-2.0
9 *
10 * Unless required by applicable law or agreed to in writing, software
11 * distributed under the License is distributed on an "AS IS" BASIS,
12 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13 * See the License for the specific language governing permissions and
14 * limitations under the License.
15 */
16
30#ifndef A11_UV_LOOP_H_
31#define A11_UV_LOOP_H_
32
33#include <atomic>
34#include <csignal>
35#include <cstdint>
36#include <cstdio>
37#include <cstdlib>
38#include <deque>
39#include <functional>
40#include <memory>
41#include <optional>
42#include <string>
43#include <string_view>
44#include <thread>
45#include <utility>
46#include <uvw.hpp>
47#include <vector>
48
49#include <absl/base/no_destructor.h>
50#include <absl/base/thread_annotations.h>
51#include <absl/container/flat_hash_map.h>
52#include <absl/container/flat_hash_set.h>
53#include <absl/log/log.h>
54#include <absl/status/status.h>
55#include <absl/status/status_macros.h>
56#include <absl/status/statusor.h>
57#include <absl/strings/numbers.h>
58#include <absl/strings/str_cat.h>
59#include <absl/time/clock.h>
60#include <absl/time/time.h>
61
63#include "thread/boost_primitives.h"
64
65namespace a11::uv {
66
68inline absl::Status UvError(int code, std::string_view operation) {
69 return absl::UnavailableError(
70 absl::StrCat(operation, " failed: ", uv_strerror(code)));
71}
72
73// One libuv loop is shared by native HTTP clients and servers. All uvw and
74// protocol mutation is serialized on this thread; fiber callbacks communicate
75// through A11 Futures and never block the loop.
77 public:
78 static UvExecutor& Instance() {
79 // The process-wide I/O scheduler lives until process exit.
80 static absl::NoDestructor<UvExecutor> executor;
81 // Wait outside the static initialization guard because the fiber-aware wait
82 // may yield and allow another fiber to enter Instance().
83 executor->EnsureStarted();
84 return *executor;
85 }
86
96 absl::Status Post(std::function<void()> work,
97 const void* order_key = nullptr) {
98 if (!work) {
99 return absl::InvalidArgumentError("uv work must be callable");
100 }
101 {
102 thread::MutexLock lock(&mu_);
103 if (!running_) {
104 return absl::FailedPreconditionError("The A11 libuv loop is stopped");
105 }
106 work_.push_back(Item{.key = order_key, .work = std::move(work)});
107 }
108 const int result = async_->send();
109 if (result != 0) {
110 return UvError(result, "uv_async_send");
111 }
112 return absl::OkStatus();
113 }
114
115 [[nodiscard]] std::shared_ptr<uvw::loop> loop() const { return loop_; }
116
117 [[nodiscard]] bool IsLoopThread() const {
118 thread::MutexLock lock(&mu_);
119 return loop_thread_id_.has_value() &&
120 *loop_thread_id_ == std::this_thread::get_id();
121 }
122
123 private:
124 friend class absl::NoDestructor<UvExecutor>;
125
139 static void IgnoreSigPipeIfDefaulted() {
140#if !defined(_WIN32)
141 // Unqualified calls on purpose: `struct sigaction` and the function of the
142 // same name both live in the global namespace, and `::sigaction(...)`
143 // parses as a functional cast to the struct.
144 struct sigaction current{};
145 if (sigaction(SIGPIPE, nullptr, &current) != 0) {
146 return;
147 }
148 if (current.sa_handler != SIG_DFL) {
149 return;
150 }
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";
157 }
158#endif
159 }
160
161 UvExecutor() {
162 {
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);
170 }
171 async_->on<uvw::async_event>(
172 [this](const uvw::async_event&, uvw::async_handle&) { Drain(); });
173 thread_ = std::thread([this]() {
174 {
175 thread::MutexLock lock(&mu_);
176 loop_thread_id_ = std::this_thread::get_id();
177 cv_.SignalAll();
178 }
179 loop_->run();
180 thread::MutexLock lock(&mu_);
181 running_ = false;
182 });
183 }
184 }
185
187 void EnsureStarted() {
188 thread::MutexLock lock(&mu_);
189 while (!loop_thread_id_.has_value()) {
190 cv_.Wait(&mu_);
191 }
192 }
193
200 void Drain() {
201 std::deque<Item> batch;
202 {
203 thread::MutexLock lock(&mu_);
204 batch.swap(work_);
205 }
206 if (DrainStatsEnabled()) {
207 RecordDrainBatch(batch);
208 // Time each item separately when drain statistics are enabled.
209 for (Item& item : batch) {
210 const absl::Time started = absl::Now();
211 item.work();
212 RecordItemDuration(absl::ToInt64Nanoseconds(absl::Now() - started));
213 }
214 return;
215 }
216 if (batch.size() <= 1 || !FairDraining()) {
217 for (Item& item : batch) {
218 item.work();
219 }
220 return;
221 }
222
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);
227 if (fresh) {
228 keys.push_back(item.key);
229 }
230 lane->second.push_back(std::move(item.work));
231 }
232 if (keys.size() == 1) {
233 for (std::function<void()>& work : lanes.begin()->second) {
234 work();
235 }
236 return;
237 }
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);
242 if (lane.empty()) {
243 continue;
244 }
245 std::function<void()> work = std::move(lane.front());
246 lane.pop_front();
247 --remaining;
248 work();
249 }
250 }
251 }
252
253 struct Item {
255 const void* key = nullptr;
256 std::function<void()> work;
257 };
258
260 static bool DrainStatsEnabled() {
261 static const bool on = [] {
262 const char* setting = std::getenv("A11_UV_DRAIN_STATS");
263 int value = 0;
264 return setting != nullptr && absl::SimpleAtoi(setting, &value) &&
265 value != 0;
266 }();
267 return on;
268 }
269
271 static void RecordItemDuration(std::int64_t nanos) {
272 struct Buckets {
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};
279 };
280
281 static absl::NoDestructor<Buckets> buckets;
282 static const bool registered = [] {
283 std::atexit([] {
284 std::fprintf(
285 stderr,
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);
294 });
295 return true;
296 }();
297 (void)registered;
298 // Plain literals, no digit separators: clang-format has been seen to read
299 // 1'000'000 as character literals when reflowing this region.
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);
311 } else {
312 buckets->over_1ms.fetch_add(1, std::memory_order_relaxed);
313 }
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)) {}
317 }
318
319 static void RecordDrainBatch(const std::deque<Item>& batch) {
320 struct Stats {
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};
326 };
327
328 static absl::NoDestructor<Stats> stats;
329 static const bool registered = [] {
330 std::atexit([] {
331 const std::uint64_t drains = stats->drains.load();
332 const double per = drains == 0 ? 1.0 : static_cast<double>(drains);
333 std::fprintf(
334 stderr,
335 "uv drain: %llu drains, %llu items (%.2f/drain), "
336 "%llu with >1 item (%.1f%%), %llu with >1 key (%.1f%%), "
337 "largest %llu\n",
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()),
342 drains == 0
343 ? 0.0
344 : 100.0 * static_cast<double>(stats->multi_item.load()) / per,
345 static_cast<unsigned long long>(stats->multi_key.load()),
346 drains == 0
347 ? 0.0
348 : 100.0 * static_cast<double>(stats->multi_key.load()) / per,
349 static_cast<unsigned long long>(stats->largest.load()));
350 });
351 return true;
352 }();
353 (void)registered;
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);
361 }
362 if (keys.size() > 1) {
363 stats->multi_key.fetch_add(1, std::memory_order_relaxed);
364 }
365 }
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)) {}
370 }
371
373 static bool FairDraining() {
374 static const bool fair = [] {
375 const char* setting = std::getenv("A11_UV_FAIR");
376 // A value that is not a number leaves fair draining on: only an explicit
377 // 0 turns it off, so a typo cannot quietly change how the loop drains.
378 int value = 0;
379 return setting == nullptr || !absl::SimpleAtoi(setting, &value) ||
380 value != 0;
381 }();
382 return fair;
383 }
384
385 mutable thread::Mutex mu_;
386 thread::CondVar cv_;
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_;
392 std::thread thread_;
393};
394
403template <typename T>
404absl::StatusOr<T> RunOnUv(std::function<absl::StatusOr<T>()> operation,
405 const void* order_key = nullptr) {
406 if (UvExecutor::Instance().IsLoopThread()) {
407 return operation();
408 }
409 auto promise = std::make_shared<a11::Promise<T>>();
410 a11::Future<T> future = promise->future();
411 ABSL_RETURN_IF_ERROR(UvExecutor::Instance().Post(
412 [promise, operation = std::move(operation)]() mutable {
413 (void)promise->SetResult(operation());
414 },
415 order_key));
416 return future.Await();
417}
418
419inline absl::Status RunStatusOnUv(std::function<absl::Status()> operation,
420 const void* order_key = nullptr) {
421 absl::StatusOr<a11::Unit> result = RunOnUv<a11::Unit>(
422 [operation =
423 std::move(operation)]() mutable -> absl::StatusOr<a11::Unit> {
424 ABSL_RETURN_IF_ERROR(operation());
425 return a11::Unit{};
426 },
427 order_key);
428 return result.status();
429}
430
431} // namespace a11::uv
432
433#endif // A11_UV_LOOP_H_
Shared handle to one asynchronous result.
Definition future.h:126
Definition loop.h:76
std::shared_ptr< uvw::loop > loop() const
Definition loop.h:115
absl::Status Post(std::function< void()> work, const void *order_key=nullptr)
Queue work for the loop thread.
Definition loop.h:96
static UvExecutor & Instance()
Definition loop.h:78
bool IsLoopThread() const
Definition loop.h:117
std::string value
Definition discover.cc:114
std::string key
Definition discover.cc:782
Completion values used by every asynchronous A11 operation.
Definition loop.h:65
absl::Status UvError(int code, std::string_view operation)
The status a failed libuv call deserves.
Definition loop.h:68
absl::StatusOr< T > RunOnUv(std::function< absl::StatusOr< T >()> operation, const void *order_key=nullptr)
Run operation on the loop thread and wait for its result.
Definition loop.h:404
absl::Status RunStatusOnUv(std::function< absl::Status()> operation, const void *order_key=nullptr)
Definition loop.h:419
std::deque< ItemPtr > items
Definition runtime.cc:641
Empty success value used by Future<Unit> operations that return no data.
Definition future.h:55