Skip to main content

wayle_network/core/connection/
monitoring.rs

1use std::sync::{Arc, Weak};
2
3use futures::StreamExt;
4use tokio_util::sync::CancellationToken;
5use tracing::debug;
6use wayle_traits::ModelMonitoring;
7
8use super::ActiveConnection;
9use crate::{
10    error::Error,
11    proxy::active_connection::ConnectionActiveProxy,
12    types::{flags::NMActivationStateFlags, states::NMActiveConnectionState},
13};
14
15impl ModelMonitoring for ActiveConnection {
16    type Error = Error;
17
18    async fn start_monitoring(self: Arc<Self>) -> Result<(), Self::Error> {
19        let proxy = ConnectionActiveProxy::new(&self.zbus_connection, self.object_path.clone())
20            .await
21            .map_err(Error::DbusError)?;
22
23        let Some(ref cancellation_token) = self.cancellation_token else {
24            return Err(Error::MissingCancellationToken);
25        };
26
27        let cancel_token = cancellation_token.clone();
28        let weak_self = Arc::downgrade(&self);
29
30        tokio::spawn(async move {
31            monitor(weak_self, proxy, cancel_token).await;
32        });
33
34        Ok(())
35    }
36}
37
38#[allow(clippy::cognitive_complexity)]
39#[allow(clippy::too_many_lines)]
40async fn monitor(
41    weak_active_connection: Weak<ActiveConnection>,
42    proxy: ConnectionActiveProxy<'static>,
43    cancellation_token: CancellationToken,
44) {
45    let mut connection_changes = proxy.receive_connection_changed().await;
46    let mut specific_object_changes = proxy.receive_specific_object_changed().await;
47    let mut id_changes = proxy.receive_id_changed().await;
48    let mut uuid_changed = proxy.receive_uuid_changed().await;
49    let mut type_changed = proxy.receive_type__changed().await;
50    let mut devices_changed = proxy.receive_devices_changed().await;
51    let mut state_changed = proxy.receive_state_changed().await;
52    let mut state_flags_changed = proxy.receive_state_flags_changed().await;
53    let mut default_changed = proxy.receive_default_changed().await;
54    let mut ip4_config_changed = proxy.receive_ip4_config_changed().await;
55    let mut dhcp4_config_changed = proxy.receive_dhcp4_config_changed().await;
56    let mut default6_changed = proxy.receive_default6_changed().await;
57    let mut ip6_config_changed = proxy.receive_ip6_config_changed().await;
58    let mut dhcp6_config_changed = proxy.receive_dhcp6_config_changed().await;
59    let mut vpn_changed = proxy.receive_vpn_changed().await;
60    let mut controller_changed = proxy.receive_controller_changed().await;
61
62    loop {
63        let Some(active_connection) = weak_active_connection.upgrade() else {
64            return;
65        };
66
67        tokio::select! {
68            _ = cancellation_token.cancelled() => {
69                debug!("ActiveConnection monitoring cancelled for {}", active_connection.object_path);
70                return;
71            }
72            Some(change) = connection_changes.next() => {
73                if let Ok(new_connection) = change.get().await {
74                    active_connection.connection_path.set(new_connection);
75                }
76            }
77            Some(change) = specific_object_changes.next() => {
78                if let Ok(new_specific_object) = change.get().await {
79                    active_connection.specific_object.set(new_specific_object);
80                }
81            }
82            Some(change) = id_changes.next() => {
83                if let Ok(new_id) = change.get().await {
84                    active_connection.id.set(new_id);
85                }
86            }
87            Some(change) = uuid_changed.next() => {
88                if let Ok(new_uuid) = change.get().await {
89                    active_connection.uuid.set(new_uuid);
90                }
91            }
92            Some(change) = type_changed.next() => {
93                if let Ok(new_type) = change.get().await {
94                    active_connection.type_.set(new_type);
95                }
96            }
97            Some(change) = devices_changed.next() => {
98                if let Ok(new_devices) = change.get().await {
99                    active_connection.devices.set(new_devices);
100                }
101            }
102            Some(change) = state_changed.next() => {
103                if let Ok(new_state) = change.get().await {
104                    let state = NMActiveConnectionState::from_u32(new_state);
105                    active_connection.state.set(state);
106                }
107            }
108            Some(change) = state_flags_changed.next() => {
109                if let Ok(new_flags) = change.get().await {
110                    let flags = NMActivationStateFlags::from_bits_truncate(new_flags);
111                    active_connection.state_flags.set(flags);
112                }
113            }
114            Some(change) = default_changed.next() => {
115                if let Ok(new_default) = change.get().await {
116                    active_connection.default.set(new_default);
117                }
118            }
119            Some(change) = ip4_config_changed.next() => {
120                if let Ok(new_ip4_config) = change.get().await {
121                    active_connection.ip4_config.set(new_ip4_config);
122                }
123            }
124            Some(change) = dhcp4_config_changed.next() => {
125                if let Ok(new_dhcp4_config) = change.get().await {
126                    active_connection.dhcp4_config.set(new_dhcp4_config);
127                }
128            }
129            Some(change) = default6_changed.next() => {
130                if let Ok(new_default6) = change.get().await {
131                    active_connection.default6.set(new_default6);
132                }
133            }
134            Some(change) = ip6_config_changed.next() => {
135                if let Ok(new_ip6_config) = change.get().await {
136                    active_connection.ip6_config.set(new_ip6_config);
137                }
138            }
139            Some(change) = dhcp6_config_changed.next() => {
140                if let Ok(new_dhcp6_config) = change.get().await {
141                    active_connection.dhcp6_config.set(new_dhcp6_config);
142                }
143            }
144            Some(change) = vpn_changed.next() => {
145                if let Ok(new_vpn) = change.get().await {
146                    active_connection.vpn.set(new_vpn);
147                }
148            }
149            Some(change) = controller_changed.next() => {
150                if let Ok(new_controller) = change.get().await {
151                    active_connection.controller.set(new_controller);
152                }
153            }
154            else => {
155                debug!("All property streams ended for active connection");
156                break;
157            }
158        }
159    }
160
161    debug!("Property monitoring ended for active connection");
162}