diff options
| author | Sebastian Holmin <sebastian.holmin@mullvad.net> | 2025-04-09 13:47:36 +0200 |
|---|---|---|
| committer | Sebastian Holmin <sebastian.holmin@mullvad.net> | 2025-05-28 13:25:27 +0200 |
| commit | 23e7acba0f8afd4d238df067c836ae649fa80b84 (patch) | |
| tree | bd658b140a9f27066d1d7ae6788d84c3f4af2eef /mullvad-daemon/src/management_interface.rs | |
| parent | c525c77d54f5c449f872094a804d77b7e85bfc55 (diff) | |
| download | mullvadvpn-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.rs | 54 |
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), |
