Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
16 kB · 508 lines
Rust
at ci
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509use 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<DropGuard>,}
#[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<FilesystemKey<'static, 'static>> for OwnedFilesystemKey { fn from(value: FilesystemKey<'static, 'static>) -> Self { Self(value) }}
impl From<OwnedFilesystemKey> for FilesystemKey<'static, 'static> { fn from(value: OwnedFilesystemKey) -> Self { value.0 }}
impl<'path, 'name> Borrow<FilesystemKey<'path, 'name>> 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<DropGuard>,}
pub struct ObjectCache { conn: zbus::Connection, pools: MultiIndexCachedPoolMap, filesystems: MultiIndexCachedFilesystemMap, any_pool_size_changed: Arc<Notify>,}
pub struct ObjectCacheNotifier { any_pool_size_changed: Arc<Notify>,}
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<str> + std::cmp::Eq + std::hash::Hash, V: Borrow<zvariant::Value<'v>>,>( path: &zvariant::ObjectPath<'_>, props: &'props HashMap<S, V>, 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<P: zbus::proxy::ProxyImpl<'static> + Send + 'static>( &self, proxy: P, token: CancellationToken, mut streams: StreamMap<ObjectEvent, BoxStream<'static, ()>>, current_name: String, get_cached_name: impl Fn(&P) -> zbus::Result<Option<String>> + 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("<unknown>"); 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<str> + std::cmp::Eq + std::hash::Hash + std::fmt::Debug, V: Borrow<zvariant::Value<'v>> + std::fmt::Debug, >( &mut self, path: Cow<'_, zvariant::ObjectPath<'_>>, props: &HashMap<S, V>, ) -> 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<str> + std::cmp::Eq + std::hash::Hash + std::fmt::Debug, V: Borrow<zvariant::Value<'v>> + std::fmt::Debug, >( &mut self, path: Cow<'_, zvariant::ObjectPath<'_>>, props: &HashMap<S, V>, ) -> 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<str> + std::cmp::Eq + std::hash::Hash + std::fmt::Debug, I: std::borrow::Borrow<str> + std::cmp::Eq + std::hash::Hash, V: Borrow<zvariant::Value<'v>> + std::fmt::Debug, >( &mut self, path: Cow<'_, zvariant::ObjectPath<'_>>, interfaces: &HashMap<I, HashMap<S, V>>, ) -> 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<Mutex<Self>>, 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<zbus::names::UniqueName<'static>>), }
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<Box<dyn OrderedStream<Ordering = _, Data = _> + 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<CachedPool> { self.pools.get_by_name(name).cloned() }
pub fn get_filesystem_by_key(&self, key: &FilesystemKey<'_, '_>) -> Option<CachedFilesystem> { self.filesystems.get_by_key(key).cloned() }
pub fn list_filesystems_in_pool( &self, pool: &zvariant::ObjectPath<'_>, ) -> Vec<CachedFilesystem> { self.filesystems .get_by_pool(&pool.as_ref()) .into_iter() .cloned() .collect() }}