summaryrefslogtreecommitdiffhomepage
path: root/mullvad-daemon/src/management_interface.rs
diff options
context:
space:
mode:
authorSebastian Holmin <sebastian.holmin@mullvad.net>2025-04-09 13:47:36 +0200
committerSebastian Holmin <sebastian.holmin@mullvad.net>2025-05-28 13:25:27 +0200
commit23e7acba0f8afd4d238df067c836ae649fa80b84 (patch)
treebd658b140a9f27066d1d7ae6788d84c3f4af2eef /mullvad-daemon/src/management_interface.rs
parentc525c77d54f5c449f872094a804d77b7e85bfc55 (diff)
downloadmullvadvpn-23e7acba0f8afd4d238df067c836ae649fa80b84.tar.xz
mullvadvpn-23e7acba0f8afd4d238df067c836ae649fa80b84.zip
Add in app upgrades to the daemon
--------- Co-authored-by: Markus Pettersson <markus.pettersson@mullvad.net>
Diffstat (limited to 'mullvad-daemon/src/management_interface.rs')
-rw-r--r--mullvad-daemon/src/management_interface.rs54
1 files changed, 28 insertions, 26 deletions
diff --git a/mullvad-daemon/src/management_interface.rs b/mullvad-daemon/src/management_interface.rs
index 96a541f87b..764d98421a 100644
--- a/mullvad-daemon/src/management_interface.rs
+++ b/mullvad-daemon/src/management_interface.rs
@@ -38,10 +38,12 @@ pub enum Error {
SetupError(#[source] mullvad_management_interface::Error),
}
+pub type AppUpgradeBroadcast = tokio::sync::broadcast::Sender<version::AppUpgradeEvent>;
+
struct ManagementServiceImpl {
daemon_tx: DaemonCommandSender,
subscriptions: Arc<Mutex<Vec<EventsListenerSender>>>,
- app_upgrade_event_subscribers: Arc<Mutex<Vec<AppUpgradeEventListenerSender>>>,
+ pub app_upgrade_broadcast: AppUpgradeBroadcast,
}
pub type ServiceResult<T> = std::result::Result<Response<T>, Status>;
@@ -49,9 +51,7 @@ type EventsListenerReceiver = UnboundedReceiverStream<Result<types::DaemonEvent,
type EventsListenerSender = tokio::sync::mpsc::UnboundedSender<Result<types::DaemonEvent, Status>>;
type AppUpgradeEventListenerReceiver =
- UnboundedReceiverStream<Result<types::AppUpgradeEvent, Status>>;
-type AppUpgradeEventListenerSender =
- tokio::sync::mpsc::UnboundedSender<Result<types::AppUpgradeEvent, Status>>;
+ Box<dyn futures::Stream<Item = Result<types::AppUpgradeEvent, Status>> + Send + Unpin>;
const INVALID_VOUCHER_MESSAGE: &str = "This voucher code is invalid";
const USED_VOUCHER_MESSAGE: &str = "This voucher code has already been used";
@@ -1131,7 +1131,9 @@ impl ManagementService for ManagementServiceImpl {
let (tx, rx) = oneshot::channel();
self.send_command_to_daemon(DaemonCommand::AppUpgrade(tx))?;
- self.wait_for_result(rx).await?.map_err(map_daemon_error)?;
+ self.wait_for_result(rx)
+ .await?
+ .map_err(map_version_check_error)?;
Ok(Response::new(()))
}
@@ -1142,7 +1144,9 @@ impl ManagementService for ManagementServiceImpl {
let (tx, rx) = oneshot::channel();
self.send_command_to_daemon(DaemonCommand::AppUpgradeAbort(tx))?;
- self.wait_for_result(rx).await?.map_err(map_daemon_error)?;
+ self.wait_for_result(rx)
+ .await?
+ .map_err(map_version_check_error)?;
Ok(Response::new(()))
}
@@ -1151,13 +1155,19 @@ impl ManagementService for ManagementServiceImpl {
&self,
_: Request<()>,
) -> ServiceResult<Self::AppUpgradeEventsListenStream> {
- let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
-
- let mut subscriptions = self.app_upgrade_event_subscribers.lock().unwrap();
- subscriptions.push(tx);
+ log::debug!("app_upgrade_events_listen");
+ let rx = self.app_upgrade_broadcast.subscribe();
+ let upgrade_event_stream =
+ tokio_stream::wrappers::BroadcastStream::new(rx).map(|result| match result {
+ Ok(event) => Ok(event.into()),
+ Err(error) => Err(Status::internal(format!(
+ "Failed to receive app upgrade event: {error}"
+ ))),
+ });
- let upgrade_event_stream = UnboundedReceiverStream::new(rx);
- Ok(Response::new(upgrade_event_stream))
+ Ok(Response::new(
+ Box::new(upgrade_event_stream) as Self::AppUpgradeEventsListenStream
+ ))
}
}
@@ -1191,18 +1201,19 @@ impl ManagementInterfaceServer {
pub fn start(
daemon_tx: DaemonCommandSender,
rpc_socket_path: impl AsRef<Path>,
+ app_upgrade_broadcast: tokio::sync::broadcast::Sender<version::AppUpgradeEvent>,
) -> Result<ManagementInterfaceServer, Error> {
let subscriptions = Arc::<Mutex<Vec<EventsListenerSender>>>::default();
- let app_upgrade_event_subscriptions =
- Arc::<Mutex<Vec<AppUpgradeEventListenerSender>>>::default();
+
// NOTE: It is important that the channel buffer size is kept at 0. When sending a signal
// to abort the gRPC server, the sender can be awaited to know when the gRPC server has
// received and started processing the shutdown signal.
let (server_abort_tx, server_abort_rx) = mpsc::channel(0);
+
let server = ManagementServiceImpl {
daemon_tx,
subscriptions: subscriptions.clone(),
- app_upgrade_event_subscribers: app_upgrade_event_subscriptions.clone(),
+ app_upgrade_broadcast,
};
let rpc_server_join_handle = mullvad_management_interface::spawn_rpc_server(
server,
@@ -1218,10 +1229,7 @@ impl ManagementInterfaceServer {
rpc_socket_path.as_ref().display()
);
- let broadcast = ManagementInterfaceEventBroadcaster {
- subscriptions,
- app_upgrade_event_subscriptions,
- };
+ let broadcast = ManagementInterfaceEventBroadcaster { subscriptions };
Ok(ManagementInterfaceServer {
rpc_server_join_handle,
@@ -1261,7 +1269,6 @@ impl ManagementInterfaceServer {
#[derive(Clone)]
pub struct ManagementInterfaceEventBroadcaster {
subscriptions: Arc<Mutex<Vec<EventsListenerSender>>>,
- app_upgrade_event_subscriptions: Arc<Mutex<Vec<AppUpgradeEventListenerSender>>>,
}
impl ManagementInterfaceEventBroadcaster {
@@ -1270,11 +1277,6 @@ impl ManagementInterfaceEventBroadcaster {
subscriptions.retain(|tx| tx.send(Ok(value.clone())).is_ok());
}
- pub(crate) fn notify_upgrade_event(&self, value: version::AppUpgradeEvent) {
- let mut subscriptions = self.app_upgrade_event_subscriptions.lock().unwrap();
- subscriptions.retain(|tx| tx.send(Ok(value.clone().into())).is_ok());
- }
-
/// Notify that the tunnel state changed.
///
/// Sends a new state update to all `new_state` subscribers of the management interface.
@@ -1313,7 +1315,7 @@ impl ManagementInterfaceEventBroadcaster {
/// Notify that info about the latest available app version changed.
/// Or some flag about the currently running version is changed.
pub(crate) fn notify_app_version(&self, app_version_info: version::AppVersionInfo) {
- log::debug!("Broadcasting new app version info");
+ log::debug!("Broadcasting app version info:\n{app_version_info}");
self.notify(types::DaemonEvent {
event: Some(daemon_event::Event::VersionInfo(
types::AppVersionInfo::from(app_version_info),