30#include "observable/observable.h"
32#include "rapidjson/document.h"
33#include "ixwebsocket/IXNetSystem.h"
34#include "ixwebsocket/IXSocketTLSOptions.h"
35#include "ixwebsocket/IXWebSocket.h"
39#include "model/geodesic.h"
43#include "wxServDisc.h"
45using namespace std::literals::chrono_literals;
47constexpr int kTimerSocket = 9006;
48constexpr int kSignalkSocketId = 5011;
49constexpr int kDogTimeoutReconnectSeconds = 10;
51constexpr double kMsToKnotFactor = 1.9438444924406;
55 IoThread(
const std::string&
iface,
const wxIPV4address& address,
56 wxEvtHandler* consumer,
const std::string& token);
65 wxIPV4address m_address;
66 wxEvtHandler* m_consumer;
67 const std::string m_iface;
70 obs::Listener m_resume_listener;
72 mutable std::mutex m_stats_mutex;
76static const wxEventTypeTag<CommDriverSignalKNet::InputEvt> SignalkEvtType(
81 addr.Hostname(params.NetworkAddress);
82 addr.Service(params.NetworkPort);
88 explicit InputEvt(std::string payload)
89 : wxEvent(0, SignalkEvtType), m_payload(std::move(payload)) {};
91 std::string GetPayload()
const {
return m_payload; }
94 wxEvent* Clone()
const override {
return new InputEvt(m_payload); };
97 const std::string m_payload;
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,
112 wxLogDebug(
"WebSocketThread: restarted");
116void CommDriverSignalKNet::IoThread::Run() {
117 using namespace std::chrono_literals;
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;
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;
138 m_ws.setUrl(wssAddress.str());
139 ix::SocketTLSOptions opt;
140 opt.disable_hostname_validation =
true;
142 m_ws.setTLSOptions(opt);
143 m_ws.setPingInterval(30);
145 auto message_callback = [&](
const ix::WebSocketMessagePtr& msg) {
146 if (msg->type == ix::WebSocketMessageType::Message) {
147 m_consumer->QueueEvent(
new InputEvt(msg->str));
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);
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());
166 m_ws.setOnMessageCallback(message_callback);
169 while (KeepGoing()) {
170 std::this_thread::sleep_for(100ms);
174 std::lock_guard lock(m_stats_mutex);
175 m_driver_stats.available =
false;
181DriverStats CommDriverSignalKNet::IoThread::GetStats()
const {
182 std::lock_guard lock(m_stats_mutex);
183 return m_driver_stats;
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) {
197 Bind(SignalkEvtType, &CommDriverSignalKNet::HandleSkSentence,
this);
199 m_socketread_watchdog_timer.SetOwner(
this, kTimerSocket);
202 m_driver_stats.driver_bus = NavAddr::Bus::Signalk;
203 m_driver_stats.driver_iface = m_params.GetStrippedDSPort();
204 m_driver_stats.available =
false;
209CommDriverSignalKNet::~CommDriverSignalKNet() { Close(); }
212 if (m_std_thread.joinable())
213 return m_io_thread->GetStats();
215 return m_driver_stats;
218void CommDriverSignalKNet::Open() {
219 wxString discoveredIP;
226 std::string service_ident =
227 std::string(
"_signalk-ws._tcp.local.");
231void CommDriverSignalKNet::Close() { CloseWebSocket(); }
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);
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;
255 wxMilliSleep(1000 * tSec / 10);
261 wxMilliSleep(1000 * tSec / 10);
267 wxMilliSleep(1000 * tSec / 10);
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!");
282 m_socketread_watchdog_timer.Start(1000, wxTIMER_ONE_SHOT);
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);
294 MESSAGE_LOG <<
"Stopped in" << stop_delay.count() <<
" msec.";
296 WARNING_LOG <<
"Not stopped after 10 sec.";
301 WARNING_LOG <<
"Thread unexpectedly died";
305void CommDriverSignalKNet::HandleSkSentence(
const InputEvt& event) {
306 rapidjson::Document root;
308 std::string msg =
event.GetPayload();
310 if (root.HasParseError()) {
312 "SignalKDataStream ERROR: the JSON document is not well-formed: %d",
313 root.GetParseError());
316 if (!root.IsObject()) {
317 wxLogMessage(
"SignalKDataStream ERROR: Message is not a JSON Object: %s",
324 if (root.HasMember(
"version")) {
325 wxString vers =
"Connected to Signal K server version: ";
326 vers << (root[
"version"].GetString());
329 if (root.HasMember(
"self")) {
330 if (strncmp(root[
"self"].GetString(),
"vessels.", 8) == 0)
331 m_self = (root[
"self"].GetString());
334 m_self = std::string(
"vessels.")
335 .append(root[
"self"].GetString());
337 if (root.HasMember(
"context") && root[
"context"].IsString()) {
338 m_context = root[
"context"].GetString();
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,
347 m_listener.
Notify(std::move(navmsg));
const std::string iface
Physical device for 0183, else a unique string.
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.
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.
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.