Skip to content

Commit 42118f8

Browse files
authored
Merge pull request #307 from sccn/fix/sample-recycling-tsan
Make sample recycling synchronization visible to ThreadSanitizer
2 parents edeab60 + c1f1958 commit 42118f8

3 files changed

Lines changed: 42 additions & 3 deletions

File tree

‎src/sample.h‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -186,8 +186,10 @@ class sample {
186186

187187
/// Decrement ref count and reclaim if unreferenced.
188188
friend void intrusive_ptr_release(sample *s) {
189-
if (s->refcount_.fetch_sub(1, std::memory_order_release) == 1) {
190-
std::atomic_thread_fence(std::memory_order_acquire);
189+
// Acquire earlier owners' releases before publishing the sample for reuse.
190+
// A release decrement followed by an acquire fence is also valid C++, but
191+
// TSan does not model that fence and reports false races on recycled data.
192+
if (s->refcount_.fetch_sub(1, std::memory_order_acq_rel) == 1) {
191193
s->factory_->reclaim_sample(s);
192194
}
193195
}

‎testing/int/samples.cpp‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -72,3 +72,32 @@ TEST_CASE("sample conversion", "[basic]") {
7272
values[1] = (double)(-buf[0]);
7373
}
7474
}
75+
76+
TEST_CASE("sample recycling acquires earlier readers", "[sample][threads]") {
77+
lsl::factory fac(cft_int32, 1, 1);
78+
for (int i = 0; i < 32; ++i) {
79+
auto sample = fac.new_sample(42.0, false);
80+
auto *address = sample.get();
81+
std::atomic<bool> released{false};
82+
double observed = 0.0;
83+
std::thread reader([copy = sample, &released, &observed]() mutable {
84+
observed = copy->timestamp();
85+
copy.reset();
86+
released.store(true, std::memory_order_relaxed);
87+
});
88+
struct join_thread {
89+
std::thread &thread;
90+
~join_thread() { if (thread.joinable()) thread.join(); }
91+
} joiner{reader};
92+
93+
// Only control the release order here. An acquire/release handshake or
94+
// joining before reuse would hide synchronization missing from refcount_.
95+
while (!released.load(std::memory_order_relaxed)) std::this_thread::yield();
96+
sample.reset(); // Last owner must acquire the earlier reader's release.
97+
auto recycled = fac.new_sample(84.0, true);
98+
reader.join();
99+
CHECK(recycled.get() == address);
100+
CHECK(observed == 42.0);
101+
CHECK(recycled->timestamp() == 84.0);
102+
}
103+
}

‎testing/int/sendbuffer.cpp‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
#include "send_buffer.h"
3030
#include <atomic>
3131
#include <catch2/catch_all.hpp>
32+
#include <mutex>
3233
#include <thread>
3334
#include <vector>
3435

@@ -136,6 +137,7 @@ TEST_CASE("multi-threaded send_buffer stress", "[queue][regression][send_buffer]
136137
const int iterations_per_thread = 200;
137138

138139
lsl::factory fac(lsl_channel_format_t::cft_float32, 4, buffer_size * 2);
140+
std::mutex factory_mut;
139141
auto sendbuf = std::make_shared<lsl::send_buffer>(buffer_size);
140142

141143
std::vector<std::shared_ptr<lsl::consumer_queue>> queues;
@@ -148,7 +150,13 @@ TEST_CASE("multi-threaded send_buffer stress", "[queue][regression][send_buffer]
148150
// Producer threads push samples concurrently
149151
auto producer = [&]() {
150152
for (int i = 0; i < iterations_per_thread; ++i) {
151-
auto sample = fac.new_sample(static_cast<double>(i), true);
153+
lsl::sample_p sample;
154+
{
155+
// The factory permits only one allocator at a time. Keep the
156+
// send_buffer pushes concurrent without racing its sample pool.
157+
std::lock_guard<std::mutex> lock(factory_mut);
158+
sample = fac.new_sample(static_cast<double>(i), true);
159+
}
152160
sendbuf->push_sample(sample);
153161
push_count.fetch_add(1, std::memory_order_relaxed);
154162
}

0 commit comments

Comments
 (0)