diff --git a/common/src/region.rs b/common/src/region.rs index f5ef662cb..1595d48cd 100644 --- a/common/src/region.rs +++ b/common/src/region.rs @@ -256,6 +256,7 @@ impl Default for RegionDefinition { } } } + impl RegionDefinition { pub fn test_default( database_read_version: usize, diff --git a/downstairs/Cargo.toml b/downstairs/Cargo.toml index 2a8893d7e..6dddfa8e5 100644 --- a/downstairs/Cargo.toml +++ b/downstairs/Cargo.toml @@ -4,6 +4,7 @@ version = "0.0.1" authors = ["Joshua M. Clulow ", "Alan Hanson extent, + region::ExtentState::Closed => panic!("dump on closed extent!"), + }; let en = e.number(); /* @@ -90,7 +95,8 @@ pub async fn dump_region( continue; } } - let inner = e.inner().await; + + let inner = e.inner.lock().await; /* * Create the ExtentMeta struct for this directory's extent diff --git a/downstairs/src/lib.rs b/downstairs/src/lib.rs index e5a5fa02f..1e3bbd7ee 100644 --- a/downstairs/src/lib.rs +++ b/downstairs/src/lib.rs @@ -864,19 +864,16 @@ where extent_id, } => { let msg = { - let mut d = ad.lock().await; + let d = ad.lock().await; debug!(d.log, "{} Close extent {}", repair_id, extent_id); - match d.region.extents.get_mut(*extent_id) { - Some(ext) => { - let (_, _, _) = ext.close().await?; - Message::RepairAckId { - repair_id: *repair_id, - } - } - None => Message::ExtentError { + match d.region.close_extent(*extent_id).await { + Ok(_) => Message::RepairAckId { + repair_id: *repair_id, + }, + Err(error) => Message::ExtentError { repair_id: *repair_id, extent_id: *extent_id, - error: CrucibleError::InvalidExtent, + error, }, } }; @@ -892,7 +889,7 @@ where dest_clients, } => { let msg = { - let mut d = ad.lock().await; + let d = ad.lock().await; debug!( d.log, "{} Repair extent {} source:[{}] {:?} dest:{:?}", @@ -926,7 +923,7 @@ where extent_id, } => { let msg = { - let mut d = ad.lock().await; + let d = ad.lock().await; debug!(d.log, "{} Reopen extent {}", repair_id, extent_id); match d.region.reopen_extent(*extent_id).await { Ok(()) => Message::RepairAckId { @@ -1803,7 +1800,7 @@ pub struct ActiveUpstairs { */ #[derive(Debug)] pub struct Downstairs { - pub region: Region, + pub region: Arc, lossy: bool, // Test flag, enables pauses and skipped jobs read_errors: bool, // Test flag write_errors: bool, // Test flag @@ -1835,7 +1832,7 @@ impl Downstairs { ))), }; Downstairs { - region, + region: Arc::new(region), lossy, read_errors, write_errors, @@ -2169,10 +2166,8 @@ impl Downstairs { let result = if !self.is_active(job.upstairs_connection) { error!(self.log, "Upstairs inactive error"); Err(CrucibleError::UpstairsInactive) - } else if let Some(ext) = self.region.extents.get_mut(*extent) { - ext.close().await } else { - Err(CrucibleError::InvalidExtent) + self.region.close_extent(*extent).await }; debug!( self.log, @@ -2216,17 +2211,10 @@ impl Downstairs { .await { Err(f_res) => Err(f_res), - Ok(_) => { - if let Some(ext) = - self.region.extents.get_mut(*extent) - { - ext.close().await - } else { - Err(CrucibleError::InvalidExtent) - } - } + Ok(_) => self.region.close_extent(*extent).await, } }; + debug!( self.log, "FlushClose:{} extent {} deps:{:?} res:{} f:{} g:{}", @@ -6014,7 +6002,7 @@ mod test { let ads = create_test_downstairs(bs, es, ec, &dir).await?; - let _ = start_downstairs( + let _jh = start_downstairs( ads, "127.0.0.1".parse().unwrap(), None, diff --git a/downstairs/src/region.rs b/downstairs/src/region.rs index a09992384..c356d7e35 100644 --- a/downstairs/src/region.rs +++ b/downstairs/src/region.rs @@ -7,7 +7,8 @@ use std::fs::{rename, File, OpenOptions}; use std::io::{BufReader, IoSlice, IoSliceMut, Read, Seek, SeekFrom, Write}; use std::net::SocketAddr; use std::path::{Path, PathBuf}; -use tokio::sync::{Mutex, MutexGuard}; +use tokio::sync::Mutex; +use tokio::task::JoinHandle; use anyhow::{anyhow, bail, Result}; use futures::TryStreamExt; @@ -32,15 +33,11 @@ pub struct Extent { block_size: u64, extent_size: Block, iov_max: usize, - /// Inner contains information about the actual extent file that holds - /// the data, and the metadata (stored in the database) about that - /// extent. - /// - /// If Some(), it means the extent file and database metadata for - /// it are opened. - /// If None, it means the extent is currently - /// closed (and possibly being updated out of band). - inner: Option>, + + /// Inner contains information about the actual extent file that holds the + /// data, the metadata (stored in the database) about that extent, and the + /// set of dirty blocks that have been written to since last flush. + pub inner: Mutex, } /// BlockContext, with the addition of block index and on_disk_hash @@ -61,6 +58,14 @@ pub struct Inner { dirty_blocks: HashSet, } +/// An extent can be Opened or Closed. If Closed, it is probably being updated +/// out of band. If Opened, then this extent is accepting operations. +#[derive(Debug)] +pub enum ExtentState { + Opened(Arc), + Closed, +} + impl Inner { pub fn gen_number(&self) -> Result { let mut stmt = self.metadb.prepare_cached( @@ -624,11 +629,11 @@ impl Extent { block_size: def.block_size(), extent_size: def.extent_size(), iov_max: Extent::get_iov_max()?, - inner: Some(Mutex::new(Inner { + inner: Mutex::new(Inner { file, metadb, dirty_blocks: HashSet::new(), - })), + }), }; // Clean out any irrelevant block contexts, which may be present if @@ -640,16 +645,21 @@ impl Extent { Ok(extent) } + #[cfg(test)] + pub async fn dirty(&self) -> bool { + self.inner.lock().await.dirty().unwrap() + } + /** * Close an extent and the metadata db files for it. */ - pub async fn close(&mut self) -> Result<(u64, u64, bool), CrucibleError> { - let inner = self.inner.as_ref().unwrap().lock().await; + pub async fn close(self) -> Result<(u64, u64, bool), CrucibleError> { + let inner = self.inner.lock().await; + let gen = inner.gen_number().unwrap(); let flush = inner.flush_number().unwrap(); let dirty = inner.dirty().unwrap(); - drop(inner); - self.inner = None; + Ok((gen, flush, dirty)) } @@ -813,11 +823,11 @@ impl Extent { block_size: def.block_size(), extent_size: def.extent_size(), iov_max: Extent::get_iov_max()?, - inner: Some(Mutex::new(Inner { + inner: Mutex::new(Inner { file, metadb, dirty_blocks: HashSet::new(), - })), + }), }) } @@ -825,10 +835,10 @@ impl Extent { * Create the copy directory for this extent. */ fn create_copy_dir>( - &self, dir: P, + eid: usize, ) -> Result { - let cp = copy_dir(dir, self.number); + let cp = copy_dir(dir, eid as u32); /* * Verify the copy directory does not exist @@ -846,12 +856,12 @@ impl Extent { * remote downstairs. */ fn create_copy_file( - &self, mut copy_dir: PathBuf, + eid: usize, extension: Option, ) -> Result { // Get the base extent name before we consider the actual Type - let name = extent_file_name(self.number, ExtentType::Data); + let name = extent_file_name(eid as u32, ExtentType::Data); copy_dir.push(name); if let Some(extension) = extension { let ext = format!("{}", extension); @@ -871,10 +881,6 @@ impl Extent { Ok(file) } - pub async fn inner(&self) -> MutexGuard<'_, Inner> { - self.inner.as_ref().unwrap().lock().await - } - pub fn number(&self) -> u32 { self.number } @@ -892,7 +898,8 @@ impl Extent { cdt::extent__read__start!(|| { (job_id, self.number, requests.len() as u64) }); - let inner = self.inner().await; + + let inner = self.inner.lock().await; // This code batches up operations for contiguous regions of // ReadRequests, so we can perform larger read syscalls and sqlite @@ -1043,12 +1050,12 @@ impl Extent { (job_id, self.number, writes.len() as u64) }); - let mut inner = self.inner().await; + let mut inner_guard = self.inner.lock().await; // I realize this looks like some nonsense but what this is doing is // borrowing the inner up-front from the MutexGuard, which will allow // us to later disjointly borrow fields. Basically, we're helping the // borrow-checker do its job. - let inner = &mut *inner; + let inner = &mut *inner_guard; for write in writes { self.check_input(write.offset, &write.data)?; @@ -1261,7 +1268,7 @@ impl Extent { job_id: u64, log: &Logger, ) -> Result<(), CrucibleError> { - let mut inner = self.inner().await; + let mut inner = self.inner.lock().await; if !inner.dirty()? { /* @@ -1444,7 +1451,7 @@ impl Extent { &self, force_override_dirty: bool, ) -> Result<(), CrucibleError> { - let mut inner = self.inner().await; + let mut inner = self.inner.lock().await; if !force_override_dirty && !inner.dirty()? { return Ok(()); @@ -1512,7 +1519,7 @@ impl Extent { pub struct Region { pub dir: PathBuf, def: RegionDefinition, - pub extents: Vec, + pub extents: Vec>>, read_only: bool, log: Logger, } @@ -1720,6 +1727,15 @@ impl Region { self.def.get_encrypted() } + async fn get_opened_extent(&self, eid: usize) -> Arc { + match &*self.extents[eid].lock().await { + ExtentState::Opened(extent) => extent.clone(), + ExtentState::Closed => { + panic!("attempting to get closed extent {}", eid) + } + } + } + /** * If our extent_count is higher than the number of populated entries * we have in our extents Vec, then open all the new extent files and @@ -1750,13 +1766,15 @@ impl Region { ) .await }; - these_extents.push(extent?); + these_extents.push(Arc::new(Mutex::new(ExtentState::Opened( + Arc::new(extent?), + )))); } self.extents.extend(these_extents); for eid in next_eid..self.def.extent_count() { - assert_eq!(self.extents[eid as usize].number, eid); + assert_eq!(self.get_opened_extent(eid as usize).await.number, eid); } assert_eq!(self.def.extent_count() as usize, self.extents.len()); @@ -1767,10 +1785,11 @@ impl Region { * Walk the list of all extents and find any that are not open. * Open any extents that are not. */ - pub async fn reopen_all_extents(&mut self) -> Result<()> { + pub async fn reopen_all_extents(&self) -> Result<()> { let mut to_open = Vec::new(); for (i, extent) in self.extents.iter().enumerate() { - if extent.inner.is_none() { + let inner = extent.lock().await; + if matches!(*inner, ExtentState::Closed) { to_open.push(i); } } @@ -1785,10 +1804,7 @@ impl Region { /** * Re open an extent that was previously closed */ - pub async fn reopen_extent( - &mut self, - eid: usize, - ) -> Result<(), CrucibleError> { + pub async fn reopen_extent(&self, eid: usize) -> Result<(), CrucibleError> { /* * Make sure the extent : * @@ -1796,8 +1812,8 @@ impl Region { * - matches our eid * - is not read-only */ - assert!(self.extents[eid].inner.is_none()); - assert_eq!(self.extents[eid].number, eid as u32); + let mut mg = self.extents[eid].lock().await; + assert!(matches!(*mg, ExtentState::Closed)); assert!(!self.read_only); let new_extent = Extent::open( @@ -1808,10 +1824,36 @@ impl Region { &self.log, ) .await?; - self.extents[eid] = new_extent; + + *mg = ExtentState::Opened(Arc::new(new_extent)); + Ok(()) } + pub async fn close_extent( + &self, + eid: usize, + ) -> Result<(u64, u64, bool), CrucibleError> { + let mut extent_state = self.extents[eid].lock().await; + + let open_extent = + std::mem::replace(&mut *extent_state, ExtentState::Closed); + + match open_extent { + ExtentState::Opened(extent) => { + // extent here is Arc, and closing only makes sense if + // the reference count is 1. + let inner = Arc::into_inner(extent); + assert!(inner.is_some()); + inner.unwrap().close().await + } + + ExtentState::Closed => { + panic!("close on closed extent {}!", eid); + } + } + } + /** * Repair an extent from another downstairs * @@ -1846,14 +1888,15 @@ impl Region { * D. Only then, open extent. */ pub async fn repair_extent( - &mut self, + &self, eid: usize, repair_addr: SocketAddr, ) -> Result<(), CrucibleError> { // Make sure the extent: // is currently closed, matches our eid, is not read-only - assert!(self.extents[eid].inner.is_none()); - assert_eq!(self.extents[eid].number, eid as u32); + let mg = self.extents[eid].lock().await; + assert!(matches!(*mg, ExtentState::Closed)); + drop(mg); assert!(!self.read_only); self.get_extent_copy(eid, repair_addr).await?; @@ -1874,12 +1917,14 @@ impl Region { * copy_dir to replace_dir. */ pub async fn get_extent_copy( - &mut self, + &self, eid: usize, repair_addr: SocketAddr, ) -> Result<(), CrucibleError> { // An extent must be closed before we replace its files. - assert!(self.extents[eid].inner.is_none()); + let mg = self.extents[eid].lock().await; + assert!(matches!(*mg, ExtentState::Closed)); + drop(mg); // Make sure copy, replace, and cleanup directories don't exist yet. // We don't need them yet, but if they do exist, then something @@ -1893,8 +1938,7 @@ impl Region { ); } - let extent = &self.extents[eid]; - let copy_dir = extent.create_copy_dir(&self.dir)?; + let copy_dir = Extent::create_copy_dir(&self.dir, eid)?; info!(self.log, "Created copy dir {:?}", copy_dir); // XXX TLS someday? Authentication? @@ -1932,7 +1976,8 @@ impl Region { } // First, copy the main extent data file. - let extent_copy = extent.create_copy_file(copy_dir.clone(), None)?; + let extent_copy = + Extent::create_copy_file(copy_dir.clone(), eid, None)?; let repair_stream = match repair_server .get_extent_file(eid as u32, FileType::Data) .await @@ -1950,8 +1995,11 @@ impl Region { save_stream_to_file(extent_copy, repair_stream.into_inner()).await?; // The .db file is also required to exist for any valid extent. - let extent_db = - extent.create_copy_file(copy_dir.clone(), Some(ExtentType::Db))?; + let extent_db = Extent::create_copy_file( + copy_dir.clone(), + eid, + Some(ExtentType::Db), + )?; let repair_stream = match repair_server .get_extent_file(eid as u32, FileType::Db) .await @@ -1973,8 +2021,9 @@ impl Region { let filename = extent_file_name(eid as u32, opt_file.clone()); if repair_files.contains(&filename) { - let extent_shm = extent.create_copy_file( + let extent_shm = Extent::create_copy_file( copy_dir.clone(), + eid, Some(opt_file.clone()), )?; let repair_stream = match repair_server @@ -2055,14 +2104,19 @@ impl Region { pub async fn flush_numbers(&self) -> Result> { let mut result = Vec::with_capacity(self.extents.len()); - for e in &self.extents { - result.push(e.inner().await.flush_number()?); + for eid in 0..self.extents.len() { + let extent = self.get_opened_extent(eid).await; + result.push(extent.inner.lock().await.flush_number()?); } if result.len() > 12 { - println!("Current flush_numbers [0..12]: {:?}", &result[0..12]); + info!( + self.log, + "Current flush_numbers [0..12]: {:?}", + &result[0..12] + ); } else { - println!("Current flush_numbers [0..12]: {:?}", result); + info!(self.log, "Current flush_numbers [0..12]: {:?}", result); } Ok(result) @@ -2070,16 +2124,18 @@ impl Region { pub async fn gen_numbers(&self) -> Result> { let mut result = Vec::with_capacity(self.extents.len()); - for e in &self.extents { - result.push(e.inner().await.gen_number()?); + for eid in 0..self.extents.len() { + let extent = self.get_opened_extent(eid).await; + result.push(extent.inner.lock().await.gen_number()?); } Ok(result) } pub async fn dirty(&self) -> Result> { let mut result = Vec::with_capacity(self.extents.len()); - for e in &self.extents { - result.push(e.inner().await.dirty()?); + for eid in 0..self.extents.len() { + let extent = self.get_opened_extent(eid).await; + result.push(extent.inner.lock().await.dirty()?); } Ok(result) } @@ -2147,7 +2203,7 @@ impl Region { cdt::os__write__start!(|| job_id); } for eid in batched_writes.keys() { - let extent = &self.extents[*eid]; + let extent = self.get_opened_extent(*eid).await; let writes = batched_writes.get(eid).unwrap(); extent .write(job_id, &writes[..], only_write_unwritten) @@ -2188,7 +2244,7 @@ impl Region { if request.eid == _eid { batched_reads.push(request); } else { - let extent = &self.extents[_eid as usize]; + let extent = self.get_opened_extent(_eid as usize).await; extent .read(job_id, &batched_reads[..], &mut responses) .await?; @@ -2205,7 +2261,7 @@ impl Region { } if let Some(_eid) = eid { - let extent = &self.extents[_eid as usize]; + let extent = self.get_opened_extent(_eid as usize).await; extent .read(job_id, &batched_reads[..], &mut responses) .await?; @@ -2235,7 +2291,7 @@ impl Region { gen_number ); - let extent = &self.extents[eid]; + let extent = self.get_opened_extent(eid).await; extent.flush(flush_number, gen_number, 0, &self.log).await?; Ok(()) @@ -2272,14 +2328,40 @@ impl Region { } None => self.def.extent_count().try_into().unwrap(), }; + + // Spawn parallel tasks for the flush + let mut join_handles: Vec>> = + Vec::with_capacity(extent_count); + for eid in 0..extent_count { - let extent = &self.extents[eid]; - extent - .flush(flush_number, gen_number, job_id, &self.log) - .await?; + let extent = self.get_opened_extent(eid).await; + let log = self.log.clone(); + let jh = tokio::spawn(async move { + extent.flush(flush_number, gen_number, job_id, &log).await + }); + join_handles.push(jh); } + + // Wait for all flushes to finish - wait until after + // cdt::os__flush__done to check the results and bail out. + let mut results = Vec::with_capacity(extent_count); + for join_handle in join_handles { + results.push( + join_handle + .await + .map_err(|e| CrucibleError::GenericError(e.to_string())), + ); + } + cdt::os__flush__done!(|| job_id); + for result in results { + // If any extent flush failed, then return that as an error. Because + // the results were all collected above, each extent flush has + // completed at this point. + result??; + } + // snapshots currently only work with ZFS if cfg!(feature = "zfs_snapshot") { if let Some(snapshot_details) = snapshot_details { @@ -2658,7 +2740,7 @@ mod test { block_size: 512, extent_size: Block::new_512(100), iov_max: Extent::get_iov_max().unwrap(), - inner: Some(Mutex::new(inn)), + inner: Mutex::new(inn), } } @@ -2796,10 +2878,9 @@ mod test { Region::create(&dir, new_region_options(), csl()).await?; region.extend(3).await?; - let ext_one = &mut region.extents[1]; let cp = copy_dir(&dir, 1); - assert!(ext_one.create_copy_dir(&dir).is_ok()); + assert!(Extent::create_copy_dir(&dir, 1).is_ok()); assert!(Path::new(&cp).exists()); assert!(remove_copy_cleanup_dir(&dir, 1).is_ok()); assert!(!Path::new(&cp).exists()); @@ -2817,9 +2898,8 @@ mod test { .unwrap(); region.extend(3).await.unwrap(); - let ext_one = &mut region.extents[1]; - ext_one.create_copy_dir(&dir).unwrap(); - let res = ext_one.create_copy_dir(&dir); + Extent::create_copy_dir(&dir, 1).unwrap(); + let res = Extent::create_copy_dir(&dir, 1); assert!(res.is_err()); Ok(()) } @@ -2833,24 +2913,25 @@ mod test { region.extend(3).await?; // Close extent 1 - let ext_one = &mut region.extents[1]; - let (gen, flush, dirty) = ext_one.close().await?; + let (gen, flush, dirty) = region.close_extent(1).await.unwrap(); // Verify inner is gone, and we returned the expected gen, flush // and dirty values for a new unwritten extent. - assert!(ext_one.inner.is_none()); + assert!(matches!( + *region.extents[1].lock().await, + ExtentState::Closed + )); assert_eq!(gen, 0); assert_eq!(flush, 0); assert!(!dirty); // Make copy directory for this extent - let cp = ext_one.create_copy_dir(&dir)?; + let cp = Extent::create_copy_dir(&dir, 1)?; // Reopen extent 1 region.reopen_extent(1).await?; // Verify extent one is valid - let ext_one = &mut region.extents[1]; - assert!(ext_one.inner.is_some()); + let ext_one = region.get_opened_extent(1).await; // Make sure the eid matches assert_eq!(ext_one.number, 1); @@ -2872,19 +2953,23 @@ mod test { region.extend(3).await?; // Close extent 1 - let ext_one = &mut region.extents[1]; - ext_one.close().await?; - assert!(ext_one.inner.is_none()); + region.close_extent(1).await.unwrap(); + assert!(matches!( + *region.extents[1].lock().await, + ExtentState::Closed + )); // Make copy directory for this extent - let cp = ext_one.create_copy_dir(&dir)?; + let cp = Extent::create_copy_dir(&dir, 1)?; // Reopen extent 1 region.reopen_extent(1).await?; // Verify extent one is valid - let ext_one = &mut region.extents[1]; - assert!(ext_one.inner.is_some()); + let ext_one = region.get_opened_extent(1).await; + + // Make sure the eid matches + assert_eq!(ext_one.number, 1); // Make sure copy directory was removed assert!(!Path::new(&cp).exists()); @@ -2903,12 +2988,14 @@ mod test { region.extend(3).await?; // Close extent 1 - let ext_one = &mut region.extents[1]; - ext_one.close().await?; - assert!(ext_one.inner.is_none()); + region.close_extent(1).await.unwrap(); + assert!(matches!( + *region.extents[1].lock().await, + ExtentState::Closed + )); // Make copy directory for this extent - let cp = ext_one.create_copy_dir(&dir)?; + let cp = Extent::create_copy_dir(&dir, 1)?; // Step through the replacement dir, but don't do any work. let rd = replace_dir(&dir, 1); @@ -2922,8 +3009,7 @@ mod test { region.reopen_extent(1).await?; // Verify extent one is valid - let ext_one = &mut region.extents[1]; - assert!(ext_one.inner.is_some()); + let _ext_one = region.get_opened_extent(1).await; // Make sure all repair directories are gone assert!(!Path::new(&cp).exists()); @@ -2945,12 +3031,14 @@ mod test { region.extend(3).await?; // Close extent 1 - let ext_one = &mut region.extents[1]; - ext_one.close().await?; - assert!(ext_one.inner.is_none()); + region.close_extent(1).await.unwrap(); + assert!(matches!( + *region.extents[1].lock().await, + ExtentState::Closed + )); // Make copy directory for this extent - let cp = ext_one.create_copy_dir(&dir)?; + let cp = Extent::create_copy_dir(&dir, 1)?; // We are simulating the copy of files from the "source" repair // extent by copying the files from extent zero into the copy @@ -2982,8 +3070,7 @@ mod test { // Reopen extent 1 region.reopen_extent(1).await?; - let ext_one = &mut region.extents[1]; - assert!(ext_one.inner.is_some()); + let _ext_one = region.get_opened_extent(1).await; // Make sure all repair directories are gone assert!(!Path::new(&cp).exists()); @@ -3011,12 +3098,14 @@ mod test { region.extend(3).await?; // Close extent 1 - let ext_one = &mut region.extents[1]; - ext_one.close().await?; - assert!(ext_one.inner.is_none()); + region.close_extent(1).await.unwrap(); + assert!(matches!( + *region.extents[1].lock().await, + ExtentState::Closed + )); // Make copy directory for this extent - let cp = ext_one.create_copy_dir(&dir)?; + let cp = Extent::create_copy_dir(&dir, 1)?; // We are simulating the copy of files from the "source" repair // extent by copying the files from extent zero into the copy @@ -3063,8 +3152,7 @@ mod test { dest_path.set_extension("db-wal"); assert!(!Path::new(&dest_path).exists()); - let ext_one = &mut region.extents[1]; - assert!(ext_one.inner.is_some()); + let _ext_one = region.get_opened_extent(1).await; // Make sure all repair directories are gone assert!(!Path::new(&cp).exists()); @@ -3092,8 +3180,8 @@ mod test { region.extend(3).await?; // Make copy directory for this extent - let ext_one = &mut region.extents[1]; - let cp = ext_one.create_copy_dir(&dir)?; + let _ext_one = region.get_opened_extent(1).await; + let cp = Extent::create_copy_dir(&dir, 1)?; // We are simulating the copy of files from the "source" repair // extent by copying the files from extent zero into the copy @@ -3114,13 +3202,12 @@ mod test { drop(region); // Open up the region read_only now. - let mut region = + let region = Region::open(&dir, new_region_options(), false, true, &csl()) .await?; // Verify extent 1 has opened again. - let ext_one = &mut region.extents[1]; - assert!(ext_one.inner.is_some()); + let _ext_one = region.get_opened_extent(1).await; // Make sure repair directory is still present assert!(Path::new(&rd).exists()); @@ -3233,31 +3320,34 @@ mod test { region.extend(5).await?; // Close extent 1 - let ext_one = &mut region.extents[1]; - ext_one.close().await?; - assert!(ext_one.inner.is_none()); + region.close_extent(1).await.unwrap(); + assert!(matches!( + *region.extents[1].lock().await, + ExtentState::Closed + )); // Close extent 4 - let ext_four = &mut region.extents[4]; - ext_four.close().await?; - assert!(ext_four.inner.is_none()); + region.close_extent(4).await.unwrap(); + assert!(matches!( + *region.extents[4].lock().await, + ExtentState::Closed + )); // Reopen all extents region.reopen_all_extents().await?; // Verify extent one is valid - let ext_one = &mut region.extents[1]; - assert!(ext_one.inner.is_some()); + let ext_one = region.get_opened_extent(1).await; // Make sure the eid matches assert_eq!(ext_one.number, 1); // Verify extent four is valid - let ext_four = &mut region.extents[4]; - assert!(ext_four.inner.is_some()); + let ext_four = region.get_opened_extent(4).await; // Make sure the eid matches assert_eq!(ext_four.number, 4); + Ok(()) } @@ -3272,7 +3362,8 @@ mod test { async fn new_existing_region() -> Result<()> { let dir = tempdir()?; let _ = Region::create(&dir, new_region_options(), csl()).await; - let _ = Region::open(&dir, new_region_options(), false, false, &csl()); + let _ = Region::open(&dir, new_region_options(), false, false, &csl()) + .await; Ok(()) } @@ -3560,8 +3651,8 @@ mod test { Region::create(&dir, new_region_options(), csl()).await?; region.extend(1).await?; - let ext = ®ion.extents[0]; - let mut inner = ext.inner().await; + let ext = region.get_opened_extent(0).await; + let mut inner = ext.inner.lock().await; // Encryption context for blocks 0 and 1 should start blank @@ -3714,8 +3805,8 @@ mod test { Region::create(&dir, new_region_options(), csl()).await?; region.extend(1).await?; - let ext = ®ion.extents[0]; - let mut inner = ext.inner().await; + let ext = region.get_opened_extent(0).await; + let mut inner = ext.inner.lock().await; assert!(inner.get_block_contexts(0, 1)?[0].is_empty()); @@ -3755,8 +3846,8 @@ mod test { Region::create(&dir, new_region_options(), csl()).await?; region.extend(1).await?; - let ext = ®ion.extents[0]; - let mut inner = ext.inner().await; + let ext = region.get_opened_extent(0).await; + let mut inner = ext.inner.lock().await; // Encryption context for blocks 0 and 1 should start blank @@ -4049,8 +4140,9 @@ mod test { // A write of some sort only wrote a block context row { - let ext = ®ion.extents[0]; - let mut inner = ext.inner().await; + let ext = region.get_opened_extent(0).await; + let mut inner = ext.inner.lock().await; + inner.set_block_contexts(&[&DownstairsBlockContext { block_context: BlockContext { encryption_context: None, @@ -4061,17 +4153,17 @@ mod test { }])?; } - // This should clear out the invalid context - for extent in &mut region.extents { - extent.close().await?; + // This should clear out the invalid contexts + for eid in 0..region.extents.len() { + region.close_extent(eid).await.unwrap(); } + region.reopen_all_extents().await?; // Verify no block context rows exist - { - let ext = ®ion.extents[0]; - let inner = ext.inner().await; + let ext = region.get_opened_extent(0).await; + let inner = ext.inner.lock().await; assert!(inner.get_block_contexts(0, 1)?[0].is_empty()); } @@ -4125,7 +4217,7 @@ mod test { Region::create(&dir, new_region_options(), csl()).await?; region.extend(1).await?; - let ext = ®ion.extents[0]; + let ext = region.get_opened_extent(0).await; // Write a block, but don't flush. { @@ -4235,12 +4327,12 @@ mod test { Region::create(&dir, new_region_options(), csl()).await?; region.extend(1).await?; - let ext = ®ion.extents[0]; + let ext = region.get_opened_extent(0).await; // Partial write, the data never hits disk, but there's a context // in the DB { - let mut inner = ext.inner().await; + let mut inner = ext.inner.lock().await; inner.set_block_contexts(&[&DownstairsBlockContext { block_context: BlockContext { encryption_context: None, @@ -4356,9 +4448,7 @@ mod test { // Verify the dirty bit is now set. // We know our EID, so we can shortcut to getting the actual extent. - let inner = ®ion.extents[eid as usize].inner.as_ref().unwrap(); - let dirty = inner.lock().await.dirty().unwrap(); - assert!(dirty); + assert!(region.get_opened_extent(eid as usize).await.dirty().await); // Now read back that block, make sure it is updated. let responses = region @@ -4475,17 +4565,13 @@ mod test { region.region_write(&writes, 0, true).await?; // Verify the dirty bit is now set. - let inner = ®ion.extents[eid as usize].inner.as_ref().unwrap(); - let dirty = inner.lock().await.dirty().unwrap(); - assert!(dirty); + assert!(region.get_opened_extent(eid as usize).await.dirty().await); // Flush extent with eid, fn, gen, job_id. region.region_flush_extent(eid as usize, 1, 1, 1).await?; // Verify the dirty bit is no longer set. - let inner = ®ion.extents[eid as usize].inner.as_ref().unwrap(); - let dirty = inner.lock().await.dirty().unwrap(); - assert!(!dirty); + assert!(!region.get_opened_extent(eid as usize).await.dirty().await); // Create a new write IO with different data. let data = BytesMut::from(&[1u8; 512][..]); @@ -4509,9 +4595,7 @@ mod test { region.region_write(&writes, 1, true).await?; // Verify the dirty bit is not set. - let inner = ®ion.extents[eid as usize].inner.as_ref().unwrap(); - let dirty = inner.lock().await.dirty().unwrap(); - assert!(!dirty); + assert!(!region.get_opened_extent(eid as usize).await.dirty().await); // Read back our block, make sure it has the first write data let responses = region @@ -5105,22 +5189,16 @@ mod test { region.region_write(&writes, 0, true).await.unwrap(); // Verify the dirty bit is now set for both extents. - let inner = ®ion.extents[0].inner.as_ref().unwrap(); - assert!(inner.lock().await.dirty().unwrap()); - - let inner = ®ion.extents[1].inner.as_ref().unwrap(); - assert!(inner.lock().await.dirty().unwrap()); + assert!(region.get_opened_extent(0).await.dirty().await); + assert!(region.get_opened_extent(1).await.dirty().await); // Call flush, but limit the flush to extent 0 region.region_flush(1, 2, &None, 3, Some(0)).await.unwrap(); // Verify the dirty bit is no longer set for 0, but still set // for extent 1. - let inner = ®ion.extents[0].inner.as_ref().unwrap(); - assert!(!inner.lock().await.dirty().unwrap()); - - let inner = ®ion.extents[1].inner.as_ref().unwrap(); - assert!(inner.lock().await.dirty().unwrap()); + assert!(!region.get_opened_extent(0).await.dirty().await); + assert!(region.get_opened_extent(1).await.dirty().await); } #[tokio::test] @@ -5142,29 +5220,22 @@ mod test { region.region_write(&writes, 2, true).await.unwrap(); // Verify the dirty bit is now set for both extents. - let inner = ®ion.extents[1].inner.as_ref().unwrap(); - assert!(inner.lock().await.dirty().unwrap()); - - let inner = ®ion.extents[2].inner.as_ref().unwrap(); - assert!(inner.lock().await.dirty().unwrap()); + assert!(region.get_opened_extent(1).await.dirty().await); + assert!(region.get_opened_extent(2).await.dirty().await); // Call flush, but limit the flush to extents < 2 region.region_flush(1, 2, &None, 3, Some(1)).await.unwrap(); // Verify the dirty bit is no longer set for 1, but still set // for extent 2. - let inner = ®ion.extents[1].inner.as_ref().unwrap(); - assert!(!inner.lock().await.dirty().unwrap()); - - let inner = ®ion.extents[2].inner.as_ref().unwrap(); - assert!(inner.lock().await.dirty().unwrap()); + assert!(!region.get_opened_extent(1).await.dirty().await); + assert!(region.get_opened_extent(2).await.dirty().await); // Now flush with no restrictions. region.region_flush(1, 2, &None, 3, None).await.unwrap(); // Extent 2 should no longer be dirty - let inner = ®ion.extents[2].inner.as_ref().unwrap(); - assert!(!inner.lock().await.dirty().unwrap()); + assert!(!region.get_opened_extent(2).await.dirty().await); } #[tokio::test] @@ -5187,8 +5258,7 @@ mod test { // Verify the dirty bit is now set for all extents. for ext in 0..10 { - let inner = ®ion.extents[ext].inner.as_ref().unwrap(); - assert!(inner.lock().await.dirty().unwrap()); + assert!(region.get_opened_extent(ext).await.dirty().await); } // Walk up the extent_limit, verify at each flush extents are @@ -5202,14 +5272,12 @@ mod test { // This ext should no longer be dirty. println!("extent {} should not be dirty now", ext); - let inner = ®ion.extents[ext].inner.as_ref().unwrap(); - assert!(!inner.lock().await.dirty().unwrap()); + assert!(!region.get_opened_extent(ext).await.dirty().await); // Any extent above the current point should still be dirty. for d_ext in ext + 1..10 { println!("verify {} still dirty", d_ext); - let inner = ®ion.extents[d_ext].inner.as_ref().unwrap(); - assert!(inner.lock().await.dirty().unwrap()); + assert!(region.get_opened_extent(d_ext).await.dirty().await); } } } @@ -5269,8 +5337,8 @@ mod test { .unwrap(); // Close extent 0 - let ext_zero = &mut region.extents[eid as usize]; - let (gen, flush, dirty) = ext_zero.close().await.unwrap(); + let (gen, flush, dirty) = + region.close_extent(eid as usize).await.unwrap(); // Verify inner is gone, and we returned the expected gen, flush // and dirty values for the write that should be flushed now. @@ -5314,8 +5382,8 @@ mod test { region.region_write(&writes, 0, true).await.unwrap(); // Close extent 0 without a flush - let ext_zero = &mut region.extents[eid as usize]; - let (gen, flush, dirty) = ext_zero.close().await.unwrap(); + let (gen, flush, dirty) = + region.close_extent(eid as usize).await.unwrap(); // Because we did not flush yet, this extent should still have // the values for an unwritten extent, except for the dirty bit. @@ -5327,8 +5395,9 @@ mod test { // the same as the previous check (testing here that dirty remains // dirty). region.reopen_extent(eid as usize).await.unwrap(); - let ext_zero = &mut region.extents[eid as usize]; - let (gen, flush, dirty) = ext_zero.close().await.unwrap(); + + let (gen, flush, dirty) = + region.close_extent(eid as usize).await.unwrap(); // Verify everything is the same, and dirty is still set. assert_eq!(gen, 0); @@ -5341,8 +5410,9 @@ mod test { .region_flush_extent(eid as usize, 4, 9, 1) .await .unwrap(); - let ext_zero = &mut region.extents[eid as usize]; - let (gen, flush, dirty) = ext_zero.close().await.unwrap(); + + let (gen, flush, dirty) = + region.close_extent(eid as usize).await.unwrap(); // Verify after flush that g,f are updated, and that dirty // is no longer set. @@ -5408,8 +5478,9 @@ mod test { // We are gonna compare against the last write iteration let last_writes = writes.last().unwrap(); - let ext = ®ion.extents[0]; - let inner = ext.inner().await; + let ext = region.get_opened_extent(0).await; + let inner = ext.inner.lock().await; + for (i, range) in ranges.iter().enumerate() { // Get the contexts for the range let ctxts = inner @@ -5487,8 +5558,9 @@ mod test { // compare against the last write iteration let last_writes = writes.last().unwrap(); - let ext = ®ion.extents[0]; - let inner = ext.inner().await; + let ext = region.get_opened_extent(0).await; + let inner = ext.inner.lock().await; + // Get the contexts for the range let ctxts = inner.get_block_contexts(0, EXTENT_SIZE)?; diff --git a/downstairs/src/repair.rs b/downstairs/src/repair.rs index b53851f0c..33512a8bb 100644 --- a/downstairs/src/repair.rs +++ b/downstairs/src/repair.rs @@ -79,14 +79,7 @@ pub async fn repair_main( .start(); let local_addr = server.local_addr(); - tokio::spawn(async move { - /* - * Wait for the server to stop. Note that there's not any code to - * shut down this server, so we should never get past this - * point. - */ - server.await - }); + tokio::spawn(server); Ok(local_addr) } @@ -339,8 +332,7 @@ mod test { Region::create(&dir, new_region_options(), csl()).await?; region.extend(3).await?; - let ext_one = &mut region.extents[1]; - ext_one.close().await?; + region.close_extent(1).await.unwrap(); // Determine the directory and name for expected extent files. let extent_dir = extent_dir(&dir, 1); diff --git a/integration_tests/src/lib.rs b/integration_tests/src/lib.rs index 4ee5da42c..24776fbed 100644 --- a/integration_tests/src/lib.rs +++ b/integration_tests/src/lib.rs @@ -290,6 +290,10 @@ mod test { pub async fn downstairs2_address(&self) -> SocketAddr { self.downstairs2.address().await } + + pub async fn downstairs3_address(&self) -> SocketAddr { + self.downstairs3.address().await + } } #[tokio::test] @@ -2458,6 +2462,135 @@ mod test { ) .await .unwrap_err(); + + Ok(()) + } + + #[tokio::test] + async fn integration_test_volume_replace_downstairs_then_takeover( + ) -> Result<()> { + // Replace a downstairs with a new one, then an Upstairs with a newer + // generation number activates. + const BLOCK_SIZE: usize = 512; + + // boot three downstairs, write some data to them + let test_downstairs_set = TestDownstairsSet::big(false).await?; + + let mut volume = Volume::new(BLOCK_SIZE as u64); + volume + .add_subvolume_create_guest( + test_downstairs_set.opts(), + volume::RegionExtentInfo { + block_size: BLOCK_SIZE as u64, + blocks_per_extent: test_downstairs_set.blocks_per_extent(), + extent_count: test_downstairs_set.extent_count(), + }, + 1, + None, + ) + .await?; + + volume.activate().await?; + + let random_buffer = { + let mut random_buffer = + vec![0u8; volume.total_size().await? as usize]; + rand::thread_rng().fill(&mut random_buffer[..]); + random_buffer + }; + + volume + .write( + Block::new(0, BLOCK_SIZE.trailing_zeros()), + Bytes::from(random_buffer.clone()), + ) + .await?; + + // Create a new downstairs, then replace one of our current + // downstairs with that new one. + let new_downstairs = test_downstairs_set.new_downstairs().await?; + + let res = volume + .replace_downstairs( + test_downstairs_set.opts().id, + test_downstairs_set.downstairs1_address().await, + new_downstairs.address().await, + ) + .await + .unwrap(); + + assert_eq!(res, ReplaceResult::Started); + + loop { + tokio::time::sleep(tokio::time::Duration::from_secs(1)).await; + match volume + .replace_downstairs( + test_downstairs_set.opts().id, + test_downstairs_set.downstairs1_address().await, + new_downstairs.address().await, + ) + .await + .unwrap() + { + ReplaceResult::StartedAlready => { + eprintln!( + "Waited for some repair work, proceeding with test" + ); + break; + } + ReplaceResult::CompletedAlready => { + // This test is invalid if the repair completed already + panic!("Downstairs replacement completed"); + } + x => { + panic!("Bad result from replace_downstairs: {:?}", x); + } + } + } + + // A new Upstairs arrives, with a newer gen number, and the updated + // target list + let mut opts = test_downstairs_set.opts(); + opts.target = vec![ + new_downstairs.address().await, + test_downstairs_set.downstairs2_address().await, + test_downstairs_set.downstairs3_address().await, + ]; + + let mut new_volume = Volume::new(BLOCK_SIZE as u64); + new_volume + .add_subvolume_create_guest( + opts, + volume::RegionExtentInfo { + block_size: BLOCK_SIZE as u64, + blocks_per_extent: test_downstairs_set.blocks_per_extent(), + extent_count: test_downstairs_set.extent_count(), + }, + 2, + None, + ) + .await?; + + new_volume.activate().await?; + + while !new_volume.query_is_active().await? { + // new_volume will repair before activating, so this waits for that + println!("Waiting for new_volume activate"); + tokio::time::sleep(tokio::time::Duration::from_secs(1)).await; + } + + // Old volume is no longer active + assert!(!volume.query_is_active().await?); + + // Read back what we wrote. + let buffer = Buffer::new(new_volume.total_size().await? as usize); + new_volume + .read(Block::new(0, BLOCK_SIZE.trailing_zeros()), buffer.clone()) + .await?; + + let buffer_vec = buffer.as_vec().await; + assert_eq!(buffer_vec[BLOCK_SIZE..], random_buffer[BLOCK_SIZE..]); + Ok(()) } @@ -2693,7 +2826,6 @@ mod test { async fn integration_test_guest_replace_many_downstairs() -> Result<()> { // Test using the guest layer to verify we can replace one // downstairs, but not another while the replace is active. - const BLOCK_SIZE: usize = 512; // Spin off three downstairs, build our Crucible struct. let tds = TestDownstairsSet::small(false).await?; @@ -3336,8 +3468,6 @@ mod test { #[tokio::test] async fn test_pantry_import_from_url_ovmf_bad_digest() { - const BLOCK_SIZE: usize = 512; - // Spin off three downstairs, build our Crucible struct. let tds = TestDownstairsSet::big(false).await.unwrap(); @@ -3491,8 +3621,6 @@ mod test { #[tokio::test] async fn test_pantry_snapshot() { - const BLOCK_SIZE: usize = 512; - // Spin off three downstairs, build our Crucible struct. let tds = TestDownstairsSet::small(false).await.unwrap(); @@ -3808,8 +3936,6 @@ mod test { #[tokio::test] async fn test_pantry_bulk_read() { - const BLOCK_SIZE: usize = 512; - // Spin off three downstairs, build our Crucible struct. let tds = TestDownstairsSet::small(false).await.unwrap(); @@ -3892,8 +4018,6 @@ mod test { #[tokio::test] async fn test_pantry_bulk_read_max_chunk_size() { - const BLOCK_SIZE: usize = 512; - // Spin off three downstairs, build our Crucible struct. let tds = TestDownstairsSet::big(false).await.unwrap(); @@ -4208,8 +4332,6 @@ mod test { // Test validating a non-block size amount fails #[tokio::test] async fn test_pantry_validate_fail() { - const BLOCK_SIZE: usize = 512; - // Spin off three downstairs, build our Crucible struct. let tds = TestDownstairsSet::small(false).await.unwrap(); diff --git a/pantry/src/server.rs b/pantry/src/server.rs index 95f9e8abe..8d53fc6a2 100644 --- a/pantry/src/server.rs +++ b/pantry/src/server.rs @@ -375,7 +375,7 @@ pub async fn run_server( let local_addr = server.local_addr(); info!(log, "listen IP: {:?}", local_addr); - let join_handle = tokio::spawn(async move { server.await }); + let join_handle = tokio::spawn(server); Ok((local_addr, join_handle)) } diff --git a/rust-toolchain.toml b/rust-toolchain.toml index fe267488f..eb8b43af5 100644 --- a/rust-toolchain.toml +++ b/rust-toolchain.toml @@ -1,3 +1,3 @@ [toolchain] -channel = "1.66" +channel = "1.70" profile = "default" diff --git a/upstairs/src/lib.rs b/upstairs/src/lib.rs index f8b4f970f..4f8a2bbf9 100644 --- a/upstairs/src/lib.rs +++ b/upstairs/src/lib.rs @@ -8520,6 +8520,7 @@ async fn test_buffer_len_after_clone() { let data = Buffer::from_slice(&[0x99; READ_SIZE]); assert_eq!(data.len(), READ_SIZE); + #[allow(clippy::redundant_clone)] let new_buffer = data.clone(); assert_eq!(new_buffer.len(), READ_SIZE); assert_eq!(data.len(), READ_SIZE); diff --git a/upstairs/src/pseudo_file.rs b/upstairs/src/pseudo_file.rs index ce951261b..ce91e0b35 100644 --- a/upstairs/src/pseudo_file.rs +++ b/upstairs/src/pseudo_file.rs @@ -264,7 +264,7 @@ impl Seek for CruciblePseudoFile { } fn stream_position(&mut self) -> IOResult { - self.seek(SeekFrom::Current(0)) + Ok(self.offset) } }