-
Notifications
You must be signed in to change notification settings - Fork 65
Expand file tree
/
Copy pathNodeSharedPrivate.hh
More file actions
460 lines (389 loc) · 17.3 KB
/
Copy pathNodeSharedPrivate.hh
File metadata and controls
460 lines (389 loc) · 17.3 KB
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
/*
* Copyright (C) 2017 Open Source Robotics Foundation
*
* 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.
*
*/
#ifndef GZ_TRANSPORT_NODESHAREDPRIVATE_HH_
#define GZ_TRANSPORT_NODESHAREDPRIVATE_HH_
#include <algorithm>
#include <atomic>
#include <cstdlib>
#include <deque>
#include <filesystem>
#include <list>
#include <map>
#include <memory>
#include <string>
#include <unordered_map>
#include <unordered_set>
#include <utility>
#include <vector>
#include <zmq.hpp>
#ifdef HAVE_ZENOH
#include <zenoh.hxx>
#endif
#include "gz/transport/config.hh"
#include "gz/transport/Exception.hh"
#include "gz/transport/Node.hh"
#include "Discovery.hh"
namespace gz::transport
{
// Inline bracket to help doxygen filtering.
inline namespace GZ_TRANSPORT_VERSION_NAMESPACE {
//
/// \brief Metadata for a publication. This is sent as part of the ZMQ
/// message for topic statistics.
class PublicationMetadata
{
/// \brief Publication timestamp.
public: uint64_t stamp = 0;
/// \brief Sequence number, used to detect dropped messages.
public: uint64_t seq = 0;
};
//
// Private data class for NodeShared.
class NodeSharedPrivate
{
/// \brief Constructor.
/// \throws gz::transport::Exception if a Zenoh session cannot be opened
/// (e.g. when using client mode without a reachable router).
public: NodeSharedPrivate()
{
// Determine implementation first, before creating any resources.
std::string gzImpl;
if (env("GZ_TRANSPORT_IMPLEMENTATION", gzImpl) && !gzImpl.empty())
{
std::transform(gzImpl.begin(), gzImpl.end(), gzImpl.begin(), ::tolower);
if (gzImpl == "zeromq" || gzImpl == "zenoh")
this->gzImplementation = gzImpl;
else
{
std::cerr << "Unrecognized value in GZ_TRANSPORT_IMPLEMENTATION. ["
<< gzImpl << "]. Ignoring this value" << std::endl;
}
}
// Create resources based on implementation.
if (this->gzImplementation == "zeromq")
{
this->context.reset(new zmq::context_t(1));
this->publisher.reset(new zmq::socket_t(*context, ZMQ_PUB));
this->subscriber.reset(new zmq::socket_t(*context, ZMQ_SUB));
this->requester.reset(new zmq::socket_t(*context, ZMQ_ROUTER));
this->responseReceiver.reset(new zmq::socket_t(*context, ZMQ_ROUTER));
this->replier.reset(new zmq::socket_t(*context, ZMQ_ROUTER));
}
#ifdef HAVE_ZENOH
else if (this->gzImplementation == "zenoh")
{
ZenohConfigSource configSource;
auto config = ZenohConfig(configSource);
if (this->verbose)
{
if (configSource == ZenohConfigSource::kFromEnvVariable)
std::cout << "Zenoh config loaded from ZENOH_CONFIG" << std::endl;
else
std::cout << "Zenoh default config loaded" << std::endl;
}
// Apply key=value overrides from GZ_TRANSPORT_ZENOH_CONFIG_OVERRIDE.
const char *overrideEnv =
std::getenv("GZ_TRANSPORT_ZENOH_CONFIG_OVERRIDE");
if (overrideEnv)
ApplyZenohConfigOverrides(config, overrideEnv, this->verbose);
try
{
this->session = std::make_shared<zenoh::Session>(
zenoh::Session::open(std::move(config)));
}
catch (const zenoh::ZException &e)
{
// Throw rather than continuing with a null session, which
// would cause segfaults downstream. This can happen when
// using client mode without a reachable router. Users can
// configure connect.timeout_ms in the Zenoh config to wait
// for the router to become available.
throw gz::transport::Exception(
std::string("Failed to open Zenoh session: ") + e.what());
}
}
#endif
}
/// \brief Initialize security
public: void SecurityInit();
/// \brief Handle new secure connections
public: void SecurityOnNewConnection();
/// \brief Access control handler for plain security.
/// This function is designed to be run in a thread.
public: void AccessControlHandler();
/// \brief Get and validate a non-negative environment variable.
/// \param[in] _envVar The name of the environment variable to get.
/// \param[in] _defaultValue The default value returned in case the
/// environment variable is invalid (e.g.: invalid number,
/// negative number).
/// \return The value read from the environment variable or the default
/// value if the validation wasn't succeed.
public: int NonNegativeEnvVar(const std::string &_envVar,
int _defaultValue) const;
#ifdef HAVE_ZENOH
/// \brief Enumeration for the source of Zenoh configuration.
public: enum class ZenohConfigSource
{
/// \brief Configuration loaded from ZENOH_CONFIG env variable.
kFromEnvVariable,
/// \brief Default Zenoh configuration.
kDefault
};
/// \brief Get the Zenoh configuration.
/// If the environment variable ZENOH_CONFIG is set, use that config file.
/// Otherwise, use the default Zenoh configuration.
/// \param[out] _configSource The source of the configuration.
/// \return The Zenoh configuration object.
public: inline zenoh::Config ZenohConfig(ZenohConfigSource &_configSource)
{
const char *zenohConfigEnv = std::getenv("ZENOH_CONFIG");
if (zenohConfigEnv)
{
std::string configPath(zenohConfigEnv);
// Check if the file exists and is a regular file.
if (std::filesystem::is_regular_file(configPath))
{
// Try to load the config file.
zenoh::ZResult result;
zenoh::Config config =
zenoh::Config::from_file(configPath, &result);
if (result == Z_OK)
{
_configSource = ZenohConfigSource::kFromEnvVariable;
return config;
}
std::cerr << "Failed to parse Zenoh config file: "
<< configPath << "\n";
}
else
{
std::cerr << "Zenoh config file not found: "
<< configPath << "\n";
}
}
// Fallback to default configuration.
_configSource = ZenohConfigSource::kDefault;
return zenoh::Config::create_default();
}
/// \brief Apply key=value config overrides to a Zenoh config.
/// Format: "key1=value1;key2=value2;..."
/// Values are passed directly to insert_json5(), so they can be any
/// valid JSON5 (strings, numbers, objects, arrays).
/// \param[in,out] _config The Zenoh config to modify.
/// \param[in] _overrides The override string to parse.
/// \param[in] _verbose Print applied overrides to stdout.
public: static inline void ApplyZenohConfigOverrides(
zenoh::Config &_config,
const std::string &_overrides,
bool _verbose = false)
{
std::string::size_type pos = 0;
while (pos < _overrides.size())
{
auto semi = _overrides.find(';', pos);
if (semi == std::string::npos)
semi = _overrides.size();
auto eq = _overrides.find('=', pos);
if (eq != std::string::npos && eq < semi)
{
std::string key = _overrides.substr(pos, eq - pos);
std::string val = _overrides.substr(eq + 1, semi - eq - 1);
// Trim leading/trailing whitespace.
auto trim = [](std::string &s)
{
auto start = s.find_first_not_of(" \t");
auto end = s.find_last_not_of(" \t");
s = (start == std::string::npos) ? "" :
s.substr(start, end - start + 1);
};
trim(key);
trim(val);
if (!key.empty())
{
zenoh::ZResult result;
_config.insert_json5(key, val, &result);
if (result != Z_OK)
{
std::cerr << "gz-transport: failed to apply Zenoh "
<< "config override: " << key << "="
<< val << std::endl;
}
else if (_verbose)
{
std::cout << "Zenoh config override: " << key
<< "=" << val << std::endl;
}
}
}
pos = semi + 1;
}
}
/// \brief Pointer to the Zenoh session.
public: std::shared_ptr<zenoh::Session> session;
#endif
//////////////////////////////////////////////////
/////// Declare here the ZMQ Context ///////
//////////////////////////////////////////////////
/// \brief 0MQ context. Always declare this object before any ZMQ socket
/// to make sure that the context is destroyed after all sockets.
public: std::unique_ptr<zmq::context_t> context;
//////////////////////////////////////////////////
/////// Declare here all ZMQ sockets ///////
//////////////////////////////////////////////////
/// \brief ZMQ socket to send topic updates.
public: std::unique_ptr<zmq::socket_t> publisher;
/// \brief ZMQ socket to receive topic updates.
public: std::unique_ptr<zmq::socket_t> subscriber;
/// \brief ZMQ socket for sending service call requests.
public: std::unique_ptr<zmq::socket_t> requester;
/// \brief ZMQ socket for receiving service call responses.
public: std::unique_ptr<zmq::socket_t> responseReceiver;
/// \brief ZMQ socket to receive service call requests.
public: std::unique_ptr<zmq::socket_t> replier;
/// \brief Thread the handle access control
public: std::thread accessControlThread;
//////////////////////////////////////////////////
/////// Declare here the discovery object ///////
//////////////////////////////////////////////////
/// \brief Discovery service (messages).
public: std::unique_ptr<MsgDiscovery> msgDiscovery;
/// \brief Discovery service (services).
public: std::unique_ptr<SrvDiscovery> srvDiscovery;
//////////////////////////////////////////////////
/////// Other private member variables ///////
//////////////////////////////////////////////////
/// \brief When true, the reception thread will finish.
public: std::atomic<bool> exit = false;
/// \brief Timeout used for receiving messages (ms.).
public: inline static const int Timeout = 250;
////////////////////////////////////////////////////////////////
/////// The following is for asynchronous publication of ///////
/////// messages to local subscribers. ///////
////////////////////////////////////////////////////////////////
/// \brief Encapsulates information needed to publish a message. An
/// instance of this class is pushed onto a publish queue, pubQueue, when
/// a message is published through Node::Publisher::Publish.
/// The pubThread processes the pubQueue in the
/// NodeSharedPrivate::PublishThread function.
///
/// A producer-consumer mechanism is used to send messages so that
/// Node::Publisher::Publish function does not block while executing
/// local subscriber callbacks.
public: struct PublishMsgDetails
{
/// \brief All the local subscription handlers.
public: std::vector<ISubscriptionHandlerPtr> localHandlers;
/// \brief All the raw handlers.
public: std::vector<RawSubscriptionHandlerPtr> rawHandlers;
/// \brief Buffer for the raw handlers.
public: std::unique_ptr<char[]> sharedBuffer = nullptr;
/// \brief Msg copy for the local handlers.
public: std::unique_ptr<ProtoMsg> msgCopy = nullptr;
/// \brief Message size.
// cppcheck-suppress unusedStructMember
public: std::size_t msgSize = 0;
/// \brief Information about the topic and type.
public: MessageInfo info;
/// \brief Publisher's node UUID.
public: std::string publisherNodeUUID;
/// \brief Fully qualified topic name. Used to enforce the
/// per topic capacity of the pubQueue.
public: std::string fullyQualifiedTopic;
};
/// \brief Type of the queue that stores messages pending delivery to
/// local subscribers.
public: using PubQueue = std::list<std::unique_ptr<PublishMsgDetails>>;
/// \brief Publish thread used to process the pubQueue.
public: std::thread pubThread;
/// \brief Mutex to protect the pubThread and pubQueue.
public: std::mutex pubThreadMutex;
/// \brief List onto which new messages are pushed. The pubThread
/// will pop off the messages and send them to local subscribers.
/// Use EnqueuePubMsg to push messages so that the localHwm capacity
/// is enforced.
public: PubQueue pubQueue;
/// \brief Iterators to the pubQueue entries of each topic, in queue
/// order. Used to track the number of queued messages per topic and to
/// drop the oldest message of a topic in constant time when the
/// localHwm capacity is reached. Protected by pubThreadMutex.
public: std::unordered_map<std::string, std::deque<PubQueue::iterator>>
pubQueueIters;
/// \brief Topics that have already dropped messages from the pubQueue.
/// Used to warn only once per overload episode: a topic is removed from
/// this set when its backlog fully drains from the pubQueue. Protected
/// by pubThreadMutex.
public: std::unordered_set<std::string> pubQueueDropWarned;
/// \brief Maximum number of messages that can be stored in the pubQueue
/// per topic. A value of 0 means no limit. Initialized from the
/// GZ_TRANSPORT_LOCAL_HWM environment variable before the pubThread
/// starts and never modified afterwards.
public: std::size_t localHwm = kDefaultLocalHwm;
/// \brief Add a message to the pubQueue so that it is asynchronously
/// delivered to local subscribers. If the topic already has localHwm
/// messages stored in the pubQueue, its oldest queued message is
/// dropped to make room for the new one.
/// \param[in] _msgDetails Details of the message to be published.
public: void EnqueuePubMsg(
std::unique_ptr<PublishMsgDetails> _msgDetails);
/// \brief used to signal when new work is available
public: std::condition_variable signalNewPub;
/// \brief Service thread used to process the srvQueue.
public: std::thread srvThread;
/// \brief Mutex to protect the srvThread and srvQueue.
public: std::mutex srvThreadMutex;
/// \brief List onto which new srv publishers are pushed.
public: std::list<ServicePublisher> srvQueue;
/// \brief Used to signal when new work is available.
public: std::condition_variable signalNewSrv;
/// \brief Handles local publication of messages on the pubQueue.
public: void PublishThread();
/// \brief Topic publication sequence numbers.
public: std::map<std::string, uint64_t> topicPubSeq;
/// \brief True if topic statistics have been enabled.
public: bool topicStatsEnabled = false;
/// \brief Statistics for a topic. The key in the map is the topic
/// name and the value contains the topic statistics.
public: std::map<std::string, TopicStatistics> topicStats;
/// \brief Set of topics that have statistics enabled.
public: std::map<std::string,
std::function<void(const TopicStatistics &_stats)>>
enabledTopicStatistics;
/// \brief A map of node UUID and its subscribed topics
public: std::unordered_map<std::string, std::unordered_set<std::string>>
topicsSubscribed;
/// \brief Service call repliers.
public: HandlerStorage<IRepHandler> repliers;
/// \brief Pending service call requests.
public: HandlerStorage<IReqHandler> requests;
/// \brief Print activity to stdout.
public: int verbose = false;
/// \brief My pub/sub address.
public: std::string myAddress;
/// \brief My requester service call address.
public: std::string myRequesterAddress;
/// \brief My replier service call address.
public: std::string myReplierAddress;
/// \brief IP address of this host.
public: std::string hostAddr;
/// \brief Underlying middleware implementation.
/// Supported values are: [zenoh, zeromq].
public: std::string gzImplementation =
std::string(GZ_TRANSPORT_DEFAULT_IMPLEMENTATION);
};
}
}
#endif