Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions src/sample.h
Original file line number Diff line number Diff line change
Expand Up @@ -186,8 +186,10 @@ class sample {

/// Decrement ref count and reclaim if unreferenced.
friend void intrusive_ptr_release(sample *s) {
if (s->refcount_.fetch_sub(1, std::memory_order_release) == 1) {
std::atomic_thread_fence(std::memory_order_acquire);
// Acquire earlier owners' releases before publishing the sample for reuse.
// A release decrement followed by an acquire fence is also valid C++, but
// TSan does not model that fence and reports false races on recycled data.
if (s->refcount_.fetch_sub(1, std::memory_order_acq_rel) == 1) {
s->factory_->reclaim_sample(s);
}
}
Expand Down
29 changes: 29 additions & 0 deletions testing/int/samples.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -72,3 +72,32 @@ TEST_CASE("sample conversion", "[basic]") {
values[1] = (double)(-buf[0]);
}
}

TEST_CASE("sample recycling acquires earlier readers", "[sample][threads]") {
lsl::factory fac(cft_int32, 1, 1);
for (int i = 0; i < 32; ++i) {
auto sample = fac.new_sample(42.0, false);
auto *address = sample.get();
std::atomic<bool> released{false};
double observed = 0.0;
std::thread reader([copy = sample, &released, &observed]() mutable {
observed = copy->timestamp();
copy.reset();
released.store(true, std::memory_order_relaxed);
});
struct join_thread {
std::thread &thread;
~join_thread() { if (thread.joinable()) thread.join(); }
} joiner{reader};

// Only control the release order here. An acquire/release handshake or
// joining before reuse would hide synchronization missing from refcount_.
while (!released.load(std::memory_order_relaxed)) std::this_thread::yield();
sample.reset(); // Last owner must acquire the earlier reader's release.
auto recycled = fac.new_sample(84.0, true);
reader.join();
CHECK(recycled.get() == address);
CHECK(observed == 42.0);
CHECK(recycled->timestamp() == 84.0);
}
}
10 changes: 9 additions & 1 deletion testing/int/sendbuffer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
#include "send_buffer.h"
#include <atomic>
#include <catch2/catch_all.hpp>
#include <mutex>
#include <thread>
#include <vector>

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

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

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