1use async_signals::Signals;
21use async_trait::async_trait;
22use cyclonedds_rs::{
23 DdsParticipant, DdsPublisher, DdsReader, DdsSubscriber, DdsWriter, PublisherBuilder,
24 ReaderBuilder, SampleBuffer, SubscriberBuilder, TopicBuilder, TopicType, WriterBuilder as CDDSWriterBuilder, DdsListener, DdsListenerBuilder,
25};
26use error::MiddlewareError;
27use futures::TryFutureExt;
28use futures_util::StreamExt;
29use qos::{QosDurability, QosHistory, QosReliability};
30use services::get_service_ids;
31use someip::{
32 tasks::ConnectionInfo, Configuration, CreateServerRequestHandler, Proxy, ProxyConstruct,
33 Server, ServerRequestHandler, ServerRequestHandlerEntry, ServiceIdentifier, ServiceVersion,
34};
35use utils::utils::create_home_directory_if_required;
36use std::{
37 ops::Deref,
38 sync::{Arc, RwLock},
39 time::Duration, marker::PhantomData,
40};
41use tokio::runtime::Builder;
42use tracing::{debug, error};
43
44use crate::{
45 cdds::{service_discovery::{service_name_to_topic_name, ServiceInfo, Transport}, cdds::{publish_options_to_cdds_qos, subscribe_options_to_cdds_qos}},
46 config::get_bind_address,
47 qos::{Qos, QosCreate},
48 services::get_config_path,
49};
50pub mod utils;
51pub mod cdds;
52pub mod config;
53pub mod error;
54pub mod qos;
55pub mod services;
56#[cfg(test)]
57mod tests;
58
59const SERVICE_MAPPING_CONFIG_PATH: &str = "/etc/sabaton/services.toml";
62
63pub trait SyncReader<T: TopicType> {
64 fn take_now(&mut self, samples: &mut Samples<T>) -> Result<usize, MiddlewareError>;
65 fn read_now(&mut self, samples: &mut Samples<T>) -> Result<usize, MiddlewareError>;
66}
67
68#[async_trait]
69pub trait AsyncReader<T: TopicType> {
70 async fn take(&mut self, samples: &mut Samples<T>) -> Result<usize, MiddlewareError>;
71 async fn read(&mut self, samples: &mut Samples<T>) -> Result<usize, MiddlewareError>;
72}
73
74#[derive(Clone)]
75pub struct InitContext;
76
77impl InitContext {
78 pub fn new() -> InitContext {
79 InitContext {}
80 }
81}
82
83impl Default for InitContext {
84 fn default() -> Self {
85 Self::new()
86 }
87}
88
89pub struct Writer<T: TopicType> {
90 writer: DdsWriter<T>,
91}
92
93impl<T> Writer<T>
94where
95 T: TopicType,
96{
97
98 pub fn publish(&mut self, msg: Arc<T>) -> Result<(), MiddlewareError> {
105 self.writer
106 .write(msg)
107 .map_err(MiddlewareError::DDSError)
108 }
109
110 pub fn loan(&mut self) -> Result<Loaned<T>, MiddlewareError> {
121 match self.writer.loan() {
122 Ok(l) => Ok(Loaned { inner: l }),
123 Err(e) => match e {
124 cyclonedds_rs::DDSError::NotEnabled => Err(MiddlewareError::SharedMemoryNotEnabled),
125 _ => Err(MiddlewareError::DDSError(e)),
126 },
127 }
128 }
129
130 pub fn return_loan(&mut self, buffer: Loaned<T>) -> Result<(), MiddlewareError> {
140 self.writer
141 .return_loan(buffer.inner)
142 .map_err(MiddlewareError::DDSError)
143 }
144}
145
146pub struct Loaned<T: TopicType> {
147 inner: cyclonedds_rs::dds_writer::Loaned<T>,
148}
149
150impl<T> Loaned<T>
151where
152 T: Sized + TopicType,
153{
154 pub fn as_mut_ptr(&mut self) -> Option<*mut T> {
157 self.inner.as_mut_ptr()
158 }
159
160 pub fn assume_init(self) -> Self {
164 Loaned {
165 inner: self.inner.assume_init(),
166 }
167 }
168}
169
170
171pub struct SampleStorage<T: TopicType> {
172 sample: cyclonedds_rs::serdes::SampleStorage<T>,
173}
174
175impl<T> Deref for SampleStorage<T>
176where
177 T: TopicType,
178{
179 type Target = T;
180
181 fn deref(&self) -> &Self::Target {
182 self.sample.deref()
183 }
184}
185
186pub struct Samples<T: TopicType> {
189 samples: SampleBuffer<T>,
190}
191
192impl<'a, T> Samples<T>
193where
194 T: TopicType,
195{
196 pub fn new(len: usize) -> Self {
203 Self {
204 samples: SampleBuffer::new(len),
205 }
206 }
207
208 pub fn iter(&self) -> impl Iterator<Item = &T> + '_ {
211 self.samples.iter()
212 }
213}
214
215unsafe impl<T> Sync for Reader<T> where T: TopicType {}
216
217pub struct Reader<T: TopicType> {
218 reader: DdsReader<T>,
219}
220
221impl<T> SyncReader<T> for Reader<T>
222where
223 T: TopicType,
224{
225 fn take_now(&mut self, samples: &mut Samples<T>) -> Result<usize, MiddlewareError> {
227 self.reader
228 .take_now(&mut samples.samples)
229 .map_err(|e| e.into())
230 }
231 fn read_now(&mut self, samples: &mut Samples<T>) -> Result<usize, MiddlewareError> {
233 self.reader
234 .read_now(&mut samples.samples)
235 .map_err(|e| e.into())
236 }
237
238}
239
240#[async_trait]
241impl<T> AsyncReader<T> for Reader<T>
242where
243 T: TopicType + std::marker::Send + std::marker::Sync,
244{
245 async fn take(&mut self, samples: &mut Samples<T>) -> Result<usize, MiddlewareError> {
247 let res = self
248 .reader
249 .take(&mut samples.samples)
250 .err_into::<MiddlewareError>()
251 .await;
252
253 res
254 }
255 async fn read(&mut self, samples: &mut Samples<T>) -> Result<usize, MiddlewareError> {
257 let res = self
258 .reader
259 .read(&mut samples.samples)
260 .err_into::<MiddlewareError>()
261 .await;
262
263 res
264 }
265}
266
267pub struct ReaderListener {
286 listener: DdsListener,
287}
288
289#[derive(Clone)]
294pub struct Node {
295 inner: Arc<RwLock<NodeInner>>,
296}
297
298pub struct NodeBuilder {
299 group: String,
300 instance: String,
301 num_workers: usize,
302 single_threaded: bool,
303 shared_memory : bool,
304 pub_sub_log_level : config::LogLevel,
305 rpc_log_level: config::LogLevel,
306 stack_size : Option<usize>,
307 trace_enabled : bool,
308 trace_level : String,
309}
310
311impl Default for NodeBuilder {
312 fn default() -> Self {
313 Self {
314 group: "default".to_owned(),
315 instance: "0".to_owned(),
316 num_workers: 1,
317 single_threaded: true,
318 shared_memory : false,
319 pub_sub_log_level : config::LogLevel::default(),
320 rpc_log_level: config::LogLevel::default(),
321 stack_size : None,
322 trace_enabled : false,
323 trace_level : "trace".to_owned()
324 }
325 }
326}
327
328impl NodeBuilder {
329 pub fn with_group(mut self, group: String) -> Self {
333 self.group = group;
334 self
335 }
336
337 pub fn with_instance(mut self, instance: String) -> Self {
340 self.instance = instance;
341 self
342 }
343
344 #[deprecated]
346 pub fn with_group_and_instance(mut self, group: String, instance: String) -> Self {
347 self.group = group;
348 self.instance = instance;
349
350 self
351 }
352
353 pub fn with_stack_size(mut self, size: usize) -> Self {
355 self.stack_size = Some(size);
356 self
357 }
358
359
360 pub fn multi_threaded(mut self) -> Self {
362 self.single_threaded = false;
363 self
364 }
365
366 pub fn with_num_workers(mut self, num_workers: usize) -> Self {
369 if !self.single_threaded {
370 self.num_workers = num_workers;
371 self
372 } else {
373 panic!("workers not allowed on single threaded runtime");
374 }
375 }
376
377 pub fn with_shared_memory(mut self, enabled : bool) -> Self {
382 self.shared_memory = enabled;
383 self
384 }
385
386 pub fn with_pubsub_log_level(mut self, log_level: config::LogLevel) -> Self {
387 self.pub_sub_log_level = log_level;
388 self
389 }
390
391 pub fn with_rpc_log_level(mut self, log_level: config::LogLevel) -> Self {
392 self.rpc_log_level = log_level;
393 self
394 }
395
396 pub fn with_trace_enabled(mut self, enabled : bool, level : &str) -> Self {
397 self.trace_enabled = enabled;
398 self.trace_level = level.to_owned();
399 self
400 }
401
402 pub fn build(self, name: String) -> Result<Node, MiddlewareError> {
404
405 cdds::cdds_config::inject_config_if_allowed(self.shared_memory,self.pub_sub_log_level, self.trace_enabled, &self.trace_level);
406
407 let participant = DdsParticipant::create(None, None, None)?;
408 let _dir_res=create_home_directory_if_required(&name);
409 let inner = NodeInner {
410 name,
411 group: self.group,
412 instance: self.instance,
413 participant,
414 maybe_publisher: None,
415 maybe_subscriber: None,
416 handlers: Vec::new(),
417 next_client_id: 0,
418 proxies: Vec::new(),
419 single_threaded: self.single_threaded,
420 num_workers: self.num_workers,
421 stack_size : self.stack_size,
422 };
423
424 Ok(Node {
425 inner: Arc::new(RwLock::new(inner)),
426 })
427 }
428}
429
430struct NodeInner {
431 name: String,
432 group: String,
433 instance: String,
434 participant: DdsParticipant,
435 maybe_publisher: Option<DdsPublisher>,
436 maybe_subscriber: Option<DdsSubscriber>,
437 handlers: Vec<ServerRequestHandlerEntry>,
438 next_client_id: u16,
439 proxies: Vec<(String, Box<dyn Proxy>, u8, u32)>,
440 single_threaded: bool,
441 num_workers: usize,
442 stack_size : Option<usize>,
443}
444
445#[derive(Default)]
446pub struct PublishOptions {
447 group: Option<String>,
448 instance: Option<String>,
449 reliability: Option<QosReliability>,
450 durability: Option<QosDurability>,
451 history: Option<QosHistory>,
452}
453
454impl PublishOptions {
455 pub fn with_group(&mut self, group: &str) -> &mut Self {
456 let _ = self.group.replace(group.to_owned());
457 self
458 }
459
460 pub fn with_instance(&mut self, instance: &str) -> &mut Self {
461 let _ = self.instance.replace(instance.to_owned());
462 self
463 }
464
465 pub fn with_reliability(&mut self, reliability: QosReliability) -> &mut Self {
466 let _ = self.reliability.replace(reliability);
467 self
468 }
469
470 pub fn with_durability(&mut self, durability: QosDurability) -> &mut Self {
471 let _ = self.durability.replace(durability);
472 self
473 }
474
475 pub fn with_history(&mut self, history: QosHistory) -> &mut Self {
476 let _ = self.history.replace(history);
477 self
478 }
479}
480
481#[derive(Default)]
482pub struct SubscribeOptions {
483 group: Option<String>,
484 instance: Option<String>,
485 reliability: Option<QosReliability>,
486 durability: Option<QosDurability>,
487 history: Option<QosHistory>,
488}
489
490impl SubscribeOptions {
491 pub fn with_group(&mut self, group: &str) -> &mut Self {
492 let _ = self.group.replace(group.to_owned());
493 self
494 }
495
496 pub fn with_instance(&mut self, instance: &str) -> &mut Self {
497 let _ = self.instance.replace(instance.to_owned());
498 self
499 }
500
501 pub fn with_reliability(&mut self, reliability: QosReliability) -> &mut Self {
502 let _ = self.reliability.replace(reliability);
503 self
504 }
505
506 pub fn with_durability(&mut self, durability: QosDurability) -> &mut Self {
507 let _ = self.durability.replace(durability);
508 self
509 }
510
511 pub fn with_history(&mut self, history: QosHistory) -> &mut Self {
512 let _ = self.history.replace(history);
513 self
514 }
515}
516
517#[derive(Default)]
518pub struct PublisherListenerBuilder {
519 listener_builder : Option<DdsListenerBuilder>,
520}
521
522impl PublisherListenerBuilder {
523 pub fn new() -> PublisherListenerBuilder {
524 Self { listener_builder: Some(DdsListenerBuilder::new())}
525 }
526
527 pub fn on_publication_matched<F>(&mut self, mut callback: F) -> &mut Self
528 where
529 F: FnMut(PublicationMatchedStatus ) + 'static,
530 {
531 self.listener_builder.as_mut().unwrap().on_publication_matched(move |_e,v|callback(v.into()));
532 self
533 }
534
535 pub fn build(&mut self) -> PublisherListener {
536 PublisherListener(self.listener_builder.take().unwrap())
537 }
538}
539
540pub struct PublisherListener(DdsListenerBuilder);
541
542#[derive(Default)]
543pub struct SubscriberListenerBuilder {
544 listener_builder : Option<DdsListenerBuilder>,
545}
546
547impl SubscriberListenerBuilder {
548 pub fn new() -> SubscriberListenerBuilder {
549 Self { listener_builder: Some(DdsListenerBuilder::new())}
550 }
551
552 pub fn on_subscription_matched<F>(&mut self, mut callback: F) -> &mut Self
553 where
554 F: FnMut(SubscriptionMatchedStatus ) + 'static,
555 {
556 self.listener_builder.as_mut().unwrap()
557 .on_subscription_matched(move |_e,v|callback(v.into()));
558 self
559 }
560
561 pub fn build(&mut self) -> SubscriberListener {
562 SubscriberListener(self.listener_builder.take().unwrap())
563 }
564}
565
566pub struct SubscriberListener(DdsListenerBuilder);
567
568pub struct PublicationMatchedStatus {
569 pub total_count: u32,
570 pub total_count_change: i32,
571 pub current_count: u32,
572 pub current_count_change: i32,
573}
574
575pub struct SubscriptionMatchedStatus {
576 pub total_count: u32,
577 pub total_count_change: i32,
578 pub current_count: u32,
579 pub current_count_change: i32,
580}
581
582impl Node {
583 fn get_topic_prefix(group: &str, instance: &str) -> Option<String> {
584 let prefix = format!("/{}/{}", group, instance);
585 Some(prefix)
586 }
587
588
589 fn advertise_internal<T>(&self, topic_path: &str, options: &PublishOptions) -> Result<Writer<T>, MiddlewareError>
590 where
591 T: TopicType,
592 {
593 if let Ok(mut inner) = self.inner.write() {
594 if inner.maybe_publisher.is_none() {
595 inner.maybe_publisher = Some(PublisherBuilder::new().create(&inner.participant)?);
596 }
597 assert!(inner.maybe_publisher.is_some());
598
599 let topic = TopicBuilder::<T>::new()
600 .with_name(topic_path.to_owned())
601 .create(&inner.participant)?;
602
603 let qos = publish_options_to_cdds_qos(options)?;
604
605 let writer = CDDSWriterBuilder::new()
606 .with_qos(qos.into())
607 .create(inner.maybe_publisher.as_ref().unwrap(), topic)?;
608 Ok(Writer { writer })
609 } else {
610 Err(MiddlewareError::InconsistentDataStructure)
611 }
612 }
613
614 pub fn advertise_with_listener<T: TopicType>(&self, options: &PublishOptions, listener : PublisherListener) -> Result<Writer<T>, MiddlewareError> {
615 self.advertise_with_maybe_listener(options, Some(listener.0))
616 }
617
618 fn advertise_with_maybe_listener<T>(&self, options: &PublishOptions, maybe_listener : Option<DdsListenerBuilder>) -> Result<Writer<T>, MiddlewareError>
619 where
620 T: TopicType,
621 {
622 if let Ok(mut inner) = self.inner.write() {
623 if inner.maybe_publisher.is_none() {
624 inner.maybe_publisher = Some(PublisherBuilder::new().create(&inner.participant)?);
625 }
626 assert!(inner.maybe_publisher.is_some());
627
628 let topic_builder = TopicBuilder::<T>::new();
629
630 let group = if let Some(group) = options.group.as_ref() {
631 group
632 } else {
633 &inner.group
634 };
635
636 let instance = if let Some(instance) = options.instance.as_ref() {
637 instance
638 } else {
639 &inner.instance
640 };
641
642 let topic_builder = if let Some(prefix) = Self::get_topic_prefix(group, instance) {
643 topic_builder.with_name_prefix(prefix)
644 } else {
645 topic_builder
646 };
647
648 let topic = topic_builder.create(&inner.participant)?;
649
650 let qos = publish_options_to_cdds_qos(options)?;
651
652
653 let writer_builder = CDDSWriterBuilder::new()
654 .with_qos(qos.into());
655
656 let writer_builder = if let Some(mut listener_builder) = maybe_listener {
657 let listener = listener_builder.build();
658 writer_builder.with_listener(listener)
659 } else { writer_builder};
660
661 let w = writer_builder.create(inner.maybe_publisher.as_ref().unwrap(), topic)?;
662 Ok(Writer { writer: w })
663 } else {
664 Err(MiddlewareError::InconsistentDataStructure)
665 }
666 }
667
668 pub fn advertise<T>(&self, options: &PublishOptions) -> Result<Writer<T>, MiddlewareError>
672 where
673 T: TopicType,
674 {
675 self.advertise_with_maybe_listener(options, None)
676 }
677
678 fn subscribe_with_maybe_listener<T>(
679 &self,
680 options: &SubscribeOptions,
681 maybe_listener : Option<SubscriberListener>
682 ) -> Result<Reader<T>, MiddlewareError>
683 where
684 T: TopicType,
685 {
686 if let Ok(mut inner) = self.inner.write() {
687 if inner.maybe_subscriber.is_none() {
688 inner.maybe_subscriber = Some(SubscriberBuilder::new().create(&inner.participant)?);
689 }
690 assert!(inner.maybe_subscriber.is_some());
691
692 let group = if let Some(group) = options.group.as_ref() {
693 group
694 } else {
695 &inner.group
696 };
697
698 let instance = if let Some(instance) = options.instance.as_ref() {
699 instance
700 } else {
701 &inner.instance
702 };
703
704 let topic_builder = TopicBuilder::<T>::new();
705
706 let topic_builder = if let Some(prefix) = Self::get_topic_prefix(group,instance) {
707 topic_builder.with_name_prefix(prefix)
708 } else {
709 topic_builder
710 };
711
712 let topic = topic_builder.create(&inner.participant)?;
713
714 let qos = subscribe_options_to_cdds_qos(options)?;
715
716 let mut reader = ReaderBuilder::new()
717 .with_qos(qos.into());
718
719 if let Some(mut listener_builder) = maybe_listener {
720 reader = reader.with_listener(listener_builder.0.build())
721 }
722 Ok(Reader { reader : reader.create(inner.maybe_subscriber.as_ref().unwrap(), topic)? })
723 } else {
724 Err(MiddlewareError::InconsistentDataStructure)
725 }
726 }
727
728 pub fn subscribe_with_listener<T>(&self, options: &SubscribeOptions, listener: SubscriberListener) -> Result<Reader<T>, MiddlewareError>
729 where T: TopicType {
730 self.subscribe_with_maybe_listener(options, Some(listener))
731 }
732
733 pub fn subscribe<T>(
736 &self,
737 options: &SubscribeOptions,
738 ) -> Result<Reader<T>, MiddlewareError>
739 where
740 T: TopicType,
741 {
742 self.subscribe_with_maybe_listener(options, None)
743 }
744
745 fn subscribe_async_internal<T>(&self, topic_path: &str,options: &SubscribeOptions,) -> Result<Reader<T>, MiddlewareError>
746 where
747 T: TopicType,
748 {
749 if let Ok(mut inner) = self.inner.write() {
750
751 if inner.maybe_subscriber.is_none() {
752 inner.maybe_subscriber = Some(SubscriberBuilder::new().create(&inner.participant)?);
753 }
754 assert!(inner.maybe_subscriber.is_some());
755
756 let topic = TopicBuilder::<T>::new()
757 .with_name(topic_path.to_owned())
758 .create(&inner.participant)?;
759
760 let qos = subscribe_options_to_cdds_qos(options)?;
761
762 let reader = ReaderBuilder::new()
763 .with_qos(qos.into())
764 .as_async()
765 .create(inner.maybe_subscriber.as_ref().unwrap(), topic)?;
766 Ok(Reader { reader })
767 } else {
768 Err(MiddlewareError::InconsistentDataStructure)
769 }
770 }
771
772 fn subscribe_async_with_maybe_listener<T>(&self, options: &SubscribeOptions, maybe_listener : Option<SubscriberListener>) -> Result<Reader<T>, MiddlewareError>
773 where
774 T: TopicType,
775 {
776 if let Ok(mut inner) = self.inner.write() {
777 if inner.maybe_subscriber.is_none() {
778 inner.maybe_subscriber = Some(SubscriberBuilder::new().create(&inner.participant)?);
779 }
780 assert!(inner.maybe_subscriber.is_some());
781
782 let group = if let Some(group) = options.group.as_ref() {
783 group
784 } else {
785 &inner.group
786 };
787
788 let instance = if let Some(instance) = options.instance.as_ref() {
789 instance
790 } else {
791 &inner.instance
792 };
793
794 let topic_builder = TopicBuilder::<T>::new();
795
796 let topic_builder = if let Some(prefix) = Self::get_topic_prefix(group,instance) {
797 topic_builder.with_name_prefix(prefix)
798 } else {
799 topic_builder
800 };
801
802 let topic = topic_builder.create(&inner.participant)?;
803
804 let qos = subscribe_options_to_cdds_qos(options)?;
805
806 let mut reader = ReaderBuilder::new()
807 .with_qos(qos.into())
808 .as_async();
809
810 if let Some(mut listener_builder) = maybe_listener {
811 reader = reader.with_listener(listener_builder.0.build())
812 }
813 let reader = reader.create(inner.maybe_subscriber.as_ref().unwrap(), topic)?;
814 Ok(Reader { reader })
815 } else {
816 Err(MiddlewareError::InconsistentDataStructure)
817 }
818 }
819
820 pub fn subscribe_async<T>(&self, options: &SubscribeOptions) -> Result<Reader<T>, MiddlewareError>
823 where
824 T: TopicType,
825 {
826 self.subscribe_async_with_maybe_listener(options,None)
827 }
828
829 pub fn create_proxy<
831 T: 'static + Proxy + ProxyConstruct + ServiceIdentifier + ServiceVersion + Clone,
832 >(
833 &self,
834 ) -> Result<T, MiddlewareError> {
835 let config = Arc::new(Configuration::default());
836 let config_path = get_config_path()?;
837 let service_ids = vec![T::service_name()];
838 let service_ids = get_service_ids(&config_path, &service_ids)?;
839 if service_ids.len() != 1 {
840 return Err(MiddlewareError::ConfigurationError);
841 }
842 if let Ok(mut inner) = self.inner.write() {
843 let proxy = T::new(service_ids[0].1, inner.next_client_id, config);
844 inner.next_client_id += 1;
845 inner.proxies.push((
846 T::service_name().to_owned(),
847 Box::new(proxy.clone()),
848 T::__major_version__(),
849 T::__minor_version__(),
850 ));
851 Ok(proxy)
852 } else {
853 Err(MiddlewareError::InconsistentDataStructure)
854 }
855 }
856
857 pub fn serve<T: CreateServerRequestHandler<Item = T>>(
859 &self,
860 server_impl: Arc<T>,
861 ) -> Result<(), MiddlewareError> {
862 let mut handlers = T::create_server_request_handler(server_impl);
863 if let Ok(mut inner) = self.inner.write() {
864 inner.handlers.append(&mut handlers);
865 Ok(())
866 } else {
867 Err(MiddlewareError::InconsistentDataStructure)
868 }
869 }
870
871 pub fn spin<F>(&self, main_function: F) -> Result<(), MiddlewareError>
874 where
875 F: 'static + Send + FnOnce(),
876 {
877 let mut builder = if self.inner.read().unwrap().single_threaded {
878 Builder::new_current_thread()
879 } else {
880 let mut builder = Builder::new_multi_thread();
881 builder.worker_threads(self.inner.read().unwrap().num_workers);
882 builder
883 };
884
885 builder.thread_name(format!("{}-worker", self.inner.read().unwrap().name));
886 builder.enable_all();
887
888 if let Some(size) = self.inner.read().unwrap().stack_size {
889 builder.thread_stack_size(size as usize);
890 }
891
892 let rt = builder
893 .build()
894 .map_err(|_e| MiddlewareError::InternalError)?;
895
896 let config = Arc::new(Configuration::default());
897
898 let maybe_services = if let Ok(inner) = self.inner.read() {
899 let service_handlers: Vec<&str> = inner.handlers.iter().map(|h| h.name).collect();
900
901 let maybe_services = if !service_handlers.is_empty() {
902 let config_path = get_config_path()?;
903 let services = crate::services::get_service_ids(&config_path, &service_handlers)?;
904
905 if services.len() != service_handlers.len() {
906 error!("Could not get all service IDs");
907 return Err(MiddlewareError::ConfigurationError);
908 }
909
910 let services: Vec<(String, u16, u8, u32, u16, Arc<dyn ServerRequestHandler>)> =
911 services
912 .into_iter()
913 .map(|(s, id)| {
914 let h = inner
915 .handlers
916 .iter()
917 .find(|h| h.name == s.as_str())
918 .unwrap();
919
920 (
921 s,
922 id,
923 h.major_version,
924 h.minor_version,
925 h.instance_id,
926 h.handler.clone(),
927 )
928 })
929 .collect();
930
931 Some(services)
932 } else {
933 None
934 };
935 maybe_services
936 } else {
937 None
938 };
939
940 let maybe_proxies = if let Ok(mut inner) = self.inner.write() {
941 let proxies: Vec<(String, Box<dyn Proxy>, u8, u32)> =
942 inner.proxies.drain(0..).collect();
943 Some(proxies)
944 } else {
945 None
946 };
947
948 debug!("starting tokio main loop");
949 rt.block_on(async {
951 if let Some(services) = maybe_services {
952 for (service, service_id, major_version, minor_version, instance_id, handler) in
953 services
954 {
955 let config = config.clone();
956 let node_name = self.inner.read().unwrap().name.clone();
957 let topic_name = service_name_to_topic_name(&service);
958 let mut sd_publisher = self
959 .advertise_internal::<ServiceInfo>(&topic_name,
960 PublishOptions::default().with_durability(QosDurability::TransientLocal))
961 .expect("Unable to create topic publisher for SD");
962
963 let (tx, mut rx) = Server::create_notify_channel(2);
964
965 tokio::spawn(async move {
967 loop {
968 if let Some(msg) = rx.recv().await {
969 match msg {
970 ConnectionInfo::ConnectionDropped(_i) => {}
971 ConnectionInfo::UdpServerSocket(_s) => {
972 }
992 ConnectionInfo::TcpServerSocket(s) => {
993 println!("Local TCP socket {:?}", s);
994
995 let service_info = ServiceInfo {
996 node: node_name.clone(),
997 major_version,
998 minor_version,
999 instance_id,
1000 socket_address: s,
1001 transport: Transport::Tcp,
1002 service_id,
1003 };
1004 println!("Going to Publish SD packet");
1005
1006 sd_publisher
1007 .publish(Arc::new(service_info))
1008 .expect("Unable to publish SD topic");
1009 println!("Published SD packet");
1010 }
1011 _ => {}
1012 }
1013 }
1014 }
1015 });
1016
1017 tokio::spawn(async move {
1019 debug!("Going to run server for {}", &service);
1022 println!(
1023 "Going to run server for {} at address {:?}",
1024 &service,
1025 get_bind_address()
1026 );
1027 let res = Server::serve(
1028 get_bind_address(),
1029 handler.clone(),
1030 config,
1031 service_id,
1032 major_version,
1033 minor_version,
1034 tx,
1035 )
1036 .await;
1037 println!("Server terminated");
1038 if let Err(e) = res {
1039 println!("Server error:{}", e);
1040 }
1041 });
1042 }
1043 }
1044
1045 let instance_id = 0;
1046 if let Some(proxies) = maybe_proxies {
1048 for (name, proxy, major_version, minor_version) in proxies {
1049 let topic_name = service_name_to_topic_name(&name);
1050 println!("Starting proxy for {} at {}", &name, &topic_name);
1051 let mut sd_subsriber = self
1052 .subscribe_async_internal::<ServiceInfo>(
1053 &topic_name,
1054 SubscribeOptions::default().with_durability(QosDurability::TransientLocal))
1055 .unwrap();
1056
1057 let mut samples = Samples::<ServiceInfo>::new(5);
1059 let client = proxy.get_dispatcher();
1060
1061 tokio::spawn(async move {
1062 let mut is_running = false;
1063 loop {
1064 if let Ok(_len) = sd_subsriber.take(&mut samples).await {
1065 for sample in samples.iter() {
1066 println!("Got sample {:?}", sample.deref());
1067
1068 let sample_socket_address = sample.socket_address;
1069
1070 if sample.major_version == major_version
1071 && sample.instance_id == instance_id
1072 && sample.minor_version == minor_version
1073 && !is_running
1074 && sample.transport == Transport::Tcp
1075 {
1076 let name = name.clone();
1077 let client = client.clone();
1078 tokio::spawn(async move {
1079 debug!(
1080 "Going to run proxy for {} connecting to {}",
1081 &name, sample_socket_address
1082 );
1083 if let Err(_res) =
1084 client.run(sample_socket_address).await
1085 {
1086 error!("Proxy run returned error");
1087 }
1088 });
1089 is_running = true;
1090 } else {
1091 }
1093 }
1094 tokio::time::sleep(Duration::from_millis(1000)).await;
1095 }
1096 }
1097 });
1098 }
1099 }
1100
1101 tokio::task::spawn_blocking( move || {
1104 main_function();
1105 });
1106
1107 let mut signals = Signals::new(vec![libc::SIGINT]).unwrap();
1109 let _signal = signals.next().await.unwrap();
1110 debug!("SIGINT received");
1111 });
1112 debug!("Tokio main loop exited");
1113 Ok(())
1114 }
1115
1116 pub fn terminate(&self) {
1119 unsafe {libc::kill(std::process::id() as i32, libc::SIGINT);}
1120 }
1121}
1122
1123#[cfg(test)]
1124mod simple_tests {
1125 use std::thread;
1126
1127 use super::*;
1128 use async_trait::async_trait;
1129 use cyclonedds_rs::cdr;
1130 use cyclonedds_rs::DDSError;
1131 use cyclonedds_rs::DdsListener;
1132 use cyclonedds_rs::DdsQos;
1133 use cyclonedds_rs::DdsTopic;
1134 use cyclonedds_rs::SampleBuffer;
1135 use cdds_derive::Topic;
1136 use interface_example::EchoResponse;
1137 use interface_example::ExampleError;
1138 use interface_example::ExampleProxy;
1139 use interface_example::ExampleStatus;
1140 use interface_example::{Example, ExampleDispatcher};
1141 use serde_derive::Deserialize;
1142 use serde_derive::Serialize;
1143 use someip::CallProperties;
1144 use someip::*;
1145 use someip_derive::*;
1146 #[test]
1147 fn it_works() {
1148 #[derive(Default, Deserialize, Serialize, Topic)]
1149 struct A {
1150 name: String,
1151 }
1152
1153 let node = NodeBuilder::default()
1154 .with_group_and_instance("group_name".to_owned(), "instance_name".to_owned())
1155 .build("nodename".to_owned())
1156 .expect("Node creation");
1157 let mut subscriber = node
1158 .subscribe::<A>(&SubscribeOptions::default())
1159 .expect("unable to subscribe");
1160
1161 let mut p = node
1162 .advertise::<A>(&PublishOptions::default())
1163 .expect("cannot create writer");
1164
1165 let a = A {
1166 name: "foo".to_owned(),
1167 };
1168
1169 let publisher_listener_builder =
1170 PublisherListenerBuilder::new()
1171 .on_publication_matched(move |s| {
1172 println!("Publication matched {}:{}", s.total_count, s.total_count_change);
1173 })
1174 .build();
1175
1176
1177 let mut wb = node.advertise_with_listener::<A>(&PublishOptions::default(), publisher_listener_builder)
1178 .expect("Cannot create writer with listener");
1179
1180 wb.publish(Arc::new(A { name : "bar".to_owned()})).expect("Cannot publish from wb");
1181
1182
1183 p.publish(Arc::new(a)).expect("Cannot publish");
1184
1185 let mut samples = Samples::<A>::new(2);
1186 subscriber.take_now(&mut samples).expect("Unable to take");
1187
1188 for sample in samples.iter() {
1189 println!("A -> {}", sample.name);
1190 }
1191
1192 }
1193
1194 #[test]
1195 fn host_service() {
1196 std::env::set_var("SERVICE_MAPPING_CONFIG_PATH", "services.toml");
1197
1198 #[service_impl(Example)]
1199 pub struct EchoServerImpl {}
1200
1201 impl ServiceInstance for EchoServerImpl {}
1202 impl ServiceVersion for EchoServerImpl {}
1203
1204 #[async_trait]
1205 impl Example for EchoServerImpl {
1206 async fn echo(&self, _data: String) -> Result<EchoResponse, ExampleError> {
1207 Err(ExampleError::Unknown)
1208 }
1209
1210 fn set_status(
1211 &self,
1212 _status: interface_example::ExampleStatus,
1213 ) -> Result<(), someip::error::FieldError> {
1214 Ok(())
1215 }
1216
1217 fn get_status(&self) -> Result<&ExampleStatus, someip::error::FieldError> {
1218 Ok(&ExampleStatus::Ready)
1219 }
1220 }
1221
1222 let node = NodeBuilder::default()
1223 .build("nodename".to_owned())
1224 .expect("Node creation");
1225
1226 let server = Arc::new(EchoServerImpl {});
1227
1228 node.serve(server).expect("Unable to serve");
1229
1230 let node_cp = node.clone();
1231
1232 node.spin(move || {println!("This is the main function");node_cp.terminate()})
1233 .expect("Unable to spin");
1234 }
1235
1236 #[test]
1237 fn client() {
1238 std::env::set_var("SERVICE_MAPPING_CONFIG_PATH", "services.toml");
1239 tracing_subscriber::fmt::init();
1240
1241 thread::spawn(|| {
1244 let client_node = NodeBuilder::default()
1245 .build("client".to_owned())
1246 .expect("Node creation");
1247 let proxy = client_node
1248 .create_proxy::<ExampleProxy>()
1249 .expect("Unable to create proxy");
1250
1251 let node_cp = client_node.clone();
1252
1253 client_node
1254 .spin(move || {
1255 let proxy = proxy.clone();
1256 tokio::spawn(async move {
1257 tokio::time::sleep(Duration::from_millis(3000)).await;
1258
1259 let call_properties =
1260 CallProperties::with_timeout(Duration::from_millis(5000));
1261
1262 match proxy.echo("Hello".to_string(), &call_properties).await {
1263 Ok(res) => {
1264 assert_eq!(res.echo.as_str(), "Hello");
1265 println!("Received echo");
1266 }
1267 Err(e) => {
1268 println!("Error:{:?}", e);
1269 panic!("Echo response failed");
1270 }
1271 }
1272 });
1273 node_cp.terminate();
1274 })
1275 .expect("Unable to spin");
1276 });
1278
1279 #[service_impl(Example)]
1280 pub struct EchoServerImpl {}
1281
1282 impl ServiceInstance for EchoServerImpl {}
1283 impl ServiceVersion for EchoServerImpl {}
1284
1285 #[async_trait]
1286 impl Example for EchoServerImpl {
1287 async fn echo(&self, data: String) -> Result<EchoResponse, ExampleError> {
1288 let response = EchoResponse { echo: data };
1289 Ok(response)
1290 }
1291
1292 fn set_status(
1293 &self,
1294 _status: interface_example::ExampleStatus,
1295 ) -> Result<(), someip::FieldError> {
1296 Ok(())
1297 }
1298
1299 fn get_status(&self) -> Result<&interface_example::ExampleStatus, someip::FieldError> {
1300 Ok(&ExampleStatus::Ready)
1301 }
1302 }
1303
1304 let node = NodeBuilder::default()
1306 .build("server".to_owned())
1307 .expect("Node creation");
1308
1309 let server = Arc::new(EchoServerImpl {});
1310
1311 let node_cp = node.clone();
1312 node.serve(server).expect("Unable to serve");
1313
1314 node.spin(|| {
1315 tokio::spawn(async move {
1316 let mut ticker = tokio::time::interval(Duration::from_millis(100));
1317
1318 let _ = ticker.tick().await;
1319 node_cp.terminate();
1320
1321 });
1322 })
1323 .expect("Unable to spin");
1324 }
1325}
1326