From dfe1303f93c10ee8cd97ea3c80ef92ebc39c080d Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Tue, 25 Aug 2026 21:41:25 +0000 Subject: [PATCH 01/18] Delete local storage in a background task Move the part of Nexus responsible for deleting local storage resources from the disk delete saga into a background task. --- dev-tools/omdb/src/bin/omdb/db.rs | 6 +- dev-tools/omdb/src/bin/omdb/nexus.rs | 36 +++ dev-tools/omdb/tests/env.out | 12 + dev-tools/omdb/tests/successes.out | 20 ++ nexus-config/src/nexus_config.rs | 16 + nexus/background-task-interface/src/init.rs | 1 + nexus/db-queries/src/db/datastore/disk.rs | 14 +- .../src/db/datastore/local_storage.rs | 59 ++++ nexus/examples/config-second.toml | 1 + nexus/examples/config.toml | 1 + nexus/src/app/background/init.rs | 13 + .../background/tasks/local_storage_delete.rs | 276 ++++++++++++++++++ nexus/src/app/background/tasks/mod.rs | 1 + nexus/src/app/disk.rs | 12 +- nexus/src/app/sagas/disk_delete.rs | 159 +--------- nexus/test-utils/src/background.rs | 105 +++++++ nexus/tests/config.test.toml | 1 + nexus/tests/integration_tests/disks.rs | 44 +-- nexus/types/src/internal_api/background.rs | 8 + smf/nexus/multi-sled/config-partial.toml | 1 + smf/nexus/single-sled/config-partial.toml | 1 + 21 files changed, 612 insertions(+), 175 deletions(-) create mode 100644 nexus/src/app/background/tasks/local_storage_delete.rs diff --git a/dev-tools/omdb/src/bin/omdb/db.rs b/dev-tools/omdb/src/bin/omdb/db.rs index 5474e0762fd..c0e9da99c44 100644 --- a/dev-tools/omdb/src/bin/omdb/db.rs +++ b/dev-tools/omdb/src/bin/omdb/db.rs @@ -2646,11 +2646,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()) @@ -2659,7 +2659,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(&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 4500580ad50..a8dc5fcac97 100644 --- a/dev-tools/omdb/src/bin/omdb/nexus.rs +++ b/dev-tools/omdb/src/bin/omdb/nexus.rs @@ -69,6 +69,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; @@ -1426,6 +1427,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: {:?} \ @@ -4317,6 +4321,38 @@ 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 { + delete_results, + deallocate_results, + errors, + } = &status; + + println!(" result of deleting local storage:"); + for result in delete_results { + println!(" > {result}"); + } + + println!(" result 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 5d39c5509e6..bceb61c98c4 100644 --- a/dev-tools/omdb/tests/env.out +++ b/dev-tools/omdb/tests/env.out @@ -154,6 +154,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 @@ -422,6 +426,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 @@ -677,6 +685,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 ea3cd79f126..d43f79da605 100644 --- a/dev-tools/omdb/tests/successes.out +++ b/dev-tools/omdb/tests/successes.out @@ -389,6 +389,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 @@ -864,6 +868,14 @@ 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 + result of deleting local storage: + result of deallocating local storage: + errors: 0 + task: "lookup_region_port" configured period: every m last completed activation: , triggered by @@ -1586,6 +1598,14 @@ 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 + result of deleting local storage: + result 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 b18a621ab9e..d18cb95f84a 100644 --- a/nexus-config/src/nexus_config.rs +++ b/nexus-config/src/nexus_config.rs @@ -477,6 +477,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] @@ -1086,6 +1088,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 { @@ -1376,6 +1386,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 @@ -1656,6 +1667,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: @@ -1772,6 +1787,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 fd2c8f30fec..14ea61e40ab 100644 --- a/nexus/background-task-interface/src/init.rs +++ b/nexus/background-task-interface/src/init.rs @@ -63,6 +63,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 7f6e31b2215..3128b177614 100644 --- a/nexus/db-queries/src/db/datastore/disk.rs +++ b/nexus/db-queries/src/db/datastore/disk.rs @@ -510,7 +510,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(&conn, disk).await } /// Return a `datastore::Disk` given a `model::Disk` @@ -518,15 +520,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. + /// 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( &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 +536,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 +554,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) diff --git a/nexus/db-queries/src/db/datastore/local_storage.rs b/nexus/db-queries/src/db/datastore/local_storage.rs index 7ba521d2807..0385a2e29aa 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; @@ -451,4 +453,61 @@ 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.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(&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. Return an error instead of panicking. + return Err(Error::internal_error(&format!( + "disk {} should be the 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 05b1546094f..9878d136191 100644 --- a/nexus/examples/config-second.toml +++ b/nexus/examples/config-second.toml @@ -218,6 +218,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 6ec9adc0eba..cfd9b643191 100644 --- a/nexus/examples/config.toml +++ b/nexus/examples/config.toml @@ -202,6 +202,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 10cfaac0f76..4f70f08f8fb 100644 --- a/nexus/src/app/background/init.rs +++ b/nexus/src/app/background/init.rs @@ -118,6 +118,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; @@ -280,6 +281,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 @@ -376,6 +378,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. @@ -1314,6 +1317,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..7898fa8516f --- /dev/null +++ b/nexus/src/app/background/tasks/local_storage_delete.rs @@ -0,0 +1,276 @@ +// 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 cleand 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 serde_json::json; +use sled_agent_client::types::LocalStorageDatasetDeleteRequest; +use slog::Logger; +use std::sync::Arc; + +pub struct LocalStorageDeleter { + datastore: Arc, + reqwest_client: reqwest::Client, +} + +/// Functions that poll for an expected change will either return that the +/// change occurred, or that the current activation of this task has to wait for +/// the change to occur in a future activations of this task. +#[derive(PartialEq)] +enum DeleteResult { + Deleted, + + WaitForNextActivation, +} + +impl LocalStorageDeleter { + pub fn new(datastore: Arc) -> Self { + LocalStorageDeleter { + datastore, + reqwest_client: reqwest::Client::new(), + } + } + + async fn delete_unencrypted_allocation( + &self, + log: &Logger, + opctx: &OpContext, + allocation: &LocalStorageUnencryptedDatasetAllocation, + status: &mut LocalStorageDeleteStatus, + ) -> 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 s = format!( + "error calling check_sled_in_service for sled \ + {sled_id}: {e}", + ); + + error!(log, "{s}"); + status.errors.push(s); + + return DeleteResult::WaitForNextActivation; + } + }; + + if !sled_in_service { + // Sled's been expunged, so consider the local storage deleted. + return DeleteResult::Deleted; + } + + 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 s = format!( + "error calling check_zpool_in_service for zpool \ + {zpool_id}: {e}", + ); + + error!(log, "{s}"); + status.errors.push(s); + + return DeleteResult::WaitForNextActivation; + } + }; + + if !zpool_in_service { + // The disk backing the zpool's been expunged, so consider the local + // storage deleted. + return DeleteResult::Deleted; + } + + // 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 s = format!( + "error calling sled_client_ext for allocation {}: {e}", + allocation.id(), + ); + + error!(log, "{s}"); + status.errors.push(s); + + return DeleteResult::WaitForNextActivation; + } + }; + + match sled_agent_client.local_storage_dataset_delete(&request).await { + Ok(_) => DeleteResult::Deleted, + + Err(e) => { + let s = + format!("error sending local_storage_dataset_delete: {e}"); + + error!(log, "{s}"); + status.errors.push(s); + + DeleteResult::WaitForNextActivation + } + } + } +} + +impl BackgroundTask for LocalStorageDeleter { + fn activate<'a>( + &'a mut self, + opctx: &'a OpContext, + ) -> BoxFuture<'a, serde_json::Value> { + async { + let log = &opctx.log; + let mut status = LocalStorageDeleteStatus::default(); + + let 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 \ + unencrypted_allocations_for_deleted_disks: {e}" + ); + + error!(log, "{s}"); + status.errors.push(s); + + return json!(status); + } + }; + + for disk in disks_needing_clean_up { + 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, + &mut status, + ) + .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::WaitForNextActivation => { + // 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_allocations(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_allocations: {e}" + ); + + error!(log, "{s}"); + status.errors.push(s); + + continue; + } + } + } + + json!(status) + } + .boxed() + } +} diff --git a/nexus/src/app/background/tasks/mod.rs b/nexus/src/app/background/tasks/mod.rs index 7533226b12b..1c97a7f3b31 100644 --- a/nexus/src/app/background/tasks/mod.rs +++ b/nexus/src/app/background/tasks/mod.rs @@ -31,6 +31,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..2116c136294 100644 --- a/nexus/src/app/disk.rs +++ b/nexus/src/app/disk.rs @@ -418,13 +418,23 @@ impl super::Nexus { let saga_params = sagas::disk_delete::Params { serialized_authn: authn::saga::Serialized::for_opctx(opctx), project_id: project.id(), - disk, + disk: disk.clone(), }; self.sagas .saga_execute::(saga_params) .await?; + match disk { + datastore::Disk::Crucible(_) => { + // For now, do nothing. Stay tuned! + } + + datastore::Disk::LocalStorage(_) => { + 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..8117b1da686 100644 --- a/nexus/src/app/sagas/disk_delete.rs +++ b/nexus/src/app/sagas/disk_delete.rs @@ -5,25 +5,16 @@ use super::ActionRegistry; use super::NexusActionContext; use super::NexusSaga; -use crate::app::InlineErrorChain; use crate::app::sagas::SagaInitError; use crate::app::sagas::declare_saga_actions; -use crate::app::sagas::sled_out_of_service_gone_check; use crate::app::sagas::volume_delete; -use crate::app::sagas::zpool_out_of_service_gone_check; use nexus_db_queries::authn; use nexus_db_queries::db; use nexus_db_queries::db::datastore; use nexus_types::saga::saga_action_failed; use omicron_common::api::external::DiskState; -use omicron_common::api::external::Error; -use omicron_common::backoff::backon_retry_policy_internal_service; -use progenitor_extras::retry::{ - GoneCheckResult, retry_operation_while_indefinitely, -}; use serde::Deserialize; use serde::Serialize; -use sled_agent_client::types::LocalStorageDatasetDeleteRequest; use steno::ActionError; use steno::Node; use uuid::Uuid; @@ -49,12 +40,6 @@ declare_saga_actions! { + sdd_account_space - sdd_account_space_undo } - DELETE_LOCAL_STORAGE -> "delete_local_storage" { - + sdd_delete_local_storage - } - DEALLOCATE_LOCAL_STORAGE -> "deallocate_local_storage" { - + sdd_deallocate_local_storage - } } // disk delete saga: definition @@ -112,12 +97,8 @@ 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 done entirely in + // the local_storage_delete background task. } } @@ -208,131 +189,6 @@ async fn sdd_account_space_undo( Ok(()) } -async fn sdd_delete_local_storage( - sagactx: NexusActionContext, -) -> Result<(), ActionError> { - let osagactx = sagactx.user_data(); - let params = sagactx.saga_params::()?; - let opctx = crate::context::op_context_for_saga_action( - &sagactx, - ¶ms.serialized_authn, - ); - - let datastore::Disk::LocalStorage(disk) = params.disk else { - unreachable!( - "check during `make_saga_dag` should have ensured disk type is \ - local storage" - ); - }; - - let Some(allocation) = disk.local_storage_dataset_allocation else { - // Nothing to do! - return Ok(()); - }; - - let sled_id = allocation.sled_id(); - let zpool_id = allocation.pool_id().upcast(); - - let request = LocalStorageDatasetDeleteRequest { - zpool_id: allocation.pool_id(), - dataset_id: allocation.id(), - encrypted_at_rest: allocation.encrypted_at_rest(), - }; - - // Get a sled agent client - - let sled_agent_client = osagactx - .nexus() - .sled_client(&sled_id) - .await - .map_err(saga_action_failed)?; - - // Ensure that the local storage is deleted - - let delete_operation = || async { - sled_agent_client.local_storage_dataset_delete(&request).await - }; - - // Bail out of the retry loop if either the disk or sled is no longer - // in-service. - let gone_check = || async { - match sled_out_of_service_gone_check( - osagactx.datastore(), - &opctx, - sled_id, - ) - .await? - { - GoneCheckResult::StillAvailable => { - // proceed to zpool check - } - - GoneCheckResult::Gone => { - return Ok(GoneCheckResult::Gone); - } - } - - zpool_out_of_service_gone_check(osagactx.datastore(), &opctx, zpool_id) - .await - }; - - let log = osagactx.log().clone(); - let result = retry_operation_while_indefinitely( - backon_retry_policy_internal_service(), - delete_operation, - gone_check, - |notification| { - slog::warn!( - log, - "failed to delete local storage dataset, retrying in {:?}", - notification.delay; - InlineErrorChain::new(¬ification.error), - ); - }, - ) - .await; - - match result { - Ok(_) => Ok(()), - - // In this case, if the particular disk hosting this local storage was - // expunged, or if the sled was expunged, then proceed with the rest of - // the saga. - Err(e) if e.is_gone() => Ok(()), - - Err(e) => Err(saga_action_failed(Error::internal_error(&format!( - "failed to delete local storage: {}", - InlineErrorChain::new(&e) - )))), - } -} - -async fn sdd_deallocate_local_storage( - sagactx: NexusActionContext, -) -> Result<(), ActionError> { - let osagactx = sagactx.user_data(); - let params = sagactx.saga_params::()?; - let opctx = crate::context::op_context_for_saga_action( - &sagactx, - ¶ms.serialized_authn, - ); - - let datastore::Disk::LocalStorage(disk) = params.disk else { - unreachable!( - "check during `make_saga_dag` should have ensured disk type is \ - local storage" - ); - }; - - osagactx - .datastore() - .delete_local_storage_dataset_allocations(&opctx, &disk) - .await - .map_err(saga_action_failed)?; - - Ok(()) -} - #[cfg(test)] pub(crate) mod test { use crate::{ @@ -353,6 +209,7 @@ pub(crate) mod test { use nexus_db_queries::authz; use nexus_db_queries::context::OpContext; use nexus_db_queries::db::datastore::Disk; + use nexus_test_utils::background::wait_for_all_local_storage_deletes; use nexus_test_utils::resource_helpers::DiskTest; use nexus_test_utils::resource_helpers::create_project; use nexus_test_utils_macros::nexus_test; @@ -641,6 +498,16 @@ pub(crate) mod test { let nexus = &self.cptestctx.server.server_context().nexus; let datastore = nexus.datastore(); + // Run the local storage delete background task to completion + + wait_for_all_local_storage_deletes( + &datastore, + &self.cptestctx.lockstep_client, + ) + .await; + + // Then check all allocations were deleted + let conn = datastore.pool_connection_for_tests().await.unwrap(); use nexus_db_schema::schema::local_storage_unencrypted_dataset_allocation::dsl; diff --git a/nexus/test-utils/src/background.rs b/nexus/test-utils/src/background.rs index 7580749b19d..996933c25e3 100644 --- a/nexus/test-utils/src/background.rs +++ b/nexus/test-utils/src/background.rs @@ -6,12 +6,16 @@ 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; use nexus_types::internal_api::background::*; use omicron_test_utils::dev::poll::{CondCheckError, wait_for_condition}; 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 @@ -581,3 +585,104 @@ pub async fn run_blueprint_rendezvous(lockstep_client: &ClientTestContext) { ) .unwrap(); } + +/// 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 c8dc72e98f0..71bbd956c82 100644 --- a/nexus/tests/config.test.toml +++ b/nexus/tests/config.test.toml @@ -238,6 +238,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..f471edaa229 100644 --- a/nexus/tests/integration_tests/disks.rs +++ b/nexus/tests/integration_tests/disks.rs @@ -19,6 +19,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 +3168,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,34 +3244,20 @@ 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"); - - tokio::time::resume(); - - assert_eq!( - cptestctx.first_sled_agent().local_storage_error_remaining(), - 0, - "not all injected errors were consumed; \ - the retry loop may not have been exercised" - ); + .expect("disk delete should succeed"); NexusRequest::new( RequestBuilder::new(client, Method::GET, &disk_url) @@ -3280,6 +3267,25 @@ 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. + + tokio::time::pause(); + + wait_for_all_local_storage_deletes_errors_ok( + &nexus.datastore(), + &cptestctx.lockstep_client, + ) + .await; + + tokio::time::resume(); + + assert_eq!( + cptestctx.first_sled_agent().local_storage_error_remaining(), + 0, + "not all injected errors were consumed; \ + the retry loop may not have been exercised" + ); } async fn disk_get(client: &ClientTestContext, disk_url: &str) -> Disk { diff --git a/nexus/types/src/internal_api/background.rs b/nexus/types/src/internal_api/background.rs index 104bd5f2226..0a9d5de3cc5 100644 --- a/nexus/types/src/internal_api/background.rs +++ b/nexus/types/src/internal_api/background.rs @@ -1502,6 +1502,14 @@ 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 delete_results: Vec, + pub deallocate_results: Vec, + pub errors: Vec, +} + #[cfg(test)] mod test { use super::TufRepoInfo; diff --git a/smf/nexus/multi-sled/config-partial.toml b/smf/nexus/multi-sled/config-partial.toml index 6981c1ee612..409c890571a 100644 --- a/smf/nexus/multi-sled/config-partial.toml +++ b/smf/nexus/multi-sled/config-partial.toml @@ -134,6 +134,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 ce56cb14814..caacec86227 100644 --- a/smf/nexus/single-sled/config-partial.toml +++ b/smf/nexus/single-sled/config-partial.toml @@ -134,6 +134,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. From 2aac5c07aa90e69bd9202221ff35c86410e48066 Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Wed, 26 Aug 2026 15:18:17 +0000 Subject: [PATCH 02/18] typo --- nexus/src/app/background/tasks/local_storage_delete.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/nexus/src/app/background/tasks/local_storage_delete.rs b/nexus/src/app/background/tasks/local_storage_delete.rs index 7898fa8516f..a34dd71c064 100644 --- a/nexus/src/app/background/tasks/local_storage_delete.rs +++ b/nexus/src/app/background/tasks/local_storage_delete.rs @@ -4,7 +4,7 @@ //! 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 cleand up. +//! local storage allocation records when those resources have been cleaned up. use crate::app::background::BackgroundTask; use futures::FutureExt; From 7433b4b8a71f7cf75849d64cd883174076ef24a1 Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Wed, 26 Aug 2026 15:18:43 +0000 Subject: [PATCH 03/18] wrong name! --- nexus/src/app/background/tasks/local_storage_delete.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/nexus/src/app/background/tasks/local_storage_delete.rs b/nexus/src/app/background/tasks/local_storage_delete.rs index a34dd71c064..3154bc13b20 100644 --- a/nexus/src/app/background/tasks/local_storage_delete.rs +++ b/nexus/src/app/background/tasks/local_storage_delete.rs @@ -169,7 +169,7 @@ impl BackgroundTask for LocalStorageDeleter { Err(e) => { let s = format!( "error calling \ - unencrypted_allocations_for_deleted_disks: {e}" + deleted_disks_with_undeleted_local_storage: {e}" ); error!(log, "{s}"); From aa4095c1d7e62bed4cffc801852b6f853dd1689d Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Wed, 26 Aug 2026 15:32:26 +0000 Subject: [PATCH 04/18] InlineErrorChain for sled_client_ext and progenitor function --- nexus/src/app/background/tasks/local_storage_delete.rs | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/nexus/src/app/background/tasks/local_storage_delete.rs b/nexus/src/app/background/tasks/local_storage_delete.rs index 3154bc13b20..7bf09e53ee9 100644 --- a/nexus/src/app/background/tasks/local_storage_delete.rs +++ b/nexus/src/app/background/tasks/local_storage_delete.rs @@ -17,6 +17,7 @@ use nexus_types::internal_api::background::LocalStorageDeleteStatus; 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 { @@ -123,8 +124,9 @@ impl LocalStorageDeleter { Err(e) => { let s = format!( - "error calling sled_client_ext for allocation {}: {e}", + "error calling sled_client_ext for allocation {}: {}", allocation.id(), + InlineErrorChain::new(&e), ); error!(log, "{s}"); @@ -138,8 +140,10 @@ impl LocalStorageDeleter { Ok(_) => DeleteResult::Deleted, Err(e) => { - let s = - format!("error sending local_storage_dataset_delete: {e}"); + let s = format!( + "error sending local_storage_dataset_delete: {}", + InlineErrorChain::new(&e), + ); error!(log, "{s}"); status.errors.push(s); From b362308ca52caee1911ec2b7c36551d5fa43b6de Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Mon, 14 Sep 2026 15:48:59 +0000 Subject: [PATCH 05/18] return as many disks for deletion as possible do not bail with an error --- nexus/db-queries/src/db/datastore/local_storage.rs | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/nexus/db-queries/src/db/datastore/local_storage.rs b/nexus/db-queries/src/db/datastore/local_storage.rs index 0385a2e29aa..4c5aaa34746 100644 --- a/nexus/db-queries/src/db/datastore/local_storage.rs +++ b/nexus/db-queries/src/db/datastore/local_storage.rs @@ -495,11 +495,16 @@ impl DataStore { // 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. Return an error instead of panicking. - return Err(Error::internal_error(&format!( - "disk {} should be the local storage, not crucible", + // 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) => { From d48805fb513b0e34bfa4657b43cd36f7fa905e41 Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Wed, 16 Sep 2026 19:19:01 +0000 Subject: [PATCH 06/18] simulated sled agent should not panic if dataset is already deleted --- sled-agent/src/sim/storage.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) 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( From 0c57be095a10ed05735c938fbdab4b9120cae017 Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Wed, 16 Sep 2026 20:13:29 +0000 Subject: [PATCH 07/18] Many more DeleteResult variants: - specific expunged variants - specific error variant (WaitForNextActivation is not an error) Adding to status done in one place now Set a short timeout for reqwest client InlineErrorChain everywhere --- .../background/tasks/local_storage_delete.rs | 133 +++++++++++++----- 1 file changed, 99 insertions(+), 34 deletions(-) diff --git a/nexus/src/app/background/tasks/local_storage_delete.rs b/nexus/src/app/background/tasks/local_storage_delete.rs index 7bf09e53ee9..b9ab6fe0778 100644 --- a/nexus/src/app/background/tasks/local_storage_delete.rs +++ b/nexus/src/app/background/tasks/local_storage_delete.rs @@ -25,21 +25,40 @@ pub struct LocalStorageDeleter { reqwest_client: reqwest::Client, } -/// Functions that poll for an expected change will either return that the -/// change occurred, or that the current activation of this task has to wait for -/// the change to occur in a future activations of this task. #[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_secs(5); + LocalStorageDeleter { datastore, - reqwest_client: reqwest::Client::new(), + // 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(), } } @@ -48,7 +67,6 @@ impl LocalStorageDeleter { log: &Logger, opctx: &OpContext, allocation: &LocalStorageUnencryptedDatasetAllocation, - status: &mut LocalStorageDeleteStatus, ) -> DeleteResult { let sled_id = allocation.sled_id(); let zpool_id = allocation.pool_id().upcast(); @@ -61,21 +79,19 @@ impl LocalStorageDeleter { Ok(sled_in_service) => sled_in_service, Err(e) => { - let s = format!( + let message = format!( "error calling check_sled_in_service for sled \ - {sled_id}: {e}", + {sled_id}: {}", + InlineErrorChain::new(&e), ); - error!(log, "{s}"); - status.errors.push(s); - - return DeleteResult::WaitForNextActivation; + return DeleteResult::Error { message }; } }; if !sled_in_service { // Sled's been expunged, so consider the local storage deleted. - return DeleteResult::Deleted; + return DeleteResult::SledExpunged; } let zpool_in_service = @@ -84,22 +100,20 @@ impl LocalStorageDeleter { Ok(zpool_in_service) => zpool_in_service, Err(e) => { - let s = format!( + let message = format!( "error calling check_zpool_in_service for zpool \ - {zpool_id}: {e}", + {zpool_id}: {}", + InlineErrorChain::new(&e), ); - error!(log, "{s}"); - status.errors.push(s); - - return DeleteResult::WaitForNextActivation; + return DeleteResult::Error { message }; } }; if !zpool_in_service { // The disk backing the zpool's been expunged, so consider the local // storage deleted. - return DeleteResult::Deleted; + return DeleteResult::ZpoolExpunged; } // Now that all checks are done, get a sled agent client and make the @@ -123,32 +137,34 @@ impl LocalStorageDeleter { Ok(client) => client, Err(e) => { - let s = format!( + let message = format!( "error calling sled_client_ext for allocation {}: {}", allocation.id(), InlineErrorChain::new(&e), ); - error!(log, "{s}"); - status.errors.push(s); - - return DeleteResult::WaitForNextActivation; + 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 s = format!( + let message = format!( "error sending local_storage_dataset_delete: {}", InlineErrorChain::new(&e), ); - error!(log, "{s}"); - status.errors.push(s); - - DeleteResult::WaitForNextActivation + DeleteResult::Error { message } } } } @@ -173,7 +189,8 @@ impl BackgroundTask for LocalStorageDeleter { Err(e) => { let s = format!( "error calling \ - deleted_disks_with_undeleted_local_storage: {e}" + deleted_disks_with_undeleted_local_storage: {}", + InlineErrorChain::new(&e), ); error!(log, "{s}"); @@ -201,7 +218,6 @@ impl BackgroundTask for LocalStorageDeleter { log, opctx, &allocation, - &mut status, ) .await { @@ -219,7 +235,56 @@ impl BackgroundTask for LocalStorageDeleter { // 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; @@ -255,6 +320,7 @@ impl BackgroundTask for LocalStorageDeleter { disk.id(), allocation.id(), ); + info!(log, "{s}"); status.deallocate_results.push(s); } @@ -262,13 +328,12 @@ impl BackgroundTask for LocalStorageDeleter { Err(e) => { let s = format!( "error calling \ - delete_local_storage_dataset_allocations: {e}" + delete_local_storage_dataset_allocations: {}", + InlineErrorChain::new(&e), ); error!(log, "{s}"); status.errors.push(s); - - continue; } } } From ba447ac2984d4e88c30a4a9a33d6b1c0703c6864 Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Wed, 16 Sep 2026 20:19:34 +0000 Subject: [PATCH 08/18] not calling sled_client_ext for an allocation --- nexus/src/app/background/tasks/local_storage_delete.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/nexus/src/app/background/tasks/local_storage_delete.rs b/nexus/src/app/background/tasks/local_storage_delete.rs index b9ab6fe0778..c374c7ee080 100644 --- a/nexus/src/app/background/tasks/local_storage_delete.rs +++ b/nexus/src/app/background/tasks/local_storage_delete.rs @@ -138,8 +138,7 @@ impl LocalStorageDeleter { Err(e) => { let message = format!( - "error calling sled_client_ext for allocation {}: {}", - allocation.id(), + "error calling sled_client_ext for sled {sled_id}: {}", InlineErrorChain::new(&e), ); From 6b0fccd7734d3eff9f95ef0141935fe512e8d3e9 Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Wed, 16 Sep 2026 20:36:13 +0000 Subject: [PATCH 09/18] separate activate and activate_impl --- .../background/tasks/local_storage_delete.rs | 316 +++++++++--------- 1 file changed, 164 insertions(+), 152 deletions(-) diff --git a/nexus/src/app/background/tasks/local_storage_delete.rs b/nexus/src/app/background/tasks/local_storage_delete.rs index c374c7ee080..270faef1761 100644 --- a/nexus/src/app/background/tasks/local_storage_delete.rs +++ b/nexus/src/app/background/tasks/local_storage_delete.rs @@ -167,177 +167,189 @@ impl LocalStorageDeleter { } } } -} -impl BackgroundTask for LocalStorageDeleter { - fn activate<'a>( - &'a mut self, - opctx: &'a OpContext, - ) -> BoxFuture<'a, serde_json::Value> { - async { - let log = &opctx.log; - let mut status = LocalStorageDeleteStatus::default(); + async fn activate_impl( + &self, + opctx: &OpContext, + ) -> LocalStorageDeleteStatus { + let log = &opctx.log; + let mut status = LocalStorageDeleteStatus::default(); + + let disks_needing_clean_up = match self + .datastore + .deleted_disks_with_undeleted_local_storage(opctx) + .await + { + Ok(v) => v, - let 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), + ); - Err(e) => { - let s = format!( - "error calling \ - deleted_disks_with_undeleted_local_storage: {}", - InlineErrorChain::new(&e), - ); + error!(log, "{s}"); + status.errors.push(s); - error!(log, "{s}"); - status.errors.push(s); + return status; + } + }; - return json!(status); - } + for disk in disks_needing_clean_up { + let Some(allocation) = &disk.local_storage_dataset_allocation + else { + // No allocation was made for this disk + continue; }; - for disk in disks_needing_clean_up { - 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; - } + // 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. } - } - 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(), - ); + 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; + } - error!(log, "{s}"); - status.errors.push(s); + DeleteResult::Error { message } => { + info!(log, "{message}"); + status.errors.push(message); - continue; + // Cannot deallocate the record until deletion + // succeeds. + continue; + } } } - match self - .datastore - .delete_local_storage_dataset_allocations(opctx, &disk) - .await - { - Ok(()) => { - let s = format!( - "deallocated disk {} allocation {}", - disk.id(), - allocation.id(), - ); - - info!(log, "{s}"); - status.deallocate_results.push(s); - } + 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(), + ); - Err(e) => { - let s = format!( - "error calling \ - delete_local_storage_dataset_allocations: {}", - InlineErrorChain::new(&e), - ); + error!(log, "{s}"); + status.errors.push(s); - error!(log, "{s}"); - status.errors.push(s); - } + continue; } } - json!(status) + match self + .datastore + .delete_local_storage_dataset_allocations(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_allocations: {}", + 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() } From 4912ae0e8eacd768873891f6e1197793ce67d9b1 Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Wed, 16 Sep 2026 20:44:55 +0000 Subject: [PATCH 10/18] change misleading function name, not plural --- nexus/db-queries/src/db/datastore/local_storage.rs | 8 ++++---- nexus/src/app/background/tasks/local_storage_delete.rs | 4 ++-- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/nexus/db-queries/src/db/datastore/local_storage.rs b/nexus/db-queries/src/db/datastore/local_storage.rs index 4c5aaa34746..6d47443b97a 100644 --- a/nexus/db-queries/src/db/datastore/local_storage.rs +++ b/nexus/db-queries/src/db/datastore/local_storage.rs @@ -282,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, @@ -298,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 { diff --git a/nexus/src/app/background/tasks/local_storage_delete.rs b/nexus/src/app/background/tasks/local_storage_delete.rs index 270faef1761..49bd8a8617b 100644 --- a/nexus/src/app/background/tasks/local_storage_delete.rs +++ b/nexus/src/app/background/tasks/local_storage_delete.rs @@ -303,7 +303,7 @@ impl LocalStorageDeleter { match self .datastore - .delete_local_storage_dataset_allocations(opctx, &disk) + .delete_local_storage_dataset_allocation(opctx, &disk) .await { Ok(()) => { @@ -320,7 +320,7 @@ impl LocalStorageDeleter { Err(e) => { let s = format!( "error calling \ - delete_local_storage_dataset_allocations: {}", + delete_local_storage_dataset_allocation: {}", InlineErrorChain::new(&e), ); From ad2f6a13d9863cef38f372bdef12f427242d9c27 Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Wed, 16 Sep 2026 20:58:20 +0000 Subject: [PATCH 11/18] append _on_connection --- dev-tools/omdb/src/bin/omdb/db.rs | 2 +- nexus/db-queries/src/db/datastore/disk.rs | 4 ++-- nexus/db-queries/src/db/datastore/local_storage.rs | 5 ++++- 3 files changed, 7 insertions(+), 4 deletions(-) diff --git a/dev-tools/omdb/src/bin/omdb/db.rs b/dev-tools/omdb/src/bin/omdb/db.rs index c0e9da99c44..170e7428301 100644 --- a/dev-tools/omdb/src/bin/omdb/db.rs +++ b/dev-tools/omdb/src/bin/omdb/db.rs @@ -2659,7 +2659,7 @@ async fn cmd_db_disk_info( .context("failed to find disk")? }; - match datastore.disk_get_with_model(&conn, 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/nexus/db-queries/src/db/datastore/disk.rs b/nexus/db-queries/src/db/datastore/disk.rs index 3128b177614..f2a472f061c 100644 --- a/nexus/db-queries/src/db/datastore/disk.rs +++ b/nexus/db-queries/src/db/datastore/disk.rs @@ -512,7 +512,7 @@ impl DataStore { let conn = self.pool_connection_authorized(opctx).await?; - self.disk_get_with_model(&conn, disk).await + self.disk_get_with_model_on_connection(&conn, disk).await } /// Return a `datastore::Disk` given a `model::Disk` @@ -522,7 +522,7 @@ impl DataStore { /// LookupPath induced permissions check and should only called from omdb. /// 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( + pub async fn disk_get_with_model_on_connection( &self, conn: &async_bb8_diesel::Connection, disk: model::Disk, diff --git a/nexus/db-queries/src/db/datastore/local_storage.rs b/nexus/db-queries/src/db/datastore/local_storage.rs index 6d47443b97a..f2a8ecf27ef 100644 --- a/nexus/db-queries/src/db/datastore/local_storage.rs +++ b/nexus/db-queries/src/db/datastore/local_storage.rs @@ -490,7 +490,10 @@ impl DataStore { let mut disks = Vec::with_capacity(found_disks.len()); for found_disk in found_disks { - match self.disk_get_with_model(&conn, found_disk).await? { + 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 From 2d9619a35db3800b88f32cfe7b17d6ff6d99bf4e Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Wed, 16 Sep 2026 21:01:46 +0000 Subject: [PATCH 12/18] specifically note everything is done in the saga --- nexus/src/app/disk.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/nexus/src/app/disk.rs b/nexus/src/app/disk.rs index 2116c136294..27a3643261d 100644 --- a/nexus/src/app/disk.rs +++ b/nexus/src/app/disk.rs @@ -427,7 +427,8 @@ impl super::Nexus { match disk { datastore::Disk::Crucible(_) => { - // For now, do nothing. Stay tuned! + // For now, do nothing - all Crucible related clean-up is done + // in the disk delete saga. Stay tuned! } datastore::Disk::LocalStorage(_) => { From 94cde1c2a02815252d9dd4e3e75c9827ddeb4b3b Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Thu, 17 Sep 2026 14:39:39 +0000 Subject: [PATCH 13/18] operate on only 128 disks per task invocation --- dev-tools/omdb/src/bin/omdb/nexus.rs | 16 ++++++++++-- .../background/tasks/local_storage_delete.rs | 25 +++++++++++++++++-- nexus/types/src/internal_api/background.rs | 2 ++ 3 files changed, 39 insertions(+), 4 deletions(-) diff --git a/dev-tools/omdb/src/bin/omdb/nexus.rs b/dev-tools/omdb/src/bin/omdb/nexus.rs index a8dc5fcac97..b0c6ebb634f 100644 --- a/dev-tools/omdb/src/bin/omdb/nexus.rs +++ b/dev-tools/omdb/src/bin/omdb/nexus.rs @@ -4330,17 +4330,29 @@ fn print_task_local_storage_delete(details: &serde_json::Value) { Ok(status) => { let LocalStorageDeleteStatus { + total_allocations_to_delete, + page_size, delete_results, deallocate_results, errors, } = &status; - println!(" result of deleting local storage:"); + 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!(" result of deallocating local storage:"); + println!(" results of deallocating local storage:"); for result in deallocate_results { println!(" > {result}"); } diff --git a/nexus/src/app/background/tasks/local_storage_delete.rs b/nexus/src/app/background/tasks/local_storage_delete.rs index 49bd8a8617b..813b143bfce 100644 --- a/nexus/src/app/background/tasks/local_storage_delete.rs +++ b/nexus/src/app/background/tasks/local_storage_delete.rs @@ -48,7 +48,7 @@ enum DeleteResult { impl LocalStorageDeleter { pub fn new(datastore: Arc) -> Self { - let duration = std::time::Duration::from_secs(5); + let duration = std::time::Duration::from_millis(250); LocalStorageDeleter { datastore, @@ -196,7 +196,28 @@ impl LocalStorageDeleter { } }; - for disk in disks_needing_clean_up { + // 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; + + 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 diff --git a/nexus/types/src/internal_api/background.rs b/nexus/types/src/internal_api/background.rs index 0a9d5de3cc5..ffdc75f39dd 100644 --- a/nexus/types/src/internal_api/background.rs +++ b/nexus/types/src/internal_api/background.rs @@ -1505,6 +1505,8 @@ pub struct PhysicalDiskAdoptionStatus { /// 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, From a65ed3be4aa19e606a9e7d4d7f695468e11eda54 Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Thu, 17 Sep 2026 14:47:41 +0000 Subject: [PATCH 14/18] randomize the list, do not wedge the task --- nexus/src/app/background/tasks/local_storage_delete.rs | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/nexus/src/app/background/tasks/local_storage_delete.rs b/nexus/src/app/background/tasks/local_storage_delete.rs index 813b143bfce..19ea90c411b 100644 --- a/nexus/src/app/background/tasks/local_storage_delete.rs +++ b/nexus/src/app/background/tasks/local_storage_delete.rs @@ -14,6 +14,7 @@ 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; @@ -175,7 +176,7 @@ impl LocalStorageDeleter { let log = &opctx.log; let mut status = LocalStorageDeleteStatus::default(); - let disks_needing_clean_up = match self + let mut disks_needing_clean_up = match self .datastore .deleted_disks_with_undeleted_local_storage(opctx) .await @@ -217,6 +218,11 @@ impl LocalStorageDeleter { 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 { From b7b9b98bcdc3bf3211e38ba8c50f4f3429cd1eb3 Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Thu, 17 Sep 2026 15:21:48 +0000 Subject: [PATCH 15/18] add missed opctx check! --- nexus/db-queries/src/db/datastore/local_storage.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/nexus/db-queries/src/db/datastore/local_storage.rs b/nexus/db-queries/src/db/datastore/local_storage.rs index f2a8ecf27ef..b616d0e9b76 100644 --- a/nexus/db-queries/src/db/datastore/local_storage.rs +++ b/nexus/db-queries/src/db/datastore/local_storage.rs @@ -459,6 +459,7 @@ impl DataStore { &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?; From 6f4abf31251e8b0bda6526429663a6427b211f4f Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Thu, 17 Sep 2026 20:38:00 +0000 Subject: [PATCH 16/18] use a single transaction to delete a disk and update the provisioning collection --- nexus/db-queries/src/db/datastore/disk.rs | 76 ++++++++++++++++++++++- nexus/src/app/disk.rs | 38 ++++++++---- 2 files changed, 101 insertions(+), 13 deletions(-) diff --git a/nexus/db-queries/src/db/datastore/disk.rs b/nexus/db-queries/src/db/datastore/disk.rs index f2a472f061c..c6e2b8aed91 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; @@ -1551,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<_> = @@ -2085,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/src/app/disk.rs b/nexus/src/app/disk.rs index 27a3643261d..3578165eef2 100644 --- a/nexus/src/app/disk.rs +++ b/nexus/src/app/disk.rs @@ -415,23 +415,37 @@ 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: disk.clone(), - }; - - self.sagas - .saga_execute::(saga_params) - .await?; - match disk { datastore::Disk::Crucible(_) => { - // For now, do nothing - all Crucible related clean-up is done - // in the disk delete saga. Stay tuned! + // 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?; } 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(); } } From e90c669ce42d3a79c755f90b1d2cf2d438648386 Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Fri, 18 Sep 2026 01:51:48 +0000 Subject: [PATCH 17/18] omdb expectorate --- dev-tools/omdb/tests/successes.out | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/dev-tools/omdb/tests/successes.out b/dev-tools/omdb/tests/successes.out index d43f79da605..8ebd3b8fde2 100644 --- a/dev-tools/omdb/tests/successes.out +++ b/dev-tools/omdb/tests/successes.out @@ -872,8 +872,10 @@ task: "local_storage_delete" configured period: every h m s last completed activation: , triggered by started at (s ago) and ran for ms - result of deleting local storage: - result of deallocating local storage: + 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" @@ -1602,8 +1604,10 @@ task: "local_storage_delete" configured period: every h m s last completed activation: , triggered by started at (s ago) and ran for ms - result of deleting local storage: - result of deallocating local storage: + 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" From 7d59f364a8dae34bebe285e85cf8c3f828611974 Mon Sep 17 00:00:00 2001 From: James MacMahon Date: Mon, 21 Sep 2026 18:31:23 +0000 Subject: [PATCH 18/18] undo all changes to the disk delete saga when Nexus is interrupted for mupdate, any existing sagas will be interrupted, including disk delete sagas for disks backed by local storage. changing the structure or semantics of the saga is a bad idea, so revert any of those changes. however no new disk delete sagas can be created for disks backed by local storage. not being able to create disk delete sagas for disks backed by local storage caused some tests to fail - move those tests from the saga crate to the integration test crate now that those support local storage, fulfilling the TODO in the ExpungeTestHarness --- nexus/src/app/sagas/disk_delete.rs | 457 +++++++-------------- nexus/src/app/sagas/instance_start.rs | 95 +---- nexus/tests/integration_tests/disks.rs | 297 +++++++++++++ nexus/tests/integration_tests/instances.rs | 114 +++++ 4 files changed, 564 insertions(+), 399 deletions(-) diff --git a/nexus/src/app/sagas/disk_delete.rs b/nexus/src/app/sagas/disk_delete.rs index 8117b1da686..0193160cc56 100644 --- a/nexus/src/app/sagas/disk_delete.rs +++ b/nexus/src/app/sagas/disk_delete.rs @@ -5,16 +5,25 @@ use super::ActionRegistry; use super::NexusActionContext; use super::NexusSaga; +use crate::app::InlineErrorChain; use crate::app::sagas::SagaInitError; use crate::app::sagas::declare_saga_actions; +use crate::app::sagas::sled_out_of_service_gone_check; use crate::app::sagas::volume_delete; +use crate::app::sagas::zpool_out_of_service_gone_check; use nexus_db_queries::authn; use nexus_db_queries::db; use nexus_db_queries::db::datastore; use nexus_types::saga::saga_action_failed; use omicron_common::api::external::DiskState; +use omicron_common::api::external::Error; +use omicron_common::backoff::backon_retry_policy_internal_service; +use progenitor_extras::retry::{ + GoneCheckResult, retry_operation_while_indefinitely, +}; use serde::Deserialize; use serde::Serialize; +use sled_agent_client::types::LocalStorageDatasetDeleteRequest; use steno::ActionError; use steno::Node; use uuid::Uuid; @@ -40,6 +49,12 @@ declare_saga_actions! { + sdd_account_space - sdd_account_space_undo } + DELETE_LOCAL_STORAGE -> "delete_local_storage" { + + sdd_delete_local_storage + } + DEALLOCATE_LOCAL_STORAGE -> "deallocate_local_storage" { + + sdd_deallocate_local_storage + } } // disk delete saga: definition @@ -97,8 +112,18 @@ impl NexusSaga for SagaDiskDelete { } datastore::Disk::LocalStorage(_) => { - // Clean up of the local storage resources is done entirely in - // the local_storage_delete background task. + // 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(), + ))); } } @@ -189,42 +214,148 @@ async fn sdd_account_space_undo( Ok(()) } +async fn sdd_delete_local_storage( + sagactx: NexusActionContext, +) -> Result<(), ActionError> { + let osagactx = sagactx.user_data(); + let params = sagactx.saga_params::()?; + let opctx = crate::context::op_context_for_saga_action( + &sagactx, + ¶ms.serialized_authn, + ); + + let datastore::Disk::LocalStorage(disk) = params.disk else { + unreachable!( + "check during `make_saga_dag` should have ensured disk type is \ + local storage" + ); + }; + + let Some(allocation) = disk.local_storage_dataset_allocation else { + // Nothing to do! + return Ok(()); + }; + + let sled_id = allocation.sled_id(); + let zpool_id = allocation.pool_id().upcast(); + + let request = LocalStorageDatasetDeleteRequest { + zpool_id: allocation.pool_id(), + dataset_id: allocation.id(), + encrypted_at_rest: allocation.encrypted_at_rest(), + }; + + // Get a sled agent client + + let sled_agent_client = osagactx + .nexus() + .sled_client(&sled_id) + .await + .map_err(saga_action_failed)?; + + // Ensure that the local storage is deleted + + let delete_operation = || async { + sled_agent_client.local_storage_dataset_delete(&request).await + }; + + // Bail out of the retry loop if either the disk or sled is no longer + // in-service. + let gone_check = || async { + match sled_out_of_service_gone_check( + osagactx.datastore(), + &opctx, + sled_id, + ) + .await? + { + GoneCheckResult::StillAvailable => { + // proceed to zpool check + } + + GoneCheckResult::Gone => { + return Ok(GoneCheckResult::Gone); + } + } + + zpool_out_of_service_gone_check(osagactx.datastore(), &opctx, zpool_id) + .await + }; + + let log = osagactx.log().clone(); + let result = retry_operation_while_indefinitely( + backon_retry_policy_internal_service(), + delete_operation, + gone_check, + |notification| { + slog::warn!( + log, + "failed to delete local storage dataset, retrying in {:?}", + notification.delay; + InlineErrorChain::new(¬ification.error), + ); + }, + ) + .await; + + match result { + Ok(_) => Ok(()), + + // In this case, if the particular disk hosting this local storage was + // expunged, or if the sled was expunged, then proceed with the rest of + // the saga. + Err(e) if e.is_gone() => Ok(()), + + Err(e) => Err(saga_action_failed(Error::internal_error(&format!( + "failed to delete local storage: {}", + InlineErrorChain::new(&e) + )))), + } +} + +async fn sdd_deallocate_local_storage( + sagactx: NexusActionContext, +) -> Result<(), ActionError> { + let osagactx = sagactx.user_data(); + let params = sagactx.saga_params::()?; + let opctx = crate::context::op_context_for_saga_action( + &sagactx, + ¶ms.serialized_authn, + ); + + let datastore::Disk::LocalStorage(disk) = params.disk else { + unreachable!( + "check during `make_saga_dag` should have ensured disk type is \ + local storage" + ); + }; + + osagactx + .datastore() + .delete_local_storage_dataset_allocation(&opctx, &disk) + .await + .map_err(saga_action_failed)?; + + Ok(()) +} + #[cfg(test)] pub(crate) mod test { use crate::{ 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::background::wait_for_all_local_storage_deletes; use nexus_test_utils::resource_helpers::DiskTest; use nexus_test_utils::resource_helpers::create_project; 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; @@ -253,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, @@ -339,277 +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(); - - // Run the local storage delete background task to completion - - wait_for_all_local_storage_deletes( - &datastore, - &self.cptestctx.lockstep_client, - ) - .await; - - // Then check all allocations were deleted - - 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 2a55b83848a..6612152eca4 100644 --- a/nexus/src/app/sagas/instance_start.rs +++ b/nexus/src/app/sagas/instance_start.rs @@ -1166,16 +1166,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}; @@ -1746,92 +1741,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/tests/integration_tests/disks.rs b/nexus/tests/integration_tests/disks.rs index f471edaa229..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; @@ -3288,6 +3293,298 @@ async fn test_delete_local_storage_disk_retries_on_transient_error( ); } +/// 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) + .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()); +} + async fn disk_get(client: &ClientTestContext, disk_url: &str) -> Disk { NexusRequest::object_get(client, disk_url) .authn_as(AuthnMode::PrivilegedUser) diff --git a/nexus/tests/integration_tests/instances.rs b/nexus/tests/integration_tests/instances.rs index e9a5e9cfc7e..6ed38d3277e 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; @@ -10083,3 +10085,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(); +}