use std::{ borrow::{Borrow, Cow}, collections::HashMap, pin::Pin, sync::Arc, time::Duration, }; use futures_util::stream::BoxStream; use multi_index_map::MultiIndexMap; use ordered_stream::OrderedStream; use tokio::sync::{Mutex, Notify}; use tokio_stream::StreamMap; use tokio_util::sync::{CancellationToken, DropGuard}; use tracing::{debug, info, warn}; use zbus::{proxy::Defaults, zvariant}; use crate::stratis::metadata::PoolMetadata; use super::dbus::{FilesystemProxy, PoolProxy, check_stratisd_version}; #[derive(MultiIndexMap, Debug, Clone)] #[multi_index_derive(Debug)] pub struct CachedPool { #[multi_index(hashed_unique)] pub path: zvariant::OwnedObjectPath, #[multi_index(hashed_unique)] pub name: String, pub proxy: PoolProxy<'static>, _properties_watcher: Arc, } #[derive(Debug, Clone, PartialEq, Eq, Hash)] pub struct FilesystemKey<'path, 'name> { pub pool: zvariant::ObjectPath<'path>, pub name: Cow<'name, str>, } #[derive(Debug, Clone, PartialEq, Eq, Hash)] pub struct OwnedFilesystemKey(pub FilesystemKey<'static, 'static>); impl From> for OwnedFilesystemKey { fn from(value: FilesystemKey<'static, 'static>) -> Self { Self(value) } } impl From for FilesystemKey<'static, 'static> { fn from(value: OwnedFilesystemKey) -> Self { value.0 } } impl<'path, 'name> Borrow> for OwnedFilesystemKey { fn borrow(&self) -> &FilesystemKey<'path, 'name> { &self.0 } } #[derive(MultiIndexMap, Debug, Clone)] #[multi_index_derive(Debug)] pub struct CachedFilesystem { #[multi_index(hashed_unique)] pub path: zvariant::OwnedObjectPath, #[multi_index(hashed_unique)] pub key: OwnedFilesystemKey, #[multi_index(hashed_non_unique)] pub pool: zvariant::OwnedObjectPath, pub proxy: FilesystemProxy<'static>, _properties_watcher: Arc, } pub struct ObjectCache { conn: zbus::Connection, pools: MultiIndexCachedPoolMap, filesystems: MultiIndexCachedFilesystemMap, any_pool_size_changed: Arc, } pub struct ObjectCacheNotifier { any_pool_size_changed: Arc, } impl ObjectCacheNotifier { // XXX: we use notify_one() which will notify a SINGLE waiter, so make this // &mut to avoid multiple waiters trampling on each other. pub async fn any_pool_size_changed(&mut self) { self.any_pool_size_changed.notified().await; } } fn get_prop< 'props, 'v, S: std::borrow::Borrow + std::cmp::Eq + std::hash::Hash, V: Borrow>, >( path: &zvariant::ObjectPath<'_>, props: &'props HashMap, name: &str, ) -> zvariant::Result<&'props zvariant::Value<'v>> { props .get(name) .map(Borrow::borrow) .ok_or_else(|| zvariant::Error::Message(format!("{} is missing {name}", path))) } #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] enum ObjectEvent { Renamed, SizeChanged, } impl ObjectCache { fn watch_object_events + Send + 'static>( &self, proxy: P, token: CancellationToken, mut streams: StreamMap>, current_name: String, get_cached_name: impl Fn(&P) -> zbus::Result> + Send + 'static, ) { let notify = self.any_pool_size_changed.clone(); tokio::task::spawn(token.run_until_cancelled_owned(async move { use futures_util::StreamExt; while let Some((event, ())) = streams.next().await { match event { ObjectEvent::Renamed => { let new_name = get_cached_name(&proxy).ok().flatten(); let new_name = new_name.as_deref().unwrap_or(""); if current_name != new_name { warn!( "unexpected rename of object {}: '{current_name}' -> '{new_name}'", proxy.inner().path() ); } } ObjectEvent::SizeChanged => { debug!( "size changed of object {}, notifying waiters", proxy.inner().path() ); notify.notify_one(); } } } info!(path = %proxy.inner().path(), "No more events available for object"); })); } async fn add_pool_from_props< 'v, S: std::borrow::Borrow + std::cmp::Eq + std::hash::Hash + std::fmt::Debug, V: Borrow> + std::fmt::Debug, >( &mut self, path: Cow<'_, zvariant::ObjectPath<'_>>, props: &HashMap, ) -> zbus::Result<()> { use futures_util::StreamExt; info!(?props, "add pool"); let name: String = get_prop(&path, props, "Name")?.try_into()?; let mver: u64 = get_prop(&path, props, "MetadataVersion")?.try_into()?; if mver != PoolMetadata::VERSION { return Err(zbus::Error::Failure(format!( "unsupported metadata version: {path} ({name}): {mver}" ))); } let proxy = PoolProxy::new(&self.conn, path.clone().into_owned().to_owned()) .await? .to_owned(); let token = CancellationToken::new(); if let Err(err) = self.pools.try_insert(CachedPool { path: path.into_owned().to_owned().into(), name: name.clone(), proxy: proxy.clone(), _properties_watcher: Arc::new(token.clone().drop_guard()), }) { let current = err.0; if !self .pools .get_by_path(¤t.path) .or_else(|| self.pools.get_by_name(¤t.name)) .map(|existing| current.path == existing.path && current.name == existing.name) .unwrap_or_default() { return Err(zbus::Error::Failure(format!( "duplicate interface: {} ({})", current.path, current.name ))); } return Ok(()); } let mut streams = StreamMap::new(); streams.insert( ObjectEvent::SizeChanged, proxy .receive_total_physical_size_changed() .await .map(|_| ()) .boxed(), ); streams.insert( ObjectEvent::SizeChanged, proxy .receive_total_physical_used_changed() .await .map(|_| ()) .boxed(), ); self.watch_object_events(proxy, token, streams, name, PoolProxy::cached_name); self.any_pool_size_changed.notify_one(); Ok(()) } async fn add_filesystem_from_props< 'v, S: std::borrow::Borrow + std::cmp::Eq + std::hash::Hash + std::fmt::Debug, V: Borrow> + std::fmt::Debug, >( &mut self, path: Cow<'_, zvariant::ObjectPath<'_>>, props: &HashMap, ) -> zbus::Result<()> { use futures_util::StreamExt; info!(?props, "add filesystem"); let name: String = get_prop(&path, props, "Name")?.try_into()?; let pool: zvariant::ObjectPath<'_> = get_prop(&path, props, "Pool")?.try_into()?; let proxy = FilesystemProxy::new(&self.conn, path.clone().into_owned().to_owned()) .await? .to_owned(); let token = CancellationToken::new(); if let Err(err) = self.filesystems.try_insert(CachedFilesystem { path: path.into_owned().to_owned().into(), pool: pool.clone().into_owned().into(), key: FilesystemKey { pool: pool.into_owned(), name: name.clone().into(), } .into(), proxy: proxy.clone(), _properties_watcher: Arc::new(token.clone().drop_guard()), }) { let current = err.0; if !self .filesystems .get_by_path(¤t.path) .or_else(|| self.filesystems.get_by_key(¤t.key)) .map(|existing| current.path == existing.path && current.key == existing.key) .unwrap_or_default() { return Err(zbus::Error::Failure(format!( "duplicate interface: {} ({})", current.path, current.key.0.name ))); } return Ok(()); } let mut streams = StreamMap::new(); streams.insert( ObjectEvent::Renamed, proxy.receive_name_changed().await.map(|_| ()).boxed(), ); streams.insert( ObjectEvent::SizeChanged, proxy.receive_used_changed().await.map(|_| ()).boxed(), ); streams.insert( ObjectEvent::SizeChanged, proxy.receive_size_limit_changed().await.map(|_| ()).boxed(), ); self.watch_object_events(proxy, token, streams, name, FilesystemProxy::cached_name); self.any_pool_size_changed.notify_one(); Ok(()) } async fn add_object_from_props< 'i, 'v, S: std::borrow::Borrow + std::cmp::Eq + std::hash::Hash + std::fmt::Debug, I: std::borrow::Borrow + std::cmp::Eq + std::hash::Hash, V: Borrow> + std::fmt::Debug, >( &mut self, path: Cow<'_, zvariant::ObjectPath<'_>>, interfaces: &HashMap>, ) -> zbus::Result<()> { if let Some(props) = interfaces.get(PoolProxy::<'static>::INTERFACE.as_ref().unwrap().as_str()) { self.add_pool_from_props(path, props).await?; } else if let Some(props) = interfaces.get( FilesystemProxy::<'static>::INTERFACE .as_ref() .unwrap() .as_str(), ) { self.add_filesystem_from_props(path, props).await?; } Ok(()) } async fn on_interfaces_added( &mut self, signal: zbus::fdo::InterfacesAdded, ) -> zbus::Result<()> { let args = signal.args()?; debug!(?args, "InterfacesAdded"); self.add_object_from_props( Cow::Borrowed(args.object_path()), args.interfaces_and_properties(), ) .await?; Ok(()) } async fn on_interfaces_removed( &mut self, signal: zbus::fdo::InterfacesRemoved, ) -> zbus::Result<()> { let args = signal.args()?; debug!(?args, "InterfacesRemoved"); for iface in args.interfaces().iter() { if iface == PoolProxy::<'static>::INTERFACE.as_ref().unwrap() { if self .pools .remove_by_path(&args.object_path().to_owned().into()) .is_none() { return Err(zbus::Error::Failure(format!( "removed unknown interface {}", args.object_path() ))); } self.any_pool_size_changed.notify_one(); break; } else if iface == FilesystemProxy::<'static>::INTERFACE.as_ref().unwrap() { if self .filesystems .remove_by_path(&args.object_path().to_owned().into()) .is_none() { return Err(zbus::Error::Failure(format!( "removed unknown interface {}", args.object_path() ))); } self.any_pool_size_changed.notify_one(); break; } } Ok(()) } async fn on_new_stratisd( &mut self, om: &zbus::fdo::ObjectManagerProxy<'_>, ) -> zbus::Result<()> { check_stratisd_version(&self.conn).await?; self.pools.clear(); self.filesystems.clear(); // Stratis can take a bit to start up, so try to update the objects a // few times. let mut previous_object_count = 0; for _ in 0..5 { let objects = om.get_managed_objects().await?; if objects.len() == previous_object_count { break; } previous_object_count = objects.len(); for (path, interfaces) in objects { self.add_object_from_props(Cow::Owned(path.into_inner()), &interfaces) .await?; } tokio::time::sleep(Duration::from_secs(1)).await; } Ok(()) } pub async fn start( conn: zbus::Connection, token: CancellationToken, ) -> zbus::Result<(Arc>, ObjectCacheNotifier)> { use ordered_stream::OrderedStreamExt; let om = zbus::fdo::ObjectManagerProxy::new( &conn, "org.storage.stratis3", "/org/storage/stratis3", ) .await?; enum Event { Added(zbus::fdo::InterfacesAdded), Removed(zbus::fdo::InterfacesRemoved), OwnerChanged(Option>), } let c = om .inner() .receive_owner_changed() .await? .map(Event::OwnerChanged); let mut events = ordered_stream::join( ordered_stream::join( om.receive_interfaces_added().await?.map(Event::Added), om.receive_interfaces_removed().await?.map(Event::Removed), ), Box::pin(c) as Pin + Send>>, ); let notify = Arc::new(Notify::new()); let mut cache = ObjectCache { conn, pools: MultiIndexCachedPoolMap::default(), filesystems: MultiIndexCachedFilesystemMap::default(), any_pool_size_changed: notify.clone(), }; for (path, interfaces) in om.get_managed_objects().await? { cache .add_object_from_props(Cow::Owned(path.into_inner()), &interfaces) .await?; } let cache = Arc::new(Mutex::new(cache)); let cache_2 = cache.clone(); tokio::task::spawn(token.run_until_cancelled_owned(async move { loop { let Some(event) = events.next().await else { panic!("no events from the bus left?"); }; if let Err(err) = match event { Event::Added(signal) => cache_2.lock().await.on_interfaces_added(signal).await, Event::Removed(signal) => { cache_2.lock().await.on_interfaces_removed(signal).await } Event::OwnerChanged(None) => { warn!("stratisd has disappeared from the bus"); Ok(()) } Event::OwnerChanged(Some(_)) => { warn!("stratisd has returned, resetting state"); cache_2.lock().await.on_new_stratisd(&om).await } } { warn!(?err, "handling bus event"); } } })); Ok(( cache, ObjectCacheNotifier { any_pool_size_changed: notify, }, )) } pub fn get_pool_by_name(&self, name: &str) -> Option { self.pools.get_by_name(name).cloned() } pub fn get_filesystem_by_key(&self, key: &FilesystemKey<'_, '_>) -> Option { self.filesystems.get_by_key(key).cloned() } pub fn list_filesystems_in_pool( &self, pool: &zvariant::ObjectPath<'_>, ) -> Vec { self.filesystems .get_by_pool(&pool.as_ref()) .into_iter() .cloned() .collect() } }