TLA Line data Source code
1 : //
2 : // Copyright (c) 2026 Steve Gerbino
3 : // Copyright (c) 2026 Michael Vandeberg
4 : //
5 : // Distributed under the Boost Software License, Version 1.0. (See accompanying
6 : // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7 : //
8 : // Official repository: https://github.com/cppalliance/corosio
9 : //
10 :
11 : #ifndef BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_SIGNAL_SERVICE_HPP
12 : #define BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_SIGNAL_SERVICE_HPP
13 :
14 : #include <boost/corosio/detail/platform.hpp>
15 :
16 : #if BOOST_COROSIO_POSIX
17 :
18 : #include <boost/corosio/native/detail/posix/posix_signal.hpp>
19 :
20 : #include <boost/corosio/detail/config.hpp>
21 : #include <boost/corosio/detail/object_pool.hpp>
22 : #include <boost/corosio/detail/object_ref.hpp>
23 : #include <boost/capy/ex/execution_context.hpp>
24 : #include <boost/corosio/detail/scheduler.hpp>
25 : #include <boost/corosio/native/detail/make_err.hpp>
26 : #include <boost/capy/error.hpp>
27 :
28 : #include <mutex>
29 : #include <tuple>
30 : #include <vector>
31 :
32 : #include <errno.h>
33 : #include <fcntl.h>
34 : #include <signal.h>
35 : #include <unistd.h>
36 :
37 : /*
38 : POSIX Signal Service
39 : ====================
40 :
41 : Concrete signal service implementation for POSIX backends. Manages signal
42 : registrations via sigaction() and dispatches completions through the
43 : scheduler. One instance per execution_context, created on first use
44 : by the public signal_set.
45 :
46 : See the block comment further down for the full architecture overview.
47 : */
48 :
49 : /*
50 : POSIX Signal Implementation
51 : ===========================
52 :
53 : This file implements signal handling for POSIX systems using sigaction().
54 : The implementation supports signal flags (SA_RESTART, etc.) and integrates
55 : with any POSIX-compatible scheduler via the abstract scheduler interface.
56 :
57 : Architecture Overview
58 : ---------------------
59 :
60 : Three layers manage signal registrations:
61 :
62 : 1. signal_state (global singleton)
63 : - Tracks the global service list and per-signal registration counts
64 : - Stores the flags used for first registration of each signal (for
65 : conflict detection when multiple signal_sets register same signal)
66 : - Owns the mutex that protects signal handler installation/removal
67 :
68 : 2. posix_signal_service (one per execution_context)
69 : - Maintains registrations_[] table indexed by signal number
70 : - Each slot is a doubly-linked list of signal_registrations for that signal
71 : - Also owns a recycling pool of all posix_signal objects it owns
72 :
73 : 3. posix_signal (one per signal_set)
74 : - Owns a singly-linked list (sorted by signal number) of signal_registrations
75 : - Contains the pending_op_ used for wait operations
76 :
77 : Signal Delivery Flow
78 : --------------------
79 :
80 : Delivery uses the self-pipe trick so the signal handler itself performs
81 : only async-signal-safe work (mirrors Boost.Asio):
82 :
83 : 1. Signal arrives -> corosio_posix_signal_handler(). The handler only
84 : write()s the signal number to the global self-pipe (write_fd) and
85 : restores errno. No locks, no allocation, no scheduler dispatch.
86 :
87 : 2. The read end of the pipe is watched by one backend's event loop
88 : (registered via scheduler::register_signal_reader on the first
89 : registration). When it becomes readable the backend drains it
90 : (drain_signal_pipe) and calls deliver_signal() in normal context.
91 :
92 : 3. deliver_signal() iterates all posix_signal_service services:
93 : - If a signal_set is waiting (impl->waiting_ == true), post the signal_op
94 : to the scheduler for immediate completion
95 : - Otherwise, increment reg->undelivered to queue the signal
96 :
97 : 4. When wait() is called via start_wait():
98 : - First check for queued signals (undelivered > 0); if found, post
99 : immediate completion without blocking
100 : - Otherwise, set waiting_ = true and call work_started() to keep
101 : the io_context alive
102 :
103 : Locking Protocol
104 : ----------------
105 :
106 : Two mutex levels exist (MUST acquire in this order to avoid deadlock):
107 : 1. signal_state::mutex - protects handler registration and service list
108 : 2. posix_signal_service::mutex_ - protects per-service registration tables
109 :
110 : Async-Signal-Safety
111 : -------------------
112 :
113 : The C signal handler (corosio_posix_signal_handler) performs only
114 : async-signal-safe operations: it reads the single global write_fd and
115 : calls write(), saving/restoring errno. It never locks a mutex, allocates
116 : memory, or dispatches through the scheduler. All of that happens in
117 : deliver_signal(), which runs in normal thread context from the backend
118 : event loop after draining the self-pipe. There is therefore no
119 : self-deadlock risk if a signal arrives while a thread holds state->mutex
120 : or service->mutex_.
121 :
122 : Flag Handling
123 : -------------
124 :
125 : - Flags are abstract values in the public API (signal_set::flags_t)
126 : - flags_supported() validates that requested flags are available on
127 : this platform; returns false if SA_NOCLDWAIT is unavailable and
128 : no_child_wait is requested
129 : - to_sigaction_flags() maps validated flags to actual SA_* constants
130 : - First registration of a signal establishes the flags; subsequent
131 : registrations must be compatible (same flags or dont_care)
132 : - Requesting unavailable flags returns operation_not_supported
133 :
134 : Work Tracking
135 : -------------
136 :
137 : When waiting for a signal:
138 : - start_wait() calls sched_->work_started() to prevent io_context::run()
139 : from returning while we wait
140 : - signal_op::svc is set to point to the service
141 : - signal_op::operator()() calls work_finished() after resuming the coroutine
142 :
143 : If a signal was already queued (undelivered > 0), no work tracking is needed
144 : because completion is posted immediately.
145 : */
146 :
147 : namespace boost::corosio {
148 :
149 : namespace detail {
150 :
151 : /** Signal service for POSIX backends.
152 :
153 : Manages signal registrations via sigaction() and dispatches signal
154 : completions through the scheduler. One instance per execution_context.
155 : */
156 : class BOOST_COROSIO_DECL posix_signal_service final
157 : : public capy::execution_context::service
158 : , public io_object::io_service
159 : {
160 : friend class posix_signal;
161 :
162 : public:
163 : using key_type = posix_signal_service;
164 :
165 : explicit posix_signal_service(capy::execution_context& ctx);
166 : ~posix_signal_service() override;
167 :
168 : posix_signal_service(posix_signal_service const&) = delete;
169 : posix_signal_service& operator=(posix_signal_service const&) = delete;
170 :
171 : io_object::implementation* construct() override;
172 :
173 HIT 252 : void destroy(io_object::implementation* p) override
174 : {
175 252 : auto& impl = static_cast<posix_signal&>(*p);
176 252 : [[maybe_unused]] auto n = impl.clear();
177 252 : impl.disarm_stop();
178 252 : impl.cancel();
179 252 : release(&impl);
180 252 : }
181 :
182 : /** Shut down the service.
183 :
184 : Destroys every implementation the service still owns and gives
185 : each of their registrations back to the process-global table.
186 : */
187 : void shutdown() override;
188 :
189 : std::error_code add_signal(
190 : posix_signal& impl, int signal_number, signal_set::flags_t flags);
191 :
192 : std::error_code remove_signal(posix_signal& impl, int signal_number);
193 :
194 : std::error_code clear_signals(posix_signal& impl);
195 :
196 : void cancel_wait(posix_signal& impl);
197 : void start_wait(posix_signal& impl, signal_op* op);
198 :
199 : /** Cancel an in-flight wait on behalf of a stop token.
200 :
201 : Identical to @ref cancel_wait except that it does not set the
202 : sticky `cancelled_` latch: a stop token scopes to one operation,
203 : so a request arriving after the wait completed must do nothing.
204 : */
205 : void cancel_wait_token(posix_signal& impl) noexcept;
206 :
207 : /** Clear the per-operation stop flag before a new wait arms.
208 :
209 : Lives here rather than on the implementation because `mutex_` is
210 : the service's; the service is a friend of `posix_signal`, not the
211 : reverse.
212 : */
213 1104 : void reset_token_cancel(posix_signal& impl) noexcept
214 : {
215 1104 : std::lock_guard lock(mutex_);
216 1104 : impl.token_cancelled_ = false;
217 1104 : }
218 :
219 : static void deliver_signal(int signal_number);
220 :
221 : void work_started() noexcept;
222 : void work_finished() noexcept;
223 : void post(signal_op* op);
224 :
225 : private:
226 : static void add_service(posix_signal_service* service);
227 : static void remove_service(posix_signal_service* service);
228 :
229 : scheduler* sched_;
230 : std::mutex mutex_;
231 :
232 : // Registers the signal self-pipe's read end with sched_ exactly once per
233 : // service, so every io_context that waits on a signal can drain the pipe.
234 : // A once_flag (not a bool under mutex_) because registration must run
235 : // without holding mutex_ or the signal-state mutex — see add_signal.
236 : std::mutex reader_mutex_;
237 : bool reader_registered_ = false;
238 :
239 : object_pool<posix_signal> pool_;
240 :
241 : // Per-signal registration table
242 : signal_registration* registrations_[max_signal_number];
243 :
244 : // Registration counts for each signal
245 : std::size_t registration_count_[max_signal_number];
246 :
247 : // Linked list of all posix_signal_service services for signal delivery
248 : posix_signal_service* next_ = nullptr;
249 : posix_signal_service* prev_ = nullptr;
250 : };
251 :
252 : } // namespace detail
253 :
254 : } // namespace boost::corosio
255 :
256 : // ---------------------------------------------------------------------------
257 : // Inline implementation
258 : // ---------------------------------------------------------------------------
259 :
260 : namespace boost::corosio {
261 :
262 : namespace detail {
263 :
264 : namespace posix_signal_detail {
265 :
266 : struct signal_state
267 : {
268 : std::mutex mutex;
269 : posix_signal_service* service_list = nullptr;
270 : std::size_t registration_count[max_signal_number] = {};
271 : signal_set::flags_t registered_flags[max_signal_number] = {};
272 :
273 : // Self-pipe used to defer signal delivery out of handler context.
274 : // The C handler writes the signal number to write_fd (async-signal-
275 : // safe); a backend event loop drains read_fd and calls deliver_signal()
276 : // in normal context. Created once (on the first signal registration) and
277 : // kept for the process lifetime. Each posix_signal_service registers the
278 : // read end with its own scheduler (see reader_once_) so every running
279 : // io_context can drain it; multiple readers on one pipe are safe because
280 : // each signal is a fixed sizeof(int) record read atomically.
281 : int read_fd = -1;
282 : int write_fd = -1;
283 : };
284 :
285 : BOOST_COROSIO_DECL signal_state* get_signal_state();
286 :
287 : // Check if requested flags are supported on this platform.
288 : // Returns true if all flags are supported, false otherwise.
289 : inline bool
290 273 : flags_supported([[maybe_unused]] signal_set::flags_t flags)
291 : {
292 : #ifndef SA_NOCLDWAIT
293 : if (flags & signal_set::no_child_wait)
294 : return false;
295 : #endif
296 273 : return true;
297 : }
298 :
299 : // Map abstract flags to sigaction() flags.
300 : // Caller must ensure flags_supported() returns true first.
301 : inline int
302 225 : to_sigaction_flags(signal_set::flags_t flags)
303 : {
304 225 : int sa_flags = 0;
305 225 : if (flags & signal_set::restart)
306 23 : sa_flags |= SA_RESTART;
307 225 : if (flags & signal_set::no_child_stop)
308 3 : sa_flags |= SA_NOCLDSTOP;
309 : #ifdef SA_NOCLDWAIT
310 225 : if (flags & signal_set::no_child_wait)
311 2 : sa_flags |= SA_NOCLDWAIT;
312 : #endif
313 225 : if (flags & signal_set::no_defer)
314 4 : sa_flags |= SA_NODEFER;
315 225 : if (flags & signal_set::reset_handler)
316 2 : sa_flags |= SA_RESETHAND;
317 225 : return sa_flags;
318 : }
319 :
320 : // Check if two flag values are compatible
321 : inline bool
322 39 : flags_compatible(signal_set::flags_t existing, signal_set::flags_t requested)
323 : {
324 : // dont_care is always compatible
325 76 : if ((existing & signal_set::dont_care) ||
326 37 : (requested & signal_set::dont_care))
327 7 : return true;
328 :
329 : // Mask out dont_care bit for comparison
330 32 : constexpr auto mask = ~signal_set::dont_care;
331 32 : return (existing & mask) == (requested & mask);
332 : }
333 :
334 : // Lazily create the global signal self-pipe. Idempotent; call under
335 : // state->mutex before installing the first signal handler so write_fd is
336 : // valid by the time the handler can fire. Both ends are non-blocking and
337 : // close-on-exec (mirrors the reactor self-pipe setup in select_scheduler).
338 : // Returns the failing call's errno and leaves the fds at -1 if creation
339 : // fails: an exhausted descriptor table and a rejected fcntl are different
340 : // problems to the caller of add().
341 : [[nodiscard]] inline std::error_code
342 273 : open_signal_pipe(signal_state* state)
343 : {
344 273 : if (state->read_fd >= 0)
345 258 : return {};
346 :
347 : int fds[2];
348 15 : if (::pipe(fds) < 0)
349 1 : return make_err(errno);
350 :
351 33 : for (int i = 0; i < 2; ++i)
352 : {
353 25 : int fl = ::fcntl(fds[i], F_GETFL, 0);
354 46 : if (fl == -1 || ::fcntl(fds[i], F_SETFL, fl | O_NONBLOCK) == -1 ||
355 21 : ::fcntl(fds[i], F_SETFD, FD_CLOEXEC) == -1)
356 : {
357 6 : auto ec = make_err(errno);
358 6 : ::close(fds[0]);
359 6 : ::close(fds[1]);
360 6 : return ec;
361 : }
362 : }
363 :
364 8 : state->read_fd = fds[0];
365 8 : state->write_fd = fds[1];
366 8 : return {};
367 : }
368 :
369 : // C signal handler. Async-signal-safe: it touches only the single global
370 : // write_fd (an int set before any handler is installed) and calls write(),
371 : // which POSIX lists as async-signal-safe. errno is saved and restored so an
372 : // interrupted foreground syscall is unaffected. A full pipe (write returns
373 : // EAGAIN) or a short write is intentionally dropped — the reactor still
374 : // coalesces because deliver_signal reports the signal to every waiting set.
375 : inline void
376 319 : corosio_posix_signal_handler(int signal_number)
377 : {
378 319 : int saved_errno = errno;
379 319 : signal_state* state = get_signal_state();
380 : [[maybe_unused]] ssize_t r =
381 319 : ::write(state->write_fd, &signal_number, sizeof(int));
382 319 : errno = saved_errno;
383 : // With sigaction(), the handler persists automatically (unlike some
384 : // signal() implementations that reset to SIG_DFL).
385 319 : }
386 :
387 : // Drain the signal self-pipe and deliver each pending signal. Runs in normal
388 : // thread context from the backend event loop, so deliver_signal()'s mutex
389 : // locking and scheduler post are safe here. Reads until EAGAIN (edge-
390 : // triggered backends require a full drain per readiness event).
391 : inline void
392 319 : drain_signal_pipe()
393 : {
394 319 : signal_state* state = get_signal_state();
395 : int signal_number;
396 638 : while (::read(state->read_fd, &signal_number, sizeof(int)) ==
397 : static_cast<ssize_t>(sizeof(int)))
398 : {
399 319 : posix_signal_service::deliver_signal(signal_number);
400 : }
401 319 : }
402 :
403 : } // namespace posix_signal_detail
404 :
405 : // signal_op implementation
406 :
407 : inline void
408 323 : signal_op::operator()()
409 : {
410 323 : if (ec_out)
411 323 : *ec_out = {};
412 323 : if (signal_out)
413 323 : *signal_out = signal_number;
414 :
415 : // Capture svc before resuming (coro may destroy us)
416 323 : auto* service = svc;
417 323 : svc = nullptr;
418 :
419 323 : cont.h = h;
420 323 : d.post(cont);
421 :
422 : // Balance the work_started() from start_wait
423 323 : if (service)
424 321 : service->work_finished();
425 323 : }
426 :
427 : inline void
428 MIS 0 : signal_op::destroy()
429 : {
430 : // No-op: signal_op is embedded in posix_signal
431 0 : }
432 :
433 : // posix_signal implementation
434 :
435 HIT 186 : inline posix_signal::posix_signal(posix_signal_service& svc) noexcept
436 186 : : svc_(svc)
437 : {
438 186 : }
439 :
440 : inline std::coroutine_handle<>
441 1138 : posix_signal::wait(
442 : std::coroutine_handle<> h,
443 : capy::executor_ref d,
444 : std::stop_token token,
445 : std::error_code* ec,
446 : int* signal_out)
447 : {
448 1138 : pending_op_.h = h;
449 1138 : pending_op_.d = d;
450 1138 : pending_op_.ec_out = ec;
451 1138 : pending_op_.signal_out = signal_out;
452 1138 : pending_op_.signal_number = 0;
453 :
454 : // Disarm any callback left over from a previous wait before doing
455 : // anything else, including the early return below: otherwise that
456 : // path leaves this object owning a callback it no longer uses.
457 : // Outside start_wait's lock on purpose: ~stop_callback blocks until a
458 : // concurrently running callback returns, and that callback takes
459 : // posix_signal_service::mutex_.
460 1138 : stop_cb_.reset();
461 :
462 1138 : if (token.stop_requested())
463 : {
464 34 : if (ec)
465 34 : *ec = make_error_code(capy::error::canceled);
466 34 : if (signal_out)
467 34 : *signal_out = 0;
468 34 : pending_op_.cont.h = h;
469 34 : d.post(pending_op_.cont);
470 : // completion is always posted to scheduler queue, never inline.
471 34 : return std::noop_coroutine();
472 : }
473 :
474 : // Clearing the flag before arming is load-bearing: reset_token_cancel
475 : // must run immediately before emplace, not before the early return
476 : // above.
477 1104 : svc_.reset_token_cancel(*this);
478 1104 : if (token.stop_possible())
479 772 : stop_cb_.emplace(token, token_canceller{this});
480 :
481 1104 : svc_.start_wait(*this, &pending_op_);
482 : // completion is always posted to scheduler queue, never inline.
483 1104 : return std::noop_coroutine();
484 : }
485 :
486 : inline std::error_code
487 277 : posix_signal::add(int signal_number, signal_set::flags_t flags)
488 : {
489 277 : return svc_.add_signal(*this, signal_number, flags);
490 : }
491 :
492 : inline std::error_code
493 26 : posix_signal::remove(int signal_number)
494 : {
495 26 : return svc_.remove_signal(*this, signal_number);
496 : }
497 :
498 : inline std::error_code
499 266 : posix_signal::clear()
500 : {
501 266 : return svc_.clear_signals(*this);
502 : }
503 :
504 : inline void
505 269 : posix_signal::cancel() noexcept
506 : {
507 269 : svc_.cancel_wait(*this);
508 269 : }
509 :
510 : // posix_signal_service implementation
511 :
512 167 : inline posix_signal_service::posix_signal_service(
513 167 : capy::execution_context& ctx)
514 167 : : sched_(&get_scheduler(ctx))
515 : {
516 10855 : for (int i = 0; i < max_signal_number; ++i)
517 : {
518 10688 : registrations_[i] = nullptr;
519 10688 : registration_count_[i] = 0;
520 : }
521 167 : add_service(this);
522 167 : }
523 :
524 334 : inline posix_signal_service::~posix_signal_service()
525 : {
526 167 : remove_service(this);
527 334 : }
528 :
529 : inline void
530 167 : posix_signal_service::shutdown()
531 : {
532 : // Collected under the locks below and released after they are
533 : // released: ~posix_signal destroys an armed stop_cb_, and
534 : // ~stop_callback blocks until a concurrently running
535 : // token_canceller returns -- which takes mutex_. Releasing while
536 : // still holding mutex_ would self-deadlock the same way
537 : // disarm_stop() would if called inside the locked loop.
538 : //
539 : // The acquire()+release() pair below nets to zero change on each
540 : // impl's refs_ -- it exists only so the lock-held loop above can
541 : // safely walk `doomed` without a concurrent recycle() tearing one
542 : // down mid-walk. It does NOT drop these impls to zero: that still
543 : // only happens when the handle's own destroy() runs, which this
544 : // context-wide shutdown does not do (shutdown() cancels/clears
545 : // registrations service-wide; it never calls destroy() per impl).
546 : // Every impl here survives at refs_ == 1 (the service's own floor
547 : // reference) until ~object_pool()'s unconditional sweep deletes
548 : // whatever is still in `live_`, bypassing refs_ entirely -- not
549 : // through this release() reaching zero.
550 167 : std::vector<posix_signal*> doomed;
551 :
552 : {
553 : posix_signal_detail::signal_state* state =
554 167 : posix_signal_detail::get_signal_state();
555 167 : std::lock_guard state_lock(state->mutex);
556 167 : std::lock_guard lock(mutex_);
557 :
558 167 : pool_.shutdown(
559 167 : [&](posix_signal* impl)
560 : {
561 6 : acquire(impl);
562 6 : doomed.push_back(impl);
563 6 : });
564 173 : for (auto* impl : doomed)
565 : {
566 12 : while (auto* reg = impl->signals_)
567 : {
568 6 : int const signal_number = reg->signal_number;
569 :
570 : // The registration table outlives every io_context, so a set
571 : // still registered here has to give its count and disposition
572 : // back the way clear() would: otherwise the signal stays
573 : // installed with these flags and the next add() of it is
574 : // refused. The per-node table unlink clear() also does is
575 : // skipped in favour of the wholesale null-out below.
576 6 : if (state->registration_count[signal_number] == 1)
577 : {
578 4 : struct sigaction sa = {};
579 4 : sa.sa_handler = SIG_DFL;
580 4 : sigemptyset(&sa.sa_mask);
581 4 : sa.sa_flags = 0;
582 4 : std::ignore = ::sigaction(signal_number, &sa, nullptr);
583 4 : state->registered_flags[signal_number] = signal_set::none;
584 : }
585 :
586 6 : --state->registration_count[signal_number];
587 6 : --registration_count_[signal_number];
588 :
589 6 : impl->signals_ = reg->next_in_set;
590 6 : delete reg;
591 6 : }
592 : }
593 :
594 : // Every live registration hung off an implementation this pool
595 : // owns, so the whole table goes stale at once and can be dropped
596 : // wholesale rather than node by node. It has to be dropped:
597 : // deliver_signal() walks this service until the destructor
598 : // unlinks it from the global list.
599 10855 : for (int i = 0; i < max_signal_number; ++i)
600 10688 : registrations_[i] = nullptr;
601 167 : }
602 :
603 173 : for (auto* impl : doomed)
604 6 : release(impl);
605 167 : }
606 :
607 : inline io_object::implementation*
608 258 : posix_signal_service::construct()
609 : {
610 258 : return pool_.acquire(*this);
611 : }
612 :
613 : inline void
614 252 : posix_signal::retire() noexcept
615 : {
616 252 : svc_.pool_.recycle(this);
617 252 : }
618 :
619 : inline std::error_code
620 277 : posix_signal_service::add_signal(
621 : posix_signal& impl, int signal_number, signal_set::flags_t flags)
622 : {
623 277 : if (signal_number < 0 || signal_number >= max_signal_number)
624 4 : return make_error_code(std::errc::invalid_argument);
625 :
626 : // Validate that requested flags are supported on this platform
627 : // (e.g., SA_NOCLDWAIT may not be available on all POSIX systems)
628 273 : if (!posix_signal_detail::flags_supported(flags))
629 MIS 0 : return make_error_code(std::errc::operation_not_supported);
630 :
631 : posix_signal_detail::signal_state* state =
632 HIT 273 : posix_signal_detail::get_signal_state();
633 :
634 : // Ensure the global self-pipe exists and this service's scheduler is
635 : // watching its read end, BEFORE taking the registration locks. The
636 : // reactor drain path locks the descriptor mutex and then the signal-state
637 : // and service mutexes; register_signal_reader locks the descriptor mutex
638 : // (via register_descriptor), so it must run holding neither of those or
639 : // the lock order would invert (a real deadlock, caught by TSan). call_once
640 : // makes the once-per-service registration safe when two signal_sets on
641 : // this context race add() from different threads.
642 : {
643 273 : std::lock_guard state_lock(state->mutex);
644 273 : if (auto ec = posix_signal_detail::open_signal_pipe(state))
645 7 : return ec;
646 273 : }
647 : {
648 : // Success-latched so a failed environmental registration
649 : // (epoll_ctl ENOMEM/ENOSPC) is retried by the next add()
650 : // instead of being lost; the code travels the return channel.
651 266 : std::lock_guard reg_lock(reader_mutex_);
652 266 : if (!reader_registered_)
653 : {
654 139 : if (auto ec = sched_->register_signal_reader(state->read_fd))
655 2 : return ec;
656 137 : reader_registered_ = true;
657 : }
658 266 : }
659 :
660 264 : std::lock_guard state_lock(state->mutex);
661 264 : std::lock_guard lock(mutex_);
662 :
663 : // Find insertion point (list is sorted by signal number)
664 264 : signal_registration** insertion_point = &impl.signals_;
665 264 : signal_registration* reg = impl.signals_;
666 287 : while (reg && reg->signal_number < signal_number)
667 : {
668 23 : insertion_point = ®->next_in_set;
669 23 : reg = reg->next_in_set;
670 : }
671 :
672 : // Already registered in this set - check flag compatibility
673 : // (same signal_set adding same signal twice with different flags)
674 264 : if (reg && reg->signal_number == signal_number)
675 : {
676 13 : if (!posix_signal_detail::flags_compatible(reg->flags, flags))
677 4 : return make_error_code(std::errc::invalid_argument);
678 9 : return {};
679 : }
680 :
681 : // Check flag compatibility with global registration
682 : // (different signal_set already registered this signal with different flags)
683 251 : if (state->registration_count[signal_number] > 0)
684 : {
685 26 : if (!posix_signal_detail::flags_compatible(
686 : state->registered_flags[signal_number], flags))
687 2 : return make_error_code(std::errc::invalid_argument);
688 : }
689 :
690 249 : auto* new_reg = new signal_registration;
691 249 : new_reg->signal_number = signal_number;
692 249 : new_reg->flags = flags;
693 249 : new_reg->owner = &impl;
694 249 : new_reg->undelivered = 0;
695 :
696 : // Install signal handler on first global registration
697 249 : if (state->registration_count[signal_number] == 0)
698 : {
699 225 : struct sigaction sa = {};
700 225 : sa.sa_handler = posix_signal_detail::corosio_posix_signal_handler;
701 225 : sigemptyset(&sa.sa_mask);
702 225 : sa.sa_flags = posix_signal_detail::to_sigaction_flags(flags);
703 :
704 225 : if (::sigaction(signal_number, &sa, nullptr) < 0)
705 : {
706 1 : delete new_reg;
707 1 : return make_error_code(std::errc::invalid_argument);
708 : }
709 :
710 : // Store the flags used for first registration
711 224 : state->registered_flags[signal_number] = flags;
712 : }
713 :
714 248 : new_reg->next_in_set = reg;
715 248 : *insertion_point = new_reg;
716 :
717 248 : new_reg->next_in_table = registrations_[signal_number];
718 248 : new_reg->prev_in_table = nullptr;
719 248 : if (registrations_[signal_number])
720 18 : registrations_[signal_number]->prev_in_table = new_reg;
721 248 : registrations_[signal_number] = new_reg;
722 :
723 248 : ++state->registration_count[signal_number];
724 248 : ++registration_count_[signal_number];
725 :
726 248 : return {};
727 264 : }
728 :
729 : inline std::error_code
730 26 : posix_signal_service::remove_signal(posix_signal& impl, int signal_number)
731 : {
732 26 : if (signal_number < 0 || signal_number >= max_signal_number)
733 2 : return make_error_code(std::errc::invalid_argument);
734 :
735 : posix_signal_detail::signal_state* state =
736 24 : posix_signal_detail::get_signal_state();
737 24 : std::lock_guard state_lock(state->mutex);
738 24 : std::lock_guard lock(mutex_);
739 :
740 24 : signal_registration** deletion_point = &impl.signals_;
741 24 : signal_registration* reg = impl.signals_;
742 26 : while (reg && reg->signal_number < signal_number)
743 : {
744 2 : deletion_point = ®->next_in_set;
745 2 : reg = reg->next_in_set;
746 : }
747 :
748 24 : if (!reg || reg->signal_number != signal_number)
749 3 : return {};
750 :
751 : // Restore default handler on last global unregistration
752 21 : if (state->registration_count[signal_number] == 1)
753 : {
754 17 : struct sigaction sa = {};
755 17 : sa.sa_handler = SIG_DFL;
756 17 : sigemptyset(&sa.sa_mask);
757 17 : sa.sa_flags = 0;
758 :
759 17 : if (::sigaction(signal_number, &sa, nullptr) < 0)
760 1 : return make_error_code(std::errc::invalid_argument);
761 :
762 : // Clear stored flags
763 16 : state->registered_flags[signal_number] = signal_set::none;
764 : }
765 :
766 20 : *deletion_point = reg->next_in_set;
767 :
768 20 : if (registrations_[signal_number] == reg)
769 18 : registrations_[signal_number] = reg->next_in_table;
770 20 : if (reg->prev_in_table)
771 2 : reg->prev_in_table->next_in_table = reg->next_in_table;
772 20 : if (reg->next_in_table)
773 2 : reg->next_in_table->prev_in_table = reg->prev_in_table;
774 :
775 20 : --state->registration_count[signal_number];
776 20 : --registration_count_[signal_number];
777 :
778 20 : delete reg;
779 20 : return {};
780 24 : }
781 :
782 : inline std::error_code
783 266 : posix_signal_service::clear_signals(posix_signal& impl)
784 : {
785 : posix_signal_detail::signal_state* state =
786 266 : posix_signal_detail::get_signal_state();
787 266 : std::lock_guard state_lock(state->mutex);
788 266 : std::lock_guard lock(mutex_);
789 :
790 266 : std::error_code first_error;
791 :
792 488 : while (signal_registration* reg = impl.signals_)
793 : {
794 222 : int signal_number = reg->signal_number;
795 :
796 222 : if (state->registration_count[signal_number] == 1)
797 : {
798 204 : struct sigaction sa = {};
799 204 : sa.sa_handler = SIG_DFL;
800 204 : sigemptyset(&sa.sa_mask);
801 204 : sa.sa_flags = 0;
802 :
803 204 : if (::sigaction(signal_number, &sa, nullptr) < 0 && !first_error)
804 1 : first_error = make_error_code(std::errc::invalid_argument);
805 :
806 : // Clear stored flags
807 204 : state->registered_flags[signal_number] = signal_set::none;
808 : }
809 :
810 222 : impl.signals_ = reg->next_in_set;
811 :
812 222 : if (registrations_[signal_number] == reg)
813 220 : registrations_[signal_number] = reg->next_in_table;
814 222 : if (reg->prev_in_table)
815 2 : reg->prev_in_table->next_in_table = reg->next_in_table;
816 222 : if (reg->next_in_table)
817 12 : reg->next_in_table->prev_in_table = reg->prev_in_table;
818 :
819 222 : --state->registration_count[signal_number];
820 222 : --registration_count_[signal_number];
821 :
822 222 : delete reg;
823 222 : }
824 :
825 266 : if (first_error)
826 1 : return first_error;
827 265 : return {};
828 266 : }
829 :
830 : inline void
831 269 : posix_signal_service::cancel_wait(posix_signal& impl)
832 : {
833 269 : bool was_waiting = false;
834 269 : signal_op* op = nullptr;
835 :
836 : {
837 269 : std::lock_guard lock(mutex_);
838 269 : impl.cancelled_ = true;
839 269 : if (impl.waiting_)
840 : {
841 7 : was_waiting = true;
842 7 : impl.waiting_ = false;
843 7 : op = &impl.pending_op_;
844 : }
845 269 : }
846 :
847 269 : if (was_waiting)
848 : {
849 7 : if (op->ec_out)
850 7 : *op->ec_out = make_error_code(capy::error::canceled);
851 7 : if (op->signal_out)
852 7 : *op->signal_out = 0;
853 7 : op->cont.h = op->h;
854 7 : op->d.post(op->cont);
855 7 : sched_->work_finished();
856 : }
857 269 : }
858 :
859 : inline void
860 768 : posix_signal_service::cancel_wait_token(posix_signal& impl) noexcept
861 : {
862 768 : bool was_waiting = false;
863 768 : signal_op* op = nullptr;
864 :
865 : {
866 768 : std::lock_guard lock(mutex_);
867 : // Persist the request even when no wait is parked yet: wait()
868 : // arms the callback before start_wait takes this lock, and
869 : // start_wait consumes this flag.
870 768 : impl.token_cancelled_ = true;
871 768 : if (impl.waiting_)
872 : {
873 726 : was_waiting = true;
874 726 : impl.waiting_ = false;
875 726 : op = &impl.pending_op_;
876 : }
877 768 : }
878 :
879 768 : if (was_waiting)
880 : {
881 726 : if (op->ec_out)
882 726 : *op->ec_out = make_error_code(capy::error::canceled);
883 726 : if (op->signal_out)
884 726 : *op->signal_out = 0;
885 726 : op->cont.h = op->h;
886 726 : op->d.post(op->cont);
887 726 : sched_->work_finished();
888 : }
889 768 : }
890 :
891 : inline void
892 768 : posix_signal::token_canceller::operator()() const noexcept
893 : {
894 768 : self->svc_.cancel_wait_token(*self);
895 768 : }
896 :
897 : inline void
898 1104 : posix_signal_service::start_wait(posix_signal& impl, signal_op* op)
899 : {
900 : {
901 1104 : std::lock_guard lock(mutex_);
902 :
903 : // Check if cancel() was called before this wait started
904 1104 : if (impl.cancelled_)
905 : {
906 2 : impl.cancelled_ = false;
907 2 : if (op->ec_out)
908 2 : *op->ec_out = make_error_code(capy::error::canceled);
909 2 : if (op->signal_out)
910 2 : *op->signal_out = 0;
911 2 : op->cont.h = op->h;
912 2 : op->d.post(op->cont);
913 2 : return;
914 : }
915 :
916 : // A stop request that arrived between wait() arming the callback
917 : // and this lock: complete now rather than parking forever.
918 1102 : if (impl.token_cancelled_)
919 : {
920 40 : impl.token_cancelled_ = false;
921 40 : if (op->ec_out)
922 40 : *op->ec_out = make_error_code(capy::error::canceled);
923 40 : if (op->signal_out)
924 40 : *op->signal_out = 0;
925 40 : op->cont.h = op->h;
926 40 : op->d.post(op->cont);
927 40 : return;
928 : }
929 :
930 : // Check for queued signals first (signal arrived before wait started)
931 1062 : signal_registration* reg = impl.signals_;
932 2126 : while (reg)
933 : {
934 1066 : if (reg->undelivered > 0)
935 : {
936 2 : --reg->undelivered;
937 2 : op->signal_number = reg->signal_number;
938 : // svc=nullptr: no work_finished needed since we never called work_started
939 2 : op->svc = nullptr;
940 2 : sched_->post(op);
941 2 : return;
942 : }
943 1064 : reg = reg->next_in_set;
944 : }
945 :
946 : // No queued signals - wait for delivery
947 1060 : impl.waiting_ = true;
948 : // svc=this: signal_op::operator() will call work_finished() to balance this
949 1060 : op->svc = this;
950 1060 : sched_->work_started();
951 1104 : }
952 : }
953 :
954 : inline void
955 319 : posix_signal_service::deliver_signal(int signal_number)
956 : {
957 319 : if (signal_number < 0 || signal_number >= max_signal_number)
958 MIS 0 : return;
959 :
960 : posix_signal_detail::signal_state* state =
961 HIT 319 : posix_signal_detail::get_signal_state();
962 319 : std::lock_guard lock(state->mutex);
963 :
964 319 : posix_signal_service* service = state->service_list;
965 638 : while (service)
966 : {
967 319 : std::lock_guard svc_lock(service->mutex_);
968 :
969 319 : signal_registration* reg = service->registrations_[signal_number];
970 642 : while (reg)
971 : {
972 323 : posix_signal* impl = static_cast<posix_signal*>(reg->owner);
973 :
974 323 : if (impl->waiting_)
975 : {
976 321 : impl->waiting_ = false;
977 321 : impl->pending_op_.signal_number = signal_number;
978 321 : service->post(&impl->pending_op_);
979 : }
980 : else
981 : {
982 2 : ++reg->undelivered;
983 : }
984 :
985 323 : reg = reg->next_in_table;
986 : }
987 :
988 319 : service = service->next_;
989 319 : }
990 319 : }
991 :
992 : inline void
993 : posix_signal_service::work_started() noexcept
994 : {
995 : sched_->work_started();
996 : }
997 :
998 : inline void
999 321 : posix_signal_service::work_finished() noexcept
1000 : {
1001 321 : sched_->work_finished();
1002 321 : }
1003 :
1004 : inline void
1005 321 : posix_signal_service::post(signal_op* op)
1006 : {
1007 321 : sched_->post(op);
1008 321 : }
1009 :
1010 : inline void
1011 167 : posix_signal_service::add_service(posix_signal_service* service)
1012 : {
1013 : posix_signal_detail::signal_state* state =
1014 167 : posix_signal_detail::get_signal_state();
1015 167 : std::lock_guard lock(state->mutex);
1016 :
1017 167 : service->next_ = state->service_list;
1018 167 : service->prev_ = nullptr;
1019 167 : if (state->service_list)
1020 8 : state->service_list->prev_ = service;
1021 167 : state->service_list = service;
1022 167 : }
1023 :
1024 : inline void
1025 167 : posix_signal_service::remove_service(posix_signal_service* service)
1026 : {
1027 : posix_signal_detail::signal_state* state =
1028 167 : posix_signal_detail::get_signal_state();
1029 167 : std::lock_guard lock(state->mutex);
1030 :
1031 167 : if (service->next_ || service->prev_ || state->service_list == service)
1032 : {
1033 167 : if (state->service_list == service)
1034 165 : state->service_list = service->next_;
1035 167 : if (service->prev_)
1036 2 : service->prev_->next_ = service->next_;
1037 167 : if (service->next_)
1038 6 : service->next_->prev_ = service->prev_;
1039 167 : service->next_ = nullptr;
1040 167 : service->prev_ = nullptr;
1041 : }
1042 167 : }
1043 :
1044 : } // namespace detail
1045 : } // namespace boost::corosio
1046 :
1047 : #endif // BOOST_COROSIO_POSIX
1048 :
1049 : #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_SIGNAL_SERVICE_HPP
|