-
Notifications
You must be signed in to change notification settings - Fork 91
/
Copy pathrmw_node.cpp
5918 lines (5367 loc) · 204 KB
/
rmw_node.cpp
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
// Copyright 2019 ADLINK Technology Limited.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
#include <cassert>
#include <cstring>
#include <mutex>
#include <unordered_map>
#include <unordered_set>
#include <algorithm>
#include <chrono>
#include <iomanip>
#include <map>
#include <set>
#include <functional>
#include <atomic>
#include <memory>
#include <vector>
#include <string>
#include <tuple>
#include <utility>
#include <regex>
#include <limits>
#include "rcutils/allocator.h"
#include "rcutils/env.h"
#include "rcutils/filesystem.h"
#include "rcutils/format_string.h"
#include "rcutils/logging_macros.h"
#include "rcutils/process.h"
#include "rcutils/strdup.h"
#include "rmw/allocators.h"
#include "rmw/convert_rcutils_ret_to_rmw_ret.h"
#include "rmw/discovery_options.h"
#include "rmw/error_handling.h"
#include "rmw/event.h"
#include "rmw/features.h"
#include "rmw/get_node_info_and_types.h"
#include "rmw/get_service_names_and_types.h"
#include "rmw/get_topic_names_and_types.h"
#include "rmw/event_callback_type.h"
#include "rmw/names_and_types.h"
#include "rmw/rmw.h"
#include "rmw/sanity_checks.h"
#include "rmw/validate_namespace.h"
#include "rmw/validate_node_name.h"
#include "fallthrough_macro.hpp"
#include "Serialization.hpp"
#include "rcpputils/scope_exit.hpp"
#include "rmw/impl/cpp/macros.hpp"
#include "rmw/impl/cpp/key_value.hpp"
#include "TypeSupport2.hpp"
#include "rmw_version_test.hpp"
#include "MessageTypeSupport.hpp"
#include "ServiceTypeSupport.hpp"
#include "rmw/get_topic_endpoint_info.h"
#include "rmw/incompatible_qos_events_statuses.h"
#include "rmw/topic_endpoint_info_array.h"
#include "rmw_dds_common/context.hpp"
#include "rmw_dds_common/graph_cache.hpp"
#include "rmw_dds_common/msg/participant_entities_info.hpp"
#include "rmw_dds_common/qos.hpp"
#include "rmw_dds_common/security.hpp"
#include "rosidl_runtime_c/type_hash.h"
#include "rosidl_typesupport_cpp/message_type_support.hpp"
#include "tracetools/tracetools.h"
#include "namespace_prefix.hpp"
#include "dds/dds.h"
#include "dds/ddsc/dds_data_allocator.h"
#include "dds/ddsc/dds_loan_api.h"
#include "serdes.hpp"
#include "serdata.hpp"
#include "demangle.hpp"
using namespace std::literals::chrono_literals;
/* Security must be enabled when compiling and requires cyclone to support QOS property lists */
#if DDS_HAS_SECURITY && DDS_HAS_PROPERTY_LIST_QOS
#define RMW_SUPPORT_SECURITY 1
#else
#define RMW_SUPPORT_SECURITY 0
#endif
#if !DDS_HAS_DDSI_SERTYPE
#define ddsi_sertype_unref(x) ddsi_sertopic_unref(x)
#endif
/* Set to > 0 for printing warnings to stderr for each messages that was taken more than this many
ms after writing */
#define REPORT_LATE_MESSAGES 0
/* Set to != 0 for periodically printing requests that have been blocked for more than 1s */
#define REPORT_BLOCKED_REQUESTS 0
#define RET_ERR_X(msg, code) do {RMW_SET_ERROR_MSG(msg); code;} while (0)
#define RET_NULL_X(var, code) do {if (!var) {RET_ERR_X(#var " is null", code);}} while (0)
#define RET_ALLOC_X(var, code) do {if (!var) {RET_ERR_X("failed to allocate " #var, code);} \
} while (0)
#define RET_ERR(msg) RET_ERR_X(msg, return RMW_RET_ERROR)
#define RET_NULL(var) RET_NULL_X(var, return RMW_RET_ERROR)
#define RET_ALLOC(var) RET_ALLOC_X(var, return RMW_RET_ERROR)
using rmw_dds_common::msg::ParticipantEntitiesInfo;
const char * const eclipse_cyclonedds_identifier = "rmw_cyclonedds_cpp";
const char * const eclipse_cyclonedds_serialization_format = "cdr";
/* instance handles are unsigned 64-bit integers carefully constructed to be as close to uniformly
distributed as possible for no other reason than making them near-perfect hash keys, hence we can
improve over the default hash function */
struct dds_instance_handle_hash
{
public:
std::size_t operator()(dds_instance_handle_t const & x) const noexcept
{
return static_cast<std::size_t>(x);
}
};
bool operator<(dds_builtintopic_guid_t const & a, dds_builtintopic_guid_t const & b)
{
return memcmp(&a, &b, sizeof(dds_builtintopic_guid_t)) < 0;
}
static rmw_ret_t discovery_thread_stop(rmw_dds_common::Context & context);
static bool dds_qos_to_rmw_qos(const dds_qos_t * dds_qos, rmw_qos_profile_t * qos_policies);
static rmw_publisher_t * create_publisher(
dds_entity_t dds_ppant, dds_entity_t dds_pub,
const rosidl_message_type_support_t * type_supports,
const char * topic_name, const rmw_qos_profile_t * qos_policies,
const rmw_publisher_options_t * publisher_options
);
static rmw_ret_t destroy_publisher(rmw_publisher_t * publisher);
static rmw_subscription_t * create_subscription(
dds_entity_t dds_ppant, dds_entity_t dds_pub,
const rosidl_message_type_support_t * type_supports,
const char * topic_name, const rmw_qos_profile_t * qos_policies,
const rmw_subscription_options_t * subscription_options
);
static rmw_ret_t destroy_subscription(rmw_subscription_t * subscription);
static rmw_guard_condition_t * create_guard_condition();
static rmw_ret_t destroy_guard_condition(rmw_guard_condition_t * gc);
struct CddsDomain;
struct CddsWaitset;
struct Cdds
{
std::mutex lock;
/* Map of domain id to per-domain state, used by create/destroy node */
std::mutex domains_lock;
std::map<dds_domainid_t, CddsDomain> domains;
/* special guard condition that gets attached to every waitset but that is never triggered:
this way, we can avoid Cyclone's behaviour of always returning immediately when no
entities are attached to a waitset */
dds_entity_t gc_for_empty_waitset;
/* set of waitsets protected by lock, used to invalidate all waitsets caches when an entity is
deleted */
std::unordered_set<CddsWaitset *> waitsets;
Cdds()
: gc_for_empty_waitset(0)
{}
};
/* Use construct-on-first-use for the global state rather than a plain global variable to
prevent its destructor from running prior to last use by some other component in the
system. E.g., some rclcpp tests (at the time of this commit) drop a guard condition in
a global destructor, but (at least) on Windows the Cyclone RMW global dtors run before
the global dtors of that test, resulting in rmw_destroy_guard_condition() attempting to
use the already destroyed "Cdds::waitsets".
The memory leak this causes is minor (an empty map of domains and an empty set of
waitsets) and by definition only once. The alternative of elimating it altogether or
tying its existence to init/shutdown is problematic because this state is used across
domains and contexts.
The only practical alternative I see is to extend Cyclone's run-time state (which is
managed correctly for these situations), but it is not Cyclone's responsibility to work
around C++ global destructor limitations. */
static Cdds & gcdds()
{
static Cdds * x = new Cdds();
return *x;
}
struct CddsEntity
{
dds_entity_t enth;
};
struct CddsDomain
{
/* This RMW implementation currently implements localhost-only by explicitly creating
domains with a configuration that consists of: (1) a hard-coded selection of
"localhost" as the network interface address; (2) followed by the contents of the
CYCLONEDDS_URI environment variable:
- the "localhost" hostname should resolve to 127.0.0.1 (or equivalent) for IPv4 and
to ::1 for IPv6, so we don't have to worry about which of IPv4 or IPv6 is used (as
would be the case with a numerical IP address), nor do we have to worry about the
name of the loopback interface;
- if the machine's configuration doesn't properly resolve "localhost", you can still
override via $CYCLONEDDS_URI.
The CddsDomain type is used to track which domains exist and how many nodes are in
it. Because the domain is instantiated with the first nodes created in that domain,
the other nodes must have the same localhost-only setting. (It bugs out if not.)
Everything resets automatically when the last node in the domain is deleted.
(It might be better still to for Cyclone to provide "loopback" or something similar
as a generic alias for a loopback interface ...)
There are a few issues with the current support for creating domains explicitly in
Cyclone, fixing those might relax alter or relax some of the above. */
rmw_discovery_options_t discovery_options;
uint32_t refcount;
/* handle of the domain entity */
dds_entity_t domain_handle;
/* Default constructor so operator[] can be safely be used to look one up */
CddsDomain()
: refcount(0), domain_handle(0)
{
discovery_options = rmw_get_zero_initialized_discovery_options();
}
~CddsDomain()
{}
};
// Definition of struct rmw_context_impl_s as declared in rmw/init.h
struct rmw_context_impl_s
{
rmw_dds_common::Context common;
dds_domainid_t domain_id;
dds_entity_t ppant;
rmw_gid_t ppant_gid;
/* handles for built-in topic readers */
dds_entity_t rd_participant;
dds_entity_t rd_subscription;
dds_entity_t rd_publication;
/* DDS publisher, subscriber used for ROS 2 publishers and subscriptions */
dds_entity_t dds_pub;
dds_entity_t dds_sub;
/* Participant reference count*/
size_t node_count{0};
std::mutex initialization_mutex;
/* Shutdown flag */
bool is_shutdown{false};
/* suffix for GUIDs to construct unique client/service ids
(protected by initialization_mutex) */
uint32_t client_service_id;
rmw_context_impl_s()
: common(), domain_id(UINT32_MAX), ppant(0), client_service_id(0)
{
/* destructor relies on these being initialized properly */
common.thread_is_running.store(false);
common.graph_guard_condition = nullptr;
common.pub = nullptr;
common.sub = nullptr;
}
// Initializes the participant, if it wasn't done already.
// node_count is increased
rmw_ret_t
init(rmw_init_options_t * options, size_t domain_id);
// Destroys the participant, when node_count reaches 0.
rmw_ret_t
fini();
~rmw_context_impl_s()
{
if (0u != this->node_count) {
RCUTILS_SAFE_FWRITE_TO_STDERR(
"Not all nodes were finished before finishing the context\n."
"Ensure `rcl_node_fini` is called for all nodes before `rcl_context_fini`,"
"to avoid leaking.\n");
}
}
private:
void
clean_up();
};
struct CddsNode
{
};
struct user_callback_data_t
{
std::mutex mutex;
rmw_event_callback_t callback {nullptr};
const void * user_data {nullptr};
size_t unread_count {0};
rmw_event_callback_t event_callback[DDS_STATUS_ID_MAX + 1] {nullptr};
const void * event_data[DDS_STATUS_ID_MAX + 1] {nullptr};
size_t event_unread_count[DDS_STATUS_ID_MAX + 1] {0};
};
struct CddsPublisher : CddsEntity
{
dds_instance_handle_t pubiid;
rmw_gid_t gid;
struct ddsi_sertype * sertype;
rosidl_message_type_support_t type_supports;
dds_data_allocator_t data_allocator;
uint32_t sample_size;
bool is_loaning_available;
user_callback_data_t user_callback_data;
};
struct CddsSubscription : CddsEntity
{
rmw_gid_t gid;
dds_entity_t rdcondh;
rosidl_message_type_support_t type_supports;
dds_data_allocator_t data_allocator;
bool is_loaning_available;
user_callback_data_t user_callback_data;
};
struct client_service_id_t
{
uint8_t data[RMW_GID_STORAGE_SIZE];
};
struct CddsCS
{
std::unique_ptr<CddsPublisher> pub;
std::unique_ptr<CddsSubscription> sub;
client_service_id_t id;
};
struct CddsClient
{
CddsCS client;
#if REPORT_BLOCKED_REQUESTS
std::mutex lock;
dds_time_t lastcheck;
std::map<int64_t, dds_time_t> reqtime;
#endif
user_callback_data_t user_callback_data;
};
struct CddsService
{
CddsCS service;
user_callback_data_t user_callback_data;
};
struct CddsGuardCondition
{
dds_entity_t gcondh;
};
struct CddsEvent : CddsEntity
{
rmw_event_type_t event_type;
};
struct CddsWaitset
{
dds_entity_t waitseth;
std::vector<dds_attach_t> trigs;
size_t nelems;
std::mutex lock;
bool inuse;
std::vector<CddsSubscription *> subs;
std::vector<CddsGuardCondition *> gcs;
std::vector<CddsClient *> cls;
std::vector<CddsService *> srvs;
std::vector<CddsEvent> evs;
};
static void clean_waitset_caches();
#if REPORT_BLOCKED_REQUESTS
static void check_for_blocked_requests(CddsClient & client);
#endif
#ifndef WIN32
/* TODO(allenh1): check for Clang */
#pragma GCC visibility push (default)
#endif
extern "C" const char * rmw_get_implementation_identifier()
{
return eclipse_cyclonedds_identifier;
}
extern "C" const char * rmw_get_serialization_format()
{
return eclipse_cyclonedds_serialization_format;
}
extern "C" rmw_ret_t rmw_set_log_severity(rmw_log_severity_t severity)
{
uint32_t mask = 0;
switch (severity) {
default:
RMW_SET_ERROR_MSG_WITH_FORMAT_STRING("%s: Invalid log severity '%d'", __func__, severity);
return RMW_RET_INVALID_ARGUMENT;
case RMW_LOG_SEVERITY_DEBUG:
mask |= DDS_LC_DISCOVERY | DDS_LC_THROTTLE | DDS_LC_CONFIG;
FALLTHROUGH;
case RMW_LOG_SEVERITY_INFO:
mask |= DDS_LC_INFO;
FALLTHROUGH;
case RMW_LOG_SEVERITY_WARN:
mask |= DDS_LC_WARNING;
FALLTHROUGH;
case RMW_LOG_SEVERITY_ERROR:
mask |= DDS_LC_ERROR;
FALLTHROUGH;
case RMW_LOG_SEVERITY_FATAL:
mask |= DDS_LC_FATAL;
}
dds_set_log_mask(mask);
return RMW_RET_OK;
}
static void dds_listener_callback(dds_entity_t entity, void * arg)
{
// Not currently used
(void)entity;
auto data = static_cast<user_callback_data_t *>(arg);
std::lock_guard<std::mutex> guard(data->mutex);
if (data->callback) {
data->callback(data->user_data, 1);
} else {
data->unread_count++;
}
}
#define MAKE_DDS_EVENT_CALLBACK_FN(event_type, EVENT_TYPE) \
static void on_ ## event_type ## _fn( \
dds_entity_t entity, \
const dds_ ## event_type ## _status_t status, \
void * arg) \
{ \
(void)status; \
(void)entity; \
auto data = static_cast<user_callback_data_t *>(arg); \
std::lock_guard<std::mutex> guard(data->mutex); \
auto cb = data->event_callback[DDS_ ## EVENT_TYPE ## _STATUS_ID]; \
if (cb) { \
cb(data->event_data[DDS_ ## EVENT_TYPE ## _STATUS_ID], 1); \
} else { \
data->event_unread_count[DDS_ ## EVENT_TYPE ## _STATUS_ID]++; \
} \
}
// Define event callback functions
MAKE_DDS_EVENT_CALLBACK_FN(requested_deadline_missed, REQUESTED_DEADLINE_MISSED)
MAKE_DDS_EVENT_CALLBACK_FN(liveliness_lost, LIVELINESS_LOST)
MAKE_DDS_EVENT_CALLBACK_FN(offered_deadline_missed, OFFERED_DEADLINE_MISSED)
MAKE_DDS_EVENT_CALLBACK_FN(requested_incompatible_qos, REQUESTED_INCOMPATIBLE_QOS)
MAKE_DDS_EVENT_CALLBACK_FN(sample_lost, SAMPLE_LOST)
MAKE_DDS_EVENT_CALLBACK_FN(offered_incompatible_qos, OFFERED_INCOMPATIBLE_QOS)
MAKE_DDS_EVENT_CALLBACK_FN(liveliness_changed, LIVELINESS_CHANGED)
MAKE_DDS_EVENT_CALLBACK_FN(inconsistent_topic, INCONSISTENT_TOPIC)
MAKE_DDS_EVENT_CALLBACK_FN(subscription_matched, SUBSCRIPTION_MATCHED)
MAKE_DDS_EVENT_CALLBACK_FN(publication_matched, PUBLICATION_MATCHED)
static void listener_set_event_callbacks(dds_listener_t * l, void * arg)
{
dds_lset_requested_deadline_missed_arg(l, on_requested_deadline_missed_fn, arg, false);
dds_lset_requested_incompatible_qos_arg(l, on_requested_incompatible_qos_fn, arg, false);
dds_lset_sample_lost_arg(l, on_sample_lost_fn, arg, false);
dds_lset_liveliness_lost_arg(l, on_liveliness_lost_fn, arg, false);
dds_lset_offered_deadline_missed_arg(l, on_offered_deadline_missed_fn, arg, false);
dds_lset_offered_incompatible_qos_arg(l, on_offered_incompatible_qos_fn, arg, false);
dds_lset_liveliness_changed_arg(l, on_liveliness_changed_fn, arg, false);
dds_lset_inconsistent_topic_arg(l, on_inconsistent_topic_fn, arg, false);
dds_lset_subscription_matched_arg(l, on_subscription_matched_fn, arg, false);
dds_lset_publication_matched_arg(l, on_publication_matched_fn, arg, false);
}
static bool get_readwrite_qos(dds_entity_t handle, rmw_qos_profile_t * rmw_qos_policies)
{
dds_qos_t * qos = dds_create_qos();
dds_return_t ret = false;
if (dds_get_qos(handle, qos) < 0) {
RMW_SET_ERROR_MSG("get_readwrite_qos: invalid handle");
} else {
ret = dds_qos_to_rmw_qos(qos, rmw_qos_policies);
}
dds_delete_qos(qos);
return ret;
}
extern "C" rmw_ret_t rmw_subscription_set_on_new_message_callback(
rmw_subscription_t * rmw_subscription,
rmw_event_callback_t callback,
const void * user_data)
{
RMW_CHECK_ARGUMENT_FOR_NULL(rmw_subscription, RMW_RET_INVALID_ARGUMENT);
auto sub = static_cast<CddsSubscription *>(rmw_subscription->data);
user_callback_data_t * data = &(sub->user_callback_data);
std::lock_guard<std::mutex> guard(data->mutex);
// Set the user callback data
data->callback = callback;
data->user_data = user_data;
if (callback && data->unread_count) {
// Push events happened before having assigned a callback,
// limiting them to the QoS depth.
rmw_qos_profile_t sub_qos;
if (!get_readwrite_qos(sub->enth, &sub_qos)) {
return RMW_RET_ERROR;
}
size_t events = std::min(data->unread_count, sub_qos.depth);
callback(user_data, events);
data->unread_count = 0;
}
return RMW_RET_OK;
}
extern "C" rmw_ret_t rmw_service_set_on_new_request_callback(
rmw_service_t * rmw_service,
rmw_event_callback_t callback,
const void * user_data)
{
RMW_CHECK_ARGUMENT_FOR_NULL(rmw_service, RMW_RET_INVALID_ARGUMENT);
auto srv = static_cast<CddsService *>(rmw_service->data);
user_callback_data_t * data = &(srv->user_callback_data);
std::lock_guard<std::mutex> guard(data->mutex);
// Set the user callback data
data->callback = callback;
data->user_data = user_data;
if (callback && data->unread_count) {
// Push events happened before having assigned a callback
callback(user_data, data->unread_count);
data->unread_count = 0;
}
return RMW_RET_OK;
}
extern "C" rmw_ret_t rmw_client_set_on_new_response_callback(
rmw_client_t * rmw_client,
rmw_event_callback_t callback,
const void * user_data)
{
RMW_CHECK_ARGUMENT_FOR_NULL(rmw_client, RMW_RET_INVALID_ARGUMENT);
auto cli = static_cast<CddsClient *>(rmw_client->data);
user_callback_data_t * data = &(cli->user_callback_data);
std::lock_guard<std::mutex> guard(data->mutex);
// Set the user callback data
data->callback = callback;
data->user_data = user_data;
if (callback && data->unread_count) {
// Push events happened before having assigned a callback
callback(user_data, data->unread_count);
data->unread_count = 0;
}
return RMW_RET_OK;
}
template<typename T>
static void event_set_callback(
T event,
dds_status_id_t status_id,
rmw_event_callback_t callback,
const void * user_data)
{
user_callback_data_t * data = &(event->user_callback_data);
std::lock_guard<std::mutex> guard(data->mutex);
// Set the user callback data
data->event_callback[status_id] = callback;
data->event_data[status_id] = user_data;
if (callback && data->event_unread_count[status_id]) {
// Push events happened before having assigned a callback
callback(user_data, data->event_unread_count[status_id]);
data->event_unread_count[status_id] = 0;
}
}
extern "C" rmw_ret_t rmw_event_set_callback(
rmw_event_t * rmw_event,
rmw_event_callback_t callback,
const void * user_data)
{
RMW_CHECK_ARGUMENT_FOR_NULL(rmw_event, RMW_RET_INVALID_ARGUMENT);
switch (rmw_event->event_type) {
case RMW_EVENT_LIVELINESS_CHANGED:
{
auto sub_event = static_cast<CddsSubscription *>(rmw_event->data);
event_set_callback(
sub_event, DDS_LIVELINESS_CHANGED_STATUS_ID,
callback, user_data);
break;
}
case RMW_EVENT_REQUESTED_DEADLINE_MISSED:
{
auto sub_event = static_cast<CddsSubscription *>(rmw_event->data);
event_set_callback(
sub_event, DDS_REQUESTED_DEADLINE_MISSED_STATUS_ID,
callback, user_data);
break;
}
case RMW_EVENT_REQUESTED_QOS_INCOMPATIBLE:
{
auto sub_event = static_cast<CddsSubscription *>(rmw_event->data);
event_set_callback(
sub_event, DDS_REQUESTED_INCOMPATIBLE_QOS_STATUS_ID,
callback, user_data);
break;
}
case RMW_EVENT_MESSAGE_LOST:
{
auto sub_event = static_cast<CddsSubscription *>(rmw_event->data);
event_set_callback(
sub_event, DDS_SAMPLE_LOST_STATUS_ID,
callback, user_data);
break;
}
case RMW_EVENT_SUBSCRIPTION_MATCHED:
{
auto sub_event = static_cast<CddsSubscription *>(rmw_event->data);
event_set_callback(
sub_event, DDS_SUBSCRIPTION_MATCHED_STATUS_ID,
callback, user_data);
break;
}
case RMW_EVENT_LIVELINESS_LOST:
{
auto pub_event = static_cast<CddsPublisher *>(rmw_event->data);
event_set_callback(
pub_event, DDS_LIVELINESS_LOST_STATUS_ID,
callback, user_data);
break;
}
case RMW_EVENT_OFFERED_DEADLINE_MISSED:
{
auto pub_event = static_cast<CddsPublisher *>(rmw_event->data);
event_set_callback(
pub_event, DDS_OFFERED_DEADLINE_MISSED_STATUS_ID,
callback, user_data);
break;
}
case RMW_EVENT_OFFERED_QOS_INCOMPATIBLE:
{
auto pub_event = static_cast<CddsPublisher *>(rmw_event->data);
event_set_callback(
pub_event, DDS_OFFERED_INCOMPATIBLE_QOS_STATUS_ID,
callback, user_data);
break;
}
case RMW_EVENT_PUBLISHER_INCOMPATIBLE_TYPE:
{
auto pub_event = static_cast<CddsPublisher *>(rmw_event->data);
event_set_callback(
pub_event, DDS_INCONSISTENT_TOPIC_STATUS_ID,
callback, user_data);
break;
}
case RMW_EVENT_SUBSCRIPTION_INCOMPATIBLE_TYPE:
{
auto sub_event = static_cast<CddsSubscription *>(rmw_event->data);
event_set_callback(
sub_event, DDS_INCONSISTENT_TOPIC_STATUS_ID,
callback, user_data);
break;
}
case RMW_EVENT_PUBLICATION_MATCHED:
{
auto pub_event = static_cast<CddsPublisher *>(rmw_event->data);
event_set_callback(
pub_event, DDS_PUBLICATION_MATCHED_STATUS_ID,
callback, user_data);
break;
}
case RMW_EVENT_INVALID:
{
return RMW_RET_INVALID_ARGUMENT;
}
}
return RMW_RET_OK;
}
extern "C" rmw_ret_t rmw_init_options_init(
rmw_init_options_t * init_options,
rcutils_allocator_t allocator)
{
RMW_CHECK_ARGUMENT_FOR_NULL(init_options, RMW_RET_INVALID_ARGUMENT);
RCUTILS_CHECK_ALLOCATOR(&allocator, return RMW_RET_INVALID_ARGUMENT);
if (NULL != init_options->implementation_identifier) {
RMW_SET_ERROR_MSG("expected zero-initialized init_options");
return RMW_RET_INVALID_ARGUMENT;
}
init_options->instance_id = 0;
init_options->implementation_identifier = eclipse_cyclonedds_identifier;
init_options->allocator = allocator;
init_options->impl = nullptr;
init_options->discovery_options = rmw_get_zero_initialized_discovery_options(),
init_options->domain_id = RMW_DEFAULT_DOMAIN_ID;
init_options->enclave = NULL;
init_options->security_options = rmw_get_zero_initialized_security_options();
return rmw_discovery_options_init(&(init_options->discovery_options), 0, &allocator);
}
extern "C" rmw_ret_t rmw_init_options_copy(const rmw_init_options_t * src, rmw_init_options_t * dst)
{
RMW_CHECK_ARGUMENT_FOR_NULL(src, RMW_RET_INVALID_ARGUMENT);
RMW_CHECK_ARGUMENT_FOR_NULL(dst, RMW_RET_INVALID_ARGUMENT);
if (NULL == src->implementation_identifier) {
RMW_SET_ERROR_MSG("expected initialized src");
return RMW_RET_INVALID_ARGUMENT;
}
RMW_CHECK_TYPE_IDENTIFIERS_MATCH(
init options copy, src->implementation_identifier,
eclipse_cyclonedds_identifier, return RMW_RET_INCORRECT_RMW_IMPLEMENTATION);
if (NULL != dst->implementation_identifier) {
RMW_SET_ERROR_MSG("expected zero-initialized dst");
return RMW_RET_INVALID_ARGUMENT;
}
const rcutils_allocator_t * allocator = &src->allocator;
rmw_init_options_t tmp = *src;
tmp.enclave = rcutils_strdup(tmp.enclave, *allocator);
if (NULL != src->enclave && NULL == tmp.enclave) {
return RMW_RET_BAD_ALLOC;
}
tmp.security_options = rmw_get_zero_initialized_security_options();
rmw_ret_t ret =
rmw_security_options_copy(&src->security_options, allocator, &tmp.security_options);
if (RMW_RET_OK != ret) {
allocator->deallocate(tmp.enclave, allocator->state);
return ret;
}
*dst = tmp;
return RMW_RET_OK;
}
extern "C" rmw_ret_t rmw_init_options_fini(rmw_init_options_t * init_options)
{
RMW_CHECK_ARGUMENT_FOR_NULL(init_options, RMW_RET_INVALID_ARGUMENT);
if (NULL == init_options->implementation_identifier) {
RMW_SET_ERROR_MSG("expected initialized init_options");
return RMW_RET_INVALID_ARGUMENT;
}
RMW_CHECK_TYPE_IDENTIFIERS_MATCH(
init options, init_options->implementation_identifier,
eclipse_cyclonedds_identifier, return RMW_RET_INCORRECT_RMW_IMPLEMENTATION);
rcutils_allocator_t * allocator = &init_options->allocator;
RCUTILS_CHECK_ALLOCATOR(allocator, return RMW_RET_INVALID_ARGUMENT);
allocator->deallocate(init_options->enclave, allocator->state);
rmw_ret_t ret = rmw_security_options_fini(&init_options->security_options, allocator);
*init_options = rmw_get_zero_initialized_init_options();
return ret;
}
static void convert_guid_to_gid(const dds_guid_t & guid, rmw_gid_t & gid)
{
static_assert(
RMW_GID_STORAGE_SIZE >= sizeof(guid),
"rmw_gid_t type too small for a Cyclone DDS GUID");
memset(&gid, 0, sizeof(gid));
gid.implementation_identifier = eclipse_cyclonedds_identifier;
memcpy(gid.data, guid.v, sizeof(guid));
}
static void get_entity_gid(dds_entity_t h, rmw_gid_t & gid)
{
dds_guid_t guid;
dds_get_guid(h, &guid);
convert_guid_to_gid(guid, gid);
}
static std::map<std::string, std::vector<uint8_t>> parse_user_data(const dds_qos_t * qos)
{
std::map<std::string, std::vector<uint8_t>> map;
void * ud;
size_t udsz;
if (dds_qget_userdata(qos, &ud, &udsz)) {
std::vector<uint8_t> udvec(static_cast<uint8_t *>(ud), static_cast<uint8_t *>(ud) + udsz);
dds_free(ud);
map = rmw::impl::cpp::parse_key_value(udvec);
}
return map;
}
static bool get_user_data_key(const dds_qos_t * qos, const std::string key, std::string & value)
{
if (qos != nullptr) {
auto map = parse_user_data(qos);
auto name_found = map.find(key);
if (name_found != map.end()) {
value = std::string(name_found->second.begin(), name_found->second.end());
return true;
}
}
return false;
}
static void handle_ParticipantEntitiesInfo(dds_entity_t reader, void * arg)
{
static_cast<void>(reader);
rmw_context_impl_t * impl = static_cast<rmw_context_impl_t *>(arg);
ParticipantEntitiesInfo msg;
bool taken;
while (rmw_take(impl->common.sub, &msg, &taken, nullptr) == RMW_RET_OK && taken) {
// locally published data is filtered because of the subscription QoS
impl->common.graph_cache.update_participant_entities(msg);
}
}
static void handle_DCPSParticipant(dds_entity_t reader, void * arg)
{
rmw_context_impl_t * impl = static_cast<rmw_context_impl_t *>(arg);
dds_sample_info_t si;
void * raw = NULL;
while (dds_take(reader, &raw, &si, 1, 1) == 1) {
auto s = static_cast<const dds_builtintopic_participant_t *>(raw);
rmw_gid_t gid;
convert_guid_to_gid(s->key, gid);
if (memcmp(&gid, &impl->common.gid, sizeof(gid)) == 0) {
// ignore the local participant
} else if (si.instance_state != DDS_ALIVE_INSTANCE_STATE) {
impl->common.graph_cache.remove_participant(gid);
} else if (si.valid_data) {
std::string enclave;
if (get_user_data_key(s->qos, "enclave", enclave)) {
impl->common.graph_cache.add_participant(gid, enclave);
}
}
dds_return_loan(reader, &raw, 1);
}
}
static void handle_builtintopic_endpoint(
dds_entity_t reader, rmw_context_impl_t * impl,
bool is_reader)
{
dds_sample_info_t si;
void * raw = NULL;
while (dds_take(reader, &raw, &si, 1, 1) == 1) {
auto s = static_cast<const dds_builtintopic_endpoint_t *>(raw);
rmw_gid_t gid;
convert_guid_to_gid(s->key, gid);
if (si.instance_state != DDS_ALIVE_INSTANCE_STATE) {
impl->common.graph_cache.remove_entity(gid, is_reader);
} else if (si.valid_data && strncmp(s->topic_name, "DCPS", 4) != 0) {
rmw_qos_profile_t qos_profile = rmw_qos_profile_unknown;
rmw_gid_t ppgid;
dds_qos_to_rmw_qos(s->qos, &qos_profile);
convert_guid_to_gid(s->participant_key, ppgid);
rosidl_type_hash_t type_hash = rosidl_get_zero_initialized_type_hash();
void * userdata;
size_t userdata_size;
if (dds_qget_userdata(s->qos, &userdata, &userdata_size)) {
RCPPUTILS_SCOPE_EXIT(dds_free(userdata));
if (RMW_RET_OK != rmw_dds_common::parse_type_hash_from_user_data(
reinterpret_cast<const uint8_t *>(userdata), userdata_size, type_hash))
{
RCUTILS_LOG_WARN_NAMED(
"rmw_cyclonedds_cpp",
"Failed to parse type hash for topic '%s' with type '%s' from USER_DATA '%*s'.",
s->topic_name, s->type_name,
static_cast<int>(userdata_size), reinterpret_cast<char *>(userdata));
type_hash = rosidl_get_zero_initialized_type_hash();
// We've handled the error, so clear it out.
rmw_reset_error();
}
}
impl->common.graph_cache.add_entity(
gid,
std::string(s->topic_name),
std::string(s->type_name),
type_hash,
ppgid,
qos_profile,
is_reader);
}
dds_return_loan(reader, &raw, 1);
}
}
static void handle_DCPSSubscription(dds_entity_t reader, void * arg)
{
rmw_context_impl_t * impl = static_cast<rmw_context_impl_t *>(arg);
handle_builtintopic_endpoint(reader, impl, true);
}
static void handle_DCPSPublication(dds_entity_t reader, void * arg)
{
rmw_context_impl_t * impl = static_cast<rmw_context_impl_t *>(arg);
handle_builtintopic_endpoint(reader, impl, false);
}
static void discovery_thread(rmw_context_impl_t * impl)
{
const CddsSubscription * sub = static_cast<const CddsSubscription *>(impl->common.sub->data);
const CddsGuardCondition * gc =
static_cast<const CddsGuardCondition *>(impl->common.listener_thread_gc->data);
dds_entity_t ws;
/* deleting ppant will delete waitset as well, so there is no real need to delete
the waitset here on error, but it is more hygienic */
if ((ws = dds_create_waitset(DDS_CYCLONEDDS_HANDLE)) < 0) {
RCUTILS_SAFE_FWRITE_TO_STDERR(
"ros discovery info listener thread: failed to create waitset, will shutdown ...\n");
return;
}
/* I suppose I could attach lambda functions one way or another, which would
definitely be more elegant, but this avoids having to deal with the C++
freakishness that is involved and works, too. */
std::vector<std::pair<dds_entity_t,
std::function<void(dds_entity_t, rmw_context_impl_t *)>>> entries = {
{gc->gcondh, nullptr},
{sub->enth, handle_ParticipantEntitiesInfo},
{impl->rd_participant, handle_DCPSParticipant},
{impl->rd_subscription, handle_DCPSSubscription},
{impl->rd_publication, handle_DCPSPublication},
};
for (size_t i = 0; i < entries.size(); i++) {
if (entries[i].second != nullptr &&
dds_set_status_mask(entries[i].first, DDS_DATA_AVAILABLE_STATUS) < 0)
{
RCUTILS_SAFE_FWRITE_TO_STDERR(
"ros discovery info listener thread: failed to set reader status masks, "
"will shutdown ...\n");