diff --git a/CMakeLists.txt b/CMakeLists.txt index d29fb27a1..369cbd968 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -146,7 +146,32 @@ if(NOT DEFINED CPM_SOURCE_CACHE AND NOT DEFINED ENV{CPM_SOURCE_CACHE}) set(CPM_SOURCE_CACHE "${CMAKE_CURRENT_SOURCE_DIR}/.cache/cpm" CACHE PATH "Where CPM keeps the sources of every fetched dependency") endif() -include(cmake/CPM.cmake) + +# CPM itself is loaded only when something has to be fetched: its bootstrap +# downloads CPM.cmake, so loading it unconditionally would put the network on +# the path of a configure whose every dependency is already installed -- the +# shape a package-manager build (vcpkg, Conan, a distribution) has, where +# reaching the network is not allowed. Each fetch site calls +# morph_use_cpm( []) before its CPMAddPackage; the first one to +# load CPM names itself in a STATUS line, so a failed bootstrap download says +# which dependency asked for it. A parent project that has loaded CPM already +# is reused, not loaded again. +# +# A macro, not a function, so CPM's own non-cache variables land in the +# caller's scope rather than vanishing with a function's. +set(MORPH_CPM_BOOTSTRAP "${CMAKE_CURRENT_SOURCE_DIR}/cmake/CPM.cmake") +macro(morph_use_cpm _morph_cpm_for) + if(NOT COMMAND CPMAddPackage) + if(${ARGC} GREATER 1) + set(_morph_cpm_why "${ARGV1}") + else() + set(_morph_cpm_why "not found installed (install it and point CMAKE_PREFIX_PATH at it to configure without the network)") + endif() + message(STATUS "morph: loading CPM to fetch ${_morph_cpm_for}: ${_morph_cpm_why}") + unset(_morph_cpm_why) + include("${MORPH_CPM_BOOTSTRAP}") + endif() +endmacro() # ── glaze ──────────────────────────────────────────────────────────────────── # On Windows vcpkg provides glaze. Elsewhere CPM fetches it, so that CI does @@ -165,6 +190,7 @@ include(cmake/CPM.cmake) set(MORPH_GLAZE_VERSION 7.4) find_package(glaze ${MORPH_GLAZE_VERSION} CONFIG QUIET) if(NOT glaze_FOUND) + morph_use_cpm(glaze) CPMAddPackage( NAME glaze GITHUB_REPOSITORY stephenberry/glaze @@ -184,40 +210,71 @@ endif() # # morph stays a header-only INTERFACE target, but core::base, core::net and # core::platform are static libraries, so every consumer of morph now builds -# them. +# them -- or links an installed core-cpp's. +# +# Found first, as glaze is, and fetched only when no installed core-cpp +# satisfies the bound. MORPH_CORE_CPP_VERSION is that bound, stated once +# because morphConfig.cmake asks for the same one; core-cpp's package accepts +# only its own minor version while it is 0.x. +# +# Except under a sanitizer: an installed core-cpp was compiled without morph's +# instrumentation, and a ThreadSanitizer that cannot see core::net's side of +# TimeoutScheduler's hand-off reports races that are not there and misses ones +# that are. A sanitizer leg therefore always builds core-cpp in-tree, so the +# block below can instrument it. # # morph's install exports morph::morph, which links core-cpp's modules, so an -# install of morph has to install core-cpp too: CMake refuses an export that -# names a target in no export set. CORE_CPP_INSTALL follows MORPH_INSTALL for -# that, and morphConfig.cmake finds the installed core-cpp package in turn. +# install of morph with a fetched core-cpp has to install core-cpp too: CMake +# refuses an export that names a target in no export set. CORE_CPP_INSTALL +# follows MORPH_INSTALL for that, and morphConfig.cmake finds the installed +# core-cpp package in turn. A found core-cpp is already installed, its targets +# are imported, and the export names them without installing anything again. # MORPH_INSTALL is described with the install rules below. EXCLUDE_FROM_ALL # is off while morph installs: CMake leaves an excluded subdirectory's install # rules out of the parent's install, so `cmake --install` would install # morph's package without the core-cpp package it depends on. +set(MORPH_CORE_CPP_VERSION 0.5) option(MORPH_INSTALL "Generate morph's install and export rules" ${PROJECT_IS_TOP_LEVEL}) -if(MORPH_INSTALL) - set(_morph_core_cpp_exclude_from_all NO) -else() - set(_morph_core_cpp_exclude_from_all YES) +if(NOT DEFINED AF_SANITIZER) + find_package(core-cpp ${MORPH_CORE_CPP_VERSION} CONFIG QUIET) +endif() +if(NOT core-cpp_FOUND) + if(MORPH_INSTALL) + set(_morph_core_cpp_exclude_from_all NO) + else() + set(_morph_core_cpp_exclude_from_all YES) + endif() + if(DEFINED AF_SANITIZER) + morph_use_cpm(core-cpp "AF_SANITIZER builds it in-tree so that it can be instrumented") + else() + morph_use_cpm(core-cpp) + endif() + CPMAddPackage( + NAME core-cpp + GITHUB_REPOSITORY contour-terminal/core-cpp + GIT_TAG v0.5.0 + VERSION 0.5.0 + SYSTEM YES + EXCLUDE_FROM_ALL ${_morph_core_cpp_exclude_from_all} + OPTIONS "CORE_CPP_TESTING OFF" "CORE_CPP_BUILD_EXAMPLES OFF" "CORE_CPP_WITH_TUI OFF" + "CORE_CPP_WITH_TLS OFF" "CORE_CPP_FETCH_DEPS OFF" "CORE_CPP_INSTALL ${MORPH_INSTALL}") + unset(_morph_core_cpp_exclude_from_all) endif() -CPMAddPackage( - NAME core-cpp - GITHUB_REPOSITORY contour-terminal/core-cpp - GIT_TAG v0.5.0 - VERSION 0.5.0 - SYSTEM YES - EXCLUDE_FROM_ALL ${_morph_core_cpp_exclude_from_all} - OPTIONS "CORE_CPP_TESTING OFF" "CORE_CPP_BUILD_EXAMPLES OFF" "CORE_CPP_WITH_TUI OFF" - "CORE_CPP_WITH_TLS OFF" "CORE_CPP_FETCH_DEPS OFF" "CORE_CPP_INSTALL ${MORPH_INSTALL}") -unset(_morph_core_cpp_exclude_from_all) # A sanitizer leg instruments core-cpp's compiled modules with morph's own -# targets: TimeoutScheduler's loop thread runs inside core::net, and a -# ThreadSanitizer that cannot see one side of a hand-off reports races that -# are not there and misses ones that are. core-cpp lists those targets in the -# CORE_CPP_TARGETS global property for exactly this. +# targets (see above for why). core-cpp lists those targets in the +# CORE_CPP_TARGETS global property for exactly this. The property is set only +# by an in-tree core-cpp; an empty one here means nothing would be +# instrumented while the configure reported success, so it stops instead -- +# e.g. a parent project that brought core-cpp in some other way. if(DEFINED AF_SANITIZER) get_property(_morph_core_cpp_targets GLOBAL PROPERTY CORE_CPP_TARGETS) + if(NOT _morph_core_cpp_targets) + message(FATAL_ERROR + "AF_SANITIZER=${AF_SANITIZER} but core-cpp lists no targets to instrument " + "(the CORE_CPP_TARGETS global property is empty): core-cpp did not come from an " + "in-tree build, so its compiled modules would run uninstrumented.") + endif() foreach(_morph_core_cpp_target IN LISTS _morph_core_cpp_targets) apply_sanitizers(${_morph_core_cpp_target} ${AF_SANITIZER}) endforeach() @@ -651,6 +708,7 @@ endif() if(MORPH_BUILD_TESTS) find_package(Catch2 CONFIG QUIET) if(NOT Catch2_FOUND) + morph_use_cpm(Catch2) CPMAddPackage( NAME Catch2 GITHUB_REPOSITORY catchorg/Catch2 diff --git a/cmake/CPM.cmake b/cmake/CPM.cmake index 729c2a6c0..4c6c0018f 100644 --- a/cmake/CPM.cmake +++ b/cmake/CPM.cmake @@ -11,6 +11,14 @@ # # With CPM_SOURCE_CACHE set (CMakeLists.txt defaults it to .cache/cpm) the # bootstrap itself lives in the cache, so a warm cache downloads nothing. +# +# morph includes this only through morph_use_cpm() (CMakeLists.txt), when a +# dependency was not found installed and has to be fetched; that call's STATUS +# line, printed just before, names the dependency. The FATAL_ERROR below is +# core-cpp's wording. At morph's top level, its "set CORE_CPP_FETCH_DEPS=OFF" +# does not apply and cmake/FetchTransferBound.cmake does not exist: the way to +# configure without this download is to install the dependency that line names +# and point CMAKE_PREFIX_PATH at it. set(_coreCppCpmBound "") if(DEFINED FASTCACHED_FETCH_SILENCE_SECONDS) set(_coreCppCpmBound INACTIVITY_TIMEOUT "${FASTCACHED_FETCH_SILENCE_SECONDS}") diff --git a/cmake/morphConfig.cmake.in b/cmake/morphConfig.cmake.in index e010aedf3..0726057e2 100644 --- a/cmake/morphConfig.cmake.in +++ b/cmake/morphConfig.cmake.in @@ -29,10 +29,11 @@ endif() find_dependency(glaze @MORPH_GLAZE_VERSION@ CONFIG) # morph::morph links core-cpp's modules (core::base, core::async, core::net and, -# natively, core::platform), which morph's install puts next to it. 0.3 is the -# minor version morph was built against; core-cpp's package accepts only that -# minor version while core-cpp is 0.x. -find_dependency(core-cpp 0.5 CONFIG) +# natively, core::platform): the core-cpp morph was built against, either +# installed by morph's own install next to it or already installed when the +# build found it. The bound is the minor version morph was built against; +# core-cpp's package accepts only that minor version while core-cpp is 0.x. +find_dependency(core-cpp @MORPH_CORE_CPP_VERSION@ CONFIG) if(NOT WIN32 AND NOT EMSCRIPTEN) find_dependency(Threads) diff --git a/docs/CMakeLists.txt b/docs/CMakeLists.txt index d74204bb2..cb6b4cdb3 100644 --- a/docs/CMakeLists.txt +++ b/docs/CMakeLists.txt @@ -2,8 +2,9 @@ find_package(Doxygen REQUIRED) message(STATUS "Doxygen found: ${DOXYGEN_EXECUTABLE}") # DOWNLOAD_ONLY: the repository is a stylesheet, and the one file used is -# `doxygen-awesome.css`; there is nothing to configure. CPM is loaded by the -# root CMakeLists.txt, and its source cache applies here as everywhere. +# `doxygen-awesome.css`; there is nothing to configure. CPM is loaded on demand by +# the root CMakeLists.txt's morph_use_cpm, and its source cache applies here as everywhere. +morph_use_cpm(doxygen-awesome-css) CPMAddPackage( NAME doxygen-awesome-css GITHUB_REPOSITORY jothepro/doxygen-awesome-css diff --git a/docs/spec/core/backend.md b/docs/spec/core/backend.md index 126c1d36d..6ddf441ea 100644 --- a/docs/spec/core/backend.md +++ b/docs/spec/core/backend.md @@ -582,7 +582,11 @@ would silently change `registerHandler`'s contract. See [Waiting for a bind — `bindWaitPolicy`](#waiting-for-a-bind--bindwaitpolicy). `executeVia` fails fast with `"handler not bound"` for a call issued before the -reply arrives; it does not queue. +reply arrives; it does not queue. A caller that wants the call held instead +names that at the call site with `BridgeHandler::executeWhenBound()`, which +chains the dispatch on `whenBound()` ([bridge.md](bridge.md#registration-readiness--isbound--whenbound)). +The default stays fail-fast so that `execute()` means the same thing whichever +backend the handler was built over. Consequences: diff --git a/docs/spec/core/bridge.md b/docs/spec/core/bridge.md index a454383dd..d0ba29771 100644 --- a/docs/spec/core/bridge.md +++ b/docs/spec/core/bridge.md @@ -16,6 +16,7 @@ that only know action names at runtime. - [The bridge's own executor](#the-bridges-own-executor) - [`BridgeHandler`](#bridgehandlermodel) - [Registration readiness — `isBound()` / `whenBound()`](#registration-readiness--isbound--whenbound) + - [`executeWhenBound()` — holding a dispatch until the bind lands](#executewhenbound--holding-a-dispatch-until-the-bind-lands) - [`ActionExecuteRegistry`](#actionexecuteregistry) - [Why the key carries the sharing policy](#why-the-key-carries-the-sharing-policy) - [`BRIDGE_REGISTER_ACTION` and `registerActionExecutorOnce`](#bridge_register_action-and-registeractionexecutoronce) @@ -756,6 +757,53 @@ it answers and does not: this binding an id, not that the transport is still up; the socket can drop the moment after. +### `executeWhenBound()` — holding a dispatch until the bind lands + +`execute()` never waits. A view model that constructs a handler and dispatches +its first action in the same breath would otherwise wrap every handler it owns +in the same gate — `whenBound()`, then the dispatch, plus a liveness guard in +case the handler goes first. `BridgeHandler::executeWhenBound(action)` is that +gate, built once: + +| Handler state when called | Result | +|---|---| +| Bound | Dispatches exactly as `execute()` does. | +| Registration in flight, then succeeds | Dispatches on the GUI executor once `whenBound()` resolves `true`. | +| Registration in flight, then fails | Rejects with the registration's error; nothing is dispatched. | +| `whenBound()` resolves `false` | Rejects with `"handler not bound"`. | +| Handler destroyed before the dispatch comes due | The held action is dropped: never dispatched, and the returned `Completion` never settles. | +| Bridge retired before the dispatch comes due | Rejects with `"bridge destroyed"`. | + +Why each choice was made: + +- **A separate method, not a constructor policy.** With a policy, one + `execute()` call would wait or fail depending on how the handler had been + built somewhere else, and a reader of the call site could not tell which. + The name shows the wait where it happens. +- **`NoSharing` only**, enforced with a `static_assert`. A shared handler's + initial binding goes through the attach path, which `whenBound()` does not + track (see the scope limits above), so on a shared handler it would resolve + `false` at once and the method would be fail-fast under another name. A + shared handler's keyed `execute()` already carries its attach. +- **The held dispatch holds the binding weakly and does not capture the + handler.** The binding owns the `whenBound()` waiter, and the waiter owns the + held dispatch; a strong capture of the binding would close that loop and + keep the binding alive after its handler is gone — the dispatch would then + reach the backend on an instance nobody holds. With a weak capture, a + handler destroyed first takes the waiter (and the held action) with it. +- **The bridge gate is read, then released, before the dispatch.** The + dispatch can settle inline, and `BridgeSink`'s settle path takes the same + gate; holding it across the call would lock it recursively. The check makes + the ordinary teardown order — bridge retired while the dispatch sat in the + GUI queue — a rejection instead of a call into a destroyed bridge. It is not + a licence to destroy the bridge concurrently with the dispatch: the bridge + must outlive the dispatch on the same terms as any other call made on the + handler. +- **The result type must be copy-constructible** (a `static_assert`). The + deferred path forwards the value from `executeVia`'s `Completion` into the + one already handed to the caller, and a `Completion`'s value is observed, + never consumed ([completion.md](completion.md)). + ## `ActionExecuteRegistry` @@ -1142,6 +1190,7 @@ make teardown order-independent.) | ctor (custom binding) | `BridgeHandler(Bridge&, IExecutor*, shared_ptr)` | Registers pre-built binding. | | dtor | `~BridgeHandler()` | Deregisters via `Bridge::deregisterHandler`, but only if the bridge's `CallbackToken` is still active; a no-op if the `Bridge` was already destroyed. | | `execute` | `Completion execute(Action)` | Typed dispatch through the bridge. For a shared handler, a payload-/result-keyed action's attach or promote step never throws synchronously — a backend refusal (e.g. `LimitPolicy::maxLiveModels`) resolves the returned `Completion` via `.onError(...)`. | +| `executeWhenBound` | `Completion executeWhenBound(Action)` | `NoSharing` only. Dispatches like `execute()` when bound; otherwise holds the action until `whenBound()` settles, then dispatches it or rejects with the registration's error (or `"handler not bound"`). Dropped if the handler is destroyed first. See [`executeWhenBound()`](#executewhenbound--holding-a-dispatch-until-the-bind-lands). | | `executeJson` | `Completion executeJson(string_view actionType, string_view bodyJson)` | Type-erased dispatch through `ActionExecuteRegistry`. | | `subscribe(cb)` | `void subscribe(function)` | Fire `cb` whenever an `R` is produced on the attached instance. | | `subscribe(scope, cb)` | `void subscribe(CallbackScope const&, function)` | As above, gated on the scope's liveness and stop state ([callback_scope.md](callback_scope.md)). Dead sinks are refused, not pruned. | diff --git a/examples/bank/CMakeLists.txt b/examples/bank/CMakeLists.txt index a281838e2..90347ea50 100644 --- a/examples/bank/CMakeLists.txt +++ b/examples/bank/CMakeLists.txt @@ -53,6 +53,7 @@ set(LIGHTWEIGHT_BUILD_BENCHMARK OFF CACHE BOOL "" FORCE) # Lightweight fetch. set(_morph_saved_skip_install_rules ${CMAKE_SKIP_INSTALL_RULES}) set(CMAKE_SKIP_INSTALL_RULES ON) +morph_use_cpm(Lightweight) CPMAddPackage( NAME Lightweight GITHUB_REPOSITORY LASTRADA-Software/Lightweight diff --git a/examples/common/CMakeLists.txt b/examples/common/CMakeLists.txt index 3de1b3fc0..fcec8f230 100644 --- a/examples/common/CMakeLists.txt +++ b/examples/common/CMakeLists.txt @@ -147,6 +147,7 @@ set(LIGHTWEIGHT_BUILD_SHARED OFF CACHE BOOL "" FORCE) # without touching Lightweight's vendored CMakeLists.txt. set(_morph_saved_skip_install_rules ${CMAKE_SKIP_INSTALL_RULES}) set(CMAKE_SKIP_INSTALL_RULES ON) +morph_use_cpm(Lightweight) CPMAddPackage( NAME Lightweight GITHUB_REPOSITORY LASTRADA-Software/Lightweight diff --git a/include/morph/core/bridge.hpp b/include/morph/core/bridge.hpp index 17aad8ba4..1c1690447 100644 --- a/include/morph/core/bridge.hpp +++ b/include/morph/core/bridge.hpp @@ -3232,6 +3232,89 @@ class BridgeHandler { } } + /// @brief Dispatches @p action once this handler's in-flight registration + /// settles, instead of failing fast with "handler not bound". + /// + /// The waiting counterpart of `execute()` for a handler built over a backend + /// that answers `BindWait::kCallerMustNotBlock`, where a handler is + /// constructed unbound and the first dispatch commonly follows on the next + /// line. `execute()` itself never waits; the wait is named at the call site. + /// + /// - Already bound: dispatches exactly as `execute()` does. + /// - A registration in flight: the action is held until `whenBound()` + /// settles, then dispatched on this handler's GUI executor. A failed + /// registration rejects the returned `Completion` with the registration's + /// error; a registration that turns out to have nothing in flight rejects + /// it with "handler not bound". + /// - Handler destroyed before the registration settles: the held action is + /// dropped, never dispatched, and the returned `Completion` never settles. + /// The deferred dispatch holds the binding weakly and does not capture + /// the handler, so destroying the handler is always safe. + /// + /// The `Bridge` must outlive the deferred dispatch on the same terms as any + /// other call made on this handler. A bridge already retired when the + /// dispatch comes due rejects it with "bridge destroyed" rather than calling + /// into it. + /// + /// Only `NoSharing` handlers: a shared handler's initial binding goes + /// through the attach path, which `whenBound()` does not track, so a wait + /// there would resolve `false` at once. A shared handler waits through the + /// `Completion` its keyed `execute()` returns instead. + /// + /// @tparam Action Concrete action type registered with `BRIDGE_REGISTER_ACTION`. + /// Its result type must be copy-constructible, because the + /// deferred path forwards the value from one `Completion` to + /// another. + /// @param action Action to execute (moved into the dispatch, or held until + /// the registration settles). + /// @return Completion that resolves on the GUI executor with the action's + /// result, or rejects with the dispatch's error, the registration's + /// error, or "handler not bound" as described above. + template + ::morph::async::Completion::Result> executeWhenBound(Action action) { + static_assert(!kShared, + "executeWhenBound() is for NoSharing handlers; a shared handler's keyed execute() already " + "carries its attach"); + using R = ::morph::model::ActionTraits::Result; + static_assert(std::is_copy_constructible_v, + "executeWhenBound() forwards the result between completions, so it must be copyable"); + if (isBound()) { + return _bridge.template executeVia(_binding, std::move(action), _guiExec); + } + auto state = std::make_shared<::morph::async::detail::CompletionState>(); + ::morph::async::Completion pending{state, _guiExec}; + auto sharedAction = std::make_shared(std::move(action)); + std::weak_ptr const weakBinding{_binding}; + whenBound() + .then([weakBinding, bridgePtr = &_bridge, gate = _bridgeGate, guiExec = _guiExec, sharedAction, + state](bool bound) { + auto binding = weakBinding.lock(); + if (!binding) { + return; // The handler is gone; drop the held action. + } + if (!bound) { + state->setException(std::make_exception_ptr(std::runtime_error("handler not bound"))); + return; + } + bool bridgeAlive = false; + { + // Released before the dispatch: executeVia can settle + // inline, and its settle path takes this same gate. + std::shared_lock const lock{gate->mtx}; + bridgeAlive = gate->alive; + } + if (!bridgeAlive) { + state->setException(std::make_exception_ptr(std::runtime_error("bridge destroyed"))); + return; + } + bridgePtr->template executeVia(binding, std::move(*sharedAction), guiExec) + .then([state](const R& value) { state->setValue(R{value}); }) + .onError([state](std::exception_ptr exc) { state->setException(exc); }); + }) + .onError([state](std::exception_ptr err) { state->setException(err); }); + return pending; + } + /// @brief Attaches (or re-points) this handler to the instance for @p key. /// /// Creates the instance if no live instance holds @p key, otherwise joins the diff --git a/include/morph/net/detail/tcp_socket.hpp b/include/morph/net/detail/tcp_socket.hpp index d9365d9cc..1b68628bb 100644 --- a/include/morph/net/detail/tcp_socket.hpp +++ b/include/morph/net/detail/tcp_socket.hpp @@ -31,6 +31,72 @@ namespace morph::net::detail { +/// @brief How a `pollUntil()` wait ended. +enum class PollOutcome : std::uint8_t { + kReady, ///< `poll()` reported the descriptor (readiness, hang-up or error bits alike). + kTimedOut, ///< The deadline passed first. + kFailed, ///< `poll()` failed with something other than `EINTR`; see `PollResult::error`. +}; + +/// @brief The result of a `pollUntil()` wait. +struct PollResult { + /// @brief How the wait ended. + PollOutcome outcome = PollOutcome::kTimedOut; + /// @brief The `errno` of the failing `poll()` when `outcome` is `kFailed`; `0` otherwise. + int error = 0; +}; + +/// @brief Waits for @p events on @p fd until @p deadline, with one budget for the whole wait. +/// +/// A signal that interrupts `poll()` (`EINTR`) does not end the wait and does +/// not re-arm the full timeout: the loop recomputes what is left of the budget +/// and polls again, so a host delivering signals steadily (a profiler's timer, +/// `SIGCHLD`) cannot stretch the wait past @p deadline. The per-call timeout is +/// clamped to `int` milliseconds, which is what `poll()` takes; an unclamped +/// multi-week budget would truncate to garbage, possibly negative, meaning +/// "wait forever". +/// +/// Reports rather than throws, so each caller keeps its own policy for a +/// timeout or a failure (a connect moves on to the next address candidate; a +/// handshake read gives up). +/// +/// @param descriptor Descriptor to wait on. +/// @param events `poll()` event mask, e.g. `POLLIN` or `POLLOUT`. +/// @param deadline Point after which the wait reports `kTimedOut`. A deadline +/// already passed reports it without polling. +/// @return The outcome, plus the failing `errno` for `kFailed`. +// A descriptor and an event mask are different things that happen to convert; the +// parameter names carry the distinction. +// NOLINTNEXTLINE(bugprone-easily-swappable-parameters) +inline PollResult pollUntil(int descriptor, short events, std::chrono::steady_clock::time_point deadline) { + pollfd pfd{}; + pfd.fd = descriptor; + pfd.events = events; + for (;;) { + // Rounded up: poll() takes whole milliseconds, and truncating would wake + // it up to a millisecond before the deadline and report a timeout early. + auto const remaining = + std::chrono::ceil(deadline - std::chrono::steady_clock::now()); + if (remaining.count() <= 0) { + return {.outcome = PollOutcome::kTimedOut}; + } + auto const waitMs = static_cast( + std::min(remaining.count(), std::numeric_limits::max())); + int const ready = ::poll(&pfd, 1, waitMs); + if (ready > 0) { + return {.outcome = PollOutcome::kReady}; + } + if (ready < 0) { + int const err = errno; + if (err != EINTR) { + return {.outcome = PollOutcome::kFailed, .error = err}; + } + } + // EINTR, or a poll that returned early with nothing ready (the + // millisecond rounding above): recompute what is left and go again. + } +} + /// @brief RAII wrapper around a POSIX (BSD sockets) TCP file descriptor. /// /// Linux/macOS only today — see `docs/spec/core/backend.md`'s `morph::net` @@ -168,34 +234,10 @@ class TcpSocket { ::close(fd); continue; } - pollfd pfd{}; - pfd.fd = fd; - pfd.events = POLLOUT; - // Retry on EINTR rather than treating a delivered signal as a - // connect failure. Every other blocking syscall in this file already - // does (accept, tryAccept, recvSome, sendAll), and the accept loop's - // own comment makes the case: "any signal the host happens to - // deliver (a profiler's timer, SIGCHLD, SIGWINCH)" must not tear the - // operation down. This poll was the one that did not honour it. - int pollRc = 0; - for (;;) { - auto const remaining = - std::chrono::duration_cast(deadline - std::chrono::steady_clock::now()); - if (remaining.count() <= 0) { - pollRc = 0; // deadline reached: treat as timeout - break; - } - // Clamped: poll takes int milliseconds, and a multi-week timeout - // would otherwise truncate to a garbage (possibly negative, - // i.e. infinite) value. - auto const waitMs = static_cast( - std::min(remaining.count(), std::numeric_limits::max())); - pollRc = ::poll(&pfd, 1, waitMs); - if (pollRc >= 0 || errno != EINTR) { - break; - } - } - if (pollRc <= 0) { + // A timeout or a poll failure on this candidate moves on to the + // next; only running out of candidates is an error. A delivered + // signal is neither (pollUntil retries across EINTR). + if (pollUntil(fd, POLLOUT, deadline).outcome != PollOutcome::kReady) { ::close(fd); continue; } diff --git a/include/morph/net/detail/ws_handshake.hpp b/include/morph/net/detail/ws_handshake.hpp index cadb6d586..d3f992b4b 100644 --- a/include/morph/net/detail/ws_handshake.hpp +++ b/include/morph/net/detail/ws_handshake.hpp @@ -3,13 +3,10 @@ #pragma once #include -#include #include -#include #include #include #include -#include #include #include #include @@ -270,36 +267,22 @@ struct HandshakeReadResult { std::string leftover; }; -/// Blocks until @p socket is readable or @p deadline passes, retrying across -/// `EINTR` by recomputing what remains rather than re-arming the full wait -- -/// mirrors `TcpSocket::connect()`'s own remaining-budget poll loop. Split out -/// of `readHttpHeaderBlock` so that function's own branching stays readable. +/// Blocks until @p socket is readable or @p deadline passes, on the same +/// single-budget wait `TcpSocket::connect()` uses (`pollUntil`). Split out of +/// `readHttpHeaderBlock` so that function's own branching stays readable. +/// @param socket Socket to wait on. +/// @param deadline Point after which the handshake counts as timed out. /// @throws std::runtime_error once @p deadline passes, or on a `poll()` error /// other than `EINTR`. inline void waitReadableUntil(const TcpSocket& socket, std::chrono::steady_clock::time_point deadline) { - for (;;) { - auto const remaining = - std::chrono::duration_cast(deadline - std::chrono::steady_clock::now()); - if (remaining.count() <= 0) { - throw std::runtime_error("readHttpHeaderBlock: handshake timed out"); - } - auto const waitMs = static_cast( - std::min(remaining.count(), std::numeric_limits::max())); - pollfd pfd{}; - pfd.fd = socket.nativeHandle(); - pfd.events = POLLIN; - int const pollRc = ::poll(&pfd, 1, waitMs); - if (pollRc > 0) { - return; // readable (or hung up) -- the caller's recvSome will not block - } - if (pollRc == 0) { - throw std::runtime_error("readHttpHeaderBlock: handshake timed out"); - } - if (errno != EINTR) { - throw std::runtime_error("readHttpHeaderBlock: poll failed: " + std::system_category().message(errno)); - } - // EINTR: loop back and recompute the remaining budget. + auto const result = pollUntil(socket.nativeHandle(), POLLIN, deadline); + if (result.outcome == PollOutcome::kTimedOut) { + throw std::runtime_error("readHttpHeaderBlock: handshake timed out"); + } + if (result.outcome == PollOutcome::kFailed) { + throw std::runtime_error("readHttpHeaderBlock: poll failed: " + std::system_category().message(result.error)); } + // kReady: readable (or hung up) -- the caller's recvSome will not block. } /// @brief Reads bytes from @p socket until the `\r\n\r\n` header terminator. diff --git a/include/morph/qt/qt_websocket_backend.hpp b/include/morph/qt/qt_websocket_backend.hpp index 702f3217a..117915931 100644 --- a/include/morph/qt/qt_websocket_backend.hpp +++ b/include/morph/qt/qt_websocket_backend.hpp @@ -52,10 +52,11 @@ struct QtWebSocketBackendConfig { /// attempt aborts the page outright. `bindModel` then sends the request and /// returns an unsettled `Completion`, which the reply settles later. That is /// a deliberate trade, not a free improvement: the caller must wait for the - /// continuation (e.g. gate the UI on `BridgeHandler::whenBound()`) before - /// firing an action through that handler, since `executeVia` fails fast - /// with "handler not bound" for an unbound binding rather than queuing or - /// blocking. + /// continuation (gate the UI on `BridgeHandler::whenBound()`, or dispatch + /// through `BridgeHandler::executeWhenBound()`, which holds the action + /// until the bind settles) before firing an action through that handler, + /// since `execute()` fails fast with "handler not bound" for an unbound + /// binding rather than queuing or blocking. /// /// This flag chooses *whether the transport blocks*, nothing more: the /// continuation exists on both paths, because `bindModel` returns a diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index cbe7e0020..28c5c1c05 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -371,11 +371,19 @@ endif() get_target_property(MORPH_VETTED_HMAC_GUARD_INCLUDE_DIRS morph INTERFACE_INCLUDE_DIRECTORIES) get_target_property(MORPH_VETTED_HMAC_GUARD_GLAZE_INCLUDE_DIRS glaze::glaze INTERFACE_INCLUDE_DIRECTORIES) list(APPEND MORPH_VETTED_HMAC_GUARD_INCLUDE_DIRS ${MORPH_VETTED_HMAC_GUARD_GLAZE_INCLUDE_DIRS}) -# core-cpp's headers, spelled out rather than read from core::base: its -# include directories are $ generator expressions, which -# a try_compile() here would receive unevaluated. `src/` holds the headers and -# the binary directory's `include/` the generated . -list(APPEND MORPH_VETTED_HMAC_GUARD_INCLUDE_DIRS "${core-cpp_SOURCE_DIR}/src" "${core-cpp_BINARY_DIR}/include") +# core-cpp's headers. An in-tree core-cpp's are spelled out rather than read +# from core::base: its include directories are $ +# generator expressions, which a try_compile() here would receive unevaluated. +# `src/` holds the headers and the binary directory's `include/` the generated +# . An installed core-cpp's imported target carries plain +# paths, which can be read directly. +if(DEFINED core-cpp_SOURCE_DIR) + list(APPEND MORPH_VETTED_HMAC_GUARD_INCLUDE_DIRS "${core-cpp_SOURCE_DIR}/src" "${core-cpp_BINARY_DIR}/include") +else() + get_target_property(_morph_core_cpp_include_dirs core::base INTERFACE_INCLUDE_DIRECTORIES) + list(APPEND MORPH_VETTED_HMAC_GUARD_INCLUDE_DIRS ${_morph_core_cpp_include_dirs}) + unset(_morph_core_cpp_include_dirs) +endif() string(REPLACE ";" "\\;" MORPH_VETTED_HMAC_GUARD_INCLUDE_DIRS "${MORPH_VETTED_HMAC_GUARD_INCLUDE_DIRS}") unset(MORPH_VETTED_HMAC_GUARD_BLOCKS_DEFAULT CACHE) diff --git a/tests/net/test_socket_backend.cpp b/tests/net/test_socket_backend.cpp index 3579ff771..bcca4ef1d 100644 --- a/tests/net/test_socket_backend.cpp +++ b/tests/net/test_socket_backend.cpp @@ -7,7 +7,6 @@ #include #include #include -#include #include #include #include @@ -248,19 +247,14 @@ class FakeWsServer { if (auto accepted = _listener.tryAccept()) { return std::move(*accepted); } - auto const remaining = - std::chrono::duration_cast(deadline - std::chrono::steady_clock::now()); - if (remaining <= std::chrono::milliseconds::zero()) { + auto const waited = morph::net::detail::pollUntil(_listener.nativeHandle(), POLLIN, deadline); + if (waited.outcome == morph::net::detail::PollOutcome::kTimedOut) { throw std::runtime_error("FakeWsServer::acceptAndHandshake: no client connected within " + std::to_string(timeout.count()) + "ms"); } - pollfd pfd{}; - pfd.fd = _listener.nativeHandle(); - pfd.events = POLLIN; - int const ready = ::poll(&pfd, 1, static_cast(remaining.count())); - if (ready < 0 && errno != EINTR) { + if (waited.outcome == morph::net::detail::PollOutcome::kFailed) { throw std::runtime_error("FakeWsServer::acceptAndHandshake: poll() failed: " + - std::system_category().message(errno)); + std::system_category().message(waited.error)); } } } diff --git a/tests/net/test_tcp_socket.cpp b/tests/net/test_tcp_socket.cpp index 5b72059b7..64db3678e 100644 --- a/tests/net/test_tcp_socket.cpp +++ b/tests/net/test_tcp_socket.cpp @@ -2,6 +2,7 @@ #include #include +#include #include #include #include @@ -12,6 +13,7 @@ #include #include #include +#include #include #include #include @@ -677,3 +679,76 @@ TEST_CASE("TcpSocket: setSendTimeout bounds a send against a peer that never rea // buffers filled costs real time too. CHECK(elapsed < std::chrono::seconds{20}); } + +// ── pollUntil: one budget for the whole wait ──────────────────────────────── + +TEST_CASE("pollUntil: a deadline already passed reports kTimedOut without waiting", "[net][tcp][poll_until]") { + auto listener = TcpSocket::listen(0); + auto const started = std::chrono::steady_clock::now(); + auto const result = morph::net::detail::pollUntil(listener.nativeHandle(), POLLIN, started); + CHECK(result.outcome == morph::net::detail::PollOutcome::kTimedOut); + CHECK(result.error == 0); + CHECK(std::chrono::steady_clock::now() - started < std::chrono::milliseconds{100}); +} + +TEST_CASE("pollUntil: an idle listener times out at the deadline", "[net][tcp][poll_until]") { + auto listener = TcpSocket::listen(0); + auto const started = std::chrono::steady_clock::now(); + auto const result = + morph::net::detail::pollUntil(listener.nativeHandle(), POLLIN, started + std::chrono::milliseconds{100}); + auto const elapsed = std::chrono::steady_clock::now() - started; + CHECK(result.outcome == morph::net::detail::PollOutcome::kTimedOut); + CHECK(elapsed >= std::chrono::milliseconds{100}); + CHECK(elapsed < std::chrono::seconds{5}); +} + +TEST_CASE("pollUntil: a pending connection reports kReady", "[net][tcp][poll_until]") { + auto listener = TcpSocket::listen(0); + auto clientSide = TcpSocket::connect("127.0.0.1", listener.boundPort(), std::chrono::milliseconds{2000}); + REQUIRE(clientSide.valid()); + auto const result = morph::net::detail::pollUntil(listener.nativeHandle(), POLLIN, + std::chrono::steady_clock::now() + std::chrono::seconds{5}); + CHECK(result.outcome == morph::net::detail::PollOutcome::kReady); +} + +namespace { +void pollUntilTestNoopHandler(int /*signal*/) {} +} // namespace + +TEST_CASE("pollUntil: signals mid-wait neither end the wait nor re-arm the full budget", "[net][tcp][poll_until]") { + // SIGUSR1 aimed at this thread every 50 ms, with no SA_RESTART, interrupts + // poll() with EINTR over and over for 1.5 s. The wait must still end at its + // own 300 ms deadline: a loop that treated EINTR as failure would report + // kFailed, and one that re-armed the full budget after each signal would + // not finish until the signals stopped. + struct sigaction action{}; + action.sa_handler = pollUntilTestNoopHandler; + sigemptyset(&action.sa_mask); + action.sa_flags = 0; + struct sigaction previous{}; + REQUIRE(::sigaction(SIGUSR1, &action, &previous) == 0); + + auto listener = TcpSocket::listen(0); + pthread_t const waiter = ::pthread_self(); + std::atomic sent{0}; + std::thread interrupter{[&] { + for (int i = 0; i < 30; ++i) { + std::this_thread::sleep_for(std::chrono::milliseconds{50}); + ::pthread_kill(waiter, SIGUSR1); + sent.fetch_add(1); + } + }}; + + auto const started = std::chrono::steady_clock::now(); + auto const result = + morph::net::detail::pollUntil(listener.nativeHandle(), POLLIN, started + std::chrono::milliseconds{300}); + auto const elapsed = std::chrono::steady_clock::now() - started; + int const sentDuringWait = sent.load(); + interrupter.join(); + ::sigaction(SIGUSR1, &previous, nullptr); + + REQUIRE(sentDuringWait > 0); + CHECK(result.outcome == morph::net::detail::PollOutcome::kTimedOut); + CHECK(elapsed >= std::chrono::milliseconds{300}); + CHECK(elapsed < std::chrono::milliseconds{1200}); +} diff --git a/tests/test_async_registration.cpp b/tests/test_async_registration.cpp index 8a5c7a09d..eeca5f55a 100644 --- a/tests/test_async_registration.cpp +++ b/tests/test_async_registration.cpp @@ -227,6 +227,7 @@ class AsyncRegisterBackend : public morph::backend::detail::IBackend { morph::exec::IExecutor* cbExec) override { auto state = std::make_shared>>(); morph::async::Completion> comp{state, cbExec}; + _executeCount.fetch_add(1); std::scoped_lock const lock{_regMtx}; auto iter = _models.find(mid.v); if (iter == _models.end()) { @@ -294,6 +295,8 @@ class AsyncRegisterBackend : public morph::backend::detail::IBackend { std::scoped_lock const lock{_pendingMtx}; return _pending.size(); } + /// The number of actions that reached execute(), whatever their outcome. + [[nodiscard]] int executeCount() const { return _executeCount.load(); } protected: /// @brief Parks one bind request until completeNext()/failNext() settles it. @@ -326,6 +329,7 @@ class AsyncRegisterBackend : public morph::backend::detail::IBackend { std::unordered_map> _models; std::vector> _assigned; uint64_t _nextId{100}; + std::atomic _executeCount{0}; }; // Settles its bind *inline* -- synchronously, from inside bindModel itself, @@ -1174,6 +1178,134 @@ TEST_CASE("Bridge::whenBound: multiple waiters on the same in-flight registratio CHECK(resolvedCount == 3); } +// ── executeWhenBound: dispatch held until the registration settles ───────── + +TEST_CASE("BridgeHandler::executeWhenBound: an action issued before the bind dispatches once it lands", + "[bridge][registration][executeWhenBound]") { + SyncExec cbExec; + auto backend = std::make_unique(); + auto* rawBackend = backend.get(); + morph::bridge::Bridge bridge{std::move(backend)}; + morph::bridge::BridgeHandler handler{bridge, &cbExec}; + REQUIRE_FALSE(handler.isBound()); + + std::atomic result{-1}; + std::atomic errored{false}; + handler.executeWhenBound(ARCount{.x = 5}) + .then([&](int v) { result.store(v); }) + .onError([&](const std::exception_ptr&) { errored.store(true); }); + + // Held, not failed and not dispatched: nothing is bound to dispatch to. + CHECK(result.load() == -1); + CHECK_FALSE(errored.load()); + CHECK(rawBackend->executeCount() == 0); + + rawBackend->completeNext(); + REQUIRE(morph::testing::waitUntil([&] { return result.load() != -1 || errored.load(); })); + CHECK(result.load() == 5); + CHECK_FALSE(errored.load()); + CHECK(rawBackend->executeCount() == 1); +} + +TEST_CASE("BridgeHandler::executeWhenBound: a failed registration rejects with the registration's error", + "[bridge][registration][executeWhenBound]") { + SyncExec cbExec; + auto backend = std::make_unique(); + auto* rawBackend = backend.get(); + morph::bridge::Bridge bridge{std::move(backend)}; + morph::bridge::BridgeHandler handler{bridge, &cbExec}; + + bool resolved = false; + std::string message; + handler.executeWhenBound(ARCount{.x = 5}) + .then([&](int) { resolved = true; }) + .onError([&](const std::exception_ptr& err) { + try { + std::rethrow_exception(err); + } catch (const std::exception& exc) { + message = exc.what(); + } + }); + CHECK(message.empty()); + + rawBackend->failNext("simulated registration failure"); + CHECK_FALSE(resolved); + CHECK(message.contains("simulated registration failure")); + CHECK(rawBackend->executeCount() == 0); +} + +TEST_CASE("BridgeHandler::executeWhenBound: a handler destroyed before the bind drops the held action", + "[bridge][registration][executeWhenBound]") { + SyncExec cbExec; + auto backend = std::make_unique(); + auto* rawBackend = backend.get(); + morph::bridge::Bridge bridge{std::move(backend)}; + + bool resolved = false; + bool errored = false; + { + morph::bridge::BridgeHandler handler{bridge, &cbExec}; + handler.executeWhenBound(ARCount{.x = 5}) + .then([&](int) { resolved = true; }) + .onError([&](const std::exception_ptr&) { errored = true; }); + } + + REQUIRE_NOTHROW(rawBackend->completeNext()); + CHECK_FALSE(resolved); + CHECK_FALSE(errored); + CHECK(rawBackend->executeCount() == 0); +} + +TEST_CASE("BridgeHandler::executeWhenBound: a handler destroyed after the bind but before the dispatch runs drops it", + "[bridge][registration][executeWhenBound]") { + // The bind lands and the dispatch is queued on the GUI executor; the handler + // goes before that queue is drained. The deferred dispatch must not reach + // the backend through a binding it kept alive on its own. + morph::testing::DeterministicExecutor guiExec; + auto backend = std::make_unique(); + auto* rawBackend = backend.get(); + morph::bridge::Bridge bridge{std::move(backend)}; + + bool resolved = false; + bool errored = false; + { + morph::bridge::BridgeHandler handler{bridge, &guiExec}; + handler.executeWhenBound(ARCount{.x = 5}) + .then([&](int) { resolved = true; }) + .onError([&](const std::exception_ptr&) { errored = true; }); + rawBackend->completeNext(); + REQUIRE(handler.isBound()); + REQUIRE(guiExec.pending() > 0); + } + + while (guiExec.pending() > 0) { + guiExec.step(); + } + CHECK_FALSE(resolved); + CHECK_FALSE(errored); + CHECK(rawBackend->executeCount() == 0); +} + +TEST_CASE("BridgeHandler::executeWhenBound: an already-bound handler dispatches immediately", + "[bridge][registration][executeWhenBound]") { + SyncExec cbExec; + auto backend = std::make_unique(); + auto* rawBackend = backend.get(); + morph::bridge::Bridge bridge{std::move(backend)}; + morph::bridge::BridgeHandler handler{bridge, &cbExec}; + rawBackend->completeNext(); + REQUIRE(handler.isBound()); + + std::atomic result{-1}; + handler.executeWhenBound(ARCount{.x = 9}) + .then([&](int v) { result.store(v); }) + .onError([](const std::exception_ptr&) {}); + REQUIRE(morph::testing::waitUntil([&] { return result.load() != -1; })); + CHECK(result.load() == 9); + CHECK(rawBackend->executeCount() == 1); + CHECK(rawBackend->pendingCount() == 0); +} + // ── assignHandlerPrimary goes through IBackend::promoteModel ────────────── // // A result-keyed action's execute() calls ensureBound() then, once the reply