10#ifndef NROS_CPP_SUBSCRIPTION_HPP
11#define NROS_CPP_SUBSCRIPTION_HPP
21#include "nros_cpp_ffi.h"
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);
43 const uint8_t* attachment,
44 size_t attachment_len,
void* ctx);
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);
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);
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);
108#if defined(NANO_ROS_SAFETY_E2E)
220 View(
View&&
o) : sub_(
o.sub_), buf_(
o.buf_), len_(
o.len_), token_(
o.token_) {
241 size_t size()
const {
return len_; }
242 bool empty()
const {
return token_ ==
nullptr; }
246 : sub_(
sub), buf_(
buf), len_(len), token_(
token) {}
250 if (token_ && sub_) {
269 void*
token =
nullptr;
311 if (initialized_ && !stream_.
is_valid()) {
329 if (initialized_ && !callback_mode_) {
332 initialized_ =
false;
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_) {
350 ::memcpy(topic_name_,
other.topic_name_,
sizeof(topic_name_));
353 other.initialized_ =
false;
358 if (
this != &
other) {
359 if (initialized_ && !callback_mode_) {
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_) {
371 ::memcpy(topic_name_,
other.topic_name_,
sizeof(topic_name_));
374 other.initialized_ =
false;
434 static void message_trampoline(
const uint8_t* data,
size_t len,
void*
ctx) {
436 if (
self ==
nullptr)
return;
438 if (M::ffi_deserialize(data, len, &
msg) != 0)
return;
439 if (
self->user_fn_ !=
nullptr) {
441 }
else if (
self->user_fn_ctx_ !=
nullptr) {
449 static void message_info_trampoline(
const uint8_t* data,
size_t len,
const uint8_t* attachment,
450 size_t attachment_len,
void* ctx) {
452 if (self ==
nullptr)
return;
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);
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) {
469 if (self ==
nullptr)
return;
471 if (M::ffi_deserialize(data, len, &msg) != 0)
return;
472 if (self->user_fn_safety_ !=
nullptr) {
473 nros_cpp_integrity_status_t status;
475 status.duplicate = duplicate;
476 status.crc_valid = crc_valid;
477 self->user_fn_safety_(msg, status);
482 alignas(8) uint8_t storage_[NROS_SUBSCRIBER_SIZE];
490 size_t sched_handle_id_ =
static_cast<size_t>(-1);
496 void* user_ctx_ =
nullptr;
497 bool callback_mode_ =
false;
498#if defined(NANO_ROS_SAFETY_E2E)
501 TypedSubscriptionSafetyFn user_fn_safety_ =
nullptr;
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;
536 out.initialized_ =
true;
557 if (!
r.ok())
return r;
563 executor_handle_,
out.sched_handle_id(),
static_cast<uint8_t>(
options.sched_context));
569 out.initialized_ =
false;
580template <
typename M,
typename F,
typename>
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;
599 out.user_fn_ctx_ =
nullptr;
600 out.user_ctx_ =
nullptr;
605 size_t handle =
static_cast<size_t>(-1);
611 out.sched_handle_id_ = handle;
612 out.callback_mode_ =
true;
613 out.initialized_ =
true;
621template <
typename M,
typename F,
typename>
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;
639 out.user_fn_ctx_ =
nullptr;
640 out.user_ctx_ =
nullptr;
645 size_t handle =
static_cast<size_t>(-1);
651 out.sched_handle_id_ = handle;
652 out.callback_mode_ =
true;
653 out.initialized_ =
true;
662template <
typename M,
typename F,
typename>
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;
679 out.user_fn_ =
nullptr;
680 out.user_fn_ctx_ =
nullptr;
681 out.user_ctx_ =
nullptr;
686 size_t handle =
static_cast<size_t>(-1);
691 out.sched_handle_id_ = handle;
692 out.callback_mode_ =
true;
693 out.initialized_ =
true;
710#if defined(NANO_ROS_SAFETY_E2E)
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) {
724 ffi_qos.
depth = qos.depth();
729 ffi_qos.
tx_express = qos.tx_express() ? 1 : 0;
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;
737 uint8_t sched = (options.sched_context == SCHED_CONTEXT_UNSET)
739 : static_cast<uint8_t>(options.sched_context);
740 size_t handle =
static_cast<size_t>(-1);
742 &handle_, topic, M::TYPE_NAME, M::TYPE_HASH, ffi_qos,
743 &Subscription<M>::message_safety_trampoline, &out, sched, &handle);
745 out.sched_handle_id_ = handle;
746 out.callback_mode_ =
true;
747 out.initialized_ =
true;
Definition result.hpp:195
ErrorCode error() const
Definition result.hpp:218
bool ok() const
Definition result.hpp:211
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
static constexpr QoS default_profile()
Default profile: RELIABLE + VOLATILE + KEEP_LAST(10).
Definition qos.hpp:170
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
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.
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)