243 virtual void async_write(std::shared_ptr<const goby::middleware::protobuf::IOData> io_msg) = 0;
252 void loop()
override;
256 std::unique_ptr<SocketType> socket_;
263 std::mutex incoming_mail_notify_mutex_;
264 std::unique_ptr<std::thread> incoming_mail_notify_thread_;
266 std::string glog_group_;
267 std::string thread_name_;
268 bool glog_group_added_{
false};
271template <
class IOThreadImplementation>
273 std::shared_ptr<const goby::middleware::protobuf::IOData> io_msg)
275 boost::asio::async_write(
276 this_thread->mutable_socket(), boost::asio::buffer(io_msg->data()),
278 [this_thread, io_msg](
const boost::system::error_code& ec, std::size_t bytes_transferred)
280 if (!ec && bytes_transferred > 0)
282 this_thread->handle_write_success(bytes_transferred);
286 this_thread->handle_write_error(ec);
300 template <
class>
class ThreadType,
bool use_indexed_groups>
302 subscribe_layer, IOConfig, SocketType, ThreadType,
303 use_indexed_groups>::try_open()
307 socket_.reset(
new SocketType(io_));
317 backoff_interval_ = min_backoff_interval_;
319 auto status = std::make_shared<protobuf::IOStatus>();
320 if (this->
index() != -1)
323 status->set_state(protobuf::IO__LINK_OPEN);
324 this->publish_in(status);
331 next_open_attempt_ = now + backoff_interval_;
333 catch (
const std::exception& e)
335 auto status = std::make_shared<protobuf::IOStatus>();
336 if (this->
index() != -1)
339 status->set_state(protobuf::IO__CRITICAL_FAILURE);
342 error.set_text(
e.what() + std::string(
": config (") + this->cfg().ShortDebugString() +
")");
343 this->publish_in(status);
346 <<
"Failed to open/configure socket/serial_port: "
347 <<
error.ShortDebugString() << std::endl;
349 if (backoff_interval_ < max_backoff_interval_)
350 backoff_interval_ *= 2.0;
353 next_open_attempt_ = now + backoff_interval_;
356 << backoff_interval_ / std::chrono::seconds(1)
357 <<
" seconds" << std::endl;
366 template <
class>
class ThreadType,
bool use_indexed_groups>
368 subscribe_layer, IOConfig, SocketType, ThreadType,
369 use_indexed_groups>::loop()
371 if (socket_ && socket_->is_open())
381 if (now > next_open_attempt_)
398 template <
class>
class ThreadType,
bool use_indexed_groups>
400 line_in_group, line_out_group, publish_layer, subscribe_layer, IOConfig, SocketType, ThreadType,
401 use_indexed_groups>::handle_read_error(
const boost::system::error_code& ec)
403 auto status = std::make_shared<protobuf::IOStatus>();
404 if (this->index() != -1)
405 status->set_index(this->index());
407 status->set_state(protobuf::IO__CRITICAL_FAILURE);
410 error.set_text(ec.message());
411 this->publish_in(status);
414 <<
"Failed to read from the socket/serial_port: "
415 << error.ShortDebugString() << std::endl;
424 template <
class>
class ThreadType,
bool use_indexed_groups>
426 line_in_group, line_out_group, publish_layer, subscribe_layer, IOConfig, SocketType, ThreadType,
427 use_indexed_groups>::handle_write_error(
const boost::system::error_code& ec)
429 auto status = std::make_shared<protobuf::IOStatus>();
430 if (this->index() != -1)
431 status->set_index(this->index());
433 status->set_state(protobuf::IO__CRITICAL_FAILURE);
436 error.set_text(ec.message());
437 this->publish_in(status);
440 <<
"Failed to write to the socket/serial_port: "
441 << error.ShortDebugString() << std::endl;