From 1df165d514ae9ef8b278ea340582af70e7ca8bdf Mon Sep 17 00:00:00 2001 From: Peter Dodd Date: Thu, 8 Oct 2026 15:58:44 +0100 Subject: [PATCH] fix(volume): keep the TTL clock across a vacuum commit instead of rescanning (#11630) * fix(volume): keep the TTL clock across a vacuum commit instead of rescanning CommitCompact reloads the swapped files while holding dataFileAccessLock, and for a vacuumed TTL volume that reload re-derived lastModifiedTsSeconds by reading every live needle's append timestamp from the .dat: two random reads per needle, with every read of the volume blocked behind them. The in-memory clock is already current at that point. Every write since the volume loaded moved it, and makeupDiff only replays writes that went through that path. Carry it across the reload instead. This also stops an over-budget scan from falling back to the new .dat's mtime and restarting an expiring volume's TTL at the commit. * volume: carry the append watermark as the TTL clock across a vacuum commit The running append watermark is the clock the reload's recovery scan recomputes, so the commit can keep it directly. Client-supplied needle modified times can run ahead of or behind the append time; keeping lastModifiedTsSeconds itself would let a forged or stale timestamp move expiry through a vacuum, where the scan it replaces used server-side append timestamps. * storage: test that a vacuum commit keeps the append clock A write's client supplied modified time can lie ahead of or behind its append time; the commit must land the TTL clock on the append watermark, the same value the recovery scan would have recomputed. * volume: carry the append watermark as the TTL clock across a vacuum commit Mirrors the Go volume server: the running append watermark is the clock the reload's recovery scan recomputes, so the commit keeps it instead of rescanning live needles under the write lock. * volume: commit carries the last-write append time, not the latest append lastAppendAtNs counts tombstone appends and is reseeded from the .dat tail at every load, so it can sit ahead of the last write -- a delete freshens the commit clock -- or behind it: a restarted vacuumed volume's tail needle is not its newest write, and the commit would move the TTL clock backward into premature expiry. Track lastWriteAppendAtNs instead, bumped only on needle appends and seeded by the recovery scan, so the commit lands the clock on the same live-write maximum the rescan would have recomputed. * volume: rescan at commit when the newest write was deleted lastWriteAppendAtNs can hold a write the index no longer holds, so carrying it extends the TTL clock past what recovery over the compacted index would compute. Remember the key behind the watermark so its tombstone or index rollback can send the reload back through recoverLastModifiedTs, landing on the newest surviving write. * volume: a tombstone retires the write rows beneath it in the last-write scan A needle deleted after the compaction copy leaves its write row followed by a tombstone in the committed index. The reverse scan skipped the tombstone row but then counted the dead write, reseeding the watermark and flag as if it were alive. Track keys whose latest row is a tombstone so their earlier write rows stop counting, and exercise the delete-inside-the-commit-window ordering in the tests. --------- Co-authored-by: Chris Lu Co-authored-by: Chris Lu --- seaweed-volume/src/storage/volume.rs | 245 +++++++++++++++++++++++-- weed/storage/volume.go | 8 +- weed/storage/volume_checking.go | 52 ++++-- weed/storage/volume_loading.go | 10 +- weed/storage/volume_ttl_expiry_test.go | 148 ++++++++++++++- weed/storage/volume_vacuum.go | 12 +- weed/storage/volume_write.go | 20 ++ 7 files changed, 450 insertions(+), 45 deletions(-) diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index d97e52300..4477a8ad7 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -1182,6 +1182,10 @@ pub struct Volume { last_modified_ts_seconds: u64, last_append_at_ns: u64, + last_write_append_at_ns: u64, // AppendAtNs of the newest write; tombstones don't move it + last_write_needle_key: NeedleId, // the write behind the watermark + last_write_deleted: bool, // that write was deleted, so recovery has to rescan + keep_last_modified_ts_on_load: bool, pub last_disk_check_ns: Arc, // for phantom volume detection cache last_compact_index_offset: u64, @@ -1281,6 +1285,10 @@ impl Volume { location_disk_space_low: Arc::new(AtomicBool::new(false)), last_modified_ts_seconds: 0, last_append_at_ns: 0, + last_write_append_at_ns: 0, + last_write_needle_key: NeedleId(0), + last_write_deleted: false, + keep_last_modified_ts_on_load: false, last_disk_check_ns: Arc::new(std::sync::atomic::AtomicI64::new(0)), last_compact_index_offset: 0, last_compact_revision: 0, @@ -1327,6 +1335,10 @@ impl Volume { location_disk_space_low: Arc::new(AtomicBool::new(false)), last_modified_ts_seconds: 0, last_append_at_ns: 0, + last_write_append_at_ns: 0, + last_write_needle_key: NeedleId(0), + last_write_deleted: false, + keep_last_modified_ts_on_load: false, last_disk_check_ns: Arc::new(std::sync::atomic::AtomicI64::new(0)), last_compact_index_offset: 0, last_compact_revision: 0, @@ -1442,12 +1454,14 @@ impl Volume { Err(e) => return Err(e.into()), } - self.last_modified_ts_seconds = metadata - .modified() - .unwrap_or(SystemTime::UNIX_EPOCH) - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_secs(); + if !self.keep_last_modified_ts_on_load { + self.last_modified_ts_seconds = metadata + .modified() + .unwrap_or(SystemTime::UNIX_EPOCH) + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_secs(); + } if metadata.len() >= SUPER_BLOCK_SIZE as u64 { already_has_super_block = true; @@ -1549,7 +1563,9 @@ impl Volume { "volumeDataIntegrityChecking failed" ); } - self.recover_last_modified_ts(); + if !self.keep_last_modified_ts_on_load { + self.recover_last_modified_ts(); + } // Structural check: no .idx entry may reference bytes past the // end of .dat. The needle map's load walk above already @@ -2339,6 +2355,13 @@ impl Volume { .collect(); } self.last_append_at_ns = last_append_at_ns; + for ((n, _), r) in run.iter().zip(&staged) { + if matches!(r, Ok(Some(_))) && n.append_at_ns > self.last_write_append_at_ns { + self.last_write_append_at_ns = n.append_at_ns; + self.last_write_needle_key = n.id; + self.last_write_deleted = false; + } + } // A durable entry that fails to publish stops the volume taking // writes, so the entries after it are refused the way a lone @@ -2521,6 +2544,9 @@ impl Volume { } self.last_append_at_ns = n.append_at_ns; + self.last_write_append_at_ns = n.append_at_ns; + self.last_write_needle_key = n.id; + self.last_write_deleted = false; self.publish_write(n, offset, fsync)?; @@ -2839,6 +2865,11 @@ impl Volume { if let Some(nm) = &mut self.nm { nm.delete(n.id, Offset::from_actual_offset(offset as i64))?; } + if n.id == self.last_write_needle_key { + self.last_write_append_at_ns = 0; + self.last_write_needle_key = NeedleId(0); + self.last_write_deleted = true; + } let checkpoint_ok = self.maybe_checkpoint_index(false); // Clear the EIO streak after a successful delete (tombstone append + @@ -3106,8 +3137,15 @@ impl Volume { return; } match self.find_last_write_append_at_ns() { - Ok(0) => {} - Ok(append_at_ns) => self.last_modified_ts_seconds = append_at_ns / 1_000_000_000, + Ok((0, _)) => {} + Ok((append_at_ns, key)) => { + self.last_modified_ts_seconds = append_at_ns / 1_000_000_000; + if append_at_ns > self.last_write_append_at_ns { + self.last_write_append_at_ns = append_at_ns; + self.last_write_needle_key = key; + self.last_write_deleted = false; + } + } Err(e) => warn!( volume_id = self.id.0, error = %e, @@ -3117,7 +3155,8 @@ impl Volume { } /// Scan the .idx backwards for the newest write — an entry that is not a - /// deletion tombstone — and return that needle's append timestamp. The .idx + /// deletion tombstone — and return that needle's append timestamp and key. + /// The .idx /// and the .dat share an order, so an append-ordered volume answers with the /// first write the scan reaches. Vacuum rewrites both in key order, which /// tracks write order only because the master issues keys increasing: an @@ -3126,19 +3165,21 @@ impl Volume { /// nothing but tombstones, when a vacuumed volume holds more needles than /// the scan budget, or for a volume older than version 3, whose needles /// carry no append timestamp. Mirrors Go's findLastWriteAppendAtNs. - fn find_last_write_append_at_ns(&self) -> Result { + fn find_last_write_append_at_ns(&self) -> Result<(u64, NeedleId), VolumeError> { let version = self.version(); if version != VERSION_3 { - return Ok(0); + return Ok((0, NeedleId(0))); } let idx_path = self.file_name(".idx"); let idx_size = fs::metadata(&idx_path).map(|m| m.len()).unwrap_or(0) as i64; if idx_size == 0 || idx_size % NEEDLE_MAP_ENTRY_SIZE as i64 != 0 { - return Ok(0); + return Ok((0, NeedleId(0))); } let scan_every_write = self.super_block.compaction_revision > 0; let mut entry_budget = Self::VACUUMED_LAST_WRITE_SCAN_ENTRIES; let mut last_write_append_at_ns = 0u64; + let mut last_write_key = NeedleId(0); + let mut dead = HashSet::new(); let mut idx_file = File::open(&idx_path)?; let mut block = vec![0u8; NEEDLE_MAP_ENTRY_SIZE * idx::ROWS_TO_READ]; let mut end = idx_size; @@ -3149,7 +3190,13 @@ impl Volume { idx_file.read_exact(entries)?; for entry in entries.as_chunks::().0.iter().rev() { let (key, offset, size) = idx_entry_from_bytes(entry); + // The first row a key presents is its latest state: a tombstone + // there retires the write rows beneath it. + if dead.contains(&key) { + continue; + } if offset.is_zero() || size.is_deleted() { + dead.insert(key); continue; } let Some(needle_offset) = @@ -3157,10 +3204,13 @@ impl Volume { else { continue; }; - last_write_append_at_ns = last_write_append_at_ns - .max(self.read_needle_append_at_ns(needle_offset, size)?); + let append_at_ns = self.read_needle_append_at_ns(needle_offset, size)?; + if append_at_ns > last_write_append_at_ns { + last_write_append_at_ns = append_at_ns; + last_write_key = key; + } if !scan_every_write { - return Ok(last_write_append_at_ns); + return Ok((last_write_append_at_ns, last_write_key)); } entry_budget -= 1; if entry_budget == 0 { @@ -3169,12 +3219,12 @@ impl Volume { budget = Self::VACUUMED_LAST_WRITE_SCAN_ENTRIES, "too many needles to scan for the last write, keeping the .dat mtime" ); - return Ok(0); + return Ok((0, NeedleId(0))); } } end = start; } - Ok(last_write_append_at_ns) + Ok((last_write_append_at_ns, last_write_key)) } /// The .dat offset holding the needle an .idx entry describes, or None when @@ -4294,6 +4344,9 @@ impl Volume { // Update lastAppendAtNs (matches Go L352: v.lastAppendAtNs = appendAtNs) self.last_append_at_ns = append_at_ns; + self.last_write_append_at_ns = append_at_ns; + self.last_write_needle_key = needle_id; + self.last_write_deleted = false; // Update needle map index let offset = Offset::from_actual_offset(dat_size); @@ -4493,8 +4546,18 @@ impl Volume { self.write_compact_commit_marker()?; self.apply_compact_swap()?; - // Reload - self.load(true, false, 0, self.version())?; + // The write watermark already equals what recover_last_modified_ts + // would rescan, so keep the clock instead of paying for the scan + // under the lock. If its write was itself deleted, the reload + // recovers the newest surviving write instead. + if self.last_write_append_at_ns != 0 { + self.last_modified_ts_seconds = self.last_write_append_at_ns / 1_000_000_000; + } + self.keep_last_modified_ts_on_load = + self.last_modified_ts_seconds != 0 && !self.last_write_deleted; + let load_result = self.load(true, false, 0, self.version()); + self.keep_last_modified_ts_on_load = false; + load_result?; Ok(()) } @@ -6381,6 +6444,7 @@ mod tests { let prior = v.nm.as_ref().unwrap().get(NeedleId(1)).unwrap().unwrap(); let dat_len_before = dat_len(&v); let last_append_before = v.last_append_at_ns; + let last_write_append_before = v.last_write_append_at_ns; let last_modified_before = v.last_modified_ts_seconds; let file_count_before = v.file_count(); @@ -6402,6 +6466,7 @@ mod tests { ); assert_eq!(dat_len(&v), dat_len_before, "the run is off the .dat"); assert_eq!(v.last_append_at_ns, last_append_before); + assert_eq!(v.last_write_append_at_ns, last_write_append_before); assert_eq!(v.last_modified_ts_seconds, last_modified_before); let now = v.nm.as_ref().unwrap().get(NeedleId(1)).unwrap().unwrap(); assert_eq!((now.offset, now.size), (prior.offset, prior.size)); @@ -7037,6 +7102,146 @@ mod tests { ); } + // Covers the reload that ends a vacuum commit: the clock must keep the + // last write's append time rather than be re-derived, and the intervening + // delete's tombstone must not freshen it. + #[test] + fn test_ttl_clock_carried_across_vacuum_commit() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let ttl = crate::storage::needle::ttl::TTL::read("5m").unwrap(); + let last_write_ns = (SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_secs() + - 2 * 60 * 60) + * 1_000_000_000; + + let mut v = make_ttl_volume(dir, ttl); + let mut written = Vec::new(); + for i in 1..=3u64 { + let data = format!("data {}", i); + let mut n = Needle { + id: NeedleId(i), + cookie: Cookie(i as u32), + data: data.as_bytes().to_vec(), + data_size: data.len() as u32, + ..Needle::default() + }; + let (offset, _, _) = v.write_needle(&mut n, true, false).unwrap(); + written.push((offset, n.size)); + } + v.delete_needle(&mut Needle { + id: NeedleId(2), + cookie: Cookie(2), + ..Needle::default() + }) + .unwrap(); + v.sync_to_disk().unwrap(); + for (offset, size) in written { + backdate_append_at_ns(&v.dat_path(), offset, size, last_write_ns); + } + // Where a restart's recovery would have left the clock. The delete + // above pushed last_append_at_ns to ~now; the commit must not use it. + v.set_last_modified_ts_for_test(last_write_ns / 1_000_000_000); + v.last_write_append_at_ns = last_write_ns; + + v.compact_by_index(0, 0, |_| true).unwrap(); + v.commit_compact().unwrap(); + + assert_eq!(v.last_modified_ts(), last_write_ns / 1_000_000_000); + assert!( + v.is_expired(v.content_size(), 1024 * 1024), + "a TTL volume whose last write is 2h old must stay expired across a vacuum commit" + ); + } + + // A write's client supplied modified time can lie ahead of or behind when + // it was appended; the commit keeps the server-side write watermark the + // recovery scan would recompute. + #[test] + fn test_ttl_clock_at_commit_uses_append_time() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let ttl = crate::storage::needle::ttl::TTL::read("5m").unwrap(); + let future = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_secs() + + 24 * 60 * 60; + + let mut v = make_ttl_volume(dir, ttl); + let data = b"data".to_vec(); + let mut n = Needle { + id: NeedleId(1), + cookie: Cookie(1), + data, + data_size: 4, + last_modified: future, + ..Needle::default() + }; + v.write_needle(&mut n, true, false).unwrap(); + assert!(v.last_modified_ts() >= future); + let append_watermark_sec = v.last_write_append_at_ns / 1_000_000_000; + + v.compact_by_index(0, 0, |_| true).unwrap(); + v.commit_compact().unwrap(); + + assert_eq!(v.last_modified_ts(), append_watermark_sec); + } + + // A vacuum commit whose newest write was deleted first: the carried + // watermark belonged to that write, so the commit has to let the reload + // rescan and land on the newest surviving write rather than keep the + // volume alive on a deleted write's time. + #[test] + fn test_ttl_clock_at_commit_skips_deleted_write() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let ttl = crate::storage::needle::ttl::TTL::read("5m").unwrap(); + let now = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_secs(); + let old_write_ns = (now - 2 * 60 * 60) * 1_000_000_000; + let new_write_ns = (now - 60 * 60) * 1_000_000_000; + + let mut v = make_ttl_volume(dir, ttl); + for (id, append_at_ns) in [(1u64, old_write_ns), (2, new_write_ns)] { + let data = format!("data {}", id); + let mut n = Needle { + id: NeedleId(id), + cookie: Cookie(id as u32), + data: data.as_bytes().to_vec(), + data_size: data.len() as u32, + ..Needle::default() + }; + let (offset, _, _) = v.write_needle(&mut n, true, false).unwrap(); + v.sync_to_disk().unwrap(); + backdate_append_at_ns(&v.dat_path(), offset, n.size, append_at_ns); + } + // The delete lands inside the commit window: makeup_diff replays its + // tombstone into the new .idx behind the write row the copy carried. + v.compact_by_index(0, 0, |_| true).unwrap(); + v.delete_needle(&mut Needle { + id: NeedleId(2), + cookie: Cookie(2), + ..Needle::default() + }) + .unwrap(); + assert!( + v.last_write_deleted, + "deleting the newest write must mark the watermark dead" + ); + v.commit_compact().unwrap(); + + assert_eq!(v.last_modified_ts(), old_write_ns / 1_000_000_000); + assert!( + !v.last_write_deleted, + "the reload's rescan must reseed the watermark off the surviving write" + ); + } + // Guard the destroy time an EC volume is reclaimed on: it was recomputed as // now+TTL every time the .vif was written, so a read-only mark, a tier // upload or an EC encode handed an already expiring volume another full TTL. diff --git a/weed/storage/volume.go b/weed/storage/volume.go index 0528b5029..c51ceb271 100644 --- a/weed/storage/volume.go +++ b/weed/storage/volume.go @@ -54,8 +54,12 @@ type Volume struct { asyncRequestsChan chan *needle.AsyncRequest asyncWorkerClosed bool - lastModifiedTsSeconds uint64 // unix time in seconds - lastAppendAtNs uint64 // unix time in nanoseconds + lastModifiedTsSeconds uint64 // unix time in seconds + lastAppendAtNs uint64 // unix time in nanoseconds + lastWriteAppendAtNs uint64 // AppendAtNs of the newest write; tombstones don't move it + lastWriteNeedleKey types.NeedleId // the write behind the watermark + lastWriteDeleted bool // that write was deleted, so recovery has to rescan + keepLastModifiedTsOnLoad bool lastCompactIndexOffset uint64 lastCompactRevision uint16 diff --git a/weed/storage/volume_checking.go b/weed/storage/volume_checking.go index 7d007f927..52cff0e74 100644 --- a/weed/storage/volume_checking.go +++ b/weed/storage/volume_checking.go @@ -264,7 +264,7 @@ func (v *Volume) recoverLastModifiedTs(indexFile *os.File) { if err != nil || indexSize == 0 { return } - appendAtNs, err := findLastWriteAppendAtNs(v, indexFile, indexSize) + appendAtNs, key, err := findLastWriteAppendAtNs(v, indexFile, indexSize) if err != nil { glog.Warningf("volume %d recover last write from %s: %v", v.Id, indexFile.Name(), err) return @@ -273,6 +273,11 @@ func (v *Volume) recoverLastModifiedTs(indexFile *os.File) { return } v.lastModifiedTsSeconds = appendAtNs / uint64(time.Second) + if appendAtNs > v.lastWriteAppendAtNs { + v.lastWriteAppendAtNs = appendAtNs + v.lastWriteNeedleKey = key + v.lastWriteDeleted = false + } } // vacuumedLastWriteScanEntries bounds the work a vacuumed volume's recovery @@ -285,25 +290,27 @@ var vacuumedLastWriteScanEntries = 1 << 16 // findLastWriteAppendAtNs scans the .idx backwards for the newest write -- an // entry that is not a deletion tombstone -- and returns that needle's append -// timestamp. The .idx and the .dat share an order, so an append-ordered volume -// answers with the first write the scan reaches. Vacuum rewrites both in key -// order, which tracks write order only because the master issues keys -// increasing: an overwrite keeps its original, lower key, so a vacuumed volume -// has to take the maximum over every write it indexes. Returns 0 when the .idx -// holds nothing but tombstones, when a vacuumed volume holds more needles than -// the scan budget, or for a volume older than version 3, whose needles carry no -// append timestamp. -func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (uint64, error) { +// timestamp and key. The .idx and the .dat share an order, so an +// append-ordered volume answers with the first write the scan reaches. Vacuum +// rewrites both in key order, which tracks write order only because the master +// issues keys increasing: an overwrite keeps its original, lower key, so a +// vacuumed volume has to take the maximum over every write it indexes. Returns +// 0 when the .idx holds nothing but tombstones, when a vacuumed volume holds +// more needles than the scan budget, or for a volume older than version 3, +// whose needles carry no append timestamp. +func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (uint64, types.NeedleId, error) { version := v.Version() if version != needle.Version3 { - return 0, nil + return 0, 0, nil } scanEveryWrite := v.SuperBlock.CompactionRevision > 0 entryBudget := vacuumedLastWriteScanEntries if scanEveryWrite && !affordableVacuumedScan(indexFile, indexSize, v.Id, v.FileName(".dat")) { - return 0, nil + return 0, 0, nil } var lastWriteAppendAtNs uint64 + var lastWriteKey types.NeedleId + dead := make(map[types.NeedleId]struct{}) block := make([]byte, types.NeedleMapEntrySize*idx.RowsToRead) for end := indexSize; end > 0; { start := max(end-int64(len(block)), 0) @@ -313,11 +320,17 @@ func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (ui err = nil } if err != nil { - return 0, fmt.Errorf("read %s at %d: %v", indexFile.Name(), start, err) + return 0, 0, fmt.Errorf("read %s at %d: %v", indexFile.Name(), start, err) } for i := len(entries) - types.NeedleMapEntrySize; i >= 0; i -= types.NeedleMapEntrySize { key, offset, size := idx.IdxFileEntry(entries[i : i+types.NeedleMapEntrySize]) + // The first row a key presents is its latest state: a tombstone + // there retires the write rows beneath it. + if _, gone := dead[key]; gone { + continue + } if offset.IsZero() || size.IsDeleted() { + dead[key] = struct{}{} continue } needleOffset := findNeedleOffset(v.DataBackend, version, offset.ToActualOffset(), key, size) @@ -326,21 +339,24 @@ func findLastWriteAppendAtNs(v *Volume, indexFile *os.File, indexSize int64) (ui } appendAtNs, err := readNeedleAppendAtNs(v.DataBackend, needleOffset, size) if err != nil { - return 0, err + return 0, 0, err + } + if appendAtNs > lastWriteAppendAtNs { + lastWriteAppendAtNs = appendAtNs + lastWriteKey = key } - lastWriteAppendAtNs = max(lastWriteAppendAtNs, appendAtNs) if !scanEveryWrite { - return lastWriteAppendAtNs, nil + return lastWriteAppendAtNs, lastWriteKey, nil } if entryBudget--; entryBudget == 0 { glog.V(0).Infof("volume %d: more than %d needles to scan for its last write, keeping the %s mtime", v.Id, vacuumedLastWriteScanEntries, v.FileName(".dat")) - return 0, nil + return 0, 0, nil } } end = start } - return lastWriteAppendAtNs, nil + return lastWriteAppendAtNs, lastWriteKey, nil } // affordableVacuumedScan reports whether a vacuumed volume's recovery scan fits diff --git a/weed/storage/volume_loading.go b/weed/storage/volume_loading.go index 9a0518b85..dc4ffa365 100644 --- a/weed/storage/volume_loading.go +++ b/weed/storage/volume_loading.go @@ -172,7 +172,7 @@ func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind return fmt.Errorf("load remote file %v: %w", v.volumeInfo, err) } // Set lastModifiedTsSeconds from remote file to prevent premature expiry on startup - if len(v.volumeInfo.GetFiles()) > 0 { + if len(v.volumeInfo.GetFiles()) > 0 && !v.keepLastModifiedTsOnLoad { remoteFileModifiedTime := v.volumeInfo.GetFiles()[0].GetModifiedTime() if remoteFileModifiedTime > 0 { v.lastModifiedTsSeconds = remoteFileModifiedTime @@ -201,7 +201,9 @@ func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind if err != nil { return datFileLoadError(v.FileName(".dat"), err) } - v.lastModifiedTsSeconds = uint64(modifiedTime.Unix()) + if !v.keepLastModifiedTsOnLoad { + v.lastModifiedTsSeconds = uint64(modifiedTime.Unix()) + } if fileSize >= super_block.SuperBlockSize { alreadyHasSuperBlock = true } @@ -292,7 +294,9 @@ func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind v.noWriteOrDelete = true glog.V(0).Infof("volumeDataIntegrityChecking failed %v", err) } - v.recoverLastModifiedTs(indexFile) + if !v.keepLastModifiedTsOnLoad { + v.recoverLastModifiedTs(indexFile) + } } // The post-load structural check below uses the in-memory needle map diff --git a/weed/storage/volume_ttl_expiry_test.go b/weed/storage/volume_ttl_expiry_test.go index 3dd18dd5d..8ebffc1c8 100644 --- a/weed/storage/volume_ttl_expiry_test.go +++ b/weed/storage/volume_ttl_expiry_test.go @@ -250,7 +250,7 @@ func TestVolumeTtlClockSkipsUnaffordableScanWithoutDatReads(t *testing.T) { defer func(budget int) { vacuumedLastWriteScanEntries = budget }(vacuumedLastWriteScanEntries) vacuumedLastWriteScanEntries = 2 - appendAtNs, err := findLastWriteAppendAtNs(v, indexFile, indexStat.Size()) + appendAtNs, _, err := findLastWriteAppendAtNs(v, indexFile, indexStat.Size()) if err != nil { t.Fatalf("recover last write: %v", err) } @@ -262,6 +262,152 @@ func TestVolumeTtlClockSkipsUnaffordableScanWithoutDatReads(t *testing.T) { } } +// TestVolumeTtlClockCarriedAcrossVacuumCommit covers the reload that ends a +// vacuum commit: the clock must keep the last write's append time rather than +// be re-derived, and the intervening delete's tombstone must not freshen it. +func TestVolumeTtlClockCarriedAcrossVacuumCommit(t *testing.T) { + dir := t.TempDir() + ttl, err := needle.ReadTTL("5m") + if err != nil { + t.Fatalf("read ttl: %v", err) + } + + v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, ttl, 0, needle.GetCurrentVersion(), 0, 0) + if err != nil { + t.Fatalf("volume creation: %v", err) + } + defer v.Close() + + lastWriteNs := uint64(time.Now().Add(-2 * time.Hour).UnixNano()) + for i := 1; i <= 3; i++ { + n := newRandomNeedle(uint64(i)) + offset, _, _, err := v.writeNeedle2(n, true, false, false) + if err != nil { + t.Fatalf("write needle %d: %v", i, err) + } + backdateAppendAtNs(t, v, int64(offset), n.Size, lastWriteNs) + } + if _, err := v.doDeleteRequest(newEmptyNeedle(2)); err != nil { + t.Fatalf("delete needle 2: %v", err) + } + // Where a restart's recovery would have left the clock. The delete above + // pushed lastAppendAtNs to ~now; the commit must not consult it. + v.lastModifiedTsSeconds = lastWriteNs / uint64(time.Second) + v.lastWriteAppendAtNs = lastWriteNs + + defer func(budget int) { vacuumedLastWriteScanEntries = budget }(vacuumedLastWriteScanEntries) + vacuumedLastWriteScanEntries = 1 + + if err := v.CompactByIndex(nil); err != nil { + t.Fatalf("compact: %v", err) + } + if err := v.CommitCompact(); err != nil { + t.Fatalf("commit compact: %v", err) + } + if v.SuperBlock.CompactionRevision == 0 { + t.Fatal("vacuum must bump CompactionRevision for this test to exercise the vacuumed path") + } + + if got, want := v.lastModifiedTsSeconds, lastWriteNs/uint64(time.Second); got != want { + t.Errorf("TTL clock after commit is %d, want the last write at %d", got, want) + } + if !v.expired(v.ContentSize(), 1024*1024) { + t.Error("a TTL volume whose last write is 2h old must stay expired across a vacuum commit") + } +} + +// TestVolumeTtlClockAtCommitUsesAppendTime covers a write whose client +// supplied modified time does not match when it was appended: the commit +// keeps the server-side append watermark the recovery scan would recompute. +func TestVolumeTtlClockAtCommitUsesAppendTime(t *testing.T) { + dir := t.TempDir() + ttl, err := needle.ReadTTL("5m") + if err != nil { + t.Fatalf("read ttl: %v", err) + } + + v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, ttl, 0, needle.GetCurrentVersion(), 0, 0) + if err != nil { + t.Fatalf("volume creation: %v", err) + } + defer v.Close() + + future := uint64(time.Now().Add(24 * time.Hour).Unix()) + n := newRandomNeedle(1) + n.LastModified = future + if _, _, _, err := v.writeNeedle2(n, true, false, false); err != nil { + t.Fatalf("write needle: %v", err) + } + if v.lastModifiedTsSeconds < future { + t.Fatalf("clock %d did not follow the needle's modified time %d", v.lastModifiedTsSeconds, future) + } + appendWatermarkSec := v.lastWriteAppendAtNs / uint64(time.Second) + + if err := v.CompactByIndex(nil); err != nil { + t.Fatalf("compact: %v", err) + } + if err := v.CommitCompact(); err != nil { + t.Fatalf("commit compact: %v", err) + } + + if got := v.lastModifiedTsSeconds; got != appendWatermarkSec { + t.Errorf("TTL clock after commit is %d, want the append watermark %d", got, appendWatermarkSec) + } +} + +// TestVolumeTtlClockAtCommitSkipsDeletedWrite covers a vacuum commit whose +// newest write was deleted first: the carried watermark belonged to that +// write, so the commit has to let the reload rescan and land on the newest +// surviving write rather than keep the volume alive on a deleted write's time. +func TestVolumeTtlClockAtCommitSkipsDeletedWrite(t *testing.T) { + dir := t.TempDir() + ttl, err := needle.ReadTTL("5m") + if err != nil { + t.Fatalf("read ttl: %v", err) + } + + v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, ttl, 0, needle.GetCurrentVersion(), 0, 0) + if err != nil { + t.Fatalf("volume creation: %v", err) + } + defer v.Close() + + oldWriteNs := uint64(time.Now().Add(-2 * time.Hour).UnixNano()) + newWriteNs := uint64(time.Now().Add(-time.Hour).UnixNano()) + for _, w := range []struct { + id uint64 + ns uint64 + }{{1, oldWriteNs}, {2, newWriteNs}} { + n := newRandomNeedle(w.id) + offset, _, _, err := v.writeNeedle2(n, true, false, false) + if err != nil { + t.Fatalf("write needle %d: %v", w.id, err) + } + backdateAppendAtNs(t, v, int64(offset), n.Size, w.ns) + } + // The delete lands inside the commit window: makeupDiff replays its + // tombstone into the new .idx behind the write row the copy carried. + if err := v.CompactByIndex(nil); err != nil { + t.Fatalf("compact: %v", err) + } + if _, err := v.doDeleteRequest(newEmptyNeedle(2)); err != nil { + t.Fatalf("delete needle 2: %v", err) + } + if !v.lastWriteDeleted { + t.Fatal("deleting the newest write must mark the watermark dead") + } + if err := v.CommitCompact(); err != nil { + t.Fatalf("commit compact: %v", err) + } + + if got, want := v.lastModifiedTsSeconds, oldWriteNs/uint64(time.Second); got != want { + t.Errorf("TTL clock after commit is %d, want the surviving write at %d", got, want) + } + if v.lastWriteDeleted { + t.Error("the reload's rescan must reseed the watermark off the surviving write") + } +} + // TestVolumeExpireAtSecCountsFromLastWrite guards the destroy time an EC volume // is reclaimed on (erasure_coding.EcVolume.IsTimeToDestroy). It was recomputed // as now+TTL on every .vif write, so a read-only mark, a tier upload or an EC diff --git a/weed/storage/volume_vacuum.go b/weed/storage/volume_vacuum.go index 287ae2651..fb98caad1 100644 --- a/weed/storage/volume_vacuum.go +++ b/weed/storage/volume_vacuum.go @@ -222,7 +222,17 @@ func (v *Volume) CommitCompact() error { //time.Sleep(20 * time.Second) glog.V(3).Infof("Loading volume %d commit file...", v.Id) - if e := v.load(true, false, v.needleMapKind, 0, v.Version()); e != nil { + // The write watermark already equals what recoverLastModifiedTs would + // rescan, so keep the clock instead of paying for the scan under the lock. + // If its write was itself deleted, the reload recovers the newest + // surviving write instead. + if v.lastWriteAppendAtNs != 0 { + v.lastModifiedTsSeconds = v.lastWriteAppendAtNs / uint64(time.Second) + } + v.keepLastModifiedTsOnLoad = v.lastModifiedTsSeconds != 0 && !v.lastWriteDeleted + e := v.load(true, false, v.needleMapKind, 0, v.Version()) + v.keepLastModifiedTsOnLoad = false + if e != nil { return e } glog.V(3).Infof("Finish committing volume %d", v.Id) diff --git a/weed/storage/volume_write.go b/weed/storage/volume_write.go index 2b03cb21c..0054b4f1a 100644 --- a/weed/storage/volume_write.go +++ b/weed/storage/volume_write.go @@ -329,6 +329,10 @@ func (v *Volume) rollbackUnflushedWrite(n *needle.Needle, offset uint64, end int } } } + if err == nil && n.Id == v.lastWriteNeedleKey { + v.lastWriteAppendAtNs, v.lastWriteNeedleKey = 0, 0 + v.lastWriteDeleted = true + } if err != nil { recoveryErr = errors.Join(recoveryErr, fmt.Errorf("roll back the index of needle %d in volume %d: %w", n.Id, v.Id, err)) @@ -409,6 +413,9 @@ func (v *Volume) doWriteRequest(n *needle.Needle, checkCookie bool, fsync bool) return } v.lastAppendAtNs = n.AppendAtNs + v.lastWriteAppendAtNs = n.AppendAtNs + v.lastWriteNeedleKey = n.Id + v.lastWriteDeleted = false // add to needle map if !ok || uint64(nv.Offset.ToActualOffset()) < offset { @@ -490,6 +497,10 @@ func (v *Volume) doDeleteRequest(n *needle.Needle) (Size, error) { if err = v.nm.Delete(n.Id, ToOffset(int64(offset))); err != nil { return size, err } + if n.Id == v.lastWriteNeedleKey { + v.lastWriteAppendAtNs, v.lastWriteNeedleKey = 0, 0 + v.lastWriteDeleted = true + } return size, err } return 0, nil @@ -524,6 +535,9 @@ func (v *Volume) processBatch(currentRequests []*needle.AsyncRequest) { } indexEnd := int64(v.nm.IndexFileSize()) batchLastAppendAtNs := v.lastAppendAtNs + batchLastWriteAppendAtNs := v.lastWriteAppendAtNs + batchLastWriteNeedleKey := v.lastWriteNeedleKey + batchLastWriteDeleted := v.lastWriteDeleted batchLastModifiedTsSeconds := v.lastModifiedTsSeconds for i := 0; i < len(currentRequests); i++ { needleID := currentRequests[i].N.Id @@ -563,6 +577,9 @@ func (v *Volume) processBatch(currentRequests []*needle.AsyncRequest) { if syncErr := v.DataBackend.Sync(); syncErr != nil { v.checkReadWriteError(syncErr) v.lastAppendAtNs = batchLastAppendAtNs + v.lastWriteAppendAtNs = batchLastWriteAppendAtNs + v.lastWriteNeedleKey = batchLastWriteNeedleKey + v.lastWriteDeleted = batchLastWriteDeleted v.lastModifiedTsSeconds = batchLastModifiedTsSeconds batchErr := syncErr if recoveryErr := v.rollbackBatch(end, indexEnd, orderedSnapshots, metricRollbacker, batchMetrics); recoveryErr != nil { @@ -702,6 +719,9 @@ func (v *Volume) WriteNeedleBlob(needleId NeedleId, needleBlob []byte, size Size return err } v.lastAppendAtNs = appendAtNs + v.lastWriteAppendAtNs = appendAtNs + v.lastWriteNeedleKey = needleId + v.lastWriteDeleted = false // add to needle map if err = v.nm.Put(needleId, ToOffset(int64(offset)), size); err != nil {