nros C++ API
Lightweight ROS 2 client for embedded real-time systems (C++ headers)
Loading...
Searching...
No Matches
subscription.hpp
Go to the documentation of this file.
1// nros-cpp: Subscription class
2// Freestanding C++ — no exceptions, no STL required
3
10#ifndef NROS_CPP_SUBSCRIPTION_HPP
11#define NROS_CPP_SUBSCRIPTION_HPP
12
13#include <cstdint>
14#include <cstddef>
15#include <string.h> // memcpy — `<cstring>` isn't in Zephyr's minimal libcpp
16
17#include "nros/config.hpp"
18#include "nros/result.hpp"
19#include "nros/stream.hpp"
20
21#include "nros_cpp_ffi.h"
22
23// Phase 189.M3.x — `nros_cpp_subscription_register` is excluded from cbindgen
24// (its Rust signature uses `RawSubscriptionCallback`, an external-crate type
25// alias cbindgen names without defining). Declare it locally with a plain
26// function-pointer typedef matching the ABI (`void(data, len, ctx)`), mirroring
27// the service.hpp callback-register treatment.
28extern "C" {
29typedef void (*nros_cpp_subscription_message_callback_t)(const uint8_t* data, size_t len,
30 void* ctx);
31
32// Phase 273 (RFC-0047): `callback_group` appended at end (NULL = default group).
33nros_cpp_ret_t nros_cpp_subscription_register(const nros_cpp_node_t* node, const char* topic,
34 const char* type_name, const char* type_hash,
37 void* context, uint8_t sched_context,
38 size_t* out_handle_id, const char* callback_group);
39
40// Phase 189.M3.4 — callback-style register that also delivers the sample's wire
41// attachment (5-arg trampoline). Same cbindgen-exclusion reason as above.
42typedef void (*nros_cpp_subscription_message_info_callback_t)(const uint8_t* data, size_t len,
43 const uint8_t* attachment,
44 size_t attachment_len, void* ctx);
45
47 const nros_cpp_node_t* node, const char* topic, const char* type_name, const char* type_hash,
49 uint8_t sched_context, size_t* out_handle_id);
50// Phase 269 W3 — callback-style subscription that delivers the sample's E2E
51// integrity status alongside the CDR bytes. Same cbindgen-exclusion reason as
52// above (takes `RawSubscriptionSafetyCallback`, an external nros-node alias
53// gated on the `safety-e2e` feature). The 6-arg trampoline unpacks the three
54// integrity scalars; `subscription.hpp` repacks them into
55// `nros_cpp_integrity_status_t` for the typed user handler. Gated on the
56// `NANO_ROS_SAFETY_E2E` build feature (lowered from `[system].features =
57// ["safety"]` via `NanoRosCapabilities.cmake`).
58#if defined(NANO_ROS_SAFETY_E2E)
59typedef void (*nros_cpp_subscription_validated_callback_t)(const uint8_t* data, size_t len,
60 int64_t gap, bool duplicate,
61 int8_t crc_valid, void* ctx);
62
63nros_cpp_ret_t nros_cpp_subscription_register_validated(
64 const nros_cpp_node_t* node, const char* topic, const char* type_name, const char* type_hash,
65 nros_cpp_qos_t qos, nros_cpp_subscription_validated_callback_t callback, void* context,
66 uint8_t sched_context, size_t* out_handle_id);
67#endif // NANO_ROS_SAFETY_E2E
68} // extern "C"
69
70namespace nros {
71
75static constexpr size_t SUBSCRIPTION_TOPIC_NAME_MAX = 256;
76
95template <typename M> class Subscription {
96 public:
100 using TypedSubscriptionFn = void (*)(const M& msg);
101 using TypedSubscriptionFnWithCtx = void (*)(const M& msg, void* ctx);
102 // Phase 189.M3.4 — callback-with-attachment handler (`bridge_origin` etc.).
104 size_t attachment_len);
105 // Phase 269 W3 — callback-with-integrity handler: receives the deserialized
106 // message plus the sample's E2E CRC/sequence status. Requires
107 // `NANO_ROS_SAFETY_E2E` (lowered from `[system].features = ["safety"]`).
108#if defined(NANO_ROS_SAFETY_E2E)
109 using TypedSubscriptionSafetyFn = void (*)(const M& msg,
111#endif // NANO_ROS_SAFETY_E2E
112
124 if (!initialized_) return Result(ErrorCode::NotInitialized);
125 uint8_t buf[M::SERIALIZED_SIZE_MAX];
126 size_t len = 0;
127 nros_cpp_ret_t ret = nros_cpp_subscription_try_recv_raw(storage_, buf, sizeof(buf), &len);
128 if (ret != 0) return Result(ret);
129 if (len == 0) return Result(ErrorCode::TryAgain);
130 if (M::ffi_deserialize(buf, len, &msg) != 0) return Result(ErrorCode::Error);
131 return Result::success();
132 }
133
150 if (!initialized_) return Result(ErrorCode::NotInitialized);
151 uint8_t buf[M::SERIALIZED_SIZE_MAX];
152 size_t len = 0;
154 nros_cpp_subscription_try_recv_validated(storage_, buf, sizeof(buf), &len, &status);
155 if (ret != 0) return Result(ret);
156 if (len == 0) return Result(ErrorCode::TryAgain);
157 if (M::ffi_deserialize(buf, len, &msg) != 0) return Result(ErrorCode::Error);
158 return Result::success();
159 }
160
171 Result try_recv_raw(uint8_t* buf, size_t capacity, size_t& out_len) {
172 if (!initialized_) {
173 out_len = 0;
175 }
177 if (ret != 0) return Result(ret);
178 if (out_len == 0) return Result(ErrorCode::TryAgain);
179 return Result::success();
180 }
181
198 uint8_t* att, size_t att_capacity, size_t& out_att_len) {
199 if (!initialized_) {
200 out_len = 0;
201 out_att_len = 0;
203 }
205 storage_, buf, capacity, &out_len, att, att_capacity, &out_att_len);
206 if (ret != 0) return Result(ret);
207 if (out_len == 0) return Result(ErrorCode::TryAgain);
208 return Result::success();
209 }
210
211 // ====================================================================
212 // Phase 124.A.7 — zero-copy receive (borrow / release)
213 // ====================================================================
214
217 class View {
218 public:
219 View() : sub_(nullptr), buf_(nullptr), len_(0), token_(nullptr) {}
220 View(View&& o) : sub_(o.sub_), buf_(o.buf_), len_(o.len_), token_(o.token_) {
221 o.sub_ = nullptr;
222 o.token_ = nullptr;
223 }
225 if (this != &o) {
226 release();
227 sub_ = o.sub_;
228 buf_ = o.buf_;
229 len_ = o.len_;
230 token_ = o.token_;
231 o.sub_ = nullptr;
232 o.token_ = nullptr;
233 }
234 return *this;
235 }
236 View(const View&) = delete;
237 View& operator=(const View&) = delete;
238 ~View() { release(); }
239
240 const uint8_t* data() const { return buf_; }
241 size_t size() const { return len_; }
242 bool empty() const { return token_ == nullptr; }
243
245 View(void* sub, const uint8_t* buf, size_t len, void* token)
246 : sub_(sub), buf_(buf), len_(len), token_(token) {}
247
248 private:
249 void release() {
250 if (token_ && sub_) {
251 nros_cpp_subscription_release(sub_, token_);
252 token_ = nullptr;
253 }
254 }
255
256 void* sub_;
257 const uint8_t* buf_;
258 size_t len_;
259 void* token_;
260 };
261
266 if (!initialized_) return Expected<View>::error(Result(ErrorCode::NotInitialized));
267 const uint8_t* buf = nullptr;
268 size_t len = 0;
269 void* token = nullptr;
270 int32_t rc = nros_cpp_subscription_borrow(storage_, &buf, &len, &token);
271 if (rc < 0) return Expected<View>::error(Result(rc));
272 if (rc == 0) return Expected<View>::ok(View{});
273 return Expected<View>::ok(View{storage_, buf, len, token});
274 }
275
289 size_t& out_count) {
290 if (!initialized_) {
291 out_count = 0;
293 }
295 storage_, buf, per_msg_cap, max_msgs, out_lens, &out_count);
296 if (ret != 0) return Result(ret);
297 return Result::success();
298 }
299
301 const char* get_topic_name() const { return initialized_ ? topic_name_ : ""; }
302
311 if (initialized_ && !stream_.is_valid()) {
312 stream_.bind(storage_, &nros_cpp_subscription_try_recv_raw);
313 }
314 return stream_;
315 }
316
317 const Stream<M>& stream() const { return stream_; }
318
320 bool is_valid() const { return initialized_; }
321
329 if (initialized_ && !callback_mode_) {
331 }
332 initialized_ = false;
333 }
334
335 // Move semantics (non-copyable). Relocation goes through the
336 // `nros_cpp_subscription_relocate` runtime call (Phase 84.C1);
337 // the `stream_` is rebound to the new storage afterwards.
338 // A callback-style subscription must NOT be moved after register — the
339 // executor arena holds `this` as the trampoline context (Phase 189.M3.x);
340 // the move only transfers bookkeeping and leaves that pointer stale. The
341 // poll-style relocation path is unchanged.
342 Subscription(Subscription&& other) : initialized_(other.initialized_) {
343 user_fn_ = other.user_fn_;
344 user_fn_ctx_ = other.user_fn_ctx_;
345 user_ctx_ = other.user_ctx_;
346 callback_mode_ = other.callback_mode_;
347 sched_handle_id_ = other.sched_handle_id_;
348 if (other.initialized_ && !other.callback_mode_) {
349 nros_cpp_subscription_relocate(other.storage_, storage_);
350 ::memcpy(topic_name_, other.topic_name_, sizeof(topic_name_));
351 stream_.bind(storage_, &nros_cpp_subscription_try_recv_raw);
352 }
353 other.initialized_ = false;
354 other.stream_ = Stream<M>();
355 }
356
358 if (this != &other) {
359 if (initialized_ && !callback_mode_) {
361 stream_ = Stream<M>();
362 }
363 initialized_ = other.initialized_;
364 user_fn_ = other.user_fn_;
365 user_fn_ctx_ = other.user_fn_ctx_;
366 user_ctx_ = other.user_ctx_;
367 callback_mode_ = other.callback_mode_;
368 sched_handle_id_ = other.sched_handle_id_;
369 if (other.initialized_ && !other.callback_mode_) {
370 nros_cpp_subscription_relocate(other.storage_, storage_);
371 ::memcpy(topic_name_, other.topic_name_, sizeof(topic_name_));
372 stream_.bind(storage_, &nros_cpp_subscription_try_recv_raw);
373 }
374 other.initialized_ = false;
375 other.stream_ = Stream<M>();
376 }
377 return *this;
378 }
379
382 Subscription() : storage_(), topic_name_{}, initialized_(false), stream_() {}
383
394 bool has_sched_handle() const { return sched_handle_id_ != static_cast<size_t>(-1); }
395 size_t sched_handle_id() const { return sched_handle_id_; }
396
397 // ====================================================================
398 // Phase 108 — status events
399 // ====================================================================
400
410
413 void* user_context = nullptr) {
414 if (!initialized_) return Result(ErrorCode::NotInitialized);
416 user_context));
417 }
418
424
425 private:
426 Subscription(const Subscription&) = delete;
427 Subscription& operator=(const Subscription&) = delete;
428
429 friend class Node;
430
434 static void message_trampoline(const uint8_t* data, size_t len, void* ctx) {
435 auto* self = static_cast<Subscription*>(ctx);
436 if (self == nullptr) return;
437 M msg;
438 if (M::ffi_deserialize(data, len, &msg) != 0) return;
439 if (self->user_fn_ != nullptr) {
440 self->user_fn_(msg);
441 } else if (self->user_fn_ctx_ != nullptr) {
442 self->user_fn_ctx_(msg, self->user_ctx_);
443 }
444 }
445
449 static void message_info_trampoline(const uint8_t* data, size_t len, const uint8_t* attachment,
450 size_t attachment_len, void* ctx) {
451 auto* self = static_cast<Subscription*>(ctx);
452 if (self == nullptr) return;
453 M msg;
454 if (M::ffi_deserialize(data, len, &msg) != 0) return;
455 if (self->user_fn_info_ != nullptr) {
456 self->user_fn_info_(msg, attachment, attachment_len);
457 }
458 }
459
460#if defined(NANO_ROS_SAFETY_E2E)
466 static void message_safety_trampoline(const uint8_t* data, size_t len, int64_t gap,
467 bool duplicate, int8_t crc_valid, void* ctx) {
468 auto* self = static_cast<Subscription*>(ctx);
469 if (self == nullptr) return;
470 M msg;
471 if (M::ffi_deserialize(data, len, &msg) != 0) return;
472 if (self->user_fn_safety_ != nullptr) {
473 nros_cpp_integrity_status_t status;
474 status.gap = gap;
475 status.duplicate = duplicate;
476 status.crc_valid = crc_valid;
477 self->user_fn_safety_(msg, status);
478 }
479 }
480#endif // NANO_ROS_SAFETY_E2E
481
482 alignas(8) uint8_t storage_[NROS_SUBSCRIBER_SIZE];
483 char topic_name_[SUBSCRIPTION_TOPIC_NAME_MAX];
484 bool initialized_;
485 Stream<M> stream_;
486 // Phase 189.M3.1 — executor HandleId for sched-context binding, or
487 // SIZE_MAX (the default) when no bindable handle exists. The poll-style
488 // thin-wrapper create path leaves this unset (see has_sched_handle());
489 // the callback-style create (Phase 189.M3.x) stores the real arena handle.
490 size_t sched_handle_id_ = static_cast<size_t>(-1);
491 // Callback-style state (Phase 189.M3.x); unused in poll mode. The executor
492 // arena owns the subscriber + dispatches `message_trampoline` during spin.
493 TypedSubscriptionFn user_fn_ = nullptr;
494 TypedSubscriptionFnWithCtx user_fn_ctx_ = nullptr;
495 TypedSubscriptionInfoFn user_fn_info_ = nullptr;
496 void* user_ctx_ = nullptr;
497 bool callback_mode_ = false;
498#if defined(NANO_ROS_SAFETY_E2E)
499 // Phase 269 W3 — handler for the integrity-carrying callback path; nullptr
500 // when not using `create_subscription_with_safety`.
501 TypedSubscriptionSafetyFn user_fn_safety_ = nullptr;
502#endif // NANO_ROS_SAFETY_E2E
503};
504
505} // namespace nros
506
507// Phase 84.G8: out-of-line definition of Node::create_subscription<M>().
508#include "nros/node.hpp"
509
510namespace nros {
511
512template <typename M>
514 if (!initialized_) return Result(ErrorCode::NotInitialized);
516 ffi_qos.reliability = static_cast<nros_cpp_qos_reliability_t>(qos.reliability_raw());
517 ffi_qos.durability = static_cast<nros_cpp_qos_durability_t>(qos.durability_raw());
518 ffi_qos.history = static_cast<nros_cpp_qos_history_t>(qos.history_raw());
519 ffi_qos.liveliness_kind = static_cast<nros_cpp_qos_liveliness_t>(qos.liveliness_raw());
520 ffi_qos.depth = qos.depth();
521 ffi_qos.deadline_ms = qos.deadline_ms();
522 ffi_qos.lifespan_ms = qos.lifespan_ms();
523 ffi_qos.liveliness_lease_ms = qos.liveliness_lease_ms();
524 ffi_qos.avoid_ros_namespace_conventions = qos.avoid_ros_namespace_conventions() ? 1 : 0;
525 ffi_qos.tx_express = qos.tx_express() ? 1 : 0;
526 nros_cpp_ret_t ret = nros_cpp_subscription_create(&handle_, topic, M::TYPE_NAME, M::TYPE_HASH,
527 ffi_qos, out.storage_);
528 if (ret == 0) {
529 // Phase 87.6: topic name lives C++-side now.
530 size_t topic_len = 0;
531 while (topic[topic_len] != '\0' && topic_len + 1 < sizeof(out.topic_name_)) {
532 out.topic_name_[topic_len] = topic[topic_len];
533 ++topic_len;
534 }
535 out.topic_name_[topic_len] = '\0';
536 out.initialized_ = true;
537 }
538 return Result(ret);
539}
540
553template <typename M>
557 if (!r.ok()) return r;
558
559 // TODO(M3.4): honour options.message_info via the with-info arena path.
560
561 if (options.sched_context != SCHED_CONTEXT_UNSET && out.has_sched_handle()) {
563 executor_handle_, out.sched_handle_id(), static_cast<uint8_t>(options.sched_context));
564 if (bind != 0) {
565 // Roll back so the caller doesn't observe a half-configured
566 // entity. Destructor-on-out would also fire, but explicit
567 // teardown keeps the returned error authoritative.
569 out.initialized_ = false;
570 return Result(bind);
571 }
572 }
573 return Result::success();
574}
575
576// Phase 189.M3.x — callback-style (arena-registered) subscription. The arena
577// owns the subscriber + dispatches `out`'s message handler during spin_once, so
578// the handle is real and `options.sched_context` is functional. Mirrors the
579// callback-style `create_service` one entity over.
580template <typename M, typename F, typename>
582 const QoS& qos, const SubscriptionOptions& options) {
583 if (!initialized_) return Result(ErrorCode::NotInitialized);
585 ffi_qos.reliability = static_cast<nros_cpp_qos_reliability_t>(qos.reliability_raw());
586 ffi_qos.durability = static_cast<nros_cpp_qos_durability_t>(qos.durability_raw());
587 ffi_qos.history = static_cast<nros_cpp_qos_history_t>(qos.history_raw());
588 ffi_qos.liveliness_kind = static_cast<nros_cpp_qos_liveliness_t>(qos.liveliness_raw());
589 ffi_qos.depth = qos.depth();
590 ffi_qos.deadline_ms = qos.deadline_ms();
591 ffi_qos.lifespan_ms = qos.lifespan_ms();
592 ffi_qos.liveliness_lease_ms = qos.liveliness_lease_ms();
593 ffi_qos.avoid_ros_namespace_conventions = qos.avoid_ros_namespace_conventions() ? 1 : 0;
594 ffi_qos.tx_express = qos.tx_express() ? 1 : 0;
595
596 // Store the user handler (compile error if F isn't convertible to the
597 // plain-fn-ptr handler type).
599 out.user_fn_ctx_ = nullptr;
600 out.user_ctx_ = nullptr;
601
602 uint8_t sched = (options.sched_context == SCHED_CONTEXT_UNSET)
603 ? 0u
604 : static_cast<uint8_t>(options.sched_context);
605 size_t handle = static_cast<size_t>(-1);
607 nros_cpp_subscription_register(&handle_, topic, M::TYPE_NAME, M::TYPE_HASH, ffi_qos,
608 &Subscription<M>::message_trampoline, &out, sched, &handle,
609 nullptr); // callback_group: nullptr = default group
610 if (ret == 0) {
611 out.sched_handle_id_ = handle;
612 out.callback_mode_ = true;
613 out.initialized_ = true;
614 }
615 return Result(ret);
616}
617
618// Phase 273 (RFC-0047) — callback-style subscription **in** a named callback group.
619// Mirrors create_subscription (callback-style) exactly but passes group.get_name()
620// as `callback_group` so the executor binds the slot via group_sched_table.
621template <typename M, typename F, typename>
623 const char* topic, F callback, const QoS& qos,
625 if (!initialized_) return Result(ErrorCode::NotInitialized);
627 ffi_qos.reliability = static_cast<nros_cpp_qos_reliability_t>(qos.reliability_raw());
628 ffi_qos.durability = static_cast<nros_cpp_qos_durability_t>(qos.durability_raw());
629 ffi_qos.history = static_cast<nros_cpp_qos_history_t>(qos.history_raw());
630 ffi_qos.liveliness_kind = static_cast<nros_cpp_qos_liveliness_t>(qos.liveliness_raw());
631 ffi_qos.depth = qos.depth();
632 ffi_qos.deadline_ms = qos.deadline_ms();
633 ffi_qos.lifespan_ms = qos.lifespan_ms();
634 ffi_qos.liveliness_lease_ms = qos.liveliness_lease_ms();
635 ffi_qos.avoid_ros_namespace_conventions = qos.avoid_ros_namespace_conventions() ? 1 : 0;
636 ffi_qos.tx_express = qos.tx_express() ? 1 : 0;
637
639 out.user_fn_ctx_ = nullptr;
640 out.user_ctx_ = nullptr;
641
642 uint8_t sched = (options.sched_context == SCHED_CONTEXT_UNSET)
643 ? 0u
644 : static_cast<uint8_t>(options.sched_context);
645 size_t handle = static_cast<size_t>(-1);
647 nros_cpp_subscription_register(&handle_, topic, M::TYPE_NAME, M::TYPE_HASH, ffi_qos,
648 &Subscription<M>::message_trampoline, &out, sched, &handle,
649 group.get_name()); // Phase 273: pass group name
650 if (ret == 0) {
651 out.sched_handle_id_ = handle;
652 out.callback_mode_ = true;
653 out.initialized_ = true;
654 }
655 return Result(ret);
656}
657
658// Phase 189.M3.4 — callback-style subscription that delivers the wire attachment.
659// Mirrors the callback `create_subscription` one step over, but stores the
660// `(const M&, attachment, att_len)` handler + registers via the with-info arena
661// path so the trampoline receives the attachment.
662template <typename M, typename F, typename>
664 const QoS& qos, const SubscriptionOptions& options) {
665 if (!initialized_) return Result(ErrorCode::NotInitialized);
667 ffi_qos.reliability = static_cast<nros_cpp_qos_reliability_t>(qos.reliability_raw());
668 ffi_qos.durability = static_cast<nros_cpp_qos_durability_t>(qos.durability_raw());
669 ffi_qos.history = static_cast<nros_cpp_qos_history_t>(qos.history_raw());
670 ffi_qos.liveliness_kind = static_cast<nros_cpp_qos_liveliness_t>(qos.liveliness_raw());
671 ffi_qos.depth = qos.depth();
672 ffi_qos.deadline_ms = qos.deadline_ms();
673 ffi_qos.lifespan_ms = qos.lifespan_ms();
674 ffi_qos.liveliness_lease_ms = qos.liveliness_lease_ms();
675 ffi_qos.avoid_ros_namespace_conventions = qos.avoid_ros_namespace_conventions() ? 1 : 0;
676 ffi_qos.tx_express = qos.tx_express() ? 1 : 0;
677
678 out.user_fn_info_ = typename Subscription<M>::TypedSubscriptionInfoFn(callback);
679 out.user_fn_ = nullptr;
680 out.user_fn_ctx_ = nullptr;
681 out.user_ctx_ = nullptr;
682
683 uint8_t sched = (options.sched_context == SCHED_CONTEXT_UNSET)
684 ? 0u
685 : static_cast<uint8_t>(options.sched_context);
686 size_t handle = static_cast<size_t>(-1);
688 &handle_, topic, M::TYPE_NAME, M::TYPE_HASH, ffi_qos,
690 if (ret == 0) {
691 out.sched_handle_id_ = handle;
692 out.callback_mode_ = true;
693 out.initialized_ = true;
694 }
695 return Result(ret);
696}
697
701template <typename M>
702inline Expected<Subscription<M>> make_subscription(Node& node, const char* topic,
703 const QoS& qos = QoS::default_profile()) {
705 Result r = node.create_subscription<M>(s, topic, qos);
706 if (!r.ok()) return Expected<Subscription<M>>::error(r);
707 return Expected<Subscription<M>>::ok(std::move(s));
708}
709
710#if defined(NANO_ROS_SAFETY_E2E)
711// Phase 269 W3 — out-of-line definition of Node::create_subscription_with_safety.
712// Mirrors `create_subscription_with_info` one overload over, but routes through
713// `nros_cpp_subscription_register_validated` so the arena dispatches the
714// `message_safety_trampoline` on each new sample.
715template <typename M, typename F, typename>
716Result Node::create_subscription_with_safety(Subscription<M>& out, const char* topic, F callback,
717 const QoS& qos, const SubscriptionOptions& options) {
718 if (!initialized_) return Result(ErrorCode::NotInitialized);
719 nros_cpp_qos_t ffi_qos;
720 ffi_qos.reliability = static_cast<nros_cpp_qos_reliability_t>(qos.reliability_raw());
721 ffi_qos.durability = static_cast<nros_cpp_qos_durability_t>(qos.durability_raw());
722 ffi_qos.history = static_cast<nros_cpp_qos_history_t>(qos.history_raw());
723 ffi_qos.liveliness_kind = static_cast<nros_cpp_qos_liveliness_t>(qos.liveliness_raw());
724 ffi_qos.depth = qos.depth();
725 ffi_qos.deadline_ms = qos.deadline_ms();
726 ffi_qos.lifespan_ms = qos.lifespan_ms();
727 ffi_qos.liveliness_lease_ms = qos.liveliness_lease_ms();
728 ffi_qos.avoid_ros_namespace_conventions = qos.avoid_ros_namespace_conventions() ? 1 : 0;
729 ffi_qos.tx_express = qos.tx_express() ? 1 : 0;
730
731 out.user_fn_safety_ = typename Subscription<M>::TypedSubscriptionSafetyFn(callback);
732 out.user_fn_ = nullptr;
733 out.user_fn_ctx_ = nullptr;
734 out.user_fn_info_ = nullptr;
735 out.user_ctx_ = nullptr;
736
737 uint8_t sched = (options.sched_context == SCHED_CONTEXT_UNSET)
738 ? 0u
739 : static_cast<uint8_t>(options.sched_context);
740 size_t handle = static_cast<size_t>(-1);
741 nros_cpp_ret_t ret = nros_cpp_subscription_register_validated(
742 &handle_, topic, M::TYPE_NAME, M::TYPE_HASH, ffi_qos,
743 &Subscription<M>::message_safety_trampoline, &out, sched, &handle);
744 if (ret == 0) {
745 out.sched_handle_id_ = handle;
746 out.callback_mode_ = true;
747 out.initialized_ = true;
748 }
749 return Result(ret);
750}
751#endif // NANO_ROS_SAFETY_E2E
752
753} // namespace nros
754
755#endif // NROS_CPP_SUBSCRIPTION_HPP
Definition result.hpp:195
ErrorCode error() const
Definition result.hpp:218
bool ok() const
Definition result.hpp:211
Definition future.hpp:40
Definition node.hpp:181
Result create_subscription_in(const CallbackGroup &group, Subscription< M > &out, const char *topic, F callback, const QoS &qos=QoS::default_profile(), const SubscriptionOptions &options={})
Definition subscription.hpp:622
Result create_subscription(Subscription< M > &out, const char *topic, const QoS &qos=QoS::default_profile())
Definition subscription.hpp:513
Result create_subscription_with_info(Subscription< M > &out, const char *topic, F callback, const QoS &qos=QoS::default_profile(), const SubscriptionOptions &options={})
Definition subscription.hpp:663
Definition qos.hpp:68
static constexpr QoS default_profile()
Default profile: RELIABLE + VOLATILE + KEEP_LAST(10).
Definition qos.hpp:170
Definition result.hpp:87
static constexpr Result success()
Named constructors.
Definition result.hpp:109
bool is_valid() const
Check if the stream is connected to a valid source.
Definition stream.hpp:105
Definition subscription.hpp:217
View & operator=(View &&o)
Definition subscription.hpp:224
View(const View &)=delete
View(View &&o)
Definition subscription.hpp:220
View()
Definition subscription.hpp:219
const uint8_t * data() const
Definition subscription.hpp:240
bool empty() const
Definition subscription.hpp:242
size_t size() const
Definition subscription.hpp:241
View(void *sub, const uint8_t *buf, size_t len, void *token)
Internal constructor — callers use Subscription::try_borrow().
Definition subscription.hpp:245
View & operator=(const View &)=delete
~View()
Definition subscription.hpp:238
Definition subscription.hpp:95
void(*)(const M &msg, const uint8_t *attachment, size_t attachment_len) TypedSubscriptionInfoFn
Definition subscription.hpp:104
Result on_liveliness_changed(nros_cpp_liveliness_changed_cb_t cb, void *user_context=nullptr)
Definition subscription.hpp:405
const char * get_topic_name() const
Get the topic name.
Definition subscription.hpp:301
Result on_requested_deadline_missed(uint32_t deadline_ms, nros_cpp_subscriber_count_cb_t cb, void *user_context=nullptr)
Register a callback for requested-deadline-missed events.
Definition subscription.hpp:412
Result try_recv_sequence(uint8_t *buf, size_t per_msg_cap, size_t max_msgs, size_t *out_lens, size_t &out_count)
Definition subscription.hpp:288
Subscription(Subscription &&other)
Definition subscription.hpp:342
Subscription()
Definition subscription.hpp:382
~Subscription()
Definition subscription.hpp:328
Expected< View > try_borrow()
Definition subscription.hpp:265
const Stream< M > & stream() const
Definition subscription.hpp:317
Stream< M > & stream()
Definition subscription.hpp:310
size_t sched_handle_id() const
Definition subscription.hpp:395
bool has_sched_handle() const
Definition subscription.hpp:394
Result try_recv(M &msg)
Definition subscription.hpp:123
Result on_message_lost(nros_cpp_subscriber_count_cb_t cb, void *user_context=nullptr)
Register a callback for message-lost events.
Definition subscription.hpp:420
Result try_recv_raw(uint8_t *buf, size_t capacity, size_t &out_len)
Definition subscription.hpp:171
void(*)(const M &msg, void *ctx) TypedSubscriptionFnWithCtx
Definition subscription.hpp:101
bool is_valid() const
Check if the subscription is initialized and valid.
Definition subscription.hpp:320
Subscription & operator=(Subscription &&other)
Definition subscription.hpp:357
Result try_recv_raw_with_attachment(uint8_t *buf, size_t capacity, size_t &out_len, uint8_t *att, size_t att_capacity, size_t &out_att_len)
Definition subscription.hpp:197
void(*)(const M &msg) TypedSubscriptionFn
Definition subscription.hpp:100
Result try_recv_validated(M &msg, nros_cpp_integrity_status_t &status)
Definition subscription.hpp:149
Inline storage-size macros for opaque entity buffers.
int nros_cpp_ret_t
Definition future.hpp:20
Definition nros.hpp:43
bool ok()
Check if the nros session is initialized.
Definition node.hpp:842
static constexpr size_t SUBSCRIPTION_TOPIC_NAME_MAX
Definition subscription.hpp:75
@ Error
Generic failure not covered by a more specific code.
@ TryAgain
Transient — no data ready yet (non-blocking take). Retry later.
nros::Node and global session helpers.
nros_cpp_qos_history_t
Definition qos.hpp:31
nros_cpp_qos_liveliness_t
Definition qos.hpp:35
nros_cpp_qos_durability_t
Definition qos.hpp:27
nros_cpp_qos_reliability_t
Definition qos.hpp:23
nros::Result, nros::ErrorCode, and the NROS_TRY macro.
nros::Stream<T> — multi-shot message receiver.
Definition qos.hpp:41
enum nros_cpp_qos_durability_t durability
Definition qos.hpp:43
uint32_t liveliness_lease_ms
Definition qos.hpp:49
uint32_t lifespan_ms
Definition qos.hpp:48
enum nros_cpp_qos_history_t history
Definition qos.hpp:44
uint8_t tx_express
Definition qos.hpp:51
enum nros_cpp_qos_liveliness_t liveliness_kind
Definition qos.hpp:45
enum nros_cpp_qos_reliability_t reliability
Definition qos.hpp:42
uint8_t avoid_ros_namespace_conventions
Definition qos.hpp:50
uint32_t deadline_ms
Definition qos.hpp:47
int depth
Definition qos.hpp:46
void(* nros_cpp_subscription_message_info_callback_t)(const uint8_t *data, size_t len, const uint8_t *attachment, size_t attachment_len, void *ctx)
Definition subscription.hpp:42
nros_cpp_ret_t nros_cpp_subscription_register(const nros_cpp_node_t *node, const char *topic, const char *type_name, const char *type_hash, nros_cpp_qos_t qos, nros_cpp_subscription_message_callback_t callback, void *context, uint8_t sched_context, size_t *out_handle_id, const char *callback_group)
void(* nros_cpp_subscription_message_callback_t)(const uint8_t *data, size_t len, void *ctx)
Definition subscription.hpp:29
nros_cpp_ret_t nros_cpp_subscription_register_with_info(const nros_cpp_node_t *node, const char *topic, const char *type_name, const char *type_hash, nros_cpp_qos_t qos, nros_cpp_subscription_message_info_callback_t callback, void *context, uint8_t sched_context, size_t *out_handle_id)