diff --git a/Cargo.lock b/Cargo.lock index 22858178..23550bb4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -427,6 +427,12 @@ dependencies = [ "serde", ] +[[package]] +name = "fnv" +version = "1.0.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3f9eec918d3f24069decb9af1554cad7c880e2da24a9afd88aca000531ab82c1" + [[package]] name = "fragile" version = "2.1.0" @@ -1207,9 +1213,9 @@ checksum = "8a7852d02fc848982e0c167ef163aaff9cd91dc640ba85e263cb1ce46fae51cd" [[package]] name = "serde" -version = "1.0.229" +version = "1.0.228" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4148590afebada386688f18773da617792bf2ef03ffc1e4cbd2b1d45b023e0ba" +checksum = "9a8e94ea7f378bd32cbbd37198a4a91436180c5bb472411e48b5ec2e2124ae9e" dependencies = [ "serde_core", "serde_derive", @@ -1217,29 +1223,29 @@ dependencies = [ [[package]] name = "serde_core" -version = "1.0.229" +version = "1.0.228" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "67dca2c9c51e58a4791a4b1ed58308b39c64224d349a935ab5039aa360942a48" +checksum = "41d385c7d4ca58e59fc732af25c3983b67ac852c1a25000afe1175de458b67ad" dependencies = [ "serde_derive", ] [[package]] name = "serde_derive" -version = "1.0.229" +version = "1.0.228" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e7a5d71263a5a7d47b41f6b3f06ba276f10cc18b0931f1799f710578e2309348" +checksum = "d540f220d3187173da220f885ab66608367b6574e925011a9353e4badda91d79" dependencies = [ "proc-macro2", "quote", - "syn 3.0.3", + "syn 2.0.119", ] [[package]] name = "serde_json" -version = "1.0.151" +version = "1.0.150" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c841b55ecdae098c80dcae9cf767f6f8a0c2cdb3416bbef72181df4d0fe73f14" +checksum = "e8014e44b4736ed0538adeecded0fce2a272f22dc9578a7eb6b2d9993c74cfb9" dependencies = [ "itoa", "memchr", @@ -1685,9 +1691,20 @@ dependencies = [ "thiserror 2.0.21", "tokio", "tracing", + "uriparse", "uuid-simd", ] +[[package]] +name = "uriparse" +version = "0.6.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0200d0fc04d809396c2ad43f3c95da3582a2556eba8d453c1087f4120ee352ff" +dependencies = [ + "fnv", + "lazy_static", +] + [[package]] name = "utf8parse" version = "0.2.2" diff --git a/Cargo.toml b/Cargo.toml index 269bffb8..9272e416 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -72,8 +72,15 @@ serde_json = { version = "1.0.150", optional = true } symphony = { version = "0.1.5", optional = true } thiserror = { version = "2.0.21" } tokio = { version = "1.53.1", default-features = false, optional = true } -tracing = { version = "0.1.44", default-features = false, features = ["log", "std"] } -uuid-simd = { version = "0.8.0", default-features = false, features = ["std", "detect"] } +tracing = { version = "0.1.44", default-features = false, features = [ + "log", + "std", +] } +uriparse = { version = "0.6.4" } +uuid-simd = { version = "0.8.0", default-features = false, features = [ + "std", + "detect", +] } [build-dependencies] protobuf-codegen = { version = "3.7.2" } @@ -163,6 +170,10 @@ required-features = ["up-l2-rpc-client", "up-l2-rpc-server", "util"] name = "symphony_target" required-features = ["symphony", "up-l2-rpc-client", "up-l2-rpc-server", "util"] +[[example]] +name = "usubscription_server" +required-features = ["usubscription", "up-l2-rpc-server", "util"] + [[test]] name = "tck_uuri" harness = false # allows Cucumber to print output instead of libtest diff --git a/build.rs b/build.rs index 97041be4..ba47fc3b 100644 --- a/build.rs +++ b/build.rs @@ -42,7 +42,7 @@ fn proto_api() -> Result<(), Box> { ), #[cfg(feature = "usubscription")] format!( - "{}uprotocol/core/usubscription/v3/usubscription.proto", + "{}uprotocol/core/usubscription/v4/usubscription.proto", UPROTOCOL_BASE_URI ), ]; diff --git a/examples/usubscription_server.rs b/examples/usubscription_server.rs new file mode 100644 index 00000000..a4c2b149 --- /dev/null +++ b/examples/usubscription_server.rs @@ -0,0 +1,151 @@ +/******************************************************************************** + * Copyright (c) 2026 Contributors to the Eclipse Foundation + * + * See the NOTICE file(s) distributed with this work for additional + * information regarding copyright ownership. + * + * This program and the accompanying materials are made available under the + * terms of the Apache License Version 2.0 which is available at + * https://www.apache.org/licenses/LICENSE-2.0 + * + * SPDX-License-Identifier: Apache-2.0 + ********************************************************************************/ + +/*! +This example illustrates how a service provider implements the server side of the +uSubscription service. + +A [`RpcClientUSubscription`] drives the example by issuing `subscribe`/`unsubscribe` +calls over an in-process [`LocalTransport`]. + */ + +use std::{ + collections::HashSet, + sync::{Arc, RwLock}, +}; + +use up_rust::{ + communication::{ + InMemoryRpcClient, InMemoryRpcServer, RequestHandler, RpcServer, ServiceInvocationError, + SubscriptionStatus, UPayload, + }, + core::usubscription::{ + extract_usubscription_request, RpcClientUSubscription, SubscribeResponse, USubscription, + USubscriptionRequest, USubscriptionResponse, RESOURCE_ID_SUBSCRIBE, + RESOURCE_ID_UNSUBSCRIBE, USUBSCRIPTION_TYPE_ID, USUBSCRIPTION_VERSION_MAJOR, + }, + local_transport::LocalTransport, + StaticUriProvider, UAttributes, UUri, +}; + +/// A minimal, in-memory implementation of the uSubscription service. +/// +/// It keeps track of the `(subscriber, topic)` pairs it has seen so far. A real +/// implementation would persist this state and enforce the uSubscription state machine. +#[derive(Default)] +struct MyUSubscriptionService { + subscriptions: RwLock>, +} + +#[async_trait::async_trait] +impl RequestHandler for MyUSubscriptionService { + async fn handle_request( + &self, + resource_id: u16, + message_attributes: &UAttributes, + request_payload: Option, + ) -> Result, ServiceInvocationError> { + // Decode the payload of uSubscription operations into a typed request. Malformed or + // unsupported requests result in a `ServiceInvocationError`. + let request = extract_usubscription_request(resource_id, request_payload)?; + + #[allow(clippy::wildcard_enum_match_arm)] + match request { + USubscriptionRequest::Subscribe(req) => { + println!( + "SUBSCRIBE subscriber={}, topic={}, expiration={:?}, sample_period={:?}", + message_attributes.source(), + req.topic, + req.expiration, + req.sample_period + ); + + self.subscriptions + .write() + .expect("couldn't acquire subscription ledger write lock") + .insert(( + message_attributes.source().to_string(), + req.topic.to_string(), + )); + + USubscriptionResponse::Subscribe(SubscribeResponse { + topic: req.topic, + status: SubscriptionStatus::Subscribed, + }) + } + USubscriptionRequest::Unsubscribe(req) => { + println!( + "UNSUBSCRIBE subscriber={}, topic={}", + message_attributes.source(), + req.topic + ); + self.subscriptions + .write() + .expect("couldn't acquire subscription ledger write lock") + .remove(&( + message_attributes.source().to_string(), + req.topic.to_string(), + )); + USubscriptionResponse::Unsubscribe(()) + } + // `USubscriptionRequest` is `#[non_exhaustive]`, so new operations can be + // added in future releases without breaking this code. + other => { + return Err(ServiceInvocationError::Unimplemented(format!( + "operation not supported by this service: {other:?}" + ))) + } + }; + + Ok(None) + } +} + +#[tokio::main] +pub async fn main() -> Result<(), Box> { + // Using the LocalTransport lets us run the service and its client in the same + // process. A real deployment would use a networked transport (e.g. MQTT5 or Zenoh). + let transport = Arc::new(LocalTransport::default()); + let service_uri_provider = Arc::new(StaticUriProvider::new( + "", + USUBSCRIPTION_TYPE_ID as u32, + USUBSCRIPTION_VERSION_MAJOR, + )?); + + // Stand up the server and register our service for the subscribe/unsubscribe methods. + let rpc_server = InMemoryRpcServer::new(transport.clone(), service_uri_provider); + let service = Arc::new(MyUSubscriptionService::default()); + rpc_server + .register_endpoint(None, RESOURCE_ID_SUBSCRIBE, service.clone()) + .await?; + rpc_server + .register_endpoint(None, RESOURCE_ID_UNSUBSCRIBE, service.clone()) + .await?; + + // Now act as a client that subscribes to and unsubscribes from a topic. The + // subscriber's identity (its source address) is taken from this URI provider and + // ends up in `SubscribeRequest::subscriber` on the server side. + let client_uri_provider = Arc::new(StaticUriProvider::new("my-vehicle", 0xABCD, 0x01)?); + let rpc_client = Arc::new(InMemoryRpcClient::new(transport, client_uri_provider).await?); + let usubscription = RpcClientUSubscription::new(rpc_client, None); + + let topic = UUri::try_from_parts("my-vehicle", 0x0000_800A, 0x01, 0x8001)?; + + let status = usubscription.subscribe(&topic, None, None).await?; + println!("client: subscribe call returned status {status:?}"); + + usubscription.unsubscribe(&topic).await?; + println!("client: unsubscribe call completed"); + + Ok(()) +} diff --git a/src/communication.rs b/src/communication.rs index 479ae0b4..5161b6d4 100644 --- a/src/communication.rs +++ b/src/communication.rs @@ -99,44 +99,6 @@ pub(crate) fn build_message( } } -/// The current status of a client's subscription to a topic. -/// -/// The status goes through different stages during its lifecycle as defined in -/// [uProtocol Specification, section 3.3.5](https://github.com/eclipse-uprotocol/up-spec/blob/v1.6.0-alpha.7/up-l3/usubscription/v3/README.adoc#usubscription-states). -#[derive(Clone, Debug, PartialEq)] -#[repr(C)] -pub enum SubscriptionStatus { - Unsubscribed, - SubscribePending, - Subscribed, - UnsubscribePending, -} - -#[cfg(all(feature = "up-core-types", feature = "usubscription"))] -mod core_types_support { - use super::{SubscriptionStatus, UCode, UStatus}; - use crate::up_core_api::usubscription::subscription_status::State; - use crate::up_core_api::usubscription::SubscriptionStatus as SubscriptionStatusProto; - - impl TryFrom<&SubscriptionStatusProto> for SubscriptionStatus { - type Error = UStatus; - - fn try_from(status_proto: &SubscriptionStatusProto) -> Result { - let state = status_proto.state.enum_value(); - match state { - Ok(State::UNSUBSCRIBED) => Ok(SubscriptionStatus::Unsubscribed), - Ok(State::SUBSCRIBE_PENDING) => Ok(SubscriptionStatus::SubscribePending), - Ok(State::SUBSCRIBED) => Ok(SubscriptionStatus::Subscribed), - Ok(State::UNSUBSCRIBE_PENDING) => Ok(SubscriptionStatus::UnsubscribePending), - Err(v) => Err(UStatus::fail_with_code( - UCode::InvalidArgument, - format!("unknown subscription status {:?}", v), - )), - } - } - } -} - /// An error indicating a problem with registering or unregistering a message listener. #[derive(Clone, Debug, Error)] pub enum RegistrationError { @@ -442,3 +404,33 @@ impl UPayload { crate::umessage::deserialize_protobuf_bytes(&self.payload, &self.payload_format) } } + +/// The current status of a client's subscription to a topic. +/// +/// The status goes through different stages during its lifecycle as defined in +/// [uProtocol Specification, section 3.3.5](https://github.com/eclipse-uprotocol/up-spec/blob/main/up-l3/usubscription/v4/README.adoc#usubscription-states). +/// +/// SubscriptionStatus is defined here because it is used by the pubsub module SubscriptionChangeHandler trait, +/// which is enabled via communication.rs by the more generic 'up-l2-api' feature. +/// So to avoid having to enable/pull in all the core (protobuf) types and code for this single trait we put +/// SubscriptionStatus here, on the same level of module/feature genericity. +#[derive(Clone, Debug, PartialEq)] +#[repr(C)] +pub enum SubscriptionStatus { + Unsubscribed, + SubscribePending, + Subscribed, + UnsubscribePending, +} + +impl std::fmt::Display for SubscriptionStatus { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let status = match self { + SubscriptionStatus::Unsubscribed => "Unsubscribed", + SubscriptionStatus::SubscribePending => "SubscribePending", + SubscriptionStatus::Subscribed => "Subscribed", + SubscriptionStatus::UnsubscribePending => "UnsubscribePending", + }; + write!(f, "{status}") + } +} diff --git a/src/communication/in_memory_subscriber.rs b/src/communication/in_memory_subscriber.rs index fbd96436..b169d358 100644 --- a/src/communication/in_memory_subscriber.rs +++ b/src/communication/in_memory_subscriber.rs @@ -161,7 +161,7 @@ impl UListener for SubscriptionChangeListener { return; } let Ok(subscription_update) = - msg.extract_protobuf::() + msg.extract_protobuf::() else { debug!("ignoring notification that does not contain subscription update"); return; @@ -227,7 +227,7 @@ impl .await .map(Arc::new)?; let usubscription_client = Arc::new( - crate::core::usubscription::RpcClientUSubscription::new(rpc_client), + crate::core::usubscription::RpcClientUSubscription::new(rpc_client, None), ); let notifier = Arc::new(crate::communication::SimpleNotifier::new( transport.clone(), @@ -262,7 +262,10 @@ impl InMemorySubscriber { }); notifier .start_listening( - &usubscription::usubscription_uri(usubscription::RESOURCE_ID_SUBSCRIPTION_CHANGE), + &usubscription::usubscription_uri( + None, + usubscription::RESOURCE_ID_SUBSCRIPTION_CHANGE, + ), subscription_change_listener.clone(), ) .await?; @@ -284,7 +287,10 @@ impl InMemorySubscriber { pub async fn stop(&self) -> Result<(), RegistrationError> { self.notifier .stop_listening( - &usubscription::usubscription_uri(usubscription::RESOURCE_ID_SUBSCRIPTION_CHANGE), + &usubscription::usubscription_uri( + None, + usubscription::RESOURCE_ID_SUBSCRIPTION_CHANGE, + ), self.subscription_change_listener.clone(), ) .await @@ -416,7 +422,6 @@ mod tests { use crate::{ communication::{notification::MockNotifier, pubsub::MockSubscriptionChangeHandler}, - up_core_api::usubscription::Update, utransport::{MockTransport, MockUListener}, UCode, UMessageBuilder, UMessageType, UStatus, UUri, UUID, }; @@ -901,14 +906,11 @@ mod tests { } fn status_update_without_topic() -> UMessage { - let status = crate::up_core_api::usubscription::SubscriptionStatus { - state: crate::up_core_api::usubscription::subscription_status::State::SUBSCRIBED.into(), - ..Default::default() - }; - let update = Update { - status: Some(status).into(), + let update = crate::up_core_api::usubscription::Subscription { + status: crate::up_core_api::usubscription::SubscriptionStatus::STATUS_SUBSCRIBED.into(), ..Default::default() }; + UMessageBuilder::notification( UUri::try_from_parts("other", 0x1a9a, 0x01, 0x8100) .expect("should have been able to create origin URI"), @@ -923,7 +925,8 @@ mod tests { let topic = UUri::try_from_parts("other", 0x1a9a, 0x01, 0x8100) .expect("should have been able to create topic URI"); let proto_topic = crate::up_core_api::uri::UUri::from(&topic); - let update = Update { + + let update = crate::up_core_api::usubscription::Subscription { topic: Some(proto_topic).into(), ..Default::default() }; @@ -962,20 +965,15 @@ mod tests { let subscriber = UUri::try_from_parts("local", 0x2000, 0x01, 0x0000) .expect("should have been able to create subscriber URI"); let topic = UUri::try_from_parts("other", 0x1a9a, 0x01, 0x8100).unwrap(); - let status_proto = crate::up_core_api::usubscription::SubscriptionStatus { - state: crate::up_core_api::usubscription::subscription_status::State::SUBSCRIBED.into(), - ..Default::default() - }; - let update_proto = Update { + let status_proto = crate::up_core_api::usubscription::SubscriptionStatus::STATUS_SUBSCRIBED; + + let update_proto = crate::up_core_api::usubscription::Subscription { topic: Some((&topic).into()).into(), - subscriber: Some(crate::up_core_api::usubscription::SubscriberInfo { - uri: Some((&subscriber).into()).into(), - ..Default::default() - }) - .into(), - status: Some(status_proto.clone()).into(), + subscriber: Some((&subscriber).into()).into(), + status: status_proto.into(), ..Default::default() }; + let subscription_change_notification = UMessageBuilder::notification( UUri::try_from_parts("local", 0x0000, 0x01, 0x8000) .expect("should have been able to create uSubscription service URI"), @@ -993,7 +991,7 @@ mod tests { topic == &expected_topic && *updated_status == SubscriptionStatus::try_from(&status_proto) - .expect("should have been able to convert status proto") + .expect("should have received a valid subscription status") }) .return_const(()); diff --git a/src/core/usubscription.rs b/src/core/usubscription.rs index 82f6a396..bd50666b 100644 --- a/src/core/usubscription.rs +++ b/src/core/usubscription.rs @@ -14,56 +14,64 @@ use async_trait::async_trait; #[cfg(test)] use mockall::automock; +use std::time::{Duration, SystemTime}; -use crate::{communication::SubscriptionStatus, UStatus, UUri}; +use crate::{communication::SubscriptionStatus, UCode, UStatus, UUri}; -#[cfg(all(feature = "up-l2-rpc-client", feature = "up-core-types"))] mod usubscription_client; -#[cfg(all(feature = "up-l2-rpc-client", feature = "up-core-types"))] pub use usubscription_client::RpcClientUSubscription; +mod usubscription_server; +pub use usubscription_server::{ + extract_usubscription_request, pack_usubscription_response, USubscriptionRequest, + USubscriptionResponse, +}; +mod usubscription_proto; -/// The uEntity (type) identifier of the uSubscription service. -pub const USUBSCRIPTION_TYPE_ID: u32 = 0x0000_0000; +// [impl->req~usubscription-uentity_id~1] +pub const USUBSCRIPTION_TYPE_ID: u16 = 0x0000_0000; /// The (latest) major version of the uSubscription service. -pub const USUBSCRIPTION_VERSION_MAJOR: u8 = 0x03; -/// The resource identifier of uSubscription's _subscribe_ operation. +pub const USUBSCRIPTION_VERSION_MAJOR: u8 = 0x04; +// [impl->req~usubscription-subscribe-method_id~1] pub const RESOURCE_ID_SUBSCRIBE: u16 = 0x0001; -/// The resource identifier of uSubscription's _unsubscribe_ operation. +// [impl->req~usubscription-unsubscribe-method_id~1] pub const RESOURCE_ID_UNSUBSCRIBE: u16 = 0x0002; -/// The resource identifier of uSubscription's _fetch subscriptions_ operation. +// [impl->req~usubscription-fetch-subscriptions-method_id~1] pub const RESOURCE_ID_FETCH_SUBSCRIPTIONS: u16 = 0x0003; -/// The resource identifier of uSubscription's _register for notifications_ operation. -pub const RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS: u16 = 0x0006; -/// The resource identifier of uSubscription's _unregister for notifications_ operation. -pub const RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS: u16 = 0x0007; -/// The resource identifier of uSubscription's _fetch subscribers_ operation. -pub const RESOURCE_ID_FETCH_SUBSCRIBERS: u16 = 0x0008; -/// The resource identifier of uSubscription's _reset_ operation. -pub const RESOURCE_ID_RESET: u16 = 0x0009; +// [impl->req~usubscription-register-notifications-method_id~1] +pub const RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS: u16 = 0x0004; +// [impl->req~usubscription-unregister-notifications-method_id~1] +pub const RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS: u16 = 0x0005; +// [impl->req~usubscription-reset-method_id~1] +pub const RESOURCE_ID_RESET: u16 = 0x0006; -/// The resource identifier of uSubscription's _subscription change_ topic. +// [impl->req~usubscription-change-notification-resource~1] pub const RESOURCE_ID_SUBSCRIPTION_CHANGE: u16 = 0x8000; +/// Information about a client-topic subscription. +/// +/// This struct represents the subscription metadata maintained by the +/// uSubscription service, including the topic, subscriber, current +/// [`SubscriptionStatus`], and optional delivery constraints. #[derive(Clone, Debug, PartialEq)] #[repr(C)] pub struct SubscriptionInfo { topic: UUri, subscriber: UUri, status: SubscriptionStatus, - expiration: Option, - min_sample_period: Option, + expiration: Option, + min_sample_period: Option, } impl SubscriptionInfo { - /// Creates a new info object. + /// Creates a new subscription info object. /// /// # Arguments /// * `topic` - The topic of the subscription. /// * `subscriber` - The uEntity that has established the subscription. /// * `status` - The status of the subscription. - /// * `expiration` - The point in time at which the subscription expires (milliseconds since Unix epoch). + /// * `expiration` - The point in time at which the subscription expires. /// If not specified, the subscription is valid until explicitly unsubscribed. - /// * `min_sample_period` - The minimum duration (in seconds) between two events that should be maintained + /// * `min_sample_period` - The minimum duration between two events that should be maintained /// for remote only topics. Device dispatchers (i.e. streamers) use this attribute to reduce the publication /// rates of events sent between devices. /// This attribute is commonly used for mobile/cloud components subscribing to vehicle topics that are published @@ -75,8 +83,8 @@ impl SubscriptionInfo { topic: UUri, subscriber: UUri, status: SubscriptionStatus, - expiration: Option, - min_sample_period: Option, + expiration: Option, + min_sample_period: Option, ) -> Self { Self { topic, @@ -103,12 +111,12 @@ impl SubscriptionInfo { } #[must_use] - pub fn expiration(&self) -> &Option { + pub fn expiration(&self) -> &Option { &self.expiration } #[must_use] - pub fn min_sample_period(&self) -> &Option { + pub fn min_sample_period(&self) -> &Option { &self.min_sample_period } @@ -136,30 +144,138 @@ impl SubscriptionInfo { } } -/// Potential reasons for resetting the uSubscription service. -#[derive(Debug, PartialEq)] -#[repr(C)] -pub enum ResetReason { - Unspecified, - FactoryReset, - CorruptedData, +/// A request to subscribe to a topic (client/subscriber UUri from the message UAttributes envelope). +#[derive(Clone, Debug, PartialEq)] +pub struct SubscribeRequest { + /// The topic to subscribe to. + pub topic: UUri, + /// The point in time at which the subscription expires. + pub expiration: Option, + /// The minimum duration between two events (before they should be forwarded by a UStreamer). + pub sample_period: Option, +} + +impl SubscribeRequest { + pub fn new( + topic: &UUri, + expiration: Option, + sample_period: Option, + ) -> Result { + Self::validate_topic(topic)?; + Self::validate_expiration(&expiration)?; + Self::validate_sample_period(&sample_period)?; + + Ok(SubscribeRequest { + topic: topic.to_owned(), + expiration, + sample_period, + }) + } + + fn validate_topic(topic: &UUri) -> Result<(), UStatus> { + topic.verify_event().map_err(|e| { + UStatus::fail_with_code( + UCode::InvalidArgument, + format!("topic URI is not a valid event source: {e}"), + ) + })?; + + Ok(()) + } + + fn validate_expiration(expiration: &Option) -> Result<(), UStatus> { + if let Some(expiration) = expiration { + if expiration <= &SystemTime::now() { + return Err(UStatus::fail_with_code( + UCode::InvalidArgument, + "expiration must be a timestamp in the future", + )); + } + } + Ok(()) + } + + fn validate_sample_period(sample_period: &Option) -> Result<(), UStatus> { + if let Some(sample_period) = sample_period { + if sample_period.is_zero() { + return Err(UStatus::fail_with_code( + UCode::InvalidArgument, + "sample_period must be greater than zero", + )); + } + } + Ok(()) + } + /// Verifies that the given subscribe request's options comply with the USubscription specification. + /// + /// # Arguments + /// + /// * `subscribe_request` - The subscribe request to verify. + /// + /// # Errors + /// + /// Returns an error if the request's topic is not a valid event source, if `expiration` is set + /// to a timestamp that is not in the future, or if `sample_period` is set to a duration of zero. + pub fn validate(&self) -> Result<(), UStatus> { + Self::validate_topic(&self.topic)?; + Self::validate_expiration(&self.expiration)?; + Self::validate_sample_period(&self.sample_period)?; + Ok(()) + } +} + +/// The response to a [`SubscribeRequest`]. +#[derive(Clone, Debug, PartialEq)] +pub struct SubscribeResponse { + /// The topic the subscription refers to. + pub topic: UUri, + /// The resulting status of the subscription. + pub status: SubscriptionStatus, +} + +/// A request to unsubscribe from a topic (client/unsubscriber UUri from the message UAttributes envelope). +#[derive(Clone, Debug, PartialEq)] +pub struct UnsubscribeRequest { + /// The topic to unsubscribe from. + pub topic: UUri, } -/// Gets a UUri referring to one of the local uSubscription service's resources. +/// A request to fetch subscription information. +#[derive(Clone, Debug, PartialEq)] +pub struct FetchSubscriptionsRequest { + /// The topic filter to fetch subscription information for. + pub topic_filter: Option, + /// The subscriber filter to fetch subscription information for. + pub subscriber_filter: Option, +} + +/// The response to a [`FetchSubscriptionsRequest`]. +#[derive(Clone, Debug, PartialEq)] +pub struct FetchSubscriptionsResponse { + /// The topic the subscription refers to. + pub subscriptions: Vec, +} + +/// Gets a UUri referring to a uSubscription service's resource. +/// +/// # Arguments +/// +/// * `authority` - (Optional) authority of USubscription service to address - will default to local authority ("") if None. +/// * `resource_id` - Desired USubscription endpoint resource ID. /// /// # Examples /// /// ```rust /// use up_rust::core::usubscription; /// -/// let uuri = usubscription::usubscription_uri(usubscription::RESOURCE_ID_SUBSCRIBE); +/// let uuri = usubscription::usubscription_uri(None, usubscription::RESOURCE_ID_SUBSCRIBE); /// assert_eq!(uuri.resource_id(), 0x0001); /// ``` #[must_use] -pub fn usubscription_uri(resource_id: u16) -> UUri { +pub fn usubscription_uri(authority: Option<&str>, resource_id: u16) -> UUri { UUri::try_from_parts( - "", - USUBSCRIPTION_TYPE_ID, + authority.unwrap_or(""), + USUBSCRIPTION_TYPE_ID as u32, USUBSCRIPTION_VERSION_MAJOR, resource_id, ) @@ -168,15 +284,8 @@ pub fn usubscription_uri(resource_id: u16) -> UUri { /// The uProtocol Application Layer client interface to the uSubscription service. /// -/// Please refer to the [uSubscription service specification](https://github.com/eclipse-uprotocol/up-spec/blob/main/up-l3/usubscription/v3/README.adoc) +/// Please refer to the [uSubscription service specification](https://github.com/eclipse-uprotocol/up-spec/blob/main/up-l3/usubscription/v4/README.adoc) /// for details. -/// -/// **Note** that in contrast to the uSubscription service specification, the functions defined in this trait only -/// support commonly used input and output parameters of the operations defined in the specification. This is mainly -/// due to the fact, that for many of the other parameters defined in the specification, it is not entirely clear if -/// and how they should be used in practice. The next version of the uSubscription service specification will -/// include a more detailed description of the operations and their parameters, which will then be reflected in the -/// next version of this trait. #[cfg_attr(test, automock)] #[async_trait] pub trait USubscription: Send + Sync { @@ -185,15 +294,23 @@ pub trait USubscription: Send + Sync { /// # Parameters /// /// * `topic` - The topic to subscribe to. - /// * `expiration` - The point in time at which the subscription expires (milliseconds since Unix epoch). + /// * `expiration` - The point in time at which the subscription expires. /// If not specified, the subscription is valid until explicitly unsubscribed. - /// * `min_sample_period` - The minimum duration (in seconds) between two events that should be maintained + /// If expiration time is set in the past, no subscription is recorded and SubscriptionStatus::Unsubscribed + /// is returned. + /// When called for a topic that the client is already subscribed to and where the expiration field value differs + /// from the current expiration time, the expiration field of the subscription is updated with the new value. + /// When called for a topic that the client is already subscribed to where the expiration field value differs + /// from the current expiration time and lies in the past, the subscription is unregistered and + /// SubscriptionStatus::Unsubscribed is returned. + /// * `min_sample_period` - The minimum duration between two events that should be maintained /// for remote only topics. Device dispatchers (i.e. streamers) use this attribute to reduce the /// publication rates of events sent between devices. /// This attribute is commonly used for mobile/cloud components subscribing to vehicle topics that are published /// at a high rate. If the desired sampling period set by the subscriber is lower than the original publisher's /// publication period, the attribute is ignored. /// If not specified, the sampling period is set by the publisher. + /// Durations used in `min_sample_period` will be clamped to [0; u32::MAX] milliseconds. /// /// # Returns /// @@ -201,8 +318,8 @@ pub trait USubscription: Send + Sync { async fn subscribe( &self, topic: &UUri, - expiration: Option, - min_sample_period: Option, + expiration: Option, + min_sample_period: Option, ) -> Result; /// Unsubscribes this client from a topic. @@ -216,87 +333,44 @@ pub trait USubscription: Send + Sync { /// Returns an error if the attempt to unsubscribe has failed. async fn unsubscribe(&self, topic: &UUri) -> Result<(), UStatus>; - /// Gets all (currently) active subscriptions for a given topic. - /// - /// # Parameters - /// - /// * `topic` - The topic to fetch subscriptions for. - /// - /// # Errors + /// Gets details about subscriptions that are currently tracked by uSubscription service. Clients + /// can provide UURIs to filter subscriptions by topic and/or subscriber. Filter UURIs may contain wildcards. /// - /// Returns an error if the attempt to retrieve the subscriptions has failed. - async fn fetch_subscriptions_by_topic( - &self, - topic: &UUri, - ) -> Result, UStatus>; - - /// Gets a uEntity's (currently) active subscriptions. + /// Topic and subscription filters can be provided individually or in combination; uSubscription service will + /// only return subscription information that match all of the provided filter criteria. /// /// # Parameters /// - /// * `subscriber` - The uEntity to get the subscriptions for. + /// * `topic_filter` - Only return subscriptions where the topic is matched by this filter. + /// * `subscriber_filter` - Only return subscriptions where the subscriber is matched by this filter. /// /// # Errors /// /// Returns an error if the attempt to retrieve the subscriptions has failed. - async fn fetch_subscriptions_by_subscriber( + async fn fetch_subscriptions( &self, - subscriber: &UUri, + topic_filter: Option, + subscriber_filter: Option, ) -> Result, UStatus>; - /// Registers this client for notifications about changes to the subscription status for a given topic. - /// - /// # Parameters - /// - /// * `topic` - The topic to receive changes to subscription status for. + /// Registers this client for notifications about changes of subscriptions managed by uSubscription service. /// /// # Errors /// /// Returns an error if the attempt to register for notifications has failed. - async fn register_for_notifications(&self, topic: &UUri) -> Result<(), UStatus>; + async fn register_for_notifications(&self) -> Result<(), UStatus>; - /// Unregisters this client from notifications about changes to the subscription status for a given topic. - /// - /// # Parameters - /// - /// * `topic` - The topic to no longer receive changes to subscription status for. + /// Unregisters this client from notifications about changes of subscriptions managed by uSubscription service. /// /// # Errors /// /// Returns an error if the attempt to unregister from notifications has failed. - async fn unregister_for_notifications(&self, topic: &UUri) -> Result<(), UStatus>; - - /// Fetches a list of subscribers that are currently subscribed to a given topic. - /// - /// # Parameters - /// - /// * `topic` - The topic to fetch subscriptions for. - /// - /// # Returns - /// - /// A list of URIs representing the uEntities that are subscribed to the given topic. - /// - /// # Errors - /// - /// Returns an error if the attempt to fetch subscribers has failed. - async fn fetch_subscribers(&self, topic: &UUri) -> Result, UStatus>; + async fn unregister_for_notifications(&self) -> Result<(), UStatus>; /// Flushes all stored subscription information, including any persistently stored subscriptions. /// - /// # Parameters - /// - /// * `reason` - The reason for the reset. - /// * `message` - An optional human-readable message providing additional context about the reset. - /// * `before` - An optional timestamp (milliseconds since Unix epoch). All subscriptions created before - /// this timestamp will be removed. - /// /// # Errors /// /// Returns an error if the attempt to reset has failed. - async fn reset( - &self, - reason: ResetReason, - message: Option, - before: Option, - ) -> Result<(), UStatus>; + async fn reset(&self) -> Result<(), UStatus>; } diff --git a/src/core/usubscription/usubscription_client.rs b/src/core/usubscription/usubscription_client.rs index a03a672e..89558faa 100644 --- a/src/core/usubscription/usubscription_client.rs +++ b/src/core/usubscription/usubscription_client.rs @@ -12,196 +12,31 @@ ********************************************************************************/ use std::sync::Arc; +use std::time::{Duration, SystemTime}; use async_trait::async_trait; -use protobuf::well_known_types::timestamp::Timestamp; use crate::{ - communication::{CallOptions, RpcClient, SubscriptionStatus}, + communication::{CallOptions, RpcClient}, core::usubscription::{ - usubscription_uri, ResetReason, SubscriptionInfo, USubscription, - RESOURCE_ID_FETCH_SUBSCRIBERS, RESOURCE_ID_FETCH_SUBSCRIPTIONS, - RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS, RESOURCE_ID_RESET, RESOURCE_ID_SUBSCRIBE, - RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS, RESOURCE_ID_UNSUBSCRIBE, + usubscription_uri, FetchSubscriptionsRequest, FetchSubscriptionsResponse, SubscribeRequest, + SubscribeResponse, SubscriptionInfo, SubscriptionStatus, USubscription, UnsubscribeRequest, + RESOURCE_ID_FETCH_SUBSCRIPTIONS, RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS, RESOURCE_ID_RESET, + RESOURCE_ID_SUBSCRIBE, RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS, RESOURCE_ID_UNSUBSCRIBE, }, up_core_api::usubscription::{ - FetchSubscribersResponse, FetchSubscriptionsRequest, FetchSubscriptionsResponse, - NotificationsResponse, ResetResponse, SubscribeAttributes, Subscription, - SubscriptionRequest, SubscriptionResponse, UnsubscribeRequest, UnsubscribeResponse, Update, + NotificationsResponse as NotificationResponseProto, ResetResponse as ResetResponseProto, + UnsubscribeResponse as UnsubscribeResponseProto, }, - UCode, UStatus, UUri, + UStatus, UUri, }; -fn unix_epoch_millis_as_protobuf_timestamp( - millis: Option, -) -> Result, UStatus> { - if let Some(milliseconds) = millis { - // this will always yield a valid Timestamp as the maximum value of u64 (2^64 - 1) divided by 1000 - // is less than the maximum number of milliseconds that can be represented in an i64 (2^63 - 1) - let seconds = (milliseconds / 1000_u64) as i64; - let nanos = (milliseconds % 1000) - .checked_mul(1_000_000) - .ok_or_else(|| { - UStatus::fail_with_code(UCode::InvalidArgument, "timestamp out of range") - }) - .and_then(|s| { - i32::try_from(s).map_err(|_| { - UStatus::fail_with_code(UCode::InvalidArgument, "timestamp out of range") - }) - })?; - Ok(Some(Timestamp { - seconds, - nanos, - ..Default::default() - })) - } else { - Ok(None) - } -} - -fn protobuf_timestamp_as_unix_epoch_milliseconds( - ts: Option<&Timestamp>, -) -> Result, UStatus> { - if let Some(ts) = ts { - let err = || { - UStatus::fail_with_code( - UCode::InvalidArgument, - "invalid timestamp: seconds value out of range", - ) - }; - if ts.nanos < 0 || ts.nanos >= 1_000_000_000 { - return Err(UStatus::fail_with_code( - UCode::InvalidArgument, - "invalid timestamp: nanos value out of range", - )); - } - u64::try_from(ts.seconds) - .ok() - .and_then(|s| s.checked_mul(1000)) - .and_then(|ms| ms.checked_add(ts.nanos as u64 / 1_000_000)) - .ok_or_else(err) - .map(Some) - } else { - Ok(None) - } -} - -impl TryFrom<&Subscription> for SubscriptionInfo { - type Error = UStatus; - - fn try_from(subscription_proto: &Subscription) -> Result { - let topic = subscription_proto - .topic - .as_ref() - .ok_or(UStatus::fail_with_code( - UCode::InvalidArgument, - "topic missing", - )) - .and_then(|t| { - UUri::try_from(t) - .map_err(|_| UStatus::fail_with_code(UCode::InvalidArgument, "invalid topic")) - })?; - let subscriber = subscription_proto - .subscriber - .as_ref() - .and_then(|s| s.uri.as_ref()) - .ok_or(UStatus::fail_with_code( - UCode::InvalidArgument, - "subscriber missing", - )) - .and_then(|s| { - UUri::try_from(s).map_err(|_| { - UStatus::fail_with_code(UCode::InvalidArgument, "invalid subscriber") - }) - })?; - let status = subscription_proto - .status - .as_ref() - .ok_or(UStatus::fail_with_code( - UCode::InvalidArgument, - "status missing", - )) - .and_then(SubscriptionStatus::try_from)?; - subscription_proto - .attributes - .as_ref() - .ok_or_else(|| UStatus::fail_with_code(UCode::InvalidArgument, "missing attributes")) - .and_then(|attributes| { - let expiration = - protobuf_timestamp_as_unix_epoch_milliseconds(attributes.expire.as_ref())?; - Ok(SubscriptionInfo::new( - topic, - subscriber, - status, - expiration, - attributes.sample_period_ms, - )) - }) - } -} - -impl TryFrom<&Update> for SubscriptionInfo { - type Error = UStatus; - fn try_from(update_proto: &Update) -> Result { - let topic = update_proto - .topic - .as_ref() - .ok_or(UStatus::fail_with_code( - UCode::InvalidArgument, - "topic missing", - )) - .and_then(|t| { - UUri::try_from(t) - .map_err(|_| UStatus::fail_with_code(UCode::InvalidArgument, "invalid topic")) - })?; - let subscriber = update_proto - .subscriber - .as_ref() - .and_then(|s| s.uri.as_ref()) - .ok_or(UStatus::fail_with_code( - UCode::InvalidArgument, - "subscriber missing", - )) - .and_then(|s| { - UUri::try_from(s).map_err(|_| { - UStatus::fail_with_code(UCode::InvalidArgument, "invalid subscriber") - }) - })?; - let status = update_proto - .status - .as_ref() - .ok_or(UStatus::fail_with_code( - UCode::InvalidArgument, - "status missing", - )) - .and_then(SubscriptionStatus::try_from)?; - let attribs = update_proto.attributes.get_or_default(); - let expiration = protobuf_timestamp_as_unix_epoch_milliseconds(attribs.expire.as_ref())?; - Ok(SubscriptionInfo::new( - topic, - subscriber, - status, - expiration, - attribs.sample_period_ms, - )) - } -} - -impl From for crate::up_core_api::usubscription::reset_request::reason::Code { - fn from(reason: ResetReason) -> Self { - match reason { - ResetReason::Unspecified => Self::UNSPECIFIED, - ResetReason::FactoryReset => Self::FACTORY_RESET, - ResetReason::CorruptedData => Self::CORRUPTED_DATA, - } - } -} - /// A [`USubscription`] client implementation for invoking operations of a local USubscription service. /// /// The client requires an [`RpcClient`] for performing the remote procedure calls. pub struct RpcClientUSubscription { rpc_client: Arc, + usubscription_authority: Option, } impl RpcClientUSubscription { @@ -210,8 +45,12 @@ impl RpcClientUSubscription { /// # Arguments /// /// * `rpc_client` - The client to use for performing the remote procedure calls on the USubscription service. - pub fn new(rpc_client: Arc) -> Self { - RpcClientUSubscription { rpc_client } + /// * `usubscription_authority` - (Optional) authority of USubscription service to address - will default to local authority ("") if None. + pub fn new(rpc_client: Arc, usubscription_authority: Option) -> Self { + RpcClientUSubscription { + rpc_client, + usubscription_authority, + } } fn default_call_options() -> CallOptions { @@ -219,79 +58,41 @@ impl RpcClientUSubscription { } } -impl RpcClientUSubscription { - async fn fetch_subscriptions( - &self, - fetch_subscriptions_request: FetchSubscriptionsRequest, - ) -> Result, UStatus> { - let response = self - .rpc_client - .invoke_proto_method::<_, FetchSubscriptionsResponse>( - usubscription_uri(RESOURCE_ID_FETCH_SUBSCRIPTIONS), - Self::default_call_options(), - fetch_subscriptions_request, - ) - .await - .map_err(UStatus::from)?; - let mut result = Vec::new(); - for subscription in &response.subscriptions { - let info = SubscriptionInfo::try_from(subscription)?; - result.push(info); - } - Ok(result) - } -} - #[async_trait] impl USubscription for RpcClientUSubscription { async fn subscribe( &self, topic: &UUri, - expiration: Option, // millis since Unix Epoch - min_sample_period: Option, + expiration: Option, + sample_period: Option, ) -> Result { - let subscription_request = SubscriptionRequest { - topic: Some(topic.into()).into(), - attributes: match (expiration, min_sample_period) { - (None, None) => None.into(), - _ => Some(SubscribeAttributes { - expire: unix_epoch_millis_as_protobuf_timestamp(expiration)?.into(), - sample_period_ms: min_sample_period, - ..Default::default() - }) - .into(), - }, - ..Default::default() - }; - self.rpc_client - .invoke_proto_method::<_, SubscriptionResponse>( - usubscription_uri(RESOURCE_ID_SUBSCRIBE), + // [impl->dsn~usubscription-unsubscribe-valid-topic-uuris~1] + let subscribe_request = SubscribeRequest::new(topic, expiration, sample_period)?; + + Ok(self + .rpc_client + .invoke_proto_method::<_, SubscribeResponse>( + usubscription_uri( + self.usubscription_authority.as_deref(), + RESOURCE_ID_SUBSCRIBE, + ), Self::default_call_options(), - subscription_request, + subscribe_request, ) - .await - .and_then(|response| { - Ok(response.status.as_ref().map_or_else( - || { - Err(UStatus::fail_with_code( - UCode::InvalidArgument, - "uSubscription returned invalid response: no subscription status", - )) - }, - SubscriptionStatus::try_from, - )?) - }) - .map_err(UStatus::from) + .await? + .status) } async fn unsubscribe(&self, topic: &UUri) -> Result<(), UStatus> { let unsubscribe_request = UnsubscribeRequest { - topic: Some(topic.into()).into(), - ..Default::default() + topic: topic.clone(), }; self.rpc_client - .invoke_proto_method::<_, UnsubscribeResponse>( - usubscription_uri(RESOURCE_ID_UNSUBSCRIBE), + .invoke_proto_method::<_, UnsubscribeResponseProto>( + usubscription_uri( + self.usubscription_authority.as_deref(), + RESOURCE_ID_UNSUBSCRIBE, + ), Self::default_call_options(), unsubscribe_request, ) @@ -300,135 +101,64 @@ impl USubscription for RpcClientUSubscription { .map_err(UStatus::from) } - async fn fetch_subscriptions_by_topic( - &self, - topic: &UUri, - ) -> Result, UStatus> { - let fetch_subscriptions_request = - crate::up_core_api::usubscription::FetchSubscriptionsRequest { - request: Some( - crate::up_core_api::usubscription::fetch_subscriptions_request::Request::Topic( - topic.into(), - ), - ), - ..Default::default() - }; - self.fetch_subscriptions(fetch_subscriptions_request).await - } - - async fn fetch_subscriptions_by_subscriber( + async fn fetch_subscriptions( &self, - subscriber: &UUri, + topic_filter: Option, + subscriber_filter: Option, ) -> Result, UStatus> { - let subscriber_info = crate::up_core_api::usubscription::SubscriberInfo { - uri: Some(subscriber.into()).into(), - ..Default::default() - }; - let fetch_subscriptions_request = - crate::up_core_api::usubscription::FetchSubscriptionsRequest { - request: Some( - crate::up_core_api::usubscription::fetch_subscriptions_request::Request::Subscriber(subscriber_info), + Ok(self + .rpc_client + .invoke_proto_method::<_, FetchSubscriptionsResponse>( + usubscription_uri( + self.usubscription_authority.as_deref(), + RESOURCE_ID_FETCH_SUBSCRIPTIONS, ), - ..Default::default() - }; - self.fetch_subscriptions(fetch_subscriptions_request).await + Self::default_call_options(), + FetchSubscriptionsRequest { + topic_filter, + subscriber_filter, + }, + ) + .await? + .subscriptions) } - async fn register_for_notifications(&self, topic: &UUri) -> Result<(), UStatus> { - let notifications_register_request = - crate::up_core_api::usubscription::NotificationsRequest { - topic: Some(topic.into()).into(), - ..Default::default() - }; + async fn register_for_notifications(&self) -> Result<(), UStatus> { self.rpc_client - .invoke_proto_method::<_, NotificationsResponse>( - usubscription_uri(RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS), + .invoke_proto_method::<_, NotificationResponseProto>( + usubscription_uri( + self.usubscription_authority.as_deref(), + RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS, + ), Self::default_call_options(), - notifications_register_request, + crate::up_core_api::usubscription::NotificationsRequest::default(), ) .await .map(|_response| ()) .map_err(UStatus::from) } - async fn unregister_for_notifications(&self, topic: &UUri) -> Result<(), UStatus> { - let notifications_unregister_request = - crate::up_core_api::usubscription::NotificationsRequest { - topic: Some(topic.into()).into(), - ..Default::default() - }; + async fn unregister_for_notifications(&self) -> Result<(), UStatus> { self.rpc_client - .invoke_proto_method::<_, NotificationsResponse>( - usubscription_uri(RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS), + .invoke_proto_method::<_, NotificationResponseProto>( + usubscription_uri( + self.usubscription_authority.as_deref(), + RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS, + ), Self::default_call_options(), - notifications_unregister_request, + crate::up_core_api::usubscription::NotificationsRequest::default(), ) .await .map(|_response| ()) .map_err(UStatus::from) } - async fn fetch_subscribers(&self, topic: &UUri) -> Result, UStatus> { - let fetch_subscribers_request = - crate::up_core_api::usubscription::FetchSubscribersRequest { - topic: Some(topic.into()).into(), - ..Default::default() - }; - let response = self - .rpc_client - .invoke_proto_method::<_, FetchSubscribersResponse>( - usubscription_uri(RESOURCE_ID_FETCH_SUBSCRIBERS), - Self::default_call_options(), - fetch_subscribers_request, - ) - .await?; - let mut result = vec![]; - for subscriber_info in &response.subscribers { - let uri = subscriber_info - .uri - .as_ref() - .ok_or_else(|| { - UStatus::fail_with_code( - UCode::InvalidArgument, - "uSubscription returned invalid response: missing subscriber URI", - ) - }) - .and_then(|uri_proto| { - UUri::try_from(uri_proto).map_err(|_e| { - UStatus::fail_with_code( - UCode::InvalidArgument, - "uSubscription returned invalid response: invalid subscriber URI", - ) - }) - })?; - result.push(uri); - } - Ok(result) - } - - async fn reset( - &self, - reason: ResetReason, - message: Option, - before: Option, // millis since Unix Epoch - ) -> Result<(), UStatus> { - let before_ts = unix_epoch_millis_as_protobuf_timestamp(before)?; - let reset_request = crate::up_core_api::usubscription::ResetRequest { - reason: Some(crate::up_core_api::usubscription::reset_request::Reason { - code: crate::up_core_api::usubscription::reset_request::reason::Code::from(reason) - .into(), - message, - ..Default::default() - }) - .into(), - before: before_ts.into(), - ..Default::default() - }; + async fn reset(&self) -> Result<(), UStatus> { self.rpc_client - .invoke_proto_method::<_, ResetResponse>( - usubscription_uri(RESOURCE_ID_RESET), + .invoke_proto_method::<_, ResetResponseProto>( + usubscription_uri(self.usubscription_authority.as_deref(), RESOURCE_ID_RESET), Self::default_call_options(), - reset_request, + crate::up_core_api::usubscription::ResetRequest::default(), ) .await .map(|_response| ()) @@ -443,75 +173,37 @@ mod tests { use super::*; use crate::{ communication::{MockRpcClient, UPayload}, - up_core_api::usubscription::{ - fetch_subscriptions_request::Request, FetchSubscribersRequest, NotificationsRequest, - ResetRequest, - }, + core::usubscription::{FetchSubscriptionsRequest, SubscribeRequest}, UCode, UUri, }; use std::sync::Arc; - #[test] - fn test_unix_epoch_millis_as_protobuf_timestamp() { - assert!( - unix_epoch_millis_as_protobuf_timestamp(Some(1_000)).is_ok_and(|ts| { - ts == Some(Timestamp { - seconds: 1, - nanos: 0, - ..Default::default() - }) - }) - ); - - assert!( - unix_epoch_millis_as_protobuf_timestamp(Some(1_234)).is_ok_and(|ts| { - ts == Some(Timestamp { - seconds: 1, - nanos: 234_000_000, - ..Default::default() - }) - }) - ); - - assert!(unix_epoch_millis_as_protobuf_timestamp(Some(u64::MAX)).is_ok()); - assert!(unix_epoch_millis_as_protobuf_timestamp(None).is_ok_and(|ts| ts.is_none())); - } - - #[test_case::test_case(10, 234_000_000 => matches Ok(Some(10_234)); "succeeds for valid timestamp")] - #[test_case::test_case(-10, 234_000_000 => matches Err(UStatus {..}); "fails for negative seconds")] - #[test_case::test_case(10, -1 => matches Err(UStatus {..}); "fails for nanos exceeding lower bound")] - #[test_case::test_case(10, 1_000_000_000 => matches Err(UStatus {..}); "fails for nanos exeeding upper bound")] - fn test_protobuf_timestamp_as_unix_epoch_milliseconds( - seconds: i64, - nanos: i32, - ) -> Result, UStatus> { - let timestamp = Timestamp { - seconds, - nanos, - ..Default::default() - }; - protobuf_timestamp_as_unix_epoch_milliseconds(Some(×tamp)) - } - + // [utest->req~usubscription-subscribe-method_id~1] + // [utest->req~usubscription-subscribe-request-signature~1] + // [utest->req~usubscription-subscribe-response-signature~1] #[tokio::test] async fn test_subscribe_invokes_rpc_client() { let topic = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap(); - let expected_request = SubscriptionRequest { - topic: Some((&topic).into()).into(), - ..Default::default() + let expected_request = SubscribeRequest { + topic: topic.clone(), + expiration: None, + sample_period: None, }; let mut rpc_client = MockRpcClient::new(); let mut seq = Sequence::new(); + rpc_client .expect_invoke_method() .once() .in_sequence(&mut seq) .withf(|method, _options, payload| { - method == &usubscription_uri(RESOURCE_ID_SUBSCRIBE) && payload.is_some() + method == &usubscription_uri(None, RESOURCE_ID_SUBSCRIBE) && payload.is_some() }) .return_const(Err(crate::communication::ServiceInvocationError::Internal( "internal error".to_string(), ))); + + let topic_clone = topic.clone(); rpc_client .expect_invoke_method() .once() @@ -520,23 +212,20 @@ mod tests { let request = payload .to_owned() .unwrap() - .extract_protobuf::() + .extract_protobuf::() .unwrap(); - request == expected_request && method == &usubscription_uri(RESOURCE_ID_SUBSCRIBE) + request == expected_request + && method == &usubscription_uri(None, RESOURCE_ID_SUBSCRIBE) }) .returning(move |_method, _options, _payload| { - let response = SubscriptionResponse { - status: Some(crate::up_core_api::usubscription::SubscriptionStatus { - state: crate::up_core_api::usubscription::subscription_status::State::SUBSCRIBED - .into(), - ..Default::default() - }).into(), - ..Default::default() + let response = SubscribeResponse { + topic: topic_clone.clone(), + status: SubscriptionStatus::Subscribed, }; Ok(Some(UPayload::try_from_protobuf(response).unwrap())) }); - let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client)); + let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client), None); assert!(usubscription_client .subscribe(&topic, None, None) @@ -548,12 +237,14 @@ mod tests { .is_ok()); } + // [utest->req~usubscription-unsubscribe-method_id~1] + // [utest->req~usubscription-unsubscribe-request-signature~1] + // [utest->req~usubscription-unsubscribe-response-signature~1] #[tokio::test] async fn test_unsubscribe_invokes_rpc_client() { let topic = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap(); let expected_request = UnsubscribeRequest { - topic: Some((&topic).into()).into(), - ..Default::default() + topic: topic.clone(), }; let mut rpc_client = MockRpcClient::new(); let mut seq = Sequence::new(); @@ -562,7 +253,7 @@ mod tests { .once() .in_sequence(&mut seq) .withf(|method, _options, payload| { - method == &usubscription_uri(RESOURCE_ID_UNSUBSCRIBE) && payload.is_some() + method == &usubscription_uri(None, RESOURCE_ID_UNSUBSCRIBE) && payload.is_some() }) .return_const(Err(crate::communication::ServiceInvocationError::Internal( "internal error".to_string(), @@ -577,16 +268,16 @@ mod tests { .unwrap() .extract_protobuf::() .unwrap(); - request == expected_request && method == &usubscription_uri(RESOURCE_ID_UNSUBSCRIBE) + request == expected_request + && method == &usubscription_uri(None, RESOURCE_ID_UNSUBSCRIBE) }) .returning(move |_method, _options, _payload| { - let response = UnsubscribeResponse { - ..Default::default() - }; - Ok(Some(UPayload::try_from_protobuf(response).unwrap())) + Ok(Some( + UPayload::try_from_protobuf(UnsubscribeResponseProto::default()).unwrap(), + )) }); - let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client)); + let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client), None); assert!(usubscription_client .unsubscribe(&topic) @@ -595,12 +286,15 @@ mod tests { assert!(usubscription_client.unsubscribe(&topic).await.is_ok()); } + // [utest->req~usubscription-fetch-subscriptions-method_id~1] + // [utest->req~usubscription-fetch-subscriptions-request-signature~1] + // [utest->req~usubscription-fetch-subscriptions-response-signature~1] #[tokio::test] async fn test_fetch_subscriptions_invokes_rpc_client() { - let topic = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap(); + let topic_filter = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap(); let expected_request = FetchSubscriptionsRequest { - request: Some(Request::Topic((&topic).into())), - ..Default::default() + topic_filter: Some(topic_filter.clone()), + subscriber_filter: None, }; let mut rpc_client = MockRpcClient::new(); let mut seq = Sequence::new(); @@ -609,7 +303,8 @@ mod tests { .once() .in_sequence(&mut seq) .withf(|method, _options, payload| { - method == &usubscription_uri(RESOURCE_ID_FETCH_SUBSCRIPTIONS) && payload.is_some() + method == &usubscription_uri(None, RESOURCE_ID_FETCH_SUBSCRIPTIONS) + && payload.is_some() }) .return_const(Err(crate::communication::ServiceInvocationError::Internal( "internal error".to_string(), @@ -626,83 +321,32 @@ mod tests { .unwrap(); request == expected_request - && method == &usubscription_uri(RESOURCE_ID_FETCH_SUBSCRIPTIONS) + && method == &usubscription_uri(None, RESOURCE_ID_FETCH_SUBSCRIPTIONS) }) .returning(move |_method, _options, _payload| { let response = FetchSubscriptionsResponse { - ..Default::default() + subscriptions: Vec::default(), }; Ok(Some(UPayload::try_from_protobuf(response).unwrap())) }); - let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client)); + let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client), None); assert!(usubscription_client - .fetch_subscriptions_by_topic(&topic) + .fetch_subscriptions(Some(topic_filter.clone()), None) .await .is_err_and(|e| e.get_code() == UCode::Internal)); assert!(usubscription_client - .fetch_subscriptions_by_topic(&topic) + .fetch_subscriptions(Some(topic_filter), None) .await .is_ok()); } - #[tokio::test] - async fn test_fetch_subscribers_invokes_rpc_client() { - let topic = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap(); - let expected_request = FetchSubscribersRequest { - topic: Some((&topic).into()).into(), - ..Default::default() - }; - let mut rpc_client = MockRpcClient::new(); - let mut seq = Sequence::new(); - rpc_client - .expect_invoke_method() - .once() - .in_sequence(&mut seq) - .withf(|method, _options, payload| { - method == &usubscription_uri(RESOURCE_ID_FETCH_SUBSCRIBERS) && payload.is_some() - }) - .return_const(Err(crate::communication::ServiceInvocationError::Internal( - "internal error".to_string(), - ))); - rpc_client - .expect_invoke_method() - .once() - .in_sequence(&mut seq) - .withf(move |method, _options, payload| { - let request = payload - .to_owned() - .unwrap() - .extract_protobuf::() - .unwrap(); - - request == expected_request - && method == &usubscription_uri(RESOURCE_ID_FETCH_SUBSCRIBERS) - }) - .returning(move |_method, _options, _payload| { - let response = FetchSubscribersResponse { - ..Default::default() - }; - Ok(Some(UPayload::try_from_protobuf(response).unwrap())) - }); - - let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client)); - - assert!(usubscription_client - .fetch_subscribers(&topic) - .await - .is_err_and(|e| e.get_code() == UCode::Internal)); - assert!(usubscription_client.fetch_subscribers(&topic).await.is_ok()); - } - + // [utest->req~usubscription-register-notifications-method_id~1] + // [utest->req~usubscription-register-notifications-request-signature~1] + // [utest->req~usubscription-register-notifications-response-signature~1] #[tokio::test] async fn test_register_for_notifications_invokes_rpc_client() { - let topic = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap(); - let expected_request = NotificationsRequest { - topic: Some((&topic).into()).into(), - ..Default::default() - }; let mut rpc_client = MockRpcClient::new(); let mut seq = Sequence::new(); rpc_client @@ -710,7 +354,7 @@ mod tests { .once() .in_sequence(&mut seq) .withf(|method, _options, payload| { - method == &usubscription_uri(RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS) + method == &usubscription_uri(None, RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS) && payload.is_some() }) .return_const(Err(crate::communication::ServiceInvocationError::Internal( @@ -720,42 +364,32 @@ mod tests { .expect_invoke_method() .once() .in_sequence(&mut seq) - .withf(move |method, _options, payload| { - let request = payload - .to_owned() - .unwrap() - .extract_protobuf::() - .unwrap(); - - request == expected_request - && method == &usubscription_uri(RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS) + .withf(move |method, _options, _payload| { + method == &usubscription_uri(None, RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS) }) .returning(move |_method, _options, _payload| { - let response = NotificationsResponse { - ..Default::default() - }; - Ok(Some(UPayload::try_from_protobuf(response).unwrap())) + Ok(Some( + UPayload::try_from_protobuf(NotificationResponseProto::default()).unwrap(), + )) }); - let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client)); + let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client), None); assert!(usubscription_client - .register_for_notifications(&topic) + .register_for_notifications() .await .is_err_and(|e| e.get_code() == UCode::Internal)); assert!(usubscription_client - .register_for_notifications(&topic) + .register_for_notifications() .await .is_ok()); } + // [utest->req~usubscription-unregister-notifications-method_id~1] + // [utest->req~usubscription-unregister-notifications-request-signature~1] + // [utest->req~usubscription-unregister-notifications-response-signature~1] #[tokio::test] async fn test_unregister_for_notifications_invokes_rpc_client() { - let topic = UUri::try_from_parts("other", 0xd5a3, 0x01, 0xd3fe).unwrap(); - let expected_request = NotificationsRequest { - topic: Some((&topic).into()).into(), - ..Default::default() - }; let mut rpc_client = MockRpcClient::new(); let mut seq = Sequence::new(); rpc_client @@ -763,7 +397,7 @@ mod tests { .once() .in_sequence(&mut seq) .withf(|method, _options, payload| { - method == &usubscription_uri(RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS) + method == &usubscription_uri(None, RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS) && payload.is_some() }) .return_const(Err(crate::communication::ServiceInvocationError::Internal( @@ -773,46 +407,32 @@ mod tests { .expect_invoke_method() .once() .in_sequence(&mut seq) - .withf(move |method, _options, payload| { - let request = payload - .to_owned() - .unwrap() - .extract_protobuf::() - .unwrap(); - - request == expected_request - && method == &usubscription_uri(RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS) + .withf(move |method, _options, _payload| { + method == &usubscription_uri(None, RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS) }) .returning(move |_method, _options, _payload| { - let response = NotificationsResponse { - ..Default::default() - }; - Ok(Some(UPayload::try_from_protobuf(response).unwrap())) + Ok(Some( + UPayload::try_from_protobuf(NotificationResponseProto::default()).unwrap(), + )) }); - let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client)); + let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client), None); assert!(usubscription_client - .unregister_for_notifications(&topic) + .unregister_for_notifications() .await .is_err_and(|e| e.get_code() == UCode::Internal)); assert!(usubscription_client - .unregister_for_notifications(&topic) + .unregister_for_notifications() .await .is_ok()); } + // [utest->req~usubscription-reset-method_id~1] + // [utest->req~usubscription-reset-request-signature~1] + // [utest->req~usubscription-reset-response-signature~1] #[tokio::test] async fn test_reset_invokes_rpc_client() { - let expected_request = ResetRequest { - reason: Some(crate::up_core_api::usubscription::reset_request::Reason { - code: crate::up_core_api::usubscription::reset_request::reason::Code::UNSPECIFIED - .into(), - ..Default::default() - }) - .into(), - ..Default::default() - }; let mut rpc_client = MockRpcClient::new(); let mut seq = Sequence::new(); rpc_client @@ -820,7 +440,7 @@ mod tests { .once() .in_sequence(&mut seq) .withf(|method, _options, payload| { - method == &usubscription_uri(RESOURCE_ID_RESET) && payload.is_some() + method == &usubscription_uri(None, RESOURCE_ID_RESET) && payload.is_some() }) .return_const(Err(crate::communication::ServiceInvocationError::Internal( "internal error".to_string(), @@ -829,31 +449,21 @@ mod tests { .expect_invoke_method() .once() .in_sequence(&mut seq) - .withf(move |method, _options, payload| { - let request = payload - .to_owned() - .unwrap() - .extract_protobuf::() - .unwrap(); - - request == expected_request && method == &usubscription_uri(RESOURCE_ID_RESET) + .withf(move |method, _options, _payload| { + method == &usubscription_uri(None, RESOURCE_ID_RESET) }) .returning(move |_method, _options, _payload| { - let response = ResetResponse { - ..Default::default() - }; - Ok(Some(UPayload::try_from_protobuf(response).unwrap())) + Ok(Some( + UPayload::try_from_protobuf(ResetResponseProto::default()).unwrap(), + )) }); - let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client)); + let usubscription_client = RpcClientUSubscription::new(Arc::new(rpc_client), None); assert!(usubscription_client - .reset(ResetReason::Unspecified, None, None) + .reset() .await .is_err_and(|e| e.get_code() == UCode::Internal)); - assert!(usubscription_client - .reset(ResetReason::Unspecified, None, None) - .await - .is_ok()); + assert!(usubscription_client.reset().await.is_ok()); } } diff --git a/src/core/usubscription/usubscription_proto.rs b/src/core/usubscription/usubscription_proto.rs new file mode 100644 index 00000000..176216a6 --- /dev/null +++ b/src/core/usubscription/usubscription_proto.rs @@ -0,0 +1,692 @@ +/******************************************************************************** + * Copyright (c) 2026 Contributors to the Eclipse Foundation + * + * See the NOTICE file(s) distributed with this work for additional + * information regarding copyright ownership. + * + * This program and the accompanying materials are made available under the + * terms of the Apache License Version 2.0 which is available at + * https://www.apache.org/licenses/LICENSE-2.0 + * + * SPDX-License-Identifier: Apache-2.0 + ********************************************************************************/ + +/*! +Conversion logic for up-rust USubscription message types into and from their corresponding protobuf types. +*/ + +use std::time::{Duration, SystemTime, UNIX_EPOCH}; + +use protobuf::{ + well_known_types::{any::Any, timestamp::Timestamp}, + EnumFull, EnumOrUnknown, Message, MessageField, +}; + +use crate::{ + core::usubscription::{ + FetchSubscriptionsRequest, FetchSubscriptionsResponse, SubscribeRequest, SubscribeResponse, + SubscriptionInfo, SubscriptionStatus, UnsubscribeRequest, + }, + up_core_api::{ + uri::UUri as UUriProto, + usubscription::{ + FetchSubscriptionsRequest as FetchSubscriptionsRequestProto, + FetchSubscriptionsResponse as FetchSubscriptionsResponseProto, + SubscribeRequest as SubscribeRequestProto, SubscribeResponse as SubscribeResponseProto, + Subscription as SubscriptionInfoProto, + SubscriptionStatus::{ + self as SubscriptionStatusProto, STATUS_SUBSCRIBED, STATUS_SUBSCRIBE_PENDING, + STATUS_UNSUBSCRIBED, STATUS_UNSUBSCRIBE_PENDING, + }, + UnsubscribeRequest as UnsubscribeRequestProto, + }, + }, + ProtobufMappable, SerializationError, UCode, UStatus, UUri, +}; + +pub(crate) fn system_time_as_protobuf_timestamp( + time: Option, +) -> Result, UStatus> { + if let Some(t) = time { + let (seconds, nanos) = match t.duration_since(UNIX_EPOCH) { + Ok(duration) => (duration.as_secs() as i64, duration.subsec_nanos() as i32), + Err(before_epoch) => { + let duration = before_epoch.duration(); + let whole_seconds = duration.as_secs() as i64; + let subsec_nanos = duration.subsec_nanos(); + // protobuf's Timestamp always uses a non-negative `nanos` offset added to + // `seconds`, so a non-zero fraction shifts one whole second into it + if subsec_nanos == 0 { + (-whole_seconds, 0) + } else { + (-whole_seconds - 1, 1_000_000_000 - subsec_nanos as i32) + } + } + }; + Ok(Some(Timestamp { + seconds, + nanos, + ..Default::default() + })) + } else { + Ok(None) + } +} + +pub(crate) fn protobuf_timestamp_as_system_time( + ts: Option<&Timestamp>, +) -> Result, UStatus> { + if let Some(ts) = ts { + let err = || { + UStatus::fail_with_code( + UCode::InvalidArgument, + "invalid timestamp: seconds value out of range", + ) + }; + if ts.nanos < 0 || ts.nanos >= 1_000_000_000 { + return Err(UStatus::fail_with_code( + UCode::InvalidArgument, + "invalid timestamp: nanos value out of range", + )); + } + // nanos already validated to be in [0, 1_000_000_000), so this cast is safe + let nanos = ts.nanos as u32; + let time = if ts.seconds >= 0 { + UNIX_EPOCH.checked_add(Duration::new(ts.seconds as u64, nanos)) + } else { + // reverse of the +1s/complementary-nanos shift applied when encoding a + // pre-epoch instant; `-(ts.seconds + 1)` cannot overflow since ts.seconds < 0 + let whole_seconds_before_epoch = (-(ts.seconds + 1)) as u64; + UNIX_EPOCH.checked_sub(Duration::new( + whole_seconds_before_epoch, + 1_000_000_000 - nanos, + )) + }; + time.ok_or_else(err).map(Some) + } else { + Ok(None) + } +} + +// SubscriptionStatus conversions +impl TryFrom<&SubscriptionStatusProto> for SubscriptionStatus { + type Error = UStatus; + + fn try_from(value: &SubscriptionStatusProto) -> Result { + match value { + STATUS_UNSUBSCRIBED => Ok(SubscriptionStatus::Unsubscribed), + STATUS_SUBSCRIBE_PENDING => Ok(SubscriptionStatus::SubscribePending), + STATUS_SUBSCRIBED => Ok(SubscriptionStatus::Subscribed), + STATUS_UNSUBSCRIBE_PENDING => Ok(SubscriptionStatus::UnsubscribePending), + _ => Err(UStatus::fail_with_code( + UCode::OutOfRange, + format!("invalid subscription status {}", value.descriptor()), + )), + } + } +} + +impl From<&SubscriptionStatus> for SubscriptionStatusProto { + fn from(value: &SubscriptionStatus) -> Self { + match value { + SubscriptionStatus::Unsubscribed => STATUS_UNSUBSCRIBED, + SubscriptionStatus::SubscribePending => STATUS_SUBSCRIBE_PENDING, + SubscriptionStatus::Subscribed => STATUS_SUBSCRIBED, + SubscriptionStatus::UnsubscribePending => STATUS_UNSUBSCRIBE_PENDING, + } + } +} + +// SubscriptionInfo conversions +impl TryFrom<&SubscriptionInfo> for SubscriptionInfoProto { + type Error = UStatus; + + fn try_from(value: &SubscriptionInfo) -> Result { + Ok(SubscriptionInfoProto { + subscriber: MessageField::some(UUriProto::from(&value.subscriber)), + topic: MessageField::some(UUriProto::from(&value.topic)), + status: SubscriptionStatusProto::from(&value.status).into(), + expiration: system_time_as_protobuf_timestamp(*value.expiration())?.into(), + sample_period: value + .min_sample_period + .map(|p| p.as_millis().clamp(0, u32::MAX as u128) as u32), + ..Default::default() + }) + } +} + +impl TryFrom<&SubscriptionInfoProto> for SubscriptionInfo { + type Error = UStatus; + + fn try_from(proto: &SubscriptionInfoProto) -> Result { + let subscriber = proto + .subscriber + .as_ref() + .ok_or(UStatus::fail_with_code( + UCode::InvalidArgument, + "subscriber missing", + )) + .and_then(|s| { + UUri::try_from(s).map_err(|_| { + UStatus::fail_with_code(UCode::InvalidArgument, "invalid subscriber") + }) + })?; + + Ok(SubscriptionInfo::new( + require_topic(proto.topic.clone())?, + subscriber, + require_status(proto.status)?, + protobuf_timestamp_as_system_time(proto.expiration.as_ref())?, + proto.sample_period.map(|p| Duration::from_millis(p as u64)), + )) + } +} + +impl ProtobufMappable for SubscriptionInfo { + fn parse_from_packed_protobuf_bytes(proto: &[u8]) -> Result { + Any::parse_from_bytes(proto) + .map_err(|err| crate::SerializationError::new(err.to_string())) + .and_then(|any| match any.unpack::() { + Ok(Some(message_proto)) => SubscriptionInfo::try_from(&message_proto) + .map_err(|e| crate::SerializationError::new(e.to_string())), + Ok(None) => Err(crate::SerializationError::new( + "protobuf Any does not contain expected type".to_string(), + )), + Err(e) => Err(crate::SerializationError::new(format!( + "protobuf Any unpack error: {e}" + ))), + }) + } + + fn parse_from_protobuf_bytes(proto: &[u8]) -> Result { + SubscriptionInfo::try_from(&SubscriptionInfoProto::parse_from_bytes(proto)?) + .map_err(|e| SerializationError::new(e.to_string())) + } + + fn write_to_packed_protobuf_bytes(&self) -> Result, crate::SerializationError> { + Any::pack(&SubscriptionInfoProto::try_from(self).map_err(|e| { + SerializationError::new(format!("failed to serialize to protobuf: {e}")) + })?) + .map_err(|e| crate::SerializationError::new(format!("failed to pack: {e}"))) + .and_then(|any| any.write_to_protobuf_bytes()) + } + + fn write_to_protobuf_bytes(&self) -> Result, crate::SerializationError> { + Ok(SubscriptionInfoProto::try_from(self) + .map_err(|e| SerializationError::new(format!("failed to serialize to protobuf: {e}")))? + .write_to_bytes()?) + } +} + +// SubscribeRequest conversions +impl TryFrom<&SubscribeRequest> for SubscribeRequestProto { + type Error = UStatus; + + fn try_from(value: &SubscribeRequest) -> Result { + Ok(SubscribeRequestProto { + topic: MessageField::some(UUriProto::from(&value.topic)), + expiration: system_time_as_protobuf_timestamp(value.expiration)?.into(), + sample_period: value + .sample_period + .map(|p| p.as_millis().clamp(0, u32::MAX as u128) as u32), + ..Default::default() + }) + } +} + +// [impl->req~usubscription-subscribe-request-signature~1] +impl TryFrom<&SubscribeRequestProto> for SubscribeRequest { + type Error = UStatus; + + fn try_from(proto: &SubscribeRequestProto) -> Result { + Ok(SubscribeRequest { + // [impl->dsn~usubscription-subscribe-valid-topic-uuris~1] + topic: require_topic(proto.topic.clone())?, + // [impl->dsn~usubscription-subscription-expiration-datatype~1] + expiration: protobuf_timestamp_as_system_time(proto.expiration.as_ref())?, + // [impl->dsn~usubscription-sample-period-datatype~1] + sample_period: proto + .sample_period + .map(|sp| Duration::from_millis(sp as u64)), + }) + } +} + +impl ProtobufMappable for SubscribeRequest { + fn parse_from_packed_protobuf_bytes(proto: &[u8]) -> Result { + Any::parse_from_bytes(proto) + .map_err(|err| crate::SerializationError::new(err.to_string())) + .and_then(|any| match any.unpack::() { + Ok(Some(message_proto)) => SubscribeRequest::try_from(&message_proto) + .map_err(|e| crate::SerializationError::new(e.to_string())), + Ok(None) => Err(crate::SerializationError::new( + "protobuf Any does not contain expected type".to_string(), + )), + Err(e) => Err(crate::SerializationError::new(format!( + "protobuf Any unpack error: {e}" + ))), + }) + } + + fn parse_from_protobuf_bytes(proto: &[u8]) -> Result { + SubscribeRequest::try_from(&SubscribeRequestProto::parse_from_bytes(proto)?) + .map_err(|e| SerializationError::new(e.to_string())) + } + + fn write_to_packed_protobuf_bytes(&self) -> Result, crate::SerializationError> { + Any::pack(&SubscribeRequestProto::try_from(self).map_err(|e| { + SerializationError::new(format!("failed to serialize to protobuf: {e}")) + })?) + .map_err(|e| crate::SerializationError::new(format!("failed to pack: {e}"))) + .and_then(|any| any.write_to_protobuf_bytes()) + } + + fn write_to_protobuf_bytes(&self) -> Result, crate::SerializationError> { + Ok(SubscribeRequestProto::try_from(self) + .map_err(|e| SerializationError::new(format!("failed to serialize to protobuf: {e}")))? + .write_to_bytes()?) + } +} + +// SubscriptionResponse conversions +impl From<&SubscribeResponse> for SubscribeResponseProto { + fn from(value: &SubscribeResponse) -> Self { + SubscribeResponseProto { + topic: MessageField::some(UUriProto::from(&value.topic)), + status: SubscriptionStatusProto::from(&value.status).into(), + ..Default::default() + } + } +} + +// [impl->req~usubscription-subscribe-response-signature~1] +impl TryFrom<&SubscribeResponseProto> for SubscribeResponse { + type Error = UStatus; + + fn try_from(proto: &SubscribeResponseProto) -> Result { + Ok(SubscribeResponse { + topic: require_topic(proto.topic.clone())?, + status: require_status(proto.status)?, + }) + } +} + +impl ProtobufMappable for SubscribeResponse { + fn parse_from_packed_protobuf_bytes(proto: &[u8]) -> Result { + Any::parse_from_bytes(proto) + .map_err(|err| crate::SerializationError::new(err.to_string())) + .and_then(|any| match any.unpack::() { + Ok(Some(message_proto)) => SubscribeResponse::try_from(&message_proto) + .map_err(|e| crate::SerializationError::new(e.to_string())), + Ok(None) => Err(crate::SerializationError::new( + "protobuf Any does not contain expected type".to_string(), + )), + Err(e) => Err(crate::SerializationError::new(format!( + "protobuf Any unpack error: {e}" + ))), + }) + } + + fn parse_from_protobuf_bytes(proto: &[u8]) -> Result { + SubscribeResponse::try_from(&SubscribeResponseProto::parse_from_bytes(proto)?) + .map_err(|e| SerializationError::new(e.to_string())) + } + + fn write_to_packed_protobuf_bytes(&self) -> Result, crate::SerializationError> { + Any::pack(&SubscribeResponseProto::from(self)) + .map_err(|e| crate::SerializationError::new(format!("failed to pack: {e}"))) + .and_then(|any| any.write_to_protobuf_bytes()) + } + + fn write_to_protobuf_bytes(&self) -> Result, crate::SerializationError> { + Ok(SubscribeResponseProto::from(self).write_to_bytes()?) + } +} + +// UnsubscribeRequest conversions +impl TryFrom<&UnsubscribeRequest> for UnsubscribeRequestProto { + type Error = UStatus; + + fn try_from(value: &UnsubscribeRequest) -> Result { + Ok(UnsubscribeRequestProto { + topic: MessageField::some(UUriProto::from(&value.topic)), + ..Default::default() + }) + } +} + +// [impl->req~usubscription-unsubscribe-request-signature~1] +impl TryFrom<&UnsubscribeRequestProto> for UnsubscribeRequest { + type Error = UStatus; + + fn try_from(proto: &UnsubscribeRequestProto) -> Result { + Ok(UnsubscribeRequest { + // [impl->dsn~usubscription-unsubscribe-valid-topic-uuris~1] + topic: require_topic(proto.topic.clone())?, + }) + } +} + +impl ProtobufMappable for UnsubscribeRequest { + fn parse_from_packed_protobuf_bytes(proto: &[u8]) -> Result { + Any::parse_from_bytes(proto) + .map_err(|err| crate::SerializationError::new(err.to_string())) + .and_then(|any| match any.unpack::() { + Ok(Some(message_proto)) => UnsubscribeRequest::try_from(&message_proto) + .map_err(|e| crate::SerializationError::new(e.to_string())), + Ok(None) => Err(crate::SerializationError::new( + "protobuf Any does not contain expected type".to_string(), + )), + Err(e) => Err(crate::SerializationError::new(format!( + "protobuf Any unpack error: {e}" + ))), + }) + } + + fn parse_from_protobuf_bytes(proto: &[u8]) -> Result { + UnsubscribeRequest::try_from(&UnsubscribeRequestProto::parse_from_bytes(proto)?) + .map_err(|e| SerializationError::new(e.to_string())) + } + + fn write_to_packed_protobuf_bytes(&self) -> Result, crate::SerializationError> { + Any::pack(&UnsubscribeRequestProto::try_from(self).map_err(|e| { + SerializationError::new(format!("failed to serialize to protobuf: {e}")) + })?) + .map_err(|e| crate::SerializationError::new(format!("failed to pack: {e}"))) + .and_then(|any| any.write_to_protobuf_bytes()) + } + + fn write_to_protobuf_bytes(&self) -> Result, crate::SerializationError> { + Ok(UnsubscribeRequestProto::try_from(self) + .map_err(|e| SerializationError::new(format!("failed to serialize to protobuf: {e}")))? + .write_to_bytes()?) + } +} + +// FetchSubscriptionsRequest conversions +impl From<&FetchSubscriptionsRequest> for FetchSubscriptionsRequestProto { + fn from(value: &FetchSubscriptionsRequest) -> Self { + FetchSubscriptionsRequestProto { + topic_filter: value.topic_filter.as_ref().map(UUriProto::from).into(), + subscriber_filter: value.subscriber_filter.as_ref().map(UUriProto::from).into(), + ..Default::default() + } + } +} + +// [impl->req~usubscription-fetch-subscriptions-request-signature~1] +impl TryFrom<&FetchSubscriptionsRequestProto> for FetchSubscriptionsRequest { + type Error = UStatus; + + fn try_from(value: &FetchSubscriptionsRequestProto) -> Result { + // [impl->dsn~usubscription-fetch-subscriptions-invalid-subscriber-filter~1] + let subscriber_filter = value + .subscriber_filter + .as_ref() + .map(UUri::try_from) + .transpose() + .map_err(|_| { + UStatus::fail_with_code(UCode::InvalidArgument, "invalid subscriber filter") + })?; + + // [impl->dsn~usubscription-fetch-subscriptions-invalid-topic-filter~1] + let topic_filter = value + .topic_filter + .as_ref() + .map(UUri::try_from) + .transpose() + .map_err(|_| UStatus::fail_with_code(UCode::InvalidArgument, "invalid topic filter"))?; + + Ok(FetchSubscriptionsRequest { + topic_filter, + subscriber_filter, + }) + } +} + +impl ProtobufMappable for FetchSubscriptionsRequest { + fn parse_from_packed_protobuf_bytes(proto: &[u8]) -> Result { + Any::parse_from_bytes(proto) + .map_err(|err| crate::SerializationError::new(err.to_string())) + .and_then(|any| match any.unpack::() { + Ok(Some(message_proto)) => FetchSubscriptionsRequest::try_from(&message_proto) + .map_err(|e| crate::SerializationError::new(e.to_string())), + Ok(None) => Err(crate::SerializationError::new( + "protobuf Any does not contain expected type".to_string(), + )), + Err(e) => Err(crate::SerializationError::new(format!( + "protobuf Any unpack error: {e}" + ))), + }) + } + + fn parse_from_protobuf_bytes(proto: &[u8]) -> Result { + FetchSubscriptionsRequest::try_from(&FetchSubscriptionsRequestProto::parse_from_bytes( + proto, + )?) + .map_err(|e| SerializationError::new(e.to_string())) + } + + fn write_to_packed_protobuf_bytes(&self) -> Result, crate::SerializationError> { + Any::pack(&FetchSubscriptionsRequestProto::from(self)) + .map_err(|e| crate::SerializationError::new(format!("failed to pack: {e}"))) + .and_then(|any| any.write_to_protobuf_bytes()) + } + + fn write_to_protobuf_bytes(&self) -> Result, crate::SerializationError> { + Ok(FetchSubscriptionsRequestProto::from(self).write_to_bytes()?) + } +} + +// FetchSubscriptionsResponse conversions +impl TryFrom<&FetchSubscriptionsResponse> for FetchSubscriptionsResponseProto { + type Error = UStatus; + + fn try_from(value: &FetchSubscriptionsResponse) -> Result { + Ok(FetchSubscriptionsResponseProto { + subscriptions: value + .subscriptions + .iter() + .map(SubscriptionInfoProto::try_from) + .collect::, _>>()?, + ..Default::default() + }) + } +} + +// [impl->req~usubscription-fetch-subscriptions-response-signature~1] +impl TryFrom<&FetchSubscriptionsResponseProto> for FetchSubscriptionsResponse { + type Error = UStatus; + + fn try_from(value: &FetchSubscriptionsResponseProto) -> Result { + Ok(FetchSubscriptionsResponse { + subscriptions: value + .subscriptions + .iter() + .map(SubscriptionInfo::try_from) + .collect::, _>>()?, + }) + } +} + +impl ProtobufMappable for FetchSubscriptionsResponse { + fn parse_from_packed_protobuf_bytes(proto: &[u8]) -> Result { + Any::parse_from_bytes(proto) + .map_err(|err| crate::SerializationError::new(err.to_string())) + .and_then( + |any| match any.unpack::() { + Ok(Some(message_proto)) => FetchSubscriptionsResponse::try_from(&message_proto) + .map_err(|e| crate::SerializationError::new(e.to_string())), + Ok(None) => Err(crate::SerializationError::new( + "protobuf Any does not contain expected type".to_string(), + )), + Err(e) => Err(crate::SerializationError::new(format!( + "protobuf Any unpack error: {e}" + ))), + }, + ) + } + + fn parse_from_protobuf_bytes(proto: &[u8]) -> Result { + FetchSubscriptionsResponse::try_from(&FetchSubscriptionsResponseProto::parse_from_bytes( + proto, + )?) + .map_err(|e| SerializationError::new(e.to_string())) + } + + fn write_to_packed_protobuf_bytes(&self) -> Result, crate::SerializationError> { + Any::pack( + &FetchSubscriptionsResponseProto::try_from(self).map_err(|e| { + SerializationError::new(format!("failed to serialize to protobuf: {e}")) + })?, + ) + .map_err(|e| crate::SerializationError::new(format!("failed to pack: {e}"))) + .and_then(|any| any.write_to_protobuf_bytes()) + } + + fn write_to_protobuf_bytes(&self) -> Result, crate::SerializationError> { + Ok(FetchSubscriptionsResponseProto::try_from(self) + .map_err(|e| SerializationError::new(format!("failed to serialize to protobuf: {e}")))? + .write_to_bytes()?) + } +} + +/// Extracts and validates the topic from a protobuf request, failing if it is +/// absent or not a well-formed URI. +fn require_topic(topic: MessageField) -> Result { + let topic = topic + .into_option() + .ok_or_else(|| UStatus::fail_with_code(UCode::InvalidArgument, "missing topic"))?; + UUri::try_from(&topic).map_err(|e| { + UStatus::fail_with_code(UCode::InvalidArgument, format!("invalid topic URI: {e}")) + }) +} + +/// Extracts and validates the status from a protobuf request, failing if it is +/// absent or invalid. +fn require_status( + status: EnumOrUnknown, +) -> Result { + SubscriptionStatus::try_from(&status.enum_value().map_err(|_| { + UStatus::fail_with_code(UCode::InvalidArgument, "subscription status missing") + })?) +} + +#[cfg(feature = "up-core-types")] +#[cfg(test)] +mod tests { + use super::*; + use protobuf::well_known_types::timestamp::Timestamp; + + #[test] + fn test_system_time_as_protobuf_timestamp() { + assert!(system_time_as_protobuf_timestamp(None).is_ok_and(|ts| ts.is_none())); + + let time = UNIX_EPOCH + Duration::new(1, 234_000_000); + assert!( + system_time_as_protobuf_timestamp(Some(time)).is_ok_and(|ts| { + ts == Some(Timestamp { + seconds: 1, + nanos: 234_000_000, + ..Default::default() + }) + }) + ); + } + + #[test] + fn test_protobuf_timestamp_as_system_time_maps_none() { + assert!(protobuf_timestamp_as_system_time(None).is_ok_and(|dt| dt.is_none())); + } + + #[test_case::test_case(10, 234_000_000 => matches Ok(Some(_)); "succeeds for valid timestamp")] + #[test_case::test_case(-10, 234_000_000 => matches Ok(Some(_)); "succeeds for timestamp before Unix epoch")] + #[test_case::test_case(10, -1 => matches Err(UStatus {..}); "fails for nanos exceeding lower bound")] + #[test_case::test_case(10, 1_000_000_000 => matches Err(UStatus {..}); "fails for nanos exceeding upper bound")] + fn test_protobuf_timestamp_as_system_time( + seconds: i64, + nanos: i32, + ) -> Result, UStatus> { + let timestamp = Timestamp { + seconds, + nanos, + ..Default::default() + }; + protobuf_timestamp_as_system_time(Some(×tamp)) + } + + // [utest->dsn~usubscription-subscription-expiration-datatype~1] + #[test] + fn test_timestamp_conversion_round_trip() { + let time = UNIX_EPOCH + Duration::new(1_700_000_000, 123_456_789); + let timestamp = system_time_as_protobuf_timestamp(Some(time)) + .expect("conversion to protobuf timestamp should succeed"); + let round_tripped = protobuf_timestamp_as_system_time(timestamp.as_ref()) + .expect("conversion back to system time should succeed"); + assert_eq!(round_tripped, Some(time)); + } + + // [utest->dsn~usubscription-fetch-subscriptions-invalid-subscriber-filter~1] + #[test] + fn test_fetch_subscriptions_request_rejects_invalid_subscriber_filter() { + let proto = FetchSubscriptionsRequestProto { + subscriber_filter: MessageField::some(UUriProto { + // construct an invalid UUri, with a too-large resource ID + resource_id: 0xFFFF_FFFF, + ..Default::default() + }), + ..Default::default() + }; + let result = FetchSubscriptionsRequest::try_from(&proto); + assert!(result.is_err()); + assert_eq!(result.unwrap_err().get_code(), UCode::InvalidArgument); + } + + // [utest->dsn~usubscription-fetch-subscriptions-invalid-topic-filter~1] + #[test] + fn test_fetch_subscriptions_request_rejects_invalid_topic_filter() { + let proto = FetchSubscriptionsRequestProto { + topic_filter: MessageField::some(UUriProto { + // construct an invalid UUri, with a too-large resource ID + resource_id: 0xFFFF_FFFF, + ..Default::default() + }), + ..Default::default() + }; + let result = FetchSubscriptionsRequest::try_from(&proto); + assert!(result.is_err()); + assert_eq!(result.unwrap_err().get_code(), UCode::InvalidArgument); + } + + // [utest->dsn~usubscription-subscribe-valid-topic-uuris~1] + #[test] + fn test_subscribe_request_rejects_missing_topic() { + let proto = SubscribeRequestProto::default(); + let result = SubscribeRequest::try_from(&proto); + assert!(result.is_err()); + assert_eq!(result.unwrap_err().get_code(), UCode::InvalidArgument); + } + + // [utest->dsn~usubscription-unsubscribe-valid-topic-uuris~1] + #[test] + fn test_unsubscribe_request_rejects_missing_topic() { + let proto = UnsubscribeRequestProto::default(); + let result = UnsubscribeRequest::try_from(&proto); + assert!(result.is_err()); + assert_eq!(result.unwrap_err().get_code(), UCode::InvalidArgument); + } + + // [utest->dsn~usubscription-sample-period-datatype~1] + #[test] + fn test_subscribe_request_preserves_sample_period() { + let proto = SubscribeRequestProto { + topic: MessageField::some(UUriProto::from( + &UUri::try_from_parts("", 0x1234, 1, 0x8000).unwrap(), + )), + sample_period: Some(500), + ..Default::default() + }; + let request = SubscribeRequest::try_from(&proto).unwrap(); + assert_eq!(request.sample_period, Some(Duration::from_millis(500))); + } +} diff --git a/src/core/usubscription/usubscription_server.rs b/src/core/usubscription/usubscription_server.rs new file mode 100644 index 00000000..c6e3ea1a --- /dev/null +++ b/src/core/usubscription/usubscription_server.rs @@ -0,0 +1,167 @@ +/******************************************************************************** + * Copyright (c) 2026 Contributors to the Eclipse Foundation + * + * See the NOTICE file(s) distributed with this work for additional + * information regarding copyright ownership. + * + * This program and the accompanying materials are made available under the + * terms of the Apache License Version 2.0 which is available at + * https://www.apache.org/licenses/LICENSE-2.0 + * + * SPDX-License-Identifier: Apache-2.0 + ********************************************************************************/ + +/*! +Convenience conversion wrappers up-rust/protobuf message conversion for uSubscription server functions. +*/ + +use crate::{ + communication::UPayload, + core::usubscription::{ + FetchSubscriptionsRequest, FetchSubscriptionsResponse, SubscribeRequest, SubscribeResponse, + UnsubscribeRequest, RESOURCE_ID_FETCH_SUBSCRIPTIONS, + RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS, RESOURCE_ID_RESET, RESOURCE_ID_SUBSCRIBE, + RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS, RESOURCE_ID_UNSUBSCRIBE, + }, + up_core_api::usubscription::{ + FetchSubscriptionsRequest as FetchSubscriptionsRequestProto, + FetchSubscriptionsResponse as FetchSubscriptionsResponseProto, + SubscribeRequest as SubscribeRequestProto, UnsubscribeRequest as UnsubscribeRequestProto, + }, + UCode, UStatus, +}; + +/// A decoded uSubscription request, tagged by the operation it belongs to. +/// +/// Returned by [`extract_usubscription_request`] so that a server can `match` on +/// the operation and know the associated up-rust message type. +#[derive(Clone, Debug, PartialEq)] +#[non_exhaustive] +pub enum USubscriptionRequest { + /// A [`RESOURCE_ID_SUBSCRIBE`] request. + Subscribe(SubscribeRequest), + /// A [`RESOURCE_ID_UNSUBSCRIBE`] request. + Unsubscribe(UnsubscribeRequest), + /// A [`RESOURCE_ID_FETCH_SUBSCRIPTIONS`] request. + FetchSubscriptions(FetchSubscriptionsRequest), + /// A [`RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS`] request. + RegisterForNotification(()), + /// A [`RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS`] request. + UnregisterForNotification(()), + /// A [`RESOURCE_ID_RESET`] request. + Reset(()), +} + +/// A decoded uSubscription response, tagged by the operation it belongs to. +/// +/// Returned by [`pack_usubscription_response`] so that a server can `match` on +/// the operation and know the associated up-rust message type. +#[derive(Clone, Debug, PartialEq)] +#[non_exhaustive] +pub enum USubscriptionResponse { + /// A [`RESOURCE_ID_SUBSCRIBE`] request. + Subscribe(SubscribeResponse), + /// A [`RESOURCE_ID_UNSUBSCRIBE`] request. + Unsubscribe(()), + /// A [`RESOURCE_ID_FETCH_SUBSCRIPTIONS`] request. + FetchSubscriptions(FetchSubscriptionsResponse), + /// A [`RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS`] request. + RegisterForNotification(()), + /// A [`RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS`] request. + UnregisterForNotification(()), + /// A [`RESOURCE_ID_RESET`] request. + Reset(()), +} + +/// Decodes a uSubscription protobuf request into its corresponding up-rust representation. +/// +/// Used by a uSubscription service to turn a raw protobuf payload into fully unpacked, native Rust data. +/// +/// # Errors +/// +/// Returns an error if `resource_id` does not identify a supported operation, if +/// the payload is missing, or if it cannot be deserialized into the expected type. +pub fn extract_usubscription_request( + resource_id: u16, + request_payload: Option, +) -> Result { + let payload = request_payload.ok_or_else(|| { + UStatus::fail_with_code(UCode::InvalidArgument, "missing request payload") + })?; + + match resource_id { + // [impl->req~usubscription-subscribe-request-signature~1] + RESOURCE_ID_SUBSCRIBE => { + let request_proto = payload + .extract_protobuf::() + .map_err(|err| UStatus::fail_with_code(UCode::Internal, err.to_string()))?; + Ok(SubscribeRequest::try_from(&request_proto).map(USubscriptionRequest::Subscribe)?) + } + // [impl->req~usubscription-unsubscribe-request-signature~1] + RESOURCE_ID_UNSUBSCRIBE => { + let request_proto = payload + .extract_protobuf::() + .map_err(|err| UStatus::fail_with_code(UCode::Internal, err.to_string()))?; + Ok(UnsubscribeRequest::try_from(&request_proto) + .map(USubscriptionRequest::Unsubscribe)?) + } + // [impl->req~usubscription-fetch-subscriptions-request-signature~1] + RESOURCE_ID_FETCH_SUBSCRIPTIONS => { + let request_proto = payload + .extract_protobuf::() + .map_err(|err| UStatus::fail_with_code(UCode::Internal, err.to_string()))?; + Ok(FetchSubscriptionsRequest::try_from(&request_proto) + .map(USubscriptionRequest::FetchSubscriptions)?) + } + // [impl->req~usubscription-register-notifications-request-signature~1] + RESOURCE_ID_REGISTER_FOR_NOTIFICATIONS => { + Ok(USubscriptionRequest::RegisterForNotification(())) + } + // [impl->req~usubscription-unregister-notifications-request-signature~1] + RESOURCE_ID_UNREGISTER_FOR_NOTIFICATIONS => { + Ok(USubscriptionRequest::UnregisterForNotification(())) + } + // [impl->req~usubscription-reset-request-signature~1] + RESOURCE_ID_RESET => Ok(USubscriptionRequest::Reset(())), + _ => Err(UStatus::fail_with_code( + UCode::Unimplemented, + format!("unsupported uSubscription resource id: {resource_id:#06x}"), + )), + } +} + +/// Encodes a up-rust uSubscription response into its corresponding protobuf representation. +/// +/// Used by a uSubscription service to turn function response messages into raw protobuf payload data +/// +/// # Errors +/// +/// Returns an error if `resource_id` does not identify a supported operation, +/// or if the response object cannot be serialized into the expected type. +pub fn pack_usubscription_response( + usubscription_response: USubscriptionResponse, +) -> Result, UStatus> { + match usubscription_response { + // [impl->req~usubscription-subscribe-response-signature~1] + USubscriptionResponse::Subscribe(response) => Ok(Some( + UPayload::try_from_protobuf(response) + .map_err(|err| UStatus::fail_with_code(UCode::Internal, err.to_string()))?, + )), + // [impl->req~usubscription-unsubscribe-response-signature~1] + USubscriptionResponse::Unsubscribe(_) => Ok(None), + // [impl->req~usubscription-fetch-subscriptions-response-signature~1] + USubscriptionResponse::FetchSubscriptions(response) => { + let r = FetchSubscriptionsResponseProto::try_from(&response) + .map_err(|err| UStatus::fail_with_code(UCode::Internal, err.to_string()))?; + Ok(Some(UPayload::try_from_protobuf(r).map_err(|err| { + UStatus::fail_with_code(UCode::Internal, err.to_string()) + })?)) + } + // [impl->req~usubscription-register-notifications-response-signature~1] + USubscriptionResponse::RegisterForNotification(_) => Ok(None), + // [impl->req~usubscription-unregister-notifications-response-signature~1] + USubscriptionResponse::UnregisterForNotification(_) => Ok(None), + // [impl->req~usubscription-reset-response-signature~1] + USubscriptionResponse::Reset(_) => Ok(None), + } +} diff --git a/up-spec b/up-spec index 347d6843..dc2748de 160000 --- a/up-spec +++ b/up-spec @@ -1 +1 @@ -Subproject commit 347d68438afd9e4813193b243dc98fb42963ad99 +Subproject commit dc2748deb48a528778d07d40e603955ab71c33ef