v2.0.0
Loading...
Searching...
No Matches
lsl_stream_inlet.cpp
Go to the documentation of this file.
1//=============================================================================================================
27
28//=============================================================================================================
29// INCLUDES
30//=============================================================================================================
31
32#include "lsl_stream_inlet.h"
33
34//=============================================================================================================
35// QT INCLUDES
36//=============================================================================================================
37
38#include <QTcpSocket>
39#include <QByteArray>
40#include <QDebug>
41
42//=============================================================================================================
43// STL INCLUDES
44//=============================================================================================================
45
46#include <cstring>
47#include <stdexcept>
48
49//=============================================================================================================
50// USED NAMESPACES
51//=============================================================================================================
52
53using namespace LSLLIB;
54
55//=============================================================================================================
56// PRIVATE IMPLEMENTATION
57//=============================================================================================================
58
64{
65public:
67 : m_info(info)
68 , m_pSocket(nullptr)
69 , m_bIsOpen(false)
70 , m_iChannelCount(info.channel_count())
71 , m_iBytesPerSample(static_cast<int>(info.channel_count() * sizeof(float)))
72 {
73 }
74
79
80 //=========================================================================================================
85 {
86 if (m_bIsOpen) {
87 return;
88 }
89
90 // Create socket (no parent, we manage lifetime ourselves)
91 m_pSocket = new QTcpSocket();
92
93 QString host = QString::fromStdString(m_info.data_host());
94 quint16 port = static_cast<quint16>(m_info.data_port());
95
96 if (host.isEmpty() || port == 0) {
97 delete m_pSocket;
98 m_pSocket = nullptr;
99 throw std::runtime_error("[lsl::stream_inlet] Invalid data host or port in stream_info");
100 }
101
102 m_pSocket->connectToHost(host, port);
103
104 if (!m_pSocket->waitForConnected(5000)) {
105 QString err = m_pSocket->errorString();
106 delete m_pSocket;
107 m_pSocket = nullptr;
108 throw std::runtime_error(std::string("[lsl::stream_inlet] Failed to connect to outlet: ") + err.toStdString());
109 }
110
111 // Read the handshake header: "LSL1" (4 bytes) + channel_count (4 bytes, little-endian)
112 if (!m_pSocket->waitForReadyRead(5000)) {
113 delete m_pSocket;
114 m_pSocket = nullptr;
115 throw std::runtime_error("[lsl::stream_inlet] Timeout waiting for handshake from outlet");
116 }
117
118 QByteArray header;
119 while (header.size() < 8) {
120 if (m_pSocket->bytesAvailable() == 0) {
121 if (!m_pSocket->waitForReadyRead(5000)) {
122 break;
123 }
124 }
125 header.append(m_pSocket->read(8 - header.size()));
126 }
127
128 if (header.size() < 8 || header.left(4) != QByteArray("LSL1", 4)) {
129 delete m_pSocket;
130 m_pSocket = nullptr;
131 throw std::runtime_error("[lsl::stream_inlet] Invalid handshake from outlet");
132 }
133
134 // Read channel count from header (little-endian int32)
135 int headerChannels = 0;
136 std::memcpy(&headerChannels, header.constData() + 4, sizeof(int));
137 if (headerChannels != m_iChannelCount) {
138 qDebug() << "[lsl::stream_inlet] Warning: outlet reports" << headerChannels
139 << "channels, expected" << m_iChannelCount << "- using outlet value";
140 m_iChannelCount = headerChannels;
141 m_iBytesPerSample = m_iChannelCount * static_cast<int>(sizeof(float));
142 }
143
144 m_bIsOpen = true;
145 }
146
147 //=========================================================================================================
152 {
153 if (m_pSocket) {
154 if (m_pSocket->state() == QAbstractSocket::ConnectedState) {
155 m_pSocket->disconnectFromHost();
156 if (m_pSocket->state() != QAbstractSocket::UnconnectedState) {
157 m_pSocket->waitForDisconnected(1000);
158 }
159 }
160 delete m_pSocket;
161 m_pSocket = nullptr;
162 }
163 m_bIsOpen = false;
164 m_rawBuffer.clear();
165 }
166
167 //=========================================================================================================
174 {
175 if (!m_bIsOpen || !m_pSocket) {
176 return false;
177 }
178
179 // Non-blocking check for available data
180 if (m_pSocket->bytesAvailable() > 0 || m_pSocket->waitForReadyRead(0)) {
181 QByteArray data = m_pSocket->readAll();
182 m_rawBuffer.append(data);
183 }
184
185 // A complete sample requires m_iBytesPerSample bytes
186 return (m_iBytesPerSample > 0) && (m_rawBuffer.size() >= m_iBytesPerSample);
187 }
188
189 //=========================================================================================================
195 std::vector<std::vector<float>> pullChunkFloat()
196 {
197 std::vector<std::vector<float>> chunk;
198
199 if (!m_bIsOpen || !m_pSocket) {
200 return chunk;
201 }
202
203 // Read any pending data
204 readPending();
205
206 if (m_iBytesPerSample <= 0) {
207 return chunk;
208 }
209
210 // Extract complete samples from the buffer
211 int nCompleteSamples = m_rawBuffer.size() / m_iBytesPerSample;
212 if (nCompleteSamples == 0) {
213 return chunk;
214 }
215
216 chunk.reserve(nCompleteSamples);
217 const char* ptr = m_rawBuffer.constData();
218
219 for (int s = 0; s < nCompleteSamples; ++s) {
220 std::vector<float> sample(m_iChannelCount);
221 std::memcpy(sample.data(), ptr + s * m_iBytesPerSample, m_iBytesPerSample);
222 chunk.push_back(std::move(sample));
223 }
224
225 // Remove consumed bytes from the buffer
226 int consumedBytes = nCompleteSamples * m_iBytesPerSample;
227 m_rawBuffer.remove(0, consumedBytes);
228
229 return chunk;
230 }
231
233 QTcpSocket* m_pSocket;
237 QByteArray m_rawBuffer;
238};
239
240//=============================================================================================================
241// DEFINE MEMBER METHODS
242//=============================================================================================================
243
245: m_pImpl(new StreamInletPrivate(info))
246{
247}
248
249//=============================================================================================================
250
252{
253 // unique_ptr destructor handles cleanup via StreamInletPrivate destructor
254}
255
256//=============================================================================================================
257
259{
260 m_pImpl->openStream();
261}
262
263//=============================================================================================================
264
266{
267 m_pImpl->closeStream();
268}
269
270//=============================================================================================================
271
273{
274 return m_pImpl->readPending();
275}
276
277//=============================================================================================================
278
279std::vector<std::vector<float>> stream_inlet::pull_chunk_float()
280{
281 return m_pImpl->pullChunkFloat();
282}
Declares stream_inlet, the client side of an LSL connection that pulls multichannel sample chunks fro...
Lab Streaming Layer (LSL) integration for real-time data exchange.
Value-type descriptor of a single LSL stream: semantic metadata plus transport endpoint,...
StreamInletPrivate(const stream_info &info)
std::vector< std::vector< float > > pullChunkFloat()
stream_inlet(const stream_info &info)
std::vector< std::vector< float > > pull_chunk_float()