Skip to main content

sabaton_mw/
lib.rs

1/*
2    Copyright (C) Sabaton Systems LLP - All Rights Reserved
3    Sojan James <sojan.james@gmail.com>, 2021
4
5    SPDX-License-Identifier: Apache-2.0 OR LicenseRef-sabaton-commercial
6*/
7
8//! The Sabaton Middleware is the interface Sabaton applications
9//! use. 
10//! 
11//! Applications can,
12//! 1. Publish topics
13//! 2. Subscribe to topics
14//! 3. Provide a Service (using SOME/IP)
15//! 4. Access a service via a proxy
16//! 
17//! The Sabaton middleware hides the underlying implementation
18//! of DDS and SOME/IP. 
19//! 
20use 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
59//pub use cdds::cdds::CddsQos as QosImpl;
60
61const 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    /// Publish data to the topic writer
99    ///
100    /// # Arguments
101    ///
102    /// * `msg` - The data to be published wrapped in an Arc<T>. 
103    ///
104    pub fn publish(&mut self, msg: Arc<T>) -> Result<(), MiddlewareError> {
105        self.writer
106            .write(msg)
107            .map_err(MiddlewareError::DDSError)
108    }
109
110    /// Loan a buffer from the writer. This may not be supported always.
111    /// The topic must of Fixed size and shared memory must be enabled in
112    /// cyclone for this to work.
113    /// 
114    /// Loaning is useful for large buffers being sent locally, like image buffers
115    /// Buffers are allocated from a shared memory pool and a reference to the buffer
116    /// is sent to the readers
117    /// 
118    /// Important:  Shared memory topics can only be published to recipients on the same
119    /// machine.
120    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    /// Return the loan that was taken.  The buffer will be published if it is marked
131    /// as initialized by the ``Loaned::assume_init`` function. If not initialized
132    /// the buffer will be simple returned to the pool.
133    /// 
134    /// # Arguments
135    /// 
136    /// * `buffer` - The buffer that was loaned.  The Loaned<T> holds an initialization
137    ///              state. If the previously loaned buffer was not initialized, the buffer
138    ///              will be returned to the pool without publishing it.
139    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    /// Access the buffer via a mutable pointer so you can
155    /// write into it
156    pub fn as_mut_ptr(&mut self) -> Option<*mut T> {
157        self.inner.as_mut_ptr()
158    }
159
160    /// Mark the loaned buffer as initialized. You will call this method
161    /// after writing the data via the pointer you got from ``Loaned::as_mut_ptr``
162    /// TODO: Perhaps this should be made unsafe
163    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
186/// A buffer to store samples. This is used for receiving
187/// one or more samples from the reader.
188pub struct Samples<T: TopicType> {
189    samples: SampleBuffer<T>,
190}
191
192impl<'a, T> Samples<T>
193where
194    T: TopicType,
195{
196    /// Create a new sample buffer with `len` elements.
197    /// 
198    /// # Arguments
199    /// 
200    /// * `len` - number of elements to store in the sample buffer
201    /// 
202    pub fn new(len: usize) -> Self {
203        Self {
204            samples: SampleBuffer::new(len),
205        }
206    }
207
208    /// Create an iterator to iterate over valid sample buffers
209    /// Invalid samples will be skipped by the iterator
210    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    /// Synchronous take. This call willl block until atleast one sample is read
226    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    /// Synchronous read. This call willl block until atleast one sample is read
232    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    /// Asynchronous take
246    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    /// Asynchronous read
256    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
267/* 
268#[derive(Default)]
269pub struct ReaderListenerBuilder{listener: DdsListenerBuilder}
270
271impl ReaderListenerBuilder {
272    pub fn new() -> Self {
273        ReaderListenerBuilder::default()
274    }
275
276    pub fn on_data_available<F>(&mut self, callback: F) -> &mut Self
277    where
278        F: FnMut() + 'static ,
279    {
280        self.listener = self.listener.on_data_available(|_|callback());
281        self
282    }
283}
284*/
285pub struct ReaderListener {
286    listener: DdsListener,
287}
288
289/// All the functionality is implemented in the Node structure. To interact with the other
290/// services of the Sabaton system, a node structure must be created.
291/// The `NodeBuilder` structure provides a builder pattern to create 
292/// the node. A node is cloneable.
293#[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    /// Set the group name of the node. This impacts how topic address are created.
330    /// You normally don't need to change the group.  The default value of the group
331    /// is "default"
332    pub fn with_group(mut self, group: String) -> Self {
333        self.group = group;
334        self
335    }
336
337    /// Set the instance of the node. This impacts how topic address are created. 
338    /// The default instance is "0".
339    pub fn with_instance(mut self, instance: String) -> Self {
340        self.instance = instance;
341        self
342    }
343
344    /// Just a convenience function to set both the group and instance.
345    #[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    // Set the stack size of the tasks in the threadpool
354    pub fn with_stack_size(mut self, size: usize) -> Self {
355        self.stack_size = Some(size);
356        self
357    }
358
359
360    /// Use a multi-threaded runtime.
361    pub fn multi_threaded(mut self) -> Self {
362        self.single_threaded = false;
363        self
364    }
365
366    /// Number of worker threads to use in the multi-threaded runtime.
367    /// This function will panic if used on a single threaded runtime
368    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    /// Enable shared memory.  Shared memory will work only
378    /// if the underlying shared memory transport is available.
379    /// This means iox-roudi must be running with enough of 
380    /// memory allocated to support the shared memory topic.
381    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    /// Create the Node structure
403    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    /// Advertise a Type to the rest of the system. This call returns a Writer<T> which you
669    /// can use to publish samples to the topic. The topic name is create from the type of T and 
670    /// the combination of the group and instance.
671    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    /// Subscribe to a topic. This call returns a Reader<T>.  You can read samples from
734    /// the reader.
735    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    /// Subscribe to a topic. This call returns a Reader<T>.  You can read samples from
821    /// the reader. The reader that is returned supports asynchronous reads.
822    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    /// create a proxy for a service
830    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    ///Hosting a service
858    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    /// The main processing loop of the node.  This function will block waiting for events and pumping
872    /// the necessary callbacks.
873    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        // blocking main loop
950        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 task to publish Service discovery topic for this service
966                    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                                        /*   We don't support UDP yet - don't publish UDP SD Message */
973                                        //println!("Local UDP socket {:?}", s);
974                                        /*
975                                        let service_info = ServiceInfo {
976                                            node: node_name.clone(),
977                                            major_version,
978                                            minor_version,
979                                            instance_id,
980                                            socket_address: s,
981                                            transport: Transport::Udp,
982                                            service_id,
983                                        };
984                                        println!("Going to Publish SD packet for Udp");
985
986                                        sd_publisher
987                                            .publish(Arc::new(service_info))
988                                            .expect("Unable to publish SD topic");
989                                        println!("Published SD packet");
990                                        */
991                                    }
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 task for this service
1018                    tokio::spawn(async move {
1019                        //let test_service : Box<dyn ServerRequestHandler + Send> = Box::new(EchoServerImpl::default());
1020                        //let handler = EchoServerImpl::create_server_request_handler(Arc::new(EchoServerImpl::default()));
1021                        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            // launch the proxies
1047            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                    // max of 5 instances for a services. TODO: this could be in a config file
1058                    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                                        // ignore
1092                                    }
1093                                }
1094                                tokio::time::sleep(Duration::from_millis(1000)).await;
1095                            }
1096                        }
1097                    });
1098                }
1099            }
1100
1101            // launch the main function
1102
1103            tokio::task::spawn_blocking( move || {
1104                main_function();
1105            });
1106
1107            // wait for SIGINT
1108            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    /// Terminate the processing of the node. The spin function hangs around waiting
1117    /// for a SIGINT.
1118    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        // Client node in separate thread
1242
1243        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            //let cloned = node.clone();
1277        });
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        // Server node
1305        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