Skip to content

Commit 3ddedd0

Browse files
authored
Merge pull request #291 from daverlon/dev
Fix have_consumers staying true after an idle consumer disconnects (#267)
2 parents 7f3d31a + 4a3e886 commit 3ddedd0

5 files changed

Lines changed: 180 additions & 3 deletions

File tree

‎src/send_buffer.cpp‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,16 @@ void send_buffer::push_sample(const sample_p &s) {
2323
}
2424

2525

26+
/**
27+
* Wake a single consumer without enqueueing real data.
28+
* Shares consumers_mut_ with push_sample() so the queue keeps its single producer.
29+
*/
30+
void send_buffer::wake_consumer(const std::shared_ptr<consumer_queue> &q) {
31+
std::lock_guard<std::mutex> lock(consumers_mut_);
32+
q->push_sample(sample_p());
33+
}
34+
35+
2636
/// Registered a new consumer.
2737
void send_buffer::register_consumer(consumer_queue *q) {
2838
{

‎src/send_buffer.h‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,15 @@ class send_buffer : public std::enable_shared_from_this<send_buffer> {
4545
/// Push a sample onto the send buffer that will subsequently be received by all consumers.
4646
void push_sample(const sample_p &s);
4747

48+
/**
49+
* Wake a single consumer's queue with the empty-sample convention.
50+
*
51+
* Takes the same lock as push_sample() so that callers on other threads (e.g. an IO handler
52+
* reacting to a peer close) do not become a second producer on a single-producer queue.
53+
* @param q The consumer to wake; held by the caller for the duration of the call.
54+
*/
55+
void wake_consumer(const std::shared_ptr<consumer_queue> &q);
56+
4857
/// Wait until some consumers are present.
4958
bool wait_for_consumers(double timeout = FOREVER);
5059

‎src/tcp_server.cpp‎

Lines changed: 37 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
#include <asio/streambuf.hpp>
1717
#include <asio/write.hpp>
1818
#include <algorithm>
19+
#include <atomic>
1920
#include <condition_variable>
2021
#include <cstdint>
2122
#include <cstring>
@@ -268,8 +269,11 @@ class client_session : public std::enable_shared_from_this<client_session> {
268269
void handle_send_feedheader_outcome(
269270
err_t err, std::size_t n, std::shared_ptr<consumer_queue> queue);
270271

272+
/// Detect peer close on the feed socket so idle consumers get unregistered (see #267).
273+
void abort_on_peer_close(const std::shared_ptr<consumer_queue> &queue);
274+
271275
/// Transfers samples from the server's send buffer into the async send queues of IO threads
272-
void transfer_samples_thread(std::shared_ptr<client_session> /*keepalive*/,
276+
void transfer_samples_thread(std::shared_ptr<client_session> keepalive,
273277
std::shared_ptr<consumer_queue> &&queue, int max_samples_per_chunk);
274278

275279
/// Handler that gets called when a sample transfer has been completed.
@@ -302,6 +306,10 @@ class client_session : public std::enable_shared_from_this<client_session> {
302306
int chunk_granularity_{0};
303307
/// maximum number of samples buffered
304308
int max_buffered_{0};
309+
/// set when the consumer closes the feed connection (or sends unexpected data)
310+
std::atomic<bool> abort_transfer_{false};
311+
/// dummy byte for the peer-close detector's pending read
312+
char peer_close_byte_{0};
305313

306314
// data exchanged between the transfer completion handler and the transfer thread
307315
/// whether the current transfer has finished (possibly with an error)
@@ -761,6 +769,8 @@ void client_session::handle_send_feedheader_outcome(
761769
// convenient for unit tests
762770
if (max_buffered_ <= 0) return;
763771

772+
abort_on_peer_close(queue);
773+
764774
// determine the maximum chunk size
765775
int max_samples_per_chunk = std::numeric_limits<int>::max();
766776
if (chunk_granularity_)
@@ -777,13 +787,31 @@ void client_session::handle_send_feedheader_outcome(
777787
}
778788
}
779789

780-
void client_session::transfer_samples_thread(std::shared_ptr<client_session> /* keepalive */,
790+
void client_session::abort_on_peer_close(const std::shared_ptr<consumer_queue> &queue) {
791+
// The feed is one-way after the handshake; a completed read means EOF or unexpected data.
792+
// Either way the consumer is gone. A weak_ptr keeps the queue from staying registered if the
793+
// transfer thread has already exited (e.g. after a failed write).
794+
std::weak_ptr<consumer_queue> weak_queue = queue;
795+
sock_.async_read_some(asio::buffer(&peer_close_byte_, 1),
796+
[shared_this = shared_from_this(), weak_queue](err_t err, std::size_t /*n*/) {
797+
if (err == asio::error::operation_aborted) return;
798+
shared_this->abort_transfer_.store(true, std::memory_order_release);
799+
// Wake through the send buffer rather than pushing here: this runs on an IO thread
800+
// and would otherwise be a second producer on the outlet thread's queue.
801+
auto q = weak_queue.lock();
802+
auto serv = shared_this->serv_.lock();
803+
if (q && serv) serv->send_buffer_->wake_consumer(q);
804+
});
805+
}
806+
807+
void client_session::transfer_samples_thread(std::shared_ptr<client_session> keepalive,
781808
std::shared_ptr<consumer_queue> &&queue, int max_samples_per_chunk) {
782809
int samples_in_current_chunk = 0;
783-
while (!serv_.expired()) {
810+
while (!serv_.expired() && !abort_transfer_.load(std::memory_order_relaxed)) {
784811
try {
785812
// get next sample from the sample queue (blocking)
786813
sample_p samp(queue->pop_sample());
814+
if (abort_transfer_.load(std::memory_order_acquire)) break;
787815

788816
// ignore blank samples (they are basically wakeup notifiers from someone's
789817
// end_serving())
@@ -816,6 +844,12 @@ void client_session::transfer_samples_thread(std::shared_ptr<client_session> /*
816844
LOG_F(WARNING, "Unexpected glitch in transfer_samples_thread: %s", e.what());
817845
}
818846
}
847+
// Unregister immediately; cancel the peer-close read so the session can be destroyed.
848+
queue.reset();
849+
post(*io_, [keepalive]() {
850+
asio::error_code ec;
851+
keepalive->sock_.cancel(ec);
852+
});
819853
}
820854

821855
void client_session::handle_chunk_transfer_outcome(err_t err, std::size_t len) {

‎testing/CMakeLists.txt‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,7 @@ add_executable(lsl_test_exported
7373
ext/inlet_open.cpp
7474
ext/inlet_reopen.cpp
7575
ext/move.cpp
76+
ext/outlet.cpp
7677
ext/streaminfo.cpp
7778
ext/sync_outlet.cpp
7879
ext/time.cpp

‎testing/ext/outlet.cpp‎

Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,123 @@
1+
#include <atomic>
2+
#include <catch2/catch_all.hpp>
3+
#include <chrono>
4+
#include <lsl_cpp.h>
5+
#include <thread>
6+
7+
// clazy:excludeall=non-pod-global-static
8+
9+
namespace {
10+
11+
bool wait_until_no_consumers(lsl::stream_outlet &outlet, double timeout_sec) {
12+
auto deadline =
13+
std::chrono::steady_clock::now() + std::chrono::duration<double>(timeout_sec);
14+
while (outlet.have_consumers() && std::chrono::steady_clock::now() < deadline)
15+
std::this_thread::sleep_for(std::chrono::milliseconds(10));
16+
return !outlet.have_consumers();
17+
}
18+
19+
TEST_CASE("have_consumers becomes false after idle inlet disconnects", "[outlet][basic]") {
20+
lsl::stream_info info("have_consumers_idle", "Markers", 1, lsl::IRREGULAR_RATE, lsl::cf_int32,
21+
"have_consumers_idle");
22+
lsl::stream_outlet outlet(info);
23+
auto found = lsl::resolve_stream("name", info.name(), 1, 2.0);
24+
REQUIRE(found.size() == 1);
25+
26+
{
27+
lsl::stream_inlet inlet(found[0]);
28+
inlet.open_stream(2);
29+
REQUIRE(outlet.wait_for_consumers(2));
30+
REQUIRE(outlet.have_consumers());
31+
// Intentionally do not push samples; disconnect must still be detected (#267).
32+
}
33+
34+
CHECK(wait_until_no_consumers(outlet, 2.0));
35+
}
36+
37+
TEST_CASE("have_consumers becomes false after disconnect during push", "[outlet][basic]") {
38+
lsl::stream_info info("have_consumers_busy", "Markers", 1, lsl::IRREGULAR_RATE, lsl::cf_int32,
39+
"have_consumers_busy");
40+
lsl::stream_outlet outlet(info);
41+
auto found = lsl::resolve_stream("name", info.name(), 1, 2.0);
42+
REQUIRE(found.size() == 1);
43+
44+
std::atomic<bool> keep_pushing{true};
45+
std::thread pusher([&outlet, &keep_pushing]() {
46+
int32_t value = 0;
47+
while (keep_pushing.load(std::memory_order_relaxed)) {
48+
outlet.push_sample(&value);
49+
++value;
50+
}
51+
});
52+
// Stop and join even if open_stream() throws or a REQUIRE aborts the test.
53+
struct stop_and_join {
54+
std::atomic<bool> &keep_pushing;
55+
std::thread &pusher;
56+
~stop_and_join() {
57+
keep_pushing.store(false, std::memory_order_relaxed);
58+
pusher.join();
59+
}
60+
} pusher_cleanup{keep_pushing, pusher};
61+
62+
{
63+
lsl::stream_inlet inlet(found[0]);
64+
inlet.open_stream(2);
65+
REQUIRE(outlet.wait_for_consumers(2));
66+
REQUIRE(outlet.have_consumers());
67+
int32_t received = 0;
68+
CHECK(inlet.pull_sample(&received, 1, 1.0) != 0.0);
69+
// Destroy the inlet while the outlet thread is still producing.
70+
std::this_thread::sleep_for(std::chrono::milliseconds(20));
71+
}
72+
73+
CHECK(wait_until_no_consumers(outlet, 2.0));
74+
}
75+
76+
TEST_CASE("have_consumers tracks explicit close and reopen", "[outlet][basic]") {
77+
lsl::stream_info info("have_consumers_reopen", "Markers", 1, lsl::IRREGULAR_RATE,
78+
lsl::cf_int32, "have_consumers_reopen");
79+
lsl::stream_outlet outlet(info);
80+
auto found = lsl::resolve_stream("name", info.name(), 1, 2.0);
81+
REQUIRE(found.size() == 1);
82+
83+
lsl::stream_inlet inlet(found[0]);
84+
for (int cycle = 0; cycle < 3; ++cycle) {
85+
CAPTURE(cycle);
86+
inlet.open_stream(2);
87+
REQUIRE(outlet.wait_for_consumers(2));
88+
// close_stream() instead of destroying the inlet; the outlet must notice either way.
89+
inlet.close_stream();
90+
REQUIRE(wait_until_no_consumers(outlet, 2.0));
91+
}
92+
}
93+
94+
TEST_CASE("have_consumers stays true when one of two inlets disconnects", "[outlet][basic]") {
95+
lsl::stream_info info("have_consumers_two", "Markers", 1, lsl::IRREGULAR_RATE, lsl::cf_int32,
96+
"have_consumers_two");
97+
lsl::stream_outlet outlet(info);
98+
auto found = lsl::resolve_stream("name", info.name(), 1, 2.0);
99+
REQUIRE(found.size() == 1);
100+
101+
lsl::stream_inlet keeper(found[0]);
102+
keeper.open_stream(2);
103+
REQUIRE(outlet.wait_for_consumers(2));
104+
105+
{
106+
lsl::stream_inlet leaver(found[0]);
107+
leaver.open_stream(2);
108+
REQUIRE(outlet.have_consumers());
109+
}
110+
std::this_thread::sleep_for(std::chrono::milliseconds(500));
111+
112+
// Unregistering one consumer must leave the other registered and still fed.
113+
CHECK(outlet.have_consumers());
114+
int32_t sent = 42, received = 0;
115+
outlet.push_sample(&sent);
116+
CHECK(keeper.pull_sample(&received, 1, 2.0) != 0.0);
117+
CHECK(received == sent);
118+
119+
keeper.close_stream();
120+
CHECK(wait_until_no_consumers(outlet, 2.0));
121+
}
122+
123+
} // namespace

0 commit comments

Comments
 (0)