wayle_network/core/device/
monitoring.rs1use std::sync::{Arc, Weak};
2
3use futures::StreamExt;
4use tokio_util::sync::CancellationToken;
5use tracing::debug;
6use wayle_traits::ModelMonitoring;
7
8use super::Device;
9use crate::{
10 error::Error,
11 proxy::devices::DeviceProxy,
12 types::{
13 connectivity::{NMConnectivityState, NMMetered},
14 device::NMDeviceType,
15 flags::{NMDeviceCapabilities, NMDeviceInterfaceFlags},
16 states::{NMDeviceState, NMDeviceStateReason},
17 },
18};
19
20impl ModelMonitoring for Device {
21 type Error = Error;
22
23 async fn start_monitoring(self: Arc<Self>) -> Result<(), Self::Error> {
24 let proxy = DeviceProxy::new(&self.connection, self.object_path.clone())
25 .await
26 .map_err(Error::DbusError)?;
27
28 let Some(ref cancellation_token) = self.cancellation_token else {
29 return Err(Error::MissingCancellationToken);
30 };
31
32 let cancel_token = cancellation_token.clone();
33 let weak_self = Arc::downgrade(&self);
34
35 tokio::spawn(async move {
36 monitor(weak_self, proxy, cancel_token).await;
37 });
38
39 Ok(())
40 }
41}
42
43#[allow(clippy::cognitive_complexity)]
44#[allow(clippy::too_many_lines)]
45async fn monitor(
46 weak_device: Weak<Device>,
47 proxy: DeviceProxy<'static>,
48 cancellation_token: CancellationToken,
49) {
50 let mut udi_changed = proxy.receive_udi_changed().await;
51 let mut udev_path_changed = proxy.receive_path_changed().await;
52 let mut interface_changed = proxy.receive_interface_changed().await;
53 let mut ip_interface_changed = proxy.receive_ip_interface_changed().await;
54 let mut driver_changed = proxy.receive_driver_changed().await;
55 let mut driver_version_changed = proxy.receive_driver_version_changed().await;
56 let mut firmware_version_changed = proxy.receive_firmware_version_changed().await;
57 let mut capabilities_changed = proxy.receive_capabilities_changed().await;
58 let mut state_changed = proxy.receive_state_changed().await;
59 let mut state_reason_changed = proxy.receive_state_reason_changed().await;
60 let mut active_connection_changed = proxy.receive_active_connection_changed().await;
61 let mut ip4_config_changed = proxy.receive_ip4_config_changed().await;
62 let mut dhcp4_config_changed = proxy.receive_dhcp4_config_changed().await;
63 let mut ip6_config_changed = proxy.receive_ip6_config_changed().await;
64 let mut dhcp6_config_changed = proxy.receive_dhcp6_config_changed().await;
65 let mut managed_changed = proxy.receive_managed_changed().await;
66 let mut autoconnect_changed = proxy.receive_autoconnect_changed().await;
67 let mut firmware_missing_changed = proxy.receive_firmware_missing_changed().await;
68 let mut nm_plugin_missing_changed = proxy.receive_nm_plugin_missing_changed().await;
69 let mut device_type_changed = proxy.receive_device_type_changed().await;
70 let mut available_connections_changed = proxy.receive_available_connections_changed().await;
71 let mut physical_port_id_changed = proxy.receive_physical_port_id_changed().await;
72 let mut mtu_changed = proxy.receive_mtu_changed().await;
73 let mut metered_changed = proxy.receive_metered_changed().await;
74 let mut real_changed = proxy.receive_real_changed().await;
75 let mut ip4_connectivity_changed = proxy.receive_ip4_connectivity_changed().await;
76 let mut ip6_connectivity_changed = proxy.receive_ip6_connectivity_changed().await;
77 let mut interface_flags_changed = proxy.receive_interface_flags_changed().await;
78 let mut hw_address_changed = proxy.receive_hw_address_changed().await;
79 let mut ports_changed = proxy.receive_ports_changed().await;
80
81 loop {
82 let Some(device) = weak_device.upgrade() else {
83 return;
84 };
85
86 tokio::select! {
87 _ = cancellation_token.cancelled() => {
88 debug!("DeviceMonitor cancelled");
89 return;
90 }
91 Some(change) = udi_changed.next() => {
92 if let Ok(value) = change.get().await {
93 device.udi.set(value);
94 }
95 }
96 Some(change) = udev_path_changed.next() => {
97 if let Ok(value) = change.get().await {
98 device.udev_path.set(value);
99 }
100 }
101 Some(change) = interface_changed.next() => {
102 if let Ok(value) = change.get().await {
103 device.interface.set(value);
104 }
105 }
106 Some(change) = ip_interface_changed.next() => {
107 if let Ok(value) = change.get().await {
108 device.ip_interface.set(value);
109 }
110 }
111 Some(change) = driver_changed.next() => {
112 if let Ok(value) = change.get().await {
113 device.driver.set(value);
114 }
115 }
116 Some(change) = driver_version_changed.next() => {
117 if let Ok(value) = change.get().await {
118 device.driver_version.set(value);
119 }
120 }
121 Some(change) = firmware_version_changed.next() => {
122 if let Ok(value) = change.get().await {
123 device.firmware_version.set(value);
124 }
125 }
126 Some(change) = capabilities_changed.next() => {
127 if let Ok(value) = change.get().await {
128 device.capabilities.set(NMDeviceCapabilities::from_bits_truncate(value));
129 }
130 }
131 Some(change) = state_changed.next() => {
132 if let Ok(value) = change.get().await {
133 device.state.set(NMDeviceState::from_u32(value));
134 }
135 }
136 Some(change) = state_reason_changed.next() => {
137 if let Ok((state, reason)) = change.get().await {
138 device.state_reason.set((
139 NMDeviceState::from_u32(state),
140 NMDeviceStateReason::from_u32(reason)
141 ));
142 }
143 }
144 Some(change) = active_connection_changed.next() => {
145 if let Ok(value) = change.get().await {
146 device.active_connection.set(value);
147 }
148 }
149 Some(change) = ip4_config_changed.next() => {
150 if let Ok(value) = change.get().await {
151 device.ip4_config.set(value);
152 }
153 }
154 Some(change) = dhcp4_config_changed.next() => {
155 if let Ok(value) = change.get().await {
156 device.dhcp4_config.set(value);
157 }
158 }
159 Some(change) = ip6_config_changed.next() => {
160 if let Ok(value) = change.get().await {
161 device.ip6_config.set(value);
162 }
163 }
164 Some(change) = dhcp6_config_changed.next() => {
165 if let Ok(value) = change.get().await {
166 device.dhcp6_config.set(value);
167 }
168 }
169 Some(change) = managed_changed.next() => {
170 if let Ok(value) = change.get().await {
171 device.managed.set(value);
172 }
173 }
174 Some(change) = autoconnect_changed.next() => {
175 if let Ok(value) = change.get().await {
176 device.autoconnect.set(value);
177 }
178 }
179 Some(change) = firmware_missing_changed.next() => {
180 if let Ok(value) = change.get().await {
181 device.firmware_missing.set(value);
182 }
183 }
184 Some(change) = nm_plugin_missing_changed.next() => {
185 if let Ok(value) = change.get().await {
186 device.nm_plugin_missing.set(value);
187 }
188 }
189 Some(change) = device_type_changed.next() => {
190 if let Ok(value) = change.get().await {
191 device.device_type.set(NMDeviceType::from_u32(value));
192 }
193 }
194 Some(change) = available_connections_changed.next() => {
195 if let Ok(value) = change.get().await {
196 device.available_connections.set(value);
197 }
198 }
199 Some(change) = physical_port_id_changed.next() => {
200 if let Ok(value) = change.get().await {
201 device.physical_port_id.set(value);
202 }
203 }
204 Some(change) = mtu_changed.next() => {
205 if let Ok(value) = change.get().await {
206 device.mtu.set(value);
207 }
208 }
209 Some(change) = metered_changed.next() => {
210 if let Ok(value) = change.get().await {
211 device.metered.set(NMMetered::from_u32(value));
212 }
213 }
214 Some(change) = real_changed.next() => {
215 if let Ok(value) = change.get().await {
216 device.real.set(value);
217 }
218 }
219 Some(change) = ip4_connectivity_changed.next() => {
220 if let Ok(value) = change.get().await {
221 device.ip4_connectivity.set(NMConnectivityState::from_u32(value));
222 }
223 }
224 Some(change) = ip6_connectivity_changed.next() => {
225 if let Ok(value) = change.get().await {
226 device.ip6_connectivity.set(NMConnectivityState::from_u32(value));
227 }
228 }
229 Some(change) = interface_flags_changed.next() => {
230 if let Ok(value) = change.get().await {
231 device.interface_flags.set(NMDeviceInterfaceFlags::from_bits_truncate(value));
232 }
233 }
234 Some(change) = hw_address_changed.next() => {
235 if let Ok(value) = change.get().await {
236 device.hw_address.set(value);
237 }
238 }
239 Some(change) = ports_changed.next() => {
240 if let Ok(value) = change.get().await {
241 device.ports.set(value);
242 }
243 }
244
245 else => {
246 debug!("All property streams ended for device");
247 break;
248 }
249 }
250 }
251
252 debug!("Property monitoring ended for device");
253}