1use 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 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
200pub 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 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>, 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(×tamp))
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}