Skip to content

Commit edeab60

Browse files
authored
Merge pull request #306 from sccn/fix/sync-idle-disconnect
Fix idle synchronous consumer disconnect detection
2 parents 3ddedd0 + cd7ad11 commit edeab60

2 files changed

Lines changed: 50 additions & 6 deletions

File tree

‎src/tcp_server.cpp‎

Lines changed: 29 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,9 @@
2727
#include <thread>
2828
#include <utility>
2929
#include <vector>
30+
#ifndef _WIN32
31+
#include <poll.h>
32+
#endif
3033

3134
#ifdef _MSC_VER
3235
// (inefficiently converting int to bool in portable_oarchive instantiation...)
@@ -157,12 +160,37 @@ class sync_write_handler {
157160
}
158161

159162
/// Check if there are any connected consumers
160-
bool have_consumers() const {
163+
bool have_consumers() {
161164
std::lock_guard<std::mutex> lock(mutex_);
165+
// Sync sockets have no running IO context. Check on demand, without changing
166+
// their native blocking mode (which SO_SNDTIMEO relies on).
167+
prune_disconnected(sockets_native_);
168+
prune_disconnected(sockets_swapped_);
162169
return !sockets_native_.empty() || !sockets_swapped_.empty();
163170
}
164171

165172
private:
173+
/// Called under mutex_, which also serializes socket writes and registration.
174+
static void prune_disconnected(std::vector<tcp_socket_p> &sockets) {
175+
sockets.erase(std::remove_if(sockets.begin(), sockets.end(), [](const tcp_socket_p &sock) {
176+
if (!sock || !sock->is_open()) return true;
177+
// After the handshake the feed is one-way: readability means EOF,
178+
// a reset, or unexpected client data. All end the feed, as in the
179+
// asynchronous peer-close watcher. Never issue a blocking receive.
180+
#ifdef _WIN32
181+
fd_set readable;
182+
FD_ZERO(&readable);
183+
FD_SET(sock->native_handle(), &readable);
184+
timeval timeout{};
185+
return ::select(0, &readable, nullptr, nullptr, &timeout) > 0;
186+
#else
187+
// poll also handles descriptors beyond select's FD_SETSIZE limit.
188+
pollfd fd{sock->native_handle(), POLLIN, 0};
189+
return ::poll(&fd, 1, 0) > 0 && (fd.revents & (POLLIN | POLLHUP | POLLERR));
190+
#endif
191+
}), sockets.end());
192+
}
193+
166194
/// Apply the configured blocking-send timeout to a socket (no-op if disabled). SO_SNDTIMEO
167195
/// only takes effect on a blocking socket, so we force blocking mode first. On expiry a
168196
/// synchronous write returns try_again/would_block/timed_out, which write_to_group treats

‎testing/ext/outlet.cpp‎

Lines changed: 21 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,9 @@ bool wait_until_no_consumers(lsl::stream_outlet &outlet, double timeout_sec) {
1919
TEST_CASE("have_consumers becomes false after idle inlet disconnects", "[outlet][basic]") {
2020
lsl::stream_info info("have_consumers_idle", "Markers", 1, lsl::IRREGULAR_RATE, lsl::cf_int32,
2121
"have_consumers_idle");
22-
lsl::stream_outlet outlet(info);
22+
const auto transport = GENERATE(transp_default, transp_sync_blocking);
23+
CAPTURE(transport);
24+
lsl::stream_outlet outlet(info, 0, 360, transport);
2325
auto found = lsl::resolve_stream("name", info.name(), 1, 2.0);
2426
REQUIRE(found.size() == 1);
2527

@@ -31,13 +33,23 @@ TEST_CASE("have_consumers becomes false after idle inlet disconnects", "[outlet]
3133
// Intentionally do not push samples; disconnect must still be detected (#267).
3234
}
3335

34-
CHECK(wait_until_no_consumers(outlet, 2.0));
36+
SECTION("have_consumers detects the disconnect") {
37+
CHECK(wait_until_no_consumers(outlet, 2.0));
38+
}
39+
SECTION("wait_for_consumers does not accept the stale connection") {
40+
auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(2);
41+
while (outlet.wait_for_consumers(0.0) && std::chrono::steady_clock::now() < deadline)
42+
std::this_thread::sleep_for(std::chrono::milliseconds(10));
43+
CHECK_FALSE(outlet.wait_for_consumers(0.05));
44+
}
3545
}
3646

3747
TEST_CASE("have_consumers becomes false after disconnect during push", "[outlet][basic]") {
3848
lsl::stream_info info("have_consumers_busy", "Markers", 1, lsl::IRREGULAR_RATE, lsl::cf_int32,
3949
"have_consumers_busy");
40-
lsl::stream_outlet outlet(info);
50+
const auto transport = GENERATE(transp_default, transp_sync_blocking);
51+
CAPTURE(transport);
52+
lsl::stream_outlet outlet(info, 0, 360, transport);
4153
auto found = lsl::resolve_stream("name", info.name(), 1, 2.0);
4254
REQUIRE(found.size() == 1);
4355

@@ -76,7 +88,9 @@ TEST_CASE("have_consumers becomes false after disconnect during push", "[outlet]
7688
TEST_CASE("have_consumers tracks explicit close and reopen", "[outlet][basic]") {
7789
lsl::stream_info info("have_consumers_reopen", "Markers", 1, lsl::IRREGULAR_RATE,
7890
lsl::cf_int32, "have_consumers_reopen");
79-
lsl::stream_outlet outlet(info);
91+
const auto transport = GENERATE(transp_default, transp_sync_blocking);
92+
CAPTURE(transport);
93+
lsl::stream_outlet outlet(info, 0, 360, transport);
8094
auto found = lsl::resolve_stream("name", info.name(), 1, 2.0);
8195
REQUIRE(found.size() == 1);
8296

@@ -94,7 +108,9 @@ TEST_CASE("have_consumers tracks explicit close and reopen", "[outlet][basic]")
94108
TEST_CASE("have_consumers stays true when one of two inlets disconnects", "[outlet][basic]") {
95109
lsl::stream_info info("have_consumers_two", "Markers", 1, lsl::IRREGULAR_RATE, lsl::cf_int32,
96110
"have_consumers_two");
97-
lsl::stream_outlet outlet(info);
111+
const auto transport = GENERATE(transp_default, transp_sync_blocking);
112+
CAPTURE(transport);
113+
lsl::stream_outlet outlet(info, 0, 360, transport);
98114
auto found = lsl::resolve_stream("name", info.name(), 1, 2.0);
99115
REQUIRE(found.size() == 1);
100116

0 commit comments

Comments
 (0)