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 {