Goby3 3.6.1
2026.09.15
Loading...
Searching...
No Matches
intervehicle.h
Go to the documentation of this file.
1// Copyright 2016-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_TRANSPORT_INTERVEHICLE_H
25#define GOBY_MIDDLEWARE_TRANSPORT_INTERVEHICLE_H
26
27#include <atomic>
28#include <functional>
29#include <map>
30#include <set>
31#include <sys/types.h>
32#include <thread>
33#include <unistd.h>
34
35#include <google/protobuf/io/zero_copy_stream_impl.h>
36
38
40#include "goby/middleware/transport/interthread.h" // used for InterVehiclePortal implementation
44#include "goby/time/convert.h"
46
47namespace goby
48{
49namespace middleware
50{
52{
53 public:
54 InvalidSubscription(const std::string& e) : Exception(e) {}
55};
56
58{
59 public:
60 InvalidPublication(const std::string& e) : Exception(e) {}
61};
62
64{
65 public:
66 InvalidUnsubscription(const std::string& e) : Exception(e) {}
67};
68
73template <typename Derived, typename InnerTransporter>
75 : public StaticTransporterInterface<InterVehicleTransporterBase<Derived, InnerTransporter>,
76 InnerTransporter>,
77 public Poller<InterVehicleTransporterBase<Derived, InnerTransporter>>
78{
79 using InterfaceType =
81 InnerTransporter>;
82
84
85 public:
87 {
90 };
91
94 {
95 // handle request from Portal to omit or include metadata on future publications for a given data type
96 this->inner()
99 [this](const protobuf::SerializerMetadataRequest& request)
100 {
101 glog.is_debug3() && glog << "Received DCCL metadata request: "
102 << request.ShortDebugString() << std::endl;
103
104 switch (request.request())
105 {
107 // resend right away if we had stopped (e.g. the portal restarted)
108 if (omit_publish_metadata_.erase(request.key().type()))
109 last_publish_metadata_time_.erase(request.key().type());
110 break;
112 omit_publish_metadata_.insert(request.key().type());
113 break;
114 }
115 });
116 }
118
119 virtual ~InterVehicleTransporterBase() = default;
120
122 template <typename Data> static constexpr int scheme()
123 {
124 static_assert(goby::middleware::scheme<typename detail::primitive_type<Data>::type>() ==
126 "Can only use DCCL messages with InterVehicleTransporters");
128 }
129
133 template <const Group& group> void check_validity()
134 {
135 static_assert(group.numeric() != Group::invalid_numeric_group,
136 "goby::middleware::Group must have non-zero numeric "
137 "value to publish on the InterVehicle layer");
138 }
139
147 template <typename Data, int scheme = goby::middleware::scheme<Data>()>
148 void publish_dynamic(const Data& data, const Group& group = Group(),
149 const Publisher<Data>& publisher = Publisher<Data>())
150 {
151 static_assert(scheme == MarshallingScheme::DCCL,
152 "Can only use DCCL messages with InterVehicleTransporters");
153
154 Data data_with_group = data;
155 publisher.set_group(data_with_group, group);
156
157 static_cast<Derived*>(this)->template _publish<Data>(data_with_group, group, publisher);
158 // publish to interprocess as both DCCL and Protobuf
159 this->inner().template publish_dynamic<Data, MarshallingScheme::DCCL>(data_with_group,
160 group, publisher);
161 this->inner().template publish_dynamic<Data, MarshallingScheme::PROTOBUF>(data_with_group,
162 group, publisher);
163 }
164
172 template <typename Data, int scheme = goby::middleware::scheme<Data>()>
173 void publish_dynamic(std::shared_ptr<const Data> data, const Group& group = Group(),
174 const Publisher<Data>& publisher = Publisher<Data>())
175 {
176 static_assert(scheme == MarshallingScheme::DCCL,
177 "Can only use DCCL messages with InterVehicleTransporters");
178 if (data)
179 {
180 // copy this way as it allows us to copy Data == google::protobuf::Message abstract base class
181 std::shared_ptr<Data> data_with_group(data->New());
182 data_with_group->CopyFrom(*data);
183
184 publisher.set_group(*data_with_group, group);
185
186 static_cast<Derived*>(this)->template _publish<Data>(*data_with_group, group,
187 publisher);
188
189 // publish to interprocess as both DCCL and Protobuf
190 this->inner().template publish_dynamic<Data, MarshallingScheme::DCCL>(data_with_group,
191 group, publisher);
192 this->inner().template publish_dynamic<Data, MarshallingScheme::PROTOBUF>(
193 data_with_group, group, publisher);
194 }
195 }
196
204 template <typename Data, int scheme = goby::middleware::scheme<Data>()>
205 void publish_dynamic(std::shared_ptr<Data> data, const Group& group = Group(),
206 const Publisher<Data>& publisher = Publisher<Data>())
207 {
208 publish_dynamic<Data, scheme>(std::shared_ptr<const Data>(data), group, publisher);
209 }
210
218 template <typename Data, int scheme = goby::middleware::scheme<Data>()>
219 void subscribe_dynamic(std::function<void(const Data&)> f, const Group& group = Group(),
220 const Subscriber<Data>& subscriber = Subscriber<Data>())
221 {
222 static_assert(scheme == MarshallingScheme::DCCL,
223 "Can only use DCCL messages with InterVehicleTransporters");
224 auto pointer_ref_lambda = [=](std::shared_ptr<const Data> d) { f(*d); };
225 static_cast<Derived*>(this)->template _subscribe<Data>(
226 pointer_ref_lambda, group, subscriber, SubscriptionAction::SUBSCRIBE);
227 }
228
236 template <typename Data, int scheme = goby::middleware::scheme<Data>()>
237 void subscribe_dynamic(std::function<void(std::shared_ptr<const Data>)> f,
238 const Group& group = Group(),
239 const Subscriber<Data>& subscriber = Subscriber<Data>())
240 {
241 static_assert(scheme == MarshallingScheme::DCCL,
242 "Can only use DCCL messages with InterVehicleTransporters");
243 static_cast<Derived*>(this)->template _subscribe<Data>(f, group, subscriber,
245 }
246
253 template <typename Data, int scheme = goby::middleware::scheme<Data>()>
255 const Subscriber<Data>& subscriber = Subscriber<Data>())
256 {
257 static_assert(scheme == MarshallingScheme::DCCL,
258 "Can only use DCCL messages with InterVehicleTransporters");
259 static_cast<Derived*>(this)->template _subscribe<Data>(
260 std::function<void(std::shared_ptr<const Data>)>(), group, subscriber,
262 }
263
264 protected:
265 template <typename Data>
266 std::shared_ptr<goby::middleware::protobuf::SerializerTransporterMessage>
267 _set_up_publish(const Data& d, const Group& group, const Publisher<Data>& publisher)
268 {
269 if (group.numeric() != Group::broadcast_group && !publisher.has_set_group_func())
270 {
271 std::stringstream ss;
272 ss << "Error: Publisher must have set_group_func in order to publish to a "
273 "non-broadcast Group ("
274 << group
275 << "). The set_group_func modifies the contents of the outgoing message to store "
276 "the group information.";
277 throw(InvalidPublication(ss.str()));
278 }
279
280 auto data = intervehicle::serialize_publication(d, group, publisher);
281
282 if (publisher.cfg().intervehicle().buffer().ack_required())
283 {
284 auto ack_handler = std::make_shared<
286 publisher.acked_func(), d);
287
288 auto expire_handler =
289 std::make_shared<PublisherCallback<Data, MarshallingScheme::DCCL,
291 publisher.expired_func(), d);
292
294 data, ack_handler, expire_handler);
295 }
296
297 const auto& type = data->key().type();
298 if (!omit_publish_metadata_.count(type))
299 {
300 auto now = goby::time::SteadyClock::now();
301 auto interval = goby::time::convert_duration<goby::time::SteadyClock::duration>(
303 auto last_it = last_publish_metadata_time_.find(type);
304 if (last_it == last_publish_metadata_time_.end() || now >= last_it->second + interval)
305 {
306 _set_protobuf_metadata<Data>(data->mutable_key()->mutable_metadata(), d);
307 last_publish_metadata_time_[type] = now;
308 }
309 }
310
312 goby::glog << "Set up publishing for: " << data->ShortDebugString() << std::endl;
313
314 return data;
315 }
316
317 template <typename Data>
318 std::shared_ptr<intervehicle::protobuf::Subscription>
319 _set_up_subscribe(std::function<void(std::shared_ptr<const Data> d)> func, const Group& group,
320 const Subscriber<Data>& subscriber, SubscriptionAction action)
321 {
323
324 switch (action)
325 {
327 {
328 if (group.numeric() != Group::broadcast_group && !subscriber.has_group_func())
329 {
330 std::stringstream ss;
331 ss << "Error: Subscriber must have group_func in order to subscribe to "
332 "non-broadcast Group ("
333 << group
334 << "). The group_func returns the appropriate Group based on the contents "
335 "of the incoming message.";
336 throw(InvalidSubscription(ss.str()));
337 }
338
339 if (subscriber.cfg().intervehicle().broadcast() &&
340 subscriber.cfg().intervehicle().buffer().ack_required())
341 {
342 std::stringstream ss;
343 ss << "Error: Broadcast subscriptions cannot have ack_required: true";
344 throw(InvalidSubscription(ss.str()));
345 }
346
347 auto subscription = std::make_shared<
349 func, group, subscriber);
350
351 this->subscriptions_[dccl_id][group] = subscription;
352 }
353 break;
355 {
356 auto sub_it = this->subscriptions_[dccl_id].find(group);
357 if (sub_it != this->subscriptions_[dccl_id].end())
358 {
359 this->subscriptions_[dccl_id].erase(sub_it);
360 }
361 else
362 {
363 std::stringstream ss;
364 ss << "Cannot unsubscribe to DCCL id: " << dccl_id
365 << " and group: " << std::string(group) << " as no subscription was found.";
366 throw(InvalidUnsubscription(ss.str()));
367 }
368 }
369 break;
370 }
371
372 auto dccl_subscription =
373 this->template _serialize_subscription<Data>(group, subscriber, action);
375 // insert pending subscription
376 auto subscription_publication = intervehicle::serialize_publication(
379
380 // overwrite timestamps to ensure mapping with driver threads
381 auto subscribe_time = dccl_subscription->time_with_units();
382 subscription_publication->mutable_key()->set_serialize_time_with_units(subscribe_time);
383
384 auto ack_handler = std::make_shared<PublisherCallback<Subscription, MarshallingScheme::DCCL,
386 subscriber.subscribed_func());
387
388 auto expire_handler =
389 std::make_shared<PublisherCallback<Subscription, MarshallingScheme::DCCL,
391 subscriber.subscribe_expired_func());
392
393 goby::glog.is_debug1() && goby::glog << "Inserting subscription ack handler for "
394 << subscription_publication->ShortDebugString()
395 << std::endl;
396
397 this->pending_ack_.insert(std::make_pair(*subscription_publication,
398 std::make_tuple(ack_handler, expire_handler)));
399
400 return dccl_subscription;
401 }
402
403 template <int tuple_index, typename AckorExpirePair>
404 void _handle_ack_or_expire(const AckorExpirePair& ack_or_expire_pair)
405 {
406 auto original = ack_or_expire_pair.serializer();
407 const auto& ack_or_expire_msg = ack_or_expire_pair.data();
408 bool is_subscription = original.key().marshalling_scheme() == MarshallingScheme::DCCL &&
409 original.key().type() ==
411
412 if (is_subscription)
413 {
414 // rewrite data to remove src()
415 auto bytes_begin = original.data().begin(), bytes_end = original.data().end();
416 decltype(bytes_begin) actual_end;
417
420 auto subscription = Helper::parse(bytes_begin, bytes_end, actual_end);
421 subscription->mutable_header()->set_src(0);
422
423 std::vector<char> bytes(Helper::serialize(*subscription));
424 std::string* sbytes = new std::string(bytes.begin(), bytes.end());
425 original.set_allocated_data(sbytes);
426 }
427
428 auto it = pending_ack_.find(original);
429 if (it != pending_ack_.end())
430 {
431 goby::glog.is_debug3() && goby::glog << ack_or_expire_msg.GetDescriptor()->name()
432 << " for: " << original.ShortDebugString() << ", "
433 << ack_or_expire_msg.ShortDebugString()
434 << std::endl;
435
436 std::get<tuple_index>(it->second)
437 ->post(original.data().begin(), original.data().end(), ack_or_expire_msg);
438 }
439 else
440 {
441 goby::glog.is_debug3() && goby::glog << "No pending Ack/Expire for "
442 << (is_subscription ? "subscription: " : "data: ")
443 << original.ShortDebugString() << std::endl;
444 }
445 }
446
448 {
449 goby::glog.is_debug3() && goby::glog << "Received DCCLForwarded data: "
450 << packets.ShortDebugString() << std::endl;
451
452 for (const auto& packet : packets.frame())
453 {
454 for (auto p : this->subscriptions_[packet.dccl_id()])
455 p.second->post(packet.data().begin(), packet.data().end(), packets.header());
456 }
457 }
458
459 template <typename Data>
460 std::shared_ptr<intervehicle::protobuf::Subscription>
462 SubscriptionAction action)
463 {
465 auto dccl_subscription = std::make_shared<intervehicle::protobuf::Subscription>();
466 dccl_subscription->mutable_header()->set_src(0);
467
468 for (auto id : subscriber.cfg().intervehicle().publisher_id())
469 dccl_subscription->mutable_header()->add_dest(id);
470
471 dccl_subscription->set_api_version(GOBY_INTERVEHICLE_API_VERSION);
472 dccl_subscription->set_dccl_id(dccl_id);
473 dccl_subscription->set_group(group.numeric());
474 dccl_subscription->set_time_with_units(
475 goby::time::SystemClock::now<goby::time::MicroTime>());
476 dccl_subscription->set_action((action == SubscriptionAction::SUBSCRIBE)
479
480 _set_protobuf_metadata<Data>(dccl_subscription->mutable_metadata());
481 *dccl_subscription->mutable_intervehicle() = subscriber.cfg().intervehicle();
482 return dccl_subscription;
483 }
484
486 int dccl_id, std::shared_ptr<goby::middleware::protobuf::SerializerTransporterMessage> data,
489 expire_handler)
490 {
491 goby::glog.is_debug3() && goby::glog << "Inserting ack handler for "
492 << data->ShortDebugString() << std::endl;
493
494 this->pending_ack_.insert(
495 std::make_pair(*data, std::make_tuple(ack_handler, expire_handler)));
496 }
497
498 protected:
499 // maps DCCL ID to map of Group->subscription
500 // only one subscription allowed per IntervehicleForwarder/Portal (new subscription overwrites old one)
501 std::unordered_map<
502 int, std::unordered_map<std::string, std::shared_ptr<const SerializationHandlerBase<
505
506 private:
507 friend PollerType;
508 int _poll(std::unique_ptr<std::unique_lock<std::mutex>>& lock)
509 {
510 _expire_pending_ack();
511
512 return static_cast<Derived*>(this)->_poll(lock);
513 }
514
515 template <typename Data> void _set_protobuf_metadata(protobuf::SerializerProtobufMetadata* meta)
516 {
518 _insert_file_desc_with_dependencies(Data::descriptor()->file(), meta);
519 }
520
521 template <typename Data>
522 void _set_protobuf_metadata(protobuf::SerializerProtobufMetadata* meta, const Data& d)
523 {
524 meta->set_protobuf_name(
526 _insert_file_desc_with_dependencies(d.GetDescriptor()->file(), meta);
527 }
528
529 // used to populated InterVehicleSubscription file_descriptor fields
530 void _insert_file_desc_with_dependencies(const google::protobuf::FileDescriptor* file_desc,
531 protobuf::SerializerProtobufMetadata* meta)
532 {
533 std::set<std::string> inserted;
534 _insert_file_desc_with_dependencies(file_desc, meta, inserted);
535 }
536
537 void _insert_file_desc_with_dependencies(const google::protobuf::FileDescriptor* file_desc,
538 protobuf::SerializerProtobufMetadata* meta,
539 std::set<std::string>& inserted)
540 {
541 if (!inserted.insert(std::string(file_desc->name())).second)
542 return;
543
544 for (int i = 0, n = file_desc->dependency_count(); i < n; ++i)
545 _insert_file_desc_with_dependencies(file_desc->dependency(i), meta, inserted);
546
547 google::protobuf::FileDescriptorProto* file_desc_proto = meta->add_file_descriptor();
548 file_desc->CopyTo(file_desc_proto);
549 }
550
551 // expire any pending_ack entries that are no longer relevant
552 void _expire_pending_ack()
553 {
554 auto now = goby::time::SystemClock::now<goby::time::MicroTime>();
555 for (auto it = pending_ack_.begin(), end = pending_ack_.end(); it != end;)
556 {
558 ->FindFieldByName("ttl")
559 ->options()
560 .GetExtension(dccl::field)
561 .max() *
563
564 decltype(now) serialize_time(it->first.key().serialize_time_with_units());
565 decltype(now) expire_time(serialize_time + max_ttl);
566
567 // time to let any expire messages from the drivers propagate through the interprocess layer before we remove this
568 const decltype(now) interprocess_wait(1.0 * boost::units::si::seconds);
569
570 // loop through pending ack, and clear any at the front that can be removed
571
572 if (now > expire_time + interprocess_wait)
573 {
574 goby::glog.is_debug3() && goby::glog << "Erasing pending ack for "
575 << it->first.ShortDebugString() << std::endl;
576 it = pending_ack_.erase(it);
577 }
578 else
579 {
580 // pending_ack_ is ordered by serialize time, so we can bail now
581 break;
582 }
583 }
584 }
585
586 private:
587 // maps data with ack_requested onto callbacks for when the data are acknowledged or expire
588 // ordered by serialize time
589 std::map<
590 protobuf::SerializerTransporterMessage,
591 std::tuple<std::shared_ptr<SerializationHandlerBase<intervehicle::protobuf::AckData>>,
592 std::shared_ptr<SerializationHandlerBase<intervehicle::protobuf::ExpireData>>>>
593 pending_ack_;
594
595 // map of Protobuf names where we can omit metadata on publication
596 std::set<std::string> omit_publish_metadata_;
597 std::map<std::string, goby::time::SteadyClock::time_point> last_publish_metadata_time_;
598};
599
604template <typename InnerTransporter>
606 : public InterVehicleTransporterBase<InterVehicleForwarder<InnerTransporter>, InnerTransporter>
607{
608 public:
609 using implementation_tag = typename InnerTransporter::implementation_tag;
610
611 using Base =
613
617 InterVehicleForwarder(InnerTransporter& inner) : Base(inner)
618 {
619 this->inner()
623 { this->_receive(msg); });
624
625 using ack_pair_type = intervehicle::protobuf::AckMessagePair;
626 this->inner().template subscribe<intervehicle::groups::modem_ack_in, ack_pair_type>(
627 [this](const ack_pair_type& ack_pair)
628 { this->template _handle_ack_or_expire<0>(ack_pair); });
629
630 using expire_pair_type = intervehicle::protobuf::ExpireMessagePair;
631 this->inner().template subscribe<intervehicle::groups::modem_expire_in, expire_pair_type>(
632 [this](const expire_pair_type& expire_pair)
633 { this->template _handle_ack_or_expire<1>(expire_pair); });
634 }
635
636 virtual ~InterVehicleForwarder() = default;
637
638 friend Base;
639
640 private:
641 template <typename Data>
642 void _publish(const Data& d, const Group& group, const Publisher<Data>& publisher)
643 {
644 this->inner().template publish<intervehicle::groups::modem_data_out>(
645 this->_set_up_publish(d, group, publisher));
646 }
647
648 template <typename Data>
649 void _subscribe(std::function<void(std::shared_ptr<const Data> d)> func, const Group& group,
650 const Subscriber<Data>& subscriber, typename Base::SubscriptionAction action)
651 {
652 try
653 {
654 this->inner()
658 this->_set_up_subscribe(func, group, subscriber, action));
659 }
660 catch (const InvalidUnsubscription& e)
661 {
662 goby::glog.is_warn() && goby::glog << e.what() << std::endl;
663 }
664 }
665
666 int _poll(std::unique_ptr<std::unique_lock<std::mutex>>& lock) { return 0; }
667};
668
672template <typename InnerTransporter>
674 : public InterVehicleTransporterBase<InterVehiclePortal<InnerTransporter>, InnerTransporter>
675{
676 // Derive the ImplementationTag from InnerTransporter so that the modem driver thread
677 // uses the same InterProcessForwarder prefix as the portal's inner transporter.
678 using implementation_tag = typename InnerTransporter::implementation_tag;
679 using modem_id_type = typename goby::middleware::intervehicle::ModemDriverThread<
680 implementation_tag>::modem_id_type;
681
682 public:
683 using Base =
685
689 InterVehiclePortal(const intervehicle::protobuf::PortalConfig& cfg) : cfg_(cfg) { _init(); }
690
696 : Base(inner), cfg_(cfg)
697 {
698 _init();
699 }
700
702 {
703 for (auto& modem_driver_data : modem_drivers_)
704 {
705 modem_driver_data->driver_thread_alive = false;
706 if (modem_driver_data->underlying_thread)
707 modem_driver_data->underlying_thread->join();
708 }
709 }
710
711 friend Base;
712
713 private:
714 template <typename Data>
715 void _publish(const Data& d, const Group& group, const Publisher<Data>& publisher)
716 {
717 this->innermost().template publish<intervehicle::groups::modem_data_out>(
718 this->_set_up_publish(d, group, publisher));
719 }
720
721 template <typename Data>
722 void _subscribe(std::function<void(std::shared_ptr<const Data> d)> func, const Group& group,
723 const Subscriber<Data>& subscriber, typename Base::SubscriptionAction action)
724 {
725 try
726 {
727 auto dccl_subscription = this->_set_up_subscribe(func, group, subscriber, action);
728
729 this->innermost().template publish<intervehicle::groups::modem_subscription_forward_tx>(
730 dccl_subscription);
731 }
732 catch (const InvalidUnsubscription& e)
733 {
734 goby::glog.is_warn() && goby::glog << e.what() << std::endl;
735 }
736 }
737
738 int _poll(std::unique_ptr<std::unique_lock<std::mutex>>& lock)
739 {
740 int items = 0;
742 while (!received_.empty())
743 {
744 this->_receive(received_.front());
745 received_.pop_front();
746 ++items;
747 if (lock)
748 lock.reset();
749 }
750 return items;
751 }
752
753 void _init()
754 {
755 // set up reception of forwarded (via acoustic) subscriptions,
756 // then re-publish to driver threads
757 {
758 using intervehicle::protobuf::Subscription;
759 auto subscribe_lambda = [=](std::shared_ptr<const Subscription> d)
760 {
761 this->innermost()
763 intervehicle::protobuf::Subscription,
765 };
766 auto subscription = std::make_shared<
767 IntervehicleSerializationSubscription<Subscription, MarshallingScheme::DCCL>>(
768 subscribe_lambda);
769
770 auto dccl_id = SerializerParserHelper<Subscription, MarshallingScheme::DCCL>::id();
771 this->subscriptions_[dccl_id].insert(
772 std::make_pair(subscription->subscribed_group(), subscription));
773 }
774
775 this->innermost().template subscribe<intervehicle::groups::modem_data_in>(
776 [this](const intervehicle::protobuf::DCCLForwardedData& msg)
777 { received_.push_back(msg); });
778
779 // a message requiring ack can be disposed by either [1] ack, [2] expire (TTL exceeded), [3] having no subscribers, [4] queue size exceeded.
780 // post the correct callback (ack for [1] and expire for [2-4])
781 // and remove the pending ack message
782 using ack_pair_type = intervehicle::protobuf::AckMessagePair;
783 this->innermost().template subscribe<intervehicle::groups::modem_ack_in, ack_pair_type>(
784 [this](const ack_pair_type& ack_pair)
785 { this->template _handle_ack_or_expire<0>(ack_pair); });
786
787 using expire_pair_type = intervehicle::protobuf::ExpireMessagePair;
788 this->innermost()
789 .template subscribe<intervehicle::groups::modem_expire_in, expire_pair_type>(
790 [this](const expire_pair_type& expire_pair)
791 { this->template _handle_ack_or_expire<1>(expire_pair); });
792
793 this->innermost().template subscribe<intervehicle::groups::modem_driver_ready, bool>(
794 [this](const bool& ready)
795 {
796 goby::glog.is_debug1() && goby::glog << "Received driver ready" << std::endl;
797 ++drivers_ready_;
798 });
799
800 // set up before drivers ready to ensure we don't miss subscriptions
801 if (cfg_.has_persist_subscriptions())
802 _set_up_persistent_subscriptions();
803
804 for (int i = 0, n = cfg_.link_size(); i < n; ++i)
805 {
806 auto* link = cfg_.mutable_link(i);
807
808 link->mutable_driver()->set_modem_id(link->modem_id());
809 link->mutable_mac()->set_modem_id(link->modem_id());
810
811 modem_drivers_.emplace_back(new ModemDriverData);
812 ModemDriverData& data = *modem_drivers_.back();
813
814 data.underlying_thread.reset(new std::thread(
815 [&data, link]()
816 {
817 try
818 {
819 data.modem_driver_thread.reset(
820 new intervehicle::ModemDriverThread<implementation_tag>(*link));
821 data.modem_driver_thread->run(data.driver_thread_alive);
822 }
823 catch (std::exception& e)
824 {
826 goby::glog << "Modem driver thread had uncaught exception: " << e.what()
827 << std::endl;
828 throw;
829 }
830 }));
831
832 if (goby::glog.buf().is_gui())
833 // allows for visual grouping of each link in the NCurses gui
834 std::this_thread::sleep_for(std::chrono::milliseconds(250));
835 }
836
837 while (drivers_ready_ < modem_drivers_.size())
838 {
839 goby::glog.is_debug1() && goby::glog << "Waiting for drivers to be ready." << std::endl;
840 this->poll();
841 std::this_thread::sleep_for(std::chrono::seconds(1));
842 }
843
844 // write subscriptions after drivers ready to ensure they aren't missed
845 if (former_sub_collection_.subscription_size() > 0)
846 {
848 goby::glog << "Begin loading subscriptions from persistent storage..." << std::endl;
849 for (const auto& sub : former_sub_collection_.subscription())
850 this->innermost()
851 .template publish<intervehicle::groups::modem_subscription_forward_rx,
852 intervehicle::protobuf::Subscription,
853 MarshallingScheme::PROTOBUF>(sub);
854 }
855 }
856
857 void _set_up_persistent_subscriptions()
858 {
859 const auto& dir = cfg_.persist_subscriptions().dir();
860 if (dir.empty())
861 goby::glog.is_die() && goby::glog << "persist_subscriptions.dir cannot be empty"
862 << std::endl;
863
864 std::stringstream file_name;
865 file_name << dir;
866 if (dir.back() != '/')
867 file_name << "/";
868 file_name << "goby_intervehicle_subscriptions_" << cfg_.persist_subscriptions().name()
869 << ".pb.txt";
870 persist_sub_file_name_ = file_name.str();
871 {
872 std::ifstream persist_sub_ifs(persist_sub_file_name_.c_str());
873 try
874 {
875 if (persist_sub_ifs.is_open())
876 {
877 google::protobuf::TextFormat::Parser parser;
878 google::protobuf::io::IstreamInputStream iis(&persist_sub_ifs);
879 parser.Parse(&iis, &former_sub_collection_);
880 }
881 else
882 {
884 goby::glog << "Could not open persistent subscriptions file: "
885 << persist_sub_file_name_
886 << ". Assuming no persistent subscriptions exist" << std::endl;
887 }
888 }
889 catch (const std::exception& e)
890 {
892 goby::glog << "Error reading persistent subscriptions file: " << e.what()
893 << std::endl;
894 }
895 }
896
897 std::ofstream persist_sub_ofs(persist_sub_file_name_.c_str());
898 if (!persist_sub_ofs.is_open())
899 {
900 goby::glog.is_die() &&
901 goby::glog << "Could not open persistent subscriptions file for writing: "
902 << persist_sub_file_name_ << std::endl;
903 }
904 remove(persist_sub_file_name_.c_str());
905
906 this->innermost().template subscribe<intervehicle::groups::subscription_report>(
907 [this](const intervehicle::protobuf::SubscriptionReport& report)
908 {
909 goby::glog.is_debug1() && goby::glog << "Received subscription report: "
910 << report.ShortDebugString() << std::endl;
911 sub_reports_[report.link_modem_id()] = report;
912 std::ofstream persist_sub_ofs(persist_sub_file_name_.c_str());
913 intervehicle::protobuf::SubscriptionPersistCollection collection;
914 collection.set_time_with_units(
915 goby::time::SystemClock::now<goby::time::MicroTime>());
916 for (auto report_p : sub_reports_)
917 {
918 for (const auto& sub : report_p.second.subscription())
919 *collection.add_subscription() = sub;
920 }
921 google::protobuf::TextFormat::Printer printer;
922 google::protobuf::io::OstreamOutputStream oos(&persist_sub_ofs);
924 goby::glog << "Collection: " << collection.ShortDebugString() << std::endl;
925 printer.Print(collection, &oos);
926 });
927 }
928
929 private:
930 intervehicle::protobuf::PortalConfig cfg_;
931
932 struct ModemDriverData
933 {
934 std::unique_ptr<std::thread> underlying_thread;
935 std::unique_ptr<intervehicle::ModemDriverThread<implementation_tag>> modem_driver_thread;
936 std::atomic<bool> driver_thread_alive{true};
937 };
938 std::vector<std::unique_ptr<ModemDriverData>> modem_drivers_;
939 unsigned drivers_ready_{0};
940
941 std::deque<intervehicle::protobuf::DCCLForwardedData> received_;
942
943 intervehicle::protobuf::SubscriptionPersistCollection former_sub_collection_;
944 std::string persist_sub_file_name_;
945 std::map<modem_id_type, intervehicle::protobuf::SubscriptionReport> sub_reports_;
946};
947} // namespace middleware
948} // namespace goby
949
950#endif
simple exception class for goby applications
Definition exception.h:35
static const ::PROTOBUF_NAMESPACE_ID::Descriptor * descriptor()
Definition buffer.pb.h:110
boost::units::unit< ttl_dimension, boost::units::si::system > ttl_unit
Definition buffer.pb.h:302
Class for grouping publications in the Goby middleware. Analogous to "topics" in ROS,...
Definition group.h:60
static constexpr std::uint32_t broadcast_group
Special group number representing the broadcast group (used when no grouping is required for a given ...
Definition group.h:63
static constexpr std::uint32_t invalid_numeric_group
Special group number representing an invalid numeric group (unsuitable for intervehicle and outer lay...
Definition group.h:65
Implements the forwarder concept for the intervehicle layer.
InterVehicleForwarder(InnerTransporter &inner)
Construct a forwarder for the intervehicle layer.
typename InnerTransporter::implementation_tag implementation_tag
Implements a portal for the intervehicle layer based on Goby Acomms.
InterVehiclePortal(InnerTransporter &inner, const intervehicle::protobuf::PortalConfig &cfg)
Instantiate a portal with the given configuration and a reference to an external inner transporter.
InterVehiclePortal(const intervehicle::protobuf::PortalConfig &cfg)
Instantiate a portal with the given configuration (with the portal owning the inner transporter)
Base class for implementing transporters (both portal and forwarder) for the intervehicle layer.
std::shared_ptr< intervehicle::protobuf::Subscription > _serialize_subscription(const Group &group, const Subscriber< Data > &subscriber, SubscriptionAction action)
static constexpr int scheme()
returns the marshalling scheme id for a given data type on this layer. Only MarshallingScheme::DCCL i...
void publish_dynamic(std::shared_ptr< Data > data, const Group &group=Group(), const Publisher< Data > &publisher=Publisher< Data >())
Publish a message using a run-time defined DynamicGroup (shared pointer to mutable data variant)....
void subscribe_dynamic(std::function< void(std::shared_ptr< const Data >)> f, const Group &group=Group(), const Subscriber< Data > &subscriber=Subscriber< Data >())
Subscribe to a specific run-time defined group and data type (shared pointer variant)....
void subscribe_dynamic(std::function< void(const Data &)> f, const Group &group=Group(), const Subscriber< Data > &subscriber=Subscriber< Data >())
Subscribe to a specific run-time defined group and data type (const reference variant)....
void _handle_ack_or_expire(const AckorExpirePair &ack_or_expire_pair)
void publish_dynamic(std::shared_ptr< const Data > data, const Group &group=Group(), const Publisher< Data > &publisher=Publisher< Data >())
Publish a message using a run-time defined DynamicGroup (shared pointer to const data variant)....
InterVehicleTransporterBase(InnerTransporter &inner)
std::shared_ptr< intervehicle::protobuf::Subscription > _set_up_subscribe(std::function< void(std::shared_ptr< const Data > d)> func, const Group &group, const Subscriber< Data > &subscriber, SubscriptionAction action)
void _receive(const intervehicle::protobuf::DCCLForwardedData &packets)
void _insert_pending_ack(int dccl_id, std::shared_ptr< goby::middleware::protobuf::SerializerTransporterMessage > data, std::shared_ptr< SerializationHandlerBase< intervehicle::protobuf::AckData > > ack_handler, std::shared_ptr< SerializationHandlerBase< intervehicle::protobuf::ExpireData > > expire_handler)
void publish_dynamic(const Data &data, const Group &group=Group(), const Publisher< Data > &publisher=Publisher< Data >())
Publish a message using a run-time defined DynamicGroup (const reference variant)....
void unsubscribe_dynamic(const Group &group=Group(), const Subscriber< Data > &subscriber=Subscriber< Data >())
Unsubscribe from a specific run-time defined group and data type. Where possible, prefer the static v...
void check_validity()
Check validity of the Group for interthread use (at compile time)
std::unordered_map< int, std::unordered_map< std::string, std::shared_ptr< const SerializationHandlerBase< intervehicle::protobuf::Header > > > > subscriptions_
std::shared_ptr< goby::middleware::protobuf::SerializerTransporterMessage > _set_up_publish(const Data &d, const Group &group, const Publisher< Data > &publisher)
Represents a subscription to a serialized data type (intervehicle layer).
InvalidPublication(const std::string &e)
InvalidSubscription(const std::string &e)
InvalidUnsubscription(const std::string &e)
int poll(const std::chrono::time_point< Clock, Duration > &timeout=std::chrono::time_point< Clock, Duration >::max())
poll for data. Blocks until a data event occurs or a timeout when a particular time has been reached
Definition interface.h:375
Utility class for allowing the various Goby middleware transporters to poll the underlying transport ...
Definition poller.h:38
Represents a callback for a published data type (e.g. acked_func or expired_func)
Class that holds additional metadata and callback functions related to a publication (and is optional...
Definition publisher.h:40
const goby::middleware::protobuf::TransporterConfig & cfg() const
Returns the metadata configuration.
Definition publisher.h:81
acked_func_type acked_func() const
Returns the acked data callback (or an empty function if none is set)
Definition publisher.h:91
bool has_set_group_func() const
Definition publisher.h:95
expired_func_type expired_func() const
Returns the expired data callback (or an empty function if none is set)
Definition publisher.h:93
Base class for handling posting callbacks for serialized data types (interprocess and outer)
Defines the common interface for publishing and subscribing data using static (constexpr) groups on G...
Definition interface.h:234
void subscribe(std::function< void(const Data &)> f, const Subscriber< Data > &subscriber=Subscriber< Data >())
Subscribe to a specific group and data type (const reference variant)
Definition interface.h:300
void publish(const Data &data, const Publisher< Data > &publisher=Publisher< Data >())
Publish a message (const reference variant)
Definition interface.h:245
Class that holds additional metadata and callback functions related to a subscription (and is optiona...
Definition subscriber.h:37
subscribed_func_type subscribed_func() const
Definition subscriber.h:92
const goby::middleware::protobuf::TransporterConfig & cfg() const
Definition subscriber.h:80
subscribe_expired_func_type subscribe_expired_func() const
Definition subscriber.h:94
Provides the modem driver thread used by InterVehiclePortal, templated on ImplementationTag so it use...
const ::goby::middleware::intervehicle::protobuf::Header & header() const
const ::goby::middleware::intervehicle::protobuf::DCCLPacket & frame(int index) const
const ::goby::middleware::intervehicle::protobuf::PortalConfig_PersistSubscriptions & persist_subscriptions() const
::goby::middleware::intervehicle::protobuf::PortalConfig_LinkConfig * mutable_link(int index)
static const ::PROTOBUF_NAMESPACE_ID::Descriptor * descriptor()
const ::goby::acomms::protobuf::DynamicBufferConfig & buffer() const
void set_protobuf_name(ArgT0 &&arg0, ArgT... args)
const ::goby::middleware::intervehicle::protobuf::TransporterConfig & intervehicle() const
#define GOBY_INTERVEHICLE_API_VERSION
goby::util::logger::GroupSetter group(std::string n)
extern ::PROTOBUF_NAMESPACE_ID::internal::ExtensionIdentifier< ::goby::acomms::protobuf::ModemReport, ::PROTOBUF_NAMESPACE_ID::internal::MessageTypeTraits< ::goby::acomms::iridium::protobuf::Report >, 11, false > report
constexpr Group modem_subscription_forward_tx
Definition groups.h:48
constexpr Group modem_subscription_forward_rx
Definition groups.h:50
std::shared_ptr< goby::middleware::protobuf::SerializerTransporterMessage > serialize_publication(const Data &d, const Group &group, const Publisher< Data > &publisher)
constexpr T e
Definition constants.h:35
The global namespace for the Goby project.
extern ::PROTOBUF_NAMESPACE_ID::internal::ExtensionIdentifier< ::PROTOBUF_NAMESPACE_ID::MessageOptions, ::PROTOBUF_NAMESPACE_ID::internal::MessageTypeTraits< ::goby::GobyMessageOptions >, 11, false > msg
util::FlexOstream glog
Access the Goby logger through this object.
Class for parsing and serializing a given marshalling scheme. Must be specialized for a particular sc...
Definition interface.h:98
static std::string type_name()
The marshalling scheme specific string name for this type.
Definition interface.h:107
static time_point now() noexcept
Returns the current steady time unless SimulatorSettings::using_sim_time == true in which case a simu...