OpenCPN Partial API docs
Loading...
Searching...
No Matches
comm_drv_signalk_net.cpp
Go to the documentation of this file.
1/***************************************************************************
2 * Copyright (C) 2022 by David Register *
3 * Copyright (C) 2022 Alec Leamas *
4 * *
5 * This program is free software; you can redistribute it and/or modify *
6 * it under the terms of the GNU General Public License as published by *
7 * the Free Software Foundation; either version 2 of the License, or *
8 * (at your option) any later version. *
9 * *
10 * This program is distributed in the hope that it will be useful, *
11 * but WITHOUT ANY WARRANTY; without even the implied warranty of *
12 * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the *
13 * GNU General Public License for more details. *
14 * *
15 * You should have received a copy of the GNU General Public License *
16 * along with this program; if not, see <https://www.gnu.org/licenses/>. *
17 **************************************************************************/
18
25#include <chrono>
26#include <mutex> // std::mutex
27
28#include <wx/socket.h>
29
30#include "observable/observable.h"
31
32#include "rapidjson/document.h"
33#include "ixwebsocket/IXNetSystem.h"
34#include "ixwebsocket/IXSocketTLSOptions.h"
35#include "ixwebsocket/IXWebSocket.h"
36
39#include "model/geodesic.h"
40#include "model/logger.h"
41#include "model/sys_events.h"
42#include "model/thread_ctrl.h"
43#include "wxServDisc.h"
44
45using namespace std::literals::chrono_literals;
46
47constexpr int kTimerSocket = 9006;
48constexpr int kSignalkSocketId = 5011;
49constexpr int kDogTimeoutReconnectSeconds = 10;
50
51constexpr double kMsToKnotFactor = 1.9438444924406;
52
54public:
55 IoThread(const std::string& iface, const wxIPV4address& address,
56 wxEvtHandler* consumer, const std::string& token);
57
58 ~IoThread() override = default;
59
60 void Run();
61
62 DriverStats GetStats() const;
63
64private:
65 wxIPV4address m_address;
66 wxEvtHandler* m_consumer;
67 const std::string m_iface;
68 std::string m_token;
69 ix::WebSocket m_ws;
70 obs::Listener m_resume_listener;
71 DriverStats m_driver_stats;
72 mutable std::mutex m_stats_mutex;
73};
74
75// i. e. wxDEFINE_EVENT(), avoiding the evil macro.
76static const wxEventTypeTag<CommDriverSignalKNet::InputEvt> SignalkEvtType(
77 wxNewEventType());
78
79static wxIPV4address ParamsIpAddress(const ConnectionParams& params) {
80 wxIPV4address addr;
81 addr.Hostname(params.NetworkAddress);
82 addr.Service(params.NetworkPort);
83 return addr;
84}
85
86class CommDriverSignalKNet::InputEvt : public wxEvent {
87public:
88 explicit InputEvt(std::string payload)
89 : wxEvent(0, SignalkEvtType), m_payload(std::move(payload)) {};
90
91 std::string GetPayload() const { return m_payload; }
92
93 // required for sending with wxPostEvent()
94 wxEvent* Clone() const override { return new InputEvt(m_payload); };
95
96private:
97 const std::string m_payload;
98};
99
100//========================================================================
101// IoThread implementation
102//
103CommDriverSignalKNet::IoThread::IoThread(const std::string& iface,
104 const wxIPV4address& address,
105 wxEvtHandler* consumer,
106 const std::string& token)
107 : m_address(address), m_consumer(consumer), m_iface(iface), m_token(token) {
108 m_resume_listener.Init(SystemEvents::GetInstance().evt_resume,
109 [&](ObservedEvt& ev) {
110 m_ws.stop();
111 m_ws.start();
112 wxLogDebug("WebSocketThread: restarted");
113 });
114}
115
116void CommDriverSignalKNet::IoThread::Run() {
117 using namespace std::chrono_literals;
118
119 {
120 std::lock_guard lock(m_stats_mutex);
121 m_driver_stats.driver_bus = NavAddr::Bus::Signalk;
122 m_driver_stats.driver_iface = m_iface;
123 m_driver_stats.available = false;
124 }
125
126 // Craft the address strings
127 wxString host = m_address.IPAddress();
128 int port = m_address.Service();
129 std::stringstream wsAddress;
130 wsAddress << "ws://" << host << ":" << port
131 << "/signalk/v1/stream?subscribe=all&sendCachedValues=false";
132 if (!m_token.empty()) wsAddress << "&token=" << m_token;
133 std::stringstream wssAddress;
134 wssAddress << "wss://" << host << ":" << port
135 << "/signalk/v1/stream?subscribe=all&sendCachedValues=false";
136 if (!m_token.empty()) wssAddress << "&token=" << m_token;
137
138 m_ws.setUrl(wssAddress.str());
139 ix::SocketTLSOptions opt;
140 opt.disable_hostname_validation = true;
141 opt.caFile = "NONE";
142 m_ws.setTLSOptions(opt);
143 m_ws.setPingInterval(30);
144
145 auto message_callback = [&](const ix::WebSocketMessagePtr& msg) {
146 if (msg->type == ix::WebSocketMessageType::Message) {
147 m_consumer->QueueEvent(new InputEvt(msg->str));
148 m_driver_stats.rx_count++;
149 } else if (msg->type == ix::WebSocketMessageType::Open) {
150 wxLogDebug("websocket: Connection to %s established",
151 m_ws.getUrl().c_str());
152 std::lock_guard lock(m_stats_mutex);
153 m_driver_stats.available = true;
154 } else if (msg->type == ix::WebSocketMessageType::Close) {
155 wxLogDebug("websocket: Connection disconnected");
156 std::lock_guard lock(m_stats_mutex);
157 m_driver_stats.available = false;
158 } else if (msg->type == ix::WebSocketMessageType::Error) {
159 std::lock_guard lock(m_stats_mutex);
160 m_driver_stats.error_count++;
161 wxLogDebug("websocket: error: %s", msg->errorInfo.reason.c_str());
162 m_ws.getUrl() == wsAddress.str() ? m_ws.setUrl(wssAddress.str())
163 : m_ws.setUrl(wsAddress.str());
164 }
165 };
166 m_ws.setOnMessageCallback(message_callback);
167
168 m_ws.start();
169 while (KeepGoing()) {
170 std::this_thread::sleep_for(100ms);
171 }
172 m_ws.stop();
173 SignalExit();
174 std::lock_guard lock(m_stats_mutex);
175 m_driver_stats.available = false;
176}
177
178//========================================================================
179// CommDriverSignalKNet implementation
180//
181DriverStats CommDriverSignalKNet::IoThread::GetStats() const {
182 std::lock_guard lock(m_stats_mutex);
183 return m_driver_stats;
184}
185
186CommDriverSignalKNet::CommDriverSignalKNet(const ConnectionParams* params,
187 DriverListener& listener)
188 : CommDriverSignalK(params->GetStrippedDSPort()),
189 m_params(*params),
190 m_listener(listener),
191 m_dog_value(kDogTimeoutSeconds),
192 m_io_thread(std::make_unique<IoThread>(params->GetStrippedDSPort(),
193 ParamsIpAddress(*params), this,
194 params->AuthToken.ToStdString())),
195 m_stats_timer(*this, 2s) {
196 // Prepare the wxEventHandler to accept events from the actual hardware thread
197 Bind(SignalkEvtType, &CommDriverSignalKNet::HandleSkSentence, this);
198
199 m_socketread_watchdog_timer.SetOwner(this, kTimerSocket);
200
201 // Dummy Driver Stats, may be polled before worker thread is active
202 m_driver_stats.driver_bus = NavAddr::Bus::Signalk;
203 m_driver_stats.driver_iface = m_params.GetStrippedDSPort();
204 m_driver_stats.available = false;
205
206 Open();
207}
208
209CommDriverSignalKNet::~CommDriverSignalKNet() { Close(); }
210
212 if (m_std_thread.joinable())
213 return m_io_thread->GetStats();
214 else
215 return m_driver_stats;
216}
217
218void CommDriverSignalKNet::Open() {
219 wxString discoveredIP;
220#if 0
221 int discoveredPort;
222#endif
223
224 // if (m_useWebSocket)
225 {
226 std::string service_ident =
227 std::string("_signalk-ws._tcp.local."); // Works for node.js server
228 OpenWebSocket();
229 }
230}
231void CommDriverSignalKNet::Close() { CloseWebSocket(); }
232
233bool CommDriverSignalKNet::DiscoverSkServer(const std::string& service_ident,
234 wxString& ip, int& port, int tSec) {
235 auto servscan = std::make_unique<wxServDisc>(
236 nullptr, wxString(service_ident.c_str()), QTYPE_PTR);
237 for (int i = 0; i < 10; i++) {
238 if (servscan->getResultCount()) {
239 auto result = servscan->getResults().at(0);
240 auto namescan =
241 std::make_unique<wxServDisc>(nullptr, result.name, QTYPE_SRV);
242 for (int j = 0; j < 10; j++) {
243 if (namescan->getResultCount()) {
244 auto namescan_result = namescan->getResults().at(0);
245 port = namescan_result.port;
246 auto addrscan = std::make_unique<wxServDisc>(
247 nullptr, namescan_result.name, QTYPE_A);
248 for (int k = 0; k < 10; k++) {
249 if (addrscan->getResultCount()) {
250 auto addrscan_result = addrscan->getResults().at(0);
251 ip = addrscan_result.ip;
252 return true;
253 } else {
254 wxYield();
255 wxMilliSleep(1000 * tSec / 10);
256 }
257 }
258 return false;
259 } else {
260 wxYield();
261 wxMilliSleep(1000 * tSec / 10);
262 }
263 }
264 return false;
265 } else {
266 wxYield();
267 wxMilliSleep(1000 * tSec / 10);
268 }
269 }
270 return false;
271}
272
273void CommDriverSignalKNet::OpenWebSocket() {
274 wxLogMessage("Opening Signal K WebSocket client: %s",
275 m_params.GetDSPort().c_str());
276 m_std_thread = std::thread([&] { m_io_thread->Run(); });
277 if (!m_std_thread.joinable()) {
278 wxLogError("Can't create WebSocketThread!");
279 return;
280 }
281 ResetWatchdog();
282 m_socketread_watchdog_timer.Start(1000, wxTIMER_ONE_SHOT);
283}
284
285void CommDriverSignalKNet::CloseWebSocket() {
286 if (m_std_thread.joinable()) {
287 if (m_io_thread->IsRunning()) {
288 wxLogMessage("Stopping Secondary SignalK Thread");
289 m_stats_timer.Stop();
290 m_io_thread->RequestStop();
291 std::chrono::milliseconds stop_delay;
292 bool stop_ok = m_io_thread->WaitUntilStopped(10s, stop_delay);
293 if (stop_ok)
294 MESSAGE_LOG << "Stopped in" << stop_delay.count() << " msec.";
295 else
296 WARNING_LOG << "Not stopped after 10 sec.";
297 }
298 wxMilliSleep(100);
299 m_std_thread.join();
300 } else {
301 WARNING_LOG << "Thread unexpectedly died";
302 }
303}
304
305void CommDriverSignalKNet::HandleSkSentence(const InputEvt& event) {
306 rapidjson::Document root;
307
308 std::string msg = event.GetPayload();
309 root.Parse(msg);
310 if (root.HasParseError()) {
311 wxLogMessage(
312 "SignalKDataStream ERROR: the JSON document is not well-formed: %d",
313 root.GetParseError());
314 return;
315 }
316 if (!root.IsObject()) {
317 wxLogMessage("SignalKDataStream ERROR: Message is not a JSON Object: %s",
318 msg.c_str());
319 return;
320 }
321
322 // Decode just enough of string to extract some identifiers
323 // such as the sK version, "self" context, and target context
324 if (root.HasMember("version")) {
325 wxString vers = "Connected to Signal K server version: ";
326 vers << (root["version"].GetString());
327 wxLogMessage(vers);
328 }
329 if (root.HasMember("self")) {
330 if (strncmp(root["self"].GetString(), "vessels.", 8) == 0)
331 m_self = (root["self"].GetString()); // for java server, and OpenPlotter
332 // node.js server 1.20
333 else
334 m_self = std::string("vessels.")
335 .append(root["self"].GetString()); // for Node.js server
336 }
337 if (root.HasMember("context") && root["context"].IsString()) {
338 m_context = root["context"].GetString();
339 }
340
341 // Notify all listeners
342 auto pos = iface.find(':');
343 std::string comm_interface;
344 if (pos != std::string::npos) comm_interface = iface.substr(pos + 1);
345 auto navmsg = std::make_shared<const SignalkMsg>(m_self, m_context, msg,
346 comm_interface);
347 m_listener.Notify(std::move(navmsg));
348}
349
350void CommDriverSignalKNet::initIXNetSystem() { ix::initNetSystem(); };
351
352void CommDriverSignalKNet::uninitIXNetSystem() { ix::uninitNetSystem(); };
const std::string iface
Physical device for 0183, else a unique string.
Definition comm_driver.h:95
static void uninitIXNetSystem()
ix::uninitIXNetSystem wrapper
DriverStats GetDriverStats() const override
Get the Driver Statistics.
static void initIXNetSystem()
ix::initIXNetSystem wrapper
static bool DiscoverSkServer(const std::string &service_ident, wxString &ip, int &port, int tSec)
Scan for a SignalK server on local network using mDNS.
Interface for handling incoming messages.
Definition comm_driver.h:50
virtual void Notify(std::shared_ptr< const NavMsg > message)=0
Handle a received message.
Custom event class for OpenCPN's notification system.
Thread mixin providing a "stop thread"/"wait until stopped" interface.
Definition thread_ctrl.h:31
SignalK IP network driver.
Raw messages layer, supports sending and recieving navmsg messages.
Enhanced logging interface on top of wx/log.h.
Driver statistics report.
unsigned rx_count
Number of bytes received since program start.
unsigned error_count
Number of detected errors since program start.
Suspend/resume and new devices events exchange point.
ThreadCtrl mixin class definition.