diff --git a/score/launch_manager/src/daemon/src/common/concurrency/BUILD b/score/launch_manager/src/daemon/src/common/concurrency/BUILD index 4fcc2200aa..afd99b178b 100644 --- a/score/launch_manager/src/daemon/src/common/concurrency/BUILD +++ b/score/launch_manager/src/daemon/src/common/concurrency/BUILD @@ -71,6 +71,7 @@ cc_library( "//score/launch_manager/src/daemon/src/osal:semaphore", "@score_baselibs//score/language/futurecpp", "@score_baselibs//score/mw/log", + "@score_baselibs//score/result", ], ) @@ -95,6 +96,7 @@ cc_library( deps = [ ":fixed_size_queue", "@score_baselibs//score/language/futurecpp", + "@score_baselibs//score/result", ], ) diff --git a/score/launch_manager/src/daemon/src/common/concurrency/concurrency_error_domain.hpp b/score/launch_manager/src/daemon/src/common/concurrency/concurrency_error_domain.hpp index edc5c1fcce..a6930cf1e6 100644 --- a/score/launch_manager/src/daemon/src/common/concurrency/concurrency_error_domain.hpp +++ b/score/launch_manager/src/daemon/src/common/concurrency/concurrency_error_domain.hpp @@ -1,84 +1,117 @@ -/******************************************************************************** - * Copyright (c) 2026 Contributors to the Eclipse Foundation - * - * See the NOTICE file(s) distributed with this work for additional - * information regarding copyright ownership. - * - * This program and the accompanying materials are made available under the - * terms of the Apache License Version 2.0 which is available at - * https://www.apache.org/licenses/LICENSE-2.0 - * - * SPDX-License-Identifier: Apache-2.0 - ********************************************************************************/ - -#ifndef CONCURRENCY_ERROR_DOMAIN_HPP_INCLUDED -#define CONCURRENCY_ERROR_DOMAIN_HPP_INCLUDED - -#include -#include - -namespace score::mw::lifecycle::internal -{ - -enum class ConcurrencyErrc : std::uint8_t -{ - /// @brief An OS call returned an error. - kOsError = 1, - - // @brief The container has overflowed. - kOverflow = 2, - - // @brief The container has stopped. - kStopped = 3, - - // @brief A timeout was triggered. - kTimeout = 4, -}; - -inline std::ostream& operator<<(std::ostream& os, ConcurrencyErrc errc) noexcept -{ - switch (errc) - { - case ConcurrencyErrc::kOsError: - return os << "kOsError"; - case ConcurrencyErrc::kOverflow: - return os << "kOverflow"; - case ConcurrencyErrc::kStopped: - return os << "kStopped"; - case ConcurrencyErrc::kTimeout: - return os << "kTimeout"; - default: - return os << static_cast(errc); - } -} - -} // namespace score::mw::lifecycle::internal - -#ifdef LC_LOG_SCORE_MW_LOG -#include "score/mw/log/logger.h" - -namespace score::mw::lifecycle::internal -{ - -inline score::mw::log::LogStream& operator<<(score::mw::log::LogStream& os, ConcurrencyErrc errc) noexcept -{ - switch (errc) - { - case ConcurrencyErrc::kOsError: - return os << "kOsError"; - case ConcurrencyErrc::kOverflow: - return os << "kOverflow"; - case ConcurrencyErrc::kStopped: - return os << "kStopped"; - case ConcurrencyErrc::kTimeout: - return os << "kTimeout"; - default: - return os << static_cast(errc); - } -} - -} // namespace score::mw::lifecycle::internal - -#endif // LC_LOG_SCORE_MW_LOG - -#endif // CONCURRENCY_ERROR_DOMAIN_HPP_INCLUDED +/******************************************************************************** + * Copyright (c) 2026 Contributors to the Eclipse Foundation + * + * See the NOTICE file(s) distributed with this work for additional + * information regarding copyright ownership. + * + * This program and the accompanying materials are made available under the + * terms of the Apache License Version 2.0 which is available at + * https://www.apache.org/licenses/LICENSE-2.0 + * + * SPDX-License-Identifier: Apache-2.0 + ********************************************************************************/ + +#ifndef CONCURRENCY_ERROR_DOMAIN_HPP_INCLUDED +#define CONCURRENCY_ERROR_DOMAIN_HPP_INCLUDED + +#include "score/result/error_domain.h" +#include "score/result/result.h" + +#include + +namespace score::mw::lifecycle::internal +{ + +enum class ConcurrencyErrc : score::result::ErrorCode +{ + /// @brief An OS call returned an error. + kOsError = 1, + + // @brief The container has overflowed. + kOverflow = 2, + + // @brief The container has stopped. + kStopped = 3, + + // @brief A timeout was triggered. + kTimeout = 4, +}; + +/// @brief Error domain for concurrency-related error codes. +class ConcurrencyErrorDomain final : public score::result::ErrorDomain +{ + public: + std::string_view MessageFor(const score::result::ErrorCode& code) const noexcept override + { + switch (static_cast(code)) + { + case ConcurrencyErrc::kOsError: + return "OS call returned an error"; + case ConcurrencyErrc::kOverflow: + return "Container has overflowed"; + case ConcurrencyErrc::kStopped: + return "Container has stopped"; + case ConcurrencyErrc::kTimeout: + return "Timeout was triggered"; + default: + return "Unknown concurrency error"; + } + } +}; + +/// @brief Global domain instance — required for ADL-based MakeError() lookup. +constexpr ConcurrencyErrorDomain kConcurrencyErrorDomain{}; + +/// @brief Creates a score::result::Error from a ConcurrencyErrc value (enables score::MakeUnexpected). +inline score::result::Error MakeError(ConcurrencyErrc code, std::string_view user_message = "") noexcept +{ + return {static_cast(code), kConcurrencyErrorDomain, user_message}; +} + +inline std::ostream& operator<<(std::ostream& os, ConcurrencyErrc errc) noexcept +{ + switch (errc) + { + case ConcurrencyErrc::kOsError: + return os << "kOsError"; + case ConcurrencyErrc::kOverflow: + return os << "kOverflow"; + case ConcurrencyErrc::kStopped: + return os << "kStopped"; + case ConcurrencyErrc::kTimeout: + return os << "kTimeout"; + default: + return os << static_cast(errc); + } +} + +} // namespace score::mw::lifecycle::internal + +#ifdef LC_LOG_SCORE_MW_LOG +#include "score/mw/log/logger.h" + +namespace score::mw::lifecycle::internal +{ + +inline score::mw::log::LogStream& operator<<(score::mw::log::LogStream& os, ConcurrencyErrc errc) noexcept +{ + switch (errc) + { + case ConcurrencyErrc::kOsError: + return os << "kOsError"; + case ConcurrencyErrc::kOverflow: + return os << "kOverflow"; + case ConcurrencyErrc::kStopped: + return os << "kStopped"; + case ConcurrencyErrc::kTimeout: + return os << "kTimeout"; + default: + return os << static_cast(errc); + } +} + +} // namespace score::mw::lifecycle::internal + +#endif // LC_LOG_SCORE_MW_LOG + +#endif // CONCURRENCY_ERROR_DOMAIN_HPP_INCLUDED diff --git a/score/launch_manager/src/daemon/src/common/concurrency/mpmc_concurrent_queue.hpp b/score/launch_manager/src/daemon/src/common/concurrency/mpmc_concurrent_queue.hpp index 3d29a6ac1f..854ff04529 100644 --- a/score/launch_manager/src/daemon/src/common/concurrency/mpmc_concurrent_queue.hpp +++ b/score/launch_manager/src/daemon/src/common/concurrency/mpmc_concurrent_queue.hpp @@ -28,7 +28,6 @@ #include "score/mw/launch_manager/common/concurrency/details/helgrind_annotations.hpp" #include "score/mw/launch_manager/osal/return_types.hpp" #include "score/mw/launch_manager/osal/semaphore.hpp" -#include namespace score::mw::lifecycle::internal { @@ -116,9 +115,7 @@ class MPMCConcurrentQueue /// @return Success if item was pushed, Error otherwise. /// Note: If the push returns false, the object is still valid for /// the user. - [[nodiscard]] score::cpp::expected_blank push( - T&& item, - std::chrono::milliseconds timeout = std::chrono::milliseconds{0}) + [[nodiscard]] score::ResultBlank push(T&& item, std::chrono::milliseconds timeout = std::chrono::milliseconds{0}) { return push_impl(std::move(item), timeout); } @@ -130,7 +127,7 @@ class MPMCConcurrentQueue /// previous consumer has finished reading it. /// @param timeout Maximum time to wait for a free slot. Zero means wait forever. /// @return Success if item was pushed, Error otherwise. - [[nodiscard]] score::cpp::expected_blank push( + [[nodiscard]] score::ResultBlank push( const T& item, std::chrono::milliseconds timeout = std::chrono::milliseconds{0}) { @@ -138,22 +135,22 @@ class MPMCConcurrentQueue } /// @brief Signals all blocked pop() callers to return with a stopped error. - [[nodiscard]] score::cpp::expected_blank stop() noexcept + [[nodiscard]] score::ResultBlank stop() noexcept { m_stopped.store(true, std::memory_order_release); // signal to consumers and publishers to wakeup if (m_items.post() != osal::OsalReturnType::kSuccess) { - return score::cpp::make_unexpected(ConcurrencyErrc::kOsError); + return score::MakeUnexpected(ConcurrencyErrc::kOsError); } if (m_spaces.post() != osal::OsalReturnType::kSuccess) { - return score::cpp::make_unexpected(ConcurrencyErrc::kOsError); + return score::MakeUnexpected(ConcurrencyErrc::kOsError); } - return score::cpp::blank{}; + return score::ResultBlank{}; } /// @brief Blocks until an item is available or stop() is called. @@ -162,32 +159,31 @@ class MPMCConcurrentQueue /// When stopped returns std::nullopt. /// @param timeout Maximum time to wait for an item. Zero means wait forever. /// @return The next item, or error. - [[nodiscard]] score::cpp::expected pop( - std::chrono::milliseconds timeout = std::chrono::milliseconds{0}) + [[nodiscard]] score::Result pop(std::chrono::milliseconds timeout = std::chrono::milliseconds{0}) { const auto wait_result = (timeout == std::chrono::milliseconds{0}) ? m_items.wait() : m_items.timedWait(timeout); if (wait_result == osal::OsalReturnType::kTimeout) { - return score::cpp::make_unexpected(ConcurrencyErrc::kTimeout); + return score::MakeUnexpected(ConcurrencyErrc::kTimeout); } else if (wait_result != osal::OsalReturnType::kSuccess) { - return score::cpp::make_unexpected(ConcurrencyErrc::kOsError); + return score::MakeUnexpected(ConcurrencyErrc::kOsError); } if (m_stopped.load(std::memory_order_acquire)) { static_cast(m_items.post()); - return score::cpp::make_unexpected(ConcurrencyErrc::kStopped); + return score::MakeUnexpected(ConcurrencyErrc::kStopped); } T item = consume_slot(m_head.fetch_add(1, std::memory_order_relaxed)); if (m_spaces.post() != osal::OsalReturnType::kSuccess) { - return score::cpp::make_unexpected(ConcurrencyErrc::kOsError); + return score::MakeUnexpected(ConcurrencyErrc::kOsError); } return item; @@ -226,25 +222,25 @@ class MPMCConcurrentQueue } template - [[nodiscard]] score::cpp::expected_blank push_impl(U&& item, std::chrono::milliseconds timeout) + [[nodiscard]] score::ResultBlank push_impl(U&& item, std::chrono::milliseconds timeout) { const auto wait_result = (timeout == std::chrono::milliseconds{0}) ? m_spaces.wait() : m_spaces.timedWait(timeout); if (wait_result == osal::OsalReturnType::kTimeout) { - return score::cpp::make_unexpected(ConcurrencyErrc::kTimeout); + return score::MakeUnexpected(ConcurrencyErrc::kTimeout); } else if (wait_result != osal::OsalReturnType::kSuccess) { - return score::cpp::make_unexpected(ConcurrencyErrc::kOsError); + return score::MakeUnexpected(ConcurrencyErrc::kOsError); } if (m_stopped.load(std::memory_order_acquire)) { // chain-wake the next blocked producer then discard the item static_cast(m_spaces.post()); - return score::cpp::make_unexpected(ConcurrencyErrc::kStopped); + return score::MakeUnexpected(ConcurrencyErrc::kStopped); } const auto tail = m_tail.fetch_add(1, std::memory_order_relaxed); @@ -273,10 +269,10 @@ class MPMCConcurrentQueue if (m_items.post() != osal::OsalReturnType::kSuccess) { - return score::cpp::make_unexpected(ConcurrencyErrc::kOsError); + return score::MakeUnexpected(ConcurrencyErrc::kOsError); } - return score::cpp::blank{}; + return score::ResultBlank{}; } /// @brief Underlying storage. diff --git a/score/launch_manager/src/daemon/src/common/concurrency/mpmc_concurrent_queue_test.cpp b/score/launch_manager/src/daemon/src/common/concurrency/mpmc_concurrent_queue_test.cpp index 16cf8af1b4..9b77e6400d 100644 --- a/score/launch_manager/src/daemon/src/common/concurrency/mpmc_concurrent_queue_test.cpp +++ b/score/launch_manager/src/daemon/src/common/concurrency/mpmc_concurrent_queue_test.cpp @@ -97,7 +97,7 @@ TEST_F(MPMCConcurrentQueueTest_Basic, PopReturnsNulloptOnSemaphoreWaitFailure) sa.sa_flags = 0; sigaction(SIGUSR1, &sa, nullptr); - score::cpp::expected result = score::cpp::make_unexpected(ConcurrencyErrc::kOsError); + score::Result result = score::MakeUnexpected(ConcurrencyErrc::kOsError); std::atomic tid_ready{false}; std::thread consumer([&] { @@ -194,7 +194,7 @@ class MPMCConcurrentQueueTest_Blocking : public ::testing::Test TEST_F(MPMCConcurrentQueueTest_Blocking, PopBlocksUntilItemAvailable) { RecordProperty("Description", "Verify that pop blocks on an empty queue until a producer pushes an item."); - score::cpp::expected result = score::cpp::make_unexpected(ConcurrencyErrc::kOsError); + score::Result result = score::MakeUnexpected(ConcurrencyErrc::kOsError); std::thread consumer([&] { result = queue8_.pop(); @@ -218,7 +218,7 @@ TEST_F(MPMCConcurrentQueueTest_Blocking, PushBlocksWhenFull) } std::atomic push_completed{false}; - score::cpp::expected_blank pushed = score::cpp::make_unexpected(ConcurrencyErrc::kOsError); + score::ResultBlank pushed = score::MakeUnexpected(ConcurrencyErrc::kOsError); std::thread producer([&] { pushed = queue4_.push(99); push_completed.store(true, std::memory_order_release); @@ -235,7 +235,7 @@ TEST_F(MPMCConcurrentQueueTest_Blocking, PushBlocksWhenFull) TEST_F(MPMCConcurrentQueueTest_Blocking, StopUnblocksBlockedConsumer) { RecordProperty("Description", "Verify that stop() unblocks a consumer thread waiting on an empty queue."); - score::cpp::expected result = score::cpp::make_unexpected(ConcurrencyErrc::kOsError); + score::Result result = score::MakeUnexpected(ConcurrencyErrc::kOsError); std::thread consumer([&] { result = queue8_.pop(); @@ -256,7 +256,7 @@ TEST_F(MPMCConcurrentQueueTest_Blocking, StopUnblocksBlockedProducer) ASSERT_TRUE(queue4_.push(i)); } - score::cpp::expected_blank pushed = score::cpp::make_unexpected(ConcurrencyErrc::kOsError); + score::ResultBlank pushed = score::MakeUnexpected(ConcurrencyErrc::kOsError); std::thread producer([&] { pushed = queue4_.push(99); }); diff --git a/score/launch_manager/src/daemon/src/common/concurrency/mpsc_bounded_queue.hpp b/score/launch_manager/src/daemon/src/common/concurrency/mpsc_bounded_queue.hpp index ef93f63c99..224c21d04a 100644 --- a/score/launch_manager/src/daemon/src/common/concurrency/mpsc_bounded_queue.hpp +++ b/score/launch_manager/src/daemon/src/common/concurrency/mpsc_bounded_queue.hpp @@ -27,7 +27,6 @@ #include "score/mw/launch_manager/common/concurrency/fixed_size_queue.hpp" #include -#include namespace score::mw::lifecycle::internal { @@ -68,12 +67,12 @@ class MpscBoundedQueue /// @brief Enqueues an item. /// @return blank on success; ConcurrencyErrc::kOverflow if the queue is full, or /// ConcurrencyErrc::kStopped if the queue has been stopped. - [[nodiscard]] score::cpp::expected_blank push(T&& item) + [[nodiscard]] score::ResultBlank push(T&& item) { return push_impl(std::move(item)); } - [[nodiscard]] score::cpp::expected_blank push(const T& item) + [[nodiscard]] score::ResultBlank push(const T& item) { return push_impl(item); } @@ -83,7 +82,7 @@ class MpscBoundedQueue /// @return blank if an item is available (caller should drain via tryPop()); /// ConcurrencyErrc::kTimeout if the timeout elapsed with none available, or /// ConcurrencyErrc::kStopped if the queue has been stopped. - [[nodiscard]] score::cpp::expected_blank wait(std::chrono::milliseconds timeout) + [[nodiscard]] score::ResultBlank wait(std::chrono::milliseconds timeout) { std::unique_lock lock(mutex_); SCORE_LANGUAGE_FUTURECPP_ASSERT_MESSAGE(ensure_single_consumer(), "Only a single consumer thread is allowed."); @@ -94,13 +93,13 @@ class MpscBoundedQueue if (stopped_) { - return score::cpp::make_unexpected(ConcurrencyErrc::kStopped); + return score::MakeUnexpected(ConcurrencyErrc::kStopped); } if (!has_item) { - return score::cpp::make_unexpected(ConcurrencyErrc::kTimeout); + return score::MakeUnexpected(ConcurrencyErrc::kTimeout); } - return {}; + return score::ResultBlank{}; } /// @brief Non-blocking pop. Never waits, regardless of whether the queue has been stopped, so @@ -146,23 +145,23 @@ class MpscBoundedQueue } template - [[nodiscard]] score::cpp::expected_blank push_impl(U&& item) + [[nodiscard]] score::ResultBlank push_impl(U&& item) { std::unique_lock lock(mutex_); if (stopped_) { - return score::cpp::make_unexpected(ConcurrencyErrc::kStopped); + return score::MakeUnexpected(ConcurrencyErrc::kStopped); } if (!queue_.push(std::forward(item))) { - return score::cpp::make_unexpected(ConcurrencyErrc::kOverflow); + return score::MakeUnexpected(ConcurrencyErrc::kOverflow); } lock.unlock(); not_empty_cv_.notify_one(); - return {}; + return score::ResultBlank{}; } mutable std::mutex mutex_{};