diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index e3e4f9514..2f20e58bb 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -1247,7 +1247,7 @@ impl VolumeServer for VolumeGrpcService { .await .map_err(|e| Status::internal(e))?; if dat_modified_ts_ns > 0 { - set_file_mtime(&dat_path, dat_modified_ts_ns); + let _ = set_file_mtime(&dat_path, dat_modified_ts_ns); } } @@ -1272,7 +1272,7 @@ impl VolumeServer for VolumeGrpcService { .await .map_err(|e| Status::internal(e))?; if idx_modified_ts_ns > 0 { - set_file_mtime(&idx_path, idx_modified_ts_ns); + let _ = set_file_mtime(&idx_path, idx_modified_ts_ns); } // Copy .vif file (ignore if not found on source) @@ -1296,7 +1296,7 @@ impl VolumeServer for VolumeGrpcService { .await .map_err(|e| Status::internal(e))?; if vif_modified_ts_ns > 0 { - set_file_mtime(&vif_path, vif_modified_ts_ns); + let _ = set_file_mtime(&vif_path, vif_modified_ts_ns); } // Remove the .note file @@ -3087,6 +3087,15 @@ impl VolumeServer for VolumeGrpcService { dat_path }; + // Store the source .dat mtime, not the upload time, so a reload computes + // TTL from real data age (matches Go VolumeTierMoveDatToRemote). + let dat_modified_secs = std::fs::metadata(&dat_path) + .and_then(|m| m.modified()) + .map_err(|e| Status::internal(format!("stat data file {}: {}", dat_path, e)))? + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_secs()) + .unwrap_or(0); + // Look up the S3 tier backend let backend = { let registry = self.state.s3_tier_registry.read().unwrap(); @@ -3139,18 +3148,13 @@ impl VolumeServer for VolumeGrpcService { { let mut store = state.store.write().unwrap(); if let Some((_, vol)) = store.find_volume_mut(vid) { - let now_unix = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap_or_default() - .as_secs(); - vol.volume_info.files.push(volume_server_pb::RemoteFile { backend_type: backend_type.clone(), backend_id: backend_id.clone(), key, offset: 0, file_size: size, - modified_time: now_unix, + modified_time: dat_modified_secs, extension: ".dat".to_string(), }); vol.refresh_remote_write_mode(); @@ -3201,7 +3205,7 @@ impl VolumeServer for VolumeGrpcService { let vid = VolumeId(req.volume_id); // Validate volume and get remote storage info - let (dat_path, storage_name, storage_key) = { + let (dat_path, storage_name, storage_key, remote_modified_secs) = { let store = self.state.store.read().unwrap(); let (_, vol) = store .find_volume(vid) @@ -3231,7 +3235,14 @@ impl VolumeServer for VolumeGrpcService { ))); } - (dat_path, storage_name, storage_key) + let remote_modified_secs = vol + .volume_info + .files + .first() + .map(|f| f.modified_time) + .unwrap_or(0); + + (dat_path, storage_name, storage_key, remote_modified_secs) }; // Look up the S3 tier backend @@ -3279,6 +3290,15 @@ impl VolumeServer for VolumeGrpcService { )) })?; + // Restore the .dat mtime so a reload computes TTL from real data age, + // not download time (matches Go VolumeTierMoveDatFromRemote). + if remote_modified_secs > 0 { + let modified_ts_ns = (remote_modified_secs as i64).saturating_mul(1_000_000_000); + if let Err(e) = set_file_mtime(&dat_path, modified_ts_ns) { + tracing::warn!("volume {} restore data file {} modified time: {}", vid, dat_path, e); + } + } + if !keep_remote { // Delete remote file backend.delete_file(&storage_key).await.map_err(|e| { @@ -4150,17 +4170,16 @@ async fn ping_filer_target( use super::grpc_client::parse_grpc_address; /// Set the modification time of a file from nanoseconds since Unix epoch. -fn set_file_mtime(path: &str, modified_ts_ns: i64) { +fn set_file_mtime(path: &str, modified_ts_ns: i64) -> std::io::Result<()> { use std::time::{Duration, SystemTime}; let ts = if modified_ts_ns >= 0 { SystemTime::UNIX_EPOCH + Duration::from_nanos(modified_ts_ns as u64) } else { SystemTime::UNIX_EPOCH }; - if let Ok(file) = std::fs::File::open(path) { - let ft = std::fs::FileTimes::new().set_accessed(ts).set_modified(ts); - let _ = file.set_times(ft); - } + let file = std::fs::File::open(path)?; + let ft = std::fs::FileTimes::new().set_accessed(ts).set_modified(ts); + file.set_times(ft) } /// Copy a file from a remote volume server via CopyFile streaming RPC. diff --git a/weed/server/volume_grpc_tier_download.go b/weed/server/volume_grpc_tier_download.go index e7a93e7bd..afde9a44f 100644 --- a/weed/server/volume_grpc_tier_download.go +++ b/weed/server/volume_grpc_tier_download.go @@ -2,8 +2,10 @@ package weed_server import ( "fmt" + "os" "time" + "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb" "github.com/seaweedfs/seaweedfs/weed/storage/backend" "github.com/seaweedfs/seaweedfs/weed/storage/needle" @@ -28,6 +30,7 @@ func (vs *VolumeServer) VolumeTierMoveDatFromRemote(req *volume_server_pb.Volume if storageName == "" || storageKey == "" { return fmt.Errorf("volume %d is already on local disk", req.VolumeId) } + remoteFileModifiedTime := v.GetVolumeInfo().GetFiles()[0].GetModifiedTime() // check whether the local .dat already exists _, ok := v.DataBackend.(*backend.DiskFile) @@ -62,6 +65,14 @@ func (vs *VolumeServer) VolumeTierMoveDatFromRemote(req *volume_server_pb.Volume if err != nil { return fmt.Errorf("backend %s copy file %s: %v", storageName, v.FileName(".dat"), err) } + if remoteFileModifiedTime > 0 { + modifiedTime := time.Unix(int64(remoteFileModifiedTime), 0) + // best-effort: a cosmetic mtime failure should not abort an otherwise + // complete tier-down, matching the EC copy path + if err := os.Chtimes(v.FileName(".dat"), modifiedTime, modifiedTime); err != nil { + glog.Warningf("volume %d restore data file %s modified time: %v", v.Id, v.FileName(".dat"), err) + } + } if req.KeepRemoteDatFile { return nil diff --git a/weed/server/volume_grpc_tier_test.go b/weed/server/volume_grpc_tier_test.go new file mode 100644 index 000000000..1329b0bf4 --- /dev/null +++ b/weed/server/volume_grpc_tier_test.go @@ -0,0 +1,220 @@ +package weed_server + +import ( + "fmt" + "io" + "os" + "path/filepath" + "testing" + "time" + + "google.golang.org/grpc" + + "github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb" + "github.com/seaweedfs/seaweedfs/weed/stats" + "github.com/seaweedfs/seaweedfs/weed/storage" + "github.com/seaweedfs/seaweedfs/weed/storage/backend" + "github.com/seaweedfs/seaweedfs/weed/storage/needle" + "github.com/seaweedfs/seaweedfs/weed/storage/types" + "github.com/seaweedfs/seaweedfs/weed/util" +) + +const tierTimestampTestBackendName = "tier_timestamp_test.default" + +type discardServerStream[T any] struct { + grpc.ServerStream +} + +func (s *discardServerStream[T]) Send(*T) error { + return nil +} + +type tierTimestampTestBackend struct { + root string +} + +func (b *tierTimestampTestBackend) ToProperties() map[string]string { + return map[string]string{"root": b.root} +} + +func (b *tierTimestampTestBackend) NewStorageFile(key string, volumeInfo *volume_server_pb.VolumeInfo) backend.BackendStorageFile { + return &tierTimestampTestBackendFile{ + path: filepath.Join(b.root, key), + volumeInfo: volumeInfo, + } +} + +func (b *tierTimestampTestBackend) CopyFile(file *os.File, fn func(progressed int64, percentage float32) error) (key string, size int64, err error) { + key = "remote.dat" + fileInfo, err := file.Stat() + if err != nil { + return "", 0, err + } + + output, err := os.Create(filepath.Join(b.root, key)) + if err != nil { + return "", 0, err + } + defer output.Close() + + size, err = io.Copy(output, io.NewSectionReader(file, 0, fileInfo.Size())) + if err == nil && fn != nil { + err = fn(size, 100) + } + return key, size, err +} + +func (b *tierTimestampTestBackend) DownloadFile(fileName string, key string, fn func(progressed int64, percentage float32) error) (size int64, err error) { + input, err := os.Open(filepath.Join(b.root, key)) + if err != nil { + return 0, err + } + defer input.Close() + + output, err := os.Create(fileName) + if err != nil { + return 0, err + } + defer output.Close() + + size, err = io.Copy(output, input) + if err == nil && fn != nil { + err = fn(size, 100) + } + return size, err +} + +func (b *tierTimestampTestBackend) DeleteFile(key string) error { + return os.Remove(filepath.Join(b.root, key)) +} + +type tierTimestampTestBackendFile struct { + path string + volumeInfo *volume_server_pb.VolumeInfo +} + +func (f *tierTimestampTestBackendFile) ReadAt(p []byte, off int64) (int, error) { + file, err := os.Open(f.path) + if err != nil { + return 0, err + } + defer file.Close() + return file.ReadAt(p, off) +} + +func (f *tierTimestampTestBackendFile) WriteAt(p []byte, off int64) (int, error) { + return 0, fmt.Errorf("remote test file is read-only") +} + +func (f *tierTimestampTestBackendFile) Truncate(off int64) error { + return fmt.Errorf("remote test file is read-only") +} + +func (f *tierTimestampTestBackendFile) Close() error { + return nil +} + +func (f *tierTimestampTestBackendFile) GetStat() (datSize int64, modTime time.Time, err error) { + files := f.volumeInfo.GetFiles() + if len(files) == 0 { + return 0, time.Time{}, fmt.Errorf("remote file info not found") + } + return int64(files[0].GetFileSize()), time.Unix(int64(files[0].GetModifiedTime()), 0), nil +} + +func (f *tierTimestampTestBackendFile) Name() string { + return f.path +} + +func (f *tierTimestampTestBackendFile) Sync() error { + return nil +} + +func TestVolumeTierMoveDatPreservesModifiedTime(t *testing.T) { + dataDir := t.TempDir() + remoteDir := t.TempDir() + testBackend := &tierTimestampTestBackend{root: remoteDir} + backend.BackendStorages[tierTimestampTestBackendName] = testBackend + t.Cleanup(func() { + delete(backend.BackendStorages, tierTimestampTestBackendName) + }) + + store := storage.NewStore( + nil, + "localhost", + 8080, + 18080, + "http://localhost:8080", + "store-id", + []string{dataDir}, + []int32{10}, + []util.MinFreeSpace{{}}, + "", + storage.NeedleMapInMemory, + []types.DiskType{types.HardDriveType}, + nil, + 0, + stats.DefaultDiskIOProbeConfig(), + ) + t.Cleanup(store.Close) + + const volumeId = needle.VolumeId(1) + if err := store.AddVolume(volumeId, "", storage.NeedleMapInMemory, "000", "", 0, needle.Version3, 0, types.HardDriveType, 0); err != nil { + t.Fatalf("add volume: %v", err) + } + + volume := store.GetVolume(volumeId) + dataFileName := volume.FileName(".dat") + sourceModifiedTime := time.Unix(1_700_000_000, 0) + if err := os.Chtimes(dataFileName, sourceModifiedTime, sourceModifiedTime); err != nil { + t.Fatalf("set source modified time: %v", err) + } + + // Re-open the data backend so the DiskFile caches the on-disk mtime, the way + // a volume freshly loaded from disk does. + volume.DataBackend.Close() + reopened, err := os.OpenFile(dataFileName, os.O_RDWR, 0644) + if err != nil { + t.Fatalf("reopen data file: %v", err) + } + volume.DataBackend = backend.NewDiskFile(reopened) + + volumeServer := &VolumeServer{store: store} + if err := volumeServer.VolumeTierMoveDatToRemote( + &volume_server_pb.VolumeTierMoveDatToRemoteRequest{ + VolumeId: uint32(volumeId), + DestinationBackendName: tierTimestampTestBackendName, + }, + &discardServerStream[volume_server_pb.VolumeTierMoveDatToRemoteResponse]{}, + ); err != nil { + t.Fatalf("move data to remote: %v", err) + } + + remoteFiles := volume.GetVolumeInfo().GetFiles() + if len(remoteFiles) != 1 { + t.Fatalf("remote file count = %d, want 1", len(remoteFiles)) + } + if got := remoteFiles[0].GetModifiedTime(); got != uint64(sourceModifiedTime.Unix()) { + t.Fatalf("remote modified time = %d, want %d", got, sourceModifiedTime.Unix()) + } + if _, err := os.Stat(dataFileName); !os.IsNotExist(err) { + t.Fatalf("local data file still exists after upload: %v", err) + } + + if err := volumeServer.VolumeTierMoveDatFromRemote( + &volume_server_pb.VolumeTierMoveDatFromRemoteRequest{ + VolumeId: uint32(volumeId), + }, + &discardServerStream[volume_server_pb.VolumeTierMoveDatFromRemoteResponse]{}, + ); err != nil { + t.Fatalf("move data from remote: %v", err) + } + + fileInfo, err := os.Stat(dataFileName) + if err != nil { + t.Fatalf("stat downloaded data file: %v", err) + } + if got := fileInfo.ModTime().Unix(); got != sourceModifiedTime.Unix() { + t.Fatalf("downloaded modified time = %d, want %d", got, sourceModifiedTime.Unix()) + } +} diff --git a/weed/server/volume_grpc_tier_upload.go b/weed/server/volume_grpc_tier_upload.go index 4f8ac028b..dee062b58 100644 --- a/weed/server/volume_grpc_tier_upload.go +++ b/weed/server/volume_grpc_tier_upload.go @@ -32,6 +32,10 @@ func (vs *VolumeServer) VolumeTierMoveDatToRemote(req *volume_server_pb.VolumeTi if !ok { return nil // already copied to remove. fmt.Errorf("volume %d is not on local disk", req.VolumeId) } + _, modTime, err := diskFile.GetStat() + if err != nil { + return fmt.Errorf("stat data file %s: %v", diskFile.Name(), err) + } // check valid storage backend type backendStorage, found := backend.BackendStorages[req.DestinationBackendName] @@ -78,7 +82,7 @@ func (vs *VolumeServer) VolumeTierMoveDatToRemote(req *volume_server_pb.VolumeTi Key: key, Offset: 0, FileSize: uint64(size), - ModifiedTime: uint64(time.Now().Unix()), + ModifiedTime: uint64(modTime.Unix()), Extension: ".dat", })