1use core::marker::PhantomData;
4
5use nros_node::rmw_type_registry::MessageForRmw;
10
11use crate::{
12 ActionTag, CallbackId, CancelResponse, EntityId, GoalId, GoalResponse, GoalStatus,
13 ParameterType, QoSProfile, RosAction, RosMessage, RosService, ServiceTag, SubscriptionTag,
14 TimerDuration,
15 heapless::Vec,
16 node_metadata::{
17 CallbackEffectKind, CallbackEffectMetadata, CallbackSlot, EntityKind, EntityMetadata,
18 EntityMetadataSpec, EntitySlot, MetadataRecorder, MetadataString, NodeId,
19 NodeMetadataError, NodeSlot, ParameterDefault, SourceLocationMetadata, copy_str,
20 entity_callback_ids, entity_metadata,
21 },
22};
23
24pub const MISSING_NODE_EXPORT_ERROR: &str = "package has no exported nros component";
34
35pub type NodeResult<T = ()> = Result<T, NodeDeclError>;
37
38#[inline]
66fn register_declared_type<M: nros_node::rmw_type_registry::MessageForRmw>() -> NodeResult<()> {
67 nros_node::rmw_type_registry::register_type::<M>().map_err(|_| NodeDeclError::Runtime)
68}
69
70#[inline]
81fn register_declared_service<S: RosService>() -> NodeResult<()>
82where
83 S::Request: nros_node::rmw_type_registry::MessageForRmw,
84 S::Reply: nros_node::rmw_type_registry::MessageForRmw,
85{
86 register_declared_type::<S::Request>()?;
87 register_declared_type::<S::Reply>()
88}
89
90#[derive(Debug, Clone, Copy, PartialEq, Eq)]
92pub enum NodeDeclError {
93 Metadata(NodeMetadataError),
95 MissingExport,
97 Runtime,
99 ExecutorFull,
104 UnknownPublisher,
115}
116
117impl NodeDeclError {
118 pub const fn message(self) -> &'static str {
120 match self {
121 Self::Metadata(NodeMetadataError::Capacity) => "component metadata capacity exceeded",
122 Self::Metadata(NodeMetadataError::NameTooLong) => "component metadata name too long",
123 Self::Metadata(NodeMetadataError::UnknownNode) => {
124 "component entity references an unknown node"
125 }
126 Self::Metadata(NodeMetadataError::UnknownEntity) => {
127 "component callback effect references an unknown entity"
128 }
129 Self::Metadata(NodeMetadataError::DuplicateId) => {
130 "component metadata contains a duplicate stable ID"
131 }
132 Self::MissingExport => MISSING_NODE_EXPORT_ERROR,
133 Self::Runtime => "component runtime rejected declaration",
134 Self::UnknownPublisher => "no publisher declared for that entity",
135 Self::ExecutorFull => {
136 "executor callback table full — raise NROS_EXECUTOR_MAX_CBS \
137 (build-time, default 4)"
138 }
139 }
140 }
141}
142
143impl From<NodeMetadataError> for NodeDeclError {
144 fn from(value: NodeMetadataError) -> Self {
145 Self::Metadata(value)
146 }
147}
148
149#[inline]
157fn register_declared_action<A: RosAction>() -> NodeResult<()>
158where
159 A::Goal: nros_node::rmw_type_registry::MessageForRmw,
160 A::Result: nros_node::rmw_type_registry::MessageForRmw,
161 A::Feedback: nros_node::rmw_type_registry::MessageForRmw,
162 A::SendGoalRequest: nros_node::rmw_type_registry::MessageForRmw,
163 A::SendGoalResponse: nros_node::rmw_type_registry::MessageForRmw,
164 A::GetResultRequest: nros_node::rmw_type_registry::MessageForRmw,
165 A::GetResultResponse: nros_node::rmw_type_registry::MessageForRmw,
166 A::FeedbackMessage: nros_node::rmw_type_registry::MessageForRmw,
167{
168 register_declared_type::<A::Goal>()?;
169 register_declared_type::<A::Result>()?;
170 register_declared_type::<A::Feedback>()?;
171 register_declared_type::<A::SendGoalRequest>()?;
172 register_declared_type::<A::SendGoalResponse>()?;
173 register_declared_type::<A::GetResultRequest>()?;
174 register_declared_type::<A::GetResultResponse>()?;
175 register_declared_type::<A::FeedbackMessage>()?;
176 A::register_protocol_types().map_err(|()| NodeDeclError::Runtime)
177}
178
179#[derive(Debug, Clone, Copy, PartialEq, Eq)]
195pub struct EntityBounds {
196 pub publishers: usize,
198 pub service_servers: usize,
200 pub service_clients: usize,
202 pub action_clients: usize,
204 pub action_servers: usize,
206}
207
208impl EntityBounds {
209 pub const fn knob_caps() -> Self {
211 Self {
212 publishers: crate::config::MAX_CELL_ENTITIES,
213 service_servers: crate::config::MAX_CELL_ENTITIES,
214 service_clients: crate::config::MAX_CELL_ENTITIES,
215 action_clients: crate::config::MAX_CELL_ENTITIES,
216 action_servers: crate::config::MAX_CELL_ENTITIES,
217 }
218 }
219
220 pub const fn exact(
223 publishers: usize,
224 service_servers: usize,
225 service_clients: usize,
226 action_clients: usize,
227 action_servers: usize,
228 ) -> Self {
229 Self {
230 publishers,
231 service_servers,
232 service_clients,
233 action_clients,
234 action_servers,
235 }
236 }
237}
238
239pub trait Node {
241 const NAME: &'static str;
243
244 const DISPATCH: crate::DispatchStrategy = crate::DispatchStrategy::Inline;
251
252 const ENTITY_BOUNDS: EntityBounds = EntityBounds::knob_caps();
257
258 fn register(context: &mut NodeContext<'_>) -> NodeResult<()>;
260}
261
262#[derive(Debug, Clone, Copy, PartialEq, Eq)]
264pub struct NodeOptions<'a> {
265 pub name: &'a str,
267 pub namespace: &'a str,
269 pub domain_id: u32,
271}
272
273#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
278pub struct Callback<'a> {
279 id: CallbackId<'a>,
280}
281
282impl<'a> Callback<'a> {
283 pub const fn as_str(self) -> &'a str {
285 self.id.as_str()
286 }
287
288 pub fn is_named(self, name: &str) -> bool {
290 self.as_str() == name
291 }
292
293 #[doc(hidden)]
295 pub const fn __from_id(id: CallbackId<'a>) -> Self {
296 Self { id }
297 }
298}
299
300impl<'a> NodeOptions<'a> {
301 pub const fn new(name: &'a str) -> Self {
303 Self {
304 name,
305 namespace: "/",
306 domain_id: 0,
307 }
308 }
309
310 pub const fn namespace(mut self, namespace: &'a str) -> Self {
312 self.namespace = namespace;
313 self
314 }
315
316 pub const fn domain_id(mut self, domain_id: u32) -> Self {
318 self.domain_id = domain_id;
319 self
320 }
321}
322
323pub trait NodeRuntime {
325 fn create_node(&mut self, id: NodeId<'_>, options: NodeOptions<'_>) -> NodeResult<()>;
327
328 fn create_entity(&mut self, metadata: EntityMetadata) -> NodeResult<()>;
330
331 fn record_callback_effect(
333 &mut self,
334 callback_id: CallbackId<'_>,
335 kind: CallbackEffectKind,
336 entity_id: EntityId<'_>,
337 ) -> NodeResult<()>;
338}
339
340impl<const MAX_NODES: usize, const MAX_ENTITIES: usize, const MAX_CALLBACKS: usize> NodeRuntime
341 for MetadataRecorder<MAX_NODES, MAX_ENTITIES, MAX_CALLBACKS>
342{
343 fn create_node(&mut self, id: NodeId<'_>, options: NodeOptions<'_>) -> NodeResult<()> {
344 self.push_node(id, options.name, options.namespace, options.domain_id)?;
345 Ok(())
346 }
347
348 fn create_entity(&mut self, metadata: EntityMetadata) -> NodeResult<()> {
349 self.push_entity(metadata)?;
350 Ok(())
351 }
352
353 fn record_callback_effect(
354 &mut self,
355 callback_id: CallbackId<'_>,
356 kind: CallbackEffectKind,
357 entity_id: EntityId<'_>,
358 ) -> NodeResult<()> {
359 self.push_callback_effect(callback_id, kind, entity_id)?;
360 Ok(())
361 }
362}
363
364pub trait DeclaredNodeRuntime {
371 type NodeHandle: Copy + Eq;
373
374 fn build_component_node(
376 &mut self,
377 id: NodeId<'_>,
378 options: NodeOptions<'_>,
379 ) -> NodeResult<Self::NodeHandle>;
380}
381
382#[derive(Debug, Clone, PartialEq, Eq)]
384pub struct RuntimeNodeRecord<H: Copy + Eq> {
385 slot: NodeSlot,
386 stable_id: MetadataString,
387 source_default_name: MetadataString,
388 handle: H,
389}
390
391impl<H: Copy + Eq> RuntimeNodeRecord<H> {
392 pub const fn slot(&self) -> NodeSlot {
394 self.slot
395 }
396
397 pub fn stable_id(&self) -> &str {
399 &self.stable_id
400 }
401
402 pub fn source_default_name(&self) -> &str {
404 &self.source_default_name
405 }
406
407 pub const fn handle(&self) -> H {
409 self.handle
410 }
411}
412
413pub struct NodeRuntimeAdapter<
415 'a,
416 R: DeclaredNodeRuntime + ?Sized,
417 const MAX_NODES: usize = { crate::node_metadata::DEFAULT_MAX_METADATA_NODES },
418 const MAX_ENTITIES: usize = { crate::node_metadata::DEFAULT_MAX_METADATA_ENTITIES },
419 const MAX_CALLBACKS: usize = { crate::node_metadata::DEFAULT_MAX_METADATA_CALLBACKS },
420> {
421 node_runtime: &'a mut R,
422 nodes: Vec<RuntimeNodeRecord<R::NodeHandle>, MAX_NODES>,
423 entities: Vec<EntityMetadata, MAX_ENTITIES>,
424 callback_effects: Vec<CallbackEffectMetadata, MAX_CALLBACKS>,
425}
426
427impl<
428 'a,
429 R: DeclaredNodeRuntime + ?Sized,
430 const MAX_NODES: usize,
431 const MAX_ENTITIES: usize,
432 const MAX_CALLBACKS: usize,
433> NodeRuntimeAdapter<'a, R, MAX_NODES, MAX_ENTITIES, MAX_CALLBACKS>
434{
435 pub fn new(node_runtime: &'a mut R) -> Self {
437 Self {
438 node_runtime,
439 nodes: Vec::new(),
440 entities: Vec::new(),
441 callback_effects: Vec::new(),
442 }
443 }
444
445 pub fn nodes(&self) -> &[RuntimeNodeRecord<R::NodeHandle>] {
447 &self.nodes
448 }
449
450 pub fn entities(&self) -> &[EntityMetadata] {
452 &self.entities
453 }
454
455 pub fn callback_effects(&self) -> &[CallbackEffectMetadata] {
457 &self.callback_effects
458 }
459
460 pub fn node_handle(&self, stable_id: NodeId<'_>) -> Option<R::NodeHandle> {
462 self.nodes
463 .iter()
464 .find(|node| node.stable_id() == stable_id.as_str())
465 .map(RuntimeNodeRecord::handle)
466 }
467
468 fn contains_node(&self, stable_id: &str) -> bool {
469 self.nodes.iter().any(|node| node.stable_id() == stable_id)
470 }
471
472 fn contains_entity(&self, stable_id: &str) -> bool {
473 self.entities
474 .iter()
475 .any(|entity| entity.id.as_str() == stable_id)
476 }
477
478 fn node_slot_for_id(&self, stable_id: &str) -> Option<NodeSlot> {
479 self.nodes
480 .iter()
481 .find(|node| node.stable_id() == stable_id)
482 .map(RuntimeNodeRecord::slot)
483 }
484
485 fn entity_slot_for_id(&self, stable_id: &str) -> Option<EntitySlot> {
486 self.entities
487 .iter()
488 .find(|entity| entity.id.as_str() == stable_id)
489 .and_then(|entity| entity.slot)
490 }
491
492 fn callback_slot_for_current_entity(
493 &self,
494 id: &str,
495 current_callbacks: &mut Vec<MetadataString, 3>,
496 next_callback_slot: &mut usize,
497 ) -> CallbackSlot {
498 if let Some(slot) = self.callback_slot_for_id(id) {
499 return slot;
500 }
501 if let Some((index, _)) = current_callbacks
502 .iter()
503 .enumerate()
504 .find(|(_, callback_id)| callback_id.as_str() == id)
505 {
506 return CallbackSlot::new(self.callback_slot_count() + index);
507 }
508 let slot = CallbackSlot::new(*next_callback_slot);
509 let _ = current_callbacks
510 .push(copy_str(id).expect("callback ID already fits metadata string capacity"));
511 *next_callback_slot += 1;
512 slot
513 }
514
515 fn callback_slot_for_id(&self, id: &str) -> Option<CallbackSlot> {
516 let mut seen = Vec::<&str, MAX_CALLBACKS>::new();
517 for entity in &self.entities {
518 for callback_id in entity_callback_ids(entity) {
519 let Some(callback_id) = callback_id else {
520 continue;
521 };
522 let callback_id = callback_id.as_str();
523 if seen.contains(&callback_id) {
524 continue;
525 }
526 if callback_id == id {
527 return Some(CallbackSlot::new(seen.len()));
528 }
529 let _ = seen.push(callback_id);
530 }
531 }
532 None
533 }
534
535 fn callback_slot_count(&self) -> usize {
536 let mut seen = Vec::<&str, MAX_CALLBACKS>::new();
537 for entity in &self.entities {
538 for callback_id in entity_callback_ids(entity) {
539 let Some(callback_id) = callback_id else {
540 continue;
541 };
542 let callback_id = callback_id.as_str();
543 if !seen.contains(&callback_id) {
544 let _ = seen.push(callback_id);
545 }
546 }
547 }
548 seen.len()
549 }
550}
551
552impl<
553 R: DeclaredNodeRuntime + ?Sized,
554 const MAX_NODES: usize,
555 const MAX_ENTITIES: usize,
556 const MAX_CALLBACKS: usize,
557> NodeRuntime for NodeRuntimeAdapter<'_, R, MAX_NODES, MAX_ENTITIES, MAX_CALLBACKS>
558{
559 fn create_node(&mut self, id: NodeId<'_>, options: NodeOptions<'_>) -> NodeResult<()> {
560 if self.contains_node(id.as_str()) {
561 return Err(NodeMetadataError::DuplicateId.into());
562 }
563 let handle = self.node_runtime.build_component_node(id, options)?;
564 let slot = NodeSlot::new(self.nodes.len());
565 self.nodes
566 .push(RuntimeNodeRecord {
567 slot,
568 stable_id: copy_str(id.as_str())?,
569 source_default_name: copy_str(options.name)?,
570 handle,
571 })
572 .map_err(|_| NodeDeclError::Metadata(NodeMetadataError::Capacity))?;
573 Ok(())
574 }
575
576 fn create_entity(&mut self, mut metadata: EntityMetadata) -> NodeResult<()> {
577 if !self.contains_node(metadata.node_id.as_str()) {
578 return Err(NodeMetadataError::UnknownNode.into());
579 }
580 if self.contains_entity(metadata.id.as_str()) {
581 return Err(NodeMetadataError::DuplicateId.into());
582 }
583 metadata.slot = Some(EntitySlot::new(self.entities.len()));
584 metadata.node_slot = self.node_slot_for_id(&metadata.node_id);
585 let mut current_callbacks = Vec::<MetadataString, 3>::new();
586 let mut next_callback_slot = self.callback_slot_count();
587 metadata.callback_slot = metadata.callback_id.as_ref().map(|callback_id| {
588 self.callback_slot_for_current_entity(
589 callback_id.as_str(),
590 &mut current_callbacks,
591 &mut next_callback_slot,
592 )
593 });
594 metadata.action_cancel_callback_slot =
595 metadata
596 .action_cancel_callback_id
597 .as_ref()
598 .map(|callback_id| {
599 self.callback_slot_for_current_entity(
600 callback_id.as_str(),
601 &mut current_callbacks,
602 &mut next_callback_slot,
603 )
604 });
605 metadata.action_accepted_callback_slot =
606 metadata
607 .action_accepted_callback_id
608 .as_ref()
609 .map(|callback_id| {
610 self.callback_slot_for_current_entity(
611 callback_id.as_str(),
612 &mut current_callbacks,
613 &mut next_callback_slot,
614 )
615 });
616 self.entities
617 .push(metadata)
618 .map_err(|_| NodeDeclError::Metadata(NodeMetadataError::Capacity))?;
619 Ok(())
620 }
621
622 fn record_callback_effect(
623 &mut self,
624 callback_id: CallbackId<'_>,
625 kind: CallbackEffectKind,
626 entity_id: EntityId<'_>,
627 ) -> NodeResult<()> {
628 if !self.contains_entity(entity_id.as_str()) {
629 return Err(NodeMetadataError::UnknownEntity.into());
630 }
631 self.callback_effects
632 .push(CallbackEffectMetadata {
633 callback_id: copy_str(callback_id.as_str())?,
634 callback_slot: self.callback_slot_for_id(callback_id.as_str()),
635 kind,
636 entity_id: copy_str(entity_id.as_str())?,
637 entity_slot: self.entity_slot_for_id(entity_id.as_str()),
638 })
639 .map_err(|_| NodeDeclError::Metadata(NodeMetadataError::Capacity))?;
640 Ok(())
641 }
642}
643
644#[cfg(feature = "rmw-cffi")]
645impl DeclaredNodeRuntime for crate::Executor<'static> {
646 type NodeHandle = nros_node::executor::NodeId;
647
648 fn build_component_node(
649 &mut self,
650 _id: NodeId<'_>,
651 options: NodeOptions<'_>,
652 ) -> NodeResult<Self::NodeHandle> {
653 self.node_builder(options.name)
654 .namespace(options.namespace)
655 .domain_id(options.domain_id)
656 .build()
657 .map_err(|_| NodeDeclError::Runtime)
658 }
659}
660
661#[cfg(feature = "rmw-cffi")]
663pub type NodeExecutorRuntime<
664 'a,
665 const MAX_NODES: usize = { crate::node_metadata::DEFAULT_MAX_METADATA_NODES },
666 const MAX_ENTITIES: usize = { crate::node_metadata::DEFAULT_MAX_METADATA_ENTITIES },
667 const MAX_CALLBACKS: usize = { crate::node_metadata::DEFAULT_MAX_METADATA_CALLBACKS },
668> = NodeRuntimeAdapter<'a, crate::Executor<'static>, MAX_NODES, MAX_ENTITIES, MAX_CALLBACKS>;
669
670pub struct NodeContext<'a, R: NodeRuntime + ?Sized = dyn NodeRuntime + 'a> {
672 component_name: &'static str,
673 runtime: &'a mut R,
674 params: &'a [(&'a str, &'a str)],
681}
682
683impl<'a, R: NodeRuntime + ?Sized> NodeContext<'a, R> {
684 pub fn new(component_name: &'static str, runtime: &'a mut R) -> Self {
686 Self {
687 component_name,
688 runtime,
689 params: &[],
690 }
691 }
692
693 pub fn set_params(&mut self, params: &'a [(&'a str, &'a str)]) {
696 self.params = params;
697 }
698
699 pub fn param(&self, name: &str) -> Option<&'a str> {
703 self.params
704 .iter()
705 .find(|(k, _)| *k == name)
706 .map(|(_, v)| *v)
707 }
708
709 pub const fn component_name(&self) -> &'static str {
711 self.component_name
712 }
713
714 #[doc(hidden)]
719 pub fn create_node_with_id<'id>(
720 &mut self,
721 id: NodeId<'id>,
722 options: NodeOptions<'_>,
723 ) -> NodeResult<DeclaredNode<'_, 'id, R>> {
724 self.runtime.create_node(id, options)?;
725 Ok(DeclaredNode {
726 runtime: self.runtime,
727 id,
728 current_group: None,
729 })
730 }
731
732 pub fn create_node<'id>(
738 &mut self,
739 options: NodeOptions<'id>,
740 ) -> NodeResult<DeclaredNode<'_, 'id, R>> {
741 self.create_node_with_id(NodeId::new(options.name), options)
742 }
743
744 #[deprecated(note = "use create_node(NodeOptions)")]
746 pub fn create_node_with_options<'id>(
747 &mut self,
748 options: NodeOptions<'id>,
749 ) -> NodeResult<DeclaredNode<'_, 'id, R>> {
750 self.create_node(options)
751 }
752
753 #[doc(hidden)]
755 pub fn callback<'id>(&mut self, id: CallbackId<'id>) -> CallbackEffects<'_, 'id, R> {
756 CallbackEffects {
757 runtime: self.runtime,
758 id,
759 }
760 }
761}
762
763pub struct DeclaredNode<'ctx, 'id, R: NodeRuntime + ?Sized = dyn NodeRuntime + 'ctx> {
765 runtime: &'ctx mut R,
766 id: NodeId<'id>,
767 current_group: Option<MetadataString>,
774}
775
776impl<'ctx, 'id, R: NodeRuntime + ?Sized> DeclaredNode<'ctx, 'id, R> {
777 #[doc(hidden)]
779 pub const fn id(&self) -> NodeId<'id> {
780 self.id
781 }
782
783 #[track_caller]
790 pub fn callback_group(&mut self, group: &str) -> NodeResult<&mut Self> {
791 self.current_group = Some(copy_str(group)?);
792 Ok(self)
793 }
794
795 fn declare_entity(&mut self, mut metadata: EntityMetadata) -> NodeResult<()> {
800 if metadata.callback_group.is_none() {
801 metadata.callback_group = self.current_group.clone();
802 }
803 self.runtime.create_entity(metadata)
804 }
805
806 #[track_caller]
808 #[doc(hidden)]
809 pub fn create_publisher<'entity, M: MessageForRmw>(
810 &mut self,
811 id: EntityId<'entity>,
812 topic: &str,
813 ) -> NodeResult<NodePublisher<'entity, M>> {
814 self.create_publisher_with_qos::<M>(id, topic, QoSProfile::default())
815 }
816
817 #[track_caller]
823 pub fn create_publisher_for_topic<'entity, M: MessageForRmw>(
824 &mut self,
825 topic: &'entity str,
826 ) -> NodeResult<NodePublisher<'entity, M>> {
827 self.create_publisher_for_topic_with_qos::<M>(topic, QoSProfile::default())
828 }
829
830 #[track_caller]
832 pub fn create_publisher_for_topic_with_qos<'entity, M: MessageForRmw>(
833 &mut self,
834 topic: &'entity str,
835 qos: QoSProfile,
836 ) -> NodeResult<NodePublisher<'entity, M>> {
837 self.create_publisher_with_qos::<M>(EntityId::new(topic), topic, qos)
838 }
839
840 #[track_caller]
842 #[doc(hidden)]
843 pub fn create_publisher_with_qos<'entity, M: MessageForRmw>(
844 &mut self,
845 id: EntityId<'entity>,
846 topic: &str,
847 qos: QoSProfile,
848 ) -> NodeResult<NodePublisher<'entity, M>> {
849 register_declared_type::<M>()?;
850 let mut metadata = entity_metadata(EntityMetadataSpec {
851 id,
852 node_id: self.id,
853 kind: EntityKind::Publisher,
854 source_name: topic,
855 type_name: <M as RosMessage>::TYPE_NAME,
858 type_hash: <M as RosMessage>::TYPE_HASH,
859 qos,
860 })?;
861 metadata.source = SourceLocationMetadata::caller()?;
862 self.declare_entity(metadata)?;
863 Ok(NodePublisher::new(id))
864 }
865
866 #[track_caller]
868 #[doc(hidden)]
869 pub fn create_subscription<'entity, 'callback, M: MessageForRmw>(
870 &mut self,
871 id: EntityId<'entity>,
872 callback_id: CallbackId<'callback>,
873 topic: &str,
874 ) -> NodeResult<NodeSubscription<'entity, M>> {
875 self.create_subscription_with_qos::<M>(id, callback_id, topic, QoSProfile::default())
876 }
877
878 #[track_caller]
883 #[doc(hidden)]
884 pub fn create_subscription_for_callback<'callback, M: MessageForRmw>(
885 &mut self,
886 callback_id: CallbackId<'callback>,
887 topic: &str,
888 ) -> NodeResult<NodeSubscription<'callback, M>> {
889 self.create_subscription_for_callback_with_qos::<M>(
890 callback_id,
891 topic,
892 QoSProfile::default(),
893 )
894 }
895
896 #[track_caller]
899 pub fn create_subscription_for_callback_name<'callback, M: MessageForRmw>(
900 &mut self,
901 callback_name: &'callback str,
902 topic: &str,
903 ) -> NodeResult<NodeSubscription<'callback, M>> {
904 self.create_subscription_for_callback::<M>(CallbackId::new(callback_name), topic)
905 }
906
907 #[track_caller]
917 pub fn create_subscription_for_callback_name_with_safety<'callback, M: MessageForRmw>(
918 &mut self,
919 callback_name: &'callback str,
920 topic: &str,
921 ) -> NodeResult<NodeSubscription<'callback, M>> {
922 let callback_id = CallbackId::new(callback_name);
923 let id = EntityId::new(callback_id.as_str());
924 register_declared_type::<M>()?;
925 let mut metadata = entity_metadata(EntityMetadataSpec {
926 id,
927 node_id: self.id,
928 kind: EntityKind::Subscription,
929 source_name: topic,
930 type_name: <M as RosMessage>::TYPE_NAME,
933 type_hash: <M as RosMessage>::TYPE_HASH,
934 qos: QoSProfile::default(),
935 })?;
936 metadata.callback_id = Some(copy_str(callback_id.as_str())?);
937 metadata.callback_source = SourceLocationMetadata::caller()?;
938 metadata.source = metadata.callback_source.clone();
939 metadata.safety = true;
940 self.declare_entity(metadata)?;
941 Ok(NodeSubscription::new(id))
942 }
943
944 #[track_caller]
946 #[doc(hidden)]
947 pub fn create_subscription_for_callback_with_qos<'callback, M: MessageForRmw>(
948 &mut self,
949 callback_id: CallbackId<'callback>,
950 topic: &str,
951 qos: QoSProfile,
952 ) -> NodeResult<NodeSubscription<'callback, M>> {
953 self.create_subscription_with_qos::<M>(
954 EntityId::new(callback_id.as_str()),
955 callback_id,
956 topic,
957 qos,
958 )
959 }
960
961 #[track_caller]
963 pub fn create_subscription_for_topic<'entity, M: MessageForRmw>(
964 &mut self,
965 topic: &'entity str,
966 ) -> NodeResult<NodeSubscription<'entity, M>> {
967 self.create_subscription_for_topic_with_qos::<M>(topic, QoSProfile::default())
968 }
969
970 #[track_caller]
972 pub fn create_subscription_for_topic_with_qos<'entity, M: MessageForRmw>(
973 &mut self,
974 topic: &'entity str,
975 qos: QoSProfile,
976 ) -> NodeResult<NodeSubscription<'entity, M>> {
977 self.create_subscription_with_qos::<M>(
978 EntityId::new(topic),
979 CallbackId::new(topic),
980 topic,
981 qos,
982 )
983 }
984
985 #[track_caller]
987 #[doc(hidden)]
988 pub fn create_subscription_with_qos<'entity, 'callback, M: MessageForRmw>(
989 &mut self,
990 id: EntityId<'entity>,
991 callback_id: CallbackId<'callback>,
992 topic: &str,
993 qos: QoSProfile,
994 ) -> NodeResult<NodeSubscription<'entity, M>> {
995 const {
1017 assert!(
1018 nros_node::rmw_type_registry::subscription_buffer_ok::<M>(),
1019 "this message type's maximum serialized size exceeds \
1020 NROS_SUBSCRIPTION_BUFFER_SIZE — every sample would be received, \
1021 ACKed and then DROPPED. Raise the knob to at least the type's \
1022 bound (`<M as nros_serdes::schema::Message>::MAX_SERIALIZED_SIZE_XCDR2`)."
1023 )
1024 }
1025 register_declared_type::<M>()?;
1026 let mut metadata = entity_metadata(EntityMetadataSpec {
1027 id,
1028 node_id: self.id,
1029 kind: EntityKind::Subscription,
1030 source_name: topic,
1031 type_name: <M as RosMessage>::TYPE_NAME,
1034 type_hash: <M as RosMessage>::TYPE_HASH,
1035 qos,
1036 })?;
1037 metadata.callback_id = Some(copy_str(callback_id.as_str())?);
1038 metadata.callback_source = SourceLocationMetadata::caller()?;
1039 metadata.source = metadata.callback_source.clone();
1040 self.declare_entity(metadata)?;
1041 Ok(NodeSubscription::new(id))
1042 }
1043
1044 #[track_caller]
1056 pub fn create_subscription_static<M: MessageForRmw>(
1057 &mut self,
1058 topic: &'static str,
1059 ) -> NodeResult<SubscriptionTag> {
1060 let id = EntityId::new(topic);
1061 let callback_id = CallbackId::new(topic);
1062 register_declared_type::<M>()?;
1063 let mut metadata = entity_metadata(EntityMetadataSpec {
1064 id,
1065 node_id: self.id,
1066 kind: EntityKind::Subscription,
1067 source_name: topic,
1068 type_name: <M as RosMessage>::TYPE_NAME,
1071 type_hash: <M as RosMessage>::TYPE_HASH,
1072 qos: QoSProfile::default(),
1073 })?;
1074 metadata.callback_id = Some(copy_str(callback_id.as_str())?);
1075 metadata.callback_source = SourceLocationMetadata::caller()?;
1076 metadata.source = metadata.callback_source.clone();
1077 self.declare_entity(metadata)?;
1078 Ok(SubscriptionTag::new(topic))
1079 }
1080
1081 #[track_caller]
1083 #[doc(hidden)]
1084 pub fn create_timer<'entity, 'callback>(
1085 &mut self,
1086 id: EntityId<'entity>,
1087 callback_id: CallbackId<'callback>,
1088 period: TimerDuration,
1089 ) -> NodeResult<NodeTimer<'entity>> {
1090 let mut metadata = entity_metadata(EntityMetadataSpec {
1091 id,
1092 node_id: self.id,
1093 kind: EntityKind::Timer,
1094 source_name: "",
1095 type_name: "",
1096 type_hash: "",
1097 qos: QoSProfile::default(),
1098 })?;
1099 metadata.callback_id = Some(copy_str(callback_id.as_str())?);
1100 metadata.callback_source = SourceLocationMetadata::caller()?;
1101 metadata.source = metadata.callback_source.clone();
1102 metadata.period_ms = Some(period.as_millis());
1103 metadata.period_us = Some(period.as_micros());
1104 self.declare_entity(metadata)?;
1105 Ok(NodeTimer::new(id))
1106 }
1107
1108 #[track_caller]
1110 #[doc(hidden)]
1111 pub fn create_timer_for_callback<'callback>(
1112 &mut self,
1113 callback_id: CallbackId<'callback>,
1114 period: TimerDuration,
1115 ) -> NodeResult<NodeTimer<'callback>> {
1116 self.create_timer(EntityId::new(callback_id.as_str()), callback_id, period)
1117 }
1118
1119 #[track_caller]
1122 pub fn create_timer_for_callback_name<'callback>(
1123 &mut self,
1124 callback_name: &'callback str,
1125 period: TimerDuration,
1126 ) -> NodeResult<NodeTimer<'callback>> {
1127 self.create_timer_for_callback(CallbackId::new(callback_name), period)
1128 }
1129
1130 #[track_caller]
1132 #[doc(hidden)]
1133 pub fn create_service_server<
1134 'entity,
1135 'callback,
1136 S: RosService<
1137 Request: nros_node::rmw_type_registry::MessageForRmw,
1138 Reply: nros_node::rmw_type_registry::MessageForRmw,
1139 >,
1140 >(
1141 &mut self,
1142 id: EntityId<'entity>,
1143 callback_id: CallbackId<'callback>,
1144 service_name: &str,
1145 ) -> NodeResult<NodeServiceServer<'entity, S>> {
1146 register_declared_service::<S>()?;
1147 let mut metadata = entity_metadata(EntityMetadataSpec {
1148 id,
1149 node_id: self.id,
1150 kind: EntityKind::ServiceServer,
1151 source_name: service_name,
1152 type_name: S::SERVICE_NAME,
1153 type_hash: S::SERVICE_HASH,
1154 qos: QoSProfile::default(),
1155 })?;
1156 metadata.callback_id = Some(copy_str(callback_id.as_str())?);
1157 metadata.callback_source = SourceLocationMetadata::caller()?;
1158 metadata.source = metadata.callback_source.clone();
1159 self.declare_entity(metadata)?;
1160 Ok(NodeServiceServer::new(id))
1161 }
1162
1163 #[track_caller]
1166 pub fn create_service_server_for_name<
1167 'entity,
1168 S: RosService<
1169 Request: nros_node::rmw_type_registry::MessageForRmw,
1170 Reply: nros_node::rmw_type_registry::MessageForRmw,
1171 >,
1172 >(
1173 &mut self,
1174 name: &'entity str,
1175 ) -> NodeResult<NodeServiceServer<'entity, S>> {
1176 self.create_service_server::<S>(EntityId::new(name), CallbackId::new(name), name)
1177 }
1178
1179 #[track_caller]
1182 pub fn create_service_server_for_name_with_callback<
1183 'entity,
1184 S: RosService<
1185 Request: nros_node::rmw_type_registry::MessageForRmw,
1186 Reply: nros_node::rmw_type_registry::MessageForRmw,
1187 >,
1188 >(
1189 &mut self,
1190 name: &'entity str,
1191 callback_name: &str,
1192 ) -> NodeResult<NodeServiceServer<'entity, S>> {
1193 self.create_service_server::<S>(EntityId::new(name), CallbackId::new(callback_name), name)
1194 }
1195
1196 #[track_caller]
1208 pub fn create_service_static<
1209 S: RosService<
1210 Request: nros_node::rmw_type_registry::MessageForRmw,
1211 Reply: nros_node::rmw_type_registry::MessageForRmw,
1212 >,
1213 >(
1214 &mut self,
1215 name: &'static str,
1216 ) -> NodeResult<ServiceTag> {
1217 self.create_service_server_for_name::<S>(name)?;
1218 Ok(ServiceTag::new(name))
1219 }
1220
1221 #[track_caller]
1223 #[doc(hidden)]
1224 pub fn create_service_client<
1225 'entity,
1226 S: RosService<
1227 Request: nros_node::rmw_type_registry::MessageForRmw,
1228 Reply: nros_node::rmw_type_registry::MessageForRmw,
1229 >,
1230 >(
1231 &mut self,
1232 id: EntityId<'entity>,
1233 service_name: &str,
1234 ) -> NodeResult<NodeServiceClient<'entity, S>> {
1235 register_declared_service::<S>()?;
1236 let mut metadata = entity_metadata(EntityMetadataSpec {
1237 id,
1238 node_id: self.id,
1239 kind: EntityKind::ServiceClient,
1240 source_name: service_name,
1241 type_name: S::SERVICE_NAME,
1242 type_hash: S::SERVICE_HASH,
1243 qos: QoSProfile::default(),
1244 })?;
1245 metadata.source = SourceLocationMetadata::caller()?;
1246 self.declare_entity(metadata)?;
1247 Ok(NodeServiceClient::new(id))
1248 }
1249
1250 #[track_caller]
1252 pub fn create_service_client_for_name<
1253 'entity,
1254 S: RosService<
1255 Request: nros_node::rmw_type_registry::MessageForRmw,
1256 Reply: nros_node::rmw_type_registry::MessageForRmw,
1257 >,
1258 >(
1259 &mut self,
1260 name: &'entity str,
1261 ) -> NodeResult<NodeServiceClient<'entity, S>> {
1262 self.create_service_client::<S>(EntityId::new(name), name)
1263 }
1264
1265 #[track_caller]
1267 #[doc(hidden)]
1268 pub fn create_action_server<
1269 'entity,
1270 'callback,
1271 A: RosAction<
1272 Goal: nros_node::rmw_type_registry::MessageForRmw,
1273 Result: nros_node::rmw_type_registry::MessageForRmw,
1274 Feedback: nros_node::rmw_type_registry::MessageForRmw,
1275 SendGoalRequest: nros_node::rmw_type_registry::MessageForRmw,
1276 SendGoalResponse: nros_node::rmw_type_registry::MessageForRmw,
1277 GetResultRequest: nros_node::rmw_type_registry::MessageForRmw,
1278 GetResultResponse: nros_node::rmw_type_registry::MessageForRmw,
1279 FeedbackMessage: nros_node::rmw_type_registry::MessageForRmw,
1280 >,
1281 >(
1282 &mut self,
1283 id: EntityId<'entity>,
1284 callback_id: CallbackId<'callback>,
1285 action_name: &str,
1286 ) -> NodeResult<NodeActionServer<'entity, A>> {
1287 self.create_action_server_with_callbacks::<A>(
1288 id,
1289 callback_id,
1290 callback_id,
1291 callback_id,
1292 action_name,
1293 )
1294 }
1295
1296 #[track_caller]
1298 #[doc(hidden)]
1299 pub fn create_action_server_with_callbacks<
1300 'entity,
1301 'goal,
1302 'cancel,
1303 'accepted,
1304 A: RosAction<
1305 Goal: nros_node::rmw_type_registry::MessageForRmw,
1306 Result: nros_node::rmw_type_registry::MessageForRmw,
1307 Feedback: nros_node::rmw_type_registry::MessageForRmw,
1308 SendGoalRequest: nros_node::rmw_type_registry::MessageForRmw,
1309 SendGoalResponse: nros_node::rmw_type_registry::MessageForRmw,
1310 GetResultRequest: nros_node::rmw_type_registry::MessageForRmw,
1311 GetResultResponse: nros_node::rmw_type_registry::MessageForRmw,
1312 FeedbackMessage: nros_node::rmw_type_registry::MessageForRmw,
1313 >,
1314 >(
1315 &mut self,
1316 id: EntityId<'entity>,
1317 goal_callback_id: CallbackId<'goal>,
1318 cancel_callback_id: CallbackId<'cancel>,
1319 accepted_callback_id: CallbackId<'accepted>,
1320 action_name: &str,
1321 ) -> NodeResult<NodeActionServer<'entity, A>> {
1322 register_declared_action::<A>()?;
1323 let mut metadata = entity_metadata(EntityMetadataSpec {
1324 id,
1325 node_id: self.id,
1326 kind: EntityKind::ActionServer,
1327 source_name: action_name,
1328 type_name: A::ACTION_NAME,
1329 type_hash: A::ACTION_HASH,
1330 qos: QoSProfile::default(),
1331 })?;
1332 metadata.callback_id = Some(copy_str(goal_callback_id.as_str())?);
1333 metadata.callback_source = SourceLocationMetadata::caller()?;
1334 metadata.action_cancel_callback_id = Some(copy_str(cancel_callback_id.as_str())?);
1335 metadata.action_cancel_source = metadata.callback_source.clone();
1336 metadata.action_accepted_callback_id = Some(copy_str(accepted_callback_id.as_str())?);
1337 metadata.action_accepted_source = metadata.callback_source.clone();
1338 metadata.source = metadata.callback_source.clone();
1339 self.declare_entity(metadata)?;
1340 Ok(NodeActionServer::new(id))
1341 }
1342
1343 #[track_caller]
1346 pub fn create_action_server_for_name<
1347 'entity,
1348 A: RosAction<
1349 Goal: nros_node::rmw_type_registry::MessageForRmw,
1350 Result: nros_node::rmw_type_registry::MessageForRmw,
1351 Feedback: nros_node::rmw_type_registry::MessageForRmw,
1352 SendGoalRequest: nros_node::rmw_type_registry::MessageForRmw,
1353 SendGoalResponse: nros_node::rmw_type_registry::MessageForRmw,
1354 GetResultRequest: nros_node::rmw_type_registry::MessageForRmw,
1355 GetResultResponse: nros_node::rmw_type_registry::MessageForRmw,
1356 FeedbackMessage: nros_node::rmw_type_registry::MessageForRmw,
1357 >,
1358 >(
1359 &mut self,
1360 name: &'entity str,
1361 ) -> NodeResult<NodeActionServer<'entity, A>> {
1362 self.create_action_server::<A>(EntityId::new(name), CallbackId::new(name), name)
1363 }
1364
1365 #[track_caller]
1368 pub fn create_action_server_for_name_with_callbacks<
1369 'entity,
1370 A: RosAction<
1371 Goal: nros_node::rmw_type_registry::MessageForRmw,
1372 Result: nros_node::rmw_type_registry::MessageForRmw,
1373 Feedback: nros_node::rmw_type_registry::MessageForRmw,
1374 SendGoalRequest: nros_node::rmw_type_registry::MessageForRmw,
1375 SendGoalResponse: nros_node::rmw_type_registry::MessageForRmw,
1376 GetResultRequest: nros_node::rmw_type_registry::MessageForRmw,
1377 GetResultResponse: nros_node::rmw_type_registry::MessageForRmw,
1378 FeedbackMessage: nros_node::rmw_type_registry::MessageForRmw,
1379 >,
1380 >(
1381 &mut self,
1382 name: &'entity str,
1383 goal_callback_name: &str,
1384 cancel_callback_name: &str,
1385 accepted_callback_name: &str,
1386 ) -> NodeResult<NodeActionServer<'entity, A>> {
1387 self.create_action_server_with_callbacks::<A>(
1388 EntityId::new(name),
1389 CallbackId::new(goal_callback_name),
1390 CallbackId::new(cancel_callback_name),
1391 CallbackId::new(accepted_callback_name),
1392 name,
1393 )
1394 }
1395
1396 #[track_caller]
1412 pub fn create_action_static<
1413 A: RosAction<
1414 Goal: nros_node::rmw_type_registry::MessageForRmw,
1415 Result: nros_node::rmw_type_registry::MessageForRmw,
1416 Feedback: nros_node::rmw_type_registry::MessageForRmw,
1417 SendGoalRequest: nros_node::rmw_type_registry::MessageForRmw,
1418 SendGoalResponse: nros_node::rmw_type_registry::MessageForRmw,
1419 GetResultRequest: nros_node::rmw_type_registry::MessageForRmw,
1420 GetResultResponse: nros_node::rmw_type_registry::MessageForRmw,
1421 FeedbackMessage: nros_node::rmw_type_registry::MessageForRmw,
1422 >,
1423 >(
1424 &mut self,
1425 name: &'static str,
1426 ) -> NodeResult<ActionTag> {
1427 self.create_action_server_for_name::<A>(name)?;
1428 Ok(ActionTag::new(name))
1429 }
1430
1431 #[track_caller]
1433 #[doc(hidden)]
1434 pub fn create_action_client<
1435 'entity,
1436 A: RosAction<
1437 Goal: nros_node::rmw_type_registry::MessageForRmw,
1438 Result: nros_node::rmw_type_registry::MessageForRmw,
1439 Feedback: nros_node::rmw_type_registry::MessageForRmw,
1440 SendGoalRequest: nros_node::rmw_type_registry::MessageForRmw,
1441 SendGoalResponse: nros_node::rmw_type_registry::MessageForRmw,
1442 GetResultRequest: nros_node::rmw_type_registry::MessageForRmw,
1443 GetResultResponse: nros_node::rmw_type_registry::MessageForRmw,
1444 FeedbackMessage: nros_node::rmw_type_registry::MessageForRmw,
1445 >,
1446 >(
1447 &mut self,
1448 id: EntityId<'entity>,
1449 action_name: &str,
1450 ) -> NodeResult<NodeActionClient<'entity, A>> {
1451 register_declared_action::<A>()?;
1452 let mut metadata = entity_metadata(EntityMetadataSpec {
1453 id,
1454 node_id: self.id,
1455 kind: EntityKind::ActionClient,
1456 source_name: action_name,
1457 type_name: A::ACTION_NAME,
1458 type_hash: A::ACTION_HASH,
1459 qos: QoSProfile::default(),
1460 })?;
1461 metadata.source = SourceLocationMetadata::caller()?;
1462 self.declare_entity(metadata)?;
1463 Ok(NodeActionClient::new(id))
1464 }
1465
1466 #[track_caller]
1468 pub fn create_action_client_for_name<
1469 'entity,
1470 A: RosAction<
1471 Goal: nros_node::rmw_type_registry::MessageForRmw,
1472 Result: nros_node::rmw_type_registry::MessageForRmw,
1473 Feedback: nros_node::rmw_type_registry::MessageForRmw,
1474 SendGoalRequest: nros_node::rmw_type_registry::MessageForRmw,
1475 SendGoalResponse: nros_node::rmw_type_registry::MessageForRmw,
1476 GetResultRequest: nros_node::rmw_type_registry::MessageForRmw,
1477 GetResultResponse: nros_node::rmw_type_registry::MessageForRmw,
1478 FeedbackMessage: nros_node::rmw_type_registry::MessageForRmw,
1479 >,
1480 >(
1481 &mut self,
1482 name: &'entity str,
1483 ) -> NodeResult<NodeActionClient<'entity, A>> {
1484 self.create_action_client::<A>(EntityId::new(name), name)
1485 }
1486
1487 #[track_caller]
1500 pub fn create_action_client_with_callbacks_for_name<
1501 'entity,
1502 A: RosAction<
1503 Goal: nros_node::rmw_type_registry::MessageForRmw,
1504 Result: nros_node::rmw_type_registry::MessageForRmw,
1505 Feedback: nros_node::rmw_type_registry::MessageForRmw,
1506 SendGoalRequest: nros_node::rmw_type_registry::MessageForRmw,
1507 SendGoalResponse: nros_node::rmw_type_registry::MessageForRmw,
1508 GetResultRequest: nros_node::rmw_type_registry::MessageForRmw,
1509 GetResultResponse: nros_node::rmw_type_registry::MessageForRmw,
1510 FeedbackMessage: nros_node::rmw_type_registry::MessageForRmw,
1511 >,
1512 >(
1513 &mut self,
1514 name: &'entity str,
1515 result_callback_name: &str,
1516 feedback_callback_name: &str,
1517 ) -> NodeResult<NodeActionClient<'entity, A>> {
1518 register_declared_action::<A>()?;
1519 let mut metadata = entity_metadata(EntityMetadataSpec {
1520 id: EntityId::new(name),
1521 node_id: self.id,
1522 kind: EntityKind::ActionClient,
1523 source_name: name,
1524 type_name: A::ACTION_NAME,
1525 type_hash: A::ACTION_HASH,
1526 qos: QoSProfile::default(),
1527 })?;
1528 metadata.callback_id = Some(copy_str(result_callback_name)?);
1529 metadata.action_accepted_callback_id = Some(copy_str(feedback_callback_name)?);
1530 metadata.callback_source = SourceLocationMetadata::caller()?;
1531 metadata.source = metadata.callback_source.clone();
1532 self.declare_entity(metadata)?;
1533 Ok(NodeActionClient::new(EntityId::new(name)))
1534 }
1535
1536 #[track_caller]
1538 #[doc(hidden)]
1539 pub fn declare_parameter<'entity>(
1540 &mut self,
1541 id: EntityId<'entity>,
1542 name: &str,
1543 parameter_type: ParameterType,
1544 ) -> NodeResult<NodeParameter<'entity>> {
1545 self.declare_parameter_with_default(id, name, ParameterDefault::for_type(parameter_type)?)
1546 }
1547
1548 #[track_caller]
1550 #[doc(hidden)]
1551 pub fn declare_parameter_with_default<'entity>(
1552 &mut self,
1553 id: EntityId<'entity>,
1554 name: &str,
1555 default: ParameterDefault,
1556 ) -> NodeResult<NodeParameter<'entity>> {
1557 let mut metadata = entity_metadata(EntityMetadataSpec {
1558 id,
1559 node_id: self.id,
1560 kind: EntityKind::Parameter,
1561 source_name: name,
1562 type_name: "",
1563 type_hash: "",
1564 qos: QoSProfile::default(),
1565 })?;
1566 metadata.parameter_type = Some(default.parameter_type());
1567 metadata.parameter_default = Some(default);
1568 metadata.source = SourceLocationMetadata::caller()?;
1569 self.declare_entity(metadata)?;
1570 Ok(NodeParameter::new(id))
1571 }
1572
1573 #[track_caller]
1575 pub fn declare_parameter_for_name<'entity>(
1576 &mut self,
1577 name: &'entity str,
1578 parameter_type: ParameterType,
1579 ) -> NodeResult<NodeParameter<'entity>> {
1580 self.declare_parameter(EntityId::new(name), name, parameter_type)
1581 }
1582
1583 #[track_caller]
1586 pub fn declare_parameter_for_name_with_default<'entity>(
1587 &mut self,
1588 name: &'entity str,
1589 default: ParameterDefault,
1590 ) -> NodeResult<NodeParameter<'entity>> {
1591 self.declare_parameter_with_default(EntityId::new(name), name, default)
1592 }
1593
1594 #[doc(hidden)]
1596 pub fn callback<'callback>(
1597 &mut self,
1598 id: CallbackId<'callback>,
1599 ) -> CallbackEffects<'_, 'callback, R> {
1600 CallbackEffects {
1601 runtime: self.runtime,
1602 id,
1603 }
1604 }
1605
1606 pub fn callback_for_name<'callback>(
1609 &mut self,
1610 name: &'callback str,
1611 ) -> CallbackEffects<'_, 'callback, R> {
1612 self.callback(CallbackId::new(name))
1613 }
1614}
1615
1616pub struct CallbackEffects<'ctx, 'id, R: NodeRuntime + ?Sized = dyn NodeRuntime + 'ctx> {
1618 runtime: &'ctx mut R,
1619 id: CallbackId<'id>,
1620}
1621
1622impl<'ctx, 'id, R: NodeRuntime + ?Sized> CallbackEffects<'ctx, 'id, R> {
1623 #[doc(hidden)]
1625 pub fn reads(self, entity_id: EntityId<'_>) -> NodeResult<Self> {
1626 self.runtime
1627 .record_callback_effect(self.id, CallbackEffectKind::Reads, entity_id)?;
1628 Ok(self)
1629 }
1630
1631 pub fn reads_entity(self, entity: &impl DeclaredEntity) -> NodeResult<Self> {
1633 self.reads(entity.entity_id())
1634 }
1635
1636 #[doc(hidden)]
1638 pub fn publishes(self, entity_id: EntityId<'_>) -> NodeResult<Self> {
1639 self.runtime
1640 .record_callback_effect(self.id, CallbackEffectKind::Publishes, entity_id)?;
1641 Ok(self)
1642 }
1643
1644 pub fn publishes_entity(self, entity: &impl DeclaredEntity) -> NodeResult<Self> {
1646 self.publishes(entity.entity_id())
1647 }
1648
1649 #[doc(hidden)]
1651 pub fn writes(self, entity_id: EntityId<'_>) -> NodeResult<Self> {
1652 self.runtime
1653 .record_callback_effect(self.id, CallbackEffectKind::Writes, entity_id)?;
1654 Ok(self)
1655 }
1656
1657 pub fn writes_entity(self, entity: &impl DeclaredEntity) -> NodeResult<Self> {
1659 self.writes(entity.entity_id())
1660 }
1661}
1662
1663#[doc(hidden)]
1665pub trait DeclaredEntity {
1666 fn entity_id(&self) -> EntityId<'_>;
1668}
1669
1670macro_rules! component_handle {
1671 ($name:ident $(, $type_param:ident)?) => {
1672 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
1674 pub struct $name<'id $(, $type_param)?> {
1675 id: EntityId<'id>,
1676 _marker: PhantomData<($($type_param,)?)>,
1677 }
1678
1679 impl<'id $(, $type_param)?> $name<'id $(, $type_param)?> {
1680 const fn new(id: EntityId<'id>) -> Self {
1681 Self {
1682 id,
1683 _marker: PhantomData,
1684 }
1685 }
1686
1687 #[doc(hidden)]
1689 pub const fn id(&self) -> EntityId<'id> {
1690 self.id
1691 }
1692 }
1693
1694 impl<'id $(, $type_param)?> DeclaredEntity for $name<'id $(, $type_param)?> {
1695 fn entity_id(&self) -> EntityId<'_> {
1696 self.id
1697 }
1698 }
1699 };
1700}
1701
1702component_handle!(NodePublisher, M);
1703component_handle!(NodeSubscription, M);
1704component_handle!(NodeServiceServer, S);
1705component_handle!(NodeServiceClient, S);
1706component_handle!(NodeActionServer, A);
1707component_handle!(NodeActionClient, A);
1708component_handle!(NodeTimer);
1709component_handle!(NodeParameter);
1710
1711pub trait PublisherResolver {
1732 fn publish_raw(&self, entity_id: &str, data: &[u8]) -> NodeResult<()>;
1736}
1737
1738struct ReplySink<'a> {
1742 buf: &'a mut [u8],
1743 written: &'a mut usize,
1744}
1745
1746enum DecisionSink<'a> {
1752 Goal(&'a mut GoalResponse),
1753 Cancel(&'a mut CancelResponse),
1754}
1755
1756const GOAL_UUID_LEN: usize = 16;
1768
1769pub struct CallbackCtx<'a> {
1770 payload: &'a [u8],
1771 publishers: &'a dyn PublisherResolver,
1772 reply: Option<ReplySink<'a>>,
1773 decision: Option<DecisionSink<'a>>,
1774 #[cfg(feature = "safety-e2e")]
1780 integrity: Option<&'a crate::IntegrityStatus>,
1781 #[cfg(feature = "param-services")]
1786 params: Option<&'a crate::ParameterServer<'a>>,
1787}
1788
1789impl<'a> CallbackCtx<'a> {
1790 pub fn new(payload: &'a [u8], publishers: &'a dyn PublisherResolver) -> Self {
1793 Self {
1794 payload,
1795 publishers,
1796 reply: None,
1797 decision: None,
1798 #[cfg(feature = "safety-e2e")]
1799 integrity: None,
1800 #[cfg(feature = "param-services")]
1801 params: None,
1802 }
1803 }
1804
1805 #[cfg(feature = "safety-e2e")]
1811 pub fn new_with_integrity(
1812 payload: &'a [u8],
1813 publishers: &'a dyn PublisherResolver,
1814 integrity: &'a crate::IntegrityStatus,
1815 ) -> Self {
1816 Self {
1817 payload,
1818 publishers,
1819 reply: None,
1820 decision: None,
1821 integrity: Some(integrity),
1822 #[cfg(feature = "param-services")]
1823 params: None,
1824 }
1825 }
1826
1827 pub fn with_reply(
1831 payload: &'a [u8],
1832 publishers: &'a dyn PublisherResolver,
1833 reply_buf: &'a mut [u8],
1834 reply_written: &'a mut usize,
1835 ) -> Self {
1836 *reply_written = 0;
1837 Self {
1838 payload,
1839 publishers,
1840 reply: Some(ReplySink {
1841 buf: reply_buf,
1842 written: reply_written,
1843 }),
1844 decision: None,
1845 #[cfg(feature = "safety-e2e")]
1846 integrity: None,
1847 #[cfg(feature = "param-services")]
1848 params: None,
1849 }
1850 }
1851
1852 pub fn with_goal_decision(
1856 payload: &'a [u8],
1857 publishers: &'a dyn PublisherResolver,
1858 out: &'a mut GoalResponse,
1859 ) -> Self {
1860 Self {
1861 payload,
1862 publishers,
1863 reply: None,
1864 decision: Some(DecisionSink::Goal(out)),
1865 #[cfg(feature = "safety-e2e")]
1866 integrity: None,
1867 #[cfg(feature = "param-services")]
1868 params: None,
1869 }
1870 }
1871
1872 pub fn with_cancel_decision(
1875 payload: &'a [u8],
1876 publishers: &'a dyn PublisherResolver,
1877 out: &'a mut CancelResponse,
1878 ) -> Self {
1879 Self {
1880 payload,
1881 publishers,
1882 reply: None,
1883 decision: Some(DecisionSink::Cancel(out)),
1884 #[cfg(feature = "safety-e2e")]
1885 integrity: None,
1886 #[cfg(feature = "param-services")]
1887 params: None,
1888 }
1889 }
1890
1891 #[cfg(feature = "param-services")]
1895 pub fn set_param_server(&mut self, params: Option<&'a crate::ParameterServer<'a>>) {
1896 self.params = params;
1897 }
1898
1899 #[cfg(feature = "param-services")]
1904 pub fn parameter<T: crate::ParameterVariant>(&self, name: &str) -> Option<T> {
1905 self.params
1906 .and_then(|server| server.get(name))
1907 .and_then(T::from_parameter_value)
1908 }
1909
1910 pub fn set_goal_response(&mut self, response: GoalResponse) -> NodeResult<()> {
1913 match &mut self.decision {
1914 Some(DecisionSink::Goal(slot)) => {
1915 **slot = response;
1916 Ok(())
1917 }
1918 _ => Err(NodeDeclError::Runtime),
1919 }
1920 }
1921
1922 pub fn set_cancel_response(&mut self, response: CancelResponse) -> NodeResult<()> {
1933 match &mut self.decision {
1934 Some(DecisionSink::Cancel(slot)) => {
1935 **slot = response;
1936 Ok(())
1937 }
1938 _ => Err(NodeDeclError::Runtime),
1939 }
1940 }
1941
1942 pub fn reply_raw(&mut self, data: &[u8]) -> NodeResult<()> {
1946 let sink = self.reply.as_mut().ok_or(NodeDeclError::Runtime)?;
1947 if data.len() > sink.buf.len() {
1948 return Err(NodeDeclError::Runtime);
1949 }
1950 sink.buf[..data.len()].copy_from_slice(data);
1951 *sink.written = data.len();
1952 Ok(())
1953 }
1954
1955 pub fn reply<M: RosMessage, const N: usize>(&mut self, msg: &M) -> NodeResult<()> {
1957 let mut buf = [0u8; N];
1958 let mut writer =
1959 crate::CdrWriter::new_with_header(&mut buf).map_err(|_| NodeDeclError::Runtime)?;
1960 msg.serialize(&mut writer)
1961 .map_err(|_| NodeDeclError::Runtime)?;
1962 let len = writer.position();
1963 self.reply_raw(&buf[..len])
1964 }
1965
1966 pub fn payload(&self) -> &[u8] {
1968 self.payload
1969 }
1970
1971 #[cfg(feature = "safety-e2e")]
1977 pub fn integrity(&self) -> Option<&crate::IntegrityStatus> {
1978 self.integrity
1979 }
1980
1981 pub fn message<M: RosMessage>(&self) -> NodeResult<M> {
1984 let mut reader =
1985 crate::CdrReader::new_with_header(self.payload).map_err(|_| NodeDeclError::Runtime)?;
1986 if matches!(self.decision, Some(DecisionSink::Goal(_))) {
1996 for _ in 0..GOAL_UUID_LEN {
1997 let _ = reader.read_u8();
1998 }
1999 }
2000 M::deserialize(&mut reader).map_err(|_| NodeDeclError::Runtime)
2001 }
2002
2003 #[doc(hidden)]
2005 pub fn publish_raw(&self, publisher: EntityId<'_>, data: &[u8]) -> NodeResult<()> {
2006 self.publishers.publish_raw(publisher.as_str(), data)
2007 }
2008
2009 #[doc(hidden)]
2013 pub fn publish<M: RosMessage, const N: usize>(
2014 &self,
2015 publisher: EntityId<'_>,
2016 msg: &M,
2017 ) -> NodeResult<()> {
2018 let mut buf = [0u8; N];
2019 let mut writer =
2020 crate::CdrWriter::new_with_header(&mut buf).map_err(|_| NodeDeclError::Runtime)?;
2021 msg.serialize(&mut writer)
2022 .map_err(|_| NodeDeclError::Runtime)?;
2023 let len = writer.position();
2024 self.publish_raw(publisher, &buf[..len])
2025 }
2026
2027 pub fn publish_to_topic<M: RosMessage, const N: usize>(
2034 &self,
2035 topic: &str,
2036 msg: &M,
2037 ) -> NodeResult<()> {
2038 self.publish::<M, N>(EntityId::new(topic), msg)
2039 }
2040}
2041
2042pub trait ActionExecutor {
2059 fn complete_goal_raw(
2061 &mut self,
2062 action_entity: &str,
2063 goal_id: &GoalId,
2064 status: GoalStatus,
2065 result: &[u8],
2066 ) -> NodeResult<()>;
2067
2068 fn publish_feedback_raw(
2070 &mut self,
2071 action_entity: &str,
2072 goal_id: &GoalId,
2073 feedback: &[u8],
2074 ) -> NodeResult<()>;
2075
2076 fn for_each_active_goal(&self, action_entity: &str, visit: &mut dyn FnMut(&GoalId, GoalStatus));
2081}
2082
2083pub trait ClientDispatch {
2098 fn call_raw(
2109 &mut self,
2110 service_entity: &str,
2111 request_cdr: &[u8],
2112 response_buf: &mut [u8],
2113 ) -> NodeResult<usize>;
2114
2115 fn send_goal_raw(&mut self, action_entity: &str, goal_cdr: &[u8]) -> NodeResult<GoalId>;
2120}
2121
2122pub struct TickCtx<'a> {
2129 publishers: &'a dyn PublisherResolver,
2130 actions: &'a mut dyn ActionExecutor,
2131 clients: &'a mut dyn ClientDispatch,
2132 #[cfg(feature = "param-services")]
2136 params: Option<&'a crate::ParameterServer<'a>>,
2137}
2138
2139impl<'a> TickCtx<'a> {
2140 pub fn new(
2142 publishers: &'a dyn PublisherResolver,
2143 actions: &'a mut dyn ActionExecutor,
2144 clients: &'a mut dyn ClientDispatch,
2145 ) -> Self {
2146 Self {
2147 publishers,
2148 actions,
2149 clients,
2150 #[cfg(feature = "param-services")]
2151 params: None,
2152 }
2153 }
2154
2155 #[cfg(feature = "param-services")]
2158 pub fn set_param_server(&mut self, params: Option<&'a crate::ParameterServer<'a>>) {
2159 self.params = params;
2160 }
2161
2162 #[cfg(feature = "param-services")]
2166 pub fn parameter<T: crate::ParameterVariant>(&self, name: &str) -> Option<T> {
2167 self.params
2168 .and_then(|server| server.get(name))
2169 .and_then(T::from_parameter_value)
2170 }
2171
2172 #[doc(hidden)]
2174 pub fn publish_raw(&self, publisher: EntityId<'_>, data: &[u8]) -> NodeResult<()> {
2175 self.publishers.publish_raw(publisher.as_str(), data)
2176 }
2177
2178 #[doc(hidden)]
2180 pub fn publish<M: RosMessage, const N: usize>(
2181 &self,
2182 publisher: EntityId<'_>,
2183 msg: &M,
2184 ) -> NodeResult<()> {
2185 let mut buf = [0u8; N];
2186 let mut writer =
2187 crate::CdrWriter::new_with_header(&mut buf).map_err(|_| NodeDeclError::Runtime)?;
2188 msg.serialize(&mut writer)
2189 .map_err(|_| NodeDeclError::Runtime)?;
2190 let len = writer.position();
2191 self.publish_raw(publisher, &buf[..len])
2192 }
2193
2194 pub fn publish_to_topic<M: RosMessage, const N: usize>(
2199 &self,
2200 topic: &str,
2201 msg: &M,
2202 ) -> NodeResult<()> {
2203 self.publish::<M, N>(EntityId::new(topic), msg)
2204 }
2205
2206 #[doc(hidden)]
2209 pub fn complete_goal<R: RosMessage, const N: usize>(
2210 &mut self,
2211 action: EntityId<'_>,
2212 goal_id: &GoalId,
2213 status: GoalStatus,
2214 result: &R,
2215 ) -> NodeResult<()> {
2216 let mut buf = [0u8; N];
2230 let mut writer = crate::CdrWriter::new(&mut buf);
2231 result
2232 .serialize(&mut writer)
2233 .map_err(|_| NodeDeclError::Runtime)?;
2234 let len = writer.position();
2235 self.actions
2236 .complete_goal_raw(action.as_str(), goal_id, status, &buf[..len])
2237 }
2238
2239 pub fn complete_goal_for_name<R: RosMessage, const N: usize>(
2245 &mut self,
2246 name: &str,
2247 goal_id: &GoalId,
2248 status: GoalStatus,
2249 result: &R,
2250 ) -> NodeResult<()> {
2251 self.complete_goal::<R, N>(EntityId::new(name), goal_id, status, result)
2252 }
2253
2254 #[doc(hidden)]
2260 pub fn for_each_active_goal(
2261 &self,
2262 action: EntityId<'_>,
2263 visit: &mut dyn FnMut(&GoalId, GoalStatus),
2264 ) {
2265 self.actions.for_each_active_goal(action.as_str(), visit);
2266 }
2267
2268 pub fn for_each_active_goal_for_name(
2270 &self,
2271 name: &str,
2272 visit: &mut dyn FnMut(&GoalId, GoalStatus),
2273 ) {
2274 self.for_each_active_goal(EntityId::new(name), visit);
2275 }
2276
2277 #[doc(hidden)]
2279 pub fn publish_feedback<F: RosMessage, const N: usize>(
2280 &mut self,
2281 action: EntityId<'_>,
2282 goal_id: &GoalId,
2283 feedback: &F,
2284 ) -> NodeResult<()> {
2285 let mut buf = [0u8; N];
2290 let mut writer = crate::CdrWriter::new(&mut buf);
2291 feedback
2292 .serialize(&mut writer)
2293 .map_err(|_| NodeDeclError::Runtime)?;
2294 let len = writer.position();
2295 self.actions
2296 .publish_feedback_raw(action.as_str(), goal_id, &buf[..len])
2297 }
2298
2299 pub fn publish_feedback_for_name<F: RosMessage, const N: usize>(
2301 &mut self,
2302 name: &str,
2303 goal_id: &GoalId,
2304 feedback: &F,
2305 ) -> NodeResult<()> {
2306 self.publish_feedback::<F, N>(EntityId::new(name), goal_id, feedback)
2307 }
2308
2309 #[doc(hidden)]
2313 pub fn call_raw(
2314 &mut self,
2315 service: EntityId<'_>,
2316 request_cdr: &[u8],
2317 response_buf: &mut [u8],
2318 ) -> NodeResult<usize> {
2319 self.clients
2320 .call_raw(service.as_str(), request_cdr, response_buf)
2321 }
2322
2323 pub fn call_raw_for_name(
2326 &mut self,
2327 name: &str,
2328 request_cdr: &[u8],
2329 response_buf: &mut [u8],
2330 ) -> NodeResult<usize> {
2331 self.call_raw(EntityId::new(name), request_cdr, response_buf)
2332 }
2333
2334 #[doc(hidden)]
2339 pub fn call<Req: RosMessage, Resp: RosMessage, const REQ_N: usize, const RESP_N: usize>(
2340 &mut self,
2341 service: EntityId<'_>,
2342 request: &Req,
2343 ) -> NodeResult<Resp> {
2344 let mut req_buf = [0u8; REQ_N];
2345 let mut writer =
2346 crate::CdrWriter::new_with_header(&mut req_buf).map_err(|_| NodeDeclError::Runtime)?;
2347 request
2348 .serialize(&mut writer)
2349 .map_err(|_| NodeDeclError::Runtime)?;
2350 let req_len = writer.position();
2351
2352 let mut resp_buf = [0u8; RESP_N];
2353 let resp_len =
2354 self.clients
2355 .call_raw(service.as_str(), &req_buf[..req_len], &mut resp_buf)?;
2356
2357 let mut reader = crate::CdrReader::new_with_header(&resp_buf[..resp_len])
2358 .map_err(|_| NodeDeclError::Runtime)?;
2359 Resp::deserialize(&mut reader).map_err(|_| NodeDeclError::Runtime)
2360 }
2361
2362 pub fn call_for_name<
2365 Req: RosMessage,
2366 Resp: RosMessage,
2367 const REQ_N: usize,
2368 const RESP_N: usize,
2369 >(
2370 &mut self,
2371 name: &str,
2372 request: &Req,
2373 ) -> NodeResult<Resp> {
2374 self.call::<Req, Resp, REQ_N, RESP_N>(EntityId::new(name), request)
2375 }
2376
2377 #[doc(hidden)]
2381 pub fn send_goal_raw(&mut self, action: EntityId<'_>, goal_cdr: &[u8]) -> NodeResult<GoalId> {
2382 self.clients.send_goal_raw(action.as_str(), goal_cdr)
2383 }
2384
2385 pub fn send_goal_raw_for_name(&mut self, name: &str, goal_cdr: &[u8]) -> NodeResult<GoalId> {
2388 self.send_goal_raw(EntityId::new(name), goal_cdr)
2389 }
2390
2391 #[doc(hidden)]
2395 pub fn send_goal<G: RosMessage, const N: usize>(
2396 &mut self,
2397 action: EntityId<'_>,
2398 goal: &G,
2399 ) -> NodeResult<GoalId> {
2400 let mut buf = [0u8; N];
2418 let mut writer = crate::CdrWriter::new(&mut buf);
2419 goal.serialize(&mut writer)
2420 .map_err(|_| NodeDeclError::Runtime)?;
2421 let len = writer.position();
2422 self.clients.send_goal_raw(action.as_str(), &buf[..len])
2423 }
2424
2425 pub fn send_goal_for_name<G: RosMessage, const N: usize>(
2428 &mut self,
2429 name: &str,
2430 goal: &G,
2431 ) -> NodeResult<GoalId> {
2432 self.send_goal::<G, N>(EntityId::new(name), goal)
2433 }
2434}
2435
2436pub trait ExecutableNode: Node {
2437 type State;
2439
2440 fn init() -> Self::State;
2442
2443 fn on_callback(state: &mut Self::State, callback: Callback<'_>, ctx: &mut CallbackCtx<'_>);
2447
2448 fn tick(_state: &mut Self::State, _ctx: &mut TickCtx<'_>) {}
2453}
2454
2455#[macro_export]
2466macro_rules! declarative_component {
2467 ($ty:ty) => {
2468 impl $crate::ExecutableNode for $ty {
2469 type State = ();
2470 fn init() -> Self::State {}
2471 fn on_callback(
2472 _state: &mut Self::State,
2473 _callback: $crate::Callback<'_>,
2474 _ctx: &mut $crate::CallbackCtx<'_>,
2475 ) {
2476 }
2477 }
2478 };
2479}
2480
2481pub fn register_node<C: Node>(runtime: &mut dyn NodeRuntime) -> NodeResult<()> {
2483 let mut context = NodeContext::new(C::NAME, runtime);
2484 C::register(&mut context)
2485}
2486
2487#[cfg(feature = "alloc")]
2494#[doc(hidden)]
2495pub fn __private_node_state_into_raw<C: ExecutableNode>(state: C::State) -> *mut () {
2496 extern crate alloc;
2497 alloc::boxed::Box::into_raw(alloc::boxed::Box::new(state)) as *mut ()
2498}
2499
2500pub fn record_node_metadata<C: Node>(recorder: &mut dyn NodeRuntime) -> NodeResult<()> {
2502 register_node::<C>(recorder)
2503}
2504
2505#[cfg(test)]
2506mod tests {
2507 use super::*;
2508 use crate::{CdrReader, CdrWriter, DeserError, SerError, SourceNameKind};
2509
2510 #[derive(Default)]
2511 struct FakeNodeRuntime {
2512 next: u8,
2513 created: Vec<MetadataString, 4>,
2514 }
2515
2516 impl DeclaredNodeRuntime for FakeNodeRuntime {
2517 type NodeHandle = u8;
2518
2519 fn build_component_node(
2520 &mut self,
2521 _id: NodeId<'_>,
2522 options: NodeOptions<'_>,
2523 ) -> NodeResult<Self::NodeHandle> {
2524 self.created
2525 .push(copy_str(options.name)?)
2526 .map_err(|_| NodeDeclError::Metadata(NodeMetadataError::Capacity))?;
2527 let handle = self.next;
2528 self.next += 1;
2529 Ok(handle)
2530 }
2531 }
2532
2533 #[derive(Debug, Clone, Copy, Default)]
2534 struct TestMsg;
2535
2536 impl crate::Serialize for TestMsg {
2537 fn serialize(&self, _writer: &mut CdrWriter) -> Result<(), SerError> {
2538 Ok(())
2539 }
2540 }
2541
2542 impl crate::Deserialize for TestMsg {
2543 fn deserialize(_reader: &mut CdrReader) -> Result<Self, DeserError> {
2544 Ok(Self)
2545 }
2546 }
2547
2548 impl RosMessage for TestMsg {
2549 const TYPE_NAME: &'static str = "test_msgs::msg::dds_::Test_";
2550 const TYPE_HASH: &'static str = "test_hash";
2551 }
2552
2553 impl nros_serdes::schema::Message for TestMsg {
2559 const TYPE_NAME: &'static str = "test_msgs/msg/Test";
2560 const FIELDS: &'static [nros_serdes::schema::Field] = &[];
2561 }
2562
2563 struct TestService;
2564
2565 impl RosService for TestService {
2566 type Request = TestMsg;
2567 type Reply = TestMsg;
2568
2569 const SERVICE_NAME: &'static str = "test_msgs::srv::dds_::Test_";
2570 const SERVICE_HASH: &'static str = "test_service_hash";
2571 }
2572
2573 struct TestAction;
2574
2575 impl RosAction for TestAction {
2576 type Goal = TestMsg;
2577 type Result = TestMsg;
2578 type Feedback = TestMsg;
2579 type SendGoalRequest = TestMsg;
2580 type SendGoalResponse = TestMsg;
2581 type GetResultRequest = TestMsg;
2582 type GetResultResponse = TestMsg;
2583 type FeedbackMessage = TestMsg;
2584
2585 const ACTION_NAME: &'static str = "test_msgs::action::dds_::Test_";
2586 const ACTION_HASH: &'static str = "test_action_hash";
2587 }
2588
2589 struct TalkerComponent;
2590
2591 impl Node for TalkerComponent {
2592 const NAME: &'static str = "talker_component";
2593
2594 fn register(context: &mut NodeContext<'_>) -> NodeResult<()> {
2595 let mut node =
2596 context.create_node_with_id(NodeId::new("node"), NodeOptions::new("talker"))?;
2597 let _publisher =
2598 node.create_publisher::<TestMsg>(EntityId::new("pub_chatter"), "chatter")?;
2599 let _subscription = node.create_subscription::<TestMsg>(
2600 EntityId::new("sub_cmd"),
2601 CallbackId::new("on_cmd"),
2602 "~/cmd",
2603 )?;
2604 let _timer = node.create_timer(
2605 EntityId::new("timer_tick"),
2606 CallbackId::new("on_tick"),
2607 TimerDuration::from_millis(10),
2608 )?;
2609 let _parameter =
2610 node.declare_parameter(EntityId::new("param_gain"), "gain", ParameterType::Double)?;
2611 node.callback(CallbackId::new("on_tick"))
2612 .publishes(EntityId::new("pub_chatter"))?
2613 .writes(EntityId::new("param_gain"))?;
2614 Ok(())
2615 }
2616 }
2617
2618 #[test]
2619 fn component_records_metadata_without_transport() {
2620 let mut recorder = MetadataRecorder::<2, 8, 4>::new();
2621 record_node_metadata::<TalkerComponent>(&mut recorder).unwrap();
2622
2623 assert_eq!(recorder.nodes().len(), 1);
2624 assert_eq!(recorder.nodes()[0].name.as_str(), "talker");
2625 assert_eq!(recorder.entities().len(), 4);
2626 assert_eq!(recorder.entities()[0].kind, EntityKind::Publisher);
2627 assert_eq!(recorder.entities()[1].source_name.as_str(), "~/cmd");
2628 assert_eq!(
2629 recorder.entities()[1]
2630 .callback_id
2631 .as_ref()
2632 .map(|id| id.as_str()),
2633 Some("on_cmd")
2634 );
2635 assert_eq!(recorder.callback_effects().len(), 2);
2636 }
2637
2638 struct SafetyComponent;
2642 impl Node for SafetyComponent {
2643 const NAME: &'static str = "safety_component";
2644 fn register(context: &mut NodeContext<'_>) -> NodeResult<()> {
2645 let mut node =
2646 context.create_node_with_id(NodeId::new("node"), NodeOptions::new("listener"))?;
2647 let _plain = node.create_subscription_for_callback_name::<TestMsg>("on_plain", "/a")?;
2648 let _safe =
2649 node.create_subscription_for_callback_name_with_safety::<TestMsg>("on_safe", "/b")?;
2650 Ok(())
2651 }
2652 }
2653
2654 #[test]
2655 fn safety_opt_in_records_metadata_flag() {
2656 let mut recorder = MetadataRecorder::<2, 8, 4>::new();
2657 record_node_metadata::<SafetyComponent>(&mut recorder).unwrap();
2658 let ents = recorder.entities();
2659 assert_eq!(ents.len(), 2);
2660 assert_eq!(ents[0].source_name.as_str(), "/a");
2662 assert!(!ents[0].safety, "plain sub must not be flagged");
2663 assert_eq!(ents[1].source_name.as_str(), "/b");
2665 assert!(ents[1].safety, "safety sub must be flagged");
2666 }
2667
2668 struct GroupedComponent;
2669
2670 impl Node for GroupedComponent {
2671 const NAME: &'static str = "grouped_component";
2672
2673 fn register(context: &mut NodeContext<'_>) -> NodeResult<()> {
2674 let mut node =
2675 context.create_node_with_id(NodeId::new("node"), NodeOptions::new("grouped"))?;
2676 let _pub = node.create_publisher::<TestMsg>(EntityId::new("pub_plain"), "plain")?;
2678 node.callback_group("control")?;
2680 let _sub = node.create_subscription::<TestMsg>(
2681 EntityId::new("sub_cmd"),
2682 CallbackId::new("on_cmd"),
2683 "~/cmd",
2684 )?;
2685 let _timer = node.create_timer(
2686 EntityId::new("timer_tick"),
2687 CallbackId::new("on_tick"),
2688 TimerDuration::from_millis(10),
2689 )?;
2690 node.callback_group("telemetry")?;
2692 let _sub2 = node.create_subscription::<TestMsg>(
2693 EntityId::new("sub_diag"),
2694 CallbackId::new("on_diag"),
2695 "~/diag",
2696 )?;
2697 Ok(())
2698 }
2699 }
2700
2701 #[test]
2702 fn sticky_callback_group_stamps_subsequent_entities() {
2703 let mut recorder = MetadataRecorder::<2, 8, 4>::new();
2704 record_node_metadata::<GroupedComponent>(&mut recorder).unwrap();
2705
2706 let group_of = |idx: usize| {
2707 recorder.entities()[idx]
2708 .callback_group
2709 .as_ref()
2710 .map(|g| g.as_str())
2711 };
2712 assert_eq!(group_of(0), None);
2714 assert_eq!(group_of(1), Some("control"));
2716 assert_eq!(group_of(2), Some("control"));
2717 assert_eq!(group_of(3), Some("telemetry"));
2719 }
2720
2721 #[test]
2722 fn runtime_adapter_maps_stable_nodes_to_runtime_handles() {
2723 let mut node_runtime = FakeNodeRuntime::default();
2724 let mut runtime = NodeRuntimeAdapter::<_, 2, 8, 4>::new(&mut node_runtime);
2725
2726 register_node::<TalkerComponent>(&mut runtime).unwrap();
2727
2728 assert_eq!(runtime.nodes().len(), 1);
2729 assert_eq!(runtime.nodes()[0].slot(), NodeSlot::new(0));
2730 assert_eq!(runtime.nodes()[0].stable_id(), "node");
2731 assert_eq!(runtime.nodes()[0].source_default_name(), "talker");
2732 assert_eq!(runtime.node_handle(NodeId::new("node")), Some(0));
2733 assert_eq!(runtime.entities().len(), 4);
2734 assert_eq!(runtime.entities()[0].slot, Some(EntitySlot::new(0)));
2735 assert_eq!(runtime.entities()[0].node_slot, Some(NodeSlot::new(0)));
2736 assert_eq!(
2737 runtime.entities()[1].callback_slot,
2738 Some(CallbackSlot::new(0))
2739 );
2740 assert_eq!(
2741 runtime.entities()[2].callback_slot,
2742 Some(CallbackSlot::new(1))
2743 );
2744 assert_eq!(runtime.callback_effects().len(), 2);
2745 assert_eq!(
2746 runtime.callback_effects()[0].callback_slot,
2747 Some(CallbackSlot::new(1))
2748 );
2749 assert_eq!(
2750 runtime.callback_effects()[0].entity_slot,
2751 Some(EntitySlot::new(0))
2752 );
2753 }
2754
2755 #[test]
2756 fn context_can_synthesize_stable_node_id_from_options_name() {
2757 let mut recorder = MetadataRecorder::<1, 0, 0>::new();
2758 let mut context = NodeContext::new("test", &mut recorder);
2759 let node = context
2760 .create_node(NodeOptions::new("talker").namespace("/demo").domain_id(42))
2761 .unwrap();
2762
2763 assert_eq!(node.id(), NodeId::new("talker"));
2764 let _ = node;
2766 let _ = context;
2767 assert_eq!(recorder.nodes().len(), 1);
2768 assert_eq!(recorder.nodes()[0].id.as_str(), "talker");
2769 assert_eq!(recorder.nodes()[0].name.as_str(), "talker");
2770 assert_eq!(recorder.nodes()[0].namespace.as_str(), "/demo");
2771 assert_eq!(recorder.nodes()[0].domain_id, 42);
2772 }
2773
2774 #[test]
2775 fn synthesized_node_ids_reject_duplicate_names() {
2776 let mut node_runtime = FakeNodeRuntime::default();
2777 let mut runtime = NodeRuntimeAdapter::<_, 2, 0, 0>::new(&mut node_runtime);
2778 {
2779 let mut context = NodeContext::new("test", &mut runtime);
2780 context.create_node(NodeOptions::new("talker")).unwrap();
2781 }
2782 let mut context = NodeContext::new("test", &mut runtime);
2783 let result = context.create_node(NodeOptions::new("talker"));
2784
2785 assert!(matches!(
2786 result,
2787 Err(NodeDeclError::Metadata(NodeMetadataError::DuplicateId))
2788 ));
2789 }
2790
2791 #[test]
2792 fn synthesized_entity_helpers_record_topic_and_callback_ids() {
2793 let mut recorder = MetadataRecorder::<1, 3, 2>::new();
2794 let mut context = NodeContext::new("test", &mut recorder);
2795 let mut node = context.create_node(NodeOptions::new("talker")).unwrap();
2796
2797 let publisher = node
2798 .create_publisher_for_topic::<TestMsg>("/chatter")
2799 .unwrap();
2800 let subscription = node
2801 .create_subscription_for_callback::<TestMsg>(CallbackId::new("on_message"), "/cmd")
2802 .unwrap();
2803 let _timer = node
2804 .create_timer_for_callback(CallbackId::new("on_tick"), TimerDuration::from_millis(10))
2805 .unwrap();
2806
2807 node.callback(CallbackId::new("on_tick"))
2808 .publishes_entity(&publisher)
2809 .unwrap();
2810 node.callback(CallbackId::new("on_message"))
2811 .reads_entity(&subscription)
2812 .unwrap();
2813
2814 assert_eq!(publisher.id(), EntityId::new("/chatter"));
2815 assert_eq!(subscription.id(), EntityId::new("on_message"));
2816 assert_eq!(recorder.entities().len(), 3);
2817
2818 let publisher = &recorder.entities()[0];
2819 assert_eq!(publisher.id.as_str(), "/chatter");
2820 assert_eq!(publisher.kind, EntityKind::Publisher);
2821 assert_eq!(publisher.source_name.as_str(), "/chatter");
2822
2823 let subscription = &recorder.entities()[1];
2824 assert_eq!(subscription.id.as_str(), "on_message");
2825 assert_eq!(subscription.kind, EntityKind::Subscription);
2826 assert_eq!(subscription.source_name.as_str(), "/cmd");
2827 assert_eq!(
2828 subscription.callback_id.as_ref().map(|id| id.as_str()),
2829 Some("on_message")
2830 );
2831
2832 let timer = &recorder.entities()[2];
2833 assert_eq!(timer.id.as_str(), "on_tick");
2834 assert_eq!(timer.kind, EntityKind::Timer);
2835 assert_eq!(
2836 timer.callback_id.as_ref().map(|id| id.as_str()),
2837 Some("on_tick")
2838 );
2839
2840 assert_eq!(recorder.callback_effects().len(), 2);
2841 assert_eq!(
2842 recorder.callback_effects()[0].entity_id.as_str(),
2843 "/chatter"
2844 );
2845 assert_eq!(
2846 recorder.callback_effects()[1].entity_id.as_str(),
2847 "on_message"
2848 );
2849 }
2850
2851 #[test]
2852 fn named_callback_helpers_avoid_manual_callback_ids() {
2853 let mut recorder = MetadataRecorder::<1, 3, 2>::new();
2854 let mut context = NodeContext::new("test", &mut recorder);
2855 let mut node = context.create_node(NodeOptions::new("listener")).unwrap();
2856
2857 let publisher = node
2858 .create_publisher_for_topic::<TestMsg>("/chatter")
2859 .unwrap();
2860 let subscription = node
2861 .create_subscription_for_callback_name::<TestMsg>("on_message", "/chatter")
2862 .unwrap();
2863 let timer = node
2864 .create_timer_for_callback_name("on_tick", TimerDuration::from_millis(10))
2865 .unwrap();
2866
2867 node.callback_for_name("on_message")
2868 .reads_entity(&subscription)
2869 .unwrap();
2870 node.callback_for_name("on_tick")
2871 .publishes_entity(&publisher)
2872 .unwrap();
2873
2874 assert_eq!(subscription.id().as_str(), "on_message");
2875 assert_eq!(timer.id().as_str(), "on_tick");
2876 assert_eq!(
2877 recorder.entities()[1]
2878 .callback_id
2879 .as_ref()
2880 .map(|id| id.as_str()),
2881 Some("on_message")
2882 );
2883 assert_eq!(
2884 recorder.entities()[2]
2885 .callback_id
2886 .as_ref()
2887 .map(|id| id.as_str()),
2888 Some("on_tick")
2889 );
2890 assert_eq!(
2891 recorder.callback_effects()[0].callback_id.as_str(),
2892 "on_message"
2893 );
2894 assert_eq!(
2895 recorder.callback_effects()[0].entity_id.as_str(),
2896 "on_message"
2897 );
2898 assert_eq!(
2899 recorder.callback_effects()[1].callback_id.as_str(),
2900 "on_tick"
2901 );
2902 assert_eq!(
2903 recorder.callback_effects()[1].entity_id.as_str(),
2904 "/chatter"
2905 );
2906 }
2907
2908 #[test]
2909 fn synthesized_entity_ids_reject_collisions() {
2910 let mut recorder = MetadataRecorder::<1, 2, 0>::new();
2911 let mut context = NodeContext::new("test", &mut recorder);
2912 let mut node = context.create_node(NodeOptions::new("talker")).unwrap();
2913
2914 node.create_publisher_for_topic::<TestMsg>("/chatter")
2915 .unwrap();
2916 let result = node.create_publisher_for_topic::<TestMsg>("/chatter");
2917
2918 assert!(matches!(
2919 result,
2920 Err(NodeDeclError::Metadata(NodeMetadataError::DuplicateId))
2921 ));
2922 }
2923
2924 #[test]
2926 fn runtime_adapter_rejects_unknown_entities() {
2927 let mut node_runtime = FakeNodeRuntime::default();
2928 let mut runtime = NodeRuntimeAdapter::<_, 1, 1, 1>::new(&mut node_runtime);
2929 runtime
2930 .create_node(NodeId::new("node"), NodeOptions::new("talker"))
2931 .unwrap();
2932
2933 assert_eq!(
2934 runtime.create_node(NodeId::new("node"), NodeOptions::new("other")),
2935 Err(NodeDeclError::Metadata(NodeMetadataError::DuplicateId))
2936 );
2937 assert_eq!(
2938 runtime.record_callback_effect(
2939 CallbackId::new("cb"),
2940 CallbackEffectKind::Reads,
2941 EntityId::new("missing")
2942 ),
2943 Err(NodeDeclError::Metadata(NodeMetadataError::UnknownEntity))
2944 );
2945 }
2946
2947 #[test]
2948 fn component_rejects_effect_for_unknown_entity() {
2949 let mut recorder = MetadataRecorder::<1, 1, 1>::new();
2950 let mut context = NodeContext::new("test", &mut recorder);
2951 let result = context
2952 .callback(CallbackId::new("cb"))
2953 .reads(EntityId::new("missing"));
2954 assert!(matches!(
2955 result,
2956 Err(NodeDeclError::Metadata(NodeMetadataError::UnknownEntity))
2957 ));
2958 }
2959
2960 #[test]
2961 fn component_missing_export_error_message_is_clear() {
2962 assert_eq!(
2963 NodeDeclError::MissingExport.message(),
2964 MISSING_NODE_EXPORT_ERROR
2965 );
2966 assert_eq!(
2967 NodeDeclError::MissingExport.message(),
2968 "package has no exported nros component"
2969 );
2970 }
2971
2972 struct RobotComponent;
2973
2974 impl Node for RobotComponent {
2975 const NAME: &'static str = "robot_component";
2976
2977 fn register(context: &mut NodeContext<'_>) -> NodeResult<()> {
2978 {
2979 let mut sensors = context.create_node_with_id(
2980 NodeId::new("node_sensors"),
2981 NodeOptions::new("sensors"),
2982 )?;
2983 let _status =
2984 sensors.create_publisher::<TestMsg>(EntityId::new("pub_status"), "~/status")?;
2985 }
2986
2987 let mut control = context
2988 .create_node_with_id(NodeId::new("node_control"), NodeOptions::new("control"))?;
2989 let _cmd = control.create_subscription::<TestMsg>(
2990 EntityId::new("sub_cmd"),
2991 CallbackId::new("cb_cmd"),
2992 "~/cmd",
2993 )?;
2994 let _reset = control.create_service_server::<TestService>(
2995 EntityId::new("srv_reset"),
2996 CallbackId::new("cb_reset"),
2997 "reset",
2998 )?;
2999 let _navigate = control.create_action_server_with_callbacks::<TestAction>(
3000 EntityId::new("act_navigate"),
3001 CallbackId::new("cb_nav_goal"),
3002 CallbackId::new("cb_nav_cancel"),
3003 CallbackId::new("cb_nav_accepted"),
3004 "~/navigate",
3005 )?;
3006 let _gain = control.declare_parameter_with_default(
3007 EntityId::new("param_gain"),
3008 "gain",
3009 ParameterDefault::Double(copy_str("1.5")?),
3010 )?;
3011
3012 control
3013 .callback(CallbackId::new("cb_cmd"))
3014 .publishes(EntityId::new("pub_status"))?
3015 .reads(EntityId::new("param_gain"))?;
3016 control
3017 .callback(CallbackId::new("cb_nav_accepted"))
3018 .writes(EntityId::new("param_gain"))?;
3019
3020 Ok(())
3021 }
3022 }
3023
3024 #[test]
3026 fn component_api_records_multi_node_services() {
3027 let mut recorder = MetadataRecorder::<4, 12, 4>::new();
3028 record_node_metadata::<RobotComponent>(&mut recorder).unwrap();
3029
3030 assert_eq!(recorder.nodes().len(), 2);
3031 assert_eq!(recorder.nodes()[0].id.as_str(), "node_sensors");
3032 assert_eq!(recorder.nodes()[1].id.as_str(), "node_control");
3033
3034 let status = recorder
3035 .entities()
3036 .iter()
3037 .find(|entity| entity.id.as_str() == "pub_status")
3038 .unwrap();
3039 assert_eq!(status.kind, EntityKind::Publisher);
3040 assert_eq!(status.source_name.as_str(), "~/status");
3041 assert_eq!(status.source_name_kind, SourceNameKind::Private);
3042
3043 let reset = recorder
3044 .entities()
3045 .iter()
3046 .find(|entity| entity.id.as_str() == "srv_reset")
3047 .unwrap();
3048 assert_eq!(reset.kind, EntityKind::ServiceServer);
3049 assert_eq!(
3050 reset.callback_id.as_ref().map(|id| id.as_str()),
3051 Some("cb_reset")
3052 );
3053
3054 let navigate = recorder
3055 .entities()
3056 .iter()
3057 .find(|entity| entity.id.as_str() == "act_navigate")
3058 .unwrap();
3059 assert_eq!(navigate.kind, EntityKind::ActionServer);
3060 assert_eq!(
3061 navigate.callback_id.as_ref().map(|id| id.as_str()),
3062 Some("cb_nav_goal")
3063 );
3064 assert_eq!(
3065 navigate
3066 .action_cancel_callback_id
3067 .as_ref()
3068 .map(|id| id.as_str()),
3069 Some("cb_nav_cancel")
3070 );
3071 assert_eq!(
3072 navigate
3073 .action_accepted_callback_id
3074 .as_ref()
3075 .map(|id| id.as_str()),
3076 Some("cb_nav_accepted")
3077 );
3078
3079 let gain = recorder
3080 .entities()
3081 .iter()
3082 .find(|entity| entity.id.as_str() == "param_gain")
3083 .unwrap();
3084 assert_eq!(gain.kind, EntityKind::Parameter);
3085 assert!(matches!(
3086 gain.parameter_default.as_ref(),
3087 Some(ParameterDefault::Double(value)) if value.as_str() == "1.5"
3088 ));
3089
3090 assert_eq!(recorder.callback_effects().len(), 3);
3091 assert!(recorder.callback_effects().iter().any(|effect| {
3092 effect.callback_id.as_str() == "cb_cmd"
3093 && effect.kind == CallbackEffectKind::Publishes
3094 && effect.entity_id.as_str() == "pub_status"
3095 }));
3096 assert!(recorder.callback_effects().iter().any(|effect| {
3097 effect.callback_id.as_str() == "cb_nav_accepted"
3098 && effect.kind == CallbackEffectKind::Writes
3099 && effect.entity_id.as_str() == "param_gain"
3100 }));
3101 }
3102
3103 #[cfg(feature = "std")]
3104 #[test]
3105 fn component_api_json_contains_planner_callback_links() {
3106 let mut recorder = MetadataRecorder::<4, 12, 4>::new();
3107 record_node_metadata::<RobotComponent>(&mut recorder).unwrap();
3108
3109 let json = recorder
3110 .to_source_metadata_json(&crate::SourceMetadataExport::new(
3111 "demo_robot",
3112 RobotComponent::NAME,
3113 ))
3114 .unwrap();
3115
3116 assert!(json.contains("\"callbacks\":["));
3117 assert!(json.contains("\"id\":\"cb_cmd\",\"declaration_slot\":0"));
3118 assert!(json.contains("\"kind\":\"subscription\""));
3119 assert!(json.contains("\"id\":\"cb_reset\",\"declaration_slot\":1"));
3120 assert!(json.contains("\"kind\":\"service\""));
3121 assert!(json.contains("\"id\":\"cb_nav_goal\",\"declaration_slot\":2"));
3122 assert!(json.contains("\"kind\":\"action_goal\""));
3123 assert!(json.contains("\"id\":\"cb_nav_cancel\",\"declaration_slot\":3"));
3124 assert!(json.contains("\"kind\":\"action_cancel\""));
3125 assert!(json.contains("\"id\":\"cb_nav_accepted\",\"declaration_slot\":4"));
3126 assert!(json.contains("\"kind\":\"action_accepted\""));
3127 assert!(json.contains("\"kind\":\"publishes\",\"entity\":\"pub_status\""));
3128 assert!(json.contains("\"kind\":\"reads_parameter\",\"entity\":\"param_gain\""));
3129 assert!(json.contains("\"kind\":\"writes_parameter\",\"entity\":\"param_gain\""));
3130 assert!(json.contains("\"goal_callback\":\"cb_nav_goal\""));
3131 assert!(json.contains("\"cancel_callback\":\"cb_nav_cancel\""));
3132 assert!(json.contains("\"accepted_callback\":\"cb_nav_accepted\""));
3133 }
3134
3135 impl ExecutableNode for TalkerComponent {
3140 type State = u32;
3141
3142 fn init() -> u32 {
3143 0
3144 }
3145
3146 fn on_callback(state: &mut u32, callback: Callback<'_>, ctx: &mut CallbackCtx<'_>) {
3147 if callback.as_str() == "on_tick" {
3148 *state += 1;
3149 let _ = ctx.publish::<TestMsg, 64>(EntityId::new("pub_chatter"), &TestMsg);
3151 }
3152 }
3153 }
3154
3155 #[test]
3156 fn executable_component_callback_publishes_and_mutates_state() {
3157 use core::cell::RefCell;
3158
3159 struct RecordingResolver {
3160 last: RefCell<Option<(MetadataString, usize)>>,
3161 }
3162 impl PublisherResolver for RecordingResolver {
3163 fn publish_raw(&self, entity_id: &str, data: &[u8]) -> NodeResult<()> {
3164 *self.last.borrow_mut() = Some((copy_str(entity_id)?, data.len()));
3165 Ok(())
3166 }
3167 }
3168
3169 let resolver = RecordingResolver {
3170 last: RefCell::new(None),
3171 };
3172 let mut state = TalkerComponent::init();
3173 let mut ctx = CallbackCtx::new(&[], &resolver);
3174
3175 TalkerComponent::on_callback(
3177 &mut state,
3178 Callback::__from_id(CallbackId::new("other")),
3179 &mut ctx,
3180 );
3181 assert_eq!(state, 0);
3182 assert!(resolver.last.borrow().is_none());
3183
3184 TalkerComponent::on_callback(
3186 &mut state,
3187 Callback::__from_id(CallbackId::new("on_tick")),
3188 &mut ctx,
3189 );
3190 assert_eq!(state, 1);
3191 let last = resolver.last.borrow();
3192 let (entity, len) = last.as_ref().expect("a publish was recorded");
3193 assert_eq!(entity.as_str(), "pub_chatter");
3194 assert_eq!(*len, 4);
3196 }
3197
3198 #[test]
3202 fn callback_ctx_reply_sink_roundtrips() {
3203 struct NoopResolver;
3204 impl PublisherResolver for NoopResolver {
3205 fn publish_raw(&self, _entity_id: &str, _data: &[u8]) -> NodeResult<()> {
3206 Ok(())
3207 }
3208 }
3209 let resolver = NoopResolver;
3210 let mut reply_buf = [0u8; 64];
3211 let mut written = 0usize;
3212 {
3213 let mut ctx = CallbackCtx::with_reply(&[], &resolver, &mut reply_buf, &mut written);
3214 ctx.reply::<TestMsg, 64>(&TestMsg).unwrap();
3215 }
3216 assert_eq!(written, 4);
3218
3219 let mut ctx2 = CallbackCtx::new(&[], &resolver);
3221 assert!(ctx2.reply_raw(&[1, 2, 3]).is_err());
3222 }
3223
3224 #[cfg(feature = "param-services")]
3227 #[test]
3228 fn callback_ctx_reads_param() {
3229 struct NoopResolver;
3230 impl PublisherResolver for NoopResolver {
3231 fn publish_raw(&self, _entity_id: &str, _data: &[u8]) -> NodeResult<()> {
3232 Ok(())
3233 }
3234 }
3235 let resolver = NoopResolver;
3236
3237 let ctx_none = CallbackCtx::new(&[], &resolver);
3239 assert_eq!(ctx_none.parameter::<i64>("speed"), None);
3240
3241 let mut storage = crate::ParameterStorage::<4>::new();
3246 let mut server = crate::ParameterServer::new_in(storage.as_table());
3247 assert!(server.declare("speed", crate::ParameterValue::Integer(7)));
3248
3249 let mut ctx = CallbackCtx::new(&[], &resolver);
3250 ctx.set_param_server(Some(&server));
3251 assert_eq!(ctx.parameter::<i64>("speed"), Some(7));
3252 assert_eq!(ctx.parameter::<bool>("speed"), None);
3254 assert_eq!(ctx.parameter::<i64>("missing"), None);
3256 }
3257
3258 #[cfg(feature = "safety-e2e")]
3262 #[test]
3263 fn callback_ctx_integrity_surface() {
3264 struct NoopResolver;
3265 impl PublisherResolver for NoopResolver {
3266 fn publish_raw(&self, _entity_id: &str, _data: &[u8]) -> NodeResult<()> {
3267 Ok(())
3268 }
3269 }
3270 let resolver = NoopResolver;
3271
3272 let ctx = CallbackCtx::new(&[], &resolver);
3274 assert!(ctx.integrity().is_none());
3275
3276 let status = crate::IntegrityStatus {
3278 gap: 2,
3279 duplicate: false,
3280 crc_valid: Some(true),
3281 };
3282 let ctx = CallbackCtx::new_with_integrity(&[], &resolver, &status);
3283 let got = ctx.integrity().expect("safety ctx carries status");
3284 assert_eq!(got.gap, 2);
3285 assert!(!got.duplicate);
3286 assert_eq!(got.crc_valid, Some(true));
3287 }
3288
3289 #[test]
3293 fn callback_ctx_decision_sink() {
3294 struct NoopResolver;
3295 impl PublisherResolver for NoopResolver {
3296 fn publish_raw(&self, _entity_id: &str, _data: &[u8]) -> NodeResult<()> {
3297 Ok(())
3298 }
3299 }
3300 let resolver = NoopResolver;
3301
3302 let mut gr = GoalResponse::Reject;
3303 {
3304 let mut ctx = CallbackCtx::with_goal_decision(&[], &resolver, &mut gr);
3305 ctx.set_goal_response(GoalResponse::AcceptAndExecute)
3306 .unwrap();
3307 assert!(ctx.set_cancel_response(CancelResponse::Accept).is_err());
3309 }
3310 assert!(matches!(gr, GoalResponse::AcceptAndExecute));
3311
3312 let mut cr = CancelResponse::Reject;
3313 {
3314 let mut ctx = CallbackCtx::with_cancel_decision(&[], &resolver, &mut cr);
3315 ctx.set_cancel_response(CancelResponse::Accept).unwrap();
3316 }
3317 assert!(matches!(cr, CancelResponse::Accept));
3318
3319 let mut ctx3 = CallbackCtx::new(&[], &resolver);
3321 assert!(ctx3.set_goal_response(GoalResponse::Reject).is_err());
3322 assert!(ctx3.set_cancel_response(CancelResponse::Accept).is_err());
3323 }
3324
3325 #[test]
3328 fn tick_ctx_publish_and_action_ops() {
3329 use core::cell::Cell;
3330 struct RecPub {
3331 published: Cell<bool>,
3332 }
3333 impl PublisherResolver for RecPub {
3334 fn publish_raw(&self, _entity_id: &str, _data: &[u8]) -> NodeResult<()> {
3335 self.published.set(true);
3336 Ok(())
3337 }
3338 }
3339 struct RecAct {
3340 completed: bool,
3341 fed: bool,
3342 visited: usize,
3343 }
3344 impl ActionExecutor for RecAct {
3345 fn complete_goal_raw(
3346 &mut self,
3347 _action_entity: &str,
3348 _goal_id: &GoalId,
3349 _status: GoalStatus,
3350 _result: &[u8],
3351 ) -> NodeResult<()> {
3352 self.completed = true;
3353 Ok(())
3354 }
3355 fn publish_feedback_raw(
3356 &mut self,
3357 _action_entity: &str,
3358 _goal_id: &GoalId,
3359 _feedback: &[u8],
3360 ) -> NodeResult<()> {
3361 self.fed = true;
3362 Ok(())
3363 }
3364 fn for_each_active_goal(
3365 &self,
3366 _action_entity: &str,
3367 visit: &mut dyn FnMut(&GoalId, GoalStatus),
3368 ) {
3369 visit(&GoalId::zero(), GoalStatus::Executing);
3371 }
3372 }
3373
3374 struct RecClients;
3375 impl ClientDispatch for RecClients {
3376 fn call_raw(
3377 &mut self,
3378 _service: &str,
3379 _req: &[u8],
3380 _resp: &mut [u8],
3381 ) -> NodeResult<usize> {
3382 Err(NodeDeclError::Runtime)
3383 }
3384 fn send_goal_raw(&mut self, _action: &str, _goal: &[u8]) -> NodeResult<GoalId> {
3385 Err(NodeDeclError::Runtime)
3386 }
3387 }
3388
3389 let pubs = RecPub {
3390 published: Cell::new(false),
3391 };
3392 let mut acts = RecAct {
3393 completed: false,
3394 fed: false,
3395 visited: 0,
3396 };
3397 let mut clients = RecClients;
3398 let goal = GoalId::zero();
3399 let mut seen = 0usize;
3400 {
3401 let mut ctx = TickCtx::new(&pubs, &mut acts, &mut clients);
3402 ctx.publish::<TestMsg, 64>(EntityId::new("pub_x"), &TestMsg)
3403 .unwrap();
3404 ctx.for_each_active_goal(EntityId::new("act"), &mut |_id, _status| seen += 1);
3406 ctx.publish_feedback::<TestMsg, 64>(EntityId::new("act"), &goal, &TestMsg)
3407 .unwrap();
3408 ctx.complete_goal::<TestMsg, 64>(
3409 EntityId::new("act"),
3410 &goal,
3411 GoalStatus::Succeeded,
3412 &TestMsg,
3413 )
3414 .unwrap();
3415 }
3416 acts.visited = seen;
3417 assert!(pubs.published.get());
3418 assert!(acts.completed);
3419 assert!(acts.fed);
3420 assert_eq!(acts.visited, 1);
3421 }
3422
3423 #[test]
3427 fn node_dispatch_default_is_inline() {
3428 struct Dummy;
3429 impl Node for Dummy {
3430 const NAME: &'static str = "dummy";
3431 fn register(_: &mut NodeContext<'_>) -> NodeResult<()> {
3432 Ok(())
3433 }
3434 }
3435 assert_eq!(Dummy::DISPATCH, crate::DispatchStrategy::Inline);
3436 }
3437
3438 #[cfg(all(feature = "alloc", feature = "rmw-cffi", feature = "macros"))]
3456 mod dispatch_probe_macro_test {
3457 use super::*;
3461
3462 pub struct DispatchProbe;
3463
3464 impl Node for DispatchProbe {
3465 const NAME: &'static str = "dispatch_probe";
3466 fn register(_: &mut NodeContext<'_>) -> NodeResult<()> {
3468 Ok(())
3469 }
3470 }
3471
3472 impl ExecutableNode for DispatchProbe {
3473 type State = ();
3474 fn init() -> Self::State {}
3475 fn on_callback(
3476 _state: &mut Self::State,
3477 _callback: Callback<'_>,
3478 _ctx: &mut CallbackCtx<'_>,
3479 ) {
3480 }
3481 }
3482
3483 nros_macros::node!(DispatchProbe);
3486 }
3487
3488 #[cfg(all(feature = "alloc", feature = "rmw-cffi", feature = "macros"))]
3492 #[test]
3493 fn node_macro_emits_dispatch_strategy_symbol() {
3494 unsafe extern "C" {
3498 fn __nros_node_nros_dispatch_strategy() -> u8;
3499 }
3500 let strategy = unsafe { __nros_node_nros_dispatch_strategy() };
3501 assert_eq!(strategy, crate::DispatchStrategy::Inline as u8);
3505 assert_eq!(strategy, 0);
3506 }
3507
3508 #[cfg(all(feature = "alloc", feature = "rmw-cffi"))]
3525 #[test]
3526 fn node_macro_emits_on_callback_symbol() {
3527 unsafe extern "C" {
3528 fn __nros_node_nros_on_callback(
3529 state: *mut core::ffi::c_void,
3530 cb_id_ptr: *const u8,
3531 cb_id_len: usize,
3532 ctx: *mut core::ffi::c_void,
3533 );
3534 }
3535 let fn_ptr: unsafe extern "C" fn(
3545 *mut core::ffi::c_void,
3546 *const u8,
3547 usize,
3548 *mut core::ffi::c_void,
3549 ) = __nros_node_nros_on_callback;
3550 core::hint::black_box(fn_ptr);
3551 }
3552
3553 #[test]
3554 fn create_subscription_static_returns_tag_matching_topic() {
3555 let mut recorder = MetadataRecorder::<1, 1, 1>::new();
3556 let mut context = NodeContext::new("test", &mut recorder);
3557 let mut node = context.create_node(NodeOptions::new("listener")).unwrap();
3558 let tag = node
3559 .create_subscription_static::<TestMsg>("/chatter")
3560 .unwrap();
3561
3562 assert_eq!(tag.as_str(), "/chatter");
3563 assert!(tag == CallbackId::new("/chatter"));
3564 assert_eq!(recorder.entities().len(), 1);
3565 let entity = &recorder.entities()[0];
3566 assert_eq!(entity.kind, EntityKind::Subscription);
3567 assert_eq!(entity.source_name.as_str(), "/chatter");
3568 assert_eq!(
3569 entity.callback_id.as_ref().map(|id| id.as_str()),
3570 Some("/chatter")
3571 );
3572 }
3573
3574 #[test]
3575 fn create_service_static_returns_tag() {
3576 let mut recorder = MetadataRecorder::<1, 1, 1>::new();
3577 let mut context = NodeContext::new("test", &mut recorder);
3578 let mut node = context.create_node(NodeOptions::new("server")).unwrap();
3579 let tag = node
3580 .create_service_static::<TestService>("/add_two_ints")
3581 .unwrap();
3582
3583 assert_eq!(tag.as_str(), "/add_two_ints");
3584 assert!(tag == CallbackId::new("/add_two_ints"));
3585 assert_eq!(recorder.entities().len(), 1);
3586 let entity = &recorder.entities()[0];
3587 assert_eq!(entity.kind, EntityKind::ServiceServer);
3588 assert_eq!(entity.source_name.as_str(), "/add_two_ints");
3589 assert_eq!(
3590 entity.callback_id.as_ref().map(|id| id.as_str()),
3591 Some("/add_two_ints")
3592 );
3593 }
3594
3595 #[test]
3596 fn create_service_helpers_use_name_as_entity_and_callback_id() {
3597 let mut recorder = MetadataRecorder::<1, 2, 1>::new();
3598 let mut context = NodeContext::new("test", &mut recorder);
3599 let mut node = context.create_node(NodeOptions::new("services")).unwrap();
3600 let server = node
3601 .create_service_server_for_name::<TestService>("/add_two_ints")
3602 .unwrap();
3603 let client = node
3604 .create_service_client_for_name::<TestService>("/reset")
3605 .unwrap();
3606
3607 assert_eq!(server.id(), EntityId::new("/add_two_ints"));
3608 assert_eq!(client.id(), EntityId::new("/reset"));
3609 assert_eq!(recorder.entities().len(), 2);
3610
3611 let server = &recorder.entities()[0];
3612 assert_eq!(server.kind, EntityKind::ServiceServer);
3613 assert_eq!(server.id.as_str(), "/add_two_ints");
3614 assert_eq!(server.source_name.as_str(), "/add_two_ints");
3615 assert_eq!(
3616 server.callback_id.as_ref().map(|id| id.as_str()),
3617 Some("/add_two_ints")
3618 );
3619
3620 let client = &recorder.entities()[1];
3621 assert_eq!(client.kind, EntityKind::ServiceClient);
3622 assert_eq!(client.id.as_str(), "/reset");
3623 assert_eq!(client.source_name.as_str(), "/reset");
3624 assert!(client.callback_id.is_none());
3625 }
3626
3627 #[test]
3628 fn create_action_static_returns_tag() {
3629 let mut recorder = MetadataRecorder::<1, 1, 1>::new();
3630 let mut context = NodeContext::new("test", &mut recorder);
3631 let mut node = context.create_node(NodeOptions::new("server")).unwrap();
3632 let tag = node
3633 .create_action_static::<TestAction>("/fibonacci")
3634 .unwrap();
3635
3636 assert_eq!(tag.as_str(), "/fibonacci");
3637 assert!(tag == CallbackId::new("/fibonacci"));
3638 assert_eq!(recorder.entities().len(), 1);
3639 let entity = &recorder.entities()[0];
3640 assert_eq!(entity.kind, EntityKind::ActionServer);
3641 assert_eq!(entity.source_name.as_str(), "/fibonacci");
3642 assert_eq!(
3643 entity.callback_id.as_ref().map(|id| id.as_str()),
3644 Some("/fibonacci")
3645 );
3646 assert_eq!(
3647 entity
3648 .action_cancel_callback_id
3649 .as_ref()
3650 .map(|id| id.as_str()),
3651 Some("/fibonacci")
3652 );
3653 assert_eq!(
3654 entity
3655 .action_accepted_callback_id
3656 .as_ref()
3657 .map(|id| id.as_str()),
3658 Some("/fibonacci")
3659 );
3660 }
3661
3662 #[test]
3663 fn create_action_helpers_use_name_as_entity_and_default_callback_id() {
3664 let mut recorder = MetadataRecorder::<1, 2, 3>::new();
3665 let mut context = NodeContext::new("test", &mut recorder);
3666 let mut node = context.create_node(NodeOptions::new("actions")).unwrap();
3667 let server = node
3668 .create_action_server_for_name::<TestAction>("/fibonacci")
3669 .unwrap();
3670 let client = node
3671 .create_action_client_for_name::<TestAction>("/navigate")
3672 .unwrap();
3673
3674 assert_eq!(server.id(), EntityId::new("/fibonacci"));
3675 assert_eq!(client.id(), EntityId::new("/navigate"));
3676 assert_eq!(recorder.entities().len(), 2);
3677
3678 let server = &recorder.entities()[0];
3679 assert_eq!(server.kind, EntityKind::ActionServer);
3680 assert_eq!(server.id.as_str(), "/fibonacci");
3681 assert_eq!(server.source_name.as_str(), "/fibonacci");
3682 assert_eq!(
3683 server.callback_id.as_ref().map(|id| id.as_str()),
3684 Some("/fibonacci")
3685 );
3686 assert_eq!(
3687 server
3688 .action_cancel_callback_id
3689 .as_ref()
3690 .map(|id| id.as_str()),
3691 Some("/fibonacci")
3692 );
3693 assert_eq!(
3694 server
3695 .action_accepted_callback_id
3696 .as_ref()
3697 .map(|id| id.as_str()),
3698 Some("/fibonacci")
3699 );
3700
3701 let client = &recorder.entities()[1];
3702 assert_eq!(client.kind, EntityKind::ActionClient);
3703 assert_eq!(client.id.as_str(), "/navigate");
3704 assert_eq!(client.source_name.as_str(), "/navigate");
3705 assert!(client.callback_id.is_none());
3706 }
3707
3708 #[test]
3714 fn node_identity_injected_wins_over_node_options() {
3715 struct CapturingRuntime {
3718 node_identity: Option<(&'static str, &'static str)>,
3719 resolved_name: MetadataString,
3720 resolved_ns: MetadataString,
3721 }
3722 impl NodeRuntime for CapturingRuntime {
3723 fn create_node(&mut self, _id: NodeId<'_>, options: NodeOptions<'_>) -> NodeResult<()> {
3724 let (name, ns) = match self.node_identity {
3726 Some((n, s)) => (n, s),
3727 None => (options.name, options.namespace),
3728 };
3729 self.resolved_name.clear();
3730 let _ = self.resolved_name.push_str(name);
3731 self.resolved_ns.clear();
3732 let _ = self.resolved_ns.push_str(ns);
3733 Ok(())
3734 }
3735 fn create_entity(&mut self, _m: EntityMetadata) -> NodeResult<()> {
3736 Ok(())
3737 }
3738 fn record_callback_effect(
3739 &mut self,
3740 _id: CallbackId<'_>,
3741 _kind: crate::node_metadata::CallbackEffectKind,
3742 _entity: EntityId<'_>,
3743 ) -> NodeResult<()> {
3744 Ok(())
3745 }
3746 }
3747
3748 let mut rt_a = CapturingRuntime {
3750 node_identity: Some(("launched", "/ns")),
3751 resolved_name: MetadataString::new(),
3752 resolved_ns: MetadataString::new(),
3753 };
3754 {
3755 let mut ctx = NodeContext::new("test_node", &mut rt_a);
3756 ctx.create_node(NodeOptions::new("default").namespace("/d"))
3757 .unwrap();
3758 }
3759 assert_eq!(rt_a.resolved_name.as_str(), "launched");
3760 assert_eq!(rt_a.resolved_ns.as_str(), "/ns");
3761
3762 let mut rt_b = CapturingRuntime {
3764 node_identity: None,
3765 resolved_name: MetadataString::new(),
3766 resolved_ns: MetadataString::new(),
3767 };
3768 {
3769 let mut ctx = NodeContext::new("test_node", &mut rt_b);
3770 ctx.create_node(NodeOptions::new("default").namespace("/d"))
3771 .unwrap();
3772 }
3773 assert_eq!(rt_b.resolved_name.as_str(), "default");
3774 assert_eq!(rt_b.resolved_ns.as_str(), "/d");
3775 }
3776
3777 #[test]
3783 fn entity_names_resolved_through_launch_remaps() {
3784 struct CapturingRuntime {
3785 node_identity: (&'static str, &'static str),
3786 remaps: &'static [(&'static str, &'static str)],
3787 resolved: MetadataString,
3788 }
3789 impl NodeRuntime for CapturingRuntime {
3790 fn create_node(&mut self, _id: NodeId<'_>, _o: NodeOptions<'_>) -> NodeResult<()> {
3791 Ok(())
3792 }
3793 fn create_entity(&mut self, m: EntityMetadata) -> NodeResult<()> {
3794 let name = match m.kind {
3796 EntityKind::Timer | EntityKind::Parameter => m.source_name.clone(),
3797 _ => crate::node_metadata::resolve_name(
3798 m.source_name.as_str(),
3799 self.node_identity.0,
3800 self.node_identity.1,
3801 self.remaps.iter().copied(),
3802 )
3803 .map_err(|_| NodeDeclError::Runtime)?,
3804 };
3805 self.resolved.clear();
3806 let _ = self.resolved.push_str(name.as_str());
3807 Ok(())
3808 }
3809 fn record_callback_effect(
3810 &mut self,
3811 _id: CallbackId<'_>,
3812 _kind: crate::node_metadata::CallbackEffectKind,
3813 _entity: EntityId<'_>,
3814 ) -> NodeResult<()> {
3815 Ok(())
3816 }
3817 }
3818
3819 let mut rt = CapturingRuntime {
3820 node_identity: ("filter", "/sensing"),
3821 remaps: &[("~/input/points", "/points_raw")],
3822 resolved: MetadataString::new(),
3823 };
3824 {
3825 let mut ctx = NodeContext::new("test_node", &mut rt);
3826 let mut node = ctx
3827 .create_node(NodeOptions::new("filter").namespace("/sensing"))
3828 .unwrap();
3829 node.create_subscription_for_callback_name::<TestMsg>("cb", "~/input/points")
3831 .unwrap();
3832 }
3833 assert_eq!(rt.resolved.as_str(), "/points_raw");
3834
3835 {
3837 let mut ctx = NodeContext::new("test_node", &mut rt);
3838 let mut node = ctx
3839 .create_node(NodeOptions::new("filter").namespace("/sensing"))
3840 .unwrap();
3841 node.create_publisher_for_topic::<TestMsg>("status")
3842 .unwrap();
3843 }
3844 assert_eq!(rt.resolved.as_str(), "/sensing/status");
3845
3846 {
3848 let mut ctx = NodeContext::new("test_node", &mut rt);
3849 let mut node = ctx
3850 .create_node(NodeOptions::new("filter").namespace("/sensing"))
3851 .unwrap();
3852 node.declare_parameter_for_name("~/input/points", crate::ParameterType::Bool)
3853 .unwrap();
3854 }
3855 assert_eq!(rt.resolved.as_str(), "~/input/points");
3856 }
3857}