wayle-network 0.1.2

WiFi and wired network management
Documentation
use std::sync::{Arc, Weak};

use tokio_stream::StreamExt;
use tokio_util::sync::CancellationToken;
use tracing::{debug, error};
use wayle_core::{Property, remove_and_cancel};
use wayle_traits::{ModelMonitoring, Reactive};
use zbus::{Connection, proxy::PropertyStream, zvariant::OwnedObjectPath};

use super::Wifi;
use crate::{
    core::{
        access_point::{
            AccessPoint,
            types::{LiveAccessPointParams, Ssid},
        },
        config::ip4_config::Ip4Config,
        device::wifi::DeviceWifi,
    },
    error::Error,
    proxy::{
        access_point::AccessPointProxy,
        devices::{DeviceProxy, wireless::DeviceWirelessProxy},
        manager::NetworkManagerProxy,
    },
    types::states::{NMDeviceState, NetworkStatus},
};

type SsidStream = PropertyStream<'static, Vec<u8>>;
type StrengthStream = PropertyStream<'static, u8>;

struct ActiveAccessPointStreams {
    ssid: Option<SsidStream>,
    strength: Option<StrengthStream>,
}

impl ModelMonitoring for Wifi {
    type Error = Error;

    async fn start_monitoring(self: Arc<Self>) -> Result<(), Self::Error> {
        let device_arc = Arc::new(self.device.clone());
        device_arc.start_monitoring().await?;

        let Some(ref cancellation_token) = self.device.core.cancellation_token else {
            return Err(Error::MissingCancellationToken);
        };

        let access_points = &self.access_points;
        let device = &self.device;

        populate_existing_access_points(
            &self.device.core.connection,
            device,
            access_points,
            cancellation_token,
        )
        .await;

        let cancel_token = cancellation_token.clone();
        let weak_self = Arc::downgrade(&self);

        let wireless_proxy = DeviceWirelessProxy::new(
            &self.device.core.connection,
            self.device.core.object_path.clone(),
        )
        .await
        .map_err(Error::DbusError)?;

        let device_proxy = DeviceProxy::new(
            &self.device.core.connection,
            self.device.core.object_path.clone(),
        )
        .await
        .map_err(Error::DbusError)?;

        let nm_proxy = NetworkManagerProxy::new(&self.device.core.connection)
            .await
            .map_err(Error::DbusError)?;

        tokio::spawn(async move {
            let _ = monitor_wifi(
                weak_self,
                wireless_proxy,
                device_proxy,
                nm_proxy,
                cancel_token,
            )
            .await;
        });

        Ok(())
    }
}

async fn populate_existing_access_points(
    connection: &Connection,
    device: &DeviceWifi,
    access_points: &Property<Vec<Arc<AccessPoint>>>,
    cancellation_token: &CancellationToken,
) {
    let existing_paths = device.access_points.get();
    let mut initial_aps = Vec::with_capacity(existing_paths.len());

    for ap_path in existing_paths {
        let Ok(path) = OwnedObjectPath::try_from(ap_path.as_str()) else {
            continue;
        };

        if let Ok(ap) = AccessPoint::get_live(LiveAccessPointParams {
            connection,
            path,
            cancellation_token,
        })
        .await
        {
            initial_aps.push(ap);
        }
    }

    if !initial_aps.is_empty() {
        access_points.set(initial_aps);
    }
}

