Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 26 additions & 9 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

15 changes: 13 additions & 2 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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" }
Expand Down Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ fn proto_api() -> Result<(), Box<dyn std::error::Error>> {
),
#[cfg(feature = "usubscription")]
format!(
"{}uprotocol/core/usubscription/v3/usubscription.proto",
"{}uprotocol/core/usubscription/v4/usubscription.proto",
UPROTOCOL_BASE_URI
),
];
Expand Down
151 changes: 151 additions & 0 deletions examples/usubscription_server.rs
Original file line number Diff line number Diff line change
@@ -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<HashSet<(String, String)>>,
}

#[async_trait::async_trait]
impl RequestHandler for MyUSubscriptionService {
async fn handle_request(
&self,
resource_id: u16,
message_attributes: &UAttributes,
request_payload: Option<UPayload>,
) -> Result<Option<UPayload>, 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<dyn std::error::Error>> {
// 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(())
}
68 changes: 30 additions & 38 deletions src/communication.rs
Original file line number Diff line number Diff line change
Expand Up @@ -99,44 +99,6 @@ pub(crate) fn build_message<S: crate::umessage::BuilderState>(
}
}

/// 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<Self, Self::Error> {
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 {
Expand Down Expand Up @@ -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}")
}
}
Loading
Loading