v2.0.0
Loading...
Searching...
No Matches
lsl_stream_outlet.cpp
Go to the documentation of this file.
1//=============================================================================================================
30
31//=============================================================================================================
32// INCLUDES
33//=============================================================================================================
34
35#include "lsl_stream_outlet.h"
36
37//=============================================================================================================
38// QT INCLUDES
39//=============================================================================================================
40
41#include <QTcpServer>
42#include <QTcpSocket>
43#include <QUdpSocket>
44#include <QHostAddress>
45#include <QDebug>
46
47//=============================================================================================================
48// STL INCLUDES
49//=============================================================================================================
50
51#include <thread>
52#include <mutex>
53#include <atomic>
54#include <cstring>
55#include <queue>
56#include <future>
57#include <chrono>
58
59//=============================================================================================================
60// USED NAMESPACES
61//=============================================================================================================
62
63using namespace LSLLIB;
64
65//=============================================================================================================
66// CONSTANTS
67//=============================================================================================================
68
69namespace
70{
71const QHostAddress DISCOVERY_MULTICAST_GROUP("239.255.172.215");
72const quint16 DISCOVERY_PORT = 16571;
73const int BROADCAST_INTERVAL_MS = 500;
74}
75
76//=============================================================================================================
77// PRIVATE IMPLEMENTATION
78//=============================================================================================================
79
90{
91public:
93 : m_info(info)
94 , m_bRunning(false)
95 {
96 }
97
99 {
100 stop();
101 }
102
103 //=========================================================================================================
107 void start()
108 {
109 m_bRunning = true;
110 std::future<void> listening = m_listening.get_future();
111 m_bgThread = std::thread(&StreamOutletPrivate::run, this);
112 listening.wait(); // m_info (with its data port) is not modified after this point
113 }
114
115 //=========================================================================================================
119 void stop()
120 {
121 m_bRunning = false;
122 if (m_bgThread.joinable()) {
123 m_bgThread.join();
124 }
125 }
126
127 //=========================================================================================================
131 void enqueueSample(const std::vector<float>& sample)
132 {
133 std::lock_guard<std::mutex> lock(m_queueMutex);
134 m_sampleQueue.push(sample);
135 }
136
137 //=========================================================================================================
141 void enqueueChunk(const std::vector<std::vector<float>>& chunk)
142 {
143 std::lock_guard<std::mutex> lock(m_queueMutex);
144 for (const auto& sample : chunk) {
145 m_sampleQueue.push(sample);
146 }
147 }
148
149 //=========================================================================================================
154 {
155 return m_info;
156 }
157
158 //=========================================================================================================
162 bool haveConsumers() const
163 {
164 return m_nClients.load() > 0;
165 }
166
167private:
168 //=========================================================================================================
177 void run()
178 {
179 // --- Set up TCP server ---
180 QTcpServer tcpServer;
181 if (!tcpServer.listen(QHostAddress::Any, 0)) {
182 qDebug() << "[lsl::stream_outlet] Failed to start TCP server:" << tcpServer.errorString();
183 m_listening.set_value();
184 return;
185 }
186
187 // Record the assigned port into stream_info
188 m_info.set_data_port(tcpServer.serverPort());
189 m_listening.set_value();
190
191 // --- Set up UDP socket for multicast discovery ---
192 QUdpSocket udpSocket;
193#ifndef Q_OS_WASM
194 udpSocket.setSocketOption(QAbstractSocket::MulticastTtlOption, 1);
195#endif
196
197 // --- Prepare the handshake header ---
198 // "LSL1" (4 bytes) + channel_count (4 bytes, little-endian int)
199 QByteArray handshake("LSL1", 4);
200 int ch = m_info.channel_count();
201 handshake.append(reinterpret_cast<const char*>(&ch), sizeof(int));
202
203 // Prepare discovery datagram (re-created each loop to include up-to-date port)
204 std::string discoveryPayload = m_info.to_string();
205
206 // Track connected client sockets (owned by this thread)
207 std::vector<QTcpSocket*> clients;
208
209 auto lastBroadcast = std::chrono::steady_clock::now() - std::chrono::seconds(10); // send immediately
210
211 // --- Main loop ---
212 while (m_bRunning) {
213 // 1. Accept new connections
214 while (tcpServer.waitForNewConnection(0)) {
215 QTcpSocket* client = tcpServer.nextPendingConnection();
216 if (client) {
217 // Send handshake to the new client
218 client->write(handshake);
219 client->flush();
220 clients.push_back(client);
221 m_nClients.store(static_cast<int>(clients.size()));
222 }
223 }
224
225 // 2. Send UDP discovery broadcast (every BROADCAST_INTERVAL_MS)
226 auto now = std::chrono::steady_clock::now();
227 auto elapsed = std::chrono::duration_cast<std::chrono::milliseconds>(now - lastBroadcast).count();
228 if (elapsed >= BROADCAST_INTERVAL_MS) {
229 QByteArray datagram(discoveryPayload.c_str(), static_cast<int>(discoveryPayload.size()));
230 udpSocket.writeDatagram(datagram, DISCOVERY_MULTICAST_GROUP, DISCOVERY_PORT);
231 // Same-host listeners still find the stream when multicast is unrouted (VPNs, CI VMs).
232 udpSocket.writeDatagram(datagram, QHostAddress(QHostAddress::LocalHost), DISCOVERY_PORT);
233 lastBroadcast = now;
234 }
235
236 // 3. Drain sample queue and write to all clients
237 {
238 std::lock_guard<std::mutex> lock(m_queueMutex);
239 while (!m_sampleQueue.empty()) {
240 const std::vector<float>& sample = m_sampleQueue.front();
241 QByteArray sampleData(reinterpret_cast<const char*>(sample.data()),
242 static_cast<int>(sample.size() * sizeof(float)));
243
244 // Write to all clients, remove disconnected ones
245 auto it = clients.begin();
246 while (it != clients.end()) {
247 QTcpSocket* client = *it;
248 if (client->state() != QAbstractSocket::ConnectedState) {
249 delete client;
250 it = clients.erase(it);
251 continue;
252 }
253 client->write(sampleData);
254 ++it;
255 }
256
257 m_sampleQueue.pop();
258 }
259 }
260
261 // Flush all clients
262 for (QTcpSocket* client : clients) {
263 client->flush();
264 }
265
266 m_nClients.store(static_cast<int>(clients.size()));
267
268 // Sleep briefly to avoid busy-waiting
269 std::this_thread::sleep_for(std::chrono::milliseconds(1));
270 }
271
272 // --- Cleanup ---
273 for (QTcpSocket* client : clients) {
274 client->disconnectFromHost();
275 delete client;
276 }
277 clients.clear();
278 m_nClients.store(0);
279
280 tcpServer.close();
281 udpSocket.close();
282 }
283
284 stream_info m_info;
285 std::atomic<bool> m_bRunning;
286 std::thread m_bgThread;
287 std::promise<void> m_listening;
288
289 std::mutex m_queueMutex;
290 std::queue<std::vector<float>> m_sampleQueue;
291 std::atomic<int> m_nClients{0};
292};
293
294//=============================================================================================================
295// DEFINE MEMBER METHODS
296//=============================================================================================================
297
299: m_pImpl(new StreamOutletPrivate(info))
300{
301 m_pImpl->start();
302}
303
304//=============================================================================================================
305
307{
308 // unique_ptr destructor handles cleanup via StreamOutletPrivate destructor
309}
310
311//=============================================================================================================
312
313void stream_outlet::push_sample(const std::vector<float>& sample)
314{
315 m_pImpl->enqueueSample(sample);
316}
317
318//=============================================================================================================
319
320void stream_outlet::push_chunk(const std::vector<std::vector<float>>& chunk)
321{
322 m_pImpl->enqueueChunk(chunk);
323}
324
325//=============================================================================================================
326
328{
329 return m_pImpl->info();
330}
331
332//=============================================================================================================
333
335{
336 return m_pImpl->haveConsumers();
337}
Declares stream_outlet, the server side of an LSL stream that publishes samples over TCP and advertis...
Lab Streaming Layer (LSL) integration for real-time data exchange.
Value-type descriptor of a single LSL stream: semantic metadata plus transport endpoint,...
int channel_count() const noexcept
Number of channels.
std::string to_string() const
Serialize stream_info into a string for network transport.
void set_data_port(int port)
Set the TCP data port (used internally during discovery / outlet creation).
StreamOutletPrivate(const stream_info &info)
void enqueueSample(const std::vector< float > &sample)
void enqueueChunk(const std::vector< std::vector< float > > &chunk)
void push_sample(const std::vector< float > &sample)
stream_outlet(const stream_info &info)
void push_chunk(const std::vector< std::vector< float > > &chunk)
stream_info info() const