Repository navigation
SDSTOR-25631: move some common headers from homestore into sisl #337
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
JacksonYao287
merged 7 commits into
eBay:dev/v14.x
from
JacksonYao287:add-common-header
Sep 15, 2026
Merged
Changes from 3 commits
Commits
Show all changes
7 commits
Select commit
Hold shift + click to select a range
76a184a
move some common headers from homestore into sisl
JacksonYao287 8bcc33a
bump conan version
JacksonYao287 dfc1f66
disable GccThreadSanitize for now
JacksonYao287 e83c73e
add UT for coro.hpp
JacksonYao287 e3d25c9
add lru_map
JacksonYao287 2ef9954
add mcmp_priority_queue
JacksonYao287 cf08343
fix
JacksonYao287 File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,64 @@ | ||
| /********************************************************************************* | ||
| * Modifications Copyright 2017-2019 eBay Inc. | ||
| * | ||
| * Licensed under the Apache License, Version 2.0 (the "License"); | ||
| * you may not use this file except in compliance with the License. | ||
| * You may obtain a copy of the License at | ||
| * https://www.apache.org/licenses/LICENSE-2.0 | ||
| * | ||
| * Unless required by applicable law or agreed to in writing, software distributed | ||
| * under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR | ||
| * CONDITIONS OF ANY KIND, either express or implied. See the License for the | ||
| * specific language governing permissions and limitations under the License. | ||
| * | ||
| *********************************************************************************/ | ||
| #pragma once | ||
|
|
||
| #include <atomic> | ||
| #include <cstddef> | ||
|
|
||
| #include <boost/lockfree/queue.hpp> | ||
|
|
||
| namespace sisl { | ||
|
|
||
| // Folly-free replacement for folly::MPMCQueue: a capacity-bounded boost::lockfree MPMC queue plus an | ||
| // approximate size counter (folly exposed this as sizeGuess()). The queue is bounded -- write() returns false | ||
| // when full -- which blkalloc relies on (the slab cache spills to the next level on a full level, and bounds | ||
| // total cached free-blocks). We use the default (pointer-freelist) boost::lockfree::queue pre-reserved to | ||
| // `capacity` and push only via bounded_push() (never grows past the reserved nodes): this gives bounded | ||
| // behavior WITHOUT boost::lockfree::fixed_sized<true>'s hard 65535-element cap (16-bit freelist indices), | ||
| // which the blkalloc free-block / slab-cache capacities exceed. boost::lockfree::queue requires a | ||
| // trivially-copyable element type (blk_num_t, blk_cache_entry both satisfy this). The size counter is the | ||
| // central accounting a bounded+queryable queue inherently needs; boost::lockfree itself keeps no size. | ||
| template < typename T > | ||
| class BoundedMPMCQueue { | ||
| public: | ||
| explicit BoundedMPMCQueue(const size_t capacity) : m_q{capacity} {} | ||
|
|
||
| // Non-blocking enqueue; returns false if the queue is full. (folly::MPMCQueue::write) | ||
| bool write(const T& value) { | ||
| if (m_q.bounded_push(value)) { | ||
| m_size.fetch_add(1, std::memory_order_relaxed); | ||
| return true; | ||
| } | ||
| return false; | ||
| } | ||
|
|
||
| // Non-blocking dequeue; returns false if the queue is empty. (folly::MPMCQueue::read) | ||
| bool read(T& out_value) { | ||
| if (m_q.pop(out_value)) { | ||
| m_size.fetch_sub(1, std::memory_order_relaxed); | ||
| return true; | ||
| } | ||
| return false; | ||
| } | ||
|
|
||
| // Approximate number of elements (racy under concurrency, like folly's). (folly::MPMCQueue::sizeGuess) | ||
| size_t sizeGuess() const { return m_size.load(std::memory_order_relaxed); } | ||
|
|
||
| private: | ||
| boost::lockfree::queue< T > m_q; | ||
| std::atomic< size_t > m_size{0}; | ||
| }; | ||
|
|
||
| } // namespace sisl |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,149 @@ | ||
| /********************************************************************************* | ||
| * Modifications Copyright 2017-2019 eBay Inc. | ||
| * | ||
| * Licensed under the Apache License, Version 2.0 (the "License"); | ||
| * you may not use this file except in compliance with the License. | ||
| * You may obtain a copy of the License at | ||
| * https://www.apache.org/licenses/LICENSE-2.0 | ||
| * | ||
| * Unless required by applicable law or agreed to in writing, software distributed | ||
| * under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR | ||
| * CONDITIONS OF ANY KIND, either express or implied. See the License for the | ||
| * specific language governing permissions and limitations under the License. | ||
| * | ||
| *********************************************************************************/ | ||
| #include <atomic> | ||
| #include <cstdint> | ||
| #include <thread> | ||
| #include <vector> | ||
|
|
||
| #include <sisl/logging/logging.h> | ||
| #include <sisl/options/options.h> | ||
|
|
||
| #include <gtest/gtest.h> | ||
|
|
||
| #include <sisl/fds/bounded_mpmc_queue.hpp> | ||
|
|
||
| using namespace sisl; | ||
|
|
||
| SISL_OPTIONS_ENABLE(logging, test_bounded_mpmc_queue) | ||
| SISL_OPTION_GROUP(test_bounded_mpmc_queue, | ||
| (num_threads, "", "num_threads", "number of producer/consumer threads", | ||
| ::cxxopts::value< uint32_t >()->default_value("4"), "number"), | ||
| (num_entries, "", "num_entries", "number of entries per producer thread", | ||
| ::cxxopts::value< uint32_t >()->default_value("5000"), "number")) | ||
|
|
||
| TEST(BoundedMPMCQueueTest, WriteAndReadSingleValue) { | ||
| BoundedMPMCQueue< int > q{4}; | ||
| EXPECT_EQ(q.sizeGuess(), 0u); | ||
|
|
||
| EXPECT_TRUE(q.write(42)); | ||
| EXPECT_EQ(q.sizeGuess(), 1u); | ||
|
|
||
| int out{0}; | ||
| EXPECT_TRUE(q.read(out)); | ||
| EXPECT_EQ(out, 42); | ||
| EXPECT_EQ(q.sizeGuess(), 0u); | ||
| } | ||
|
|
||
| TEST(BoundedMPMCQueueTest, ReadFailsWhenEmpty) { | ||
| BoundedMPMCQueue< int > q{4}; | ||
| int out{0}; | ||
| EXPECT_FALSE(q.read(out)); | ||
| } | ||
|
|
||
| TEST(BoundedMPMCQueueTest, WriteFailsWhenFull) { | ||
| constexpr size_t capacity{4}; | ||
| BoundedMPMCQueue< int > q{capacity}; | ||
|
|
||
| for (size_t i{0}; i < capacity; ++i) { | ||
| EXPECT_TRUE(q.write(static_cast< int >(i))); | ||
| } | ||
| EXPECT_EQ(q.sizeGuess(), capacity); | ||
| EXPECT_FALSE(q.write(999)); | ||
| EXPECT_EQ(q.sizeGuess(), capacity); | ||
|
|
||
| int out{0}; | ||
| EXPECT_TRUE(q.read(out)); | ||
| EXPECT_EQ(out, 0); | ||
| EXPECT_TRUE(q.write(999)); | ||
| EXPECT_EQ(q.sizeGuess(), capacity); | ||
| } | ||
|
|
||
| TEST(BoundedMPMCQueueTest, PreservesFifoOrder) { | ||
| BoundedMPMCQueue< int > q{8}; | ||
| for (int i{0}; i < 8; ++i) { | ||
| EXPECT_TRUE(q.write(i)); | ||
| } | ||
|
|
||
| for (int i{0}; i < 8; ++i) { | ||
| int out{-1}; | ||
| EXPECT_TRUE(q.read(out)); | ||
| EXPECT_EQ(out, i); | ||
| } | ||
| } | ||
|
|
||
| TEST(BoundedMPMCQueueTest, ConcurrentMultiProducerMultiConsumer) { | ||
| auto const num_threads = SISL_OPTIONS["num_threads"].as< uint32_t >(); | ||
| auto const num_entries = SISL_OPTIONS["num_entries"].as< uint32_t >(); | ||
| auto const total_entries = num_threads * num_entries; | ||
|
|
||
| BoundedMPMCQueue< uint64_t > q{16}; | ||
| std::vector< std::atomic< bool > > received(total_entries); | ||
| for (auto& r : received) { | ||
| r.store(false); | ||
| } | ||
| std::atomic< uint32_t > consumed_count{0}; | ||
|
|
||
| std::vector< std::thread > producers; | ||
| for (uint32_t t{0}; t < num_threads; ++t) { | ||
| producers.emplace_back([&q, t, num_entries]() { | ||
| for (uint32_t i{0}; i < num_entries; ++i) { | ||
| const uint64_t value{static_cast< uint64_t >(t) * num_entries + i}; | ||
| while (!q.write(value)) { | ||
| std::this_thread::yield(); | ||
| } | ||
| } | ||
| }); | ||
| } | ||
|
|
||
| std::vector< std::thread > consumers; | ||
| for (uint32_t t{0}; t < num_threads; ++t) { | ||
| consumers.emplace_back([&q, &received, &consumed_count, total_entries]() { | ||
| uint64_t value{0}; | ||
| while (consumed_count.load(std::memory_order_relaxed) < total_entries) { | ||
| if (q.read(value)) { | ||
| ASSERT_LT(value, total_entries); | ||
| ASSERT_FALSE(received[value].exchange(true)); | ||
| consumed_count.fetch_add(1, std::memory_order_relaxed); | ||
| } else { | ||
| std::this_thread::yield(); | ||
| } | ||
| } | ||
| }); | ||
| } | ||
|
|
||
| for (auto& thr : producers) { | ||
| thr.join(); | ||
| } | ||
| for (auto& thr : consumers) { | ||
| thr.join(); | ||
| } | ||
|
|
||
| EXPECT_EQ(consumed_count.load(), total_entries); | ||
| EXPECT_EQ(q.sizeGuess(), 0u); | ||
| for (const auto& r : received) { | ||
| EXPECT_TRUE(r.load()); | ||
| } | ||
| } | ||
|
|
||
| int main(int argc, char* argv[]) { | ||
| int parsed_argc{argc}; | ||
| ::testing::InitGoogleTest(&parsed_argc, argv); | ||
| SISL_OPTIONS_LOAD(parsed_argc, argv, logging, test_bounded_mpmc_queue); | ||
|
|
||
| sisl::logging::SetLogger("test_bounded_mpmc_queue"); | ||
| spdlog::set_pattern("[%D %T%z] [%^%l%$] [%n] [%t] %v"); | ||
|
|
||
| return RUN_ALL_TESTS(); | ||
| } |
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Purpose?
Uh oh!
There was an error while loading. Please reload this page.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I see all the UT hit this error
ThreadSanitizer: unexpected memory mapping
it is a known issue caused by the incompatibility between the TSan runtime and the Linux kernel’s ASLR entropy settings. pls refer to : https://bugs.launchpad.net/ubuntu/+source/linux/+bug/2056762
I disabled GccThreadSanitize to unblock github CI. or do you have any other suggestion?