From 768516c69e2f944a2f4a727ca3405bc3c6e4d52e Mon Sep 17 00:00:00 2001 From: Travis Downs Date: Thu, 30 Jul 2026 12:01:52 -0400 Subject: [PATCH 01/89] ragel: return consumption_result from the parser call operator ragel_parser_base::operator() returns the legacy optional consumer result, so every parser derived from it is an obsolete input_stream consumer. Passing such a parser to input_stream::consume() converts that optional through the deprecated consumption_result(optional) constructor. Return consumption_result directly instead: a completed parse becomes stop_consuming carrying the unconsumed remainder, an incomplete one becomes continue_consuming. That is exactly the mapping consume() was applying implicitly, so parser behaviour is unchanged. Update the places that call a parser directly rather than through consume(). The two in content_source.hh already produce consumption_result, so they lose a conversion step. request_parser_test and chunk_parsers_test only ask whether the parse reached a terminal state, which is now a check for stop_consuming rather than for an engaged optional. The unconsumed_remainder alias goes away with the old return type. input_stream::unconsumed_remainder still names the optional type for anyone who needs it. --- include/seastar/core/ragel.hh | 8 ++++---- include/seastar/http/internal/content_source.hh | 14 +++++++------- tests/unit/chunk_parsers_test.cc | 7 +++++-- tests/unit/request_parser_test.cc | 4 +++- 4 files changed, 19 insertions(+), 14 deletions(-) diff --git a/include/seastar/core/ragel.hh b/include/seastar/core/ragel.hh index 53f96c1994a..867db7f1c5d 100644 --- a/include/seastar/core/ragel.hh +++ b/include/seastar/core/ragel.hh @@ -21,6 +21,7 @@ #pragma once +#include #include #include #include @@ -124,17 +125,16 @@ protected: return std::move(_builder).get(); } public: - using unconsumed_remainder = std::optional>; - future operator()(temporary_buffer buf) { + future> operator()(temporary_buffer buf) { char* p = buf.get_write(); char* pe = p + buf.size(); char* eof = buf.empty() ? pe : nullptr; char* parsed = static_cast(this)->parse(p, pe, eof); if (parsed) { buf.trim_front(parsed - p); - return make_ready_future(std::move(buf)); + return make_ready_future>(stop_consuming{std::move(buf)}); } - return make_ready_future(); + return make_ready_future>(continue_consuming{}); } }; diff --git a/include/seastar/http/internal/content_source.hh b/include/seastar/http/internal/content_source.hh index ae053549e08..6d3ab548371 100644 --- a/include/seastar/http/internal/content_source.hh +++ b/include/seastar/http/internal/content_source.hh @@ -138,8 +138,8 @@ class chunked_source_impl : public data_source_impl { switch (_ps) { // "data" buffer is non-empty case parsing_state::size_and_ext: - return _size_and_ext_parser(std::move(data)).then([this] (std::optional> res) { - if (res.has_value()) { + return _size_and_ext_parser(std::move(data)).then([this] (consumption_result_type res) { + if (auto* stop = std::get_if>(&res.get())) { if (_size_and_ext_parser.failed()) { return make_exception_future(bad_request_exception("Can't parse chunk size and extensions")); } @@ -164,10 +164,10 @@ class chunked_source_impl : public data_source_impl { } else { _ps = parsing_state::body; } - if (res->empty()) { + if (stop->get_buffer().empty()) { return make_ready_future(continue_consuming{}); } - return this->operator()(std::move(res.value())); + return this->operator()(std::move(stop->get_buffer())); } else { return make_ready_future(continue_consuming{}); } @@ -211,8 +211,8 @@ class chunked_source_impl : public data_source_impl { } return this->operator()(std::move(data)); case parsing_state::trailer_part: - return _trailer_parser(std::move(data)).then([this] (std::optional> res) { - if (res.has_value()) { + return _trailer_parser(std::move(data)).then([this] (consumption_result_type res) { + if (auto* stop = std::get_if>(&res.get())) { if (_trailer_parser.failed()) { return make_exception_future(bad_request_exception("Can't parse chunked request trailer")); } @@ -220,7 +220,7 @@ class chunked_source_impl : public data_source_impl { _trailing_headers = _trailer_parser.get_parsed_headers(); _end_of_request = true; _remaining_bytes = 0; - return make_ready_future(stop_consuming(std::move(*res))); + return make_ready_future(std::move(*stop)); } else { return make_ready_future(continue_consuming{}); } diff --git a/tests/unit/chunk_parsers_test.cc b/tests/unit/chunk_parsers_test.cc index d4a211e5bd9..34df2f6316e 100644 --- a/tests/unit/chunk_parsers_test.cc +++ b/tests/unit/chunk_parsers_test.cc @@ -29,6 +29,7 @@ #include #include #include +#include #include using namespace seastar; @@ -64,7 +65,8 @@ SEASTAR_TEST_CASE(test_size_and_extensions_parsing) { http_chunk_size_and_ext_parser parser; for (auto& tset : tests) { parser.init(); - BOOST_REQUIRE(parser(tset.buf()).get().has_value()); + // the parser is done once it hands back the unconsumed remainder + BOOST_REQUIRE(std::holds_alternative>(parser(tset.buf()).get().get())); BOOST_REQUIRE_NE(parser.failed(), tset.parsable); if (tset.parsable) { BOOST_REQUIRE_EQUAL(parser.get_size(), std::move(tset.size)); @@ -109,7 +111,8 @@ SEASTAR_TEST_CASE(test_trailer_headers_parsing) { http_chunk_trailer_parser parser; for (auto& tset : tests) { parser.init(); - BOOST_REQUIRE(parser(tset.buf()).get().has_value()); + // the parser is done once it hands back the unconsumed remainder + BOOST_REQUIRE(std::holds_alternative>(parser(tset.buf()).get().get())); BOOST_REQUIRE_NE(parser.failed(), tset.parsable); if (tset.parsable) { auto heads = parser.get_parsed_headers(); diff --git a/tests/unit/request_parser_test.cc b/tests/unit/request_parser_test.cc index 13b45a7be1a..1cf8655f305 100644 --- a/tests/unit/request_parser_test.cc +++ b/tests/unit/request_parser_test.cc @@ -28,6 +28,7 @@ #include #include #include +#include #include using namespace seastar; @@ -63,7 +64,8 @@ SEASTAR_TEST_CASE(test_header_parsing) { http_request_parser parser; for (auto& tset : tests) { parser.init(); - BOOST_REQUIRE(parser(tset.buf()).get().has_value()); + // the parser is done once it hands back the unconsumed remainder + BOOST_REQUIRE(std::holds_alternative>(parser(tset.buf()).get().get())); BOOST_REQUIRE_NE(parser.failed(), tset.parsable); if (tset.parsable) { auto req = parser.get_parsed_request(); From 90cedbf25e39def2037de523ce2056dab4bdc932 Mon Sep 17 00:00:00 2001 From: Travis Downs Date: Thu, 30 Jul 2026 12:01:52 -0400 Subject: [PATCH 02/89] demos, tests: migrate the remaining obsolete input_stream consumers The line count demo's reader and the fstream test's consumer both return the legacy optional result, which input_stream::consume() converts through the deprecated consumption_result(optional) constructor. Return consumption_result directly. With the ragel parsers converted, these were the last obsolete consumers in the tree. --- demos/line_count_demo.cc | 8 ++++---- tests/unit/fstream_test.cc | 6 +++--- 2 files changed, 7 insertions(+), 7 deletions(-) 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/tests/unit/fstream_test.cc b/tests/unit/fstream_test.cc index b673a63e3cf..26df5630631 100644 --- a/tests/unit/fstream_test.cc +++ b/tests/unit/fstream_test.cc @@ -243,16 +243,16 @@ future<> test_consume_until_end(uint64_t size) { }).then([size, f] (size_t real_size) { BOOST_REQUIRE_EQUAL(size, real_size); }).then([size, f] { - auto consumer = [offset = uint64_t(0), size] (temporary_buffer buf) mutable -> future::unconsumed_remainder> { + auto consumer = [offset = uint64_t(0), size] (temporary_buffer buf) mutable -> future> { if (!buf) { - return make_ready_future::unconsumed_remainder>(temporary_buffer()); + return make_ready_future>(stop_consuming{temporary_buffer()}); } BOOST_REQUIRE(offset + buf.size() <= size); std::vector expected(buf.size()); std::iota(expected.begin(), expected.end(), offset); offset += buf.size(); BOOST_REQUIRE(std::equal(buf.begin(), buf.end(), expected.begin())); - return make_ready_future::unconsumed_remainder>(std::nullopt); + return make_ready_future>(continue_consuming{}); }; return do_with(make_file_input_stream(f), std::move(consumer), [] (input_stream& in, auto& consumer) { return in.consume(consumer).then([&in] { From 2ebea85bd68f3bfdec0e560fa68864618fe8e906 Mon Sep 17 00:00:00 2001 From: Kefu Chai Date: Fri, 26 Jun 2026 22:50:34 +0800 Subject: [PATCH 03/89] core/reactor: don't read the thread-local shard count in smp::qs_deleter smp::qs_deleter::operator() sizes its destruction loops with this_smp_shard_count(), which dereferences the thread_local _this_smp. The deleter runs from ~smp(), but _this_smp is only set on reactor threads, so it is null when the smp is destroyed on another thread. That happens when the reactor runs on a dedicated thread and the app_template is destroyed elsewhere, for example at exit on the main thread; the deleter then dereferences null and faults under ASan/UBSan. The queue array is allocated from _shard_count, so store that count in the deleter and use it at teardown instead of querying the thread-local. Signed-off-by: Kefu Chai --- include/seastar/core/smp.hh | 1 + src/core/reactor.cc | 6 +++--- 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/include/seastar/core/smp.hh b/include/seastar/core/smp.hh index 379c89c753f..8af98d9e029 100644 --- a/include/seastar/core/smp.hh +++ b/include/seastar/core/smp.hh @@ -319,6 +319,7 @@ class smp : public std::enable_shared_from_this { std::optional> _all_event_loops_done; std::unique_ptr _prefaulter; struct qs_deleter { + unsigned shard_count; void operator()(smp_message_queue** qs) const; }; std::unique_ptr _qs_owner; diff --git a/src/core/reactor.cc b/src/core/reactor.cc index fdd2f6c55ed..a7d22c3cbd1 100644 --- a/src/core/reactor.cc +++ b/src/core/reactor.cc @@ -4320,8 +4320,8 @@ void install_oneshot_signal_handler() { #endif void smp::qs_deleter::operator()(smp_message_queue** qs) const { - for (unsigned i = 0; i < this_smp_shard_count(); i++) { - for (unsigned j = 0; j < this_smp_shard_count(); j++) { + for (unsigned i = 0; i < shard_count; i++) { + for (unsigned j = 0; j < shard_count; j++) { qs[i][j].~smp_message_queue(); } ::operator delete[](qs[i], std::align_val_t(alignof(smp_message_queue)) @@ -4730,7 +4730,7 @@ void smp::configure(const smp_options& smp_opts, const reactor_options& reactor_ seastar_logger.info("Reactor backend: {}", backend_selector); - _qs_owner = decltype(smp::_qs_owner){new smp_message_queue* [_shard_count], qs_deleter{}}; + _qs_owner = decltype(smp::_qs_owner){new smp_message_queue* [_shard_count], qs_deleter{_shard_count}}; auto allocate_qs_owner = [this] (unsigned i) { // smp_message_queue has members with hefty alignment requirements. From d9e8305dbbbb9ff82ea432c7b64e3f229a03a0b9 Mon Sep 17 00:00:00 2001 From: Stephan Dollberg Date: Tue, 28 Jul 2026 17:14:36 -0400 Subject: [PATCH 04/89] log: avoid deprecated fmt12 format string conversion warnings fmt 12 deprecates the implicit basic_format_string -> basic_string_view conversion in favor of basic_format_string::get(), causing warnings at every affected logger instantiation. Use get() before passing the string to failed_to_log() or fmt::runtime(). The accessor was only introduced in fmt 10, so retain the existing conversion for fmt 9 and for builds without SEASTAR_LOGGER_COMPILE_TIME_FMT. Both version checks are hidden in a single internal::format_string_view() helper, so the call sites need no #ifdef of their own. The helper is a single template deducing the format string type, so it serves both the compile-time fmt::format_string and the runtime std::string_view. Passing the already-built fmt.format member (rather than a raw string) avoids any conversion, so fmt's compile-time format-string checking is preserved. Co-authored-by: Travis Downs --- include/seastar/util/log.hh | 29 ++++++++++++++++++++++++++--- 1 file changed, 26 insertions(+), 3 deletions(-) diff --git a/include/seastar/util/log.hh b/include/seastar/util/log.hh index 6f789c17ad6..1c42d055a45 100644 --- a/include/seastar/util/log.hh +++ b/include/seastar/util/log.hh @@ -47,6 +47,29 @@ namespace seastar { class logger; class logger_registry; +namespace internal { + +// Get a format_info's format string as a string_view. +// +// fmt 12 deprecates the implicit basic_format_string -> basic_string_view +// conversion in favour of basic_format_string::get(), which was only added +// in fmt 10. Keeping the version check in here means the call sites need no +// #ifdef of their own. +// +// The format string is a deduced template parameter so one template serves +// both the compile-time fmt::format_string and the runtime +// std::string_view without the caller needing to know which is in play. +template +fmt::string_view format_string_view(const FormatString& format) noexcept { +#if defined(SEASTAR_LOGGER_COMPILE_TIME_FMT) && FMT_VERSION >= 100000 + return format.get(); +#else + return fmt::string_view(format); +#endif +} + +} // namespace internal + /// \brief Logger class for ostream or syslog. /// /// Java style api for logging. @@ -268,7 +291,7 @@ public: }); do_log(level, writer); } catch (...) { - failed_to_log(std::current_exception(), fmt::string_view(fmt.format), fmt.loc); + failed_to_log(std::current_exception(), internal::format_string_view(fmt.format), fmt.loc); } } } @@ -295,11 +318,11 @@ public: if (rl.has_dropped_messages()) { it = fmt::format_to(it, "(rate limiting dropped {} similar messages) ", rl.get_and_reset_dropped_messages()); } - return fmt::format_to(it, fmt::runtime(fmt.format), std::forward(args)...); + return fmt::format_to(it, fmt::runtime(internal::format_string_view(fmt.format)), std::forward(args)...); }); do_log(level, writer); } catch (...) { - failed_to_log(std::current_exception(), fmt::string_view(fmt.format), fmt.loc); + failed_to_log(std::current_exception(), internal::format_string_view(fmt.format), fmt.loc); } } } From fb8a6919a101fc179e7440bfe7af019079a4978b Mon Sep 17 00:00:00 2001 From: Pavel Emelyanov Date: Fri, 31 Jul 2026 19:36:39 +0300 Subject: [PATCH 05/89] util/backtrace: Remove deprecated simple_backtrace constructors Remove the two constructors taking a delimiter character that were deprecated in 3133ecdd65 ("util/backtrace: Optimize formatter to reduce memory allocation overhead") back in 2024 Signed-off-by: Pavel Emelyanov Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- include/seastar/util/backtrace.hh | 2 -- 1 file changed, 2 deletions(-) diff --git a/include/seastar/util/backtrace.hh b/include/seastar/util/backtrace.hh index 36795999a4f..0557a264d6e 100644 --- a/include/seastar/util/backtrace.hh +++ b/include/seastar/util/backtrace.hh @@ -107,8 +107,6 @@ private: public: simple_backtrace(vector_type f) noexcept : _frames(std::move(f)), _hash(calculate_hash()) {} simple_backtrace() noexcept = default; - [[deprecated]] simple_backtrace(vector_type f, char delimeter) : _frames(std::move(f)), _hash(calculate_hash()), _delimeter(delimeter) {} - [[deprecated]] simple_backtrace(char delimeter) : _delimeter(delimeter) {} size_t hash() const noexcept { return _hash; } char delimeter() const noexcept { return _delimeter; } From 23d54bd70a2620121d0d1c303f7b0958d2ed55f3 Mon Sep 17 00:00:00 2001 From: Emil Maskovsky Date: Thu, 30 Jul 2026 17:31:28 +0000 Subject: [PATCH 06/89] rpc: don't abort the process on missing client_info auxiliary data client_info::retrieve_auxiliary() asserts that the requested key is present in user_data. That is the wrong contract: an assert documents a programming error, and a missing key is not one - it is a data-dependent outcome the local process does not control. SEASTAR_ASSERT also fires in all build modes, so the miss aborts the whole process in production. A miss is reachable: auxiliary data is typically attached by an initialization verb the peer sends right after connecting. If that handler fails before calling attach_auxiliary() - e.g. an allocation failure under memory pressure - the exception is logged and swallowed for one-way verbs and the connection stays up with empty user_data, permanently. Every later verb that retrieves the key then takes the server down. A transient std::bad_alloc on the peer's initialization path is thereby converted into a hard crash of a server that is otherwise healthy. Instead, log a warning, abort the connection, and throw the new missing_auxiliary_error. Aborting (rather than only throwing) matters because the initialization verb is never re-sent on an established connection, so a connection that missed it can never heal: wait verbs would keep failing one by one and one-way verbs would be swallowed silently, forever. Tearing the connection down makes the failure self-healing - the client's next request opens a fresh connection and runs initialization again. The abort happens before the throw, so the reply path sees connection::error() and does not serialize the exception back over a connection already known to be broken; the warning is logged by report_missing_auxiliary() itself through the connection's logger (client_info is made a friend of rpc::server for the connection lookup), because that reply path logs nothing itself. No rate limiting is needed: after the first miss the connection is dead, so log volume is bounded at roughly one line per broken connection. retrieve_auxiliary_opt() already provides a non-aborting lookup, but leaving retrieve_auxiliary() as it is and asking callers to switch is not equivalent: every existing and future call site has to remember to, and any one that does not reintroduces the process abort. Fixing the accessor makes the safe behaviour the default. Ref SCYLLADB-3348 Closes scylladb/seastar#3581 --- include/seastar/rpc/rpc.hh | 1 + include/seastar/rpc/rpc_types.hh | 22 ++++++++- src/rpc/rpc.cc | 10 ++++ tests/unit/rpc_test.cc | 78 ++++++++++++++++++++++++++++++++ 4 files changed, 109 insertions(+), 2 deletions(-) diff --git a/include/seastar/rpc/rpc.hh b/include/seastar/rpc/rpc.hh index 7293d258129..15626a65ad8 100644 --- a/include/seastar/rpc/rpc.hh +++ b/include/seastar/rpc/rpc.hh @@ -764,6 +764,7 @@ public: } friend connection; friend client; + friend struct client_info; }; using rpc_handler_func = std::function (shared_ptr, std::optional timeout, int64_t msgid, diff --git a/include/seastar/rpc/rpc_types.hh b/include/seastar/rpc/rpc_types.hh index 4a54786b93e..dae54632fa4 100644 --- a/include/seastar/rpc/rpc_types.hh +++ b/include/seastar/rpc/rpc_types.hh @@ -31,7 +31,6 @@ #include #include #include -#include #include #include #include @@ -104,10 +103,22 @@ struct client_info { void attach_auxiliary(const sstring& key, T&& object) { user_data.emplace(key, std::any(std::forward(object))); } + // Logs a warning, aborts the connection identified by conn_id and throws + // missing_auxiliary_error. Defined out of line because rpc::server is an + // incomplete type in this header. + [[noreturn]] void report_missing_auxiliary(const sstring& key) const; + /// Returns a reference to the auxiliary object attached under \c key. + /// + /// If nothing was attached under \c key (for example because the verb + /// that was supposed to call attach_auxiliary() failed before doing so), + /// the connection is aborted and missing_auxiliary_error is thrown. + /// Use retrieve_auxiliary_opt() to probe for a key without side effects. template T& retrieve_auxiliary(const sstring& key) { auto it = user_data.find(key); - SEASTAR_ASSERT(it != user_data.end()); + if (it == user_data.end()) { + report_missing_auxiliary(key); + } return std::any_cast(it->second); } template @@ -178,6 +189,13 @@ class remote_verb_error : public error { using error::error; }; +/// Thrown by client_info::retrieve_auxiliary() when no auxiliary object is +/// attached under the requested key. The connection is aborted before the +/// exception is thrown. +class missing_auxiliary_error : public error { + using error::error; +}; + struct no_wait_type {}; // return this from a callback if client does not want to waiting for a reply diff --git a/src/rpc/rpc.cc b/src/rpc/rpc.cc index aa145097ede..75c3beea50c 100644 --- a/src/rpc/rpc.cc +++ b/src/rpc/rpc.cc @@ -1400,6 +1400,16 @@ future<> server::stop() { ).discard_result(); } +[[noreturn]] void client_info::report_missing_auxiliary(const sstring& key) const { + server.abort_connection(conn_id); + auto msg = format("missing auxiliary data '{}' on connection {}", key, conn_id); + auto it = server._conns.find(conn_id); + if (it != server._conns.end()) { + it->second->get_logger()(*this, log_level::warn, msg); + } + throw missing_auxiliary_error(msg); +} + void server::abort_connection(connection_id id) { auto it = _conns.find(id); if (it == _conns.end()) { diff --git a/tests/unit/rpc_test.cc b/tests/unit/rpc_test.cc index c7a32c5839a..a514ae60646 100644 --- a/tests/unit/rpc_test.cc +++ b/tests/unit/rpc_test.cc @@ -1510,6 +1510,9 @@ SEASTAR_TEST_CASE(test_client_info) { BOOST_REQUIRE_EQUAL(info.retrieve_auxiliary_opt("missing"), nullptr); BOOST_REQUIRE_EQUAL(const_info.retrieve_auxiliary_opt("missing"), nullptr); + BOOST_REQUIRE_THROW(info.retrieve_auxiliary("missing"), rpc::missing_auxiliary_error); + BOOST_REQUIRE_THROW(const_info.retrieve_auxiliary("missing"), rpc::missing_auxiliary_error); + return make_ready_future<>(); }); } @@ -1574,6 +1577,81 @@ SEASTAR_TEST_CASE(test_rpc_abort_connection) { }); } +// A handler retrieving a never-attached auxiliary object must not reach the +// code after the retrieval, and the connection must be aborted: the caller +// of the verb sees closed_error, and so do all subsequent calls. +SEASTAR_TEST_CASE(test_retrieve_missing_auxiliary_wait_handler) { + return rpc_test_env<>::do_with_thread(rpc_test_config(), [] (rpc_test_env<>& env) { + test_rpc_proto::client c1(env.proto(), {}, env.make_socket(), ipv4_addr()); + bool reached_after_retrieve = false; + env.register_handler(1, [&] (rpc::client_info& cinfo, int) { + auto& v = cinfo.retrieve_auxiliary("never_attached"); + reached_after_retrieve = true; + return v; + }).get(); + auto f = env.proto().make_client(1); + BOOST_REQUIRE_THROW(f(c1, 0).get(), rpc::closed_error); + BOOST_REQUIRE(!reached_after_retrieve); + // The connection was aborted server-side; further calls fail too. + BOOST_REQUIRE_THROW(f(c1, 1).get(), rpc::closed_error); + c1.stop().get(); + }); +} + +// Same, for a one-way (no_wait) handler: the exception is logged and +// swallowed by the reply path, the connection is aborted, and the process +// survives. +SEASTAR_TEST_CASE(test_retrieve_missing_auxiliary_no_wait_handler) { + using namespace std::chrono_literals; + static seastar::logger log("rpc"); + return rpc_test_env<>::do_with_thread(rpc_test_config(), [] (rpc_test_env<>& env) { + // Attach a logger so the swallowed exception shows up in the test + // output. No assertion on the log contents; this is a smoke capture. + env.proto().set_logger(&log); + auto reset_logger = defer([&env] () noexcept { env.proto().set_logger(nullptr); }); + test_rpc_proto::client c1(env.proto(), {}, env.make_socket(), ipv4_addr()); + bool handler_entered = false; + env.register_handler(1, [&] (rpc::client_info& cinfo, int) { + handler_entered = true; + cinfo.retrieve_auxiliary("never_attached"); + return rpc::no_wait; + }).get(); + env.register_handler(2, [] (int x) { return x; }).get(); + auto one_way = env.proto().make_client(1); + auto echo = env.proto().make_client(2); + // Sanity check: the connection works before the poisoned verb. + BOOST_REQUIRE_EQUAL(echo(c1, 7).get(), 7); + one_way(c1, 0).get(); // resolves once the request is sent + // The abort is asynchronous wrt. the client; poll (bounded) until + // it is observed. + bool closed = false; + for (int i = 0; i < 500 && !closed; ++i) { + try { + echo(c1, i).get(); + seastar::sleep(10ms).get(); + } catch (rpc::closed_error&) { + closed = true; + } + } + BOOST_REQUIRE(handler_entered); + BOOST_REQUIRE(closed); + c1.stop().get(); + }); +} + +// Retrieving a missing auxiliary object outside of any connection (a +// connection id unknown to the server) must still throw; aborting an unknown +// connection is a no-op. +SEASTAR_TEST_CASE(test_retrieve_missing_auxiliary_standalone) { + return rpc_test_env<>::do_with(rpc_test_config(), [] (rpc_test_env<>& env) { + rpc::client_info info{.server{env.server()}, .conn_id{0}}; + const rpc::client_info& const_info = info; + BOOST_REQUIRE_THROW(info.retrieve_auxiliary("missing"), rpc::missing_auxiliary_error); + BOOST_REQUIRE_THROW(const_info.retrieve_auxiliary("missing"), rpc::missing_auxiliary_error); + return make_ready_future<>(); + }); +} + SEASTAR_THREAD_TEST_CASE(test_rpc_metric_domains) { auto do_one_echo = [] (rpc_test_env<>& env, test_rpc_proto::client& cln, int nr_calls) { env.register_handler(1, [] (int v) { return make_ready_future(v); }).get(); From 585a90ea9d5b63ecdd12e23bddcb1d86554762aa Mon Sep 17 00:00:00 2001 From: Michal Maslanka Date: Thu, 1 Apr 2021 09:04:04 +0200 Subject: [PATCH 07/89] httpd: added listener index to http request Seastar http server implementation supports multiple listeners. It may be required for the handler logic to know which listener the connection is coming from. Added listener_idx field to `httpd::request` to allow handler recognize listener. Signed-off-by: Michal Maslanka --- include/seastar/http/httpd.hh | 13 ++++++++----- include/seastar/http/request.hh | 8 ++++++++ src/http/httpd.cc | 9 +++++---- 3 files changed, 21 insertions(+), 9 deletions(-) diff --git a/include/seastar/http/httpd.hh b/include/seastar/http/httpd.hh index c8c2df15612..b6bee5767e2 100644 --- a/include/seastar/http/httpd.hh +++ b/include/seastar/http/httpd.hh @@ -77,26 +77,29 @@ class connection : public boost::intrusive::list_base_hook<> { queue> _replies { 10 }; bool _done = false; const bool _tls; + int _listener_idx; public: - connection(http_server& server, connected_socket&& fd, bool tls) + connection(http_server& server, connected_socket&& fd, bool tls, int listener_idx) : _server(server) , _fd(std::move(fd)) , _read_buf(_fd.input()) , _write_buf(_fd.output()) , _client_addr(_fd.remote_address()) , _server_addr(_fd.local_address()) - , _tls(tls) { + , _tls(tls) + , _listener_idx(listener_idx) { on_new_connection(); } connection(http_server& server, connected_socket&& fd, - socket_address client_addr, socket_address server_addr, bool tls) + socket_address client_addr, socket_address server_addr, bool tls, int listener_idx) : _server(server) , _fd(std::move(fd)) , _read_buf(_fd.input()) , _write_buf(_fd.output()) , _client_addr(std::move(client_addr)) , _server_addr(std::move(server_addr)) - , _tls(tls) { + , _tls(tls) + , _listener_idx(listener_idx) { on_new_connection(); } ~connection(); @@ -206,7 +209,7 @@ public: static sstring http_date(); private: future<> do_accept_one(int which, bool with_tls); - future<> do_process_connection(connected_socket conn_fd, socket_address remote_address, bool tls); + future<> do_process_connection(connected_socket conn_fd, socket_address remote_address, bool tls, int listener_idx); boost::intrusive::list _connections; friend class seastar::httpd::connection; friend class http_server_tester; diff --git a/include/seastar/http/request.hh b/include/seastar/http/request.hh index bbce49f0453..0d446f06961 100644 --- a/include/seastar/http/request.hh +++ b/include/seastar/http/request.hh @@ -92,6 +92,7 @@ struct request { /// a request. const std::vector* tls_san = nullptr; http::body_writer_type body_writer; // for client + int listener_idx; using query_parameters_type = std::unordered_map, seastar::internal::string_view_hash, std::equal_to<>>; private: @@ -254,6 +255,13 @@ public: bool is_form_post() const { return content_type_class == ctclass::app_x_www_urlencoded; } + /** + * Get index of listener which accepted connection receiving this request + * @return position of listener in server _listeners vector + */ + int get_listener_idx() const { + return listener_idx; + } bool should_keep_alive() const { if (_version == "0.9") { diff --git a/src/http/httpd.cc b/src/http/httpd.cc index be9db4d83d0..6ba3ae7b9bc 100644 --- a/src/http/httpd.cc +++ b/src/http/httpd.cc @@ -212,6 +212,7 @@ future<> connection::read_one() { req->_server_address = this->_server_addr; req->_client_address = this->_client_addr; + req->listener_idx = _listener_idx; if (_tls) { req->protocol_name = "https"; @@ -506,18 +507,18 @@ future<> http_server::do_accept_one(int which, bool tls) { // lambda would be destroyed while its frame is still running. (void)try_with_gate(_task_gate, [this, conn_fd = std::move(ar.connection), - remote_address = std::move(ar.remote_address), tls]() mutable { - return do_process_connection(std::move(conn_fd), std::move(remote_address), tls); + remote_address = std::move(ar.remote_address), tls, which]() mutable { + return do_process_connection(std::move(conn_fd), std::move(remote_address), tls, which); }).handle_exception_type([] (const gate_closed_exception& e) {}); } // Named member coroutine for per-connection processing, called from the // non-coroutine lambda in try_with_gate inside do_accept_one(). Parameters // are passed by value so they live safely in the coroutine frame. -future<> http_server::do_process_connection(connected_socket conn_fd, socket_address remote_address, bool tls) { +future<> http_server::do_process_connection(connected_socket conn_fd, socket_address remote_address, bool tls, int listener_idx) { auto local_address = conn_fd.local_address(); auto conn = std::make_unique(*this, std::move(conn_fd), - std::move(remote_address), std::move(local_address), tls); + std::move(remote_address), std::move(local_address), tls, listener_idx); try { co_await conn->prepare(); } catch (...) { From 2c8818bfa925999b98abe2459d00cd6fcdea34fd Mon Sep 17 00:00:00 2001 From: John Spray Date: Thu, 9 Dec 2021 14:10:25 +0000 Subject: [PATCH 08/89] http: enable specifying a content type on exceptions Since an exception carries some text for the response body text, the raising site might like to specify the content type if it's e.g. json. Signed-off-by: John Spray --- include/seastar/http/exception.hh | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/include/seastar/http/exception.hh b/include/seastar/http/exception.hh index 55116d8f44f..f41a06dfcce 100644 --- a/include/seastar/http/exception.hh +++ b/include/seastar/http/exception.hh @@ -44,6 +44,16 @@ public: : _msg(msg), _status(status) { } + /** + * A base_exception with a content_type is specifying a full response body, whereas + * a base_exception with only a _status is specifying a string that may be wrapped + * in e.g. a json_exception. + */ + base_exception(const std::string& msg, http::reply::status_type status, const std::string &content_type) + : _msg(msg), _status(status), _content_type(content_type) { + } + + virtual const char* what() const noexcept { return _msg.c_str(); } @@ -55,9 +65,14 @@ public: virtual const std::string& str() const { return _msg; } + + virtual const std::string& content_type() const { + return _content_type; + } private: std::string _msg; http::reply::status_type _status; + std::string _content_type; }; From f3c518a07423416f59f7af94118f48b33ddb74ee Mon Sep 17 00:00:00 2001 From: John Spray Date: Thu, 9 Dec 2021 14:11:07 +0000 Subject: [PATCH 09/89] http: don't jsonize exception if has content type This enables throwing a base_exception from a json request handler with a json payload inside it. Signed-off-by: John Spray --- src/http/routes.cc | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/src/http/routes.cc b/src/http/routes.cc index d06e5649e32..3022cc19cd3 100644 --- a/src/http/routes.cc +++ b/src/http/routes.cc @@ -77,7 +77,12 @@ std::unique_ptr routes::exception_reply(std::exception_ptr eptr) { } catch (const redirect_exception& _e) { *rep = _e.to_reply(); } catch (const base_exception& e) { - rep->set_status(e.status(), internal::to_json(e)); + if (e.content_type().size()) { + rep->set_status(e.status(), e.str()); + rep->set_content_type(e.content_type()); + } else { + rep->set_status(e.status(), internal::to_json(e)); + } } catch (...) { rep->set_status(http::reply::status_type::internal_server_error, internal::to_json(std::current_exception())); From 68aac4c672f5bde69c9eadba8c11c928ba99ee7d Mon Sep 17 00:00:00 2001 From: John Spray Date: Thu, 9 Dec 2021 14:29:14 +0000 Subject: [PATCH 10/89] http: use base_exception content type in non-json errors Signed-off-by: John Spray --- include/seastar/http/httpd.hh | 2 +- src/http/httpd.cc | 7 +++++-- 2 files changed, 6 insertions(+), 3 deletions(-) diff --git a/include/seastar/http/httpd.hh b/include/seastar/http/httpd.hh index b6bee5767e2..8a05392d73a 100644 --- a/include/seastar/http/httpd.hh +++ b/include/seastar/http/httpd.hh @@ -118,7 +118,7 @@ public: future<> start_response(); future generate_reply(std::unique_ptr req); - void generate_error_reply_and_close(std::unique_ptr req, http::reply::status_type status, const sstring& msg); + void generate_error_reply_and_close(std::unique_ptr req, http::reply::status_type status, const sstring& msg, const sstring &content_type={}); output_stream& out(); }; diff --git a/src/http/httpd.cc b/src/http/httpd.cc index 6ba3ae7b9bc..1e56620bb45 100644 --- a/src/http/httpd.cc +++ b/src/http/httpd.cc @@ -190,11 +190,14 @@ static void set_header_connection(http::reply& resp, bool keep_alive) { } } -void connection::generate_error_reply_and_close(std::unique_ptr req, http::reply::status_type status, const sstring& msg) { +void connection::generate_error_reply_and_close(std::unique_ptr req, http::reply::status_type status, const sstring& msg, const sstring &content_type) { auto resp = std::make_unique(); // TODO: Handle HTTP/2.0 when it releases resp->set_version(req->_version); resp->set_status(status, msg); + if (!content_type.empty()) { + resp->set_content_type(content_type); + } set_header_connection(*resp, false); _done = true; _replies.push(std::move(resp)); @@ -283,7 +286,7 @@ future<> connection::read_one() { // before passing the request to handler - when we were parsing chunks auto err_req = std::make_unique(); err_req->_version = version; - generate_error_reply_and_close(std::move(err_req), e.status(), e.str()); + generate_error_reply_and_close(std::move(err_req), e.status(), e.str(), e.content_type()); }); }); }); From c5dd8301f44f9e291bac8366472b175515334ef6 Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Wed, 1 Jun 2022 13:35:42 +0100 Subject: [PATCH 11/89] metrics: allow multiple metrics::impl instances Prior to this patch seastar only exposes one global metrics::impl::impl object which holds all metric related data for one application. This patch changes the implementation details such that multiple metrics::impl::impl objects can exist for any given application. Said objects are stored into a map on each shard and created dinamically whenever requested. A metrics::impl::impl is identified by an integer handle that acts as the key for the storage map. Implementation note: in order to avoid issues caused by the ordering of static thread_local objects I had to declare the storage in reactor.cc. (cherry picked from commit 585a8af6a586e161bbfc93a087d06c9c30175ca0) --- include/seastar/core/metrics_api.hh | 7 ++++++- src/core/metrics.cc | 19 ++++++++++++++----- src/core/reactor.cc | 8 ++++++++ 3 files changed, 28 insertions(+), 6 deletions(-) diff --git a/include/seastar/core/metrics_api.hh b/include/seastar/core/metrics_api.hh index e47fc035c80..6d8968e1c8d 100644 --- a/include/seastar/core/metrics_api.hh +++ b/include/seastar/core/metrics_api.hh @@ -41,6 +41,9 @@ namespace impl { using internalized_labels_ref = lw_shared_ptr; + +int default_handle(); + } } } @@ -230,6 +233,8 @@ inline bool operator<(const internalized_holder& lhs, const internalized_holder& class impl; +using metric_implementations = std::unordered_map>; +metric_implementations& get_metric_implementations(); class registered_metric final { metric_info _info; @@ -510,7 +515,7 @@ using values_reference = shared_ptr; foreign_ptr get_values(); -shared_ptr get_local_impl(); +shared_ptr get_local_impl(int handle = default_handle()); void unregister_metric(const metric_id & id); diff --git a/src/core/metrics.cc b/src/core/metrics.cc index 1e3bf662177..e6f041a4239 100644 --- a/src/core/metrics.cc +++ b/src/core/metrics.cc @@ -403,13 +403,17 @@ bool metric_id::operator==( return as_tuple() == id2.as_tuple(); } -// Unfortunately, metrics_impl can not be shared because it -// need to be available before the first users (reactor) will call it +shared_ptr get_local_impl(int handle) { + auto& impls = get_metric_implementations(); + auto [it, inserted] = impls.try_emplace(handle); -shared_ptr get_local_impl() { - static thread_local auto the_impl = ::seastar::make_shared(); - return the_impl; + if (inserted) { + it->second = ::seastar::make_shared(); + } + + return it->second; } + void impl::remove_registration(const metric_id& id) { auto i = get_value_map().find(id.full_name()); if (i != get_value_map().end()) { @@ -616,6 +620,11 @@ void impl::set_metric_family_configs(const std::vector& fa } } } + +int default_handle() { + return 0; +} + } const bool metric_disabled = false; diff --git a/src/core/reactor.cc b/src/core/reactor.cc index a7d22c3cbd1..01bcfd71c85 100644 --- a/src/core/reactor.cc +++ b/src/core/reactor.cc @@ -127,6 +127,7 @@ #include #include #include +#include #include #include #include @@ -175,6 +176,7 @@ #include "core/reactor_backend.hh" #include "core/syscall_result.hh" #include "core/thread_pool.hh" +#include "core/scollectd-impl.hh" #include "syscall_work_queue.hh" #include "cgroup.hh" #ifdef SEASTAR_HAVE_DPDK @@ -4113,6 +4115,12 @@ smp_options::smp_options(program_options::option_group* parent_group) { } +thread_local metrics::impl::metric_implementations metric_impls; + +metrics::impl::metric_implementations& metrics::impl::get_metric_implementations() { + return metric_impls; +} + struct reactor_deleter { void operator()(reactor* p) { p->~reactor(); From f9ae9d7ed267ee9d751afb28c05d61a4450e5886 Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Wed, 22 Jun 2022 17:31:25 +0100 Subject: [PATCH 12/89] metrics: expose metric impl handle to internal api This patch extends the metrics internal apis to use a specific metrics::impl::impl object identified by its integer handle. (cherry picked from commit 6ee4af7) --- include/seastar/core/metrics_api.hh | 14 ++++--- include/seastar/core/metrics_registration.hh | 3 ++ src/core/metrics.cc | 39 +++++++++++--------- 3 files changed, 33 insertions(+), 23 deletions(-) diff --git a/include/seastar/core/metrics_api.hh b/include/seastar/core/metrics_api.hh index 6d8968e1c8d..a954850cd6b 100644 --- a/include/seastar/core/metrics_api.hh +++ b/include/seastar/core/metrics_api.hh @@ -275,10 +275,11 @@ using metric_instances = std::map using metrics_registration = std::vector; class metric_groups_impl : public metric_groups_def { + int _handle; metrics_registration _registration; shared_ptr _impl; // keep impl alive while metrics are registered public: - metric_groups_impl(); + explicit metric_groups_impl(int handle = default_handle()); ~metric_groups_impl(); metric_groups_impl(const metric_groups_impl&) = delete; metric_groups_impl(metric_groups_impl&&) = default; @@ -510,14 +511,15 @@ private: bool apply_relabeling(const relabel_config& rc, metric_info& info); }; -const value_map& get_value_map(); +const value_map& get_value_map(int handle = default_handle()); using values_reference = shared_ptr; -foreign_ptr get_values(); +foreign_ptr get_values(int handle = default_handle()); shared_ptr get_local_impl(int handle = default_handle()); -void unregister_metric(const metric_id & id); + +void unregister_metric(const metric_id & id, int handle = default_handle()); /*! * \brief initialize metric group @@ -525,7 +527,7 @@ void unregister_metric(const metric_id & id); * Create a metric_group_def. * No need to use it directly. */ -std::unique_ptr create_metric_groups(); +std::unique_ptr create_metric_groups(int handle = default_handle()); } @@ -542,7 +544,7 @@ struct options : public program_options::option_group { /*! * \brief set the metrics configuration */ -future<> configure(const options& opts); +future<> configure(const options& opts, int handle = default_handle()); /*! * \brief Perform relabeling and operation on metrics dynamically. diff --git a/include/seastar/core/metrics_registration.hh b/include/seastar/core/metrics_registration.hh index 52ef7a30e35..fe3d74c7a71 100644 --- a/include/seastar/core/metrics_registration.hh +++ b/include/seastar/core/metrics_registration.hh @@ -51,12 +51,15 @@ namespace seastar { namespace metrics { namespace impl { +int default_handle(); class metric_groups_def; struct metric_definition_impl; class metric_groups_impl; } +int default_handle(); + using group_name_type = sstring; /*!< A group of logically related metrics */ class metric_groups; diff --git a/src/core/metrics.cc b/src/core/metrics.cc index e6f041a4239..1806045dc3d 100644 --- a/src/core/metrics.cc +++ b/src/core/metrics.cc @@ -39,6 +39,10 @@ namespace seastar { extern seastar::logger seastar_logger; namespace metrics { +int default_handle() { + return impl::default_handle(); +}; + double_registration::double_registration(std::string what): std::runtime_error(what) {} metric_groups::metric_groups() noexcept : _impl(impl::create_metric_groups()) { @@ -110,11 +114,11 @@ options::options(program_options::option_group* parent_group) { } -future<> configure(const options& opts) { +future<> configure(const options& opts, int handle) { impl::config c; c.hostname = opts.metrics_hostname.get_value(); - return smp::invoke_on_all([c] { - impl::get_local_impl()->set_config(c); + return smp::invoke_on_all([c, handle] { + impl::get_local_impl(handle)->set_config(c); }); } @@ -330,15 +334,15 @@ metric_definition_impl& metric_definition_impl::set_skip_when_empty(bool skip) n return *this; } -std::unique_ptr create_metric_groups() { - return std::make_unique(); +std::unique_ptr create_metric_groups(int handle) { + return std::make_unique(handle); } -metric_groups_impl::metric_groups_impl() {} +metric_groups_impl::metric_groups_impl(int handle) : _handle(handle) {} metric_groups_impl::~metric_groups_impl() { for (const auto& i : _registration) { - unregister_metric(i->info().id); + unregister_metric(i->info().id, _handle); } } @@ -356,14 +360,15 @@ metric_groups_impl& metric_groups_impl::add_metric(group_name_type name, const m // than where the actual metrics are added. // Hence, the shared_ptr owning shard check would fail so we do it only here. if (_impl == nullptr) { - _impl = get_local_impl(); + _impl = get_local_impl(_handle); } - auto internalized_labels = get_local_impl()->internalize_labels(md._impl->labels); + auto internalized_labels = get_local_impl(_handle)->internalize_labels(md._impl->labels); metric_id id(name, md._impl->name, internalized_labels); - auto reg = get_local_impl()->add_registration(id, md._impl->type, md._impl->f, md._impl->d, md._impl->enabled, md._impl->_skip_when_empty, md._impl->aggregate_labels); + auto reg = get_local_impl(_handle)->add_registration( + id, md._impl->type, md._impl->f, md._impl->d, md._impl->enabled, md._impl->_skip_when_empty, md._impl->aggregate_labels); _registration.push_back(std::move(reg)); return *this; @@ -429,20 +434,20 @@ void impl::remove_registration(const metric_id& id) { } } -void unregister_metric(const metric_id & id) { - get_local_impl()->remove_registration(id); +void unregister_metric(const metric_id & id, int handle) { + get_local_impl(handle)->remove_registration(id); } -const value_map& get_value_map() { - return get_local_impl()->get_value_map(); +const value_map& get_value_map(int handle) { + return get_local_impl(handle)->get_value_map(); } -foreign_ptr get_values() { +foreign_ptr get_values(int handle) { shared_ptr res_ref = ::seastar::make_shared(); auto& res = *(res_ref.get()); auto& mv = res.values; - res.metadata = get_local_impl()->metadata(); - auto & functions = get_local_impl()->functions(); + res.metadata = get_local_impl(handle)->metadata(); + auto & functions = get_local_impl(handle)->functions(); for (auto&& i : functions) { value_vector values; for (auto&& v : i) { From 4dfa81da2d4c14ed42554bde910cc761e8f530cf Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Mon, 18 Jul 2022 18:58:56 +0100 Subject: [PATCH 13/89] metrics: expose handle in metric_groups_impl Add a public method to metric_groups_impl that exposes the handle of the internal implementation it is using. This is required in order for the metric_groups class to be able to reset itself to the configured implementation handle. --- include/seastar/core/metrics.hh | 1 + include/seastar/core/metrics_api.hh | 1 + src/core/metrics.cc | 4 ++++ 3 files changed, 6 insertions(+) diff --git a/include/seastar/core/metrics.hh b/include/seastar/core/metrics.hh index 93cd400a704..0803886bfb7 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(); diff --git a/include/seastar/core/metrics_api.hh b/include/seastar/core/metrics_api.hh index a954850cd6b..584cc29632c 100644 --- a/include/seastar/core/metrics_api.hh +++ b/include/seastar/core/metrics_api.hh @@ -286,6 +286,7 @@ public: metric_groups_impl& add_metric(group_name_type name, const metric_definition& md); metric_groups_impl& add_group(group_name_type name, const std::initializer_list& l); metric_groups_impl& add_group(group_name_type name, const std::vector& l); + int get_handle() const; }; class metric_family { diff --git a/src/core/metrics.cc b/src/core/metrics.cc index 1806045dc3d..53c36708721 100644 --- a/src/core/metrics.cc +++ b/src/core/metrics.cc @@ -388,6 +388,10 @@ metric_groups_impl& metric_groups_impl::add_group(group_name_type name, const st return *this; } +int metric_groups_impl::get_handle() const { + return _handle; +} + bool metric_id::operator<( const metric_id& id2) const { return as_tuple() < id2.as_tuple(); From a5eaa6d38297a7fb0147dc334cb69ecc8e80c56e Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Wed, 22 Jun 2022 17:34:52 +0100 Subject: [PATCH 14/89] metrics: expose metric impl handle to external api This patch extends the metrics user facing apis to use a specific metrics::impl::impl object identified by its integer handle. Note that the constructor of 'metric_groups' is marked explicit in this patch and updates two call sites where the constructor was used implicitly. --- include/seastar/core/metrics_registration.hh | 8 ++++---- src/core/io_queue.cc | 2 +- src/core/metrics.cc | 12 ++++++------ src/core/reactor.cc | 2 +- 4 files changed, 12 insertions(+), 12 deletions(-) diff --git a/include/seastar/core/metrics_registration.hh b/include/seastar/core/metrics_registration.hh index fe3d74c7a71..46a37cd5fec 100644 --- a/include/seastar/core/metrics_registration.hh +++ b/include/seastar/core/metrics_registration.hh @@ -93,7 +93,7 @@ public: class metric_groups { std::unique_ptr _impl; public: - metric_groups() noexcept; + explicit metric_groups(int handle = default_handle()) noexcept; metric_groups(metric_groups&&) = default; virtual ~metric_groups(); metric_groups& operator=(metric_groups&&) = default; @@ -102,7 +102,7 @@ public: * * combine the constructor with the add_group functionality. */ - metric_groups(std::initializer_list mg); + metric_groups(std::initializer_list mg, int handle = default_handle()); /*! * \brief Add metrics belonging to the same group. @@ -158,7 +158,7 @@ public: */ class metric_group : public metric_groups { public: - metric_group() noexcept; + explicit metric_group(int handle = default_handle()) noexcept; metric_group(const metric_group&) = delete; metric_group(metric_group&&) = default; virtual ~metric_group(); @@ -169,7 +169,7 @@ public: * * */ - metric_group(const group_name_type& name, std::initializer_list l); + metric_group(const group_name_type& name, std::initializer_list l, int handle = default_handle()); }; diff --git a/src/core/io_queue.cc b/src/core/io_queue.cc index c3468102a90..ff29a883bd2 100644 --- a/src/core/io_queue.cc +++ b/src/core/io_queue.cc @@ -861,7 +861,7 @@ void io_queue::register_stats(sstring name, priority_class_data& pc) { } new_metrics.add_group("io_queue", std::move(metrics)); - pc.metric_groups = std::exchange(new_metrics, {}); + pc.metric_groups = std::exchange(new_metrics, sm::metric_groups{}); } io_queue::priority_class_data& io_queue::find_or_create_class(internal::priority_class pc) { diff --git a/src/core/metrics.cc b/src/core/metrics.cc index 53c36708721..183fcd267a7 100644 --- a/src/core/metrics.cc +++ b/src/core/metrics.cc @@ -45,14 +45,15 @@ int default_handle() { double_registration::double_registration(std::string what): std::runtime_error(what) {} -metric_groups::metric_groups() noexcept : _impl(impl::create_metric_groups()) { +metric_groups::metric_groups(int handle) noexcept : _impl(impl::create_metric_groups(handle)) { } void metric_groups::clear() { - _impl = impl::create_metric_groups(); + const auto current_handle = _impl->get_handle(); + _impl = impl::create_metric_groups(current_handle); } -metric_groups::metric_groups(std::initializer_list mg) : _impl(impl::create_metric_groups()) { +metric_groups::metric_groups(std::initializer_list mg, int handle) : _impl(impl::create_metric_groups(handle)) { for (auto&& i : mg) { add_group(i.name, i.metrics); } @@ -65,10 +66,9 @@ metric_groups& metric_groups::add_group(const group_name_type& name, const std:: _impl->add_group(name, l); return *this; } -metric_group::metric_group() noexcept = default; +metric_group::metric_group(int handle) noexcept : metric_groups(handle) {} metric_group::~metric_group() = default; -metric_group::metric_group(const group_name_type& name, std::initializer_list l) { - add_group(name, l); +metric_group::metric_group(const group_name_type& name, std::initializer_list l, int handle) : metric_groups({metric_group_definition(name, l)}, handle) { } metric_group_definition::metric_group_definition(const group_name_type& name, std::initializer_list l) : name(name), metrics(l) { diff --git a/src/core/reactor.cc b/src/core/reactor.cc index 01bcfd71c85..a2a8fa6ea8d 100644 --- a/src/core/reactor.cc +++ b/src/core/reactor.cc @@ -1062,7 +1062,7 @@ reactor::task_queue::register_stats() { register_net_metrics_for_scheduling_group(new_metrics, _id, group_label); - _metrics = std::exchange(new_metrics, {}); + _metrics = std::exchange(new_metrics, sm::metric_groups{}); } void From af59bf818213c1917d21132393abcc80258196d7 Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Wed, 22 Jun 2022 14:06:39 +0100 Subject: [PATCH 15/89] metrics: Use handle for impl object This patch removes two subsequent calls to `get_local_impl` and reuses the returned handle in that scope. --- src/core/metrics.cc | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/src/core/metrics.cc b/src/core/metrics.cc index 183fcd267a7..9b97b57a6e9 100644 --- a/src/core/metrics.cc +++ b/src/core/metrics.cc @@ -450,8 +450,9 @@ foreign_ptr get_values(int handle) { shared_ptr res_ref = ::seastar::make_shared(); auto& res = *(res_ref.get()); auto& mv = res.values; - res.metadata = get_local_impl(handle)->metadata(); - auto & functions = get_local_impl(handle)->functions(); + auto impl = get_local_impl(handle); + res.metadata = impl->metadata(); + auto & functions = impl->functions(); for (auto&& i : functions) { value_vector values; for (auto&& v : i) { From 11652f60747376fa13fa6f355dfe0c76f99fae8a Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Wed, 1 Jun 2022 13:42:25 +0100 Subject: [PATCH 16/89] prometheus: support multiple metric impls This patch extends the user facing prometheus apis allowing the user to specify the internal metrics implementation to be used through a handle. Additionally, 'add_prometheus_routes' now takes an argument that specifies the route on which to advertise the metrics. This enables different metrics "namespaces" to be served by different endpoints in isolation. (cherry picked from commit 6189522fd8247894fae1bd664a7305484cd22724) --- include/seastar/core/prometheus.hh | 7 +++++-- src/core/prometheus.cc | 15 +++++++-------- 2 files changed, 12 insertions(+), 10 deletions(-) diff --git a/include/seastar/core/prometheus.hh b/include/seastar/core/prometheus.hh index 6144daa87c1..f9e225ce723 100644 --- a/include/seastar/core/prometheus.hh +++ b/include/seastar/core/prometheus.hh @@ -55,12 +55,15 @@ struct config { std::optional label; //!< A label that will be added to all metrics, we advice not to use it and set it on the prometheus server sstring prefix = "seastar"; //!< a prefix that will be added to metric names bool allow_protobuf = false; // protobuf support is experimental and off by default + int handle = metrics::default_handle(); //!< Handle that specifies which metric implementation to query + sstring route = "/metrics"; //!< Name of the route on which to expose the metrics }; future<> start(httpd::http_server_control& http_server, config ctx); -/// \defgroup add_prometheus_routes adds a /metrics endpoint that returns prometheus metrics -/// both in txt format and in protobuf according to the prometheus spec +/// \defgroup add_prometheus_routes adds a specified endpoint (defaults to /metrics) that returns prometheus metrics +/// in txt format format and in protobuf according to the prometheus spec + /// @{ future<> add_prometheus_routes(sharded& server, config ctx); future<> add_prometheus_routes(httpd::http_server& server, config ctx); diff --git a/src/core/prometheus.cc b/src/core/prometheus.cc index 31a9cb59cb1..e691f500f36 100644 --- a/src/core/prometheus.cc +++ b/src/core/prometheus.cc @@ -561,17 +561,15 @@ class metrics_families_per_shard { /** @} */ }; -static future get_map_value() { - metrics_families_per_shard vec; +static future<> get_map_value(metrics_families_per_shard& vec, int handle) { vec.resize(this_smp_shard_count()); - co_await parallel_for_each(std::views::iota(0u, this_smp_shard_count()), [&vec] (auto cpu) { - return smp::submit_to(cpu, [] { - return mi::get_values(); + co_await parallel_for_each(std::views::iota(0u, this_smp_shard_count()), [handle, &vec] (auto cpu) { + return smp::submit_to(cpu, [handle] { + return mi::get_values(handle); }).then([&vec, cpu] (auto res) { vec[cpu] = std::move(res); }); }); - co_return vec; } /*! @@ -1111,7 +1109,8 @@ class metrics_handler : public httpd::handler_base { future<> write_body(write_body_args args, output_stream&& out_stream) { auto s = std::move(out_stream); - auto families = co_await get_map_value(); + metrics_families_per_shard families; + co_await get_map_value(families, _ctx.handle); bool use_protobuf = args.use_protobuf_format; write_context context{ @@ -1138,7 +1137,7 @@ std::function metrics_handler::_true_function = [] }; future<> add_prometheus_routes(httpd::http_server& server, config ctx) { - server._routes.put(httpd::GET, "/metrics", new metrics_handler(ctx)); + server._routes.put(httpd::GET, ctx.route, new metrics_handler(ctx)); return make_ready_future<>(); } From 5eb6836691e210418bedfda4d4aecafe8c75c715 Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Wed, 1 Jun 2022 13:42:57 +0100 Subject: [PATCH 17/89] scollectd: select internal metrics implementation This patch extends the scollectd apis with the ability to select the internal metrics implementation to be used by providing a handle. (cherry picked from commit d4331d148af54213890cf1d2fe05e8df05cdd504) --- include/seastar/core/scollectd.hh | 10 +++++----- src/core/scollectd.cc | 12 ++++++------ 2 files changed, 11 insertions(+), 11 deletions(-) diff --git a/include/seastar/core/scollectd.hh b/include/seastar/core/scollectd.hh index ef59e4b6aee..f4e3e241c97 100644 --- a/include/seastar/core/scollectd.hh +++ b/include/seastar/core/scollectd.hh @@ -371,7 +371,7 @@ struct options : public program_options::option_group { /// \endcond }; -void configure(const options&); +void configure(const options&, int handle = seastar::metrics::default_handle()); void remove_polled_metric(const type_instance_id &); class plugin_instance_metrics; @@ -390,8 +390,8 @@ class plugin_instance_metrics; */ struct registration { registration() = default; - registration(const type_instance_id& id); - registration(type_instance_id&& id); + registration(const type_instance_id& id, int handle = seastar::metrics::default_handle()); + registration(type_instance_id&& id, int handle = seastar::metrics::default_handle()); registration(const registration&) = delete; registration(registration&&) = default; ~registration(); @@ -779,8 +779,8 @@ seastar::metrics::impl::metric_id to_metrics_id(const type_instance_id & id); */ template [[deprecated("Use the metrics layer")]] type_instance_id add_polled_metric(const type_instance_id & id, description d, - Arg&& arg, bool enabled = true) { - seastar::metrics::impl::get_local_impl()->add_registration(to_metrics_id(id), arg.type, seastar::metrics::impl::make_function(arg.value, arg.type), d, enabled); + Arg&& arg, bool enabled = true, int handle = seastar::metrics::default_handle()) { + seastar::metrics::impl::get_local_impl(handle)->add_registration(to_metrics_id(id), arg.type, seastar::metrics::impl::make_function(arg.value, arg.type), d, enabled); return id; } /*! diff --git a/src/core/scollectd.cc b/src/core/scollectd.cc index 7e20dccc2ab..fc5d82de8cd 100644 --- a/src/core/scollectd.cc +++ b/src/core/scollectd.cc @@ -84,12 +84,12 @@ registration::~registration() { unregister(); } -registration::registration(const type_instance_id& id) -: _id(id), _impl(seastar::metrics::impl::get_local_impl()) { +registration::registration(const type_instance_id& id, int handle) +: _id(id), _impl(seastar::metrics::impl::get_local_impl(handle)) { } -registration::registration(type_instance_id&& id) -: _id(std::move(id)), _impl(seastar::metrics::impl::get_local_impl()) { +registration::registration(type_instance_id&& id, int handle) +: _id(std::move(id)), _impl(seastar::metrics::impl::get_local_impl(handle)) { } seastar::metrics::impl::metric_id to_metrics_id(const type_instance_id & id) { @@ -542,7 +542,7 @@ future<> send_metric(const type_instance_id & id, return get_impl().send_metric(id, values); } -void configure(const options& opts) { +void configure(const options& opts, int handle) { bool enable = opts.collectd.get_value(); if (!enable) { return; @@ -551,7 +551,7 @@ void configure(const options& opts) { auto period = std::chrono::milliseconds(opts.collectd_poll_period.get_value()); auto host = (opts.collectd_hostname.get_value() == "") - ? seastar::metrics::impl::get_local_impl()->get_config().hostname + ? seastar::metrics::impl::get_local_impl(handle)->get_config().hostname : sstring(opts.collectd_hostname.get_value()); // Now create send loops on each cpu From 33b8c9884dbeeaa9c78b3dd99a4fef0dae5ec9dd Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Thu, 7 Jul 2022 11:52:21 +0100 Subject: [PATCH 18/89] metrics: Expose 'skip_when_empty' for metrics This patch adds a 'get_skip_when_empy' getter to the 'registered_metric' class. It is used by follow-up patches in order to replicate metrics. --- include/seastar/core/metrics_api.hh | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/include/seastar/core/metrics_api.hh b/include/seastar/core/metrics_api.hh index 584cc29632c..c4d66179208 100644 --- a/include/seastar/core/metrics_api.hh +++ b/include/seastar/core/metrics_api.hh @@ -255,6 +255,11 @@ public: void set_skip_when_empty(skip_when_empty skip) noexcept { _info.should_skip_when_empty = skip; } + + skip_when_empty get_skip_when_empty() const { + return _info.should_skip_when_empty; + } + const metric_id& get_id() const { return _info.id; } From 3363acbf324594f22012f35c72de83a1759a3517 Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Tue, 12 Jul 2022 12:24:50 +0100 Subject: [PATCH 19/89] metrics: add helpers for creation of replicas This patch adds private methods to the 'metrics::impl' class that deal with the creation of replicated metrics. They will be used to build the public api in future commits. --- include/seastar/core/metrics_api.hh | 9 ++++++ src/core/metrics.cc | 45 +++++++++++++++++++++++++++++ 2 files changed, 54 insertions(+) diff --git a/include/seastar/core/metrics_api.hh b/include/seastar/core/metrics_api.hh index c4d66179208..dde9083d018 100644 --- a/include/seastar/core/metrics_api.hh +++ b/include/seastar/core/metrics_api.hh @@ -463,6 +463,7 @@ class impl { std::vector _relabel_configs; std::vector _metric_family_configs; internalized_set _internalized_labels; + std::unordered_multimap _metric_families_to_replicate; public: value_map& get_value_map() { return _value_map; @@ -504,6 +505,7 @@ public: const std::vector& get_relabel_configs() const noexcept { return _relabel_configs; } + const std::vector& get_metric_family_configs() const noexcept { return _metric_family_configs; } @@ -515,6 +517,13 @@ public: private: void gc_internalized_labels(); bool apply_relabeling(const relabel_config& rc, metric_info& info); + void replicate_metric_family(const seastar::sstring& name, + int destination_handle) const; + void replicate_metric_if_required(const shared_ptr& metric) const; + void replicate_metric(const shared_ptr& metric, + const metric_family& family, + const shared_ptr& destination, + int destination_handle) const; }; const value_map& get_value_map(int handle = default_handle()); diff --git a/src/core/metrics.cc b/src/core/metrics.cc index 9b97b57a6e9..61d3da3a8c5 100644 --- a/src/core/metrics.cc +++ b/src/core/metrics.cc @@ -478,6 +478,51 @@ void impl::gc_internalized_labels() { } } +void impl::replicate_metric_family(const seastar::sstring& name, + int destination_handle) const { + const auto& entry = _value_map.find(name); + + if (entry == _value_map.end()) { + return; + } + + const auto& metric_family = entry->second; + auto destination = get_local_impl(destination_handle); + for (const auto& [labels, metric_ptr]: metric_family) { + replicate_metric(metric_ptr, metric_family, destination, destination_handle); + } +} + +void impl::replicate_metric_if_required(const shared_ptr& metric) const { + auto full_name = metric->get_id().full_name(); + auto [begin, end]= _metric_families_to_replicate.equal_range(full_name); + + for (; begin != end; ++begin) { + const auto& [name, destination_handle] = *begin; + const auto& metric_family = _value_map.at(name); + + auto destination = get_local_impl(destination_handle); + replicate_metric(metric, metric_family, destination, destination_handle); + } +} + +void impl::replicate_metric(const shared_ptr& metric, + const metric_family& family, + const shared_ptr& destination, + int destination_handle) const { + const auto& family_info = family.info(); + metric_type type = { .base_type = family_info.type, + .type_name = family_info.inherit_type }; + + destination->add_registration(metric->get_id(), + type, + metric->get_function(), + family_info.d, + metric->is_enabled(), + metric->get_skip_when_empty(), + family_info.aggregate_labels); +} + void impl::update_metrics_if_needed() { if (_dirty) { // Forcing the metadata to an empty initialization From f09fb2bde15bed4c067d6147d6d9e5eedb3dcaee Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Mon, 11 Jul 2022 12:43:58 +0100 Subject: [PATCH 20/89] metrics: allow for removal of replicated metrics This patch adds private helpers to 'metrics::impl' that deal with the removal of replicated metric families from their destintation implementation. These methods will be used in subsequent commits to manage the lifetime of replicated metrics. --- include/seastar/core/metrics_api.hh | 6 ++++++ src/core/metrics.cc | 30 +++++++++++++++++++++++++++++ 2 files changed, 36 insertions(+) diff --git a/include/seastar/core/metrics_api.hh b/include/seastar/core/metrics_api.hh index dde9083d018..aea384154b0 100644 --- a/include/seastar/core/metrics_api.hh +++ b/include/seastar/core/metrics_api.hh @@ -524,6 +524,12 @@ private: const metric_family& family, const shared_ptr& destination, int destination_handle) const; + + void remove_metric_replica_family(const seastar::sstring& name, + int destination_handle) const; + void remove_metric_replica(const metric_id& id, + const shared_ptr& destination) const; + void remove_metric_replica_if_required(const metric_id& id) const; }; const value_map& get_value_map(int handle = default_handle()); diff --git a/src/core/metrics.cc b/src/core/metrics.cc index 61d3da3a8c5..70af470767a 100644 --- a/src/core/metrics.cc +++ b/src/core/metrics.cc @@ -438,6 +438,36 @@ void impl::remove_registration(const metric_id& id) { } } +void impl::remove_metric_replica_family(const seastar::sstring& name, + int destination_handle) const { + auto entry = _value_map.find(name); + + if (entry == _value_map.end()) { + return; + } + + auto destination = get_local_impl(destination_handle); + for (const auto& metric_instance: entry->second) { + const auto& registered_metric = metric_instance.second; + remove_metric_replica(registered_metric->get_id(), + destination); + } +} + +void impl::remove_metric_replica(const metric_id& id, + const shared_ptr& destination) const { + destination->remove_registration(id); +} + +void impl::remove_metric_replica_if_required(const metric_id& id) const { + auto [begin, end] = _metric_families_to_replicate.equal_range(id.full_name()); + + for (; begin != end; ++begin) { + auto destination = get_local_impl(begin->second); + remove_metric_replica(id, destination); + } +} + void unregister_metric(const metric_id & id, int handle) { get_local_impl(handle)->remove_registration(id); } From b90e2e784899ebb9940257305324850c6e2138b5 Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Tue, 12 Jul 2022 12:27:41 +0100 Subject: [PATCH 21/89] metrics: add metric replication internal interface This patch adds a public method to the 'metrics::impl' class: 'set_metric_families_to_replicate'. When this method is called the families that match any of the specifications will be replicated on the specified destinations. --- include/seastar/core/metrics_api.hh | 16 ++++++++++++++++ src/core/metrics.cc | 16 ++++++++++++++++ 2 files changed, 32 insertions(+) diff --git a/include/seastar/core/metrics_api.hh b/include/seastar/core/metrics_api.hh index aea384154b0..54ddfd8324b 100644 --- a/include/seastar/core/metrics_api.hh +++ b/include/seastar/core/metrics_api.hh @@ -514,6 +514,22 @@ public: void update_aggregate(metric_family_info& mf) const noexcept; + // Set the metrics families to be replicated from this metrics::impl. + // All metrics families that match one of the keys of + // the 'metric_families_to_replicate' argument will be replicated + // on the metrics::impl identified by the corresponding value. + // + // If this function was called previously, any previously + // replicated metrics will be removed before the provided ones are + // replicated. + // + // Metric replication spans the full life cycle of this class. + // Newly registered metrics that belong to a replicated family + // be replicated too and unregistering a replicated metric will + // unregister the replica. + void set_metric_families_to_replicate( + std::unordered_multimap metric_families_to_replicate); + private: void gc_internalized_labels(); bool apply_relabeling(const relabel_config& rc, metric_info& info); diff --git a/src/core/metrics.cc b/src/core/metrics.cc index 70af470767a..0987e909e53 100644 --- a/src/core/metrics.cc +++ b/src/core/metrics.cc @@ -508,6 +508,22 @@ void impl::gc_internalized_labels() { } } +void +impl::set_metric_families_to_replicate( + std::unordered_multimap metric_families_to_replicate) { + // Remove all previous metric replica families + for (const auto& [name, destination]: _metric_families_to_replicate) { + remove_metric_replica_family(name, destination); + } + + // Replicate the specified metric families. + for (const auto& [name, destination]: metric_families_to_replicate) { + replicate_metric_family(name, destination); + } + + _metric_families_to_replicate = std::move(metric_families_to_replicate); +} + void impl::replicate_metric_family(const seastar::sstring& name, int destination_handle) const { const auto& entry = _value_map.find(name); From 34504f0919055bcae76dc4cd5fcf0354661f9c24 Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Mon, 11 Jul 2022 12:44:22 +0100 Subject: [PATCH 22/89] metrics: register replicated metrics dinamically This patch extends the metric registration and unregistration processes to make them aware of metric replication. In the case of metric registration, if the new metric belongs to a family that matches one of the replication specs, then a replicated metric is created accordingly. For unregistration of a metric, the replicated metric is unregistered too if one exists. --- src/core/metrics.cc | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/src/core/metrics.cc b/src/core/metrics.cc index 0987e909e53..e2e6fec9986 100644 --- a/src/core/metrics.cc +++ b/src/core/metrics.cc @@ -424,6 +424,8 @@ shared_ptr get_local_impl(int handle) { } void impl::remove_registration(const metric_id& id) { + remove_metric_replica_if_required(id); + auto i = get_value_map().find(id.full_name()); if (i != get_value_map().end()) { auto j = i->second.find(id.labels()); @@ -645,6 +647,8 @@ register_ref impl::add_registration(const metric_id& id, const metric_type& type } dirty(); + replicate_metric_if_required(rm); + return rm; } From 004df15bee074546727d2293ef2e5e79afa0d44c Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Tue, 12 Jul 2022 12:41:42 +0100 Subject: [PATCH 23/89] metrics: public family replication interface This patch exposes a method in the public interface of the metrics module ('replicate_metric_families'), which enables metric replication internally for the requested metric families. --- include/seastar/core/metrics_api.hh | 7 +++++++ src/core/metrics.cc | 11 +++++++++++ 2 files changed, 18 insertions(+) diff --git a/include/seastar/core/metrics_api.hh b/include/seastar/core/metrics_api.hh index 54ddfd8324b..2daf5457f3d 100644 --- a/include/seastar/core/metrics_api.hh +++ b/include/seastar/core/metrics_api.hh @@ -659,5 +659,12 @@ void set_metric_family_configs(const std::vector& metrics_ * This function returns a vector of the current metrics family config */ const std::vector& get_metric_family_configs(); + +/*! + * \brief replicate metric families accross internal metrics implementations + */ +future<> +replicate_metric_families(int source_handle, std::unordered_multimap metric_families_to_replicate); + } } diff --git a/src/core/metrics.cc b/src/core/metrics.cc index e2e6fec9986..743f2745080 100644 --- a/src/core/metrics.cc +++ b/src/core/metrics.cc @@ -201,6 +201,17 @@ bool impl::impl::apply_relabeling(const relabel_config& rc, metric_info& info) { return true; } +future<> +replicate_metric_families( + int source_handle, + std::unordered_multimap metric_families_to_replicate) { + return smp::invoke_on_all([source_handle, metric_families_to_replicate] { + auto source_impl = impl::get_local_impl(source_handle); + source_impl->set_metric_families_to_replicate( + std::move(metric_families_to_replicate)); + }); +} + bool label_instance::operator!=(const label_instance& id2) const { auto& id1 = *this; return !(id1 == id2); From 214d26c6c74727a36aff18df6a751492d8cbf067 Mon Sep 17 00:00:00 2001 From: Vlad Lazar Date: Mon, 11 Jul 2022 12:49:33 +0100 Subject: [PATCH 24/89] tests: add metrics replication unit tests --- tests/unit/CMakeLists.txt | 3 + tests/unit/metric_family_replication_test.cc | 96 ++++++++++++++++++++ 2 files changed, 99 insertions(+) create mode 100644 tests/unit/metric_family_replication_test.cc diff --git a/tests/unit/CMakeLists.txt b/tests/unit/CMakeLists.txt index eb451be2f7a..e51e1ceec97 100644 --- a/tests/unit/CMakeLists.txt +++ b/tests/unit/CMakeLists.txt @@ -463,6 +463,9 @@ seastar_add_test (lowres_clock seastar_add_test (metrics SOURCES metrics_test.cc) +seastar_add_test (metrics_family_replication + SOURCES metric_family_replication_test.cc) + seastar_add_test (net_config KIND BOOST SOURCES net_config_test.cc) diff --git a/tests/unit/metric_family_replication_test.cc b/tests/unit/metric_family_replication_test.cc new file mode 100644 index 00000000000..584ba13f479 --- /dev/null +++ b/tests/unit/metric_family_replication_test.cc @@ -0,0 +1,96 @@ +#include +#include +#include +#include +#include +#include + +namespace sm = seastar::metrics; +namespace smi = seastar::metrics::impl; + +bool metric_family_exists(int handle, const seastar::sstring& name) { + return smi::get_value_map(handle).contains(name); +} + +void assert_metric_families_equivalent(int source, int destination, + const seastar::sstring& name) { + const auto& source_value_map = smi::get_value_map(source); + const auto& destination_value_map = smi::get_value_map(destination); + + BOOST_REQUIRE(source_value_map.contains(name)); + BOOST_REQUIRE(destination_value_map.contains(name)); + + const auto& source_family = source_value_map.at(name); + const auto& destination_family = destination_value_map.at(name); + for (const auto& [labels, source_metric]: source_family) { + auto replica_iter = destination_family.find(labels.labels()); + BOOST_REQUIRE(replica_iter != destination_family.end()); + + const auto& replica_metric = replica_iter->second; + BOOST_REQUIRE(source_metric->get_id() == replica_metric->get_id()); + + auto source_current_value = source_metric->get_function()().i(); + auto replica_current_value = replica_metric->get_function()().i(); + BOOST_REQUIRE(source_current_value == replica_current_value); + } +} + +SEASTAR_THREAD_TEST_CASE(replicate_metrics_test) { + int foo_handle = sm::default_handle(); + sm::metric_groups foo(foo_handle); + foo.add_group("a", { + sm::make_gauge( + "gauge", + [] { return 0; })}); + + int bar_handle = sm::default_handle() + 1; + int baz_handle = sm::default_handle() + 2; + sm::replicate_metric_families(foo_handle, { + {"a_gauge", bar_handle}, + {"a_gauge", baz_handle} + }).get(); + + assert_metric_families_equivalent(foo_handle, bar_handle, "a_gauge"); + assert_metric_families_equivalent(foo_handle, baz_handle, "a_gauge"); +} + +SEASTAR_THREAD_TEST_CASE(replicate_same_metric_test) { + int foo_handle = sm::default_handle(); + + sm::metric_groups foo(foo_handle); + foo.add_group("a", { + sm::make_gauge( + "x", + [] { return 0; }, + sm::description("a_x_description"), + {sm::label("id")("1")}), + sm::make_gauge( + "x", + [] { return 0; }, + sm::description("a_x_description"), + {sm::label("id")("2")}), + sm::make_gauge( + "y", + [] { return 0; }, + sm::description("a_y_description"), + {sm::label("id")("1")}), + }); + + int bar_handle = sm::default_handle() + 1; + sm::metric_groups bar(bar_handle); + + // Test that subsequent attempts to replicate the same metric + // family are ignored. + sm::replicate_metric_families(foo_handle, {{"a_x", bar_handle}}).get(); + assert_metric_families_equivalent(foo_handle, bar_handle, "a_x"); + sm::replicate_metric_families(foo_handle, {{"a_x", bar_handle}}).get(); + assert_metric_families_equivalent(foo_handle, bar_handle, "a_x"); + + + // Ensure that when the set of replicated metric families is changed + // the replicas that are not in the new set are removed. + sm::replicate_metric_families(foo_handle, {{"a_y", bar_handle}}).get(); + assert_metric_families_equivalent(foo_handle, bar_handle, "a_y"); + + BOOST_REQUIRE(metric_family_exists(bar_handle, "a_y")); +} From 1ae3229171dc490b71b0fb4de28b3fb8c178217e Mon Sep 17 00:00:00 2001 From: Stephan Dollberg Date: Thu, 19 Oct 2023 10:30:21 +0100 Subject: [PATCH 25/89] metrics: Add update_aggregate_labels() Extends the metrics api to allow changing the aggregation labels of a metrics family. Otherwise one had to un-register every single metric instance in a metric family and then re-register with the changed aggregation labels. For metric families with thousands of instances (e.g.: histograms with lots of different labels) this is quite expensive. With this change we avoid the full reconstruction of the metrics family and all its metrics. Only the work associated with marking the metrics `dirty()` is needed then. --- include/seastar/core/metrics.hh | 7 +++++++ include/seastar/core/metrics_api.hh | 1 + include/seastar/core/metrics_registration.hh | 1 + src/core/metrics.cc | 19 +++++++++++++++++++ 4 files changed, 28 insertions(+) diff --git a/include/seastar/core/metrics.hh b/include/seastar/core/metrics.hh index 0803886bfb7..5e97e1e7889 100644 --- a/include/seastar/core/metrics.hh +++ b/include/seastar/core/metrics.hh @@ -650,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