Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
25 kB · 752 lines
Rust
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753use std::{borrow::Cow, sync::Arc, time::Duration};
use futures_util::stream::BoxStream;use tokio::sync::{Mutex, MutexGuard, watch};use tokio_stream::wrappers::WatchStream;use tokio_util::sync::{CancellationToken, DropGuard};use tracing::{debug, info, warn};use zbus::zvariant;
use crate::{ config::{Config, DeviceClass}, stratis::{ dbus::{Bytes, FilesystemProxy, FilesystemSpec, OperationResult, check_stratisd_version}, object_cache::{CachedFilesystem, CachedPool, FilesystemKey, ObjectCache}, },};
pub mod lvmd { tonic::include_proto!("proto");}
const BYTES_PER_GIB: u64 = 1024 * 1024 * 1024;
trait IntoTonicResult<T> { fn into_tonic_result(self) -> tonic::Result<T>;}
impl<T> IntoTonicResult<T> for zbus::Result<T> { fn into_tonic_result(self) -> tonic::Result<T> { self.map_err(|e| tonic::Status::unknown(format!("zbus: {e}"))) }}
impl<T: zvariant::Type> IntoTonicResult<T> for OperationResult<T> { fn into_tonic_result(self) -> tonic::Result<T> { if self.return_code > 0 { Err(tonic::Status::unknown(format!( "stratis: {}", self.return_message ))) } else { Ok(self.result) } }}
fn major_minor_from_stat(stat: rustix::fs::Stat) -> (u32, u32) { ( rustix::fs::major(stat.st_rdev), rustix::fs::minor(stat.st_rdev), )}
fn devnode_major_minor(path: &str) -> tonic::Result<(u32, u32)> { let stat = rustix::fs::stat(path) .map_err(|e| tonic::Status::internal(format!("stat {path}: {e:?}")))?; Ok(major_minor_from_stat(stat))}
async fn devnode_major_minor_retry_missing(path: &str) -> tonic::Result<(u32, u32)> { for _ in 0..50 { match rustix::fs::stat(path) { Ok(stat) => return Ok(major_minor_from_stat(stat)), Err(e) if e.kind() == std::io::ErrorKind::NotFound => { tokio::time::sleep(Duration::from_millis(100)).await; } Err(e) => return Err(tonic::Status::internal(format!("stat {path}: {e:?}"))), } }
Err(tonic::Status::deadline_exceeded(format!( "waiting for devnode {path}" )))}
fn get_size_bytes_unsigned(size_bytes: i64) -> tonic::Result<u64> { size_bytes .try_into() .map_err(|_| tonic::Status::invalid_argument(format!("size_bytes {} is <0", size_bytes)))}
async fn check_size_fits(proxy: &FilesystemProxy<'_>, size: u128) -> tonic::Result<()> { let Bytes(current_size) = proxy.size().await.into_tonic_result()?; if current_size > size { Err(tonic::Status::failed_precondition(format!( "current size of {current_size} bytes is greater than requested {size}" ))) } else { Ok(()) }}
#[tracing::instrument(skip_all, fields(pool = pool.name))]async fn get_pool_free_total_bytes( pool: &CachedPool, object_cache: &ObjectCache, device_class: &DeviceClass,) -> tonic::Result<(Bytes, Bytes)> { // Find out how much space it "reserved" by filesystem limits. let reserved: u128 = futures_util::future::join_all( object_cache .list_filesystems_in_pool(&pool.path) .into_iter() .map(|fs| async move { let reserved = match fs.proxy.size_limit().await?.into_option() { Some(limit) => limit.0, // This shouldn't really happen, but try to be smart about // it just in case. None => fs.proxy.size().await?.0, }; debug!(fs = fs.key.0.name.as_ref(), reserved); Ok::<_, zbus::Error>(reserved) }), ) .await .into_iter() .filter_map(|reserved| { reserved .map_err(|err| { warn!(?err, "processing logical volume"); err }) .ok() }) .sum();
let total = pool.proxy.total_physical_size().await.into_tonic_result()?;
// TODO: cache the metadata based on *something*? let metadata = pool .proxy .metadata(true) .await .into_tonic_result()? .into_tonic_result()? .into_inner();
let total_data_allocated: u128 = metadata .flex_devs .thin_data_dev .iter() .map(|region| Bytes::from(region.length).0) .sum();
// Stratis won't allocate any more than 50GiB from a pool at first, so we // need to combine the data device size with the amount of space that hasn't // been allocated at all. Additionally, data sizes are always a multiple of // the data block size, so round this down to match. const BYTES_PER_DATA_BLOCK: u128 = 2 * 1024 * 512; let total_unallocated = (total.0 - metadata .backstore .data_tier .blockdev .allocs .iter() .flatten() .map(|a| Bytes::from(a.length).0) .sum::<u128>() - metadata .backstore .data_tier .blockdev .devs .iter() .flat_map(|d| &d.integrity_meta_allocs) .map(|r| Bytes::from(r.length).0) .sum::<u128>() // 8Ki sectors (4MiB) metadata at the very front. - (8 * 1024 * 512)) / BYTES_PER_DATA_BLOCK * BYTES_PER_DATA_BLOCK;
let total_data = total_data_allocated + total_unallocated;
let spare = u128::from(device_class.spare_gb) * u128::from(BYTES_PER_GIB);
let free = total_data.saturating_sub(reserved).saturating_sub(spare);
debug!( reserved, spare, free, total_data_allocated, total_unallocated, total = total.0, );
Ok((Bytes(free), total))}
fn get_cached_pool_for_device_class( object_cache: &ObjectCache, device_class: &DeviceClass,) -> tonic::Result<CachedPool> { object_cache .get_pool_by_name(&device_class.pool) .ok_or_else(|| { tonic::Status::failed_precondition(format!( "device class uses missing pool '{}'", device_class.pool )) })}
fn lookup_cached_filesystem( object_cache: &ObjectCache, device_class: &DeviceClass, pool: zvariant::ObjectPath<'_>, name: &str,) -> tonic::Result<CachedFilesystem> { object_cache .get_filesystem_by_key(&FilesystemKey { pool, name: Cow::Borrowed(name), }) .ok_or_else(|| { tonic::Status::not_found(format!( "name={name} (for device_class={})", device_class.pool )) })}
pub struct Service { config: Config, object_cache: Arc<Mutex<ObjectCache>>, watch_sender: watch::Sender<lvmd::WatchResponse>, _running_watcher: DropGuard,}
impl Service { pub async fn start(config: Config, conn: zbus::Connection) -> zbus::Result<Arc<Self>> { check_stratisd_version(&conn).await?;
let token = CancellationToken::new(); let (object_cache, mut notifier) = ObjectCache::start(conn, token.clone()).await?;
let (watch_sender, _) = watch::channel(lvmd::WatchResponse { free_bytes: 0, items: vec![], });
let service = Arc::new(Self { config, object_cache, watch_sender, _running_watcher: token.clone().drop_guard(), });
let service_2 = service.clone(); tokio::task::spawn(token.run_until_cancelled_owned(async move { loop { notifier.any_pool_size_changed().await; service_2.on_any_pool_size_changed().await; } }));
Ok(service) }
async fn on_any_pool_size_changed(&self) { let object_cache = self.object_cache.lock().await;
let mut items: Vec<lvmd::WatchItem> = vec![]; for (name, device_class) in self.config.device_classes.by_name.iter() { let Ok(pool) = get_cached_pool_for_device_class(&object_cache, device_class) else { warn!( pool = device_class.pool, "skipping missing pool in WatchItem" ); continue; };
match async { let (free, total) = get_pool_free_total_bytes(&pool, &object_cache, device_class).await?; Ok::<_, tonic::Status>(lvmd::WatchItem { free_bytes: free.0.try_into().map_err(|_| { tonic::Status::out_of_range(format!("free bytes {} is >u64", free.0)) })?, device_class: name.to_owned(), size_bytes: total.0.try_into().map_err(|_| { tonic::Status::out_of_range(format!("total bytes {} is >u64", total.0)) })?, thin_pool: None, }) } .await { Ok(item) => items.push(item), Err(err) => { warn!(?pool, ?err, "processing pool"); } } }
let resp = lvmd::WatchResponse { free_bytes: self .config .device_classes .default .as_ref() .and_then(|default| items.iter().find(|item| item.device_class == *default)) .map(|item| item.free_bytes) .unwrap_or_default(), items, }; debug!(?resp, "replacing current watch state"); self.watch_sender.send_replace(resp); }
fn lookup_device_class(&self, name: &str) -> tonic::Result<&DeviceClass> { if !name.is_empty() { self.config .device_classes .by_name .get(name) .ok_or_else(|| tonic::Status::not_found(format!("device_class={name}"))) } else { self.config.device_classes.default().ok_or_else(|| { tonic::Status::invalid_argument("device_class not given, and no default set") }) } }
async fn wait_for_filesystem_presence( &self, mut object_cache: MutexGuard<'_, ObjectCache>, device_class: &DeviceClass, pool: zvariant::ObjectPath<'_>, name: &str, ) -> tonic::Result<CachedFilesystem> { for _ in 0..50 { if let Some(fs) = object_cache.get_filesystem_by_key(&FilesystemKey { pool: pool.as_ref(), name: Cow::Borrowed(name), }) { return Ok(fs); }
std::mem::drop(object_cache); tokio::time::sleep(Duration::from_millis(100)).await; object_cache = self.object_cache.lock().await; }
Err(tonic::Status::deadline_exceeded(format!( "waiting for name={name} (for device_class={})", device_class.pool ))) }
async fn wait_for_filesystem_removal( &self, mut object_cache: MutexGuard<'_, ObjectCache>, device_class: &DeviceClass, pool: zvariant::ObjectPath<'_>, name: &str, ) -> tonic::Result<()> { for _ in 0..50 { if object_cache .get_filesystem_by_key(&FilesystemKey { pool: pool.as_ref(), name: Cow::Borrowed(name), }) .is_none() { return Ok(()); }
std::mem::drop(object_cache); tokio::time::sleep(Duration::from_millis(100)).await; object_cache = self.object_cache.lock().await; }
Err(tonic::Status::deadline_exceeded(format!( "waiting for removal of name={name} (for device_class={})", device_class.pool ))) }}
#[tonic::async_trait]impl lvmd::lv_service_server::LvService for Service { #[tracing::instrument(skip(self), err)] async fn create_lv( &self, request: tonic::Request<lvmd::CreateLvRequest>, ) -> tonic::Result<tonic::Response<lvmd::CreateLvResponse>> { let request = request.into_inner();
if !request.tags.is_empty() { return Err(tonic::Status::unimplemented("tags")); } if !request.lvcreate_option_class.is_empty() { return Err(tonic::Status::unimplemented("lvcreate_option_class")); }
let device_class = self.lookup_device_class(&request.device_class)?;
#[allow(deprecated)] let size = get_size_bytes_unsigned(request.size_bytes)?;
let object_cache = self.object_cache.lock().await; let pool = get_cached_pool_for_device_class(&object_cache, device_class)?;
info!("asking stratis to create LV"); let ret = pool .proxy .create_filesystems(&[FilesystemSpec { name: request.name.clone(), size: Some(Bytes(size.into())).into(), size_limit: Some(Bytes(size.into())).into(), }]) .await .into_tonic_result() .and_then(IntoTonicResult::into_tonic_result)? .into_option() .unwrap_or_default();
if ret.len() != 1 { return Err(tonic::Status::internal(format!( "expected 1 created filesystem, got: {ret:?}" ))); }
let fs = self .wait_for_filesystem_presence( object_cache, device_class, pool.path.as_ref(), &request.name, ) .await?;
info!("retrieving devnode information"); let devnode = fs.proxy.devnode().await.into_tonic_result()?; let (dev_major, dev_minor) = devnode_major_minor_retry_missing(&devnode).await?;
#[allow(deprecated)] Ok(tonic::Response::new(lvmd::CreateLvResponse { volume: Some(lvmd::LogicalVolume { name: request.name, dev_major, dev_minor, tags: vec![], size_bytes: request.size_bytes, path: devnode, // lvmd doesn't set this in CreateLV attr: "".to_owned(), }), })) }
#[tracing::instrument(skip(self), err)] async fn remove_lv( &self, request: tonic::Request<lvmd::RemoveLvRequest>, ) -> tonic::Result<tonic::Response<lvmd::Empty>> { let request = request.into_inner();
let device_class = self.lookup_device_class(&request.device_class)?;
let object_cache = self.object_cache.lock().await; let pool = get_cached_pool_for_device_class(&object_cache, device_class)?; let filesystem = lookup_cached_filesystem( &object_cache, device_class, pool.path.as_ref(), &request.name, )?;
info!("asking stratis to remove LV"); let ret = pool .proxy .destroy_filesystems(&[filesystem.path.as_ref()]) .await .into_tonic_result()? .into_tonic_result()? .into_option() .unwrap_or_default();
if ret.len() != 1 { return Err(tonic::Status::internal(format!( "expected 1 destroyed filesystem, got: {ret:?}" ))); }
self.wait_for_filesystem_removal( object_cache, device_class, pool.path.as_ref(), &request.name, ) .await?;
Ok(tonic::Response::new(lvmd::Empty {})) }
#[tracing::instrument(skip(self), err)] async fn resize_lv( &self, request: tonic::Request<lvmd::ResizeLvRequest>, ) -> tonic::Result<tonic::Response<lvmd::ResizeLvResponse>> { let request = request.into_inner();
let device_class = self.lookup_device_class(&request.device_class)?;
#[allow(deprecated)] let size = get_size_bytes_unsigned(request.size_bytes)?;
let object_cache = self.object_cache.lock().await; let pool = get_cached_pool_for_device_class(&object_cache, device_class)?; let filesystem = lookup_cached_filesystem( &object_cache, device_class, pool.path.as_ref(), &request.name, )?;
check_size_fits(&filesystem.proxy, size.into()).await?;
info!("updating current size"); filesystem .proxy .set_size_limit(Some(Bytes(size.into())).into()) .await .into_tonic_result()?;
Ok(tonic::Response::new(lvmd::ResizeLvResponse { size_bytes: request.size_bytes, })) }
#[tracing::instrument(skip(self), err)] async fn create_lv_snapshot( &self, request: tonic::Request<lvmd::CreateLvSnapshotRequest>, ) -> tonic::Result<tonic::Response<lvmd::CreateLvSnapshotResponse>> { let request = request.into_inner();
if !request.tags.is_empty() { return Err(tonic::Status::unimplemented("tags")); }
// XXX: we lie about the access type, because topolvm always gives us // "ro". if request.access_type != "rw" { warn!( access_type = request.access_type, "unsupported; assuming 'rw", ); }
let device_class = self.lookup_device_class(&request.device_class)?;
#[allow(deprecated)] let size = get_size_bytes_unsigned(request.size_bytes)?;
let object_cache = self.object_cache.lock().await; let pool = get_cached_pool_for_device_class(&object_cache, device_class)?; let source = lookup_cached_filesystem( &object_cache, device_class, pool.path.as_ref(), &request.source_volume, )?;
if object_cache .get_filesystem_by_key(&FilesystemKey { pool: pool.path.as_ref(), name: Cow::Borrowed(&request.name), }) .is_some() { return Err(tonic::Status::already_exists(request.name)); }
// TODO: will the snapshot size match the source size? is this actually // useful? check_size_fits(&source.proxy, size.into()).await?;
info!("creating snapshot"); let ret = pool .proxy .snapshot_filesystem(source.path.as_ref(), &request.name) .await .into_tonic_result()? .into_tonic_result()? .into_option();
if ret.is_none() { return Err(tonic::Status::internal( "expected a created snapshot, got nothing", )); };
info!("waiting for snapshot object"); let snapshot = self .wait_for_filesystem_presence( object_cache, device_class, pool.path.as_ref(), &request.name, ) .await?;
info!("updating snapshot size limit");
snapshot .proxy .set_size_limit(Some(Bytes(size.into())).into()) .await .into_tonic_result()?;
info!("retrieving devnode information"); let devnode = snapshot.proxy.devnode().await.into_tonic_result()?; let (dev_major, dev_minor) = devnode_major_minor_retry_missing(&devnode).await?;
#[allow(deprecated)] Ok(tonic::Response::new(lvmd::CreateLvSnapshotResponse { snapshot: Some(lvmd::LogicalVolume { name: request.name, dev_major, dev_minor, tags: vec![], size_bytes: request.size_bytes, path: devnode, // lvmd doesn't set this in CreateLVSnapshot attr: "".to_owned(), }), })) }}
#[tonic::async_trait]impl lvmd::vg_service_server::VgService for Service { #[tracing::instrument(skip(self), err)] async fn get_lv_list( &self, request: tonic::Request<lvmd::GetLvListRequest>, ) -> tonic::Result<tonic::Response<lvmd::GetLvListResponse>> { let request = request.into_inner();
let device_class = self.lookup_device_class(&request.device_class)?;
let object_cache = self.object_cache.lock().await; let pool = get_cached_pool_for_device_class(&object_cache, device_class)?; let filesystems = object_cache.list_filesystems_in_pool(&pool.path);
let volumes = futures_util::future::join_all(filesystems.into_iter().map(|fs| { let pool_name = pool.name.clone(); async move { let devnode = fs.proxy.devnode().await.into_tonic_result()?; let (dev_major, dev_minor) = devnode_major_minor(&devnode)?; let size = fs .proxy .size_limit() .await .into_tonic_result()? .into_option() .ok_or_else(|| { tonic::Status::unknown(format!( "{pool_name}/{} has no size limit", fs.key.0.name )) })?;
let size_bytes = size.0.try_into().map_err(|_| { tonic::Status::out_of_range(format!( "{pool_name}/{} size {} is >i64", fs.key.0.name, size.0 )) })?;
#[allow(deprecated)] Ok::<_, tonic::Status>(lvmd::LogicalVolume { name: fs.key.0.name.into_owned(), dev_major, dev_minor, tags: vec![], size_bytes, path: devnode, // https://man7.org/linux/man-pages/man8/lvs.8.html#NOTES // really topolvm only cares about this for health check // purposes, but we try to fill in plausible values anyway. // TODO: is there anything that could be better filled by // asking stratis for more info? // - : volume type = normal // w : permissions = writable // a : allocation policy = anywhere // - : minor = no (idk what this even means) // a : active = yes // o : open = yes // - : target = none // - : zero = no // - : volume health = ok // - : skip activation = no attr: "-wa-ao----".to_owned(), }) } })) .await .into_iter() .filter_map(|lv| { lv.map_err(|err| { warn!(?err, "processing logical volume"); err }) .ok() }) .collect();
Ok(tonic::Response::new(lvmd::GetLvListResponse { volumes })) }
#[tracing::instrument(skip(self), err)] async fn get_free_bytes( &self, request: tonic::Request<lvmd::GetFreeBytesRequest>, ) -> tonic::Result<tonic::Response<lvmd::GetFreeBytesResponse>> { let request = request.into_inner();
let device_class = self.lookup_device_class(&request.device_class)?;
let object_cache = self.object_cache.lock().await; let pool = get_cached_pool_for_device_class(&object_cache, device_class)?; let (free, _) = get_pool_free_total_bytes(&pool, &object_cache, device_class).await?;
Ok(tonic::Response::new(lvmd::GetFreeBytesResponse { free_bytes: free.0.try_into().map_err(|_| { tonic::Status::out_of_range(format!("free bytes {} is >u64", free.0)) })?, })) }
type WatchStream = BoxStream<'static, tonic::Result<lvmd::WatchResponse>>;
#[tracing::instrument(skip(self), err)] async fn watch( &self, _request: tonic::Request<lvmd::Empty>, ) -> tonic::Result<tonic::Response<Self::WatchStream>> { use futures_util::StreamExt;
Ok(tonic::Response::new( WatchStream::new(self.watch_sender.subscribe()) .map(Ok) .boxed(), )) }}