110 std::future<void> listening = m_listening.get_future();
111 m_bgThread = std::thread(&StreamOutletPrivate::run,
this);
122 if (m_bgThread.joinable()) {
133 std::lock_guard<std::mutex> lock(m_queueMutex);
134 m_sampleQueue.push(sample);
143 std::lock_guard<std::mutex> lock(m_queueMutex);
144 for (
const auto& sample : chunk) {
145 m_sampleQueue.push(sample);
164 return m_nClients.load() > 0;
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();
189 m_listening.set_value();
192 QUdpSocket udpSocket;
194 udpSocket.setSocketOption(QAbstractSocket::MulticastTtlOption, 1);
199 QByteArray handshake(
"LSL1", 4);
201 handshake.append(
reinterpret_cast<const char*
>(&ch),
sizeof(
int));
204 std::string discoveryPayload = m_info.
to_string();
207 std::vector<QTcpSocket*> clients;
209 auto lastBroadcast = std::chrono::steady_clock::now() - std::chrono::seconds(10);
214 while (tcpServer.waitForNewConnection(0)) {
215 QTcpSocket* client = tcpServer.nextPendingConnection();
218 client->write(handshake);
220 clients.push_back(client);
221 m_nClients.store(
static_cast<int>(clients.size()));
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);
232 udpSocket.writeDatagram(datagram, QHostAddress(QHostAddress::LocalHost), DISCOVERY_PORT);
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)));
245 auto it = clients.begin();
246 while (it != clients.end()) {
247 QTcpSocket* client = *it;
248 if (client->state() != QAbstractSocket::ConnectedState) {
250 it = clients.erase(it);
253 client->write(sampleData);
262 for (QTcpSocket* client : clients) {
266 m_nClients.store(
static_cast<int>(clients.size()));
269 std::this_thread::sleep_for(std::chrono::milliseconds(1));
273 for (QTcpSocket* client : clients) {
274 client->disconnectFromHost();
285 std::atomic<bool> m_bRunning;
286 std::thread m_bgThread;
287 std::promise<void> m_listening;
289 std::mutex m_queueMutex;
290 std::queue<std::vector<float>> m_sampleQueue;
291 std::atomic<int> m_nClients{0};