1use core::marker::PhantomData;
4
5use nros_core::{RosAction, RosMessage, RosService};
6use nros_rmw::{ActionInfo, QoSProfile, ServiceInfo, Session as _, TopicInfo, TransportError};
7
8use crate::{
9 rmw_type_registry::{MessageForRmw, register_type},
10 session,
11};
12
13use super::{
14 handles::{
15 ActionClient, ActionClientCallback, ActionServer, EmbeddedPublisher, EmbeddedServiceClient,
16 EmbeddedServiceServer, ServiceClientCallback, Subscription,
17 },
18 types::NodeError,
19};
20
21pub struct NodeHandle<'a> {
27 name: heapless::String<64>,
28 namespace: heapless::String<64>,
29 session: &'a mut session::ConcreteSession,
30 domain_id: u32,
31 qos_overrides: &'static [nros_rmw::QoSOverride],
38 monitors: &'static [crate::executor::monitor::MonitorSpec],
42 age_monitors: &'static [crate::executor::monitor::AgeMonitorSpec],
45 epoch_us_fn: Option<fn() -> u64>,
46}
47
48impl<'a> NodeHandle<'a> {
49 pub(crate) fn new(
51 name: heapless::String<64>,
52 namespace: heapless::String<64>,
53 session: &'a mut session::ConcreteSession,
54 domain_id: u32,
55 ) -> Self {
56 Self {
57 name,
58 namespace,
59 session,
60 domain_id,
61 qos_overrides: &[],
62 monitors: &[],
63 age_monitors: &[],
64 epoch_us_fn: None,
65 }
66 }
67
68 pub fn set_qos_overrides(&mut self, overrides: &'static [nros_rmw::QoSOverride]) {
76 self.qos_overrides = overrides;
77 }
78
79 pub fn set_monitors(&mut self, monitors: &'static [crate::executor::monitor::MonitorSpec]) {
82 self.monitors = monitors;
83 }
84
85 pub fn set_age_monitors(
88 &mut self,
89 table: &'static [crate::executor::monitor::AgeMonitorSpec],
90 epoch_us: Option<fn() -> u64>,
91 ) {
92 self.age_monitors = table;
93 self.epoch_us_fn = epoch_us;
94 }
95
96 #[must_use]
98 pub fn qos_overrides(&self) -> &'static [nros_rmw::QoSOverride] {
99 self.qos_overrides
100 }
101
102 pub fn name(&self) -> &str {
104 &self.name
105 }
106
107 pub fn serialization_format(&self) -> &'static str {
119 nros_rmw::Session::serialization_format(&*self.session)
120 }
121
122 #[must_use]
145 pub fn logger(&self) -> &'static nros_log::Logger {
146 nros_log::get_logger(self.name())
147 }
148
149 pub fn domain_id(&self) -> u32 {
151 self.domain_id
152 }
153
154 pub fn set_domain_id(&mut self, domain_id: u32) {
156 self.domain_id = domain_id;
157 }
158
159 pub fn session_mut(&mut self) -> &mut session::ConcreteSession {
161 self.session
162 }
163
164 fn topic_info<'b>(
184 domain_id: u32,
185 node_name: &'b str,
186 namespace: &'b str,
187 topic_name: &'b str,
188 type_name: &'b str,
189 type_hash: &'b str,
190 ) -> TopicInfo<'b> {
191 TopicInfo::new(topic_name, type_name, type_hash)
192 .with_domain(domain_id)
193 .with_node_name(node_name)
194 .with_namespace(namespace)
195 }
196
197 fn service_info<'b>(
198 domain_id: u32,
199 node_name: &'b str,
200 namespace: &'b str,
201 service_name: &'b str,
202 type_name: &'b str,
203 type_hash: &'b str,
204 ) -> ServiceInfo<'b> {
205 ServiceInfo::new(service_name, type_name, type_hash)
206 .with_domain(domain_id)
207 .with_node_name(node_name)
208 .with_namespace(namespace)
209 }
210
211 fn action_info<'b>(
212 domain_id: u32,
213 action_name: &'b str,
214 type_name: &'b str,
215 type_hash: &'b str,
216 ) -> ActionInfo<'b> {
217 ActionInfo::new(action_name, type_name, type_hash).with_domain(domain_id)
221 }
222
223 pub fn create_publisher<M: MessageForRmw>(
227 &mut self,
228 topic_name: &str,
229 ) -> Result<EmbeddedPublisher<M>, NodeError> {
230 self.create_publisher_with_qos::<M>(topic_name, QoSProfile::default())
231 }
232
233 pub fn create_publisher_with_qos<M: MessageForRmw>(
235 &mut self,
236 topic_name: &str,
237 qos: QoSProfile,
238 ) -> Result<EmbeddedPublisher<M>, NodeError> {
239 crate::format_check::assert_message_format::<M>();
243 register_type::<M>()?;
247 let qos = qos.apply_overrides(
251 topic_name,
252 nros_rmw::QoSOverrideRole::Publisher,
253 self.qos_overrides,
254 );
255 qos.validate_against(nros_rmw::Session::supported_qos_policies(self.session))
258 .map_err(NodeError::Transport)?;
259 let topic = Self::topic_info(
260 self.domain_id,
261 &self.name,
262 &self.namespace,
263 topic_name,
264 <M as RosMessage>::TYPE_NAME,
265 <M as RosMessage>::TYPE_HASH,
266 );
267 let handle = self
268 .session
269 .create_publisher(&topic, qos)
270 .map_err(|_| NodeError::Transport(TransportError::PublisherCreationFailed))?;
271 let monitor = self
274 .monitors
275 .iter()
276 .find(|m| m.topic == topic_name)
277 .map(|m| m.cell);
278 Ok(EmbeddedPublisher {
279 handle,
280 event_regs: crate::executor::handles::empty_event_regs(),
281 monitor,
282 epoch: self.epoch_us_fn,
283 _phantom: PhantomData,
284 })
285 }
286
287 pub fn create_publisher_raw(
293 &mut self,
294 topic_name: &str,
295 type_name: &str,
296 type_hash: &str,
297 ) -> Result<crate::executor::handles::EmbeddedRawPublisher, NodeError> {
298 self.create_publisher_raw_with_qos(topic_name, type_name, type_hash, QoSProfile::default())
299 }
300
301 pub fn create_publisher_raw_with_qos(
303 &mut self,
304 topic_name: &str,
305 type_name: &str,
306 type_hash: &str,
307 qos: QoSProfile,
308 ) -> Result<crate::executor::handles::EmbeddedRawPublisher, NodeError> {
309 let qos = qos.apply_overrides(
311 topic_name,
312 nros_rmw::QoSOverrideRole::Publisher,
313 self.qos_overrides,
314 );
315 qos.validate_against(nros_rmw::Session::supported_qos_policies(self.session))
316 .map_err(NodeError::Transport)?;
317 let topic = Self::topic_info(
318 self.domain_id,
319 &self.name,
320 &self.namespace,
321 topic_name,
322 type_name,
323 type_hash,
324 );
325 let handle = self
326 .session
327 .create_publisher(&topic, qos)
328 .map_err(|_| NodeError::Transport(TransportError::PublisherCreationFailed))?;
329 Ok(crate::executor::handles::EmbeddedRawPublisher {
330 handle,
331 arena: crate::executor::handles::TxArena::new(),
332 event_regs: crate::executor::handles::empty_event_regs(),
333 })
334 }
335
336 pub fn publisher<'t>(&mut self, topic: &'t str) -> PublisherBuilder<'_, 'a, 't> {
342 PublisherBuilder {
343 node: self,
344 topic,
345 qos: QoSProfile::default(),
346 }
347 }
348
349 pub fn create_subscription<M: MessageForRmw>(
353 &mut self,
354 topic_name: &str,
355 ) -> Result<Subscription<M>, NodeError> {
356 self.create_subscription_sized::<M, { crate::config::DEFAULT_RX_BUF_SIZE }>(topic_name)
357 }
358
359 pub fn create_subscription_sized<M: MessageForRmw, const RX_BUF: usize>(
361 &mut self,
362 topic_name: &str,
363 ) -> Result<Subscription<M, RX_BUF>, NodeError> {
364 self.create_subscription_with_qos::<M, RX_BUF>(topic_name, QoSProfile::default())
365 }
366
367 pub fn create_subscription_with_qos<M: MessageForRmw, const RX_BUF: usize>(
369 &mut self,
370 topic_name: &str,
371 qos: QoSProfile,
372 ) -> Result<Subscription<M, RX_BUF>, NodeError> {
373 crate::format_check::assert_message_format::<M>();
377 register_type::<M>()?;
379 let qos = qos.apply_overrides(
381 topic_name,
382 nros_rmw::QoSOverrideRole::Subscription,
383 self.qos_overrides,
384 );
385 qos.validate_against(nros_rmw::Session::supported_qos_policies(self.session))
386 .map_err(NodeError::Transport)?;
387 let topic = Self::topic_info(
388 self.domain_id,
389 &self.name,
390 &self.namespace,
391 topic_name,
392 <M as RosMessage>::TYPE_NAME,
393 <M as RosMessage>::TYPE_HASH,
394 );
395 let handle = self
396 .session
397 .create_subscription(&topic, qos)
398 .map_err(NodeError::Transport)?;
399 let age_mon = match (<M as RosMessage>::STAMP_OFFSET, self.epoch_us_fn) {
402 (Some(_), Some(epoch)) => self
403 .age_monitors
404 .iter()
405 .find(|a| a.topic == topic_name)
406 .map(|a| (a.cell, epoch)),
407 _ => None,
408 };
409 Ok(Subscription {
410 handle,
411 buffer: [0u8; RX_BUF],
412 event_regs: crate::executor::handles::empty_event_regs(),
413 age_mon,
414 _phantom: PhantomData,
415 })
416 }
417
418 pub fn create_subscription_raw(
420 &mut self,
421 topic_name: &str,
422 type_name: &str,
423 type_hash: &str,
424 ) -> Result<crate::executor::handles::RawSubscription, NodeError> {
425 self.create_subscription_raw_sized::<{ crate::config::DEFAULT_RX_BUF_SIZE }>(
426 topic_name, type_name, type_hash,
427 )
428 }
429
430 pub fn create_subscription_raw_sized<const RX_BUF: usize>(
432 &mut self,
433 topic_name: &str,
434 type_name: &str,
435 type_hash: &str,
436 ) -> Result<crate::executor::handles::RawSubscription<RX_BUF>, NodeError> {
437 let qos = QoSProfile::default().apply_overrides(
442 topic_name,
443 nros_rmw::QoSOverrideRole::Subscription,
444 self.qos_overrides,
445 );
446 qos.validate_against(nros_rmw::Session::supported_qos_policies(self.session))
447 .map_err(NodeError::Transport)?;
448 let topic = Self::topic_info(
449 self.domain_id,
450 &self.name,
451 &self.namespace,
452 topic_name,
453 type_name,
454 type_hash,
455 );
456 let handle = self
457 .session
458 .create_subscription(&topic, qos)
459 .map_err(NodeError::Transport)?;
460 Ok(crate::executor::handles::RawSubscription {
461 handle,
462 buffer: [0u8; RX_BUF],
463 event_regs: crate::executor::handles::empty_event_regs(),
464 })
465 }
466
467 pub fn create_service<Svc: RosService>(
471 &mut self,
472 service_name: &str,
473 ) -> Result<EmbeddedServiceServer<Svc>, NodeError>
474 where
475 Svc::Request: MessageForRmw,
476 Svc::Reply: MessageForRmw,
477 {
478 self.create_service_sized::<Svc, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }>(service_name, QoSProfile::services_default())
479 }
480
481 pub fn create_service_with_qos<Svc: RosService>(
484 &mut self,
485 service_name: &str,
486 qos: QoSProfile,
487 ) -> Result<EmbeddedServiceServer<Svc>, NodeError>
488 where
489 Svc::Request: MessageForRmw,
490 Svc::Reply: MessageForRmw,
491 {
492 self.create_service_sized::<Svc, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }>(service_name, qos)
493 }
494
495 pub fn create_service_sized<Svc: RosService, const REQ_BUF: usize, const REPLY_BUF: usize>(
497 &mut self,
498 service_name: &str,
499 qos: QoSProfile,
500 ) -> Result<EmbeddedServiceServer<Svc, REQ_BUF, REPLY_BUF>, NodeError>
501 where
502 Svc::Request: MessageForRmw,
503 Svc::Reply: MessageForRmw,
504 {
505 register_type::<Svc::Request>()?;
508 register_type::<Svc::Reply>()?;
509 qos.validate_against(nros_rmw::Session::supported_qos_policies(self.session))
514 .map_err(NodeError::Transport)?;
515 let info = Self::service_info(
516 self.domain_id,
517 &self.name,
518 &self.namespace,
519 service_name,
520 Svc::SERVICE_NAME,
521 Svc::SERVICE_HASH,
522 );
523 let handle = self
524 .session
525 .create_service(&info, qos)
526 .map_err(NodeError::Transport)?;
527 Ok(EmbeddedServiceServer {
528 handle,
529 req_buffer: [0u8; REQ_BUF],
530 reply_buffer: [0u8; REPLY_BUF],
531 _phantom: PhantomData,
532 })
533 }
534
535 pub fn create_client<Svc: RosService>(
537 &mut self,
538 service_name: &str,
539 ) -> Result<EmbeddedServiceClient<Svc>, NodeError>
540 where
541 Svc::Request: MessageForRmw,
542 Svc::Reply: MessageForRmw,
543 {
544 self.create_client_sized::<Svc, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }>(service_name, QoSProfile::services_default())
545 }
546
547 pub fn create_client_with_qos<Svc: RosService>(
549 &mut self,
550 service_name: &str,
551 qos: QoSProfile,
552 ) -> Result<EmbeddedServiceClient<Svc>, NodeError>
553 where
554 Svc::Request: MessageForRmw,
555 Svc::Reply: MessageForRmw,
556 {
557 self.create_client_sized::<Svc, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }>(service_name, qos)
558 }
559
560 pub fn create_client_sized<Svc: RosService, const REQ_BUF: usize, const REPLY_BUF: usize>(
562 &mut self,
563 service_name: &str,
564 qos: QoSProfile,
565 ) -> Result<EmbeddedServiceClient<Svc, REQ_BUF, REPLY_BUF>, NodeError>
566 where
567 Svc::Request: MessageForRmw,
568 Svc::Reply: MessageForRmw,
569 {
570 register_type::<Svc::Request>()?;
572 register_type::<Svc::Reply>()?;
573 qos.validate_against(nros_rmw::Session::supported_qos_policies(self.session))
576 .map_err(NodeError::Transport)?;
577 let info = Self::service_info(
578 self.domain_id,
579 &self.name,
580 &self.namespace,
581 service_name,
582 Svc::SERVICE_NAME,
583 Svc::SERVICE_HASH,
584 );
585 let handle = self
586 .session
587 .create_client(&info, qos)
588 .map_err(|_| NodeError::Transport(TransportError::ServiceClientCreationFailed))?;
589 Ok(EmbeddedServiceClient {
590 handle,
591 req_buffer: [0u8; REQ_BUF],
592 reply_buffer: [0u8; REPLY_BUF],
593 in_flight: false,
594 _phantom: PhantomData,
595 })
596 }
597
598 pub fn create_service_raw(
603 &mut self,
604 service_name: &str,
605 type_name: &str,
606 type_hash: &str,
607 ) -> Result<crate::executor::handles::RawServiceServer, NodeError> {
608 self.create_service_raw_sized::<
609 { crate::config::DEFAULT_RX_BUF_SIZE },
610 { crate::config::DEFAULT_RX_BUF_SIZE },
611 >(service_name, type_name, type_hash)
612 }
613
614 pub fn create_service_raw_sized<const REQ_BUF: usize, const RESP_BUF: usize>(
616 &mut self,
617 service_name: &str,
618 type_name: &str,
619 type_hash: &str,
620 ) -> Result<crate::executor::handles::RawServiceServer<REQ_BUF, RESP_BUF>, NodeError> {
621 let info = Self::service_info(
622 self.domain_id,
623 &self.name,
624 &self.namespace,
625 service_name,
626 type_name,
627 type_hash,
628 );
629 let handle = self
630 .session
631 .create_service(&info, QoSProfile::services_default())
632 .map_err(NodeError::Transport)?;
633 Ok(crate::executor::handles::RawServiceServer::new(handle))
634 }
635
636 pub fn create_client_raw(
638 &mut self,
639 service_name: &str,
640 type_name: &str,
641 type_hash: &str,
642 ) -> Result<crate::executor::handles::RawServiceClient, NodeError> {
643 self.create_client_raw_sized::<
644 { crate::config::DEFAULT_RX_BUF_SIZE },
645 { crate::config::DEFAULT_RX_BUF_SIZE },
646 >(service_name, type_name, type_hash)
647 }
648
649 pub fn create_client_raw_sized<const REQ_BUF: usize, const REPLY_BUF: usize>(
651 &mut self,
652 service_name: &str,
653 type_name: &str,
654 type_hash: &str,
655 ) -> Result<crate::executor::handles::RawServiceClient<REQ_BUF, REPLY_BUF>, NodeError> {
656 let info = Self::service_info(
657 self.domain_id,
658 &self.name,
659 &self.namespace,
660 service_name,
661 type_name,
662 type_hash,
663 );
664 let handle = self
665 .session
666 .create_client(&info, QoSProfile::services_default())
667 .map_err(|_| NodeError::Transport(TransportError::ServiceClientCreationFailed))?;
668 Ok(crate::executor::handles::RawServiceClient::new(handle))
669 }
670
671 pub fn create_action_server_raw(
681 &mut self,
682 action_name: &str,
683 type_name: &str,
684 type_hash: &str,
685 ) -> Result<
686 super::action_core::ActionServerCore<
687 { crate::config::DEFAULT_RX_BUF_SIZE },
688 { crate::config::DEFAULT_RX_BUF_SIZE },
689 { crate::config::DEFAULT_RX_BUF_SIZE },
690 4,
691 >,
692 NodeError,
693 > {
694 self.create_action_server_raw_sized::<
695 { crate::config::DEFAULT_RX_BUF_SIZE },
696 { crate::config::DEFAULT_RX_BUF_SIZE },
697 { crate::config::DEFAULT_RX_BUF_SIZE },
698 4,
699 >(action_name, type_name, type_hash)
700 }
701
702 pub fn create_action_server_raw_sized<
704 const GOAL_BUF: usize,
705 const RESULT_BUF: usize,
706 const FEEDBACK_BUF: usize,
707 const MAX_GOALS: usize,
708 >(
709 &mut self,
710 action_name: &str,
711 type_name: &str,
712 type_hash: &str,
713 ) -> Result<
714 super::action_core::ActionServerCore<GOAL_BUF, RESULT_BUF, FEEDBACK_BUF, MAX_GOALS>,
715 NodeError,
716 > {
717 let action_info = Self::action_info(self.domain_id, action_name, type_name, type_hash);
718
719 let send_goal_type: heapless::String<256> =
725 super::action_core::action_channel_type(type_name, "SendGoal");
726 let send_goal_keyexpr: heapless::String<256> = action_info.send_goal_key();
727 let send_goal_info = Self::service_info(
728 self.domain_id,
729 &self.name,
730 &self.namespace,
731 &send_goal_keyexpr,
732 &send_goal_type,
733 type_hash,
734 );
735 let send_goal_server = self
736 .session
737 .create_service(&send_goal_info, QoSProfile::services_default())
738 .map_err(|_| NodeError::ActionCreationFailed)?;
739
740 let cancel_goal_keyexpr: heapless::String<256> = action_info.cancel_goal_key();
741 let cancel_goal_info = Self::service_info(
742 self.domain_id,
743 &self.name,
744 &self.namespace,
745 &cancel_goal_keyexpr,
746 "action_msgs::srv::dds_::CancelGoal_",
747 type_hash,
748 );
749 let cancel_goal_server = self
750 .session
751 .create_service(&cancel_goal_info, QoSProfile::services_default())
752 .map_err(|_| NodeError::ActionCreationFailed)?;
753
754 let get_result_type: heapless::String<256> =
755 super::action_core::action_channel_type(type_name, "GetResult");
756 let get_result_keyexpr: heapless::String<256> = action_info.get_result_key();
757 let get_result_info = Self::service_info(
758 self.domain_id,
759 &self.name,
760 &self.namespace,
761 &get_result_keyexpr,
762 &get_result_type,
763 type_hash,
764 );
765 let get_result_server = self
766 .session
767 .create_service(&get_result_info, QoSProfile::services_default())
768 .map_err(|_| NodeError::ActionCreationFailed)?;
769
770 let feedback_type: heapless::String<256> =
771 super::action_core::action_channel_type(type_name, "FeedbackMessage");
772 let feedback_keyexpr: heapless::String<256> = action_info.feedback_key();
773 let feedback_topic = Self::topic_info(
774 self.domain_id,
775 &self.name,
776 &self.namespace,
777 &feedback_keyexpr,
778 &feedback_type,
779 type_hash,
780 );
781 let feedback_publisher = self
782 .session
783 .create_publisher(&feedback_topic, QoSProfile::QOS_PROFILE_DEFAULT)
784 .map_err(|_| NodeError::ActionCreationFailed)?;
785
786 let status_keyexpr: heapless::String<256> = action_info.status_key();
787 let status_topic = Self::topic_info(
788 self.domain_id,
789 &self.name,
790 &self.namespace,
791 &status_keyexpr,
792 "action_msgs::msg::dds_::GoalStatusArray_",
793 type_hash,
794 );
795 let status_publisher = self
796 .session
797 .create_publisher(&status_topic, QoSProfile::QOS_PROFILE_ACTION_STATUS_DEFAULT)
798 .map_err(|_| NodeError::ActionCreationFailed)?;
799
800 Ok(super::action_core::ActionServerCore {
801 send_goal_server,
802 cancel_goal_server,
803 get_result_server,
804 feedback_publisher,
805 status_publisher,
806 active_goals: heapless::Vec::new(),
807 completed_results: heapless::Vec::new(),
808 pending_get_results: heapless::Vec::new(),
809 result_slab: [0u8; RESULT_BUF],
810 result_slab_used: 0,
811 goal_buffer: [0u8; GOAL_BUF],
812 feedback_buffer: [0u8; FEEDBACK_BUF],
813 cancel_buffer: [0u8; 256],
814 })
815 }
816
817 pub fn create_action_client_raw(
821 &mut self,
822 action_name: &str,
823 type_name: &str,
824 type_hash: &str,
825 ) -> Result<
826 super::action_core::ActionClientCore<
827 { crate::config::DEFAULT_RX_BUF_SIZE },
828 { crate::config::DEFAULT_RX_BUF_SIZE },
829 { crate::config::DEFAULT_RX_BUF_SIZE },
830 >,
831 NodeError,
832 > {
833 self.create_action_client_raw_sized::<
834 { crate::config::DEFAULT_RX_BUF_SIZE },
835 { crate::config::DEFAULT_RX_BUF_SIZE },
836 { crate::config::DEFAULT_RX_BUF_SIZE },
837 >(action_name, type_name, type_hash)
838 }
839
840 pub fn create_action_client_raw_sized<
842 const GOAL_BUF: usize,
843 const RESULT_BUF: usize,
844 const FEEDBACK_BUF: usize,
845 >(
846 &mut self,
847 action_name: &str,
848 type_name: &str,
849 type_hash: &str,
850 ) -> Result<super::action_core::ActionClientCore<GOAL_BUF, RESULT_BUF, FEEDBACK_BUF>, NodeError>
851 {
852 let action_info = Self::action_info(self.domain_id, action_name, type_name, type_hash);
853
854 let send_goal_type: heapless::String<256> =
860 super::action_core::action_channel_type(type_name, "SendGoal");
861 let send_goal_keyexpr: heapless::String<256> = action_info.send_goal_key();
862 let send_goal_info = Self::service_info(
863 self.domain_id,
864 &self.name,
865 &self.namespace,
866 &send_goal_keyexpr,
867 &send_goal_type,
868 type_hash,
869 );
870 let send_goal_client = self
871 .session
872 .create_client(&send_goal_info, QoSProfile::services_default())
873 .map_err(|_| NodeError::ActionCreationFailed)?;
874
875 let cancel_goal_keyexpr: heapless::String<256> = action_info.cancel_goal_key();
876 let cancel_goal_info = Self::service_info(
877 self.domain_id,
878 &self.name,
879 &self.namespace,
880 &cancel_goal_keyexpr,
881 "action_msgs::srv::dds_::CancelGoal_",
882 type_hash,
883 );
884 let cancel_goal_client = self
885 .session
886 .create_client(&cancel_goal_info, QoSProfile::services_default())
887 .map_err(|_| NodeError::ActionCreationFailed)?;
888
889 let get_result_type: heapless::String<256> =
890 super::action_core::action_channel_type(type_name, "GetResult");
891 let get_result_keyexpr: heapless::String<256> = action_info.get_result_key();
892 let get_result_info = Self::service_info(
893 self.domain_id,
894 &self.name,
895 &self.namespace,
896 &get_result_keyexpr,
897 &get_result_type,
898 type_hash,
899 );
900 let get_result_client = self
901 .session
902 .create_client(&get_result_info, QoSProfile::services_default())
903 .map_err(|_| NodeError::ActionCreationFailed)?;
904
905 let feedback_type: heapless::String<256> =
906 super::action_core::action_channel_type(type_name, "FeedbackMessage");
907 let feedback_keyexpr: heapless::String<256> = action_info.feedback_key();
908 let feedback_topic = Self::topic_info(
909 self.domain_id,
910 &self.name,
911 &self.namespace,
912 &feedback_keyexpr,
913 &feedback_type,
914 type_hash,
915 );
916 let feedback_subscriber = self
917 .session
918 .create_subscription(&feedback_topic, QoSProfile::BEST_EFFORT)
919 .map_err(|_| NodeError::ActionCreationFailed)?;
920
921 Ok(super::action_core::ActionClientCore::new(
922 send_goal_client,
923 cancel_goal_client,
924 get_result_client,
925 feedback_subscriber,
926 ))
927 }
928
929 pub fn create_action_server<A: RosAction>(
931 &mut self,
932 action_name: &str,
933 ) -> Result<ActionServer<A>, NodeError>
934 where
935 A::Goal: MessageForRmw,
936 A::Result: MessageForRmw,
937 A::Feedback: MessageForRmw,
938 A::SendGoalRequest: MessageForRmw,
939 A::SendGoalResponse: MessageForRmw,
940 A::GetResultRequest: MessageForRmw,
941 A::GetResultResponse: MessageForRmw,
942 A::FeedbackMessage: MessageForRmw,
943 {
944 self.create_action_server_sized::<A, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }, 4>(action_name)
945 }
946
947 pub fn create_action_server_sized<
949 A: RosAction,
950 const GOAL_BUF: usize,
951 const RESULT_BUF: usize,
952 const FEEDBACK_BUF: usize,
953 const MAX_GOALS: usize,
954 >(
955 &mut self,
956 action_name: &str,
957 ) -> Result<ActionServer<A, GOAL_BUF, RESULT_BUF, FEEDBACK_BUF, MAX_GOALS>, NodeError>
958 where
959 A::Goal: MessageForRmw,
960 A::Result: MessageForRmw,
961 A::Feedback: MessageForRmw,
962 A::SendGoalRequest: MessageForRmw,
963 A::SendGoalResponse: MessageForRmw,
964 A::GetResultRequest: MessageForRmw,
965 A::GetResultResponse: MessageForRmw,
966 A::FeedbackMessage: MessageForRmw,
967 {
968 register_type::<A::Goal>()?;
976 register_type::<A::Result>()?;
977 register_type::<A::Feedback>()?;
978 register_type::<A::SendGoalRequest>()?;
979 register_type::<A::SendGoalResponse>()?;
980 register_type::<A::GetResultRequest>()?;
981 register_type::<A::GetResultResponse>()?;
982 register_type::<A::FeedbackMessage>()?;
983 A::register_protocol_types().map_err(|()| NodeError::ActionCreationFailed)?;
994 let action_info =
995 Self::action_info(self.domain_id, action_name, A::ACTION_NAME, A::ACTION_HASH);
996
997 let send_goal_type = super::action_core::action_service_base_type(
1006 <A::SendGoalRequest as RosMessage>::TYPE_NAME,
1007 A::ACTION_NAME,
1008 );
1009 let get_result_type = super::action_core::action_service_base_type(
1010 <A::GetResultRequest as RosMessage>::TYPE_NAME,
1011 A::ACTION_NAME,
1012 );
1013 let feedback_type = <A::FeedbackMessage as RosMessage>::TYPE_NAME;
1014
1015 let send_goal_keyexpr: heapless::String<256> = action_info.send_goal_key();
1016 let send_goal_info = Self::service_info(
1017 self.domain_id,
1018 &self.name,
1019 &self.namespace,
1020 &send_goal_keyexpr,
1021 send_goal_type,
1022 A::SEND_GOAL_SERVICE_HASH,
1023 );
1024 let send_goal_server = self
1025 .session
1026 .create_service(&send_goal_info, QoSProfile::services_default())
1027 .map_err(|_| NodeError::ActionCreationFailed)?;
1028
1029 let cancel_goal_keyexpr: heapless::String<256> = action_info.cancel_goal_key();
1030 let cancel_goal_info = Self::service_info(
1031 self.domain_id,
1032 &self.name,
1033 &self.namespace,
1034 &cancel_goal_keyexpr,
1035 "action_msgs::srv::dds_::CancelGoal_",
1036 A::ACTION_HASH,
1037 );
1038 let cancel_goal_server = self
1039 .session
1040 .create_service(&cancel_goal_info, QoSProfile::services_default())
1041 .map_err(|_| NodeError::ActionCreationFailed)?;
1042
1043 let get_result_keyexpr: heapless::String<256> = action_info.get_result_key();
1044 let get_result_info = Self::service_info(
1045 self.domain_id,
1046 &self.name,
1047 &self.namespace,
1048 &get_result_keyexpr,
1049 get_result_type,
1050 A::GET_RESULT_SERVICE_HASH,
1051 );
1052 let get_result_server = self
1053 .session
1054 .create_service(&get_result_info, QoSProfile::services_default())
1055 .map_err(|_| NodeError::ActionCreationFailed)?;
1056
1057 let feedback_keyexpr: heapless::String<256> = action_info.feedback_key();
1058 let feedback_topic = Self::topic_info(
1059 self.domain_id,
1060 &self.name,
1061 &self.namespace,
1062 &feedback_keyexpr,
1063 feedback_type,
1064 <A::FeedbackMessage as RosMessage>::TYPE_HASH,
1065 );
1066 let feedback_publisher = self
1067 .session
1068 .create_publisher(&feedback_topic, QoSProfile::QOS_PROFILE_DEFAULT)
1069 .map_err(|_| NodeError::ActionCreationFailed)?;
1070
1071 let status_keyexpr: heapless::String<256> = action_info.status_key();
1072 let status_topic = Self::topic_info(
1073 self.domain_id,
1074 &self.name,
1075 &self.namespace,
1076 &status_keyexpr,
1077 "action_msgs::msg::dds_::GoalStatusArray_",
1078 A::ACTION_HASH,
1079 );
1080 let status_publisher = self
1081 .session
1082 .create_publisher(&status_topic, QoSProfile::QOS_PROFILE_ACTION_STATUS_DEFAULT)
1083 .map_err(|_| NodeError::ActionCreationFailed)?;
1084
1085 Ok(ActionServer {
1086 core: super::action_core::ActionServerCore {
1087 send_goal_server,
1088 cancel_goal_server,
1089 get_result_server,
1090 feedback_publisher,
1091 status_publisher,
1092 active_goals: heapless::Vec::new(),
1093 completed_results: heapless::Vec::new(),
1094 pending_get_results: heapless::Vec::new(),
1095 result_slab: [0u8; RESULT_BUF],
1096 result_slab_used: 0,
1097 goal_buffer: [0u8; GOAL_BUF],
1098 feedback_buffer: [0u8; FEEDBACK_BUF],
1099 cancel_buffer: [0u8; 256],
1100 },
1101 typed_goals: heapless::Vec::new(),
1102 completed_goals: heapless::Vec::new(),
1103 })
1104 }
1105
1106 pub fn create_action_client<A: RosAction>(
1108 &mut self,
1109 action_name: &str,
1110 ) -> Result<ActionClient<A>, NodeError>
1111 where
1112 A::Goal: MessageForRmw,
1113 A::Result: MessageForRmw,
1114 A::Feedback: MessageForRmw,
1115 A::SendGoalRequest: MessageForRmw,
1116 A::SendGoalResponse: MessageForRmw,
1117 A::GetResultRequest: MessageForRmw,
1118 A::GetResultResponse: MessageForRmw,
1119 A::FeedbackMessage: MessageForRmw,
1120 {
1121 self.create_action_client_sized::<A, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }, { crate::config::DEFAULT_RX_BUF_SIZE }>(action_name)
1122 }
1123
1124 pub fn create_action_client_sized<
1126 A: RosAction,
1127 const GOAL_BUF: usize,
1128 const RESULT_BUF: usize,
1129 const FEEDBACK_BUF: usize,
1130 >(
1131 &mut self,
1132 action_name: &str,
1133 ) -> Result<ActionClient<A, GOAL_BUF, RESULT_BUF, FEEDBACK_BUF>, NodeError>
1134 where
1135 A::Goal: MessageForRmw,
1136 A::Result: MessageForRmw,
1137 A::Feedback: MessageForRmw,
1138 A::SendGoalRequest: MessageForRmw,
1139 A::SendGoalResponse: MessageForRmw,
1140 A::GetResultRequest: MessageForRmw,
1141 A::GetResultResponse: MessageForRmw,
1142 A::FeedbackMessage: MessageForRmw,
1143 {
1144 register_type::<A::Goal>()?;
1146 register_type::<A::Result>()?;
1147 register_type::<A::Feedback>()?;
1148 register_type::<A::SendGoalRequest>()?;
1149 register_type::<A::SendGoalResponse>()?;
1150 register_type::<A::GetResultRequest>()?;
1151 register_type::<A::GetResultResponse>()?;
1152 register_type::<A::FeedbackMessage>()?;
1153 A::register_protocol_types().map_err(|()| NodeError::ActionCreationFailed)?;
1159 let action_info =
1160 Self::action_info(self.domain_id, action_name, A::ACTION_NAME, A::ACTION_HASH);
1161
1162 let send_goal_type = super::action_core::action_service_base_type(
1170 <A::SendGoalRequest as RosMessage>::TYPE_NAME,
1171 A::ACTION_NAME,
1172 );
1173 let get_result_type = super::action_core::action_service_base_type(
1174 <A::GetResultRequest as RosMessage>::TYPE_NAME,
1175 A::ACTION_NAME,
1176 );
1177 let feedback_type = <A::FeedbackMessage as RosMessage>::TYPE_NAME;
1178
1179 let send_goal_keyexpr: heapless::String<256> = action_info.send_goal_key();
1180 let send_goal_info = Self::service_info(
1181 self.domain_id,
1182 &self.name,
1183 &self.namespace,
1184 &send_goal_keyexpr,
1185 send_goal_type,
1186 A::ACTION_HASH,
1187 );
1188 let send_goal_client = self
1189 .session
1190 .create_client(&send_goal_info, QoSProfile::services_default())
1191 .map_err(|_| NodeError::ActionCreationFailed)?;
1192
1193 let cancel_goal_keyexpr: heapless::String<256> = action_info.cancel_goal_key();
1194 let cancel_goal_info = Self::service_info(
1195 self.domain_id,
1196 &self.name,
1197 &self.namespace,
1198 &cancel_goal_keyexpr,
1199 "action_msgs::srv::dds_::CancelGoal_",
1200 A::ACTION_HASH,
1201 );
1202 let cancel_goal_client = self
1203 .session
1204 .create_client(&cancel_goal_info, QoSProfile::services_default())
1205 .map_err(|_| NodeError::ActionCreationFailed)?;
1206
1207 let get_result_keyexpr: heapless::String<256> = action_info.get_result_key();
1208 let get_result_info = Self::service_info(
1209 self.domain_id,
1210 &self.name,
1211 &self.namespace,
1212 &get_result_keyexpr,
1213 get_result_type,
1214 A::ACTION_HASH,
1215 );
1216 let get_result_client = self
1217 .session
1218 .create_client(&get_result_info, QoSProfile::services_default())
1219 .map_err(|_| NodeError::ActionCreationFailed)?;
1220
1221 let feedback_keyexpr: heapless::String<256> = action_info.feedback_key();
1222 let feedback_topic = Self::topic_info(
1223 self.domain_id,
1224 &self.name,
1225 &self.namespace,
1226 &feedback_keyexpr,
1227 feedback_type,
1228 A::ACTION_HASH,
1229 );
1230 let feedback_subscriber = self
1231 .session
1232 .create_subscription(&feedback_topic, QoSProfile::BEST_EFFORT)
1233 .map_err(|_| NodeError::ActionCreationFailed)?;
1234
1235 Ok(ActionClient {
1236 core: super::action_core::ActionClientCore {
1237 send_goal_client,
1238 cancel_goal_client,
1239 get_result_client,
1240 feedback_subscriber,
1241 goal_buffer: [0u8; GOAL_BUF],
1242 result_buffer: [0u8; RESULT_BUF],
1243 feedback_buffer: [0u8; FEEDBACK_BUF],
1244 goal_counter: 0,
1245 in_flight_send_goal: false,
1246 in_flight_cancel: false,
1247 in_flight_get_result: false,
1248 },
1249 _phantom: PhantomData,
1250 })
1251 }
1252}
1253
1254pub struct PublisherBuilder<'n, 'a, 't> {
1261 node: &'n mut NodeHandle<'a>,
1262 topic: &'t str,
1263 qos: QoSProfile,
1264}
1265
1266impl<'n, 'a, 't> PublisherBuilder<'n, 'a, 't> {
1267 pub fn qos(mut self, qos: QoSProfile) -> Self {
1269 self.qos = qos;
1270 self
1271 }
1272
1273 pub fn tx_express(mut self, express: bool) -> Self {
1277 self.qos.tx_express = express;
1278 self
1279 }
1280
1281 pub fn typed<M: MessageForRmw>(self) -> TypedPublisherBuilder<'n, 'a, 't, M> {
1283 TypedPublisherBuilder {
1284 node: self.node,
1285 topic: self.topic,
1286 qos: self.qos,
1287 _phantom: PhantomData,
1288 }
1289 }
1290
1291 pub fn generic(
1294 self,
1295 type_name: &'t str,
1296 type_hash: &'t str,
1297 ) -> GenericPublisherBuilder<'n, 'a, 't> {
1298 GenericPublisherBuilder {
1299 node: self.node,
1300 topic: self.topic,
1301 type_name,
1302 type_hash,
1303 qos: self.qos,
1304 }
1305 }
1306}
1307
1308pub struct TypedPublisherBuilder<'n, 'a, 't, M> {
1310 node: &'n mut NodeHandle<'a>,
1311 topic: &'t str,
1312 qos: QoSProfile,
1313 _phantom: PhantomData<M>,
1314}
1315
1316impl<'n, 'a, 't, M: MessageForRmw> TypedPublisherBuilder<'n, 'a, 't, M> {
1317 pub fn qos(mut self, qos: QoSProfile) -> Self {
1318 self.qos = qos;
1319 self
1320 }
1321
1322 pub fn tx_express(mut self, express: bool) -> Self {
1324 self.qos.tx_express = express;
1325 self
1326 }
1327
1328 pub fn build(self) -> Result<EmbeddedPublisher<M>, NodeError> {
1329 self.node
1330 .create_publisher_with_qos::<M>(self.topic, self.qos)
1331 }
1332}
1333
1334pub struct GenericPublisherBuilder<'n, 'a, 't> {
1336 node: &'n mut NodeHandle<'a>,
1337 topic: &'t str,
1338 type_name: &'t str,
1339 type_hash: &'t str,
1340 qos: QoSProfile,
1341}
1342
1343impl<'n, 'a, 't> GenericPublisherBuilder<'n, 'a, 't> {
1344 pub fn qos(mut self, qos: QoSProfile) -> Self {
1345 self.qos = qos;
1346 self
1347 }
1348
1349 pub fn tx_express(mut self, express: bool) -> Self {
1351 self.qos.tx_express = express;
1352 self
1353 }
1354
1355 pub fn build(self) -> Result<crate::executor::handles::EmbeddedRawPublisher, NodeError> {
1356 self.node.create_publisher_raw_with_qos(
1357 self.topic,
1358 self.type_name,
1359 self.type_hash,
1360 self.qos,
1361 )
1362 }
1363}
1364
1365pub struct CallbackGroup {
1381 name: heapless::String<32>,
1382}
1383
1384impl CallbackGroup {
1385 pub fn name(&self) -> &str {
1387 &self.name
1388 }
1389}
1390
1391pub struct NodeCtx<'e, 's> {
1398 executor: &'e mut super::spin::Executor<'s>,
1399 node_id: super::node_record::NodeId,
1400}
1401
1402impl<'e, 's> NodeCtx<'e, 's> {
1403 pub(crate) fn new(
1404 executor: &'e mut super::spin::Executor<'s>,
1405 node_id: super::node_record::NodeId,
1406 ) -> Self {
1407 Self { executor, node_id }
1408 }
1409
1410 pub fn subscription<'t>(&mut self, topic: &'t str) -> SubscriptionBuilder<'_, 'e, 't, 's> {
1413 SubscriptionBuilder {
1414 ctx: self,
1415 topic,
1416 qos: QoSProfile::default(),
1417 }
1418 }
1419
1420 pub fn publisher<'t>(&mut self, topic: &'t str) -> CtxPublisherBuilder<'_, 'e, 't, 's> {
1427 CtxPublisherBuilder {
1428 ctx: self,
1429 topic,
1430 qos: QoSProfile::default(),
1431 }
1432 }
1433
1434 pub fn create_publisher<M: MessageForRmw>(
1436 &mut self,
1437 topic: &str,
1438 ) -> Result<EmbeddedPublisher<M>, NodeError> {
1439 self.executor
1440 .create_publisher_on::<M>(self.node_id, topic, QoSProfile::default())
1441 }
1442
1443 pub fn create_generic_publisher(
1445 &mut self,
1446 topic: &str,
1447 type_name: &str,
1448 type_hash: &str,
1449 ) -> Result<crate::executor::handles::EmbeddedRawPublisher, NodeError> {
1450 self.create_generic_publisher_with_qos(topic, type_name, type_hash, QoSProfile::default())
1451 }
1452
1453 pub fn create_generic_publisher_with_qos(
1459 &mut self,
1460 topic: &str,
1461 type_name: &str,
1462 type_hash: &str,
1463 qos: QoSProfile,
1464 ) -> Result<crate::executor::handles::EmbeddedRawPublisher, NodeError> {
1465 self.executor
1466 .create_publisher_raw_on(self.node_id, topic, type_name, type_hash, qos)
1467 }
1468
1469 pub fn create_subscription<M, F>(
1472 &mut self,
1473 topic: &str,
1474 callback: F,
1475 ) -> Result<super::types::HandleId, NodeError>
1476 where
1477 M: MessageForRmw + 'static,
1478 F: FnMut(&M) + 'static,
1479 {
1480 self.executor
1481 .register_subscription_buffered_on::<M, F, { crate::config::DEFAULT_RX_BUF_SIZE }>(
1482 self.node_id,
1483 topic,
1484 QoSProfile::default(),
1485 callback,
1486 None, None, )
1489 }
1490
1491 #[cfg(all(feature = "sim-time", any(has_rmw, test)))]
1513 pub fn install_ros_time_source(&mut self) -> Result<super::types::HandleId, NodeError> {
1514 self.install_ros_time_source_on(crate::time_source::CLOCK_TOPIC)
1515 }
1516
1517 #[cfg(all(feature = "sim-time", any(has_rmw, test)))]
1521 pub fn install_ros_time_source_on(
1522 &mut self,
1523 topic: &str,
1524 ) -> Result<super::types::HandleId, NodeError> {
1525 self.executor.install_ros_time_source(self.node_id, topic)
1529 }
1530
1531 pub fn create_callback_group(&self, name: &str) -> CallbackGroup {
1550 let mut s = heapless::String::<32>::new();
1551 for ch in name.chars() {
1554 if s.push(ch).is_err() {
1555 break;
1556 }
1557 }
1558 CallbackGroup { name: s }
1559 }
1560
1561 pub fn create_timer_in<F>(
1568 &mut self,
1569 group: &CallbackGroup,
1570 period: crate::timer::TimerDuration,
1571 callback: F,
1572 ) -> Result<super::types::HandleId, NodeError>
1573 where
1574 F: FnMut() + 'static,
1575 {
1576 self.executor
1577 .register_timer_on(Some(self.node_id), period, callback, Some(group.name()))
1578 }
1579
1580 pub fn create_subscription_in<M, F>(
1586 &mut self,
1587 group: &CallbackGroup,
1588 topic: &str,
1589 callback: F,
1590 ) -> Result<super::types::HandleId, NodeError>
1591 where
1592 M: MessageForRmw + 'static,
1593 F: FnMut(&M) + 'static,
1594 {
1595 self.executor
1596 .register_subscription_buffered_on::<M, F, { crate::config::DEFAULT_RX_BUF_SIZE }>(
1597 self.node_id,
1598 topic,
1599 QoSProfile::default(),
1600 callback,
1601 Some(group.name()),
1602 None, )
1604 }
1605
1606 pub fn create_publisher_in<M: MessageForRmw>(
1614 &mut self,
1615 _group: &CallbackGroup,
1616 topic: &str,
1617 ) -> Result<EmbeddedPublisher<M>, NodeError> {
1618 self.executor
1620 .create_publisher_on::<M>(self.node_id, topic, QoSProfile::default())
1621 }
1622
1623 pub fn create_client_with_callback<Svc, F>(
1629 &mut self,
1630 service_name: &str,
1631 callback: F,
1632 ) -> Result<ServiceClientCallback<Svc>, NodeError>
1633 where
1634 Svc: RosService + 'static,
1635 Svc::Request: MessageForRmw,
1636 Svc::Reply: MessageForRmw,
1637 F: FnMut(&Svc::Reply) + 'static,
1638 {
1639 self.create_client_with_callback_sized::<
1640 Svc,
1641 F,
1642 { crate::config::DEFAULT_RX_BUF_SIZE },
1643 { crate::config::DEFAULT_RX_BUF_SIZE },
1644 >(service_name, callback)
1645 }
1646
1647 pub fn create_client_with_callback_sized<Svc, F, const REQ_BUF: usize, const REPLY_BUF: usize>(
1649 &mut self,
1650 service_name: &str,
1651 callback: F,
1652 ) -> Result<ServiceClientCallback<Svc, REQ_BUF, REPLY_BUF>, NodeError>
1653 where
1654 Svc: RosService + 'static,
1655 Svc::Request: MessageForRmw,
1656 Svc::Reply: MessageForRmw,
1657 F: FnMut(&Svc::Reply) + 'static,
1658 {
1659 register_type::<Svc::Request>()?;
1660 register_type::<Svc::Reply>()?;
1661 let (_id, hdr) = self
1662 .executor
1663 .register_service_client_callback::<Svc, F, REPLY_BUF>(
1664 Some(self.node_id),
1665 service_name,
1666 Svc::SERVICE_NAME,
1667 Svc::SERVICE_HASH,
1668 QoSProfile::services_default(),
1669 callback,
1670 )?;
1671 Ok(ServiceClientCallback::new(hdr))
1672 }
1673
1674 #[allow(clippy::type_complexity)]
1682 pub fn create_action_client_with_callbacks<A, GRespF, FbF, ResF>(
1683 &mut self,
1684 action_name: &str,
1685 on_goal_response: GRespF,
1686 on_feedback: FbF,
1687 on_result: ResF,
1688 ) -> Result<ActionClientCallback<A>, NodeError>
1689 where
1690 A: RosAction + 'static,
1691 A::Goal: MessageForRmw,
1692 A::Result: MessageForRmw,
1693 A::Feedback: MessageForRmw,
1694 GRespF: FnMut(&nros_core::GoalId, bool) + 'static,
1695 FbF: FnMut(&nros_core::GoalId, &A::Feedback) + 'static,
1696 ResF: FnMut(&nros_core::GoalId, nros_core::GoalStatus, &A::Result) + 'static,
1697 {
1698 self.create_action_client_with_callbacks_sized::<
1699 A,
1700 GRespF,
1701 FbF,
1702 ResF,
1703 { crate::config::DEFAULT_RX_BUF_SIZE },
1704 { crate::config::DEFAULT_RX_BUF_SIZE },
1705 { crate::config::DEFAULT_RX_BUF_SIZE },
1706 >(action_name, on_goal_response, on_feedback, on_result)
1707 }
1708
1709 #[allow(clippy::type_complexity)]
1711 pub fn create_action_client_with_callbacks_sized<
1712 A,
1713 GRespF,
1714 FbF,
1715 ResF,
1716 const GOAL_BUF: usize,
1717 const RESULT_BUF: usize,
1718 const FEEDBACK_BUF: usize,
1719 >(
1720 &mut self,
1721 action_name: &str,
1722 on_goal_response: GRespF,
1723 on_feedback: FbF,
1724 on_result: ResF,
1725 ) -> Result<ActionClientCallback<A, GOAL_BUF, RESULT_BUF, FEEDBACK_BUF>, NodeError>
1726 where
1727 A: RosAction + 'static,
1728 A::Goal: MessageForRmw,
1729 A::Result: MessageForRmw,
1730 A::Feedback: MessageForRmw,
1731 GRespF: FnMut(&nros_core::GoalId, bool) + 'static,
1732 FbF: FnMut(&nros_core::GoalId, &A::Feedback) + 'static,
1733 ResF: FnMut(&nros_core::GoalId, nros_core::GoalStatus, &A::Result) + 'static,
1734 {
1735 register_type::<A::Goal>()?;
1736 register_type::<A::Result>()?;
1737 register_type::<A::Feedback>()?;
1738 let (_id, core) = self
1739 .executor
1740 .register_action_client_callback::<A, GRespF, FbF, ResF, GOAL_BUF, RESULT_BUF, FEEDBACK_BUF>(
1741 Some(self.node_id),
1742 action_name,
1743 A::ACTION_NAME,
1744 A::ACTION_HASH,
1745 8u16,
1748 on_goal_response,
1749 on_feedback,
1750 on_result,
1751 )?;
1752 Ok(ActionClientCallback::new(core))
1753 }
1754
1755 pub fn create_generic_subscription<F>(
1757 &mut self,
1758 topic: &str,
1759 type_name: &str,
1760 type_hash: &str,
1761 callback: F,
1762 ) -> Result<super::types::HandleId, NodeError>
1763 where
1764 F: FnMut(&[u8]) + 'static,
1765 {
1766 self.create_generic_subscription_with_qos(
1767 topic,
1768 type_name,
1769 type_hash,
1770 QoSProfile::default(),
1771 callback,
1772 )
1773 }
1774
1775 pub fn create_generic_subscription_with_qos<F>(
1778 &mut self,
1779 topic: &str,
1780 type_name: &str,
1781 type_hash: &str,
1782 qos: QoSProfile,
1783 callback: F,
1784 ) -> Result<super::types::HandleId, NodeError>
1785 where
1786 F: FnMut(&[u8]) + 'static,
1787 {
1788 self.executor
1789 .register_subscription_buffered_raw_on::<F, { crate::config::DEFAULT_RX_BUF_SIZE }>(
1790 self.node_id,
1791 topic,
1792 type_name,
1793 type_hash,
1794 qos,
1795 callback,
1796 )
1797 }
1798
1799 #[cfg(feature = "safety-e2e")]
1806 pub fn create_generic_subscription_with_integrity<F>(
1807 &mut self,
1808 topic: &str,
1809 type_name: &str,
1810 type_hash: &str,
1811 callback: F,
1812 ) -> Result<super::types::HandleId, NodeError>
1813 where
1814 F: FnMut(&[u8], &nros_rmw::IntegrityStatus) + 'static,
1815 {
1816 self.executor
1817 .register_subscription_buffered_raw_safety_on::<F, { crate::config::DEFAULT_RX_BUF_SIZE }>(
1818 self.node_id,
1819 topic,
1820 type_name,
1821 type_hash,
1822 QoSProfile::default(),
1823 callback,
1824 )
1825 }
1826
1827 #[deprecated(
1835 since = "0.5.0",
1836 note = "renamed to `create_subscription_viewable` (phase-390: RFC-0033 \
1837 `borrowed` mode is now `view`)"
1838 )]
1839 pub fn create_subscription_borrowed<B, F>(
1840 &mut self,
1841 topic: &str,
1842 callback: F,
1843 ) -> Result<super::types::HandleId, NodeError>
1844 where
1845 B: nros_core::ViewableMessage + 'static,
1846 F: for<'a> FnMut(&B::View<'a>) + 'static,
1847 {
1848 self.create_subscription_viewable::<B, F>(topic, callback)
1849 }
1850
1851 pub fn create_subscription_viewable<B, F>(
1871 &mut self,
1872 topic: &str,
1873 callback: F,
1874 ) -> Result<super::types::HandleId, NodeError>
1875 where
1876 B: nros_core::ViewableMessage + 'static,
1877 F: for<'a> FnMut(&B::View<'a>) + 'static,
1878 {
1879 self.executor
1880 .register_subscription_buffered_borrowed_on::<B, F, { crate::config::DEFAULT_RX_BUF_SIZE }>(
1881 self.node_id,
1882 topic,
1883 QoSProfile::default().keep_last(1),
1884 callback,
1885 )
1886 }
1887
1888 pub fn service<'t>(&mut self, name: &'t str) -> CtxServiceBuilder<'_, 'e, 't, 's> {
1892 CtxServiceBuilder {
1893 ctx: self,
1894 name,
1895 qos: QoSProfile::services_default(),
1896 }
1897 }
1898
1899 pub fn create_service<Svc, F>(
1902 &mut self,
1903 name: &str,
1904 callback: F,
1905 ) -> Result<super::types::HandleId, NodeError>
1906 where
1907 Svc: RosService + 'static,
1908 Svc::Request: crate::rmw_type_registry::MessageForRmw,
1909 Svc::Reply: crate::rmw_type_registry::MessageForRmw,
1910 F: FnMut(&Svc::Request) -> Svc::Reply + 'static,
1911 {
1912 self.executor.register_service_sized_on::<
1913 Svc,
1914 F,
1915 { crate::config::DEFAULT_RX_BUF_SIZE },
1916 { crate::config::DEFAULT_RX_BUF_SIZE },
1917 >(self.node_id, name, QoSProfile::services_default(), callback)
1918 }
1919}
1920
1921pub struct CtxServiceBuilder<'c, 'e, 't, 's> {
1923 ctx: &'c mut NodeCtx<'e, 's>,
1924 name: &'t str,
1925 qos: QoSProfile,
1926}
1927
1928impl<'c, 'e, 't, 's> CtxServiceBuilder<'c, 'e, 't, 's> {
1929 pub fn qos(mut self, qos: QoSProfile) -> Self {
1932 self.qos = qos;
1933 self
1934 }
1935
1936 pub fn build<Svc, F>(self, callback: F) -> Result<super::types::HandleId, NodeError>
1937 where
1938 Svc: RosService + 'static,
1939 Svc::Request: crate::rmw_type_registry::MessageForRmw,
1940 Svc::Reply: crate::rmw_type_registry::MessageForRmw,
1941 F: FnMut(&Svc::Request) -> Svc::Reply + 'static,
1942 {
1943 self.ctx.executor.register_service_sized_on::<
1944 Svc,
1945 F,
1946 { crate::config::DEFAULT_RX_BUF_SIZE },
1947 { crate::config::DEFAULT_RX_BUF_SIZE },
1948 >(self.ctx.node_id, self.name, self.qos, callback)
1949 }
1950}
1951
1952pub struct CtxPublisherBuilder<'c, 'e, 't, 's> {
1954 ctx: &'c mut NodeCtx<'e, 's>,
1955 topic: &'t str,
1956 qos: QoSProfile,
1957}
1958
1959impl<'c, 'e, 't, 's> CtxPublisherBuilder<'c, 'e, 't, 's> {
1960 pub fn qos(mut self, qos: QoSProfile) -> Self {
1961 self.qos = qos;
1962 self
1963 }
1964
1965 pub fn typed<M: MessageForRmw>(self) -> CtxTypedPublisherBuilder<'c, 'e, 't, 's, M> {
1967 CtxTypedPublisherBuilder {
1968 ctx: self.ctx,
1969 topic: self.topic,
1970 qos: self.qos,
1971 _phantom: PhantomData,
1972 }
1973 }
1974
1975 pub fn generic(
1977 self,
1978 type_name: &'t str,
1979 type_hash: &'t str,
1980 ) -> CtxGenericPublisherBuilder<'c, 'e, 't, 's> {
1981 CtxGenericPublisherBuilder {
1982 ctx: self.ctx,
1983 topic: self.topic,
1984 type_name,
1985 type_hash,
1986 qos: self.qos,
1987 }
1988 }
1989}
1990
1991pub struct CtxTypedPublisherBuilder<'c, 'e, 't, 's, M> {
1993 ctx: &'c mut NodeCtx<'e, 's>,
1994 topic: &'t str,
1995 qos: QoSProfile,
1996 _phantom: PhantomData<M>,
1997}
1998
1999impl<'c, 'e, 't, 's, M: MessageForRmw> CtxTypedPublisherBuilder<'c, 'e, 't, 's, M> {
2000 pub fn qos(mut self, qos: QoSProfile) -> Self {
2001 self.qos = qos;
2002 self
2003 }
2004
2005 pub fn build(self) -> Result<EmbeddedPublisher<M>, NodeError> {
2006 self.ctx
2007 .executor
2008 .create_publisher_on::<M>(self.ctx.node_id, self.topic, self.qos)
2009 }
2010}
2011
2012pub struct CtxGenericPublisherBuilder<'c, 'e, 't, 's> {
2014 ctx: &'c mut NodeCtx<'e, 's>,
2015 topic: &'t str,
2016 type_name: &'t str,
2017 type_hash: &'t str,
2018 qos: QoSProfile,
2019}
2020
2021impl<'c, 'e, 't, 's> CtxGenericPublisherBuilder<'c, 'e, 't, 's> {
2022 pub fn qos(mut self, qos: QoSProfile) -> Self {
2023 self.qos = qos;
2024 self
2025 }
2026
2027 pub fn build(self) -> Result<crate::executor::handles::EmbeddedRawPublisher, NodeError> {
2028 self.ctx.executor.create_publisher_raw_on(
2029 self.ctx.node_id,
2030 self.topic,
2031 self.type_name,
2032 self.type_hash,
2033 self.qos,
2034 )
2035 }
2036}
2037
2038pub struct SubscriptionBuilder<'c, 'e, 't, 's> {
2040 ctx: &'c mut NodeCtx<'e, 's>,
2041 topic: &'t str,
2042 qos: QoSProfile,
2043}
2044
2045impl<'c, 'e, 't, 's> SubscriptionBuilder<'c, 'e, 't, 's> {
2046 pub fn qos(mut self, qos: QoSProfile) -> Self {
2047 self.qos = qos;
2048 self
2049 }
2050
2051 pub fn typed<M: MessageForRmw + 'static>(self) -> TypedSubscriptionBuilder<'c, 'e, 't, 's, M> {
2053 TypedSubscriptionBuilder {
2054 ctx: self.ctx,
2055 topic: self.topic,
2056 qos: self.qos,
2057 sched: None,
2058 _phantom: PhantomData,
2059 }
2060 }
2061
2062 pub fn generic(
2064 self,
2065 type_name: &'t str,
2066 type_hash: &'t str,
2067 ) -> GenericSubscriptionBuilder<'c, 'e, 't, 's> {
2068 GenericSubscriptionBuilder {
2069 ctx: self.ctx,
2070 topic: self.topic,
2071 type_name,
2072 type_hash,
2073 qos: self.qos,
2074 sched: None,
2075 }
2076 }
2077}
2078
2079pub struct TypedSubscriptionBuilder<
2082 'c,
2083 'e,
2084 't,
2085 's,
2086 M,
2087 const RX: usize = { crate::config::DEFAULT_RX_BUF_SIZE },
2088> {
2089 ctx: &'c mut NodeCtx<'e, 's>,
2090 topic: &'t str,
2091 qos: QoSProfile,
2092 sched: Option<super::sched_context::SchedContextId>,
2093 _phantom: PhantomData<M>,
2094}
2095
2096impl<'c, 'e, 't, 's, M: MessageForRmw + 'static, const RX: usize>
2097 TypedSubscriptionBuilder<'c, 'e, 't, 's, M, RX>
2098{
2099 pub fn qos(mut self, qos: QoSProfile) -> Self {
2100 self.qos = qos;
2101 self
2102 }
2103
2104 pub fn sched_context(mut self, sc: super::sched_context::SchedContextId) -> Self {
2106 self.sched = Some(sc);
2107 self
2108 }
2109
2110 pub fn rx_buffer<const N: usize>(self) -> TypedSubscriptionBuilder<'c, 'e, 't, 's, M, N> {
2112 TypedSubscriptionBuilder {
2113 ctx: self.ctx,
2114 topic: self.topic,
2115 qos: self.qos,
2116 sched: self.sched,
2117 _phantom: PhantomData,
2118 }
2119 }
2120
2121 pub fn rx_buffer_from_type(self) -> TypedSubBoundBuilder<'c, 'e, 't, 's, M, RX>
2167 where
2168 M: nros_serdes::schema::Message,
2169 {
2170 const {
2174 assert!(
2175 crate::rmw_type_registry::subscription_rx_bytes::<M>(RX).is_some(),
2176 "this message type has NO maximum serialized size, so its receive \
2177 buffer cannot be sized from it. Every message type must carry a \
2178 bound: give the unbounded member one in the `.msg` \
2179 (`string<=64`, `int32[<=8]`) or a `cap` in `nros-codegen.toml`. \
2180 The generated C header for this type names the member that costs \
2181 it the bound (`NROS_UNBOUNDED__<type>__field_<member>`)."
2182 )
2183 }
2184 TypedSubBoundBuilder {
2185 ctx: self.ctx,
2186 topic: self.topic,
2187 qos: self.qos,
2188 sched: self.sched,
2189 _phantom: PhantomData,
2190 }
2191 }
2192
2193 pub fn message_info(self) -> TypedSubInfoBuilder<'c, 'e, 't, 's, M, RX> {
2198 TypedSubInfoBuilder {
2199 ctx: self.ctx,
2200 topic: self.topic,
2201 qos: self.qos,
2202 sched: self.sched,
2203 _phantom: PhantomData,
2204 }
2205 }
2206
2207 #[cfg(feature = "safety-e2e")]
2210 pub fn safety(self) -> TypedSubSafetyBuilder<'c, 'e, 't, 's, M, RX> {
2211 TypedSubSafetyBuilder {
2212 ctx: self.ctx,
2213 topic: self.topic,
2214 qos: self.qos,
2215 sched: self.sched,
2216 _phantom: PhantomData,
2217 }
2218 }
2219
2220 pub fn build<F: FnMut(&M) + 'static>(
2221 self,
2222 callback: F,
2223 ) -> Result<super::types::HandleId, NodeError> {
2224 let handle = self
2225 .ctx
2226 .executor
2227 .register_subscription_buffered_on::<M, F, RX>(
2228 self.ctx.node_id,
2229 self.topic,
2230 self.qos,
2231 callback,
2232 None, None, )?;
2235 if let Some(sc) = self.sched {
2236 self.ctx.executor.bind_handle_to_sched_context(handle, sc)?;
2237 }
2238 Ok(handle)
2239 }
2240}
2241
2242pub struct TypedSubBoundBuilder<
2250 'c,
2251 'e,
2252 't,
2253 's,
2254 M,
2255 const RX: usize = { crate::config::DEFAULT_RX_BUF_SIZE },
2256> {
2257 ctx: &'c mut NodeCtx<'e, 's>,
2258 topic: &'t str,
2259 qos: QoSProfile,
2260 sched: Option<super::sched_context::SchedContextId>,
2261 _phantom: PhantomData<M>,
2262}
2263
2264impl<'c, 'e, 't, 's, M: MessageForRmw + nros_serdes::schema::Message + 'static, const RX: usize>
2265 TypedSubBoundBuilder<'c, 'e, 't, 's, M, RX>
2266{
2267 pub fn qos(mut self, qos: QoSProfile) -> Self {
2268 self.qos = qos;
2269 self
2270 }
2271
2272 pub fn sched_context(mut self, sc: super::sched_context::SchedContextId) -> Self {
2274 self.sched = Some(sc);
2275 self
2276 }
2277
2278 pub fn build<F: FnMut(&M) + 'static>(
2279 self,
2280 callback: F,
2281 ) -> Result<super::types::HandleId, NodeError> {
2282 let rx_bytes = crate::rmw_type_registry::subscription_rx_bytes::<M>(RX)
2285 .expect("rx_buffer_from_type asserts the bound exists at build time");
2286 let handle = self
2287 .ctx
2288 .executor
2289 .register_subscription_buffered_on::<M, F, RX>(
2290 self.ctx.node_id,
2291 self.topic,
2292 self.qos,
2293 callback,
2294 None,
2295 Some(rx_bytes),
2296 )?;
2297 if let Some(sc) = self.sched {
2298 self.ctx.executor.bind_handle_to_sched_context(handle, sc)?;
2299 }
2300 Ok(handle)
2301 }
2302}
2303
2304pub struct TypedSubInfoBuilder<
2307 'c,
2308 'e,
2309 't,
2310 's,
2311 M,
2312 const RX: usize = { crate::config::DEFAULT_RX_BUF_SIZE },
2313> {
2314 ctx: &'c mut NodeCtx<'e, 's>,
2315 topic: &'t str,
2316 qos: QoSProfile,
2317 sched: Option<super::sched_context::SchedContextId>,
2318 _phantom: PhantomData<M>,
2319}
2320
2321impl<'c, 'e, 't, 's, M: MessageForRmw + 'static, const RX: usize>
2322 TypedSubInfoBuilder<'c, 'e, 't, 's, M, RX>
2323{
2324 pub fn qos(mut self, qos: QoSProfile) -> Self {
2325 self.qos = qos;
2326 self
2327 }
2328
2329 pub fn sched_context(mut self, sc: super::sched_context::SchedContextId) -> Self {
2330 self.sched = Some(sc);
2331 self
2332 }
2333
2334 pub fn rx_buffer<const N: usize>(self) -> TypedSubInfoBuilder<'c, 'e, 't, 's, M, N> {
2335 TypedSubInfoBuilder {
2336 ctx: self.ctx,
2337 topic: self.topic,
2338 qos: self.qos,
2339 sched: self.sched,
2340 _phantom: PhantomData,
2341 }
2342 }
2343
2344 pub fn build<F: FnMut(&M, Option<&nros_core::MessageInfo>) + 'static>(
2345 self,
2346 callback: F,
2347 ) -> Result<super::types::HandleId, NodeError> {
2348 let handle = self
2349 .ctx
2350 .executor
2351 .register_subscription_with_info_sized_inner::<M, F, RX>(
2352 Some(self.ctx.node_id),
2353 self.topic,
2354 self.qos,
2355 callback,
2356 )?;
2357 if let Some(sc) = self.sched {
2358 self.ctx.executor.bind_handle_to_sched_context(handle, sc)?;
2359 }
2360 Ok(handle)
2361 }
2362}
2363
2364#[cfg(feature = "safety-e2e")]
2367pub struct TypedSubSafetyBuilder<
2368 'c,
2369 'e,
2370 't,
2371 's,
2372 M,
2373 const RX: usize = { crate::config::DEFAULT_RX_BUF_SIZE },
2374> {
2375 ctx: &'c mut NodeCtx<'e, 's>,
2376 topic: &'t str,
2377 qos: QoSProfile,
2378 sched: Option<super::sched_context::SchedContextId>,
2379 _phantom: PhantomData<M>,
2380}
2381
2382#[cfg(feature = "safety-e2e")]
2383impl<'c, 'e, 't, 's, M: MessageForRmw + 'static, const RX: usize>
2384 TypedSubSafetyBuilder<'c, 'e, 't, 's, M, RX>
2385{
2386 pub fn qos(mut self, qos: QoSProfile) -> Self {
2387 self.qos = qos;
2388 self
2389 }
2390
2391 pub fn sched_context(mut self, sc: super::sched_context::SchedContextId) -> Self {
2392 self.sched = Some(sc);
2393 self
2394 }
2395
2396 pub fn rx_buffer<const N: usize>(self) -> TypedSubSafetyBuilder<'c, 'e, 't, 's, M, N> {
2397 TypedSubSafetyBuilder {
2398 ctx: self.ctx,
2399 topic: self.topic,
2400 qos: self.qos,
2401 sched: self.sched,
2402 _phantom: PhantomData,
2403 }
2404 }
2405
2406 pub fn build<F: FnMut(&M, &nros_rmw::IntegrityStatus) + 'static>(
2407 self,
2408 callback: F,
2409 ) -> Result<super::types::HandleId, NodeError> {
2410 let handle = self
2411 .ctx
2412 .executor
2413 .register_subscription_with_safety_sized_inner::<M, F, RX>(
2414 Some(self.ctx.node_id),
2415 self.topic,
2416 self.qos,
2417 callback,
2418 )?;
2419 if let Some(sc) = self.sched {
2420 self.ctx.executor.bind_handle_to_sched_context(handle, sc)?;
2421 }
2422 Ok(handle)
2423 }
2424}
2425
2426pub struct GenericSubscriptionBuilder<
2428 'c,
2429 'e,
2430 't,
2431 's,
2432 const RX: usize = { crate::config::DEFAULT_RX_BUF_SIZE },
2433> {
2434 ctx: &'c mut NodeCtx<'e, 's>,
2435 topic: &'t str,
2436 type_name: &'t str,
2437 type_hash: &'t str,
2438 qos: QoSProfile,
2439 sched: Option<super::sched_context::SchedContextId>,
2440}
2441
2442impl<'c, 'e, 't, 's, const RX: usize> GenericSubscriptionBuilder<'c, 'e, 't, 's, RX> {
2443 pub fn qos(mut self, qos: QoSProfile) -> Self {
2444 self.qos = qos;
2445 self
2446 }
2447
2448 pub fn sched_context(mut self, sc: super::sched_context::SchedContextId) -> Self {
2449 self.sched = Some(sc);
2450 self
2451 }
2452
2453 pub fn rx_buffer<const N: usize>(self) -> GenericSubscriptionBuilder<'c, 'e, 't, 's, N> {
2454 GenericSubscriptionBuilder {
2455 ctx: self.ctx,
2456 topic: self.topic,
2457 type_name: self.type_name,
2458 type_hash: self.type_hash,
2459 qos: self.qos,
2460 sched: self.sched,
2461 }
2462 }
2463
2464 pub fn message_info(self) -> GenericSubInfoBuilder<'c, 'e, 't, 's, RX> {
2468 GenericSubInfoBuilder {
2469 ctx: self.ctx,
2470 topic: self.topic,
2471 type_name: self.type_name,
2472 type_hash: self.type_hash,
2473 qos: self.qos,
2474 sched: self.sched,
2475 }
2476 }
2477
2478 pub fn build<F: FnMut(&[u8]) + 'static>(
2479 self,
2480 callback: F,
2481 ) -> Result<super::types::HandleId, NodeError> {
2482 let handle = self
2483 .ctx
2484 .executor
2485 .register_subscription_buffered_raw_on::<F, RX>(
2486 self.ctx.node_id,
2487 self.topic,
2488 self.type_name,
2489 self.type_hash,
2490 self.qos,
2491 callback,
2492 )?;
2493 if let Some(sc) = self.sched {
2494 self.ctx.executor.bind_handle_to_sched_context(handle, sc)?;
2495 }
2496 Ok(handle)
2497 }
2498}
2499
2500pub struct GenericSubInfoBuilder<
2503 'c,
2504 'e,
2505 't,
2506 's,
2507 const RX: usize = { crate::config::DEFAULT_RX_BUF_SIZE },
2508> {
2509 ctx: &'c mut NodeCtx<'e, 's>,
2510 topic: &'t str,
2511 type_name: &'t str,
2512 type_hash: &'t str,
2513 qos: QoSProfile,
2514 sched: Option<super::sched_context::SchedContextId>,
2515}
2516
2517impl<'c, 'e, 't, 's, const RX: usize> GenericSubInfoBuilder<'c, 'e, 't, 's, RX> {
2518 pub fn qos(mut self, qos: QoSProfile) -> Self {
2519 self.qos = qos;
2520 self
2521 }
2522
2523 pub fn sched_context(mut self, sc: super::sched_context::SchedContextId) -> Self {
2524 self.sched = Some(sc);
2525 self
2526 }
2527
2528 pub fn rx_buffer<const N: usize>(self) -> GenericSubInfoBuilder<'c, 'e, 't, 's, N> {
2529 GenericSubInfoBuilder {
2530 ctx: self.ctx,
2531 topic: self.topic,
2532 type_name: self.type_name,
2533 type_hash: self.type_hash,
2534 qos: self.qos,
2535 sched: self.sched,
2536 }
2537 }
2538
2539 pub fn build<F: FnMut(&[u8], &nros_core::RawMessageInfo) + 'static>(
2540 self,
2541 callback: F,
2542 ) -> Result<super::types::HandleId, NodeError> {
2543 let handle = self
2544 .ctx
2545 .executor
2546 .register_subscription_buffered_raw_info_on::<F, RX>(
2547 self.ctx.node_id,
2548 self.topic,
2549 self.type_name,
2550 self.type_hash,
2551 self.qos,
2552 callback,
2553 )?;
2554 if let Some(sc) = self.sched {
2555 self.ctx.executor.bind_handle_to_sched_context(handle, sc)?;
2556 }
2557 Ok(handle)
2558 }
2559}
2560
2561#[cfg(all(test, feature = "std", not(feature = "rmw-cffi")))]
2567mod builder_tests {
2568 use super::*;
2569 use crate::{executor::Executor, mock::MockSession};
2570 use nros_core::{CdrReader, CdrWriter, DeserError, Deserialize, SerError, Serialize};
2571
2572 struct TestMsg;
2573 impl RosMessage for TestMsg {
2574 const TYPE_NAME: &'static str = "test/msg/TestMsg";
2575 const TYPE_HASH: &'static str = "test_hash";
2576 }
2577 impl Serialize for TestMsg {
2578 fn serialize(&self, _w: &mut CdrWriter) -> Result<(), SerError> {
2579 Ok(())
2580 }
2581 }
2582 impl Deserialize for TestMsg {
2583 fn deserialize(_r: &mut CdrReader) -> Result<Self, DeserError> {
2584 Ok(Self)
2585 }
2586 }
2587 impl nros_serdes::schema::Message for TestMsg {
2596 const TYPE_NAME: &'static str = "test/msg/TestMsg";
2597 const FIELDS: &'static [nros_serdes::schema::Field] = &[nros_serdes::schema::Field {
2598 name: "data",
2599 ty: nros_serdes::schema::FieldType::Uint8,
2600 offset: 0,
2601 }];
2602 }
2603
2604 fn s(v: &str) -> heapless::String<64> {
2605 heapless::String::try_from(v).unwrap()
2606 }
2607
2608 #[test]
2609 fn publisher_builder_typed_and_generic() {
2610 let mut session = MockSession::new();
2611 let mut node = NodeHandle::new(s("n"), s("/"), &mut session, 0);
2612
2613 let _typed = node
2615 .publisher("/chatter")
2616 .typed::<TestMsg>()
2617 .qos(QoSProfile::default().keep_last(5))
2618 .build()
2619 .expect("typed publisher builds");
2620
2621 let _generic = node
2623 .publisher("/chatter")
2624 .qos(QoSProfile::default())
2625 .generic("std_msgs/msg/Int32", "hash")
2626 .build()
2627 .expect("generic publisher builds");
2628 }
2629
2630 #[test]
2631 fn subscription_builder_and_convenient() {
2632 let mut exec: Executor = Executor::from_session(MockSession::new());
2633 let id = exec.node_builder("n").build().expect("node");
2634
2635 let _h = exec
2637 .node_mut(id)
2638 .subscription("/chatter")
2639 .typed::<TestMsg>()
2640 .qos(QoSProfile::default().keep_last(5))
2641 .build(|_m: &TestMsg| {})
2642 .expect("typed subscription builds");
2643
2644 let _g = exec
2646 .node_mut(id)
2647 .subscription("/raw")
2648 .generic("std_msgs/msg/Int32", "hash")
2649 .build(|_b: &[u8]| {})
2650 .expect("generic subscription builds");
2651
2652 let sc = exec.default_sched_context_id();
2654 let _s = exec
2655 .node_mut(id)
2656 .subscription("/sized")
2657 .typed::<TestMsg>()
2658 .rx_buffer::<64>()
2659 .sched_context(sc)
2660 .build(|_m: &TestMsg| {})
2661 .expect("sized + sched subscription builds");
2662
2663 let _c = exec
2665 .node_mut(id)
2666 .create_subscription::<TestMsg, _>("/conv", |_m: &TestMsg| {})
2667 .expect("convenient typed subscription builds");
2668 }
2669
2670 #[test]
2671 fn generic_message_info_builder() {
2672 let mut exec: Executor = Executor::from_session(MockSession::new());
2675 let id = exec.node_builder("n").build().expect("node");
2676
2677 let _i = exec
2678 .node_mut(id)
2679 .subscription("/info")
2680 .generic("std_msgs/msg/Int32", "hash")
2681 .message_info()
2682 .rx_buffer::<256>()
2683 .build(|_payload: &[u8], info: &nros_core::RawMessageInfo| {
2684 let _ = info.attachment();
2685 })
2686 .expect("generic + message_info subscription builds");
2687 }
2688
2689 #[test]
2690 fn typed_message_info_builder() {
2691 let mut exec: Executor = Executor::from_session(MockSession::new());
2694 let id = exec.node_builder("n").build().expect("node");
2695 let _h = exec
2696 .node_mut(id)
2697 .subscription("/chatter")
2698 .typed::<TestMsg>()
2699 .qos(QoSProfile::default().keep_last(5))
2700 .message_info()
2701 .build(|_m: &TestMsg, _info: Option<&nros_core::MessageInfo>| {})
2702 .expect("typed + message_info subscription builds");
2703 }
2704
2705 #[cfg(feature = "safety-e2e")]
2706 #[test]
2707 fn typed_safety_builder() {
2708 let mut exec: Executor = Executor::from_session(MockSession::new());
2710 let id = exec.node_builder("n").build().expect("node");
2711 let _h = exec
2712 .node_mut(id)
2713 .subscription("/chatter")
2714 .typed::<TestMsg>()
2715 .safety()
2716 .build(|_m: &TestMsg, _status: &nros_rmw::IntegrityStatus| {})
2717 .expect("typed + safety subscription builds");
2718 }
2719
2720 #[test]
2721 fn generator_emitted_chain_compiles() {
2722 let mut exec: Executor = Executor::from_session(MockSession::new());
2725 let id = exec.node_builder("n").build().expect("node");
2726 let _h = exec
2727 .node_mut(id)
2728 .subscription("/topic")
2729 .generic("std_msgs/msg/Int32", "hash")
2730 .qos(QoSProfile::default().keep_last(1))
2731 .rx_buffer::<1024>()
2732 .build(|_data: &[u8]| {})
2733 .expect("generator-shape subscription builds");
2734 }
2735
2736 #[test]
2737 fn nodectx_publisher_and_bridge_shape() {
2738 let mut exec: Executor = Executor::from_session(MockSession::new());
2742 let id = exec.node_builder("n").build().expect("node");
2743
2744 let _p = exec
2746 .node_mut(id)
2747 .create_publisher::<TestMsg>("/p")
2748 .expect("ctx convenient publisher");
2749 let dest_pub = exec
2750 .node_mut(id)
2751 .publisher("/fwd")
2752 .generic("std_msgs/msg/Int32", "hash")
2753 .build()
2754 .expect("ctx generic publisher builds"); let _s = exec
2758 .node_mut(id)
2759 .subscription("/src")
2760 .generic("std_msgs/msg/Int32", "hash")
2761 .message_info()
2762 .build(move |payload: &[u8], _info: &nros_core::RawMessageInfo| {
2763 let _ = dest_pub.publish_raw(payload);
2764 })
2765 .expect("bridge-shape source subscription builds");
2766 }
2767}