30#include <wx/msw/winundef.h>
36#include <netinet/tcp.h>
44#include <wx/datetime.h>
48#include <wx/chartype.h>
49#include <wx/sckaddr.h>
51#include "observable/observable.h"
62using namespace std::literals::chrono_literals;
64#define N_DOG_TIMEOUT 8
67static bool IsBroadcastAddr(
unsigned addr,
unsigned netmask_bits) {
68 assert(netmask_bits <= 32);
69#if defined(_MSC_VER) || __BYTE_ORDER__ == __ORDER_LITTLE_ENDIAN__
70 uint32_t netmask = 0xffffffff >> (32 - netmask_bits);
72 uint32_t netmask = 0xffffffff << (32 - netmask_bits);
74 uint32_t host_mask = ~netmask;
75 return (addr & host_mask) == host_mask;
81 void SetMrqAddr(
unsigned int addr) {
82 m_mrq.imr_multiaddr.s_addr = addr;
83 m_mrq.imr_interface.s_addr = INADDR_ANY;
87static bool SetOutputSocketOptions(wxSocketBase* tsock) {
97 ret = tsock->SetOption(IPPROTO_TCP, TCP_NODELAY, &nagleDisable,
98 sizeof(nagleDisable));
104 unsigned long outbuf_size = 1024;
105 return (tsock->SetOption(SOL_SOCKET, SO_SNDBUF, &outbuf_size,
106 sizeof(outbuf_size)) &&
118 m_listener(listener),
121 m_socket_server(nullptr),
122 m_is_multicast(false),
123 m_stats_timer(*this, 2s),
126 m_rx_connect_event(false),
127 m_socket_timer(*this),
128 m_socketread_watchdog_timer(*this),
130 m_is_conn_err_reported(false) {
131 m_addr.Hostname(params->NetworkAddress);
132 m_addr.Service(params->NetworkPort);
133 this->attributes[
"netAddress"] = params->NetworkAddress.ToStdString();
134 this->attributes[
"netPort"] = std::to_string(params->NetworkPort);
135 this->attributes[
"userComment"] = params->UserComment.ToStdString();
137 m_driver_stats.driver_bus = NavAddr::Bus::N0183;
138 m_driver_stats.driver_iface = params->GetStrippedDSPort();
140 m_mrq_container = std::make_unique<MrqContainer>();
143 resume_listener.Init(SystemEvents::GetInstance().evt_resume,
145 Bind(wxEVT_SOCKET, &CommDriverN0183Net::OnSocketEvent,
this, DS_SOCKET_ID);
146 Bind(wxEVT_SOCKET, &CommDriverN0183Net::OnServerSocketEvent,
this,
152CommDriverN0183Net::~CommDriverN0183Net() { Close(); }
154void CommDriverN0183Net::HandleN0183Msg(
const std::string& sentence) {
156 m_driver_stats.
rx_count += sentence.size();
160void CommDriverN0183Net::Open() {
163 ((
struct sockaddr_in*)m_addr.GetAddressData())->sin_addr.s_addr;
165 unsigned int addr = inet_addr(m_addr.IPAddress().mb_str());
168 switch (m_params.NetProtocol) {
174 OpenNetworkTcp(addr);
178 OpenNetworkUdp(addr);
187void CommDriverN0183Net::OpenNetworkUdp(
unsigned int addr) {
188 if (m_params.is_server) {
191 wxIPV4address conn_addr;
192 conn_addr.Service(std::to_string(m_params.NetworkPort));
193 conn_addr.AnyAddress();
194 conn_addr.AnyAddress();
196 new wxDatagramSocket(conn_addr, wxSOCKET_NOWAIT | wxSOCKET_REUSEADDR);
199 if ((ntohl(addr) & 0xf0000000) == 0xe0000000) {
200 m_is_multicast =
true;
201 m_mrq_container->SetMrqAddr(addr);
202 m_sock->SetOption(IPPROTO_IP, IP_ADD_MEMBERSHIP, &m_mrq_container->m_mrq,
203 sizeof(m_mrq_container->m_mrq));
206 m_sock->SetEventHandler(*
this, DS_SOCKET_ID);
208 m_sock->SetNotify(wxSOCKET_CONNECTION_FLAG | wxSOCKET_INPUT_FLAG |
210 m_sock->Notify(TRUE);
211 m_sock->SetTimeout(1);
212 m_driver_stats.available =
true;
216 if (m_params.IOSelect != DS_TYPE_INPUT) {
217 wxIPV4address tconn_addr;
218 tconn_addr.Service(0);
219 tconn_addr.AnyAddress();
221 new wxDatagramSocket(tconn_addr, wxSOCKET_NOWAIT | wxSOCKET_REUSEADDR);
226 if (!m_is_multicast && IsBroadcastAddr(addr, g_netmask_bits)) {
227 int broadcastEnable = 1;
228 m_tsock->SetOption(SOL_SOCKET, SO_BROADCAST, &broadcastEnable,
229 sizeof(broadcastEnable));
230 m_driver_stats.available =
true;
235 m_connect_time = std::chrono::steady_clock::now();
238void CommDriverN0183Net::OpenNetworkTcp(
unsigned int addr) {
239 if (addr == INADDR_ANY) {
240 MESSAGE_LOG <<
"Listening for TCP connections on " << INADDR_ANY;
241 m_socket_server =
new wxSocketServer(m_addr, wxSOCKET_REUSEADDR);
242 m_socket_server->SetEventHandler(*
this, DS_SERVERSOCKET_ID);
243 m_socket_server->SetNotify(wxSOCKET_CONNECTION_FLAG);
244 m_socket_server->Notify(TRUE);
245 m_socket_server->SetTimeout(1);
246 m_driver_stats.available = m_socket_server->IsOk();
248 MESSAGE_LOG <<
"Opening TCP connection to " << m_params.NetworkAddress
249 <<
":" << m_params.NetworkPort;
250 m_sock =
new wxSocketClient();
251 m_sock->SetEventHandler(*
this, DS_SOCKET_ID);
252 int notify_flags = (wxSOCKET_CONNECTION_FLAG | wxSOCKET_LOST_FLAG);
253 if (m_params.IOSelect != DS_TYPE_INPUT)
254 notify_flags |= wxSOCKET_OUTPUT_FLAG;
255 if (m_params.IOSelect != DS_TYPE_OUTPUT)
256 notify_flags |= wxSOCKET_INPUT_FLAG;
257 m_sock->SetNotify(notify_flags);
258 m_sock->Notify(
true);
259 m_sock->SetTimeout(1);
261 m_rx_connect_event =
false;
262 m_socket_timer.Start(100, wxTIMER_ONE_SHOT);
263 m_driver_stats.available = m_sock->IsOk();
266 m_connect_time = std::chrono::steady_clock::now();
269void CommDriverN0183Net::OpenNetworkGpsd() {
270 m_sock =
new wxSocketClient();
271 m_sock->SetEventHandler(*
this, DS_SOCKET_ID);
272 m_sock->SetNotify(wxSOCKET_CONNECTION_FLAG | wxSOCKET_INPUT_FLAG |
274 m_sock->Notify(TRUE);
275 m_sock->SetTimeout(1);
277 auto* tcp_socket =
dynamic_cast<wxSocketClient*
>(m_sock);
278 tcp_socket->Connect(m_addr,
false);
279 m_rx_connect_event =
false;
282void CommDriverN0183Net::OnSocketReadWatchdogTimer() {
285 if (m_dog_value <= 0) {
286 if (GetParams().NoDataReconnect) {
288 if (m_params.NetProtocol == TCP) {
289 auto* tcp_socket =
dynamic_cast<wxSocketClient*
>(m_sock);
290 if (tcp_socket) tcp_socket->Close();
292 int n_reconnect_delay = wxMax(N_DOG_TIMEOUT - 2, 2);
293 wxLogMessage(
"Reconnection scheduled in %d seconds.",
295 m_socket_timer.Start(n_reconnect_delay * 1000, wxTIMER_ONE_SHOT);
298 m_socketread_watchdog_timer.Stop();
304void CommDriverN0183Net::OnTimerSocket() {
306 using namespace std::chrono;
307 auto* tcp_socket =
dynamic_cast<wxSocketClient*
>(m_sock);
309 if (tcp_socket->IsDisconnected()) {
310 m_driver_stats.available =
false;
311 wxLogDebug(
"Attempting reconnection...");
312 m_rx_connect_event =
false;
314 m_socketread_watchdog_timer.Stop();
315 tcp_socket->Connect(m_addr,
false);
318 int n_reconnect_delay = N_DOG_TIMEOUT;
319 m_socket_timer.Start(n_reconnect_delay * 1000, wxTIMER_ONE_SHOT);
322 if (m_connect_time == time_point<steady_clock>())
return;
323 auto since_connect = steady_clock::now() - m_connect_time;
324 if (since_connect > 10s && !m_is_conn_err_reported) {
325 std::stringstream ss;
326 ss << _(
"Cannot connect to remote server ") << m_params.NetworkAddress
327 <<
":" << m_params.NetworkPort;
328 CommDriverRegistry::GetInstance().
evt_driver_msg.Notify(ss.str());
329 m_is_conn_err_reported =
true;
333 m_driver_stats.available = tcp_socket->IsOk();
337void CommDriverN0183Net::HandleResume() {
339 auto* tcp_socket =
dynamic_cast<wxSocketClient*
>(m_sock);
341 m_socketread_watchdog_timer.Stop();
346 int n_reconnect_delay = wxMax(N_DOG_TIMEOUT - 2, 2);
347 wxLogMessage(
"Reconnection scheduled in %d seconds.", n_reconnect_delay);
349 m_socket_timer.Start(n_reconnect_delay * 1000, wxTIMER_ONE_SHOT);
353bool CommDriverN0183Net::SendMessage(std::shared_ptr<const NavMsg> msg,
354 std::shared_ptr<const NavAddr> addr) {
355 auto msg_0183 = std::dynamic_pointer_cast<const Nmea0183Msg>(msg);
356 std::string payload(msg_0183->payload);
358 m_driver_stats.
tx_count += payload.size();
359 return SendSentenceNetwork(payload.c_str());
362void CommDriverN0183Net::OnSocketEvent(wxSocketEvent& event) {
363#define RD_BUF_SIZE 4096
367 switch (event.GetSocketEvent()) {
382 uint8_t buff[RD_BUF_SIZE + 1];
383 event.GetSocket()->Read(buff, RD_BUF_SIZE);
384 if (!event.GetSocket()->Error()) {
385 unsigned count =
event.GetSocket()->LastCount();
386 for (
unsigned i = 0; i < count; i += 1) n0183_buffer.
Put(buff[i]);
388 HandleN0183Msg(n0183_buffer.
GetSentence() +
"\r\n");
391 m_dog_value = N_DOG_TIMEOUT;
395 case wxSOCKET_LOST: {
396 m_driver_stats.available = GetSock()->IsOk();
397 using namespace std::chrono;
398 if (m_params.NetProtocol == TCP || m_params.NetProtocol == GPSD) {
399 if (m_rx_connect_event) {
400 MESSAGE_LOG <<
"NetworkDataStream connection lost: "
401 << m_params.GetDSPort();
403 if (m_socket_server) {
408 auto since_connect = 10s;
410 auto now = steady_clock::now();
411 if (m_connect_time != time_point<steady_clock>())
412 since_connect = duration_cast<seconds>(now - m_connect_time);
414 auto retry_time = 5s;
418 if (!m_rx_connect_event && (since_connect < 5s)) retry_time = 10s;
420 m_socketread_watchdog_timer.Stop();
423 m_socket_timer.Start(duration_cast<milliseconds>(retry_time).count(),
429 case wxSOCKET_CONNECTION: {
430 if (m_params.NetProtocol == GPSD) {
435 char cmd[] = R
"--(?WATCH={"class":"WATCH", "nmea":true})--";
436 m_sock->Write(cmd, strlen(cmd));
437 } else if (m_params.NetProtocol == TCP) {
438 MESSAGE_LOG <<
"TCP NetworkDataStream connection established: "
439 << m_params.GetDSPort();
441 m_dog_value = N_DOG_TIMEOUT;
442 if (m_params.IOSelect != DS_TYPE_OUTPUT) {
444 if (GetParams().NoDataReconnect)
445 m_socketread_watchdog_timer.Start(1000);
448 if (m_params.IOSelect != DS_TYPE_INPUT && GetSock()->IsOk())
449 (void)SetOutputSocketOptions(m_sock);
450 m_socket_timer.Stop();
451 m_rx_connect_event =
true;
454 m_driver_stats.available =
true;
455 m_connect_time = std::chrono::steady_clock::now();
464void CommDriverN0183Net::OnServerSocketEvent(wxSocketEvent& event) {
465 switch (event.GetSocketEvent()) {
466 case wxSOCKET_CONNECTION: {
467 m_sock = m_socket_server->Accept(
false);
470 m_sock->SetTimeout(2);
472 m_sock->SetEventHandler(*
this, DS_SOCKET_ID);
473 int notify_flags = (wxSOCKET_CONNECTION_FLAG | wxSOCKET_LOST_FLAG);
474 if (m_params.IOSelect != DS_TYPE_INPUT) {
475 notify_flags |= wxSOCKET_OUTPUT_FLAG;
476 (void)SetOutputSocketOptions(m_sock);
478 if (m_params.IOSelect != DS_TYPE_OUTPUT)
479 notify_flags |= wxSOCKET_INPUT_FLAG;
480 m_sock->SetNotify(notify_flags);
481 m_sock->Notify(
true);
490bool CommDriverN0183Net::SendSentenceNetwork(
const wxString& payload) {
497 wxDatagramSocket* udp_socket;
498 switch (m_params.NetProtocol) {
500 if (GetSock() && GetSock()->IsOk()) {
501 m_sock->Write(payload.mb_str(), strlen(payload.mb_str()));
502 m_dog_value = N_DOG_TIMEOUT;
503 if (GetSock()->Error()) {
504 if (m_socket_server) {
508 auto* tcp_socket =
dynamic_cast<wxSocketClient*
>(m_sock);
509 if (tcp_socket) tcp_socket->Close();
510 if (!m_socket_timer.IsRunning())
511 m_socket_timer.Start(5000, wxTIMER_ONE_SHOT);
513 m_socketread_watchdog_timer.Stop();
524 udp_socket =
dynamic_cast<wxDatagramSocket*
>(m_tsock);
525 if (udp_socket && udp_socket->IsOk()) {
526 udp_socket->SendTo(m_addr, payload.mb_str(), payload.size());
527 m_dog_value = N_DOG_TIMEOUT;
528 if (udp_socket->Error()) ret =
false;
532 m_driver_stats.available = ret;
544void CommDriverN0183Net::Close() {
545 MESSAGE_LOG <<
"Closing NMEA NetworkDataStream " << m_params.NetworkPort;
546 m_stats_timer.Stop();
550 m_sock->SetOption(IPPROTO_IP, IP_DROP_MEMBERSHIP, &m_mrq_container->m_mrq,
551 sizeof(m_mrq_container->m_mrq));
552 m_sock->Notify(FALSE);
557 m_tsock->Notify(FALSE);
561 if (m_socket_server) {
562 m_socket_server->Notify(FALSE);
563 m_socket_server->Destroy();
566 m_socket_timer.Stop();
567 m_socketread_watchdog_timer.Stop();
568 m_driver_stats.available =
false;
NMEA0183 basic parsing common parts:
void SendToListener(const std::string &payload, DriverListener &listener, const ConnectionParams ¶ms)
Wrap argument string in NavMsg pointer, forward to listener.
obs::EventVar evt_driver_msg
Notified for messages from drivers.
Interface for handling incoming messages.
bool HasSentence() const
Return true if a sentence is available to be returned by GetSentence()
std::string GetSentence()
Retrieve a sentence from buffer.
void Put(uint8_t ch)
Add a single character, possibly making a sentence available.
Where messages are sent to or received from.
Custom event class for OpenCPN's notification system.
Driver registration container, a singleton.
Raw messages layer, supports sending and recieving navmsg messages.
Global variables stored in configuration file.
std::string DsPortTypeToString(dsPortType type)
Return textual representation for use in driver ioDirection attribute.
GUI constant definitions.
Enhanced logging interface on top of wx/log.h.
bool endswith(const std::string &str, const std::string &suffix)
Return true if s ends with given suffix.
unsigned tx_count
Number of bytes sent since program start.
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.