10#ifndef NROS_CPP_SUBSCRIPTION_HPP
11#define NROS_CPP_SUBSCRIPTION_HPP
19#include "nros/size_bound.hpp"
24#include "nros_cpp_ffi.h"
39 const char* type_name,
const char* type_hash,
42 void* context,
size_t* out_handle_id,
43 const nros_cpp_subscription_options_t* options);
48 const uint8_t* attachment,
49 size_t attachment_len,
void* ctx);
52 const nros_cpp_node_t* node,
const char* topic,
const char* type_name,
const char* type_hash,
54 size_t* out_handle_id,
const nros_cpp_subscription_options_t* options);
63#if defined(NANO_ROS_SAFETY_E2E)
64typedef void (*nros_cpp_subscription_validated_callback_t)(
const uint8_t* data,
size_t len,
65 int64_t gap,
bool duplicate,
66 int8_t crc_valid,
void* ctx);
69 const nros_cpp_node_t* node,
const char* topic,
const char* type_name,
const char* type_hash,
70 nros_cpp_qos_t qos, nros_cpp_subscription_validated_callback_t callback,
void* context,
71 size_t* out_handle_id,
const nros_cpp_subscription_options_t* options);
109 size_t attachment_len);
113#if defined(NANO_ROS_SAFETY_E2E)
114 using TypedSubscriptionSafetyFn = void (*)(
const M& msg,
115 const nros_cpp_integrity_status_t& integrity);
128 Result take(M& msg) {
return take_sized<::nros::rx_buffer_capacity<M>::value>(msg); }
143 nros_cpp_subscription_take_serialized(storage_, buf,
sizeof(buf), &len);
144 if (ret != 0)
return Result(ret);
166 return take_validated_sized<::nros::rx_buffer_capacity<M>::value>(msg, status);
176 nros_cpp_subscription_take_validated(storage_, buf,
sizeof(buf), &len, &status);
177 if (ret != 0)
return Result(ret);
199 nros_cpp_subscription_take_serialized(storage_, buf, capacity, &out_len);
200 if (ret != 0)
return Result(ret);
221 uint8_t* att,
size_t att_capacity,
size_t& out_att_len) {
227 nros_cpp_ret_t ret = nros_cpp_subscription_take_serialized_with_attachment(
228 storage_, buf, capacity, &out_len, att, att_capacity, &out_att_len);
229 if (ret != 0)
return Result(ret);
242 View() : sub_(nullptr), buf_(nullptr), len_(0), token_(nullptr) {}
243 View(
View&& o) : sub_(o.sub_), buf_(o.buf_), len_(o.len_), token_(o.token_) {
263 const uint8_t*
data()
const {
return buf_; }
264 size_t size()
const {
return len_; }
265 bool empty()
const {
return token_ ==
nullptr; }
268 View(
void* sub,
const uint8_t* buf,
size_t len,
void* token)
269 : sub_(sub), buf_(buf), len_(len), token_(token) {}
273 if (token_ && sub_) {
274 nros_cpp_subscription_release(sub_, token_);
290 const uint8_t* buf =
nullptr;
292 void* token =
nullptr;
293 int32_t rc = nros_cpp_subscription_borrow(storage_, &buf, &len, &token);
317 nros_cpp_ret_t ret = nros_cpp_subscription_take_sequence(storage_, buf, per_msg_cap,
318 max_msgs, out_lens, &out_count);
319 if (ret != 0)
return Result(ret);
339 [[deprecated(
"Subscription::try_recv is deprecated; use Subscription::take")]]
Result
349 template <
size_t Cap>
350 [[deprecated(
"Subscription::try_recv_sized is deprecated; use "
351 "Subscription::take_sized")]]
Result
353 return take_sized<Cap>(msg);
357 template <
size_t Cap>
358 [[deprecated(
"Subscription::try_recv_validated_sized is deprecated; use "
359 "Subscription::take_validated_sized")]]
Result
361 return take_validated_sized<Cap>(msg, status);
365 [[deprecated(
"Subscription::try_recv_validated is deprecated; use "
366 "Subscription::take_validated")]]
Result
372 [[deprecated(
"Subscription::try_recv_raw is deprecated; use "
373 "Subscription::take_serialized")]]
Result
379 [[deprecated(
"Subscription::try_recv_raw_with_attachment is deprecated; use "
380 "Subscription::take_serialized_with_attachment")]]
Result
382 size_t att_capacity,
size_t& out_att_len) {
388 [[deprecated(
"Subscription::try_recv_sequence is deprecated; use "
389 "Subscription::take_sequence")]]
Result
392 return take_sequence(buf, per_msg_cap, max_msgs, out_lens, out_count);
406 if (initialized_ && !stream_.
is_valid()) {
407 stream_.bind(storage_, &nros_cpp_subscription_take_serialized);
424 if (initialized_ && !callback_mode_) {
425 nros_cpp_subscription_destroy(storage_);
427 initialized_ =
false;
438 user_fn_ = other.user_fn_;
439 user_fn_ctx_ = other.user_fn_ctx_;
440 user_ctx_ = other.user_ctx_;
441 callback_mode_ = other.callback_mode_;
442 sched_handle_id_ = other.sched_handle_id_;
443 if (other.initialized_ && !other.callback_mode_) {
444 nros_cpp_subscription_relocate(other.storage_, storage_);
445 ::memcpy(topic_name_, other.topic_name_,
sizeof(topic_name_));
446 stream_.bind(storage_, &nros_cpp_subscription_take_serialized);
448 other.initialized_ =
false;
453 if (
this != &other) {
454 if (initialized_ && !callback_mode_) {
455 nros_cpp_subscription_destroy(storage_);
458 initialized_ = other.initialized_;
459 user_fn_ = other.user_fn_;
460 user_fn_ctx_ = other.user_fn_ctx_;
461 user_ctx_ = other.user_ctx_;
462 callback_mode_ = other.callback_mode_;
463 sched_handle_id_ = other.sched_handle_id_;
464 if (other.initialized_ && !other.callback_mode_) {
465 nros_cpp_subscription_relocate(other.storage_, storage_);
466 ::memcpy(topic_name_, other.topic_name_,
sizeof(topic_name_));
467 stream_.bind(storage_, &nros_cpp_subscription_take_serialized);
469 other.initialized_ =
false;
477 Subscription() : storage_(), topic_name_{}, initialized_(false), stream_() {}
501 void* user_context =
nullptr) {
503 return Result(nros_cpp_subscription_set_liveliness_changed(storage_, cb, user_context));
508 void* user_context =
nullptr) {
510 return Result(nros_cpp_subscription_set_requested_deadline_missed(storage_, deadline_ms, cb,
517 return Result(nros_cpp_subscription_set_message_lost(storage_, cb, user_context));
529 static void message_trampoline(
const uint8_t* data,
size_t len,
void* ctx) {
531 if (self ==
nullptr)
return;
533 if (M::ffi_deserialize(data, len, &msg) != 0)
return;
534 if (self->user_fn_ !=
nullptr) {
536 }
else if (self->user_fn_ctx_ !=
nullptr) {
537 self->user_fn_ctx_(msg, self->user_ctx_);
544 static void message_info_trampoline(
const uint8_t* data,
size_t len,
const uint8_t* attachment,
545 size_t attachment_len,
void* ctx) {
547 if (self ==
nullptr)
return;
549 if (M::ffi_deserialize(data, len, &msg) != 0)
return;
550 if (self->user_fn_info_ !=
nullptr) {
551 self->user_fn_info_(msg, attachment, attachment_len);
555#if defined(NANO_ROS_SAFETY_E2E)
561 static void message_safety_trampoline(
const uint8_t* data,
size_t len, int64_t gap,
562 bool duplicate, int8_t crc_valid,
void* ctx) {
564 if (self ==
nullptr)
return;
566 if (M::ffi_deserialize(data, len, &msg) != 0)
return;
567 if (self->user_fn_safety_ !=
nullptr) {
568 nros_cpp_integrity_status_t status;
570 status.duplicate = duplicate;
571 status.crc_valid = crc_valid;
572 self->user_fn_safety_(msg, status);
577 alignas(8) uint8_t storage_[NROS_SUBSCRIBER_SIZE];
585 size_t sched_handle_id_ =
static_cast<size_t>(-1);
591 void* user_ctx_ =
nullptr;
592 bool callback_mode_ =
false;
593#if defined(NANO_ROS_SAFETY_E2E)
596 TypedSubscriptionSafetyFn user_fn_safety_ =
nullptr;
614 nros_cpp_ret_t ret = nros_cpp_subscription_create(&handle_, topic, M::TYPE_NAME, M::TYPE_HASH,
615 ffi_qos, out.storage_);
618 size_t topic_len = 0;
619 while (topic[topic_len] !=
'\0' && topic_len + 1 <
sizeof(out.topic_name_)) {
620 out.topic_name_[topic_len] = topic[topic_len];
623 out.topic_name_[topic_len] =
'\0';
624 out.initialized_ =
true;
643 const SubscriptionOptions& options) {
644 Result r = create_subscription<M>(out, topic, qos);
645 if (!r.
ok())
return r;
649 if (options.sched_context != SCHED_CONTEXT_UNSET && out.
has_sched_handle()) {
651 executor_handle_, out.
sched_handle_id(),
static_cast<uint8_t
>(options.sched_context));
656 nros_cpp_subscription_destroy(out.storage_);
657 out.initialized_ =
false;
668template <
typename M,
typename F,
typename>
670 const QoS& qos,
const SubscriptionOptions& options) {
680 out.user_fn_ctx_ =
nullptr;
681 out.user_ctx_ =
nullptr;
683 uint8_t sched = (options.sched_context == SCHED_CONTEXT_UNSET)
685 :
static_cast<uint8_t
>(options.sched_context);
686 size_t handle =
static_cast<size_t>(-1);
689 nros_cpp_subscription_options_t ffi_options = nros_cpp_subscription_default_options();
690 ffi_options.sched_context = sched;
693 &out, &handle, &ffi_options);
695 out.sched_handle_id_ = handle;
696 out.callback_mode_ =
true;
697 out.initialized_ =
true;
705template <
typename M,
typename F,
typename>
707 const char* topic, F callback,
const QoS& qos,
708 const SubscriptionOptions& options) {
716 out.user_fn_ctx_ =
nullptr;
717 out.user_ctx_ =
nullptr;
719 uint8_t sched = (options.sched_context == SCHED_CONTEXT_UNSET)
721 :
static_cast<uint8_t
>(options.sched_context);
722 size_t handle =
static_cast<size_t>(-1);
724 nros_cpp_subscription_options_t ffi_options = nros_cpp_subscription_default_options();
725 ffi_options.sched_context = sched;
726 ffi_options.callback_group = group.get_name();
729 &out, &handle, &ffi_options);
731 out.sched_handle_id_ = handle;
732 out.callback_mode_ =
true;
733 out.initialized_ =
true;
742template <
typename M,
typename F,
typename>
744 const QoS& qos,
const SubscriptionOptions& options) {
752 out.user_fn_ =
nullptr;
753 out.user_fn_ctx_ =
nullptr;
754 out.user_ctx_ =
nullptr;
756 uint8_t sched = (options.sched_context == SCHED_CONTEXT_UNSET)
758 :
static_cast<uint8_t
>(options.sched_context);
759 size_t handle =
static_cast<size_t>(-1);
761 nros_cpp_subscription_options_t ffi_options = nros_cpp_subscription_default_options();
762 ffi_options.sched_context = sched;
764 &handle_, topic, M::TYPE_NAME, M::TYPE_HASH, ffi_qos,
767 out.sched_handle_id_ = handle;
768 out.callback_mode_ =
true;
769 out.initialized_ =
true;
786#if defined(NANO_ROS_SAFETY_E2E)
791template <
typename M,
typename F,
typename>
792Result Node::create_subscription_with_safety(Subscription<M>& out,
const char* topic, F callback,
793 const QoS& qos,
const SubscriptionOptions& options) {
800 out.user_fn_safety_ =
typename Subscription<M>::TypedSubscriptionSafetyFn(callback);
801 out.user_fn_ =
nullptr;
802 out.user_fn_ctx_ =
nullptr;
803 out.user_fn_info_ =
nullptr;
804 out.user_ctx_ =
nullptr;
806 uint8_t sched = (options.sched_context == SCHED_CONTEXT_UNSET)
808 : static_cast<uint8_t>(options.sched_context);
809 size_t handle =
static_cast<size_t>(-1);
811 nros_cpp_subscription_options_t ffi_options = nros_cpp_subscription_default_options();
812 ffi_options.sched_context = sched;
814 &handle_, topic, M::TYPE_NAME, M::TYPE_HASH, ffi_qos,
815 &Subscription<M>::message_safety_trampoline, &out, &handle, &ffi_options);
817 out.sched_handle_id_ = handle;
818 out.callback_mode_ =
true;
819 out.initialized_ =
true;
Definition result.hpp:198
ErrorCode error() const
Definition result.hpp:221
bool ok() const
Definition result.hpp:214
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:706
Result create_subscription(Subscription< M > &out, const char *topic, const QoS &qos=QoS::default_profile())
Definition subscription.hpp:608
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:743
static constexpr QoS default_profile()
Default profile: RELIABLE + VOLATILE + KEEP_LAST(10).
Definition qos.hpp:328
static constexpr Result success()
Named constructors.
Definition result.hpp:112
bool ok() const
Returns true if the operation succeeded.
Definition result.hpp:100
bool is_valid() const
Check if the stream is connected to a valid source.
Definition stream.hpp:119
Definition subscription.hpp:240
View & operator=(View &&o)
Definition subscription.hpp:247
View(const View &)=delete
View(View &&o)
Definition subscription.hpp:243
View()
Definition subscription.hpp:242
const uint8_t * data() const
Definition subscription.hpp:263
bool empty() const
Definition subscription.hpp:265
size_t size() const
Definition subscription.hpp:264
View(void *sub, const uint8_t *buf, size_t len, void *token)
Internal constructor — callers use Subscription::try_borrow().
Definition subscription.hpp:268
View & operator=(const View &)=delete
~View()
Definition subscription.hpp:261
Definition subscription.hpp:100
void(*)(const M &msg, const uint8_t *attachment, size_t attachment_len) TypedSubscriptionInfoFn
Definition subscription.hpp:109
Result on_liveliness_changed(nros_cpp_liveliness_changed_cb_t cb, void *user_context=nullptr)
Definition subscription.hpp:500
Result take_serialized(uint8_t *buf, size_t capacity, size_t &out_len)
Definition subscription.hpp:193
const char * get_topic_name() const
Get the topic name.
Definition subscription.hpp:396
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:507
Result try_recv_validated_sized(M &msg, nros_cpp_integrity_status_t &status)
Definition subscription.hpp:360
Result take_sequence(uint8_t *buf, size_t per_msg_cap, size_t max_msgs, size_t *out_lens, size_t &out_count)
Definition subscription.hpp:311
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:390
Subscription(Subscription &&other)
Definition subscription.hpp:437
Subscription()
Definition subscription.hpp:477
Result try_recv_sized(M &msg)
Definition subscription.hpp:352
~Subscription()
Definition subscription.hpp:423
Expected< View > try_borrow()
Definition subscription.hpp:288
const Stream< M > & stream() const
Definition subscription.hpp:412
Stream< M > & stream()
Definition subscription.hpp:405
size_t sched_handle_id() const
Definition subscription.hpp:490
bool has_sched_handle() const
Definition subscription.hpp:489
Result take_serialized_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:220
Result try_recv(M &msg)
Definition subscription.hpp:340
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:515
Result try_recv_raw(uint8_t *buf, size_t capacity, size_t &out_len)
Definition subscription.hpp:374
void(*)(const M &msg, void *ctx) TypedSubscriptionFnWithCtx
Definition subscription.hpp:106
bool is_valid() const
Check if the subscription is initialized and valid.
Definition subscription.hpp:415
Result take(M &msg)
Definition subscription.hpp:128
Subscription & operator=(Subscription &&other)
Definition subscription.hpp:452
Result take_sized(M &msg)
Definition subscription.hpp:138
Result take_validated_sized(M &msg, nros_cpp_integrity_status_t &status)
Definition subscription.hpp:171
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:381
Result take_validated(M &msg, nros_cpp_integrity_status_t &status)
Definition subscription.hpp:165
void(*)(const M &msg) TypedSubscriptionFn
Definition subscription.hpp:105
Result try_recv_validated(M &msg, nros_cpp_integrity_status_t &status)
Definition subscription.hpp:367
Inline storage-size macros for opaque entity buffers.
int nros_cpp_ret_t
Definition future.hpp:21
bool ok()
Check if the nros session is initialized.
Definition node.hpp:997
static constexpr size_t SUBSCRIPTION_TOPIC_NAME_MAX
Definition subscription.hpp:80
@ 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::Result, nros::ErrorCode, and the NROS_TRY macro.
nros::Stream<T> — multi-shot message receiver.
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, size_t *out_handle_id, const nros_cpp_subscription_options_t *options)
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:47
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, size_t *out_handle_id, const nros_cpp_subscription_options_t *options)
void(* nros_cpp_subscription_message_callback_t)(const uint8_t *data, size_t len, void *ctx)
Definition subscription.hpp:32