wayle_network/core/connection/
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::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}