Preview of the proposed up-rust native-frame-model branch (up-rust b6b99c6d, up-spec f0e9b17) — not released documentation. branch · write-up

up_rust/core/usubscription/
usubscription_client.rs

1/********************************************************************************
2 * Copyright (c) 2024 Contributors to the Eclipse Foundation
3 *
4 * See the NOTICE file(s) distributed with this work for additional
5 * information regarding copyright ownership.
6 *
7 * This program and the accompanying materials are made available under the
8 * terms of the Apache License Version 2.0 which is available at
9 * https://www.apache.org/licenses/LICENSE-2.0
10 *
11 * SPDX-License-Identifier: Apache-2.0
12 ********************************************************************************/
13
14use std::sync::Arc;
15
16use async_trait::async_trait;
17use protobuf::well_known_types::timestamp::Timestamp;
18
19use crate::{
20    communication::{CallOptions, RpcClient, SubscriptionStatus},
21    core::usubscription::{
22        usubscription_uri, ResetReason, SubscriptionInfo, USubscription,
23        RESOURCE_ID_FETCH_SUBSCRIBERS, RESOURCE_ID_FETCH_SUBSCRIPTIONS,
24        RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS, RESOURCE_ID_RESET, RESOURCE_ID_SUBSCRIBE,
25        RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS, RESOURCE_ID_UNSUBSCRIBE,
26    },
27    up_core_api::usubscription::{
28        FetchSubscribersResponse, FetchSubscriptionsRequest, FetchSubscriptionsResponse,
29        NotificationsResponse, ResetResponse, SubscribeAttributes, Subscription,
30        SubscriptionRequest, SubscriptionResponse, UnsubscribeRequest, UnsubscribeResponse, Update,
31    },
32    UCode, UStatus, UUri,
33};
34
35fn unix_epoch_millis_as_protobuf_timestamp(
36    millis: Option<u64>,
37) -> Result<Option<Timestamp>, UStatus> {
38    if let Some(milliseconds) = millis {
39        // this will always yield a valid Timestamp as the maximum value of u64 (2^64 - 1) divided by 1000
40        // is less than the maximum number of milliseconds that can be represented in an i64 (2^63 - 1)
41        let seconds = (milliseconds / 1000_u64) as i64;
42        let nanos = (milliseconds % 1000)
43            .checked_mul(1_000_000)
44            .ok_or_else(|| {
45                UStatus::fail_with_code(UCode::InvalidArgument, "timestamp out of range")
46            })
47            .and_then(|s| {
48                i32::try_from(s).map_err(|_| {
49                    UStatus::fail_with_code(UCode::InvalidArgument, "timestamp out of range")
50                })
51            })?;
52        Ok(Some(Timestamp {
53            seconds,
54            nanos,
55            ..Default::default()
56        }))
57    } else {
58        Ok(None)
59    }
60}
61
62fn protobuf_timestamp_as_unix_epoch_milliseconds(
63    ts: Option<&Timestamp>,
64) -> Result<Option<u64>, UStatus> {
65    if let Some(ts) = ts {
66        let err = || {
67            UStatus::fail_with_code(
68                UCode::InvalidArgument,
69                "invalid timestamp: seconds value out of range",
70            )
71        };
72        if ts.nanos < 0 || ts.nanos >= 1_000_000_000 {
73            return Err(UStatus::fail_with_code(
74                UCode::InvalidArgument,
75                "invalid timestamp: nanos value out of range",
76            ));
77        }
78        u64::try_from(ts.seconds)
79            .ok()
80            .and_then(|s| s.checked_mul(1000))
81            .and_then(|ms| ms.checked_add(ts.nanos as u64 / 1_000_000))
82            .ok_or_else(err)
83            .map(Some)
84    } else {
85        Ok(None)
86    }
87}
88
89impl TryFrom<&Subscription> for SubscriptionInfo {
90    type Error = UStatus;
91
92    fn try_from(subscription_proto: &Subscription) -> Result<Self, Self::Error> {
93        let topic = subscription_proto
94            .topic
95            .as_ref()
96            .ok_or(UStatus::fail_with_code(
97                UCode::InvalidArgument,
98                "topic missing",
99            ))
100            .and_then(|t| {
101                UUri::try_from(t)
102                    .map_err(|_| UStatus::fail_with_code(UCode::InvalidArgument, "invalid topic"))
103            })?;
104        let subscriber = subscription_proto
105            .subscriber
106            .as_ref()
107            .and_then(|s| s.uri.as_ref())
108            .ok_or(UStatus::fail_with_code(
109                UCode::InvalidArgument,
110                "subscriber missing",
111            ))
112            .and_then(|s| {
113                UUri::try_from(s).map_err(|_| {
114                    UStatus::fail_with_code(UCode::InvalidArgument, "invalid subscriber")
115                })
116            })?;
117        let status = subscription_proto
118            .status
119            .as_ref()
120            .ok_or(UStatus::fail_with_code(
121                UCode::InvalidArgument,
122                "status missing",
123            ))
124            .and_then(SubscriptionStatus::try_from)?;
125        subscription_proto
126            .attributes
127            .as_ref()
128            .ok_or_else(|| UStatus::fail_with_code(UCode::InvalidArgument, "missing attributes"))
129            .and_then(|attributes| {
130                let expiration =
131                    protobuf_timestamp_as_unix_epoch_milliseconds(attributes.expire.as_ref())?;
132                Ok(SubscriptionInfo::new(
133                    topic,
134                    subscriber,
135                    status,
136                    expiration,
137                    attributes.sample_period_ms,
138                ))
139            })
140    }
141}
142
143impl TryFrom<&Update> for SubscriptionInfo {
144    type Error = UStatus;
145    fn try_from(update_proto: &Update) -> Result<Self, Self::Error> {
146        let topic = update_proto
147            .topic
148            .as_ref()
149            .ok_or(UStatus::fail_with_code(
150                UCode::InvalidArgument,
151                "topic missing",
152            ))
153            .and_then(|t| {
154                UUri::try_from(t)
155                    .map_err(|_| UStatus::fail_with_code(UCode::InvalidArgument, "invalid topic"))
156            })?;
157        let subscriber = update_proto
158            .subscriber
159            .as_ref()
160            .and_then(|s| s.uri.as_ref())
161            .ok_or(UStatus::fail_with_code(
162                UCode::InvalidArgument,
163                "subscriber missing",
164            ))
165            .and_then(|s| {
166                UUri::try_from(s).map_err(|_| {
167                    UStatus::fail_with_code(UCode::InvalidArgument, "invalid subscriber")
168                })
169            })?;
170        let status = update_proto
171            .status
172            .as_ref()
173            .ok_or(UStatus::fail_with_code(
174                UCode::InvalidArgument,
175                "status missing",
176            ))
177            .and_then(SubscriptionStatus::try_from)?;
178        let attribs = update_proto.attributes.get_or_default();
179        let expiration = protobuf_timestamp_as_unix_epoch_milliseconds(attribs.expire.as_ref())?;
180        Ok(SubscriptionInfo::new(
181            topic,
182            subscriber,
183            status,
184            expiration,
185            attribs.sample_period_ms,
186        ))
187    }
188}
189
190impl From<ResetReason> for crate::up_core_api::usubscription::reset_request::reason::Code {
191    fn from(reason: ResetReason) -> Self {
192        match reason {
193            ResetReason::Unspecified => Self::UNSPECIFIED,
194            ResetReason::FactoryReset => Self::FACTORY_RESET,
195            ResetReason::CorruptedData => Self::CORRUPTED_DATA,
196        }
197    }
198}
199
200/// A [`USubscription`] client implementation for invoking operations of a local USubscription service.
201///
202/// The client requires an [`RpcClient`] for performing the remote procedure calls.
203pub struct RpcClientUSubscription {
204    rpc_client: Arc<dyn RpcClient>,
205}
206
207impl core::fmt::Debug for RpcClientUSubscription {
208    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
209        f.debug_struct("RpcClientUSubscription")
210            .finish_non_exhaustive()
211    }
212}
213
214impl RpcClientUSubscription {
215    /// Creates a new Notifier for a given transport.
216    ///
217    /// # Arguments
218    ///
219    /// * `rpc_client` - The client to use for performing the remote procedure calls on the USubscription service.
220    pub fn new(rpc_client: Arc<dyn RpcClient>) -> Self {
221        RpcClientUSubscription { rpc_client }
222    }
223
224    fn default_call_options() -> CallOptions {
225        CallOptions::for_rpc_request(5_000, None, None, None)
226    }
227}
228
229impl RpcClientUSubscription {
230    async fn fetch_subscriptions(
231        &self,
232        fetch_subscriptions_request: FetchSubscriptionsRequest,
233    ) -> Result<Vec<SubscriptionInfo>, UStatus> {
234        let response = self
235            .rpc_client
236            .invoke_proto_method::<_, FetchSubscriptionsResponse>(
237                usubscription_uri(RESOURCE_ID_FETCH_SUBSCRIPTIONS),
238                Self::default_call_options(),
239                fetch_subscriptions_request,
240            )
241            .await
242            .map_err(UStatus::from)?;
243        let mut result = Vec::new();
244        for subscription in &response.subscriptions {
245            let info = SubscriptionInfo::try_from(subscription)?;
246            result.push(info);
247        }
248        Ok(result)
249    }
250}
251
252#[async_trait]
253impl USubscription for RpcClientUSubscription {
254    async fn subscribe(
255        &self,
256        topic: &UUri,
257        expiration: Option<u64>, // millis since Unix Epoch
258        min_sample_period: Option<u32>,
259    ) -> Result<SubscriptionStatus, UStatus> {
260        let subscription_request = SubscriptionRequest {
261            topic: Some(topic.into()).into(),
262            attributes: match (expiration, min_sample_period) {
263                (None, None) => None.into(),
264                _ => Some(SubscribeAttributes {
265                    expire: unix_epoch_millis_as_protobuf_timestamp(expiration)?.into(),
266                    sample_period_ms: min_sample_period,
267                    ..Default::default()
268                })
269                .into(),
270            },
271            ..Default::default()
272        };
273        self.rpc_client
274            .invoke_proto_method::<_, SubscriptionResponse>(
275                usubscription_uri(RESOURCE_ID_SUBSCRIBE),
276                Self::default_call_options(),
277                subscription_request,
278            )
279            .await
280            .and_then(|response| {
281                Ok(response.status.as_ref().map_or_else(
282                    || {
283                        Err(UStatus::fail_with_code(
284                            UCode::InvalidArgument,
285                            "uSubscription returned invalid response: no subscription status",
286                        ))
287                    },
288                    SubscriptionStatus::try_from,
289                )?)
290            })
291            .map_err(UStatus::from)
292    }
293
294    async fn unsubscribe(&self, topic: &UUri) -> Result<(), UStatus> {
295        let unsubscribe_request = UnsubscribeRequest {
296            topic: Some(topic.into()).into(),
297            ..Default::default()
298        };
299        self.rpc_client
300            .invoke_proto_method::<_, UnsubscribeResponse>(
301                usubscription_uri(RESOURCE_ID_UNSUBSCRIBE),
302                Self::default_call_options(),
303                unsubscribe_request,
304            )
305            .await
306            .map(|_response| ())
307            .map_err(UStatus::from)
308    }
309
310    async fn fetch_subscriptions_by_topic(
311        &self,
312        topic: &UUri,
313    ) -> Result<Vec<SubscriptionInfo>, UStatus> {
314        let fetch_subscriptions_request =
315            crate::up_core_api::usubscription::FetchSubscriptionsRequest {
316                request: Some(
317                    crate::up_core_api::usubscription::fetch_subscriptions_request::Request::Topic(
318                        topic.into(),
319                    ),
320                ),
321                ..Default::default()
322            };
323        self.fetch_subscriptions(fetch_subscriptions_request).await
324    }
325
326    async fn fetch_subscriptions_by_subscriber(
327        &self,
328        subscriber: &UUri,
329    ) -> Result<Vec<SubscriptionInfo>, UStatus> {
330        let subscriber_info = crate::up_core_api::usubscription::SubscriberInfo {
331            uri: Some(subscriber.into()).into(),
332            ..Default::default()
333        };
334        let fetch_subscriptions_request =
335            crate::up_core_api::usubscription::FetchSubscriptionsRequest {
336                request: Some(
337                    crate::up_core_api::usubscription::fetch_subscriptions_request::Request::Subscriber(subscriber_info),
338                ),
339                ..Default::default()
340            };
341        self.fetch_subscriptions(fetch_subscriptions_request).await
342    }
343
344    async fn register_for_notifications(&self, topic: &UUri) -> Result<(), UStatus> {
345        let notifications_register_request =
346            crate::up_core_api::usubscription::NotificationsRequest {
347                topic: Some(topic.into()).into(),
348                ..Default::default()
349            };
350        self.rpc_client
351            .invoke_proto_method::<_, NotificationsResponse>(
352                usubscription_uri(RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS),
353                Self::default_call_options(),
354                notifications_register_request,
355            )
356            .await
357            .map(|_response| ())
358            .map_err(UStatus::from)
359    }
360
361    async fn unregister_for_notifications(&self, topic: &UUri) -> Result<(), UStatus> {
362        let notifications_unregister_request =
363            crate::up_core_api::usubscription::NotificationsRequest {
364                topic: Some(topic.into()).into(),
365                ..Default::default()
366            };
367        self.rpc_client
368            .invoke_proto_method::<_, NotificationsResponse>(
369                usubscription_uri(RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS),
370                Self::default_call_options(),
371                notifications_unregister_request,
372            )
373            .await
374            .map(|_response| ())
375            .map_err(UStatus::from)
376    }
377
378    async fn fetch_subscribers(&self, topic: &UUri) -> Result<Vec<UUri>, UStatus> {
379        let fetch_subscribers_request =
380            crate::up_core_api::usubscription::FetchSubscribersRequest {
381                topic: Some(topic.into()).into(),
382                ..Default::default()
383            };
384        let response = self
385            .rpc_client
386            .invoke_proto_method::<_, FetchSubscribersResponse>(
387                usubscription_uri(RESOURCE_ID_FETCH_SUBSCRIBERS),
388                Self::default_call_options(),
389                fetch_subscribers_request,
390            )
391            .await?;
392        let mut result = vec![];
393        for subscriber_info in &response.subscribers {
394            let uri = subscriber_info
395                .uri
396                .as_ref()
397                .ok_or_else(|| {
398                    UStatus::fail_with_code(
399                        UCode::InvalidArgument,
400                        "uSubscription returned invalid response: missing subscriber URI",
401                    )
402                })
403                .and_then(|uri_proto| {
404                    UUri::try_from(uri_proto).map_err(|_e| {
405                        UStatus::fail_with_code(
406                            UCode::InvalidArgument,
407                            "uSubscription returned invalid response: invalid subscriber URI",
408                        )
409                    })
410                })?;
411            result.push(uri);
412        }
413        Ok(result)
414    }
415
416    async fn reset(&self, reason: ResetReason, message: Option<String>) -> Result<(), UStatus> {
417        let reset_request = crate::up_core_api::usubscription::ResetRequest {
418            reason: Some(crate::up_core_api::usubscription::reset_request::Reason {
419                code: crate::up_core_api::usubscription::reset_request::reason::Code::from(reason)
420                    .into(),
421                message,
422                ..Default::default()
423            })
424            .into(),
425            ..Default::default()
426        };
427        self.rpc_client
428            .invoke_proto_method::<_, ResetResponse>(
429                usubscription_uri(RESOURCE_ID_RESET),
430                Self::default_call_options(),
431                reset_request,
432            )
433            .await
434            .map(|_response| ())
435            .map_err(UStatus::from)
436    }
437}
438
439#[cfg(test)]
440mod tests {
441    use mockall::Sequence;
442
443    use super::*;
444    use crate::{
445        communication::{MockRpcClient, UPayload},
446        up_core_api::usubscription::{
447            fetch_subscriptions_request::Request, FetchSubscribersRequest, NotificationsRequest,
448            ResetRequest,
449        },
450        UCode, UUri,
451    };
452    use std::sync::Arc;
453
454    #[test]
455    fn test_unix_epoch_millis_as_protobuf_timestamp() {
456        assert!(
457            unix_epoch_millis_as_protobuf_timestamp(Some(1_000)).is_ok_and(|ts| {
458                ts == Some(Timestamp {
459                    seconds: 1,
460                    nanos: 0,
461                    ..Default::default()
462                })
463            })
464        );
465
466        assert!(
467            unix_epoch_millis_as_protobuf_timestamp(Some(1_234)).is_ok_and(|ts| {
468                ts == Some(Timestamp {
469                    seconds: 1,
470                    nanos: 234_000_000,
471                    ..Default::default()
472                })
473            })
474        );
475
476        assert!(unix_epoch_millis_as_protobuf_timestamp(Some(u64::MAX)).is_ok());
477        assert!(unix_epoch_millis_as_protobuf_timestamp(None).is_ok_and(|ts| ts.is_none()));
478    }
479
480    #[test_case::test_case(10, 234_000_000 => matches Ok(Some(10_234)); "succeeds for valid timestamp")]
481    #[test_case::test_case(-10, 234_000_000 => matches Err(UStatus {..}); "fails for negative seconds")]
482    #[test_case::test_case(10, -1 => matches Err(UStatus {..}); "fails for nanos exceeding lower bound")]
483    #[test_case::test_case(10, 1_000_000_000 => matches Err(UStatus {..}); "fails for nanos exeeding upper bound")]
484    fn test_protobuf_timestamp_as_unix_epoch_milliseconds(
485        seconds: i64,
486        nanos: i32,
487    ) -> Result<Option<u64>, UStatus> {
488        let timestamp = Timestamp {
489            seconds,
490            nanos,
491            ..Default::default()
492        };
493        protobuf_timestamp_as_unix_epoch_milliseconds(Some(&timestamp))
494    }
495
496    #[tokio::test]
497    async fn test_subscribe_invokes_rpc_client() {
498        let topic = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap();
499        let expected_request = SubscriptionRequest {
500            topic: Some((&topic).into()).into(),
501            ..Default::default()
502        };
503        let mut rpc_client = MockRpcClient::new();
504        let mut seq = Sequence::new();
505        rpc_client
506            .expect_invoke_method()
507            .once()
508            .in_sequence(&mut seq)
509            .withf(|method, _options, payload| {
510                method == &usubscription_uri(RESOURCE_ID_SUBSCRIBE) && payload.is_some()
511            })
512            .return_const(Err(crate::communication::ServiceInvocationError::Internal(
513                "internal error".to_string(),
514            )));
515        rpc_client
516            .expect_invoke_method()
517            .once()
518            .in_sequence(&mut seq)
519            .withf(move |method, _options, payload| {
520                let request = payload
521                    .to_owned()
522                    .unwrap()
523                    .extract_protobuf::<SubscriptionRequest>()
524                    .unwrap();
525                request == expected_request && method == &usubscription_uri(RESOURCE_ID_SUBSCRIBE)
526            })
527            .returning(move |_method, _options, _payload| {
528                let response = SubscriptionResponse {
529                    status: Some(crate::up_core_api::usubscription::SubscriptionStatus {
530                        state: crate::up_core_api::usubscription::subscription_status::State::SUBSCRIBED
531                            .into(),
532                        ..Default::default()
533                    }).into(),
534                    ..Default::default()
535                };
536                Ok(Some(UPayload::try_from_protobuf(response).unwrap()))
537            });
538
539        let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client));
540
541        assert!(usubscription_client
542            .subscribe(&topic, None, None)
543            .await
544            .is_err_and(|e| e.code() == UCode::Internal));
545        assert!(usubscription_client
546            .subscribe(&topic, None, None)
547            .await
548            .is_ok());
549    }
550
551    #[tokio::test]
552    async fn test_unsubscribe_invokes_rpc_client() {
553        let topic = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap();
554        let expected_request = UnsubscribeRequest {
555            topic: Some((&topic).into()).into(),
556            ..Default::default()
557        };
558        let mut rpc_client = MockRpcClient::new();
559        let mut seq = Sequence::new();
560        rpc_client
561            .expect_invoke_method()
562            .once()
563            .in_sequence(&mut seq)
564            .withf(|method, _options, payload| {
565                method == &usubscription_uri(RESOURCE_ID_UNSUBSCRIBE) && payload.is_some()
566            })
567            .return_const(Err(crate::communication::ServiceInvocationError::Internal(
568                "internal error".to_string(),
569            )));
570        rpc_client
571            .expect_invoke_method()
572            .once()
573            .in_sequence(&mut seq)
574            .withf(move |method, _options, payload| {
575                let request = payload
576                    .to_owned()
577                    .unwrap()
578                    .extract_protobuf::<UnsubscribeRequest>()
579                    .unwrap();
580                request == expected_request && method == &usubscription_uri(RESOURCE_ID_UNSUBSCRIBE)
581            })
582            .returning(move |_method, _options, _payload| {
583                let response = UnsubscribeResponse {
584                    ..Default::default()
585                };
586                Ok(Some(UPayload::try_from_protobuf(response).unwrap()))
587            });
588
589        let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client));
590
591        assert!(usubscription_client
592            .unsubscribe(&topic)
593            .await
594            .is_err_and(|e| e.code() == UCode::Internal));
595        assert!(usubscription_client.unsubscribe(&topic).await.is_ok());
596    }
597
598    #[tokio::test]
599    async fn test_fetch_subscriptions_invokes_rpc_client() {
600        let topic = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap();
601        let expected_request = FetchSubscriptionsRequest {
602            request: Some(Request::Topic((&topic).into())),
603            ..Default::default()
604        };
605        let mut rpc_client = MockRpcClient::new();
606        let mut seq = Sequence::new();
607        rpc_client
608            .expect_invoke_method()
609            .once()
610            .in_sequence(&mut seq)
611            .withf(|method, _options, payload| {
612                method == &usubscription_uri(RESOURCE_ID_FETCH_SUBSCRIPTIONS) && payload.is_some()
613            })
614            .return_const(Err(crate::communication::ServiceInvocationError::Internal(
615                "internal error".to_string(),
616            )));
617        rpc_client
618            .expect_invoke_method()
619            .once()
620            .in_sequence(&mut seq)
621            .withf(move |method, _options, payload| {
622                let request = payload
623                    .to_owned()
624                    .unwrap()
625                    .extract_protobuf::<FetchSubscriptionsRequest>()
626                    .unwrap();
627
628                request == expected_request
629                    && method == &usubscription_uri(RESOURCE_ID_FETCH_SUBSCRIPTIONS)
630            })
631            .returning(move |_method, _options, _payload| {
632                let response = FetchSubscriptionsResponse {
633                    ..Default::default()
634                };
635                Ok(Some(UPayload::try_from_protobuf(response).unwrap()))
636            });
637
638        let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client));
639
640        assert!(usubscription_client
641            .fetch_subscriptions_by_topic(&topic)
642            .await
643            .is_err_and(|e| e.code() == UCode::Internal));
644        assert!(usubscription_client
645            .fetch_subscriptions_by_topic(&topic)
646            .await
647            .is_ok());
648    }
649
650    #[tokio::test]
651    async fn test_fetch_subscribers_invokes_rpc_client() {
652        let topic = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap();
653        let expected_request = FetchSubscribersRequest {
654            topic: Some((&topic).into()).into(),
655            ..Default::default()
656        };
657        let mut rpc_client = MockRpcClient::new();
658        let mut seq = Sequence::new();
659        rpc_client
660            .expect_invoke_method()
661            .once()
662            .in_sequence(&mut seq)
663            .withf(|method, _options, payload| {
664                method == &usubscription_uri(RESOURCE_ID_FETCH_SUBSCRIBERS) && payload.is_some()
665            })
666            .return_const(Err(crate::communication::ServiceInvocationError::Internal(
667                "internal error".to_string(),
668            )));
669        rpc_client
670            .expect_invoke_method()
671            .once()
672            .in_sequence(&mut seq)
673            .withf(move |method, _options, payload| {
674                let request = payload
675                    .to_owned()
676                    .unwrap()
677                    .extract_protobuf::<FetchSubscribersRequest>()
678                    .unwrap();
679
680                request == expected_request
681                    && method == &usubscription_uri(RESOURCE_ID_FETCH_SUBSCRIBERS)
682            })
683            .returning(move |_method, _options, _payload| {
684                let response = FetchSubscribersResponse {
685                    ..Default::default()
686                };
687                Ok(Some(UPayload::try_from_protobuf(response).unwrap()))
688            });
689
690        let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client));
691
692        assert!(usubscription_client
693            .fetch_subscribers(&topic)
694            .await
695            .is_err_and(|e| e.code() == UCode::Internal));
696        assert!(usubscription_client.fetch_subscribers(&topic).await.is_ok());
697    }
698
699    #[tokio::test]
700    async fn test_register_for_notifications_invokes_rpc_client() {
701        let topic = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap();
702        let expected_request = NotificationsRequest {
703            topic: Some((&topic).into()).into(),
704            ..Default::default()
705        };
706        let mut rpc_client = MockRpcClient::new();
707        let mut seq = Sequence::new();
708        rpc_client
709            .expect_invoke_method()
710            .once()
711            .in_sequence(&mut seq)
712            .withf(|method, _options, payload| {
713                method == &usubscription_uri(RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS)
714                    && payload.is_some()
715            })
716            .return_const(Err(crate::communication::ServiceInvocationError::Internal(
717                "internal error".to_string(),
718            )));
719        rpc_client
720            .expect_invoke_method()
721            .once()
722            .in_sequence(&mut seq)
723            .withf(move |method, _options, payload| {
724                let request = payload
725                    .to_owned()
726                    .unwrap()
727                    .extract_protobuf::<NotificationsRequest>()
728                    .unwrap();
729
730                request == expected_request
731                    && method == &usubscription_uri(RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS)
732            })
733            .returning(move |_method, _options, _payload| {
734                let response = NotificationsResponse {
735                    ..Default::default()
736                };
737                Ok(Some(UPayload::try_from_protobuf(response).unwrap()))
738            });
739
740        let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client));
741
742        assert!(usubscription_client
743            .register_for_notifications(&topic)
744            .await
745            .is_err_and(|e| e.code() == UCode::Internal));
746        assert!(usubscription_client
747            .register_for_notifications(&topic)
748            .await
749            .is_ok());
750    }
751
752    #[tokio::test]
753    async fn test_unregister_for_notifications_invokes_rpc_client() {
754        let topic = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap();
755        let expected_request = NotificationsRequest {
756            topic: Some((&topic).into()).into(),
757            ..Default::default()
758        };
759        let mut rpc_client = MockRpcClient::new();
760        let mut seq = Sequence::new();
761        rpc_client
762            .expect_invoke_method()
763            .once()
764            .in_sequence(&mut seq)
765            .withf(|method, _options, payload| {
766                method == &usubscription_uri(RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS)
767                    && payload.is_some()
768            })
769            .return_const(Err(crate::communication::ServiceInvocationError::Internal(
770                "internal error".to_string(),
771            )));
772        rpc_client
773            .expect_invoke_method()
774            .once()
775            .in_sequence(&mut seq)
776            .withf(move |method, _options, payload| {
777                let request = payload
778                    .to_owned()
779                    .unwrap()
780                    .extract_protobuf::<NotificationsRequest>()
781                    .unwrap();
782
783                request == expected_request
784                    && method == &usubscription_uri(RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS)
785            })
786            .returning(move |_method, _options, _payload| {
787                let response = NotificationsResponse {
788                    ..Default::default()
789                };
790                Ok(Some(UPayload::try_from_protobuf(response).unwrap()))
791            });
792
793        let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client));
794
795        assert!(usubscription_client
796            .unregister_for_notifications(&topic)
797            .await
798            .is_err_and(|e| e.code() == UCode::Internal));
799        assert!(usubscription_client
800            .unregister_for_notifications(&topic)
801            .await
802            .is_ok());
803    }
804
805    #[tokio::test]
806    async fn test_reset_invokes_rpc_client() {
807        let expected_request = ResetRequest {
808            reason: Some(crate::up_core_api::usubscription::reset_request::Reason {
809                code: crate::up_core_api::usubscription::reset_request::reason::Code::UNSPECIFIED
810                    .into(),
811                ..Default::default()
812            })
813            .into(),
814            ..Default::default()
815        };
816        let mut rpc_client = MockRpcClient::new();
817        let mut seq = Sequence::new();
818        rpc_client
819            .expect_invoke_method()
820            .once()
821            .in_sequence(&mut seq)
822            .withf(|method, _options, payload| {
823                method == &usubscription_uri(RESOURCE_ID_RESET) && payload.is_some()
824            })
825            .return_const(Err(crate::communication::ServiceInvocationError::Internal(
826                "internal error".to_string(),
827            )));
828        rpc_client
829            .expect_invoke_method()
830            .once()
831            .in_sequence(&mut seq)
832            .withf(move |method, _options, payload| {
833                let request = payload
834                    .to_owned()
835                    .unwrap()
836                    .extract_protobuf::<ResetRequest>()
837                    .unwrap();
838
839                request == expected_request && method == &usubscription_uri(RESOURCE_ID_RESET)
840            })
841            .returning(move |_method, _options, _payload| {
842                let response = ResetResponse {
843                    ..Default::default()
844                };
845                Ok(Some(UPayload::try_from_protobuf(response).unwrap()))
846            });
847
848        let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client));
849
850        assert!(usubscription_client
851            .reset(ResetReason::Unspecified, None)
852            .await
853            .is_err_and(|e| e.code() == UCode::Internal));
854        assert!(usubscription_client
855            .reset(ResetReason::Unspecified, None)
856            .await
857            .is_ok());
858    }
859}