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>
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->network_address);
132 m_addr.Service(params->network_port);
133 this->attributes[
"netAddress"] = params->network_address.ToStdString();
134 this->attributes[
"netPort"] = std::to_string(params->network_port);
135 this->attributes[
"userComment"] = params->user_comment.ToStdString();
137 m_driver_stats.driver_bus = NavAddr::Bus::N0183;
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.net_protocol) {
174 OpenNetworkTcp(addr);
178 OpenNetworkUdp(addr);
187void CommDriverN0183Net::OpenNetworkUdp(
unsigned int addr) {
188 if (m_params.direction != PortDirection::kOutput &&
189 m_params.direction != PortDirection::kUpload) {
192 wxIPV4address conn_addr;
193 conn_addr.Service(std::to_string(m_params.network_port));
194 conn_addr.AnyAddress();
195 conn_addr.AnyAddress();
197 new wxDatagramSocket(conn_addr, wxSOCKET_NOWAIT | wxSOCKET_REUSEADDR);
200 if ((ntohl(addr) & 0xf0000000) == 0xe0000000) {
201 m_is_multicast =
true;
202 m_mrq_container->SetMrqAddr(addr);
203 m_sock->SetOption(IPPROTO_IP, IP_ADD_MEMBERSHIP, &m_mrq_container->m_mrq,
204 sizeof(m_mrq_container->m_mrq));
207 m_sock->SetEventHandler(*
this, DS_SOCKET_ID);
209 m_sock->SetNotify(wxSOCKET_CONNECTION_FLAG | wxSOCKET_INPUT_FLAG |
211 m_sock->Notify(TRUE);
212 m_sock->SetTimeout(1);
213 m_driver_stats.available =
true;
217 if (m_params.direction != PortDirection::kInput) {
218 wxIPV4address tconn_addr;
219 tconn_addr.Service(0);
220 tconn_addr.AnyAddress();
222 new wxDatagramSocket(tconn_addr, wxSOCKET_NOWAIT | wxSOCKET_REUSEADDR);
227 if (!m_is_multicast && IsBroadcastAddr(addr, g_netmask_bits)) {
228 int broadcastEnable = 1;
229 m_tsock->SetOption(SOL_SOCKET, SO_BROADCAST, &broadcastEnable,
230 sizeof(broadcastEnable));
231 m_driver_stats.available =
true;
236 m_connect_time = std::chrono::steady_clock::now();
239void CommDriverN0183Net::OpenNetworkTcp(
unsigned int addr) {
240 if (addr == INADDR_ANY) {
241 MESSAGE_LOG <<
"Listening for TCP connections on " << INADDR_ANY;
242 m_socket_server =
new wxSocketServer(m_addr, wxSOCKET_REUSEADDR);
243 m_socket_server->SetEventHandler(*
this, DS_SERVERSOCKET_ID);
244 m_socket_server->SetNotify(wxSOCKET_CONNECTION_FLAG);
245 m_socket_server->Notify(TRUE);
246 m_socket_server->SetTimeout(1);
247 m_driver_stats.available = m_socket_server->IsOk();
249 MESSAGE_LOG <<
"Opening TCP connection to " << m_params.network_address
250 <<
":" << m_params.network_port;
251 m_sock =
new wxSocketClient();
252 m_sock->SetEventHandler(*
this, DS_SOCKET_ID);
253 int notify_flags = (wxSOCKET_CONNECTION_FLAG | wxSOCKET_LOST_FLAG);
254 if (m_params.direction != PortDirection::kInput)
255 notify_flags |= wxSOCKET_OUTPUT_FLAG;
256 if (m_params.direction != PortDirection::kOutput)
257 notify_flags |= wxSOCKET_INPUT_FLAG;
258 m_sock->SetNotify(notify_flags);
259 m_sock->Notify(
true);
260 m_sock->SetTimeout(1);
262 m_rx_connect_event =
false;
263 m_socket_timer.Start(100, wxTIMER_ONE_SHOT);
264 m_driver_stats.available = m_sock->IsOk();
267 m_connect_time = std::chrono::steady_clock::now();
270void CommDriverN0183Net::OpenNetworkGpsd() {
271 m_sock =
new wxSocketClient();
272 m_sock->SetEventHandler(*
this, DS_SOCKET_ID);
273 m_sock->SetNotify(wxSOCKET_CONNECTION_FLAG | wxSOCKET_INPUT_FLAG |
275 m_sock->Notify(TRUE);
276 m_sock->SetTimeout(1);
278 auto* tcp_socket =
dynamic_cast<wxSocketClient*
>(m_sock);
279 tcp_socket->Connect(m_addr,
false);
280 m_rx_connect_event =
false;
283void CommDriverN0183Net::OnSocketReadWatchdogTimer() {
286 if (m_dog_value <= 0) {
287 if (GetParams().no_data_reconnect) {
289 if (m_params.net_protocol == TCP) {
290 auto* tcp_socket =
dynamic_cast<wxSocketClient*
>(m_sock);
291 if (tcp_socket) tcp_socket->Close();
293 int n_reconnect_delay = wxMax(N_DOG_TIMEOUT - 2, 2);
294 wxLogMessage(
"Reconnection scheduled in %d seconds.",
296 m_socket_timer.Start(n_reconnect_delay * 1000, wxTIMER_ONE_SHOT);
299 m_socketread_watchdog_timer.Stop();
305void CommDriverN0183Net::OnTimerSocket() {
307 using namespace std::chrono;
308 auto* tcp_socket =
dynamic_cast<wxSocketClient*
>(m_sock);
310 if (tcp_socket->IsDisconnected()) {
311 m_driver_stats.available =
false;
312 wxLogDebug(
"Attempting reconnection...");
313 m_rx_connect_event =
false;
315 m_socketread_watchdog_timer.Stop();
316 tcp_socket->Connect(m_addr,
false);
319 int n_reconnect_delay = N_DOG_TIMEOUT;
320 m_socket_timer.Start(n_reconnect_delay * 1000, wxTIMER_ONE_SHOT);
323 if (m_connect_time == time_point<steady_clock>())
return;
324 auto since_connect = steady_clock::now() - m_connect_time;
325 if (since_connect > 10s && !m_is_conn_err_reported) {
326 std::stringstream ss;
327 ss << _(
"Cannot connect to remote server ") << m_params.network_address
328 <<
":" << m_params.network_port;
330 m_is_conn_err_reported =
true;
334 m_driver_stats.available = tcp_socket->IsOk();
338void CommDriverN0183Net::HandleResume() {
340 auto* tcp_socket =
dynamic_cast<wxSocketClient*
>(m_sock);
342 m_socketread_watchdog_timer.Stop();
347 int n_reconnect_delay = wxMax(N_DOG_TIMEOUT - 2, 2);
348 wxLogMessage(
"Reconnection scheduled in %d seconds.", n_reconnect_delay);
350 m_socket_timer.Start(n_reconnect_delay * 1000, wxTIMER_ONE_SHOT);
354bool CommDriverN0183Net::SendMessage(std::shared_ptr<const NavMsg> msg,
355 std::shared_ptr<const NavAddr> addr) {
356 auto msg_0183 = std::dynamic_pointer_cast<const Nmea0183Msg>(msg);
357 std::string payload(msg_0183->payload);
359 m_driver_stats.
tx_count += payload.size();
360 return SendSentenceNetwork(payload.c_str());
363void CommDriverN0183Net::OnSocketEvent(wxSocketEvent& event) {
364#define RD_BUF_SIZE 4096
368 switch (event.GetSocketEvent()) {
383 uint8_t buff[RD_BUF_SIZE + 1];
384 event.GetSocket()->Read(buff, RD_BUF_SIZE);
385 if (!event.GetSocket()->Error()) {
386 unsigned count =
event.GetSocket()->LastCount();
387 for (
unsigned i = 0; i < count; i += 1) n0183_buffer.
Put(buff[i]);
389 HandleN0183Msg(n0183_buffer.
GetSentence() +
"\r\n");
392 m_dog_value = N_DOG_TIMEOUT;
396 case wxSOCKET_LOST: {
397 m_driver_stats.available = GetSock()->IsOk();
398 using namespace std::chrono;
399 if (m_params.net_protocol == TCP || m_params.net_protocol == GPSD) {
400 if (m_rx_connect_event) {
401 MESSAGE_LOG <<
"NetworkDataStream connection lost: "
404 if (m_socket_server) {
409 auto since_connect = 10s;
411 auto now = steady_clock::now();
412 if (m_connect_time != time_point<steady_clock>())
413 since_connect = duration_cast<seconds>(now - m_connect_time);
415 auto retry_time = 5s;
419 if (!m_rx_connect_event && (since_connect < 5s)) retry_time = 10s;
421 m_socketread_watchdog_timer.Stop();
424 m_socket_timer.Start(duration_cast<milliseconds>(retry_time).count(),
430 case wxSOCKET_CONNECTION: {
431 if (m_params.net_protocol == GPSD) {
436 char cmd[] = R
"--(?WATCH={"class":"WATCH", "nmea":true})--";
437 m_sock->Write(cmd, strlen(cmd));
438 } else if (m_params.net_protocol == TCP) {
439 MESSAGE_LOG <<
"TCP NetworkDataStream connection established: "
442 m_dog_value = N_DOG_TIMEOUT;
443 if (m_params.direction != PortDirection::kOutput) {
445 if (GetParams().no_data_reconnect)
446 m_socketread_watchdog_timer.Start(1000);
449 if (m_params.direction != PortDirection::kInput && GetSock()->IsOk())
450 (void)SetOutputSocketOptions(m_sock);
451 m_socket_timer.Stop();
452 m_rx_connect_event =
true;
455 m_driver_stats.available =
true;
456 m_connect_time = std::chrono::steady_clock::now();
465void CommDriverN0183Net::OnServerSocketEvent(wxSocketEvent& event) {
466 switch (event.GetSocketEvent()) {
467 case wxSOCKET_CONNECTION: {
468 m_sock = m_socket_server->Accept(
false);
471 m_sock->SetTimeout(2);
473 m_sock->SetEventHandler(*
this, DS_SOCKET_ID);
474 int notify_flags = (wxSOCKET_CONNECTION_FLAG | wxSOCKET_LOST_FLAG);
475 if (m_params.direction != PortDirection::kInput) {
476 notify_flags |= wxSOCKET_OUTPUT_FLAG;
477 (void)SetOutputSocketOptions(m_sock);
479 if (m_params.direction != PortDirection::kOutput)
480 notify_flags |= wxSOCKET_INPUT_FLAG;
481 m_sock->SetNotify(notify_flags);
482 m_sock->Notify(
true);
491bool CommDriverN0183Net::SendSentenceNetwork(
const wxString& payload) {
498 wxDatagramSocket* udp_socket;
499 switch (m_params.net_protocol) {
501 if (GetSock() && GetSock()->IsOk()) {
502 m_sock->Write(payload.mb_str(), strlen(payload.mb_str()));
503 m_dog_value = N_DOG_TIMEOUT;
504 if (GetSock()->Error()) {
505 if (m_socket_server) {
509 auto* tcp_socket =
dynamic_cast<wxSocketClient*
>(m_sock);
510 if (tcp_socket) tcp_socket->Close();
511 if (!m_socket_timer.IsRunning())
512 m_socket_timer.Start(5000, wxTIMER_ONE_SHOT);
514 m_socketread_watchdog_timer.Stop();
525 udp_socket =
dynamic_cast<wxDatagramSocket*
>(m_tsock);
526 if (udp_socket && udp_socket->IsOk()) {
527 udp_socket->SendTo(m_addr, payload.mb_str(), payload.size());
528 m_dog_value = N_DOG_TIMEOUT;
529 if (udp_socket->Error()) ret =
false;
533 m_driver_stats.available = ret;
545void CommDriverN0183Net::Close() {
546 MESSAGE_LOG <<
"Closing NMEA NetworkDataStream " << m_params.network_port;
547 m_stats_timer.Stop();
551 m_sock->SetOption(IPPROTO_IP, IP_DROP_MEMBERSHIP, &m_mrq_container->m_mrq,
552 sizeof(m_mrq_container->m_mrq));
553 m_sock->Notify(FALSE);
558 m_tsock->Notify(FALSE);
562 if (m_socket_server) {
563 m_socket_server->Notify(FALSE);
564 m_socket_server->Destroy();
567 m_socket_timer.Stop();
568 m_socketread_watchdog_timer.Stop();
569 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.
Connection data container close to a POD struct.
std::string GetStrippedDSPort() const
Return port string with possible windows extra data removed, in some cases empty.
wxString GetDSPort() const
Return port description including for example serial port or network address/port,...
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.
void Notify() override
Notify all listeners, no data supplied.
Driver registration container, a singleton.
Raw messages layer, supports sending and recieving navmsg messages.
Global variables stored in configuration file.
std::string PortDirectionToString(PortDirection pd)
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.
General observable pattern implementation built on top of wxWidgets event handling.
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.