diff --git a/.github/ai-opt-out b/.github/ai-opt-out new file mode 100644 index 00000000000..f2bf078d222 --- /dev/null +++ b/.github/ai-opt-out @@ -0,0 +1 @@ +opt-out: true diff --git a/.github/workflows/gen-matrix.py b/.github/workflows/gen-matrix.py index 8d549aef7b4..cca565165c1 100755 --- a/.github/workflows/gen-matrix.py +++ b/.github/workflows/gen-matrix.py @@ -30,7 +30,7 @@ Uses all-pairs (pairwise) coverage over compiler x standard x mode x arch for the regular builds, plus an explicit list of special-purpose -jobs (dpdk, cxx-modules, fuzz). All-pairs guarantees that every pair +jobs (dual TLS, OpenSSL TLS, fuzz). All-pairs guarantees that every pair of parameter values is exercised by at least one job, with far fewer combinations than the full cartesian product. Excluded pairs (see EXCLUDED_PAIRS) drop out of regular coverage. The file is round-tripped @@ -105,19 +105,17 @@ def _keep_row(row: list[Any]) -> bool: "compiler": "clang++-22", "standard": 23, "arch": "x86", - "mode": "release", - "enables": "--enable-dpdk", - "options": "--cook dpdk --dpdk-machine corei7-avx", - "info": "dpdk, ", + "mode": "debug", + "options": "--tls-mode=both", + "info": "dual TLS, ", }, { "compiler": "clang++-22", "standard": 23, "arch": "x86", "mode": "debug", - "enables": "--enable-cxx-modules", - "enable-ccache": False, - "info": "modules, ", + "options": "--tls-mode=openssl", + "info": "OpenSSL TLS, ", }, { "compiler": "clang++-22", diff --git a/.github/workflows/tests.yaml b/.github/workflows/tests.yaml index c5acef3c174..77f978f7dc7 100644 --- a/.github/workflows/tests.yaml +++ b/.github/workflows/tests.yaml @@ -32,8 +32,8 @@ jobs: - {compiler: clang++-21, standard: 26, mode: sanitize, arch: x86} - {compiler: clang++-21, standard: 26, mode: debug, arch: arm} - {compiler: clang++-22, standard: 23, arch: x86, mode: dev} - - {compiler: clang++-22, standard: 23, arch: x86, mode: release, enables: --enable-dpdk, options: --cook dpdk --dpdk-machine corei7-avx, info: 'dpdk, '} - - {compiler: clang++-22, standard: 23, arch: x86, mode: debug, enables: --enable-cxx-modules, enable-ccache: false, info: 'modules, '} + - {compiler: clang++-22, standard: 23, arch: x86, mode: debug, options: --tls-mode=both, info: 'dual TLS, '} + - {compiler: clang++-22, standard: 23, arch: x86, mode: debug, options: --tls-mode=openssl, info: 'OpenSSL TLS, '} - {compiler: clang++-22, standard: 23, arch: x86, mode: fuzz, test-args: -- -R 'Seastar.fuzz.'} with: compiler: ${{ matrix.compiler }} diff --git a/CMakeLists.txt b/CMakeLists.txt index 4b05745dbe8..4d82f05f4cd 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -421,6 +421,35 @@ set (Seastar_GEN_BINARY_DIR ${Seastar_BINARY_DIR}/gen) include (SeastarDependencies) seastar_find_dependencies () +# unordered_dense is header-only and not reliably packaged by distributions. +# When seastar_find_dependencies() did not locate an installed copy (e.g. in +# CI), fetch the pinned version so the build is self-contained. This lives in +# the Seastar build only, not in SeastarDependencies, so consumers that use +# find_package(Seastar) never trigger a download. +if (NOT TARGET unordered_dense::unordered_dense) + include (FetchContent) + FetchContent_Declare (unordered_dense + URL https://github.com/martinus/unordered_dense/archive/f30ed41b58af8c79788e8581fe57a6faf856258e.tar.gz + URL_HASH MD5=0370e4a35c1e573aa6639fde97f0c93f) + FetchContent_GetProperties (unordered_dense) + if (NOT unordered_dense_POPULATED) + # Populate (download) only; we deliberately do not add it as a subproject + # (see below). CMP0169 OLD keeps the bare Populate() form available. + if (POLICY CMP0169) + cmake_policy (SET CMP0169 OLD) + endif () + FetchContent_Populate (unordered_dense) + endif () + # Expose the fetched headers as an IMPORTED target rather than building it as + # a subproject, so that install(EXPORT) for Seastar (which PUBLIC-links it) + # does not require exporting a fetched target. Downstream consumers resolve + # unordered_dense themselves via SeastarDependencies. + add_library (unordered_dense_fetched INTERFACE IMPORTED GLOBAL) + set_target_properties (unordered_dense_fetched PROPERTIES + INTERFACE_INCLUDE_DIRECTORIES "${unordered_dense_SOURCE_DIR}/include") + add_library (unordered_dense::unordered_dense ALIAS unordered_dense_fetched) +endif () + # Private build dependencies not visible to consumers find_package (ragel 6.10 REQUIRED) find_package (Threads REQUIRED) @@ -529,11 +558,13 @@ seastar_generate_protobuf ( IN_FILE ${CMAKE_CURRENT_SOURCE_DIR}/src/proto/metrics2.proto OUT_DIR ${Seastar_GEN_BINARY_DIR}/src/proto) -set_option_if_package_is_found (Seastar_GNUTLS GnuTLS) -set_option_if_package_is_found (Seastar_OPENSSL OpenSSL) +option (Seastar_GNUTLS "Enable the GnuTLS-based TLS backend" ON) +option (Seastar_OPENSSL "Enable the OpenSSL-based TLS backend" OFF) if (NOT Seastar_GNUTLS AND NOT Seastar_OPENSSL) - message (FATAL_ERROR "At least one TLS/crypto backend is required. Install GnuTLS or OpenSSL development packages.") + message (FATAL_ERROR "At least one TLS backend must be enabled. " + "Pass -DSeastar_GNUTLS=ON and/or -DSeastar_OPENSSL=ON, " + "or use configure.py --tls-mode=gnutls|openssl|both.") endif () add_library (seastar @@ -552,6 +583,9 @@ add_library (seastar include/seastar/core/cacheline.hh include/seastar/core/checked_ptr.hh include/seastar/core/chunked_fifo.hh + include/seastar/core/chunked_hash_map.hh + include/seastar/core/chunked_vector.hh + include/seastar/core/chunked_vector_async.hh include/seastar/core/circular_buffer.hh include/seastar/core/circular_buffer_fixed_capacity.hh include/seastar/core/condition-variable.hh @@ -725,6 +759,7 @@ add_library (seastar src/core/reactor_backend.cc src/core/thread_pool.cc src/core/app-template.cc + src/core/cpu_profiler.cc src/core/disk_params.cc src/core/dpdk_rte.cc src/core/exception_hacks.cc @@ -758,6 +793,7 @@ add_library (seastar src/core/semaphore.cc src/core/condition-variable.cc src/core/crypto.cc + src/core/signal_mutex.cc src/http/api_docs.cc src/http/common.cc src/http/file_handler.cc @@ -920,6 +956,8 @@ target_link_libraries (seastar c-ares::cares fmt::fmt lz4::lz4 + unordered_dense::unordered_dense + absl::hash PRIVATE ${CMAKE_DL_LIBS} StdAtomic::atomic @@ -1143,6 +1181,15 @@ if (Seastar_OPENSSL) PRIVATE OpenSSL::SSL OpenSSL::Crypto) endif () +if (Seastar_GNUTLS AND Seastar_OPENSSL) + # Public marker: both TLS backends are compiled in, so the active backend is + # selected at reactor startup. Code that needs to handle the no-reactor case + # (e.g. static initializers, unit tests without a reactor) can use this to + # distinguish from the single-backend builds where the backend is fixed at + # compile time and available unconditionally. + target_compile_definitions (seastar PUBLIC SEASTAR_TLS_DUAL_BACKEND) +endif () + set_option_if_package_is_found (Seastar_IO_URING LibUring) if (Seastar_IO_URING) target_compile_definitions (seastar diff --git a/apps/memcached/memcache.cc b/apps/memcached/memcache.cc index f6f0e1cd430..3ce1bfde65e 100644 --- a/apps/memcached/memcache.cc +++ b/apps/memcached/memcache.cc @@ -893,7 +893,7 @@ class ascii_protocol { private: static void append(std::vector>& bufs, const char* buf, size_t size) { if (size) { - bufs.emplace_back(const_cast(buf), size, deleter()); + bufs.push_back(temporary_buffer::maybe_unsafe_from_deleter(const_cast(buf), size, deleter())); } } @@ -917,7 +917,7 @@ class ascii_protocol { append(bufs, msg_crlf); append(bufs, item->value()); - bufs.emplace_back(const_cast(msg_crlf), strlen(msg_crlf), make_deleter([item = std::move(item)]{})); + bufs.push_back(temporary_buffer::maybe_unsafe_from_deleter(const_cast(msg_crlf), strlen(msg_crlf), make_deleter([item = std::move(item)]{}))); } template diff --git a/cmake/SeastarDependencies.cmake b/cmake/SeastarDependencies.cmake index e62e1961df1..d13da4e240f 100644 --- a/cmake/SeastarDependencies.cmake +++ b/cmake/SeastarDependencies.cmake @@ -99,6 +99,11 @@ macro (seastar_find_dependencies) "to build a newer fmt locally.") endif () seastar_find_dep (lz4 1.7.3 REQUIRED) + # Not REQUIRED: unordered_dense is header-only and not reliably packaged by + # distributions. When it is not installed, Seastar's own build fetches it + # (see CMakeLists.txt); consumers that lack it must provide it themselves. + seastar_find_dep (unordered_dense) + seastar_find_dep (absl CONFIG REQUIRED) seastar_find_dep (GnuTLS 3.7.4) seastar_find_dep (OpenSSL 3.0) if (Seastar_IO_URING) diff --git a/configure.py b/configure.py index 9004c1fa17c..bd9d667b48b 100755 --- a/configure.py +++ b/configure.py @@ -173,6 +173,9 @@ def resolve_compilers_for_compiler_cache(args, compiler_cache): arg_parser.add_argument('--verbose', dest='verbose', action='store_true', help='Make configure output more verbose.') arg_parser.add_argument('--scheduling-groups-count', action='store', dest='scheduling_groups_count', default='16', help='Number of available scheduling groups in the reactor') +arg_parser.add_argument('--tls-mode', action='store', dest='tls_mode', + choices=['gnutls', 'openssl', 'both'], default='gnutls', + help='TLS backend(s) to enable: gnutls (default), openssl, or both') add_tristate( arg_parser, @@ -294,6 +297,8 @@ def configure_mode(mode): '-DBUILD_SHARED_LIBS={}'.format('yes' if mode in ('debug', 'dev') else 'no'), '-DSeastar_API_LEVEL={}'.format(args.api_level), '-DSeastar_SCHEDULING_GROUPS_COUNT={}'.format(args.scheduling_groups_count), + '-DSeastar_GNUTLS={}'.format('ON' if args.tls_mode in ('gnutls', 'both') else 'OFF'), + '-DSeastar_OPENSSL={}'.format('ON' if args.tls_mode in ('openssl', 'both') else 'OFF'), tr(args.exclude_tests, 'EXCLUDE_TESTS_FROM_ALL'), tr(args.exclude_apps, 'EXCLUDE_APPS_FROM_ALL'), tr(args.exclude_demos, 'EXCLUDE_DEMOS_FROM_ALL'), diff --git a/cooking_recipe.cmake b/cooking_recipe.cmake index 8253835e33b..e07ff24645c 100644 --- a/cooking_recipe.cmake +++ b/cooking_recipe.cmake @@ -336,3 +336,21 @@ cooking_ingredient (lz4 CONFIGURE_COMMAND BUILD_COMMAND INSTALL_COMMAND ${make_command} PREFIX= install) + +# Header-only library used by chunked_hash_map. Pinned to the same commit +# (v4.4.0) consumed by Redpanda. +cooking_ingredient (unordered_dense + EXTERNAL_PROJECT_ARGS + URL https://github.com/martinus/unordered_dense/archive/f30ed41b58af8c79788e8581fe57a6faf856258e.tar.gz + URL_MD5 0370e4a35c1e573aa6639fde97f0c93f) + +# Used by chunked_hash_map for absl::Hash support. Pinned to the same LTS +# release (20250814.1) consumed by Redpanda. +cooking_ingredient (absl + EXTERNAL_PROJECT_ARGS + URL https://github.com/abseil/abseil-cpp/releases/download/20250814.1/abseil-cpp-20250814.1.tar.gz + URL_MD5 d4d3c25f78e28d61ad83e54cd1116933 + CMAKE_ARGS + -DABSL_PROPAGATE_CXX_STD=ON + -DABSL_ENABLE_INSTALL=ON + -DBUILD_TESTING=OFF) diff --git a/demos/line_count_demo.cc b/demos/line_count_demo.cc index eab4609ecd4..c1414f28368 100644 --- a/demos/line_count_demo.cc +++ b/demos/line_count_demo.cc @@ -42,14 +42,14 @@ struct reader { size_t count = 0; // for input_stream::consume(): - using unconsumed_remainder = std::optional>; - future operator()(temporary_buffer data) { + using consumption_result_type = consumption_result; + future operator()(temporary_buffer data) { if (data.empty()) { - return make_ready_future(std::move(data)); + return make_ready_future(stop_consuming{std::move(data)}); } else { count += std::count(data.begin(), data.end(), '\n'); // FIXME: last line without \n? - return make_ready_future(); + return make_ready_future(continue_consuming{}); } } }; diff --git a/demos/udp_zero_copy_demo.cc b/demos/udp_zero_copy_demo.cc index 79dae99444f..c59b238138a 100644 --- a/demos/udp_zero_copy_demo.cc +++ b/demos/udp_zero_copy_demo.cc @@ -107,7 +107,7 @@ class server { if (_copy) { bufs.emplace_back(temporary_buffer(chunk, _chunk_size)); } else { - bufs.emplace_back(temporary_buffer(chunk, _chunk_size, deleter())); + bufs.emplace_back(temporary_buffer::maybe_unsafe_from_deleter(chunk, _chunk_size, deleter())); } chunk += _chunk_size; } diff --git a/include/seastar/core/chunked_hash_map.hh b/include/seastar/core/chunked_hash_map.hh new file mode 100644 index 00000000000..4e0f54ff583 --- /dev/null +++ b/include/seastar/core/chunked_hash_map.hh @@ -0,0 +1,279 @@ +/* + * This file is open source software, licensed to you under the terms + * of the Apache License, Version 2.0 (the "License"). See the NOTICE file + * distributed with this work for additional information regarding copyright + * ownership. You may not use this file except in compliance with the License. + * + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +/* + * Copyright 2024 Redpanda Data, Inc. + */ +#pragma once + +#include + +#include +#include +#include +#include + +#include +#include +#include +#include +#include +#include + +namespace seastar { + +namespace internal { + +template +concept has_absl_hash = requires(T val) { + { AbslHashValue(std::declval(), val) }; +}; + +/// Wrapper around absl::Hash that disables the extra hash mixing in +/// unordered_dense +template +struct avalanching_absl_hash { + // absl always hash mixes itself so no need to do it again + using is_avalanching = void; + + auto operator()(const T& x) const noexcept -> uint64_t { + return absl::Hash()(x); + } +}; + +} // namespace internal + +/** + * @brief A hash map that uses a chunked vector as the underlying storage. + * + * Use when the hash map is expected to have a large number of elements (e.g.: + * scales with partitions or topics). Performance wise it's equal to the abseil + * hashmaps. + * + * NB: References and iterators are not stable across insertions and deletions. + * + * Both std::hash and abseil's AbslHashValue are supported. We dispatch to the + * latter if available. Given AbslHashValue also supports std::hash we could + * also unconditionally dispatch to it. However, absl's hash mixing seems more + * extensive (and hence less performant) so we only do that when needed. + * + * For more info please see + * https://github.com/martinus/unordered_dense/?tab=readme-ov-file#1-overview + */ +template< + typename Key, + typename Value, + typename Hash = std::conditional_t< + internal::has_absl_hash, + internal::avalanching_absl_hash, + ankerl::unordered_dense::hash>, + typename EqualTo = std::equal_to> +using chunked_hash_map = ankerl::unordered_dense::segmented_map< + Key, + Value, + Hash, + EqualTo, + chunked_vector>, + ankerl::unordered_dense::bucket_type::standard, + chunked_vector>; + +/** + * @brief A set counterpart of chunked_hash_map (uses a chunked vector as the + * underlying storage). + */ +template< + typename Key, + typename Hash = std::conditional_t< + internal::has_absl_hash, + internal::avalanching_absl_hash, + ankerl::unordered_dense::hash>, + typename EqualTo = std::equal_to> +using chunked_hash_set = ankerl::unordered_dense::segmented_set< + Key, + Hash, + EqualTo, + chunked_vector, + ankerl::unordered_dense::bucket_type::standard, + chunked_vector>; + +namespace internal { +template +struct chunked_hash_map_from_range_impl { + using value_t = std::ranges::range_value_t>; + using first_t = typename value_t::first_type; + using second_t = typename value_t::second_type; + using ret_t = chunked_hash_map, second_t>; +}; +} // namespace internal + +// reserves if range size is known +template +requires std::ranges::input_range + && std::convertible_to, typename TargetTable::value_type> +TargetTable chunked_table_from_range(Range&& range) { + size_t size = 0; + if constexpr (std::ranges::sized_range) { + size = std::ranges::size(range); + } + return {std::ranges::begin(range), std::ranges::end(range), size}; +} + +// reserves if range size is known +template +auto chunked_hash_map_from_range(Range&& range) { + return chunked_table_from_range< + typename internal::chunked_hash_map_from_range_impl::ret_t>( + std::forward(range)); +} + +// reserves if range size is known +template +auto chunked_hash_set_from_range(Range&& range) { + return chunked_table_from_range< + chunked_hash_set>>>( + std::forward(range)); +} + +/// Returns a lower bound on the memory currently being held by `m`. +template< + typename K, + typename V, + typename Hash = std::conditional_t< + internal::has_absl_hash, + internal::avalanching_absl_hash, + ankerl::unordered_dense::hash>, + typename EqualTo = std::equal_to> +size_t +memory_usage_lower_bound(const chunked_hash_map& m) { + return m.bucket_count() + * sizeof(typename chunked_hash_map::bucket_type) + + m.values().capacity() * sizeof(m.values()[0]); +} + +} // namespace seastar + +// Disable fmt's range formatter for chunked_hash_map/set to use our custom +// formatters instead. +template< + typename K, + typename V, + typename H, + typename E, + typename AllocatorOrContainer, + typename Bucket, + typename BucketContainer, + bool IsSegmented, + typename Char> +struct fmt::range_format_kind< + ankerl::unordered_dense::detail::table< + K, + V, + H, + E, + AllocatorOrContainer, + Bucket, + BucketContainer, + IsSegmented>, + Char> + : std::integral_constant {}; + +template< + typename K, + typename V, + typename H, + typename E, + typename AllocatorOrContainer, + typename Bucket, + typename BucketContainer, + bool IsSegmented> +struct fmt::formatter> { + using type = ankerl::unordered_dense::detail::table< + K, + V, + H, + E, + AllocatorOrContainer, + Bucket, + BucketContainer, + IsSegmented>; + + constexpr auto parse(format_parse_context& ctx) const { + return ctx.begin(); + } + + template + typename FormatContext::iterator + format(const type& map, FormatContext& ctx) const { + auto out = ctx.out(); + out = fmt::format_to(out, "["); + auto it = map.begin(); + if (it != map.end()) { + out = fmt::format_to(out, "{{{} -> {}}}", it->first, it->second); + for (++it; it != map.end(); ++it) { + out = fmt::format_to( + out, ", {{{} -> {}}}", it->first, it->second); + } + } + return fmt::format_to(out, "]"); + } +}; + +template< + typename K, + typename H, + typename E, + typename AllocatorOrContainer, + typename Bucket, + typename BucketContainer, + bool IsSegmented> +struct fmt::formatter> { + using type = ankerl::unordered_dense::detail::table< + K, + void, + H, + E, + AllocatorOrContainer, + Bucket, + BucketContainer, + IsSegmented>; + + constexpr auto parse(format_parse_context& ctx) const { + return ctx.begin(); + } + + template + typename FormatContext::iterator + format(const type& set, FormatContext& ctx) const { + return fmt::format_to(ctx.out(), "[{}]", fmt::join(set, ",")); + } +}; diff --git a/include/seastar/core/chunked_vector.hh b/include/seastar/core/chunked_vector.hh new file mode 100644 index 00000000000..2f2193c188e --- /dev/null +++ b/include/seastar/core/chunked_vector.hh @@ -0,0 +1,624 @@ +/* + * This file is open source software, licensed to you under the terms + * of the Apache License, Version 2.0 (the "License"). See the NOTICE file + * distributed with this work for additional information regarding copyright + * ownership. You may not use this file except in compliance with the License. + * + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +/* + * Copyright 2020 Redpanda Data, Inc. + */ +#pragma once + +#include + +#include +#include + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +namespace seastar { + +template +class future; + +template +class chunked_vector; + +// Defined in ; declared here so they can +// be befriended for direct fragment access. +template +future chunked_vector_fill_async(chunked_vector& vec, const T& value); +template +future chunked_vector_clear_async(chunked_vector& vec); + +/** + * A chunked vector is a container that provides random access like a + * vector, but does not store its data in contiguous memory. + * + * Instead the allocations are broken up across many different individual + * vectors, but the exposed view is of a single container. + * + * Additionally the allocation strategy is like a "normal" vector until the + * first chunk is full, after which we will then only allocate full chunks. + * + * The iterator implementation works for a few things like std::lower_bound, + * upper_bound, distance, etc... see chunked_vector_test. + */ +template +class chunked_vector { + static constexpr size_t max_allocation_size = 128UL * 1024; + + // calculate the maximum number of elements per fragment while + // keeping the element count a power of two + static consteval size_t calc_elems_per_frag() { + size_t max = max_allocation_size / sizeof(T); + return std::bit_floor(max); + } + + /** + * The maximum number of bytes per fragment as specified in + * as part of the type. Note that for most types, the true + * number of bytes in a full fragment may as low as half + * of this amount (+1) since the number of elements is restricted + * to a power of two. + */ + static consteval size_t calc_max_frag_bytes() { + return calc_elems_per_frag() * sizeof(T); + } + +public: + using this_type = chunked_vector; + using backing_type = std::vector>; + using value_type = T; + using reference = std::conditional_t, bool, T&>; + using const_reference + = std::conditional_t, bool, const T&>; + using size_type = size_t; + using allocator_type = backing_type::allocator_type; + using difference_type = backing_type::difference_type; + using pointer = T*; + using const_pointer = const T*; + + chunked_vector() noexcept = default; + explicit chunked_vector(allocator_type alloc) + : _frags(alloc) {} + chunked_vector& operator=(const chunked_vector&) noexcept = delete; + chunked_vector(chunked_vector&& other) noexcept { + *this = std::move(other); + } + + /** + * @brief Create a vector from a begin, end iterator pair. + * + * This has the same semantics as the corresponding std::vector + * constructor. + */ + template + requires std::input_iterator + chunked_vector(Iter begin, Iter end) + : chunked_vector() { + if constexpr (std::random_access_iterator) { + reserve(std::distance(begin, end)); + } + // Improvement: Write a more efficient implementation for + // std::contiguous_iterator + for (auto it = begin; it != end; ++it) { + push_back(*it); + } + } + + /** + * @brief Construct a new vector using an initializer list + * + * In the same manner as the corresponding std::vector method. + */ + chunked_vector(std::initializer_list elems) + : chunked_vector(elems.begin(), elems.end()) {} + +#ifdef __cpp_lib_containers_ranges + /** + * @brief Construct a new vector from a range + * + * This constructor will copy or move from the range depending on the value + * category of the elements NOT the one of the range. I.e. + * `chunked_vector(std::move(src))` will not necessarily invoke move + * constructor on the elements of `src`. For example, an rvalue `std::span` + * is a non-owning view, and moving from its elements would be a bug. + * Similar for most of the standard library views. + * + * To ensure move semantics from the range elements, use + * `std::views::as_rvalue`. + * + * https://en.cppreference.com/w/cpp/ranges/as_rvalue_view.html + */ + template + requires(std::ranges::range) + // NOLINTNEXTLINE(cppcoreguidelines-missing-std-forward) + chunked_vector(std::from_range_t, Range &&range) : chunked_vector() { + if constexpr (std::ranges::sized_range) { + reserve(std::ranges::size(range)); + } + std::copy(std::ranges::begin(range), std::ranges::end(range), + std::back_inserter(*this)); + } +#endif + + chunked_vector& operator=(chunked_vector&& other) noexcept { + if (this != &other) { + this->_size = other._size; + this->_capacity = other._capacity; + this->_frags = std::move(other._frags); + // Move compatibility with std::vector that post move + // the vector is empty(). + other._size = other._capacity = 0; + other.update_generation(); + update_generation(); + } + return *this; + } + ~chunked_vector() noexcept = default; + + chunked_vector copy() const noexcept { return *this; } + + auto get_allocator() const { return _frags.get_allocator(); } + + void swap(chunked_vector& other) noexcept { + std::swap(_size, other._size); + std::swap(_capacity, other._capacity); + std::swap(_frags, other._frags); + other.update_generation(); + update_generation(); + } + + template + void push_back(E&& elem) { + maybe_add_capacity(); + _frags.back().push_back(std::forward(elem)); + ++_size; + update_generation(); + } + + template + T& emplace_back(Args&&... args) { + maybe_add_capacity(); + T& emplaced = _frags.back().emplace_back(std::forward(args)...); + ++_size; + update_generation(); + return emplaced; + } + + void pop_back() { + SEASTAR_ASSERT(_size > 0 && "Cannot pop from empty container"); + _frags.back().pop_back(); + --_size; + if (_frags.back().empty()) { + _frags.pop_back(); + _capacity -= std::min(calc_elems_per_frag(), _capacity); + } + update_generation(); + } + + /* + * Replacement for `erase(some_it, end())` but more efficient than n + * `pop_back()s` + */ + void pop_back_n(size_t n) { + SEASTAR_ASSERT( + _size >= n && "Cannot pop more than size() elements in container"); + + if (_size == n) { + clear(); + return; + } + + _size -= n; + + while (n >= _frags.back().size()) { + n -= _frags.back().size(); + _frags.pop_back(); + _capacity -= calc_elems_per_frag(); + } + + for (size_t i = 0; i < n; ++i) { + _frags.back().pop_back(); + } + update_generation(); + } + + const_reference at(size_t index) const { + static constexpr size_t elems_per_frag = calc_elems_per_frag(); + return _frags.at(index / elems_per_frag).at(index % elems_per_frag); + } + + reference at(size_t index) { + static constexpr size_t elems_per_frag = calc_elems_per_frag(); + return _frags.at(index / elems_per_frag).at(index % elems_per_frag); + } + + const_reference operator[](size_t index) const { + static constexpr size_t elems_per_frag = calc_elems_per_frag(); + return _frags[index / elems_per_frag][index % elems_per_frag]; + } + + reference operator[](size_t index) { + static constexpr size_t elems_per_frag = calc_elems_per_frag(); + return _frags[index / elems_per_frag][index % elems_per_frag]; + } + + const_reference front() const { return _frags.front().front(); } + const_reference back() const { return _frags.back().back(); } + reference front() { return _frags.front().front(); } + reference back() { return _frags.back().back(); } + bool empty() const noexcept { return _size == 0; } + size_t size() const noexcept { return _size; } + size_t capacity() const noexcept { return _capacity; } + + /** + * Requests the removal of unused capacity. + * + * If reallocation occurs, all iterators (including the end() iterator) and + * all references to the elements are invalidated. + */ + void shrink_to_fit() { + // Calling shrink to fix then modifying the container could result + // in allocations that overshoot our max_frag_bytes, except when + // we're managing the dynamic size of the first fragment. + if (_frags.size() == 1) { + auto& front = _frags.front(); + front.shrink_to_fit(); + _capacity = front.capacity(); + } + update_generation(); + } + + /** + * Increase the capacity of the vector (the total number of elements that + * the vector can hold without requiring reallocation) to a value that's + * greater or equal to new_cap. + * + * This method has the same guarantees as `std::vector::reserve`, that being + * calling reserve doesn't preserve pointer or iterator stability when + * new_cap is larger than the current capacity. There maybe cases where this + * class invalidates less than `std::vector`, but that should not be relied + * upon. + * + * Performance wise if you know the intended size of this structure is + * desired to call this to prevent costly reallocations when inserting + * elements. + */ + void reserve(size_t new_cap) { + static constexpr size_t elems_per_frag = calc_elems_per_frag(); + if (new_cap > _capacity) { + if (_frags.empty()) { + auto& frag = _frags.emplace_back(); + frag.reserve(std::min(elems_per_frag, new_cap)); + _capacity = frag.capacity(); + } else if (_frags.size() == 1) { + auto& frag = _frags.front(); + frag.reserve(std::min(elems_per_frag, new_cap)); + _capacity = frag.capacity(); + } + // We only reserve the first fragment as all fragments after the + // first are allocated at the maximum size, so we don't save + // anything in terms of reallocs after fully allocating the + // first fragment. In addition, due to cache locality, it's + // better to delay the allocations of those other fragments + // until they're going to be used. + } + update_generation(); + } + + bool operator==(const chunked_vector& o) const noexcept { + return o._frags == _frags; + } + + /** + * Returns the approximate in-memory size of this vector in bytes. + */ + size_t memory_size() const { + return (_frags.size() * sizeof(_frags[0])) + (_capacity * sizeof(T)); + } + + /** + * Returns the (maximum) number of elements in each fragment of this vector. + */ + static constexpr size_t elements_per_fragment() { + return calc_elems_per_frag(); + } + + static constexpr size_t max_frag_bytes() { return calc_max_frag_bytes(); } + + /** + * Remove all elements from the vector. + * + * Unlike std::vector, this also releases all the memory from + * the vector (since this vector already the same pointer + * and iterator stability guarantees that std::vector provides + * based on non-reallocation and capacity()). + */ + void clear() { + // do the swap dance to actually clear the memory held by the vector + std::vector>{}.swap(_frags); + _size = 0; + _capacity = 0; + update_generation(); + } + + template + class iter { + public: + using iterator_category = std::random_access_iterator_tag; + using iterator_concept = std::random_access_iterator_tag; + using value_type = T; + using difference_type = std::ptrdiff_t; + using pointer = std::conditional_t< + std::is_same_v, + bool, + std::conditional_t>; + using reference = std::conditional_t< + std::is_same_v, + bool, + std::conditional_t>; + + iter() = default; + + /** + * Conversion operator allowing iterator to be converted to + * const_iterator, as required by the general iterator contract. + */ + operator iter() const { // NOLINT(hicpp-explicit-conversions) + check_generation(); + iter ret; + ret._vec = _vec; + ret._index = _index; +#ifndef NDEBUG + ret._my_generation = _my_generation; +#endif + return ret; + } + + reference operator*() const { + check_generation(); + return _vec->operator[](_index); + } + pointer operator->() const { + check_generation(); + return &_vec->operator[](_index); + } + + iter& operator+=(ssize_t n) { + check_generation(); + _index += n; + return *this; + } + + iter& operator-=(ssize_t n) { + check_generation(); + _index -= n; + return *this; + } + + iter& operator++() { + check_generation(); + ++_index; + return *this; + } + + iter& operator--() { + check_generation(); + --_index; + return *this; + } + + iter operator++(int) { + check_generation(); + auto tmp = *this; + ++*this; + return tmp; + } + + iter operator--(int) { + check_generation(); + auto tmp = *this; + --*this; + return tmp; + } + + iter operator+(difference_type offset) const { + check_generation(); + return iter{*this} += offset; + } + + iter operator-(difference_type offset) const { + check_generation(); + return iter{*this} -= offset; + } + + reference operator[](difference_type offset) const { + check_generation(); + return *(*this + offset); + } + + bool operator==(const iter& o) const { + check_generation(); +#ifndef NDEBUG + SEASTAR_ASSERT( + _vec == o._vec + && "iterator compared with different chunked_vector"); +#endif + return std::tie(_index, _vec) == std::tie(o._index, o._vec); + }; + auto operator<=>(const iter& o) const { + check_generation(); +#ifndef NDEBUG + SEASTAR_ASSERT( + _vec == o._vec + && "iterator compared with different chunked_vector"); +#endif + return std::tie(_index, _vec) <=> std::tie(o._index, o._vec); + }; + + friend ssize_t operator-(const iter& a, const iter& b) { + return a._index - b._index; + } + + friend iter operator+(const difference_type& offset, const iter& i) { + return i + offset; + } + + private: + friend class chunked_vector; + using vec_type = std::conditional_t; + + iter(vec_type* vec, size_t index) + : _index(index) +#ifndef NDEBUG + , _my_generation(vec->_generation) +#endif + , _vec(vec) { + } + + inline void check_generation() const { +#ifndef NDEBUG + SEASTAR_ASSERT( + _vec->_generation == _my_generation + && "Attempting to use an invalidated iterator. The corresponding " + "chunked_vector container has been mutated since this " + "iterator was constructed."); +#endif + } + + size_t _index{}; +#ifndef NDEBUG + size_t _my_generation{}; +#endif + vec_type* _vec{}; + }; + + using const_iterator = iter; + using iterator = iter; + + using const_reverse_iterator = std::reverse_iterator; + using reverse_iterator = std::reverse_iterator; + + iterator begin() { return iterator(this, 0); } + iterator end() { return iterator(this, _size); } + + const_iterator begin() const { return const_iterator(this, 0); } + const_iterator end() const { return const_iterator(this, _size); } + + const_iterator cbegin() const { return const_iterator(this, 0); } + const_iterator cend() const { return const_iterator(this, _size); } + + reverse_iterator rbegin() { return reverse_iterator(end()); } + reverse_iterator rend() { return reverse_iterator(begin()); } + + const_reverse_iterator rbegin() const { + return const_reverse_iterator(end()); + } + const_reverse_iterator rend() const { + return const_reverse_iterator(begin()); + } + + const_reverse_iterator crbegin() const { + return const_reverse_iterator(end()); + } + const_reverse_iterator crend() const { + return const_reverse_iterator(begin()); + } + + /** + * @brief Erases all elements from begin to the end of the vector. + */ + void erase_to_end(const_iterator begin) { pop_back_n(cend() - begin); } + + template + static chunked_vector single(Args&&... args) { + chunked_vector v; + v.emplace_back(std::forward(args)...); + return v; + } + + friend std::ostream& + operator<<(std::ostream& os, const chunked_vector& v) { + fmt::print(os, "{}", v); + return os; + } + +private: + [[gnu::always_inline]] void maybe_add_capacity() { + if (_size == _capacity) [[unlikely]] { + add_capacity(); + } + } + + void add_capacity() { + static constexpr size_t elems_per_frag = calc_elems_per_frag(); + static_assert( + calc_max_frag_bytes() <= max_allocation_size, + "max size of a fragment must be <= 128KiB"); + if (_frags.size() == 1 && _frags.back().capacity() < elems_per_frag) { + auto& frag = _frags.back(); + auto new_cap = std::min(elems_per_frag, frag.capacity() * 2); + frag.reserve(new_cap); + _capacity = new_cap; + return; + } else if (_frags.empty()) { + // At least one element or 32 bytes worth of elements for small + // items. + static constexpr size_t initial_cap = std::max( + 1UL, 32UL / sizeof(T)); + _capacity = initial_cap; + _frags.emplace_back(_frags.get_allocator()).reserve(_capacity); + return; + } + _frags.emplace_back(_frags.get_allocator()).reserve(elems_per_frag); + _capacity += elems_per_frag; + } + + inline void update_generation() { +#ifndef NDEBUG + ++_generation; +#endif + } + +private: + friend class chunked_vector_validator; + template + friend future + chunked_vector_fill_async(chunked_vector& vec, const U& value); + template + friend future chunked_vector_clear_async(chunked_vector& vec); + chunked_vector(const chunked_vector&) noexcept = default; + + size_t _size{0}; + size_t _capacity{0}; + backing_type _frags; +#ifndef NDEBUG + // Add a generation number that is incremented on every mutation to catch + // invalidated iterator accesses. + size_t _generation{0}; +#endif +}; + +} // namespace seastar diff --git a/include/seastar/core/chunked_vector_async.hh b/include/seastar/core/chunked_vector_async.hh new file mode 100644 index 00000000000..7e85e68b83e --- /dev/null +++ b/include/seastar/core/chunked_vector_async.hh @@ -0,0 +1,78 @@ +/* + * This file is open source software, licensed to you under the terms + * of the Apache License, Version 2.0 (the "License"). See the NOTICE file + * distributed with this work for additional information regarding copyright + * ownership. You may not use this file except in compliance with the License. + * + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +/* + * Copyright 2025 Redpanda Data, Inc. + */ +#pragma once + +#include +#include +#include +#include +#include +#include + +#include + +// Async helpers for chunked_vector. +// Not in the main header to avoid pulling in all of seastar if you just need +// chunked_vector. + +namespace seastar { + +/** + * A futurized version of std::fill optimized for fragmented vector. It is + * futurized to allow for large vectors to be filled without incurring reactor + * stalls. It is optimized by circumventing the indexing indirection incurred by + * using the fragmented vector interface directly. + */ +template +future<> chunked_vector_fill_async(chunked_vector& vec, const T& value) { + auto remaining = vec._size; + for (auto& frag : vec._frags) { + const auto n = std::min(frag.size(), remaining); + if (n == 0) { + break; + } + std::fill_n(frag.begin(), n, value); + remaining -= n; + if (need_preempt()) { + co_await yield(); + } + } + SEASTAR_ASSERT( + remaining == 0 + && "chunked_vector_fill_async: fragmented vector inconsistency"); +} + +/** + * A futurized version of chunked_vector::clear that allows clearing a large + * vector without incurring a reactor stall. + */ +template +future<> chunked_vector_clear_async(chunked_vector& vec) { + while (!vec._frags.empty()) { + vec._frags.pop_back(); + if (need_preempt()) { + co_await yield(); + } + } + vec.clear(); +} + +} // namespace seastar diff --git a/include/seastar/core/deleter.hh b/include/seastar/core/deleter.hh index 371459e242f..9db9e0c61ee 100644 --- a/include/seastar/core/deleter.hh +++ b/include/seastar/core/deleter.hh @@ -21,13 +21,26 @@ #pragma once +#include #include #include #include #include #include +// The forward declarations of classes below are used for +// friending by the deleter. +struct test_deleter_append_does_not_free_shared_object; +struct test_deleter_append_same_shared_object_twice; + namespace seastar { +namespace net { + class packet; +}; +namespace internal { + struct wrapped_iovecs; +} +class pipe_data_sink_impl; /// \addtogroup memory-module /// @{ @@ -82,11 +95,21 @@ public: this->~deleter(); new (this) deleter(i); } +private: /// \endcond /// Appends another deleter to this deleter. When this deleter is /// destroyed, both encapsulated actions will be carried out. + /// + /// This operation is not thread-safe and therefore not made public + /// except for a few manually verified uses that are marked as freinds + /// below. void append(deleter d); -private: + friend class ::seastar::net::packet; + friend struct ::test_deleter_append_does_not_free_shared_object; + friend struct ::test_deleter_append_same_shared_object_twice; + friend struct ::seastar::internal::wrapped_iovecs; + friend class ::seastar::pipe_data_sink_impl; + static bool is_raw_object(impl* i) noexcept { auto x = reinterpret_cast(i); return x & 1; @@ -109,7 +132,9 @@ private: /// \cond internal struct deleter::impl { - unsigned refs = 1; + // The memory ordering on operations to this counter is similar to + // std::shared_ptr. + std::atomic refs = 1; deleter next; impl(deleter next) : next(std::move(next)) {} virtual ~impl() {} @@ -122,7 +147,7 @@ deleter::~deleter() { std::free(to_raw_object()); return; } - if (_impl && --_impl->refs == 0) { + if (_impl && _impl->refs.fetch_sub(1, std::memory_order_acq_rel) == 1) { delete _impl; } } @@ -203,7 +228,7 @@ deleter::share() { if (is_raw_object()) { _impl = new free_deleter_impl(to_raw_object()); } - ++_impl->refs; + _impl->refs.fetch_add(1, std::memory_order_relaxed); return deleter(_impl); } diff --git a/include/seastar/core/disk_params.hh b/include/seastar/core/disk_params.hh index 27542e34eaf..9503d95e507 100644 --- a/include/seastar/core/disk_params.hh +++ b/include/seastar/core/disk_params.hh @@ -45,6 +45,7 @@ struct disk_params { std::optional physical_block_size; // Override for disks that lie about their physical block size bool duplex = false; float rate_factor = 1.0; + bool max_cost_function = true; }; class disk_config_params { diff --git a/include/seastar/core/file.hh b/include/seastar/core/file.hh index 1d149079163..9b03aca3afa 100644 --- a/include/seastar/core/file.hh +++ b/include/seastar/core/file.hh @@ -651,7 +651,7 @@ public: future> dma_read_bulk(uint64_t offset, size_t range_size, io_intent* intent = nullptr) noexcept { return dma_read_bulk_impl(offset, range_size, intent).then([] (temporary_buffer t) { - return temporary_buffer(reinterpret_cast(t.get_write()), t.size(), t.release()); + return temporary_buffer::maybe_unsafe_from_deleter(reinterpret_cast(t.get_write()), t.size(), t.release()); }); } diff --git a/include/seastar/core/internal/cpu_profiler.hh b/include/seastar/core/internal/cpu_profiler.hh new file mode 100644 index 00000000000..3ed88cc3a1c --- /dev/null +++ b/include/seastar/core/internal/cpu_profiler.hh @@ -0,0 +1,185 @@ +/* + * This file is open source software, licensed to you under the terms + * of the Apache License, Version 2.0 (the "License"). See the NOTICE file + * distributed with this work for additional information regarding copyright + * ownership. You may not use this file except in compliance with the License. + * + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +/* + * Copyright (C) 2023 ScyllaDB + */ + +#pragma once + +#include +#include +#include +#include +#include + +#include + +#include +#include +#include +#include +#include +#include + +namespace seastar { + +class reactor; + +struct cpu_profiler_trace { + using kernel_trace_vec = boost::container::static_vector; + simple_backtrace user_backtrace; + kernel_trace_vec kernel_backtrace; + // The scheduling group active at the time the same was taken. Note that + // non-task reactor work (such as polling) ends up the in the default + // scheduling group (with name "main"). + scheduling_group sg; +}; + +constexpr size_t max_number_of_traces = 128; + +namespace internal { + +// Temporarily enable/disable the CPU profiler from taking stacktraces on this thread, +// but don't disable the profiler completely. This can be used disable the profiler +// for cases when taking a backtrace isn't valid (IE JIT generated code). +void profiler_drop_stacktraces(bool) noexcept; + +// A small RAII object to disable profiling temporarily +// +// This is not reentrant. +class scoped_disable_profile_temporarily { +public: + scoped_disable_profile_temporarily() noexcept { + profiler_drop_stacktraces(true); + } + ~scoped_disable_profile_temporarily() noexcept { + profiler_drop_stacktraces(false); + } +}; + +struct cpu_profiler_config { + bool enabled; + std::chrono::nanoseconds period; +}; + +struct cpu_profiler_stats { + unsigned dropped_samples_from_manual_disablement{0}; + unsigned dropped_samples_from_exceptions{0}; + unsigned dropped_samples_from_buffer_full{0}; + unsigned dropped_samples_from_mutex_contention{0}; + unsigned dropped_samples_from_context_switches{0}; + + void clear_dropped() { + dropped_samples_from_manual_disablement = 0; + dropped_samples_from_exceptions = 0; + dropped_samples_from_buffer_full = 0; + dropped_samples_from_mutex_contention = 0; + dropped_samples_from_context_switches = 0; + } + + unsigned sum_dropped() const { + return dropped_samples_from_manual_disablement + + dropped_samples_from_buffer_full + + dropped_samples_from_exceptions + + dropped_samples_from_mutex_contention + + dropped_samples_from_context_switches; + } +}; + +class cpu_profiler { +private: + circular_buffer_fixed_capacity _traces; + // The operations in `_traces` are not reentrant. Therefore mutex is used to ensure + // that an interrupt cannot access `_traces` if the interrupted thread was already + // accessing it. + signal_mutex _traces_mutex; + cpu_profiler_config _cfg; + std::chrono::nanoseconds _last_set_timeout; + cpu_profiler_stats _stats; + bool _is_stopped{true}; + + + bool is_enabled() const; + std::chrono::nanoseconds period() const; + std::chrono::nanoseconds get_next_timeout(); + +protected: + friend reactor; + +public: + static int signal_number() { return SIGRTMIN + 2; } + + cpu_profiler(cpu_profiler_config cfg) : _cfg(cfg) {} + + // Allows for the sampling period of the profiler to be adjusted + // and the profiler to be enabled and disabled. + void update_config(cpu_profiler_config cfg); + // Stops the profiler if running and prevents it from starting until + // `start()` is explicitly called. + void stop(); + // Allows to profiler to run when it's enabled via the `cpu_profiler_config`. + void start(); + void on_signal(); + size_t results(std::vector& results_buffer); + + virtual ~cpu_profiler() = default; + virtual void arm_timer(std::chrono::nanoseconds) = 0; + virtual void disarm_timer() = 0; + virtual bool is_spurious_signal() { return false; } + virtual std::optional + try_get_kernel_backtrace() { return std::nullopt; } +}; + +class cpu_profiler_posix_timer : public cpu_profiler { + posix_timer _timer; +public: + cpu_profiler_posix_timer(cpu_profiler_config cfg) + : cpu_profiler(cfg) + // CLOCK_MONOTONIC is used here in place of CLOCK_THREAD_CPUTIME_ID. + // This is since for intervals of ~5ms or less CLOCK_THREAD_CPUTIME_ID + // fires 200-600% after it's configured time. Therefore it is not granular + // enough for cases where the reactor is configured to sleep when idle and + // is only active for short intervals. CLOCK_MONOTONIC doesn't suffer from + // this issue. + , _timer({signal_number()}, CLOCK_MONOTONIC) {} + + virtual ~cpu_profiler_posix_timer() override = default; + virtual void arm_timer(std::chrono::nanoseconds) override; + virtual void disarm_timer() override; +}; + +class cpu_profiler_linux_perf_event : public cpu_profiler { + linux_perf_event _perf_event; +public: + static std::unique_ptr try_make(cpu_profiler_config); + cpu_profiler_linux_perf_event(linux_perf_event perf_event, cpu_profiler_config cfg) + : cpu_profiler(cfg) + , _perf_event(std::move(perf_event)) {} + + virtual ~cpu_profiler_linux_perf_event() override = default; + virtual void arm_timer(std::chrono::nanoseconds) override; + virtual void disarm_timer() override; + virtual bool is_spurious_signal() override; + virtual std::optional + try_get_kernel_backtrace() override; +}; + +std::unique_ptr make_cpu_profiler(cpu_profiler_config cfg = {false, std::chrono::milliseconds(100)}); + +} +} diff --git a/include/seastar/core/internal/signal_mutex.hh b/include/seastar/core/internal/signal_mutex.hh new file mode 100644 index 00000000000..e48841cc772 --- /dev/null +++ b/include/seastar/core/internal/signal_mutex.hh @@ -0,0 +1,52 @@ +/* + * This file is open source software, licensed to you under the terms + * of the Apache License, Version 2.0 (the "License"). See the NOTICE file + * distributed with this work for additional information regarding copyright + * ownership. You may not use this file except in compliance with the License. + * + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +/* + * Copyright (C) 2025 ScyllaDB + */ + +#pragma once + +#include +#include + +namespace seastar::internal { + +/// A lightweight mutex designed to work with interrupts +/// utilizing only compiler barriers. +class signal_mutex { +public: + class guard { + private: + signal_mutex* _mutex; + guard(signal_mutex* m) : _mutex(m) {} + friend class signal_mutex; + public: + guard(guard&& o) : _mutex(o._mutex) { o._mutex = nullptr; } + ~guard(); + }; + + // Returns a `guard` if the lock was acquired. + // Otherwise returns a nullopt. + std::optional try_lock(); + +private: + friend class guard; + std::atomic_bool _mutex; +}; + +} // namespace seastar::internal diff --git a/include/seastar/core/internal/stall_detector.hh b/include/seastar/core/internal/stall_detector.hh index 8976142f5ca..8c7fcac2e85 100644 --- a/include/seastar/core/internal/stall_detector.hh +++ b/include/seastar/core/internal/stall_detector.hh @@ -32,6 +32,7 @@ #include #include #include +#include namespace seastar { @@ -94,85 +95,25 @@ public: }; class cpu_stall_detector_posix_timer : public cpu_stall_detector { - timer_t _timer; + posix_timer _timer; public: explicit cpu_stall_detector_posix_timer(cpu_stall_detector_config cfg = {}); - virtual ~cpu_stall_detector_posix_timer() override; + virtual ~cpu_stall_detector_posix_timer() override = default; private: virtual void arm_timer() override; virtual void start_sleep() override; }; class cpu_stall_detector_linux_perf_event : public cpu_stall_detector { - file_desc _fd; - bool _enabled = false; - uint64_t _current_period = 0; - struct ::perf_event_mmap_page* _mmap; - char* _data_area; - size_t _data_area_mask; - // after the detector has been armed (i.e., _enabled is true), this - // is the moment at or after which the next signal is expected to occur - // and can be used for detecting spurious signals - sched_clock::time_point _next_signal_time{}; -private: - class data_area_reader { - cpu_stall_detector_linux_perf_event& _p; - const char* _data_area; - size_t _data_area_mask; - uint64_t _head; - uint64_t _tail; - public: - explicit data_area_reader(cpu_stall_detector_linux_perf_event& p) - : _p(p) - , _data_area(p._data_area) - , _data_area_mask(p._data_area_mask) { - _head = _p._mmap->data_head; - _tail = _p._mmap->data_tail; - std::atomic_thread_fence(std::memory_order_acquire); // required after reading data_head - } - ~data_area_reader() { - std::atomic_thread_fence(std::memory_order_release); // not documented, but probably required before writing data_tail - _p._mmap->data_tail = _tail; - } - uint64_t read_u64() { - uint64_t ret; - // We cannot wrap around if the 8-byte unit is aligned - std::copy_n(_data_area + (_tail & _data_area_mask), 8, reinterpret_cast(&ret)); - _tail += 8; - return ret; - } - template - S read_struct() { - static_assert(sizeof(S) % 8 == 0); - S ret; - char* p = reinterpret_cast(&ret); - for (size_t i = 0; i != sizeof(S); i += 8) { - uint64_t w = read_u64(); - std::copy_n(reinterpret_cast(&w), 8, p + i); - } - return ret; - } - void skip(uint64_t bytes_to_skip) { - _tail += bytes_to_skip; - } - // skip all the remaining data in the buffer, as-if calling read until - // have_data returns false (but much faster) - void skip_all() { - _tail = _head; - } - bool have_data() const { - return _head != _tail; - } - }; - - virtual void maybe_report_kernel_trace(backtrace_buffer& buf) override; + linux_perf_event _perf_event; public: static std::unique_ptr try_make(cpu_stall_detector_config cfg = {}); - explicit cpu_stall_detector_linux_perf_event(file_desc fd, cpu_stall_detector_config cfg = {}); - ~cpu_stall_detector_linux_perf_event(); + explicit cpu_stall_detector_linux_perf_event(linux_perf_event perf_event, cpu_stall_detector_config cfg = {}); + virtual ~cpu_stall_detector_linux_perf_event() override = default; virtual void arm_timer() override; virtual void start_sleep() override; virtual bool is_spurious_signal() override; + virtual void maybe_report_kernel_trace(backtrace_buffer& buf) override; }; std::unique_ptr make_cpu_stall_detector(cpu_stall_detector_config cfg = {}); diff --git a/include/seastar/core/internal/timers.hh b/include/seastar/core/internal/timers.hh new file mode 100644 index 00000000000..ec4995bfe34 --- /dev/null +++ b/include/seastar/core/internal/timers.hh @@ -0,0 +1,146 @@ +/* + * This file is open source software, licensed to you under the terms + * of the Apache License, Version 2.0 (the "License"). See the NOTICE file + * distributed with this work for additional information regarding copyright + * ownership. You may not use this file except in compliance with the License. + * + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +/* + * Copyright (C) 2023 ScyllaDB + */ + +#pragma once + +#ifndef SEASTAR_MODULE +#include +#include +#include +#include +#include + +#include +#endif + +#include +#include + +namespace seastar { +namespace internal { + +struct timer_cfg { + int signal_number; +}; + +class posix_timer { + timer_t _timer; +public: + explicit posix_timer(timer_cfg cfg, clockid_t clock_id = CLOCK_THREAD_CPUTIME_ID); + virtual ~posix_timer(); + void arm_timer(std::chrono::nanoseconds); + void disarm_timer(); +}; + +class linux_perf_event { + file_desc _fd; + bool _enabled = false; + uint64_t _current_period = 0; + struct ::perf_event_mmap_page* _mmap; + char* _data_area; + size_t _data_area_mask; + // after the detector has been armed (i.e., _enabled is true), this + // is the moment at or after which the next signal is expected to occur + // and can be used for detecting spurious signals + sched_clock::time_point _next_signal_time{}; +private: + class data_area_reader { + std::reference_wrapper _p; + const char* _data_area; + size_t _data_area_mask; + uint64_t _head; + uint64_t _tail; + public: + explicit data_area_reader(linux_perf_event& p) + : _p(p) + , _data_area(p._data_area) + , _data_area_mask(p._data_area_mask) { + _head = _p.get()._mmap->data_head; + _tail = _p.get()._mmap->data_tail; + std::atomic_thread_fence(std::memory_order_acquire); // required after reading data_head + } + data_area_reader(data_area_reader&& o) + : _p(o._p) + , _data_area(o._data_area) + , _data_area_mask(o._data_area_mask) + , _head(o._head) + , _tail(o._tail) { + o._data_area = nullptr; + } + ~data_area_reader() { + if(_data_area != nullptr) { + std::atomic_thread_fence(std::memory_order_release); // not documented, but probably required before writing data_tail + _p.get()._mmap->data_tail = _tail; + } + } + uint64_t read_u64() { + + uint64_t ret; + // We cannot wrap around if the 8-byte unit is aligned + std::copy_n(_data_area + (_tail & _data_area_mask), 8, reinterpret_cast(&ret)); + _tail += 8; + return ret; + } + template + S read_struct() { + static_assert(sizeof(S) % 8 == 0); + S ret; + char* p = reinterpret_cast(&ret); + for (size_t i = 0; i != sizeof(S); i += 8) { + uint64_t w = read_u64(); + std::copy_n(reinterpret_cast(&w), 8, p + i); + } + return ret; + } + void skip(uint64_t bytes_to_skip) { + _tail += bytes_to_skip; + } + // skip all the remaining data in the buffer, as-if calling read until + // have_data returns false (but much faster) + void skip_all() { + _tail = _head; + } + bool have_data() const { + return _head != _tail; + } + }; + + explicit linux_perf_event(file_desc fd); +public: + + class kernel_backtrace { + data_area_reader _reader; + public: + kernel_backtrace(data_area_reader reader) : _reader(std::move(reader)) {} + void read_backtrace(std::function); + }; + + linux_perf_event(linux_perf_event&&); + static linux_perf_event try_make(timer_cfg cfg); + ~linux_perf_event(); + void arm_timer(std::chrono::nanoseconds); + void disarm_timer(); + bool is_spurious_signal(); + std::optional try_get_kernel_backtrace(); +}; + +} +} diff --git a/include/seastar/core/io_queue.hh b/include/seastar/core/io_queue.hh index 69411e901f4..d07db79feb4 100644 --- a/include/seastar/core/io_queue.hh +++ b/include/seastar/core/io_queue.hh @@ -23,6 +23,7 @@ #include #include +#include #include #include #include @@ -192,6 +193,13 @@ public: std::chrono::milliseconds stall_threshold = std::chrono::milliseconds(100); std::chrono::microseconds tau = std::chrono::milliseconds(5); std::optional physical_block_size; // Override for disks that lie about their physical block size + + // Original values of io-properties (if available) + size_t read_bytes_rate = std::numeric_limits::max(); + size_t write_bytes_rate = std::numeric_limits::max(); + size_t read_req_rate = std::numeric_limits::max(); + size_t write_req_rate = std::numeric_limits::max(); + bool max_cost_function = true; }; io_queue(io_group_ptr group, internal::io_sink& sink); diff --git a/include/seastar/core/memory.hh b/include/seastar/core/memory.hh index 43a0350512d..40aa315eb97 100644 --- a/include/seastar/core/memory.hh +++ b/include/seastar/core/memory.hh @@ -186,6 +186,17 @@ public: } }; +// Within the scope of this object, the allocator will not abort (when it +// would normally do so, i.e., when abort_on_allocation_failure is true), +// but rather fall back to the system allocator. +struct scoped_system_alloc_fallback { + scoped_system_alloc_fallback() noexcept; + ~scoped_system_alloc_fallback() noexcept; + + scoped_system_alloc_fallback(const scoped_system_alloc_fallback&) = delete; + scoped_system_alloc_fallback& operator=(const scoped_system_alloc_fallback&) = delete; +}; + // Disables heap profiling as long as this object is alive. // Can be nested, in which case the profiling is re-enabled when all // the objects go out of scope. @@ -302,6 +313,7 @@ class statistics { uint64_t _reclaims; uint64_t _large_allocs; uint64_t _failed_allocs; + uint64_t _fallback_allocs; uint64_t _foreign_mallocs; uint64_t _foreign_frees; @@ -309,11 +321,11 @@ class statistics { private: statistics(uint64_t mallocs, uint64_t frees, uint64_t cross_cpu_frees, uint64_t total_memory, uint64_t free_memory, uint64_t total_bytes_allocated, uint64_t reclaims, - uint64_t large_allocs, uint64_t failed_allocs, + uint64_t large_allocs, uint64_t failed_allocs, uint64_t fallback_allocs, uint64_t foreign_mallocs, uint64_t foreign_frees, uint64_t foreign_cross_frees) : _mallocs(mallocs), _frees(frees), _cross_cpu_frees(cross_cpu_frees) , _total_memory(total_memory), _free_memory(free_memory), _total_bytes_allocated(total_bytes_allocated), _reclaims(reclaims) - , _large_allocs(large_allocs), _failed_allocs(failed_allocs) + , _large_allocs(large_allocs), _failed_allocs(failed_allocs), _fallback_allocs(fallback_allocs) , _foreign_mallocs(foreign_mallocs), _foreign_frees(foreign_frees) , _foreign_cross_frees(foreign_cross_frees) {} public: @@ -339,13 +351,18 @@ public: /// Number of allocations which violated the large allocation threshold uint64_t large_allocations() const { return _large_allocs; } /// Number of allocations which failed, i.e., where the required memory could not be obtained - /// even after reclaim was attempted + /// even after reclaim was attempted and which did not fallback (see fallback_allocations()) uint64_t failed_allocations() const { return _failed_allocs; } - /// Number of foreign allocations + /// Number of allocations which fell back to the system allocator, i.e., because they were + /// in a fallback allocation scope. These are not counted in failed_allocations. + uint64_t fallback_allocations() const { return _fallback_allocs; } + /// Number of foreign allocations, which are all allocations which use the system allocator. + /// These include allocations on alien threads, allocations (even on reactor threads) before + /// the allocator is initialized and allocations in a fallback allocation scope. uint64_t foreign_mallocs() const { return _foreign_mallocs; } - /// Number of foreign frees + /// Number of foreign frees (frees of non-seastar-heap pointers) on alien threads uint64_t foreign_frees() const { return _foreign_frees; } - /// Number of foreign frees on reactor threads + /// Number of foreign frees (frees of non-seastar-heap pointers) on reactor threads uint64_t foreign_cross_frees() const { return _foreign_cross_frees; } friend statistics stats(); }; diff --git a/include/seastar/core/metrics.hh b/include/seastar/core/metrics.hh index 93cd400a704..5e97e1e7889 100644 --- a/include/seastar/core/metrics.hh +++ b/include/seastar/core/metrics.hh @@ -435,6 +435,7 @@ public: virtual metric_groups_def& add_metric(group_name_type name, const metric_definition& md) = 0; virtual metric_groups_def& add_group(group_name_type name, const std::initializer_list& l) = 0; virtual metric_groups_def& add_group(group_name_type name, const std::vector& l) = 0; + virtual int get_handle() const = 0; }; escaped_string shard(); @@ -649,6 +650,13 @@ impl::metric_definition_impl make_total_operations(metric_name_type name, return make_counter(name, std::forward(val), d, labels).set_type("total_operations"); } +/*! + * \brief Update the aggregation labels of a metric family + */ +void update_aggregate_labels(const group_name_type& group_name, + const metric_name_type& metric_name, + const std::vector