#[allow(clippy::too_many_lines)]
async fn monitor_wifi(
    weak_wifi: Weak<Wifi>,
    wireless_proxy: DeviceWirelessProxy<'static>,
    device_proxy: DeviceProxy<'static>,
    nm_proxy: NetworkManagerProxy<'static>,
    cancellation_token: CancellationToken,
) -> Result<(), Error> {
    let mut ap_added = wireless_proxy
        .receive_access_point_added()
        .await
        .map_err(Error::DbusError)?;

    let mut ap_removed = wireless_proxy
        .receive_access_point_removed()
        .await
        .map_err(Error::DbusError)?;

    let mut enabled_changed = nm_proxy.receive_wireless_enabled_changed().await;
    let mut access_point_changed = wireless_proxy.receive_active_access_point_changed().await;
    let mut connectivity_changed = device_proxy.receive_state_changed().await;
    let mut ip4_config_changed = device_proxy.receive_ip4_config_changed().await;

    let ActiveAccessPointStreams {
        ssid: mut ap_ssid_stream,
        strength: mut ap_strength_stream,
    } = {
        let Some(wifi) = weak_wifi.upgrade() else {
            error!("cannot upgrade weak wifi reference, access point monitoring may be degraded");
            return Err(Error::OperationFailed {
                operation: "monitor wifi",
                source: "weak reference dropped".into(),
            });
        };

        handle_access_point_changed(
            &wifi.device.core.connection,
            wifi.device.active_access_point.get(),
            &wifi.ssid,
            &wifi.strength,
            &wifi.frequency,
        )
        .await
    };

    tokio::spawn(async move {
        loop {
            let Some(wifi) = weak_wifi.upgrade() else {
                return;
            };

            tokio::select! {
                _ = cancellation_token.cancelled() => {
                    debug!("WifiMonitor cancelled");
                    return;
                }

                Some(added) = ap_added.next() => {
                    if let Ok(args) = added.args() {
                        handle_ap_added(
                            &wifi.device.core.connection,
                            args.access_point,
                            &wifi.access_points,
                            &cancellation_token
                        ).await;
                    }
                }

                Some(removed) = ap_removed.next() => {
                    if let Ok(args) = removed.args() {
                        handle_ap_removed(&args.access_point, &wifi.access_points);
                    }
                }

                Some(change) = enabled_changed.next() => {
                    if let Ok(new_state) = change.get().await {
                        wifi.enabled.set(new_state);
                    }
                }

                Some(change) = access_point_changed.next() => {
                    let Ok(new_ap_path) = change.get().await else {
                        continue;
                    };

                    let streams = handle_access_point_changed(
                        &wifi.device.core.connection,
                        new_ap_path,
                        &wifi.ssid,
                        &wifi.strength,
                        &wifi.frequency,
                    ).await;

                    ap_ssid_stream = streams.ssid;
                    ap_strength_stream = streams.strength;
                }

                Some(change) = async { ap_ssid_stream.as_mut()?.next().await } => {
                    if let Ok(new_ssid) = change.get().await {
                        wifi.ssid.set(Some(Ssid::new(new_ssid).to_string()));
                    }
                }

                Some(change) = async { ap_strength_stream.as_mut()?.next().await } => {
                    if let Ok(new_strength) = change.get().await {
                        wifi.strength.set(Some(new_strength));
                    }
                }

                Some(change) = connectivity_changed.next() => {
                    handle_connectivity_changed(change, &wifi).await;
                }

                Some(change) = ip4_config_changed.next() => {
                    handle_ip4_config_changed(change, &wifi).await;
                }

                else => {
                    break;
                }
            }
        }
    });

    Ok(())
}

async fn handle_connectivity_changed(change: zbus::proxy::PropertyChanged<'_, u32>, wifi: &Wifi) {
    let Ok(new_connectivity) = change.get().await else {
        return;
    };

    let device_state = NMDeviceState::from_u32(new_connectivity);
    let status = NetworkStatus::from_device_state(device_state);

    wifi.connectivity.set(status);

    if status == NetworkStatus::Connected {
        let ip = Ip4Config::resolve_address(
            &wifi.device.core.connection,
            wifi.device.core.ip4_config.get(),
        )
        .await;

        wifi.ip4_address.set(ip);
    }
}

async fn handle_ip4_config_changed(
    change: zbus::proxy::PropertyChanged<'_, OwnedObjectPath>,
    wifi: &Wifi,
) {
    let Ok(ip4_path) = change.get().await else {
        return;
    };

    let ip = Ip4Config::resolve_address(&wifi.device.core.connection, ip4_path).await;

    wifi.ip4_address.set(ip);
}

async fn handle_ap_added(
    connection: &Connection,
    ap_path: OwnedObjectPath,
    access_points: &Property<Vec<Arc<AccessPoint>>>,
    cancellation_token: &CancellationToken,
) {
    if let Ok(new_ap) = AccessPoint::get_live(LiveAccessPointParams {
        connection,
        path: ap_path,
        cancellation_token,
    })
    .await
    {
        let mut aps = access_points.get();
        aps.push(new_ap);

        access_points.set(aps);
    }
}

fn handle_ap_removed(ap_path: &OwnedObjectPath, access_points: &Property<Vec<Arc<AccessPoint>>>) {
    remove_and_cancel!(access_points, ap_path.clone());
}

async fn handle_access_point_changed(
    connection: &Connection,
    new_ap_path: OwnedObjectPath,
    ssid_prop: &Property<Option<String>>,
    strength_prop: &Property<Option<u8>>,
    frequency_prop: &Property<Option<u32>>,
) -> ActiveAccessPointStreams {
    if new_ap_path.is_empty() || new_ap_path == OwnedObjectPath::default() {
        ssid_prop.set(None);
        strength_prop.set(None);
        frequency_prop.set(None);

        return ActiveAccessPointStreams {
            ssid: None,
            strength: None,
        };
    }

    match AccessPointProxy::new(connection, new_ap_path).await {
        Ok(ap_proxy) => {
            if let Ok(raw_ssid) = ap_proxy.ssid().await {
                ssid_prop.set(Some(Ssid::new(raw_ssid).to_string()));
            }

            if let Ok(strength) = ap_proxy.strength().await {
                strength_prop.set(Some(strength));
            }

            if let Ok(freq) = ap_proxy.frequency().await {
                frequency_prop.set(Some(freq));
            }

            ActiveAccessPointStreams {
                ssid: Some(ap_proxy.receive_ssid_changed().await),
                strength: Some(ap_proxy.receive_strength_changed().await),
            }
        }
        Err(_) => {
            ssid_prop.set(None);
            strength_prop.set(None);
            frequency_prop.set(None);

            ActiveAccessPointStreams {
                ssid: None,
                strength: None,
            }
        }
    }
}