Goby3 3.6.1
2026.09.15
Loading...
Searching...
No Matches
io_interface.h
Go to the documentation of this file.
1// Copyright 2019-2026:
2// GobySoft, LLC (2013-)
3// Community contributors (see AUTHORS file)
4// File authors:
5// Toby Schneider <toby@gobysoft.org>
6//
7//
8// This file is part of the Goby Underwater Autonomy Project Libraries
9// ("The Goby Libraries").
10//
11// The Goby Libraries are free software: you can redistribute them and/or modify
12// them under the terms of the GNU Lesser General Public License as published by
13// the Free Software Foundation, either version 2.1 of the License, or
14// (at your option) any later version.
15//
16// The Goby Libraries are distributed in the hope that they will be useful,
17// but WITHOUT ANY WARRANTY; without even the implied warranty of
18// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
19// GNU Lesser General Public License for more details.
20//
21// You should have received a copy of the GNU Lesser General Public License
22// along with Goby. If not, see <http://www.gnu.org/licenses/>.
23
24#ifndef GOBY_MIDDLEWARE_IO_DETAIL_IO_INTERFACE_H
25#define GOBY_MIDDLEWARE_IO_DETAIL_IO_INTERFACE_H
26
27#include <chrono> // for seconds
28#include <exception> // for exception
29#include <memory> // for shared_ptr
30#include <mutex> // for mutex, lock_...
31#include <ostream> // for endl, size_t
32#include <string> // for string, oper...
33#include <thread> // for thread
34#include <unistd.h> // for usleep
35
36#include <boost/asio/write.hpp> // for async_write
37#include <boost/asio/post.hpp>
38#include <boost/system/error_code.hpp> // for error_code
39
40#include "goby/exception.h" // for Exception
41#include "goby/middleware/application/multi_thread.h" // for SimpleThread
42#include "goby/middleware/common.h" // for thread_id
43#include "goby/middleware/io/groups.h" // for status
44#include "goby/middleware/protobuf/io.pb.h" // for IOError, IOS...
45#include "goby/time/steady_clock.h" // for SteadyClock
47#include "goby/util/debug_logger.h" // for glog
48
49#include "io_transporters.h"
50
51namespace goby
52{
53namespace middleware
54{
55class Group;
56class InterThreadTransporter;
57template <typename InnerTransporter, typename ImplementationTag> class InterProcessForwarder;
58namespace io
59{
60enum class ThreadState
61{
63};
64
65namespace detail
66{
67template <typename ProtobufEndpoint, typename ASIOEndpoint>
68ProtobufEndpoint endpoint_convert(const ASIOEndpoint& asio_ep)
69{
70 ProtobufEndpoint pb_ep;
71 pb_ep.set_addr(asio_ep.address().to_string());
72 pb_ep.set_port(asio_ep.port());
73 return pb_ep;
74}
75
76template <const goby::middleware::Group& line_in_group,
77 const goby::middleware::Group& line_out_group, PubSubLayer publish_layer,
78 PubSubLayer subscribe_layer, typename IOConfig, typename SocketType,
79 template <class> class ThreadType, bool use_indexed_groups = false>
81 : public ThreadType<IOConfig>,
83 IOThread<line_in_group, line_out_group, publish_layer, subscribe_layer, IOConfig,
84 SocketType, ThreadType, use_indexed_groups>,
85 line_in_group, publish_layer,
86 typename ThreadType<IOConfig>::Transporter::implementation_tag, use_indexed_groups>,
88 IOThread<line_in_group, line_out_group, publish_layer, subscribe_layer, IOConfig,
89 SocketType, ThreadType, use_indexed_groups>,
90 line_out_group, subscribe_layer,
91 typename ThreadType<IOConfig>::Transporter::implementation_tag, use_indexed_groups>
92{
93 public:
98 IOThread(const IOConfig& config, int index, std::string glog_group = "i/o")
99 : ThreadType<IOConfig>(config, this->loop_max_frequency(), index),
101 IOThread<line_in_group, line_out_group, publish_layer, subscribe_layer, IOConfig,
102 SocketType, ThreadType, use_indexed_groups>,
103 line_in_group, publish_layer,
104 typename ThreadType<IOConfig>::Transporter::implementation_tag, use_indexed_groups>(
105 index),
107 IOThread<line_in_group, line_out_group, publish_layer, subscribe_layer, IOConfig,
108 SocketType, ThreadType, use_indexed_groups>,
109 line_out_group, subscribe_layer,
110 typename ThreadType<IOConfig>::Transporter::implementation_tag, use_indexed_groups>(
111 index),
112 glog_group_(glog_group + " / t" + std::to_string(goby::middleware::gettid())),
113 thread_name_(glog_group)
114 {
115 auto data_out_callback =
116 [this](std::shared_ptr<const goby::middleware::protobuf::IOData> io_msg)
117 {
118 if (!io_msg->has_index() || io_msg->index() == this->index())
119 {
120 write(io_msg);
121 }
122 };
123
124 this->template subscribe_out<goby::middleware::protobuf::IOData>(data_out_callback);
125
126 if (!glog_group_added_)
127 {
129 glog_group_added_ = true;
130 }
131 }
132
133 void initialize() override
134 {
135 // thread to handle synchonization between boost::asio and goby condition_variable signaling
136 incoming_mail_notify_thread_.reset(new std::thread(
137 [this]()
138 {
139 while (this->alive())
140 {
141 std::unique_lock<std::mutex> lock(incoming_mail_notify_mutex_);
142 this->interthread().cv()->wait(lock);
143 // post empty handler to cause loop() to return and allow incoming mail to be handled
144 boost::asio::post(io_, []() {});
145 }
146 }));
147
148 this->set_name(thread_name_);
149 }
150
151 void finalize() override
152 {
153 // join incoming mail thread
154 {
155 std::lock_guard<std::mutex> l(incoming_mail_notify_mutex_);
156 this->interthread().cv()->notify_all();
157 }
158 incoming_mail_notify_thread_->join();
159 incoming_mail_notify_thread_.reset();
160 }
161
162 virtual ~IOThread()
163 {
164 socket_.reset();
165
166 // for non clean shutdown, avoid abort
167 if (incoming_mail_notify_thread_)
168 incoming_mail_notify_thread_->detach();
169
170 auto status = std::make_shared<protobuf::IOStatus>();
171 status->set_state(protobuf::IO__LINK_CLOSED);
172
173 this->publish_in(status);
174 this->template unsubscribe_out<goby::middleware::protobuf::IOData>();
175 }
176
177 template <class IOThreadImplementation>
178 friend void basic_async_write(IOThreadImplementation* this_thread,
179 std::shared_ptr<const goby::middleware::protobuf::IOData> io_msg);
180
181 protected:
182 void write(std::shared_ptr<const goby::middleware::protobuf::IOData> io_msg)
183 {
185 goby::glog << group(glog_group_) << "(" << io_msg->data().size() << "B) <"
186 << ((this->index() == -1) ? std::string() : std::to_string(this->index()))
187 << " " << io_msg->ShortDebugString() << std::endl;
188 if (io_msg->data().empty())
189 return;
190 if (!socket_ || !socket_->is_open())
191 return;
192
193 this->async_write(io_msg);
194 }
195
196 void handle_read_success(std::size_t bytes_transferred, const std::string& bytes)
197 {
198 auto io_msg = std::make_shared<goby::middleware::protobuf::IOData>();
199 *io_msg->mutable_data() = bytes;
200
201 handle_read_success(bytes_transferred, io_msg);
202 }
203
204 void handle_read_success(std::size_t bytes_transferred,
205 std::shared_ptr<goby::middleware::protobuf::IOData> io_msg)
206 {
207 if (this->index() != -1)
208 io_msg->set_index(this->index());
209
211 goby::glog << group(glog_group_) << "(" << bytes_transferred << "B) >"
212 << ((this->index() == -1) ? std::string() : std::to_string(this->index()))
213 << " " << io_msg->ShortDebugString() << std::endl;
214
215 this->publish_in(io_msg);
216 }
217
218 void handle_write_success(std::size_t bytes_transferred) {}
219 void handle_read_error(const boost::system::error_code& ec);
220 void handle_write_error(const boost::system::error_code& ec);
221
223 SocketType& mutable_socket()
224 {
225 if (socket_)
226 return *socket_;
227 else
228 throw goby::Exception("Attempted to access null socket/serial_port");
229 }
230
232
234 bool socket_is_open() { return socket_ && socket_->is_open(); }
235
237 virtual void open_socket() = 0;
238
240 virtual void async_read() = 0;
241
243 virtual void async_write(std::shared_ptr<const goby::middleware::protobuf::IOData> io_msg) = 0;
244
245 const std::string& glog_group() { return glog_group_; }
246
247 private:
249 void try_open();
250
252 void loop() override;
253
254 private:
256 std::unique_ptr<SocketType> socket_;
257
258 const goby::time::SteadyClock::duration min_backoff_interval_{std::chrono::seconds(1)};
259 const goby::time::SteadyClock::duration max_backoff_interval_{std::chrono::seconds(128)};
260 goby::time::SteadyClock::duration backoff_interval_{min_backoff_interval_};
262
263 std::mutex incoming_mail_notify_mutex_;
264 std::unique_ptr<std::thread> incoming_mail_notify_thread_;
265
266 std::string glog_group_;
267 std::string thread_name_;
268 bool glog_group_added_{false};
269};
270
271template <class IOThreadImplementation>
272void basic_async_write(IOThreadImplementation* this_thread,
273 std::shared_ptr<const goby::middleware::protobuf::IOData> io_msg)
274{
275 boost::asio::async_write(
276 this_thread->mutable_socket(), boost::asio::buffer(io_msg->data()),
277 // capture io_msg in callback to ensure write buffer exists until async_write is done
278 [this_thread, io_msg](const boost::system::error_code& ec, std::size_t bytes_transferred)
279 {
280 if (!ec && bytes_transferred > 0)
281 {
282 this_thread->handle_write_success(bytes_transferred);
283 }
284 else
285 {
286 this_thread->handle_write_error(ec);
287 }
288 });
289}
290
291} // namespace detail
292} // namespace io
293} // namespace middleware
294} // namespace goby
295
296template <const goby::middleware::Group& line_in_group,
297 const goby::middleware::Group& line_out_group,
299 goby::middleware::io::PubSubLayer subscribe_layer, typename IOConfig, typename SocketType,
300 template <class> class ThreadType, bool use_indexed_groups>
301void goby::middleware::io::detail::IOThread<line_in_group, line_out_group, publish_layer,
302 subscribe_layer, IOConfig, SocketType, ThreadType,
303 use_indexed_groups>::try_open()
304{
305 try
306 {
307 socket_.reset(new SocketType(io_));
308 open_socket();
309
310 // messages read from the socket
311 this->async_read();
312
313 // reset io_context, which ran out of work
314 io_.restart();
315
316 // successful, reset backoff
317 backoff_interval_ = min_backoff_interval_;
318
319 auto status = std::make_shared<protobuf::IOStatus>();
320 if (this->index() != -1)
321 status->set_index(this->index());
322
323 status->set_state(protobuf::IO__LINK_OPEN);
324 this->publish_in(status);
325
326 goby::glog.is_debug2() && goby::glog << group(glog_group_) << "Successfully opened socket"
327 << std::endl;
328
329 // update to avoid thrashing on open success but read/write failure
330 decltype(next_open_attempt_) now(goby::time::SteadyClock::now());
331 next_open_attempt_ = now + backoff_interval_;
332 }
333 catch (const std::exception& e)
334 {
335 auto status = std::make_shared<protobuf::IOStatus>();
336 if (this->index() != -1)
337 status->set_index(this->index());
338
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);
344
345 goby::glog.is_warn() && goby::glog << group(glog_group_)
346 << "Failed to open/configure socket/serial_port: "
347 << error.ShortDebugString() << std::endl;
348
349 if (backoff_interval_ < max_backoff_interval_)
350 backoff_interval_ *= 2.0;
351
352 decltype(next_open_attempt_) now(goby::time::SteadyClock::now());
353 next_open_attempt_ = now + backoff_interval_;
354
355 goby::glog.is_warn() && goby::glog << group(glog_group_) << "Will retry in "
356 << backoff_interval_ / std::chrono::seconds(1)
357 << " seconds" << std::endl;
358 socket_.reset();
359 }
360}
361
362template <const goby::middleware::Group& line_in_group,
363 const goby::middleware::Group& line_out_group,
365 goby::middleware::io::PubSubLayer subscribe_layer, typename IOConfig, typename SocketType,
366 template <class> class ThreadType, bool use_indexed_groups>
367void goby::middleware::io::detail::IOThread<line_in_group, line_out_group, publish_layer,
368 subscribe_layer, IOConfig, SocketType, ThreadType,
369 use_indexed_groups>::loop()
370{
371 if (socket_ && socket_->is_open())
372 {
373 // run the io service (blocks until either we read something
374 // from the socket or a subscription is available
375 // as signaled from an empty handler in the incoming_mail_notify_thread)
376 io_.run_one();
377 }
378 else
379 {
380 decltype(next_open_attempt_) now(goby::time::SteadyClock::now());
381 if (now > next_open_attempt_)
382 {
383 try_open();
384 }
385 else
386 {
387 // poll in case we need to quit
388 io_.poll();
389 usleep(100000); // avoid pegging CPU while waiting to attempt reopening socket
390 }
391 }
392}
393
394template <const goby::middleware::Group& line_in_group,
395 const goby::middleware::Group& line_out_group,
397 goby::middleware::io::PubSubLayer subscribe_layer, typename IOConfig, typename SocketType,
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)
402{
403 auto status = std::make_shared<protobuf::IOStatus>();
404 if (this->index() != -1)
405 status->set_index(this->index());
406
407 status->set_state(protobuf::IO__CRITICAL_FAILURE);
408 goby::middleware::protobuf::IOError& error = *status->mutable_error();
410 error.set_text(ec.message());
411 this->publish_in(status);
412
413 goby::glog.is_warn() && goby::glog << group(glog_group_)
414 << "Failed to read from the socket/serial_port: "
415 << error.ShortDebugString() << std::endl;
416
417 socket_.reset();
418}
419
420template <const goby::middleware::Group& line_in_group,
421 const goby::middleware::Group& line_out_group,
423 goby::middleware::io::PubSubLayer subscribe_layer, typename IOConfig, typename SocketType,
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)
428{
429 auto status = std::make_shared<protobuf::IOStatus>();
430 if (this->index() != -1)
431 status->set_index(this->index());
432
433 status->set_state(protobuf::IO__CRITICAL_FAILURE);
434 goby::middleware::protobuf::IOError& error = *status->mutable_error();
436 error.set_text(ec.message());
437 this->publish_in(status);
438
439 goby::glog.is_warn() && goby::glog << group(glog_group_)
440 << "Failed to write to the socket/serial_port: "
441 << error.ShortDebugString() << std::endl;
442 socket_.reset();
443}
444
445#endif
simple exception class for goby applications
Definition exception.h:35
Class for grouping publications in the Goby middleware. Analogous to "topics" in ROS,...
Definition group.h:60
void write(std::shared_ptr< const goby::middleware::protobuf::IOData > io_msg)
void handle_read_success(std::size_t bytes_transferred, const std::string &bytes)
virtual void async_write(std::shared_ptr< const goby::middleware::protobuf::IOData > io_msg)=0
Starts an asynchronous write from data published.
bool socket_is_open()
Does the socket exist and is it open?
void handle_read_error(const boost::system::error_code &ec)
void handle_read_success(std::size_t bytes_transferred, std::shared_ptr< goby::middleware::protobuf::IOData > io_msg)
void handle_write_error(const boost::system::error_code &ec)
virtual void async_read()=0
Starts an asynchronous read on the socket.
IOThread(const IOConfig &config, int index, std::string glog_group="i/o")
Constructs the thread.
SocketType & mutable_socket()
Access the (mutable) socket (or serial_port) object.
virtual void open_socket()=0
Opens the newly created socket/serial_port.
friend void basic_async_write(IOThreadImplementation *this_thread, std::shared_ptr< const goby::middleware::protobuf::IOData > io_msg)
boost::asio::io_context & mutable_io()
void handle_write_success(std::size_t bytes_transferred)
static constexpr ErrorCode IO__INIT_FAILURE
Definition io.pb.h:1933
static constexpr ErrorCode IO__WRITE_FAILURE
Definition io.pb.h:1937
static constexpr ErrorCode IO__READ_FAILURE
Definition io.pb.h:1935
void add_group(const std::string &name, Colors::Color color=Colors::nocolor, const std::string &description="")
Add another group to the logger. A group provides related manipulator for categorizing log messages.
goby::util::logger::GroupSetter group(std::string n)
detail namespace with internal helper functions
Definition json.hpp:263
@ error
throw a parse_error exception in case of a tag
void basic_async_write(IOThreadImplementation *this_thread, std::shared_ptr< const goby::middleware::protobuf::IOData > io_msg)
ProtobufEndpoint endpoint_convert(const ASIOEndpoint &asio_ep)
std::string to_string(goby::middleware::protobuf::Layer layer)
Definition common.h:44
middleware::InterProcessForwarder< InnerTransporter, detail::InterProcessTag > InterProcessForwarder
constexpr T e
Definition constants.h:35
The global namespace for the Goby project.
util::FlexOstream glog
Access the Goby logger through this object.
uint32_t index(const std::array< int8_t, 256 > &rdata, char symbol)
Definition base.h:140
STL namespace.
std::chrono::time_point< SteadyClock > time_point
static time_point now() noexcept
Returns the current steady time unless SimulatorSettings::using_sim_time == true in which case a simu...
std::chrono::microseconds duration
Duration type.