OSSIA
Open Scenario System for Interactive Application
Loading...
Searching...
No Matches
qml_tcp_inbound_socket.hpp
1#pragma once
2#include <ossia/detail/variant.hpp>
3#include <ossia/network/context.hpp>
4#include <ossia/network/sockets/configuration.hpp>
5#include <ossia/network/sockets/encoding.hpp>
6#include <ossia/network/sockets/cobs_framing.hpp>
7#include <ossia/network/sockets/fixed_length_framing.hpp>
8#include <ossia/network/sockets/line_framing.hpp>
9#include <ossia/network/sockets/no_framing.hpp>
10#include <ossia/network/sockets/size_prefix_framing.hpp>
11#include <ossia/network/sockets/slip_framing.hpp>
12#include <ossia/network/sockets/stx_etx_framing.hpp>
13#include <ossia/network/sockets/tcp_socket.hpp>
14#include <ossia/network/sockets/var_size_prefix_framing.hpp>
15
16#include <ossia-qt/protocols/utils.hpp>
17
18#include <QJSValue>
19#include <QObject>
20#include <QQmlEngine>
21
22#include <nano_observer.hpp>
23
24#include <verdigris>
25
26namespace ossia::qt
27{
28class qml_tcp_connection
29 : public QObject
30 , public Nano::Observer
31{
32 W_OBJECT(qml_tcp_connection)
33public:
34 using socket_t = boost::asio::ip::tcp::socket;
35 using decoder_type = ossia::slow_variant<
36 ossia::net::no_framing::decoder<socket_t>,
37 ossia::net::slip_decoder<socket_t>,
38 ossia::net::size_prefix_decoder<socket_t>,
39 ossia::net::line_framing_decoder<socket_t>,
40 ossia::net::cobs_decoder<socket_t>,
41 ossia::net::stx_etx_framing::decoder<socket_t>,
42 ossia::net::size_prefix_1byte_framing::decoder<socket_t>,
43 ossia::net::size_prefix_2byte_be_framing::decoder<socket_t>,
44 ossia::net::size_prefix_2byte_le_framing::decoder<socket_t>,
45 ossia::net::size_prefix_4byte_le_framing::decoder<socket_t>,
46 ossia::net::fixed_length_decoder<socket_t>>;
47
48 struct state
49 {
50 boost::asio::io_context& context;
51 ossia::net::tcp_listener listener;
52 std::atomic_bool alive{true};
53 ossia::net::framing framing{ossia::net::framing::none};
54 ossia::net::encoding enc{ossia::net::encoding::none};
55 char line_delimiter[8] = {};
56 decoder_type decoder;
57
58 state(
59 ossia::net::tcp_listener l, boost::asio::io_context& ctx,
60 ossia::net::framing f = ossia::net::framing::none,
61 const std::string& delim = {},
62 ossia::net::encoding e = ossia::net::encoding::none)
63 : context{ctx}
64 , listener{std::move(l)}
65 , decoder{ossia::in_place_index<0>, listener.m_socket}
66 {
67 framing = f;
68 enc = e;
69 if(!delim.empty())
70 {
71 auto sz = std::min(delim.size(), (size_t)7);
72 std::copy_n(delim.begin(), sz, line_delimiter);
73 }
74
75 switch(f)
76 {
77 default:
78 case ossia::net::framing::none:
79 break;
80 case ossia::net::framing::slip:
81 decoder.template emplace<1>(listener.m_socket);
82 break;
83 case ossia::net::framing::size_prefix:
84 decoder.template emplace<2>(listener.m_socket);
85 break;
86 case ossia::net::framing::line_delimiter:
87 decoder.template emplace<3>(listener.m_socket);
88 {
89 auto& dec = ossia::get<3>(decoder);
90 std::copy_n(line_delimiter, 8, dec.delimiter);
91 }
92 break;
93 case ossia::net::framing::cobs:
94 decoder.template emplace<4>(listener.m_socket);
95 break;
96 case ossia::net::framing::stx_etx:
97 decoder.template emplace<5>(listener.m_socket);
98 break;
99 case ossia::net::framing::size_prefix_1byte:
100 decoder.template emplace<6>(listener.m_socket);
101 break;
102 case ossia::net::framing::size_prefix_2byte_be:
103 decoder.template emplace<7>(listener.m_socket);
104 break;
105 case ossia::net::framing::size_prefix_2byte_le:
106 decoder.template emplace<8>(listener.m_socket);
107 break;
108 case ossia::net::framing::size_prefix_4byte_le:
109 decoder.template emplace<9>(listener.m_socket);
110 break;
111 case ossia::net::framing::fixed_length:
112 decoder.template emplace<10>(listener.m_socket);
113 if(!delim.empty())
114 ossia::get<10>(decoder).frame_size = std::stoul(delim);
115 break;
116 }
117 }
118
119 void write_encoded(const char* data, std::size_t sz)
120 {
121 switch(framing)
122 {
123 default:
124 case ossia::net::framing::none:
125 listener.write(boost::asio::const_buffer(data, sz));
126 break;
127 case ossia::net::framing::slip:
128 ossia::net::slip_encoder<socket_t>{listener.m_socket}.write(data, sz);
129 break;
130 case ossia::net::framing::size_prefix:
131 ossia::net::size_prefix_encoder<socket_t>{listener.m_socket}.write(data, sz);
132 break;
133 case ossia::net::framing::line_delimiter: {
134 ossia::net::line_framing_encoder<socket_t> enc{listener.m_socket};
135 std::copy_n(line_delimiter, 8, enc.delimiter);
136 enc.write(data, sz);
137 break;
138 }
139 case ossia::net::framing::cobs:
140 ossia::net::cobs_encoder<socket_t>{listener.m_socket}.write(data, sz);
141 break;
142 case ossia::net::framing::stx_etx:
143 ossia::net::stx_etx_framing::encoder<socket_t>{listener.m_socket}.write(
144 data, sz);
145 break;
146 case ossia::net::framing::size_prefix_1byte:
147 ossia::net::size_prefix_1byte_framing::encoder<socket_t>{listener.m_socket}
148 .write(data, sz);
149 break;
150 case ossia::net::framing::size_prefix_2byte_be:
151 ossia::net::size_prefix_2byte_be_framing::encoder<socket_t>{listener.m_socket}
152 .write(data, sz);
153 break;
154 case ossia::net::framing::size_prefix_2byte_le:
155 ossia::net::size_prefix_2byte_le_framing::encoder<socket_t>{listener.m_socket}
156 .write(data, sz);
157 break;
158 case ossia::net::framing::size_prefix_4byte_le:
159 ossia::net::size_prefix_4byte_le_framing::encoder<socket_t>{listener.m_socket}
160 .write(data, sz);
161 break;
162 case ossia::net::framing::fixed_length:
163 ossia::net::fixed_length_encoder<socket_t>{listener.m_socket}.write(data, sz);
164 break;
165 }
166 }
167 };
168
169 struct receive_callback
170 {
171 std::shared_ptr<state> st;
172 QPointer<qml_tcp_connection> self;
173
174 void operator()(const unsigned char* data, std::size_t sz) const
175 {
176 if(!st->alive)
177 return;
178 auto buf = apply_decoding(st->enc, data, sz);
179 // A frame that decodes to nothing is encoding metadata, not a message:
180 // the EOF record of Intel HEX / S-record firmware streams carries no
181 // payload. Unencoded frames are passed on as they are, empty or not.
182 if(buf.isEmpty() && st->enc != ossia::net::encoding::none)
183 return;
184 ossia::qt::run_async(
185 self.get(),
186 [self = self, buf] {
187 if(!self.get())
188 return;
189 if(self->onBytes.isCallable())
190 {
191 auto engine = qjsEngine(self.get());
192 if(engine)
193 self->onBytes.call({engine->toScriptValue(buf)});
194 }
195 },
196 Qt::AutoConnection);
197 }
198
199 bool validate_stream(boost::system::error_code ec) const
200 {
201 if(ec == boost::asio::error::operation_aborted)
202 return false;
203 if(ec == boost::asio::error::eof)
204 {
205 ossia::qt::run_async(
206 self.get(),
207 [self = self] {
208 if(self)
209 {
210 if(self->onClose.isCallable())
211 self->onClose.call();
212 // Connection is dead — schedule cleanup.
213 // Parent (server) child list is updated automatically by Qt.
214 self->deleteLater();
215 }
216 },
217 Qt::AutoConnection);
218 return false;
219 }
220 return true;
221 }
222 };
223
224 explicit qml_tcp_connection(
225 ossia::net::tcp_listener listener, boost::asio::io_context& ctx,
226 ossia::net::framing f = ossia::net::framing::none,
227 const std::string& delim = {},
228 ossia::net::encoding e = ossia::net::encoding::none)
229 : m_state{std::make_shared<state>(std::move(listener), ctx, f, delim, e)}
230 {
231 }
232
233 ~qml_tcp_connection()
234 {
235 if(m_state)
236 {
237 m_state->alive = false;
238 close({});
239 }
240 }
241
242 bool isOpen() const noexcept { return m_state && m_state->alive; }
243
244 inline boost::asio::io_context& context() noexcept { return m_state->context; }
245
246 void write(QByteArray buffer)
247 {
248 if(!m_state)
249 return;
250 auto st = m_state;
251 if(st->enc != ossia::net::encoding::none)
252 buffer = apply_encoding(st->enc, buffer);
253 boost::asio::dispatch(st->context, [st, buffer = std::move(buffer)] {
254 if(st->alive)
255 st->write_encoded(buffer.data(), buffer.size());
256 });
257 }
258 W_SLOT(write)
259
260 void close(QByteArray)
261 {
262 if(!m_state)
263 return;
264 auto st = m_state;
265 boost::asio::dispatch(st->context, [st] { st->listener.close(); });
266 }
267 W_SLOT(close)
268
269 void receive(QJSValue v)
270 {
271 onBytes = v;
272 auto st = m_state;
273 auto self = QPointer{this};
274 ossia::visit(
275 [cb = receive_callback{st, self}](auto& decoder) mutable {
276 decoder.receive(std::move(cb));
277 },
278 st->decoder);
279 }
280 W_SLOT(receive)
281
282 QJSValue onBytes;
283 QJSValue onClose;
284 W_PROPERTY(QJSValue, onClose W_MEMBER onClose);
285
286private:
287 std::shared_ptr<state> m_state;
288};
289
290class qml_tcp_inbound_socket
291 : public QObject
292 , public Nano::Observer
293{
294 W_OBJECT(qml_tcp_inbound_socket)
295public:
296 struct state
297 {
298 ossia::net::tcp_server server;
299 std::atomic_bool alive{true};
300 std::atomic_bool open{false};
301 ossia::net::framing framing{ossia::net::framing::none};
302 std::string framing_delimiter;
303 ossia::net::encoding enc{ossia::net::encoding::none};
304
305 state(
306 const ossia::net::inbound_socket_configuration& conf,
307 boost::asio::io_context& ctx,
308 ossia::net::framing f = ossia::net::framing::none,
309 std::string delim = {},
310 ossia::net::encoding e = ossia::net::encoding::none)
311 : server{conf, ctx}
312 , framing{f}
313 , framing_delimiter{std::move(delim)}
314 , enc{e}
315 {
316 }
317 };
318
319 qml_tcp_inbound_socket() { }
320
321 ~qml_tcp_inbound_socket()
322 {
323 if(m_state)
324 {
325 m_state->alive = false;
326 close();
327 }
328 }
329
330 bool isOpen() const noexcept { return m_state && m_state->open; }
331
332 void open(
333 const ossia::net::inbound_socket_configuration& conf,
334 boost::asio::io_context& ctx,
335 ossia::net::framing f = ossia::net::framing::none,
336 const std::string& delim = {},
337 ossia::net::encoding e = ossia::net::encoding::none)
338 {
339 m_state = std::make_shared<state>(conf, ctx, f, delim, e);
340 m_state->open = true;
341 accept_impl(m_state, QPointer{this});
342 if(onOpen.isCallable())
343 onOpen.call({qjsEngine(this)->newQObject(this)});
344 }
345
346 void close()
347 {
348 if(!m_state)
349 return;
350 m_state->open = false;
351 m_state->server.m_acceptor.close();
352 if(onClose.isCallable())
353 onClose.call();
354 }
355 W_SLOT(close)
356
357 void on_close()
358 {
359 if(!m_state)
360 return;
361 m_state->open = false;
362 ossia::qt::run_async(this, [=, this] { onClose.call(); }, Qt::AutoConnection);
363 }
364
365 QJSValue onOpen;
366 QJSValue onClose;
367 QJSValue onError;
368 QJSValue onConnection;
369
370private:
371 static void accept_impl(
372 std::shared_ptr<state> st, QPointer<qml_tcp_inbound_socket> self)
373 {
374 st->server.m_acceptor.async_accept(
375 [self, st](
376 boost::system::error_code ec, ossia::net::tcp_server::proto::socket socket) {
377 if(!st->alive || !st->open)
378 return;
379 if(!ec)
380 {
381 ossia::qt::run_async(
382 self.get(),
383 [self, st, socket = std::move(socket)]() mutable {
384 if(!self.get())
385 return;
386 auto conn = new qml_tcp_connection{
387 ossia::net::tcp_listener{std::move(socket)}, st->server.m_context,
388 st->framing, st->framing_delimiter, st->enc};
389
390 // Parent to the server so Qt uses CppOwnership (prevents QML GC)
391 // and automatically deletes all connections when the server is destroyed.
392 conn->setParent(self.get());
393
394 if(self->onConnection.isCallable())
395 {
396 self->onConnection.call({qjsEngine(self.get())->newQObject(conn)});
397 }
398 },
399 Qt::AutoConnection);
400 accept_impl(st, self);
401 }
402 });
403 }
404
405 std::shared_ptr<state> m_state;
406};
407
408}
Definition qml_device.cpp:43
Definition git_info.h:7