diff --git a/dev-tools/omdb/src/bin/omdb/db.rs b/dev-tools/omdb/src/bin/omdb/db.rs index d02b22d6ce1..3f6b7328715 100644 --- a/dev-tools/omdb/src/bin/omdb/db.rs +++ b/dev-tools/omdb/src/bin/omdb/db.rs @@ -2649,11 +2649,11 @@ async fn cmd_db_disk_info( datastore: &DataStore, args: &DiskInfoArgs, ) -> Result<(), anyhow::Error> { + let conn = datastore.pool_connection_for_tests().await?; + let disk = { use nexus_db_schema::schema::disk::dsl; - let conn = datastore.pool_connection_for_tests().await?; - dsl::disk .filter(dsl::id.eq(args.uuid)) .select(nexus_db_model::Disk::as_select()) @@ -2662,7 +2662,7 @@ async fn cmd_db_disk_info( .context("failed to find disk")? }; - match datastore.disk_get_with_model(opctx, disk).await? { + match datastore.disk_get_with_model_on_connection(&conn, disk).await? { Disk::Crucible(disk) => { crucible_disk_info(opctx, datastore, disk).await } diff --git a/dev-tools/omdb/src/bin/omdb/nexus.rs b/dev-tools/omdb/src/bin/omdb/nexus.rs index 43bc8b86e55..df8b4ef3963 100644 --- a/dev-tools/omdb/src/bin/omdb/nexus.rs +++ b/dev-tools/omdb/src/bin/omdb/nexus.rs @@ -71,6 +71,7 @@ use nexus_types::internal_api::background::IncompleteBootstoreConfigReport; use nexus_types::internal_api::background::InstanceReincarnationStatus; use nexus_types::internal_api::background::InstanceUpdaterStatus; use nexus_types::internal_api::background::InventoryLoadStatus; +use nexus_types::internal_api::background::LocalStorageDeleteStatus; use nexus_types::internal_api::background::LookupRegionPortStatus; use nexus_types::internal_api::background::PhysicalDiskAdoptionStatus; use nexus_types::internal_api::background::ProbeDistributorStatus; @@ -1433,6 +1434,9 @@ fn print_task_details(bgtask: &BackgroundTask, details: &serde_json::Value) { "switch_port_config_manager" => { print_task_switch_port_settings_manager(details); } + "local_storage_delete" => { + print_task_local_storage_delete(details); + } _ => { println!( "warning: unknown background task: {:?} \ @@ -4386,6 +4390,50 @@ fn print_task_physical_disk_adoption(details: &serde_json::Value) { } } +fn print_task_local_storage_delete(details: &serde_json::Value) { + match serde_json::from_value::(details.clone()) { + Err(error) => eprintln!( + "warning: failed to interpret task details: {:?}: {:?}", + error, details + ), + + Ok(status) => { + let LocalStorageDeleteStatus { + total_allocations_to_delete, + page_size, + delete_results, + deallocate_results, + errors, + } = &status; + + println!( + " total allocations left to delete: \ + {total_allocations_to_delete}" + ); + + println!( + " number of allocations deleted per invoked task: \ + {page_size}" + ); + + println!(" results of deleting local storage:"); + for result in delete_results { + println!(" > {result}"); + } + + println!(" results of deallocating local storage:"); + for result in deallocate_results { + println!(" > {result}"); + } + + println!(" errors: {}", errors.len()); + for error in errors { + println!(" > {error}"); + } + } + } +} + const ERRICON: &str = "/!\\"; fn warn_if_nonzero(n: usize) -> &'static str { diff --git a/dev-tools/omdb/tests/env.out b/dev-tools/omdb/tests/env.out index e271a2f3d11..6b8fcbb17ad 100644 --- a/dev-tools/omdb/tests/env.out +++ b/dev-tools/omdb/tests/env.out @@ -158,6 +158,10 @@ task: "inventory_loader" loads the latest inventory collection from the DB +task: "local_storage_delete" + delete resources for disks backed by local storage + + task: "lookup_region_port" fill in missing ports for region records @@ -430,6 +434,10 @@ task: "inventory_loader" loads the latest inventory collection from the DB +task: "local_storage_delete" + delete resources for disks backed by local storage + + task: "lookup_region_port" fill in missing ports for region records @@ -689,6 +697,10 @@ task: "inventory_loader" loads the latest inventory collection from the DB +task: "local_storage_delete" + delete resources for disks backed by local storage + + task: "lookup_region_port" fill in missing ports for region records diff --git a/dev-tools/omdb/tests/successes.out b/dev-tools/omdb/tests/successes.out index 8d72e5dfb3c..553656f8a79 100644 --- a/dev-tools/omdb/tests/successes.out +++ b/dev-tools/omdb/tests/successes.out @@ -402,6 +402,10 @@ task: "inventory_loader" loads the latest inventory collection from the DB +task: "local_storage_delete" + delete resources for disks backed by local storage + + task: "lookup_region_port" fill in missing ports for region records @@ -883,6 +887,16 @@ task: "inventory_loader" loaded latest inventory collection as of : collection ....................., taken at +task: "local_storage_delete" + configured period: every h m s + last completed activation: , triggered by + started at (s ago) and ran for ms + total allocations left to delete: 0 + number of allocations deleted per invoked task: 128 + results of deleting local storage: + results of deallocating local storage: + errors: 0 + task: "lookup_region_port" configured period: every m last completed activation: , triggered by @@ -1611,6 +1625,16 @@ task: "inventory_loader" loaded latest inventory collection as of : collection ....................., taken at +task: "local_storage_delete" + configured period: every h m s + last completed activation: , triggered by + started at (s ago) and ran for ms + total allocations left to delete: 0 + number of allocations deleted per invoked task: 128 + results of deleting local storage: + results of deallocating local storage: + errors: 0 + task: "lookup_region_port" configured period: every m last completed activation: , triggered by diff --git a/nexus-config/src/nexus_config.rs b/nexus-config/src/nexus_config.rs index 4dc88fb73fa..6c459612ee8 100644 --- a/nexus-config/src/nexus_config.rs +++ b/nexus-config/src/nexus_config.rs @@ -486,6 +486,8 @@ pub struct BackgroundTaskConfig { pub audit_log_cleanup: AuditLogCleanupConfig, /// configuration for populate switch ports task pub populate_switch_ports: PopulateSwitchPortsConfig, + /// configuration for local storage delete task + pub local_storage_delete: LocalStorageDeleteConfig, } #[serde_as] @@ -1100,6 +1102,14 @@ pub struct TrustQuorumConfig { pub period_secs: Duration, } +#[serde_as] +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +pub struct LocalStorageDeleteConfig { + /// period (in seconds) for periodic activations of this background task + #[serde_as(as = "DurationSeconds")] + pub period_secs: Duration, +} + /// Configuration for a nexus server #[derive(Clone, Debug, Deserialize, PartialEq, Serialize)] pub struct PackageConfig { @@ -1393,6 +1403,7 @@ mod test { audit_log_cleanup.retention_days = 90 audit_log_cleanup.max_deleted_per_activation = 10000 populate_switch_ports.period_secs = 31 + local_storage_delete.period_secs = 30 [default_region_allocation_strategy] type = "random" seed = 0 @@ -1679,6 +1690,10 @@ mod test { populate_switch_ports: PopulateSwitchPortsConfig { period_secs: Duration::from_secs(31), }, + local_storage_delete: + LocalStorageDeleteConfig { + period_secs: Duration::from_secs(30), + }, }, multicast: MulticastConfig { enabled: false }, default_region_allocation_strategy: @@ -1796,6 +1811,7 @@ mod test { audit_log_cleanup.retention_days = 90 audit_log_cleanup.max_deleted_per_activation = 10000 populate_switch_ports.period_secs = 31 + local_storage_delete.period_secs = 30 [default_region_allocation_strategy] type = "random" diff --git a/nexus/background-task-interface/src/init.rs b/nexus/background-task-interface/src/init.rs index fdf9447d466..2ec43e083e5 100644 --- a/nexus/background-task-interface/src/init.rs +++ b/nexus/background-task-interface/src/init.rs @@ -64,6 +64,7 @@ pub struct BackgroundTasks { pub task_attached_subnet_manager: Activator, pub task_session_cleanup: Activator, pub task_populate_switch_ports: Activator, + pub task_local_storage_delete: Activator, // Handles to activate background tasks that do not get used by Nexus // at-large. These background tasks are implementation details as far as diff --git a/nexus/db-queries/src/db/datastore/disk.rs b/nexus/db-queries/src/db/datastore/disk.rs index 9204f5870d1..167820cb189 100644 --- a/nexus/db-queries/src/db/datastore/disk.rs +++ b/nexus/db-queries/src/db/datastore/disk.rs @@ -34,6 +34,8 @@ use crate::db::model::Volume; use crate::db::model::to_db_typed_uuid; use crate::db::pagination::paginated; use crate::db::queries::disk::DiskSetClauseForAttach; +use crate::db::queries::virtual_provisioning_collection_update; +use crate::db::queries::virtual_provisioning_collection_update::*; use crate::db::update_and_check::UpdateAndCheck; use crate::db::update_and_check::UpdateStatus; use async_bb8_diesel::AsyncRunQueryDsl; @@ -52,6 +54,7 @@ use nexus_types::identity::Asset; use omicron_common::api; use omicron_common::api::external; use omicron_common::api::external::CreateResult; +use omicron_common::api::external::DeleteResult; use omicron_common::api::external::Error; use omicron_common::api::external::ListResultVec; use omicron_common::api::external::LookupResult; @@ -510,7 +513,9 @@ impl DataStore { let (.., disk) = LookupPath::new(opctx, self).disk_id(disk_id).fetch().await?; - self.disk_get_with_model(opctx, disk).await + let conn = self.pool_connection_authorized(opctx).await?; + + self.disk_get_with_model_on_connection(&conn, disk).await } /// Return a `datastore::Disk` given a `model::Disk` @@ -518,15 +523,15 @@ impl DataStore { /// Note: basically all of Nexus should _not_ be using this, and should be /// using `disk_get` instead: this version of the function bypasses the /// LookupPath induced permissions check and should only called from omdb. - pub async fn disk_get_with_model( + /// Code that is looking up deleted disks should also use this method, as + /// `LookupPath` will not return deleted resources. + pub async fn disk_get_with_model_on_connection( &self, - opctx: &OpContext, + conn: &async_bb8_diesel::Connection, disk: model::Disk, ) -> LookupResult { let disk_id = disk.id(); - let conn = self.pool_connection_authorized(opctx).await?; - let disk = match disk.disk_type { db::model::DiskType::Crucible => { use nexus_db_schema::schema::disk_type_crucible::dsl; @@ -534,7 +539,7 @@ impl DataStore { let disk_type_crucible = dsl::disk_type_crucible .filter(dsl::disk_id.eq(disk_id)) .select(DiskTypeCrucible::as_select()) - .first_async(&*conn) + .first_async(conn) .await .map_err(|e| { public_error_from_diesel(e, ErrorHandler::Server) @@ -552,7 +557,7 @@ impl DataStore { let disk_type_local_storage = dsl::disk_type_local_storage .filter(dsl::disk_id.eq(disk_id)) .select(DiskTypeLocalStorage::as_select()) - .first_async(&*conn) + .first_async(conn) .await .map_err(|e| { public_error_from_diesel(e, ErrorHandler::Server) @@ -1549,8 +1554,22 @@ impl DataStore { disk_id: &Uuid, ok_to_delete_states: &[api::external::DiskState], ) -> Result { - use nexus_db_schema::schema::disk::dsl; let conn = self.pool_connection_unauthorized().await?; + + Self::project_delete_disk_no_auth_on_connection( + &conn, + disk_id, + ok_to_delete_states, + ) + .await + } + + async fn project_delete_disk_no_auth_on_connection( + conn: &async_bb8_diesel::Connection, + disk_id: &Uuid, + ok_to_delete_states: &[api::external::DiskState], + ) -> Result { + use nexus_db_schema::schema::disk::dsl; let now = Utc::now(); let ok_to_delete_state_labels: Vec<_> = @@ -2083,6 +2102,63 @@ impl DataStore { Ok(disk) } + + /// In a single transaction, set time_deleted for a disk and delete the + /// storage from the appropriate virtual provisioning collection. + pub async fn delete_disk_and_update_provisioning_collection( + &self, + opctx: &OpContext, + project: &authz::Project, + disk: &Disk, + ok_to_delete_states: &[api::external::DiskState], + ) -> DeleteResult { + let err = OptionalError::new(); + let conn = self.pool_connection_authorized(opctx).await?; + + let provisions = self + .transaction_retry_wrapper( + "delete_disk_and_update_provisioning_collection", + ) + .transaction(&conn, |conn| { + let err = err.clone(); + async move { + Self::project_delete_disk_no_auth_on_connection( + &conn, + &disk.id(), + ok_to_delete_states, + ) + .await + .map_err(|e| err.bail(e))?; + + let provisions = + VirtualProvisioningCollectionUpdate::new_delete_storage( + disk.id(), + disk.size(), + project.id(), + ) + .get_results_async(&conn) + .await + .map_err(|e| err.bail( + virtual_provisioning_collection_update::from_diesel(e) + ))?; + + Ok(provisions) + } + }) + .await + .map_err(|e| { + if let Some(err) = err.take() { + err + } else { + public_error_from_diesel(e, ErrorHandler::Server) + } + })?; + + self.virtual_provisioning_collection_producer + .append_disk_metrics(&provisions)?; + + Ok(()) + } } #[cfg(test)] diff --git a/nexus/db-queries/src/db/datastore/local_storage.rs b/nexus/db-queries/src/db/datastore/local_storage.rs index 7ba521d2807..b616d0e9b76 100644 --- a/nexus/db-queries/src/db/datastore/local_storage.rs +++ b/nexus/db-queries/src/db/datastore/local_storage.rs @@ -9,10 +9,12 @@ use crate::authz; use crate::context::OpContext; use crate::db::collection_insert::AsyncInsertError; use crate::db::collection_insert::DatastoreCollection; +use crate::db::datastore; use crate::db::datastore::DbConnection; use crate::db::datastore::LocalStorageAllocation; use crate::db::datastore::LocalStorageDisk; use crate::db::datastore::SQL_BATCH_SIZE; +use crate::db::model; use crate::db::model::LocalStorageDatasetAllocation; use crate::db::model::LocalStorageUnencryptedDatasetAllocation; use crate::db::model::RendezvousLocalStorageDataset; @@ -280,9 +282,9 @@ impl DataStore { Ok(()) } - /// Mark the local storage dataset allocations as deleted, and re-compute - /// the appropriate dataset size_used columns. - pub async fn delete_local_storage_dataset_allocations( + /// Mark the local storage dataset allocation backing this disk as deleted, + /// and re-compute the appropriate dataset size_used columns. + pub async fn delete_local_storage_dataset_allocation( &self, opctx: &OpContext, local_storage_disk: &LocalStorageDisk, @@ -296,7 +298,7 @@ impl DataStore { let conn = self.pool_connection_authorized(opctx).await?; self.transaction_retry_wrapper( - "delete_local_storage_dataset_allocations", + "delete_local_storage_dataset_allocation", ) .transaction(&conn, |conn| async move { match local_storage_dataset_allocation { @@ -451,4 +453,70 @@ impl DataStore { }) .map_err(|e| public_error_from_diesel(e, ErrorHandler::Server)) } + + /// Return all deleted disks that have undeleted local storage allocations + pub async fn deleted_disks_with_undeleted_local_storage( + &self, + opctx: &OpContext, + ) -> Result, Error> { + opctx.authorize(authz::Action::Delete, &authz::FLEET).await?; + opctx.check_complex_operations_allowed()?; + + let conn = self.pool_connection_authorized(opctx).await?; + + use nexus_db_schema::schema::disk::dsl; + use nexus_db_schema::schema::disk_type_local_storage::dsl as dtls_dsl; + use nexus_db_schema::schema::local_storage_unencrypted_dataset_allocation::dsl as lsuda_dsl; + + // Find all deleted disks where the unencrypted local storage allocation + // is not yet deleted. + let found_disks: Vec = dsl::disk + .inner_join( + dtls_dsl::disk_type_local_storage + .on(dsl::id.eq(dtls_dsl::disk_id)), + ) + .inner_join( + lsuda_dsl::local_storage_unencrypted_dataset_allocation.on( + dtls_dsl::local_storage_unencrypted_dataset_allocation_id + .eq(lsuda_dsl::id.nullable()), + ), + ) + .filter(lsuda_dsl::time_deleted.is_null()) + .filter(dsl::time_deleted.is_not_null()) + .select(model::Disk::as_select()) + .load_async(&*conn) + .await + .map_err(|e| public_error_from_diesel(e, ErrorHandler::Server))?; + + let mut disks = Vec::with_capacity(found_disks.len()); + + for found_disk in found_disks { + match self + .disk_get_with_model_on_connection(&conn, found_disk) + .await? + { + datastore::Disk::Crucible(crucible_disk) => { + // The query above joins the disk table with the + // disk_type_local_storage table, meaning the higher level + // Disk can never be the Crucible type, unless there's a + // serious problem. Log an error but return whatever disks + // we can for deletion. + // + // TODO surface this to operators, support intervention is + // likely required. + error!( + self.log, + "disk {} should be local storage, not crucible", + crucible_disk.id(), + ); + } + + datastore::Disk::LocalStorage(local_storage_disk) => { + disks.push(local_storage_disk); + } + } + } + + Ok(disks) + } } diff --git a/nexus/examples/config-second.toml b/nexus/examples/config-second.toml index bf0cf0b3f1f..e9c7996704c 100644 --- a/nexus/examples/config-second.toml +++ b/nexus/examples/config-second.toml @@ -219,6 +219,7 @@ audit_log_cleanup.period_secs = 600 audit_log_cleanup.retention_days = 90 audit_log_cleanup.max_deleted_per_activation = 10000 populate_switch_ports.period_secs = 30 +local_storage_delete.period_secs = 30 [default_region_allocation_strategy] # allocate region on 3 random distinct zpools, on 3 random distinct sleds. diff --git a/nexus/examples/config.toml b/nexus/examples/config.toml index 3161542d409..36cd4bd7fa1 100644 --- a/nexus/examples/config.toml +++ b/nexus/examples/config.toml @@ -210,6 +210,7 @@ audit_log_cleanup.period_secs = 600 audit_log_cleanup.retention_days = 90 audit_log_cleanup.max_deleted_per_activation = 10000 populate_switch_ports.period_secs = 30 +local_storage_delete.period_secs = 30 [default_region_allocation_strategy] # allocate region on 3 random distinct zpools, on 3 random distinct sleds. diff --git a/nexus/src/app/background/init.rs b/nexus/src/app/background/init.rs index 77bddee6c01..c589871bf46 100644 --- a/nexus/src/app/background/init.rs +++ b/nexus/src/app/background/init.rs @@ -119,6 +119,7 @@ use super::tasks::instance_updater; use super::tasks::instance_watcher; use super::tasks::inventory_collection; use super::tasks::inventory_load; +use super::tasks::local_storage_delete::*; use super::tasks::lookup_region_port; use super::tasks::metrics_producer_gc; use super::tasks::multicast::MulticastGroupReconciler; @@ -282,6 +283,7 @@ impl BackgroundTasksInitializer { task_attached_subnet_manager: Activator::new(), task_session_cleanup: Activator::new(), task_populate_switch_ports: Activator::new(), + task_local_storage_delete: Activator::new(), // Handles to activate background tasks that do not get used by Nexus // at-large. These background tasks are implementation details as far as @@ -379,6 +381,7 @@ impl BackgroundTasksInitializer { task_audit_log_timeout_incomplete, task_audit_log_cleanup, task_populate_switch_ports, + task_local_storage_delete, // Add new background tasks here. Be sure to use this binding in a // call to `Driver::register()` below. That's what actually wires // up the Activator to the corresponding background task. @@ -1336,6 +1339,16 @@ impl BackgroundTasksInitializer { activator: task_populate_switch_ports, }); + driver.register(TaskDefinition { + name: "local_storage_delete", + description: "delete resources for disks backed by local storage", + period: config.local_storage_delete.period_secs, + task_impl: Box::new(LocalStorageDeleter::new(datastore.clone())), + opctx: opctx.child(BTreeMap::new()), + watchers: vec![], + activator: task_local_storage_delete, + }); + driver } } diff --git a/nexus/src/app/background/tasks/local_storage_delete.rs b/nexus/src/app/background/tasks/local_storage_delete.rs new file mode 100644 index 00000000000..19ea90c411b --- /dev/null +++ b/nexus/src/app/background/tasks/local_storage_delete.rs @@ -0,0 +1,383 @@ +// This Source Code Form is subject to the terms of the Mozilla Public +// License, v. 2.0. If a copy of the MPL was not distributed with this +// file, You can obtain one at https://mozilla.org/MPL/2.0/. + +//! When the higher level disks backed by local storage are deleted, this +//! background task will delete any allocated local storage, and will delete +//! local storage allocation records when those resources have been cleaned up. + +use crate::app::background::BackgroundTask; +use futures::FutureExt; +use futures::future::BoxFuture; +use nexus_db_queries::context::OpContext; +use nexus_db_queries::db::DataStore; +use nexus_db_queries::db::datastore::LocalStorageAllocation; +use nexus_db_queries::db::model::LocalStorageUnencryptedDatasetAllocation; +use nexus_types::internal_api::background::LocalStorageDeleteStatus; +use rand::prelude::*; +use serde_json::json; +use sled_agent_client::types::LocalStorageDatasetDeleteRequest; +use slog::Logger; +use slog_error_chain::InlineErrorChain; +use std::sync::Arc; + +pub struct LocalStorageDeleter { + datastore: Arc, + reqwest_client: reqwest::Client, +} + +#[derive(PartialEq)] +enum DeleteResult { + /// This invocation of the task deleted the local storage allocation + Deleted, + + /// No delete is required as the sled hosting the local storage allocation + /// was expunged. + SledExpunged, + + /// No delete is required as the zpool hosting the local storage allocation + /// was expunged. + ZpoolExpunged, + + /// This invocation of the task requested deletion but needs to wait + WaitForNextActivation, + + Error { + message: String, + }, +} + +impl LocalStorageDeleter { + pub fn new(datastore: Arc) -> Self { + let duration = std::time::Duration::from_millis(250); + + LocalStorageDeleter { + datastore, + // Create a client with a _short_ timeout, don't block this task + // waiting for request responses. + reqwest_client: reqwest::ClientBuilder::new() + .connect_timeout(duration) + .timeout(duration) + .build() + .unwrap(), + } + } + + async fn delete_unencrypted_allocation( + &self, + log: &Logger, + opctx: &OpContext, + allocation: &LocalStorageUnencryptedDatasetAllocation, + ) -> DeleteResult { + let sled_id = allocation.sled_id(); + let zpool_id = allocation.pool_id().upcast(); + + // Check if either the sled or disk backing the zpool was expunged. If + // we can't determine then bail and wait for the next task activation. + + let sled_in_service = + match self.datastore.check_sled_in_service(&opctx, sled_id).await { + Ok(sled_in_service) => sled_in_service, + + Err(e) => { + let message = format!( + "error calling check_sled_in_service for sled \ + {sled_id}: {}", + InlineErrorChain::new(&e), + ); + + return DeleteResult::Error { message }; + } + }; + + if !sled_in_service { + // Sled's been expunged, so consider the local storage deleted. + return DeleteResult::SledExpunged; + } + + let zpool_in_service = + match self.datastore.check_zpool_in_service(&opctx, zpool_id).await + { + Ok(zpool_in_service) => zpool_in_service, + + Err(e) => { + let message = format!( + "error calling check_zpool_in_service for zpool \ + {zpool_id}: {}", + InlineErrorChain::new(&e), + ); + + return DeleteResult::Error { message }; + } + }; + + if !zpool_in_service { + // The disk backing the zpool's been expunged, so consider the local + // storage deleted. + return DeleteResult::ZpoolExpunged; + } + + // Now that all checks are done, get a sled agent client and make the + // delete request. + + let request = LocalStorageDatasetDeleteRequest { + zpool_id: allocation.pool_id(), + dataset_id: allocation.id(), + encrypted_at_rest: false, + }; + + let sled_agent_client = match nexus_networking::sled_client_ext( + &self.datastore, + opctx, + sled_id, + log, + self.reqwest_client.clone(), + ) + .await + { + Ok(client) => client, + + Err(e) => { + let message = format!( + "error calling sled_client_ext for sled {sled_id}: {}", + InlineErrorChain::new(&e), + ); + + return DeleteResult::Error { message }; + } + }; + + match sled_agent_client.local_storage_dataset_delete(&request).await { + Ok(_) => DeleteResult::Deleted, + + Err(progenitor_client::Error::CommunicationError(e)) + if e.is_timeout() => + { + // The request timed out but is still being processed by the + // remote sled-agent. + DeleteResult::WaitForNextActivation + } + + Err(e) => { + let message = format!( + "error sending local_storage_dataset_delete: {}", + InlineErrorChain::new(&e), + ); + + DeleteResult::Error { message } + } + } + } + + async fn activate_impl( + &self, + opctx: &OpContext, + ) -> LocalStorageDeleteStatus { + let log = &opctx.log; + let mut status = LocalStorageDeleteStatus::default(); + + let mut disks_needing_clean_up = match self + .datastore + .deleted_disks_with_undeleted_local_storage(opctx) + .await + { + Ok(v) => v, + + Err(e) => { + let s = format!( + "error calling \ + deleted_disks_with_undeleted_local_storage: {}", + InlineErrorChain::new(&e), + ); + + error!(log, "{s}"); + status.errors.push(s); + + return status; + } + }; + + // Report the remaining work to do in the status. + + status.total_allocations_to_delete = disks_needing_clean_up.len(); + + // Operate on a maximum of 128 disks per task invocation. If users + // create disks very quickly during the time between this task's + // periodic invocation (and then deletes them all), we could be faced + // with many disks to delete. Each Nexus that then activates this task + // would fetch all those disks for deletion, potentially causing each of + // the tasks to take a long time. + // + // The timeout for the reqwest client created by this task is 250 + // milliseconds, so the maximum time (assuming the requests to all + // sled-agents time out) is 128 * 250 ms = 32 seconds, roughly + // approximating the periodic task wakeup time. Note that each local + // storage disk delete will activate this task, so the latency between + // the delete request and the actual deletion should remain low, + // assuming there aren't too many lingering problematic allocations. + + status.page_size = 128; + + // Randomize the list: if there are lingering problematic allocations + // at the beginning of `disks_needing_clean_up` it could wedge the task. + + disks_needing_clean_up.shuffle(&mut rand::rng()); + + for disk in disks_needing_clean_up.into_iter().take(status.page_size) { + let Some(allocation) = &disk.local_storage_dataset_allocation + else { + // No allocation was made for this disk + continue; + }; + + // Attempt deleting the local storage before removing the + // database record. If the delete does not succeed, try again in + // the next task activation. + + match allocation { + LocalStorageAllocation::Unencrypted(allocation) => { + match self + .delete_unencrypted_allocation(log, opctx, &allocation) + .await + { + DeleteResult::Deleted => { + let s = format!( + "deleted disk {} allocation {}", + disk.id(), + allocation.id(), + ); + + info!(log, "{s}"); + status.delete_results.push(s); + + // Drop through to deallocation once deletion + // succeeds. + } + + DeleteResult::SledExpunged => { + let s = format!( + "disk {} allocation {} sled expunged, \ + considering deleted", + disk.id(), + allocation.id(), + ); + + info!(log, "{s}"); + status.delete_results.push(s); + + // Drop through to deallocation, deletion not + // required. + } + + DeleteResult::ZpoolExpunged => { + let s = format!( + "disk {} allocation {} zpool expunged, \ + considering deleted", + disk.id(), + allocation.id(), + ); + + info!(log, "{s}"); + status.delete_results.push(s); + + // Drop through to deallocation, deletion not + // required. + } + + DeleteResult::WaitForNextActivation => { + let s = format!( + "requested deletion of disk {} allocation \ + {}", + disk.id(), + allocation.id(), + ); + + info!(log, "{s}"); + status.delete_results.push(s); + + // Cannot deallocate the record until deletion + // succeeds. + continue; + } + + DeleteResult::Error { message } => { + info!(log, "{message}"); + status.errors.push(message); + + // Cannot deallocate the record until deletion + // succeeds. + continue; + } + } + } + + LocalStorageAllocation::Encrypted(allocation) => { + // Until encrypted local storage is supported, seeing a + // request to clean up disks of that type should be + // noted as a error. + let s = format!( + "request to delete disk {} encrypted allocation {}", + disk.id(), + allocation.id(), + ); + + error!(log, "{s}"); + status.errors.push(s); + + continue; + } + } + + match self + .datastore + .delete_local_storage_dataset_allocation(opctx, &disk) + .await + { + Ok(()) => { + let s = format!( + "deallocated disk {} allocation {}", + disk.id(), + allocation.id(), + ); + + info!(log, "{s}"); + status.deallocate_results.push(s); + } + + Err(e) => { + let s = format!( + "error calling \ + delete_local_storage_dataset_allocation: {}", + InlineErrorChain::new(&e), + ); + + error!(log, "{s}"); + status.errors.push(s); + } + } + } + + status + } +} + +impl BackgroundTask for LocalStorageDeleter { + fn activate<'a>( + &'a mut self, + opctx: &'a OpContext, + ) -> BoxFuture<'a, serde_json::Value> { + async { + let status = self.activate_impl(opctx).await; + match serde_json::to_value(status) { + Ok(val) => val, + Err(e) => json!({ + "error": format!( + "could not serialize task status: {}", + InlineErrorChain::new(&e), + ) + }), + } + } + .boxed() + } +} diff --git a/nexus/src/app/background/tasks/mod.rs b/nexus/src/app/background/tasks/mod.rs index b78a217346e..fe2007dc056 100644 --- a/nexus/src/app/background/tasks/mod.rs +++ b/nexus/src/app/background/tasks/mod.rs @@ -32,6 +32,7 @@ pub mod instance_updater; pub mod instance_watcher; pub mod inventory_collection; pub mod inventory_load; +pub mod local_storage_delete; pub mod lookup_region_port; pub mod metrics_producer_gc; pub mod multicast; diff --git a/nexus/src/app/disk.rs b/nexus/src/app/disk.rs index 75172078f0c..3578165eef2 100644 --- a/nexus/src/app/disk.rs +++ b/nexus/src/app/disk.rs @@ -415,15 +415,40 @@ impl super::Nexus { let disk = self.datastore().disk_get(opctx, authz_disk.id()).await?; - let saga_params = sagas::disk_delete::Params { - serialized_authn: authn::saga::Serialized::for_opctx(opctx), - project_id: project.id(), - disk, - }; + match disk { + datastore::Disk::Crucible(_) => { + // For now, all Crucible related clean-up is done in the disk + // delete saga. Stay tuned! + + let saga_params = sagas::disk_delete::Params { + serialized_authn: authn::saga::Serialized::for_opctx(opctx), + project_id: project.id(), + disk: disk.clone(), + }; - self.sagas - .saga_execute::(saga_params) - .await?; + self.sagas + .saga_execute::( + saga_params, + ) + .await?; + } + + datastore::Disk::LocalStorage(_) => { + // Disks backed by local storage are marked for deletion and + // then cleaned up in a background task. + + self.datastore() + .delete_disk_and_update_provisioning_collection( + opctx, + &project, + &disk, + &[DiskState::Detached, DiskState::Faulted], + ) + .await?; + + self.background_tasks.task_local_storage_delete.activate(); + } + } Ok(()) } diff --git a/nexus/src/app/sagas/disk_delete.rs b/nexus/src/app/sagas/disk_delete.rs index 271852fe0ca..0193160cc56 100644 --- a/nexus/src/app/sagas/disk_delete.rs +++ b/nexus/src/app/sagas/disk_delete.rs @@ -112,12 +112,18 @@ impl NexusSaga for SagaDiskDelete { } datastore::Disk::LocalStorage(_) => { - // Attempt deleting the local storage before removing the - // database record. If the delete does not succeed, at least the - // user can re-request the deletion. - - builder.append(delete_local_storage_action()); - builder.append(deallocate_local_storage_action()); + // Clean up of the local storage resources is now done entirely + // in the `local_storage_delete` background task. The structure + // of this saga has been left untouched (to handle the case + // where Nexus is mupdated to a new version containing this + // change, but will resume any existing interrupted sagas) but + // it is now an error to create a new version of this saga for + // disks backed by local storage. + + return Err(SagaInitError::InvalidParameter(format!( + "disk {} is of type local storage", + params.disk.id(), + ))); } } @@ -326,7 +332,7 @@ async fn sdd_deallocate_local_storage( osagactx .datastore() - .delete_local_storage_dataset_allocations(&opctx, &disk) + .delete_local_storage_dataset_allocation(&opctx, &disk) .await .map_err(saga_action_failed)?; @@ -339,18 +345,7 @@ pub(crate) mod test { app::saga::create_saga_dag, app::sagas::disk_delete::Params, app::sagas::disk_delete::SagaDiskDelete, }; - use async_bb8_diesel::AsyncRunQueryDsl; - use chrono::Utc; - use diesel::ExpressionMethods; - use diesel::QueryDsl; - use diesel::SelectableHelper; - use nexus_db_lookup::LookupPath; - use nexus_db_model::LocalStorageUnencryptedDatasetAllocation; - use nexus_db_model::PhysicalDiskPolicy; - use nexus_db_model::RendezvousLocalStorageUnencryptedDataset; - use nexus_db_model::to_db_typed_uuid; use nexus_db_queries::authn::saga::Serialized; - use nexus_db_queries::authz; use nexus_db_queries::context::OpContext; use nexus_db_queries::db::datastore::Disk; use nexus_test_utils::resource_helpers::DiskTest; @@ -358,16 +353,9 @@ pub(crate) mod test { use nexus_test_utils_macros::nexus_test; use nexus_types::external_api::disk; use nexus_types::external_api::project; - use nexus_types::identity::Asset; use omicron_common::api::external::ByteCount; use omicron_common::api::external::IdentityMetadataCreateParams; use omicron_common::api::external::Name; - use omicron_uuid_kinds::BlueprintUuid; - use omicron_uuid_kinds::DatasetUuid; - use omicron_uuid_kinds::ExternalZpoolUuid; - use omicron_uuid_kinds::GenericUuid; - use omicron_uuid_kinds::ZpoolUuid; - use uuid::Uuid; type ControlPlaneTestContext = nexus_test_utils::ControlPlaneTestContext; @@ -396,17 +384,6 @@ pub(crate) mod test { } } - pub fn new_local_disk_create_params() -> disk::DiskCreate { - disk::DiskCreate { - identity: IdentityMetadataCreateParams { - name: "local".parse().expect("Invalid disk name"), - description: "My disk".to_string(), - }, - disk_backend: disk::DiskBackend::Local {}, - size: ByteCount::from_gibibytes_u32(1), - } - } - pub async fn create_disk( cptestctx: &ControlPlaneTestContext, create_params: F, @@ -482,267 +459,4 @@ pub(crate) mod test { ) .await; } - - pub struct ExpungeTestHarness<'a> { - cptestctx: &'a ControlPlaneTestContext, - zpool_id: ZpoolUuid, - allocation: LocalStorageUnencryptedDatasetAllocation, - } - - impl<'a> ExpungeTestHarness<'a> { - pub async fn setup( - cptestctx: &'a ControlPlaneTestContext, - disk_id: Uuid, - zpool_id: ZpoolUuid, - ) -> Self { - // TODO normally these allocations are only performed during sled - // reservation so manually inserting records bypasses the need for - // also creating an instance in this test. Further, Nexus created - // for unit and integration tests using the `nexus_test` does _not_ - // perform reconfigurator execution, which means local storage - // datasets are not created, which necessitates extra manual steps - // here to create the relevant rendezvous table entries. Once this - // occurs, this test (and others!) will need to be updated. - - let nexus = &cptestctx.server.server_context().nexus; - let datastore = nexus.datastore(); - - let opctx = OpContext::for_tests( - cptestctx.logctx.log.new(o!()), - datastore.clone(), - ); - - // Manually insert into rendezvous table - - let local_storage_dataset = datastore - .local_storage_unencrypted_dataset_insert_if_not_exists( - &opctx, - RendezvousLocalStorageUnencryptedDataset::new( - DatasetUuid::new_v4(), - zpool_id, - BlueprintUuid::new_v4(), - ), - ) - .await - .unwrap() - .unwrap(); - - // Manually populate both the `disk_type_local_storage` table's - // allocation id, and the allocation table. - - let (.., db_zpool) = LookupPath::new(&opctx, datastore) - .zpool_id(zpool_id) - .fetch() - .await - .unwrap(); - - let allocation = - LocalStorageUnencryptedDatasetAllocation::new_for_tests_only( - DatasetUuid::new_v4(), - Utc::now(), - local_storage_dataset.id(), - ExternalZpoolUuid::from_untyped_uuid( - db_zpool.id().into_untyped_uuid(), - ), - db_zpool.sled_id(), - ByteCount::from_gibibytes_u32(1).into(), - ); - - { - let conn = datastore.pool_connection_for_tests().await.unwrap(); - - use nexus_db_schema::schema::disk_type_local_storage::dsl; - - diesel::update(dsl::disk_type_local_storage) - .filter(dsl::disk_id.eq(disk_id)) - .set( - dsl::local_storage_unencrypted_dataset_allocation_id - .eq(to_db_typed_uuid(allocation.id())), - ) - .execute_async(&*conn) - .await - .unwrap(); - } - - { - let conn = datastore.pool_connection_for_tests().await.unwrap(); - - use nexus_db_schema::schema::local_storage_unencrypted_dataset_allocation::dsl; - - diesel::insert_into( - dsl::local_storage_unencrypted_dataset_allocation, - ) - .values(allocation.clone()) - .execute_async(&*conn) - .await - .unwrap(); - } - - ExpungeTestHarness { cptestctx, zpool_id, allocation } - } - - pub async fn expunge_sled(&self) { - // Set the sled backing the zpool's policy to expunged - - let nexus = &self.cptestctx.server.server_context().nexus; - let datastore = nexus.datastore(); - - let opctx = OpContext::for_tests( - self.cptestctx.logctx.log.new(o!()), - datastore.clone(), - ); - - let (.., db_zpool) = LookupPath::new(&opctx, datastore) - .zpool_id(self.zpool_id) - .fetch() - .await - .unwrap(); - - let (.., authz_sled) = LookupPath::new(&opctx, datastore) - .sled_id(db_zpool.sled_id()) - .lookup_for(authz::Action::Modify) - .await - .unwrap(); - - datastore - .sled_set_policy_to_expunged(&opctx, &authz_sled) - .await - .unwrap(); - } - - pub async fn expunge_disk(&self) { - // Set the physical disk backing the zpool's policy to expunged - - let nexus = &self.cptestctx.server.server_context().nexus; - let datastore = nexus.datastore(); - - let opctx = OpContext::for_tests( - self.cptestctx.logctx.log.new(o!()), - datastore.clone(), - ); - - let (.., db_zpool) = LookupPath::new(&opctx, datastore) - .zpool_id(self.zpool_id) - .fetch() - .await - .unwrap(); - - datastore - .physical_disk_update_policy( - &opctx, - db_zpool.physical_disk_id(), - PhysicalDiskPolicy::Expunged, - ) - .await - .unwrap(); - } - - pub async fn validate_allocation_deleted(&self) { - let nexus = &self.cptestctx.server.server_context().nexus; - let datastore = nexus.datastore(); - - let conn = datastore.pool_connection_for_tests().await.unwrap(); - - use nexus_db_schema::schema::local_storage_unencrypted_dataset_allocation::dsl; - - let allocation = dsl::local_storage_unencrypted_dataset_allocation - .filter(dsl::id.eq(to_db_typed_uuid(self.allocation.id()))) - .select(LocalStorageUnencryptedDatasetAllocation::as_select()) - .get_result_async(&*conn) - .await - .unwrap(); - - assert!(allocation.time_deleted.is_some()); - } - } - - #[nexus_test(server = crate::Server)] - async fn test_delete_local_disk_backed_by_expunged_sled( - cptestctx: &ControlPlaneTestContext, - ) { - let disk_test = DiskTest::new(cptestctx).await; - - let client = &cptestctx.external_client; - let nexus = &cptestctx.server.server_context().nexus; - let datastore = nexus.datastore(); - let opctx = OpContext::for_tests( - cptestctx.logctx.log.new(o!()), - datastore.clone(), - ); - - let project_id = create_project(client, PROJECT_NAME).await.identity.id; - let disk = create_disk(&cptestctx, new_local_disk_create_params).await; - - // Create a local storage disk, manually create an allocation for that - // local storage on a pool, then expunge the sled backing the pool. - - let zpool = disk_test.zpools().next().unwrap(); - - let harness = - ExpungeTestHarness::setup(cptestctx, disk.id(), zpool.id).await; - - harness.expunge_sled().await; - - // Now that the disk' local storage specific information has been - // populated, re-fetch it. - - let disk = datastore.disk_get(&opctx, disk.id()).await.unwrap(); - - // Build the saga DAG with the provided test parameters - let opctx = test_opctx(&cptestctx); - let params = Params { - serialized_authn: Serialized::for_opctx(&opctx), - project_id, - disk, - }; - - nexus.sagas.saga_execute::(params).await.unwrap(); - - harness.validate_allocation_deleted().await; - } - - #[nexus_test(server = crate::Server)] - async fn test_delete_local_disk_backed_by_expunged_zpool( - cptestctx: &ControlPlaneTestContext, - ) { - let disk_test = DiskTest::new(cptestctx).await; - - let client = &cptestctx.external_client; - let nexus = &cptestctx.server.server_context().nexus; - let datastore = nexus.datastore(); - let opctx = OpContext::for_tests( - cptestctx.logctx.log.new(o!()), - datastore.clone(), - ); - - let project_id = create_project(client, PROJECT_NAME).await.identity.id; - let disk = create_disk(&cptestctx, new_local_disk_create_params).await; - - // Create a local storage disk, manually create an allocation for that - // local storage on a pool, then expunge the sled backing the pool. - - let zpool = disk_test.zpools().next().unwrap(); - - let harness = - ExpungeTestHarness::setup(cptestctx, disk.id(), zpool.id).await; - - harness.expunge_disk().await; - - // Now that the disk' local storage specific information has been - // populated, re-fetch it. - - let disk = datastore.disk_get(&opctx, disk.id()).await.unwrap(); - - // Build the saga DAG with the provided test parameters - let opctx = test_opctx(&cptestctx); - let params = Params { - serialized_authn: Serialized::for_opctx(&opctx), - project_id, - disk, - }; - - nexus.sagas.saga_execute::(params).await.unwrap(); - - harness.validate_allocation_deleted().await; - } } diff --git a/nexus/src/app/sagas/instance_start.rs b/nexus/src/app/sagas/instance_start.rs index 91d094a301e..5e2f5d6f962 100644 --- a/nexus/src/app/sagas/instance_start.rs +++ b/nexus/src/app/sagas/instance_start.rs @@ -1173,16 +1173,11 @@ mod test { use core::time::Duration; use nexus_types::identity::Asset as _; - use crate::app::sagas::disk_delete::test::ExpungeTestHarness; - use crate::app::sagas::disk_delete::test::create_disk; - use crate::app::sagas::disk_delete::test::new_local_disk_create_params; use crate::app::{saga::create_saga_dag, sagas::test_helpers}; use dropshot::test_util::ClientTestContext; use nexus_db_queries::authn; - use nexus_test_utils::resource_helpers::DiskTest; use nexus_test_utils::resource_helpers::{ - attach_disk_to_instance, create_default_ip_pools, create_project, - object_create, + create_default_ip_pools, create_project, object_create, }; use nexus_types::external_api::instance::InstanceCpuCount; use nexus_types::external_api::{instance as instance_types, networking}; @@ -1753,92 +1748,4 @@ mod test { cptestctx.teardown().await; } - - #[tokio::test] - async fn test_cannot_start_local_storage_disk_gone() { - let cptestctx = test_helpers::instance_saga_test_builder( - "test_cannot_start_local_storage_disk_gone", - ) - .start::() - .await; - - let client = &cptestctx.external_client; - let nexus = &cptestctx.server.server_context().nexus; - let _project_id = setup_test_project(&client).await; - let opctx = test_helpers::test_opctx(&cptestctx); - - let instance = create_instance(client).await; - - // Create a local storage disk, manually create an allocation for that - // local storage on a pool, attach the disk to the instance, then - // expunge the disk backing the pool. - - let disk_test = DiskTest::new(&cptestctx).await; - let zpool = disk_test.zpools().next().unwrap(); - let disk = create_disk(&cptestctx, new_local_disk_create_params).await; - let harness = - ExpungeTestHarness::setup(&cptestctx, disk.id(), zpool.id).await; - - attach_disk_to_instance( - client, - PROJECT_NAME, - INSTANCE_NAME, - disk.name().as_str(), - ) - .await; - - harness.expunge_disk().await; - - // Run the saga and make sure that the instance does not start - - let instance_id = InstanceUuid::from_untyped_uuid(instance.identity.id); - let db_instance = test_helpers::instance_fetch(&cptestctx, instance_id) - .await - .instance() - .clone(); - - let params = Params { - serialized_authn: authn::saga::Serialized::for_opctx(&opctx), - db_instance, - reason: Reason::User, - }; - - let dag = create_saga_dag::(params).unwrap(); - - let runnable_saga = nexus.sagas.saga_prepare(dag).await.unwrap(); - let saga_result = runnable_saga - .run_to_completion() - .await - .expect("saga execution should have started") - .into_raw_result(); - - let saga_error = saga_result - .kind - .expect_err("saga should fail when ensuring local storage"); - - assert_eq!( - saga_error.error_node_name.as_ref(), - "ensure_local_storage_0", - ); - - let db_instance = - test_helpers::instance_fetch(&cptestctx, instance_id).await; - - assert_eq!( - db_instance.instance().nexus_state, - nexus_db_model::InstanceState::NoVmm - ); - assert!(db_instance.vmm().is_none()); - - assert!( - test_helpers::no_virtual_provisioning_resource_records_exist( - &cptestctx - ) - .await - ); - - assert!(test_helpers::no_virtual_provisioning_collection_records_using_instances(&cptestctx).await); - - cptestctx.teardown().await; - } } diff --git a/nexus/test-utils/src/background.rs b/nexus/test-utils/src/background.rs index 433bc82b878..ba6af10cc53 100644 --- a/nexus/test-utils/src/background.rs +++ b/nexus/test-utils/src/background.rs @@ -6,6 +6,8 @@ use crate::http_testing::NexusRequest; use dropshot::test_util::ClientTestContext; +use nexus_db_queries::context::OpContext; +use nexus_db_queries::db::DataStore; use nexus_lockstep_client::types::BackgroundTask; use nexus_lockstep_client::types::CurrentStatus; use nexus_lockstep_client::types::LastResult; @@ -13,6 +15,8 @@ use nexus_types::internal_api::background::*; use omicron_test_utils::dev::poll::{CondCheckError, wait_for_condition}; use omicron_uuid_kinds::CollectionUuid; use slog::info; +use slog::o; +use std::sync::Arc; use std::time::Duration; /// Given the name of a background task, wait for it to complete if it's @@ -679,3 +683,104 @@ pub async fn run_blueprint_rendezvous( status } + +/// Run the local_storage_delete background task, and assert that there are no +/// reported errors. +pub async fn run_local_storage_delete(internal_client: &ClientTestContext) { + let status = run_local_storage_delete_return_status(internal_client).await; + assert!(status.errors.is_empty()); +} + +/// Run the local_storage_delete background task and return the status. +pub async fn run_local_storage_delete_return_status( + internal_client: &ClientTestContext, +) -> LocalStorageDeleteStatus { + let last_background_task = + activate_background_task(&internal_client, "local_storage_delete") + .await; + + let LastResult::Completed(last_result_completed) = + last_background_task.last + else { + panic!( + "unexpected {:?} returned from volume_delete task", + last_background_task.last, + ); + }; + + serde_json::from_value::( + last_result_completed.details, + ) + .unwrap() +} + +pub async fn wait_for_all_local_storage_deletes( + datastore: &Arc, + lockstep_client: &ClientTestContext, +) { + wait_for_all_local_storage_deletes_impl(datastore, lockstep_client, false) + .await +} + +pub async fn wait_for_all_local_storage_deletes_errors_ok( + datastore: &Arc, + lockstep_client: &ClientTestContext, +) { + wait_for_all_local_storage_deletes_impl(datastore, lockstep_client, true) + .await +} + +async fn wait_for_all_local_storage_deletes_impl( + datastore: &Arc, + lockstep_client: &ClientTestContext, + errors_ok: bool, +) { + wait_for_condition( + || { + let datastore = datastore.clone(); + let opctx = OpContext::for_tests( + lockstep_client.client_log.new(o!()), + datastore.clone(), + ); + + async move { + // Trigger the local storage delete background task. Bail out of + // this loop only when there's no more allocations to clean up. + // + // Be careful not to check if the background tasks performed any + // actions: the fixed point that we're waiting for is for all + // resources to be cleaned up. + + if errors_ok { + let _ = + run_local_storage_delete_return_status(lockstep_client) + .await; + } else { + run_local_storage_delete(lockstep_client).await; + } + + let disks_requiring_work = datastore + .deleted_disks_with_undeleted_local_storage(&opctx) + .await + .unwrap(); + + if !disks_requiring_work.is_empty() { + info!( + &lockstep_client.client_log, + "wait_for_all_local_storage_deletes: {} disks \ + requiring work left", + disks_requiring_work.len(), + ); + + return Err(CondCheckError::<()>::NotYet { status: None }); + } + + Ok(()) + } + }, + &std::time::Duration::from_millis(50), + &std::time::Duration::from_secs(260), + ) + .await + .expect("all deletes finished"); +} diff --git a/nexus/tests/config.test.toml b/nexus/tests/config.test.toml index 6113a8960cc..9840ca7f9a9 100644 --- a/nexus/tests/config.test.toml +++ b/nexus/tests/config.test.toml @@ -239,6 +239,7 @@ audit_log_cleanup.period_secs = 600 audit_log_cleanup.retention_days = 90 audit_log_cleanup.max_deleted_per_activation = 10000 populate_switch_ports.period_secs = 30 +local_storage_delete.period_secs = 10000 [multicast] # Enable multicast functionality for tests (disabled by default in production) diff --git a/nexus/tests/integration_tests/disks.rs b/nexus/tests/integration_tests/disks.rs index dbc6f0ecc66..1fc779b8fcd 100644 --- a/nexus/tests/integration_tests/disks.rs +++ b/nexus/tests/integration_tests/disks.rs @@ -5,13 +5,18 @@ //! Tests basic disk support in the API use super::instances::instance_wait_for_state; +use async_bb8_diesel::AsyncRunQueryDsl; +use diesel::prelude::*; use dropshot::HttpErrorResponseBody; use dropshot::test_util::ClientTestContext; use http::StatusCode; use http::method::Method; use nexus_config::RegionAllocationStrategy; use nexus_db_lookup::LookupPath; +use nexus_db_model::LocalStorageUnencryptedDatasetAllocation; use nexus_db_model::PhysicalDiskPolicy; +use nexus_db_model::to_db_typed_uuid; +use nexus_db_queries::authz; use nexus_db_queries::context::OpContext; use nexus_db_queries::db::datastore; use nexus_db_queries::db::datastore::REGION_REDUNDANCY_THRESHOLD; @@ -19,6 +24,7 @@ use nexus_db_queries::db::datastore::RegionAllocationFor; use nexus_db_queries::db::datastore::RegionAllocationParameters; use nexus_db_queries::db::fixed_data::FLEET_ID; use nexus_test_utils::SLED_AGENT_UUID; +use nexus_test_utils::background::wait_for_all_local_storage_deletes_errors_ok; use nexus_test_utils::http_testing::AuthnMode; use nexus_test_utils::http_testing::NexusRequest; use nexus_test_utils::http_testing::RequestBuilder; @@ -3167,12 +3173,12 @@ async fn test_read_only_disk_different_vcr( /// Test that deleting a local storage disk retries through transient sled /// agent errors. /// -/// This exercises the retry loop in `sdd_delete_local_storage` by: +/// This exercises the `local_storage_delete` background task by: /// /// 1. Creating a local storage disk and starting an instance with it. /// 2. Stopping the instance and detaching the disk. /// 3. Injecting transient 503 errors into the simulated sled agent. -/// 4. Deleting the disk and verifying the saga retries and succeeds. +/// 4. Deleting the disk and verifying the task retries and succeeds. /// /// Time is paused so that exponential backoff sleeps resolve quickly. #[nexus_test] @@ -3243,25 +3249,39 @@ async fn test_delete_local_storage_disk_retries_on_transient_error( disk_post(client, &url_instance_detach_disk, local_disk_name.clone()).await; // Inject transient failures into local storage operations on the sled - // agent. The disk delete saga's `sdd_delete_local_storage` action will + // agent. The background task responsible for deleting local storage will // retry through these. 8 errors comfortably exceeds backon's default // max_times of 3 (catching the regression fixed by #9993), while keeping // virtual time low enough that background task timers don't cause excessive // real I/O during auto-advance. cptestctx.first_sled_agent().set_local_storage_error_count(8); - // Pausing the timer turns on Tokio's auto-advance behavior, making backoff - // sleeps resolve instantly. - tokio::time::pause(); - - // Delete the disk. This triggers the disk delete saga, which must retry - // in the face of injected errors. + // Delete the disk. let disk_url = get_disk_url("local-disk"); NexusRequest::object_delete(client, &disk_url) .authn_as(AuthnMode::PrivilegedUser) .execute() .await - .expect("disk delete should succeed after retrying transient errors"); + .expect("disk delete should succeed"); + + NexusRequest::new( + RequestBuilder::new(client, Method::GET, &disk_url) + .expect_status(Some(StatusCode::NOT_FOUND)), + ) + .authn_as(AuthnMode::PrivilegedUser) + .execute() + .await + .expect("disk should no longer exist"); + + // Ensure the background task eventually deletes the local storage. + + tokio::time::pause(); + + wait_for_all_local_storage_deletes_errors_ok( + &nexus.datastore(), + &cptestctx.lockstep_client, + ) + .await; tokio::time::resume(); @@ -3271,6 +3291,263 @@ async fn test_delete_local_storage_disk_retries_on_transient_error( "not all injected errors were consumed; \ the retry loop may not have been exercised" ); +} + +/// Test that deleting a local storage disk succeeds even if the backing zpool +/// is expunged +#[nexus_test] +async fn test_delete_local_disk_backed_by_expunged_zpool( + cptestctx: &ControlPlaneTestContext, +) { + let client = &cptestctx.external_client; + DiskTest::new(&cptestctx).await; + create_project_and_pool(client).await; + let nexus = &cptestctx.server.server_context().nexus; + + let local_disk_name: Name = "local-disk".parse().unwrap(); + let instance_name = "local-disk-instance"; + + // Create a local storage disk. + let disks_url = get_disks_url(); + NexusRequest::new( + RequestBuilder::new(client, Method::POST, &disks_url) + .body(Some(&disk::DiskCreate { + identity: IdentityMetadataCreateParams { + name: local_disk_name.clone(), + description: "local storage disk".to_string(), + }, + disk_backend: disk::DiskBackend::Local {}, + size: ByteCount::from_gibibytes_u32(1), + })) + .expect_status(Some(StatusCode::CREATED)), + ) + .authn_as(AuthnMode::PrivilegedUser) + .execute() + .await + .expect("created local storage disk"); + + // Create an instance with the local disk attached and start it. Starting + // the instance triggers `sled_reservation_create`, which allocates a + // dataset for the local storage disk. Without this allocation, the disk + // delete saga short-circuits and never reaches the retry loop. + let instance = create_instance_with( + client, + PROJECT_NAME, + instance_name, + &instance::InstanceNetworkInterfaceAttachment::DefaultIpv4, + vec![instance::InstanceDiskAttachment::Attach( + instance::InstanceDiskAttach { name: local_disk_name.clone() }, + )], + Vec::::new(), + true, + Default::default(), + None, + Vec::new(), + ) + .await; + let instance_id = InstanceUuid::from_untyped_uuid(instance.identity.id); + + // Simulate the instance transitioning to Running so the start saga + // completes (including local storage allocation). + instance_simulate(nexus, &instance_id).await; + instance_wait_for_state(client, instance_id, InstanceState::Running).await; + + // Stop the instance. + set_instance_state(client, instance_name, "stop").await; + instance_simulate(nexus, &instance_id).await; + instance_wait_for_state(client, instance_id, InstanceState::Stopped).await; + + // Expunge the zpool backing the allocation + + let datastore = nexus.datastore(); + let opctx = + OpContext::for_tests(cptestctx.logctx.log.new(o!()), datastore.clone()); + + let allocations: Vec<_> = { + let conn = datastore.pool_connection_for_tests().await.unwrap(); + + use nexus_db_schema::schema::local_storage_unencrypted_dataset_allocation::dsl; + + dsl::local_storage_unencrypted_dataset_allocation + .filter(dsl::time_deleted.is_null()) + .select(LocalStorageUnencryptedDatasetAllocation::as_select()) + .load_async(&*conn) + .await + .unwrap() + }; + + assert_eq!(allocations.len(), 1); + + let (.., db_zpool) = LookupPath::new(&opctx, datastore) + .zpool_id(allocations[0].pool_id().upcast()) + .fetch() + .await + .unwrap(); + + datastore + .physical_disk_update_policy( + &opctx, + db_zpool.physical_disk_id(), + PhysicalDiskPolicy::Expunged, + ) + .await + .unwrap(); + + // Detach the disk then delete it. + + let url_instance_detach_disk = + get_disk_detach_url(&instance.identity.id.into()); + disk_post(client, &url_instance_detach_disk, local_disk_name.clone()).await; + + let disk_url = get_disk_url("local-disk"); + NexusRequest::object_delete(client, &disk_url) + .authn_as(AuthnMode::PrivilegedUser) + .execute() + .await + .expect("disk delete should succeed"); + + NexusRequest::new( + RequestBuilder::new(client, Method::GET, &disk_url) + .expect_status(Some(StatusCode::NOT_FOUND)), + ) + .authn_as(AuthnMode::PrivilegedUser) + .execute() + .await + .expect("disk should no longer exist"); + + // Ensure the background task eventually deletes the local storage. + + wait_for_all_local_storage_deletes_errors_ok( + &datastore, + &cptestctx.lockstep_client, + ) + .await; + + // The allocation is only marked deleted after the sled-agent DELETE returns + // ok. + + let allocation = { + let conn = datastore.pool_connection_for_tests().await.unwrap(); + + use nexus_db_schema::schema::local_storage_unencrypted_dataset_allocation::dsl; + + dsl::local_storage_unencrypted_dataset_allocation + .filter(dsl::id.eq(to_db_typed_uuid(allocations[0].id()))) + .select(LocalStorageUnencryptedDatasetAllocation::as_select()) + .get_result_async(&*conn) + .await + .unwrap() + }; + + assert!(allocation.time_deleted.is_some()); +} + +/// Test that deleting a local storage disk succeeds even if the backing sled is +/// expunged +#[nexus_test] +async fn test_delete_local_disk_backed_by_expunged_sled( + cptestctx: &ControlPlaneTestContext, +) { + let client = &cptestctx.external_client; + DiskTest::new(&cptestctx).await; + create_project_and_pool(client).await; + let nexus = &cptestctx.server.server_context().nexus; + + let local_disk_name: Name = "local-disk".parse().unwrap(); + let instance_name = "local-disk-instance"; + + // Create a local storage disk. + let disks_url = get_disks_url(); + NexusRequest::new( + RequestBuilder::new(client, Method::POST, &disks_url) + .body(Some(&disk::DiskCreate { + identity: IdentityMetadataCreateParams { + name: local_disk_name.clone(), + description: "local storage disk".to_string(), + }, + disk_backend: disk::DiskBackend::Local {}, + size: ByteCount::from_gibibytes_u32(1), + })) + .expect_status(Some(StatusCode::CREATED)), + ) + .authn_as(AuthnMode::PrivilegedUser) + .execute() + .await + .expect("created local storage disk"); + + // Create an instance with the local disk attached and start it. Starting + // the instance triggers `sled_reservation_create`, which allocates a + // dataset for the local storage disk. Without this allocation, the disk + // delete saga short-circuits and never reaches the retry loop. + let instance = create_instance_with( + client, + PROJECT_NAME, + instance_name, + &instance::InstanceNetworkInterfaceAttachment::DefaultIpv4, + vec![instance::InstanceDiskAttachment::Attach( + instance::InstanceDiskAttach { name: local_disk_name.clone() }, + )], + Vec::::new(), + true, + Default::default(), + None, + Vec::new(), + ) + .await; + let instance_id = InstanceUuid::from_untyped_uuid(instance.identity.id); + + // Simulate the instance transitioning to Running so the start saga + // completes (including local storage allocation). + instance_simulate(nexus, &instance_id).await; + instance_wait_for_state(client, instance_id, InstanceState::Running).await; + + // Stop the instance so we can detach and delete the disk. + set_instance_state(client, instance_name, "stop").await; + instance_simulate(nexus, &instance_id).await; + instance_wait_for_state(client, instance_id, InstanceState::Stopped).await; + + // Detach the disk. + let url_instance_detach_disk = + get_disk_detach_url(&instance.identity.id.into()); + disk_post(client, &url_instance_detach_disk, local_disk_name.clone()).await; + + // For the sled backing the allocation, set its sled policy to expunged + + let datastore = nexus.datastore(); + let opctx = + OpContext::for_tests(cptestctx.logctx.log.new(o!()), datastore.clone()); + + let allocations: Vec<_> = { + let conn = datastore.pool_connection_for_tests().await.unwrap(); + + use nexus_db_schema::schema::local_storage_unencrypted_dataset_allocation::dsl; + + dsl::local_storage_unencrypted_dataset_allocation + .filter(dsl::time_deleted.is_null()) + .select(LocalStorageUnencryptedDatasetAllocation::as_select()) + .load_async(&*conn) + .await + .unwrap() + }; + + assert_eq!(allocations.len(), 1); + + let (.., authz_sled) = LookupPath::new(&opctx, datastore) + .sled_id(allocations[0].sled_id()) + .lookup_for(authz::Action::Modify) + .await + .unwrap(); + + datastore.sled_set_policy_to_expunged(&opctx, &authz_sled).await.unwrap(); + + // Delete the disk + + let disk_url = get_disk_url("local-disk"); + NexusRequest::object_delete(client, &disk_url) + .authn_as(AuthnMode::PrivilegedUser) + .execute() + .await + .expect("disk delete should succeed"); NexusRequest::new( RequestBuilder::new(client, Method::GET, &disk_url) @@ -3280,6 +3557,32 @@ async fn test_delete_local_storage_disk_retries_on_transient_error( .execute() .await .expect("disk should no longer exist"); + + // Ensure the background task eventually deletes the local storage. + + wait_for_all_local_storage_deletes_errors_ok( + &datastore, + &cptestctx.lockstep_client, + ) + .await; + + // The allocation is only marked deleted after the sled-agent DELETE returns + // ok. + + let allocation = { + let conn = datastore.pool_connection_for_tests().await.unwrap(); + + use nexus_db_schema::schema::local_storage_unencrypted_dataset_allocation::dsl; + + dsl::local_storage_unencrypted_dataset_allocation + .filter(dsl::id.eq(to_db_typed_uuid(allocations[0].id()))) + .select(LocalStorageUnencryptedDatasetAllocation::as_select()) + .get_result_async(&*conn) + .await + .unwrap() + }; + + assert!(allocation.time_deleted.is_some()); } async fn disk_get(client: &ClientTestContext, disk_url: &str) -> Disk { diff --git a/nexus/tests/integration_tests/instances.rs b/nexus/tests/integration_tests/instances.rs index adc8accf6d4..3e4ef795618 100644 --- a/nexus/tests/integration_tests/instances.rs +++ b/nexus/tests/integration_tests/instances.rs @@ -16,6 +16,8 @@ use itertools::Itertools; use nexus_auth::authz::Action; use nexus_db_lookup::AsyncConnection; use nexus_db_lookup::LookupPath; +use nexus_db_model::LocalStorageUnencryptedDatasetAllocation; +use nexus_db_model::PhysicalDiskPolicy; use nexus_db_model::SledResourceVmm; use nexus_db_model::to_db_typed_uuid; use nexus_db_queries::context::OpContext; @@ -10018,3 +10020,115 @@ async fn can_create_instance_with_multiple_nics_and_ephemeral_ip( ) .await; } + +#[nexus_test] +async fn test_cannot_start_local_storage_disk_gone( + cptestctx: &ControlPlaneTestContext, +) { + let client = &cptestctx.external_client; + DiskTest::new(&cptestctx).await; + create_project_and_pool(client).await; + let nexus = &cptestctx.server.server_context().nexus; + + let local_disk_name: Name = "local-disk".parse().unwrap(); + let instance_name = "local-disk-instance"; + + // Create a local storage disk. + let disks_url = get_disks_url(); + NexusRequest::new( + RequestBuilder::new(client, Method::POST, &disks_url) + .body(Some(&disk::DiskCreate { + identity: IdentityMetadataCreateParams { + name: local_disk_name.clone(), + description: "local storage disk".to_string(), + }, + disk_backend: disk::DiskBackend::Local {}, + size: ByteCount::from_gibibytes_u32(1), + })) + .expect_status(Some(StatusCode::CREATED)), + ) + .authn_as(AuthnMode::PrivilegedUser) + .execute() + .await + .expect("created local storage disk"); + + // Create an instance with the local disk attached and start it. Starting + // the instance triggers `sled_reservation_create`, which allocates a + // dataset for the local storage disk. Without this allocation, the disk + // delete saga short-circuits and never reaches the retry loop. + let instance = create_instance_with( + client, + PROJECT_NAME, + instance_name, + &instance::InstanceNetworkInterfaceAttachment::DefaultIpv4, + vec![instance::InstanceDiskAttachment::Attach( + instance::InstanceDiskAttach { name: local_disk_name.clone() }, + )], + Vec::::new(), + true, + Default::default(), + None, + Vec::new(), + ) + .await; + let instance_id = InstanceUuid::from_untyped_uuid(instance.identity.id); + + // Simulate the instance transitioning to Running so the start saga + // completes (including local storage allocation). + instance_simulate(nexus, &instance_id).await; + instance_wait_for_state(client, instance_id, InstanceState::Running).await; + + // Stop the instance. + instance_post(client, instance_name, InstanceOp::Stop).await; + instance_simulate(nexus, &instance_id).await; + instance_wait_for_state(client, instance_id, InstanceState::Stopped).await; + + // Expunge the zpool backing the allocation + + let datastore = nexus.datastore(); + let opctx = + OpContext::for_tests(cptestctx.logctx.log.new(o!()), datastore.clone()); + + let allocations: Vec<_> = { + let conn = datastore.pool_connection_for_tests().await.unwrap(); + + use nexus_db_schema::schema::local_storage_unencrypted_dataset_allocation::dsl; + + dsl::local_storage_unencrypted_dataset_allocation + .filter(dsl::time_deleted.is_null()) + .select(LocalStorageUnencryptedDatasetAllocation::as_select()) + .load_async(&*conn) + .await + .unwrap() + }; + + assert_eq!(allocations.len(), 1); + + let (.., db_zpool) = LookupPath::new(&opctx, datastore) + .zpool_id(allocations[0].pool_id().upcast()) + .fetch() + .await + .unwrap(); + + datastore + .physical_disk_update_policy( + &opctx, + db_zpool.physical_disk_id(), + PhysicalDiskPolicy::Expunged, + ) + .await + .unwrap(); + + // The instance should no longer be able to be started + + NexusRequest::expect_failure( + client, + StatusCode::INTERNAL_SERVER_ERROR, + Method::POST, + get_instance_start_url(instance_name).as_str(), + ) + .authn_as(AuthnMode::PrivilegedUser) + .execute() + .await + .unwrap(); +} diff --git a/nexus/types/src/internal_api/background.rs b/nexus/types/src/internal_api/background.rs index 20e1f38c366..228699a7189 100644 --- a/nexus/types/src/internal_api/background.rs +++ b/nexus/types/src/internal_api/background.rs @@ -1641,6 +1641,16 @@ pub struct PhysicalDiskAdoptionStatus { pub errors: Vec, } +/// The status of a `local_storage_delete` background task activation +#[derive(Serialize, Deserialize, Default, Debug, PartialEq, Eq)] +pub struct LocalStorageDeleteStatus { + pub total_allocations_to_delete: usize, + pub page_size: usize, + pub delete_results: Vec, + pub deallocate_results: Vec, + pub errors: Vec, +} + #[cfg(test)] mod test { use super::TufRepoInfo; diff --git a/sled-agent/src/sim/storage.rs b/sled-agent/src/sim/storage.rs index 8c98973b8ed..05160bef2e7 100644 --- a/sled-agent/src/sim/storage.rs +++ b/sled-agent/src/sim/storage.rs @@ -1238,7 +1238,8 @@ impl Zpool { } pub fn drop_dataset(&mut self, id: DatasetUuid) { - let _ = self.datasets.remove(&id).expect("Failed to get the dataset"); + // Must remain idempotent to repeated requests + let _ = self.datasets.remove(&id); } fn insert_local_storage_unencrypted_dataset( diff --git a/smf/nexus/multi-sled/config-partial.toml b/smf/nexus/multi-sled/config-partial.toml index 9b59caea431..822067e2442 100644 --- a/smf/nexus/multi-sled/config-partial.toml +++ b/smf/nexus/multi-sled/config-partial.toml @@ -135,6 +135,7 @@ audit_log_cleanup.period_secs = 600 audit_log_cleanup.retention_days = 90 audit_log_cleanup.max_deleted_per_activation = 10000 populate_switch_ports.period_secs = 30 +local_storage_delete.period_secs = 30 [default_region_allocation_strategy] # by default, allocate across 3 distinct sleds diff --git a/smf/nexus/single-sled/config-partial.toml b/smf/nexus/single-sled/config-partial.toml index e9826da7435..b7d6606bbc4 100644 --- a/smf/nexus/single-sled/config-partial.toml +++ b/smf/nexus/single-sled/config-partial.toml @@ -135,6 +135,7 @@ audit_log_cleanup.period_secs = 600 audit_log_cleanup.retention_days = 90 audit_log_cleanup.max_deleted_per_activation = 10000 populate_switch_ports.period_secs = 30 +local_storage_delete.period_secs = 30 [default_region_allocation_strategy] # by default, allocate without requirement for distinct sleds.